From 553bb41138c2c12ef07042c09177efeedf320f4c Mon Sep 17 00:00:00 2001 From: mptyl Date: Mon, 3 Aug 2026 22:33:39 +0200 Subject: [PATCH] fix: harden workspace registry refresh and snapshots --- backend/src/workspaces/git-repository.ts | 103 ++++++++++++-- backend/src/workspaces/registry.ts | 134 ++++++++++++++---- backend/test/workspace-registry.test.ts | 83 ++++++++++- .../test/workspaces-git-repository.test.ts | 34 ++++- 4 files changed, 312 insertions(+), 42 deletions(-) diff --git a/backend/src/workspaces/git-repository.ts b/backend/src/workspaces/git-repository.ts index d2fe8e98..fef70512 100644 --- a/backend/src/workspaces/git-repository.ts +++ b/backend/src/workspaces/git-repository.ts @@ -1,6 +1,8 @@ import { execFile } from "node:child_process"; -import { constants, lstatSync, mkdirSync, openSync, closeSync, unlinkSync } from "node:fs"; -import { access, lstat, mkdir } from "node:fs/promises"; +import { + closeSync, constants, lstatSync, mkdirSync, openSync, readFileSync, unlinkSync, writeFileSync, +} from "node:fs"; +import { mkdir } from "node:fs/promises"; import { basename, isAbsolute, join } from "node:path"; import { promisify } from "node:util"; import type { WorkspaceErrorCode, WorkspaceRegistryConfig } from "./types.js"; @@ -83,9 +85,14 @@ export class GitWorkspaceRepository { } async ensureLayout(): Promise { - for (const path of [this.root, this.snapshotsPath, this.statePath, this.locksPath, this.hooksPath]) { - await mkdir(path, { recursive: true, mode: 0o700 }); - assertDirectory(path); + try { + for (const path of [this.root, this.snapshotsPath, this.statePath, this.locksPath, this.hooksPath]) { + await mkdir(path, { recursive: true, mode: 0o700 }); + assertDirectory(path); + } + } catch (error) { + if (error instanceof WorkspaceRegistryError) throw error; + throw new WorkspaceRegistryError("git_unavailable", "Workspace registry storage is unavailable"); } } @@ -162,11 +169,25 @@ export class GitWorkspaceRepository { } private async refresh(): Promise { + if ((await this.git(["status", "--porcelain"])).trim() !== "") { + throw new WorkspaceRegistryError("workspace_stale", "Workspace checkout has local changes"); + } if (this.config.remoteUrl) { await this.git(["remote", "set-url", "origin", "--", this.config.remoteUrl]); } await this.git(["fetch", "--no-tags", "origin", this.config.branch]); - await this.git(["merge", "--ff-only", "FETCH_HEAD"]); + const remoteHead = (await this.git(["rev-parse", "FETCH_HEAD"])).trim(); + const localHead = (await this.git(["rev-parse", "HEAD"])).trim(); + if (localHead !== remoteHead) { + const commonAncestor = (await this.git(["merge-base", "HEAD", "FETCH_HEAD"])).trim(); + if (commonAncestor !== localHead) { + throw new WorkspaceRegistryError("git_non_fast_forward", "Workspace checkout diverged from remote"); + } + await this.git(["merge", "--ff-only", "FETCH_HEAD"]); + } + if ((await this.git(["rev-parse", "HEAD"])).trim() !== remoteHead) { + throw new WorkspaceRegistryError("git_non_fast_forward", "Workspace checkout does not match remote"); + } } private async git(args: string[]): Promise { @@ -207,18 +228,21 @@ export class WorkspaceRepositoryLock { this.queue = new Promise((resolve) => { releaseQueue = resolve; }); await previous; - mkdirSync(this.locksPath, { recursive: true, mode: 0o700 }); - assertDirectory(this.locksPath); let descriptor: number | undefined; const lockPath = join(this.locksPath, "repository.lock"); try { - descriptor = openSync(lockPath, constants.O_CREAT | constants.O_EXCL | constants.O_WRONLY, 0o600); - return await operation(); + mkdirSync(this.locksPath, { recursive: true, mode: 0o700 }); + assertDirectory(this.locksPath); } catch (error) { - if (typeof error === "object" && error !== null && "code" in error && error.code === "EEXIST") { - throw new WorkspaceRegistryError("workspace_stale", "Workspace registry is busy"); - } - throw error; + throw new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable"); + } + try { + descriptor = this.acquire(lockPath); + } catch (error) { + throw this.lockError(error); + } + try { + return await operation(); } finally { if (descriptor !== undefined) closeSync(descriptor); if (descriptor !== undefined) { @@ -227,4 +251,55 @@ export class WorkspaceRepositoryLock { releaseQueue(); } } + + private acquire(lockPath: string): number { + try { + return this.createProcessLock(lockPath); + } catch (error) { + if (!this.recoverDeadProcessLock(lockPath, error)) throw error; + return this.createProcessLock(lockPath); + } + } + + private createProcessLock(lockPath: string): number { + const descriptor = openSync(lockPath, constants.O_CREAT | constants.O_EXCL | constants.O_WRONLY, 0o600); + try { + writeFileSync(descriptor, JSON.stringify({ pid: process.pid }), "utf8"); + return descriptor; + } catch (error) { + closeSync(descriptor); + try { unlinkSync(lockPath); } catch { /* incomplete lock is never treated as recoverable */ } + throw error; + } + } + + private recoverDeadProcessLock(lockPath: string, error: unknown): boolean { + if (!(typeof error === "object" && error !== null && "code" in error && error.code === "EEXIST")) { + return false; + } + try { + const record = JSON.parse(readFileSync(lockPath, "utf8")) as { pid?: unknown }; + const pid = record.pid; + if (typeof pid !== "number" || !Number.isSafeInteger(pid) || pid <= 0 || pid === process.pid) return false; + try { + process.kill(pid, 0); + return false; + } catch (probeError) { + if (!(typeof probeError === "object" && probeError !== null && "code" in probeError && probeError.code === "ESRCH")) { + return false; + } + } + unlinkSync(lockPath); + return true; + } catch { + return false; + } + } + + private lockError(error: unknown): WorkspaceRegistryError { + if (typeof error === "object" && error !== null && "code" in error && error.code === "EEXIST") { + return new WorkspaceRegistryError("workspace_stale", "Workspace registry is busy"); + } + return new WorkspaceRegistryError("git_unavailable", "Workspace registry lock is unavailable"); + } } diff --git a/backend/src/workspaces/registry.ts b/backend/src/workspaces/registry.ts index cd6efe3f..727b2ea0 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -1,4 +1,4 @@ -import { randomUUID } from "node:crypto"; +import { createHash, randomUUID } from "node:crypto"; import { lstatSync } from "node:fs"; import { mkdir, readFile, rename, rm, writeFile } from "node:fs/promises"; import { isAbsolute, join } from "node:path"; @@ -31,6 +31,10 @@ interface ActiveState { revisions: WorkspaceRevision[]; } +interface SnapshotManifest extends ActiveState { + files: Record; +} + function workspacePath(id: string): string { if (!/^[a-z][a-z0-9-]{2,62}$/.test(id)) { throw new WorkspaceRegistryError("workspace_invalid", "Workspace ID is invalid"); @@ -45,6 +49,17 @@ function safeCommit(commit: string): string { return commit; } +function safeBlob(blob: string): string { + if (!/^[0-9a-f]{40}$/.test(blob)) { + throw new WorkspaceRegistryError("workspace_invalid", "Workspace snapshot blob is invalid"); + } + return blob; +} + +function digest(contents: string | Buffer): string { + return createHash("sha256").update(contents).digest("hex"); +} + function workspaceError(error: unknown): WorkspaceRegistryError { if (error instanceof WorkspaceRegistryError) return error; return new WorkspaceRegistryError("workspace_invalid", "Workspace repository content is invalid"); @@ -136,20 +151,32 @@ export class WorkspaceRegistry { } const snapshotDirectory = join(this.repository.snapshotsPath, safeHead); - if (!this.pathExists(snapshotDirectory)) { + const revisions = snapshots.map((snapshot) => ({ + id: snapshot.id, + commit: safeHead, + blob: snapshot.blob, + snapshotPath: this.snapshotPath(safeHead, snapshot.id), + })); + if (this.pathExists(snapshotDirectory)) { + await this.assertSnapshotIntegrity({ head: safeHead, revisions }); + } else { const staging = join(this.repository.snapshotsPath, `.staging-${randomUUID()}`); await mkdir(staging, { mode: 0o700 }); try { - const revisions: WorkspaceRevision[] = []; + const files: Record = {}; for (const snapshot of snapshots) { - const path = join(staging, `${snapshot.id}.yaml`); + const yamlName = `${snapshot.id}.yaml`; + const envName = `${snapshot.id}.env.example`; + const docsName = `${snapshot.id}.md`; const docs = renderWorkspaceDocs(snapshot.workspace); - await writeFile(path, snapshot.source, { encoding: "utf8", mode: 0o400 }); - await writeFile(join(staging, `${snapshot.id}.env.example`), docs.envExample, { encoding: "utf8", mode: 0o400 }); - await writeFile(join(staging, `${snapshot.id}.md`), docs.markdown, { encoding: "utf8", mode: 0o400 }); - revisions.push({ id: snapshot.id, commit: safeHead, blob: snapshot.blob, snapshotPath: this.snapshotPath(safeHead, snapshot.id) }); + await writeFile(join(staging, yamlName), snapshot.source, { encoding: "utf8", mode: 0o400 }); + await writeFile(join(staging, envName), docs.envExample, { encoding: "utf8", mode: 0o400 }); + await writeFile(join(staging, docsName), docs.markdown, { encoding: "utf8", mode: 0o400 }); + files[yamlName] = digest(snapshot.source); + files[envName] = digest(docs.envExample); + files[docsName] = digest(docs.markdown); } - await writeFile(join(staging, "snapshot.json"), JSON.stringify({ head: safeHead, revisions }), { + await writeFile(join(staging, "snapshot.json"), JSON.stringify({ head: safeHead, revisions, files }), { encoding: "utf8", mode: 0o400, }); await rename(staging, snapshotDirectory); @@ -159,12 +186,6 @@ export class WorkspaceRegistry { } } - const revisions = snapshots.map((snapshot) => ({ - id: snapshot.id, - commit: safeHead, - blob: snapshot.blob, - snapshotPath: this.snapshotPath(safeHead, snapshot.id), - })); await this.writeActiveState({ head: safeHead, revisions }); } @@ -193,18 +214,13 @@ export class WorkspaceRegistry { const file = join(this.repository.statePath, "active.json"); try { const state = JSON.parse(await readFile(file, "utf8")) as ActiveState; - safeCommit(state.head); - if (!Array.isArray(state.revisions) || state.revisions.length === 0) throw new Error("bad state"); - for (const revision of state.revisions) { - safeCommit(revision.commit); - workspacePath(revision.id); - if (!isAbsolute(revision.snapshotPath) || revision.snapshotPath !== this.snapshotPath(revision.commit, revision.id)) { - throw new Error("bad snapshot path"); - } - } + this.assertActiveState(state); + await this.assertSnapshotIntegrity(state); return state; - } catch { - return undefined; + } catch (error) { + if (this.pathIsMissing(file)) return undefined; + if (error instanceof WorkspaceRegistryError) throw error; + throw new WorkspaceRegistryError("workspace_invalid", "Workspace active snapshot is invalid"); } } @@ -215,6 +231,63 @@ export class WorkspaceRegistry { await rename(staging, target); } + private assertActiveState(state: ActiveState): void { + safeCommit(state.head); + if (!Array.isArray(state.revisions) || state.revisions.length === 0) throw new Error("bad state"); + const ids = new Set(); + for (const revision of state.revisions) { + safeCommit(revision.commit); + safeBlob(revision.blob); + if (revision.commit !== state.head || ids.has(revision.id)) throw new Error("bad revision"); + ids.add(revision.id); + workspacePath(revision.id); + if (!isAbsolute(revision.snapshotPath) || revision.snapshotPath !== this.snapshotPath(revision.commit, revision.id)) { + throw new Error("bad snapshot path"); + } + } + } + + private async assertSnapshotIntegrity(state: ActiveState): Promise { + const directory = join(this.repository.snapshotsPath, state.head); + const manifestPath = join(directory, "snapshot.json"); + try { + const manifest = JSON.parse(await readFile(manifestPath, "utf8")) as SnapshotManifest; + this.assertActiveState(manifest); + if (manifest.head !== state.head || !this.sameRevisions(manifest.revisions, state.revisions)) { + throw new Error("manifest revisions do not match active state"); + } + const expected = state.revisions.flatMap((revision) => [ + `${revision.id}.yaml`, `${revision.id}.env.example`, `${revision.id}.md`, + ]); + if (Object.keys(manifest.files).length !== expected.length || !expected.every((name) => ( + /^[0-9a-f]{64}$/.test(manifest.files[name] ?? "") + ))) throw new Error("manifest files are invalid"); + for (const name of expected) { + const path = join(directory, name); + const entry = lstatSync(path); + if (!entry.isFile() || entry.isSymbolicLink()) throw new Error("snapshot file is invalid"); + const contents = await readFile(path); + if (digest(contents) !== manifest.files[name]) throw new Error("snapshot file does not match manifest"); + if (name.endsWith(".yaml")) { + const workspace = parseWorkspaceYaml(contents.toString("utf8")); + if (workspace.workspace.id !== name.slice(0, -".yaml".length)) throw new Error("snapshot workspace is invalid"); + } + } + } catch (error) { + if (error instanceof WorkspaceRegistryError) throw error; + throw new WorkspaceRegistryError("workspace_invalid", "Workspace snapshot integrity check failed"); + } + } + + private sameRevisions(left: WorkspaceRevision[], right: WorkspaceRevision[]): boolean { + return left.length === right.length && left.every((revision, index) => { + const candidate = right[index]; + return candidate !== undefined + && candidate.id === revision.id && candidate.commit === revision.commit + && candidate.blob === revision.blob && candidate.snapshotPath === revision.snapshotPath; + }); + } + private pathExists(path: string): boolean { try { const entry = lstatSync(path); @@ -227,4 +300,13 @@ export class WorkspaceRegistry { return false; } } + + private pathIsMissing(path: string): boolean { + try { + lstatSync(path); + return false; + } catch { + return true; + } + } } diff --git a/backend/test/workspace-registry.test.ts b/backend/test/workspace-registry.test.ts index 46e31c29..86a0dd4e 100644 --- a/backend/test/workspace-registry.test.ts +++ b/backend/test/workspace-registry.test.ts @@ -1,5 +1,7 @@ import { execFile } from "node:child_process"; -import { existsSync, mkdtempSync, mkdirSync, rmSync, symlinkSync, writeFileSync } from "node:fs"; +import { + chmodSync, existsSync, mkdtempSync, mkdirSync, readFileSync, rmSync, symlinkSync, writeFileSync, +} from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { promisify } from "node:util"; @@ -132,3 +134,82 @@ test("rejects a symbolic-link registry root before creating a lock below it", as await expect(registry.bootstrap()).rejects.toMatchObject({ code: "git_unavailable" }); expect(existsSync(join(target, "locks"))).toBe(false); }); + +test("rejects a locally-ahead checkout instead of activating local-only content", async () => { + const remote = await fixture(); + const root = join(remote.root, "registry"); + const registry = new WorkspaceRegistry(config(root, remote.remote)); + await registry.bootstrap(); + const checkout = join(root, "repo"); + writeFileSync(join(checkout, "workspaces", "psd-clinical.yaml"), validYaml.replace( + "name: Policlinico San Donato", "name: Local only workspace", + )); + await git(checkout, ["config", "user.name", "Workspace Registry Test"]); + await git(checkout, ["config", "user.email", "workspace-registry@example.invalid"]); + await git(checkout, ["add", "workspaces/psd-clinical.yaml"]); + await git(checkout, ["commit", "-m", "Local-only workspace"]); + + await expect(registry.pull()).rejects.toMatchObject({ code: "git_non_fast_forward" }); + await expect(registry.read("psd-clinical")).resolves.toMatchObject({ + revision: { commit: remote.initialCommit }, + workspace: { workspace: { name: "Policlinico San Donato" } }, + }); +}); + +test("recovers a dead-process advisory lock while preserving active snapshot safety", async () => { + const remote = await fixture(); + const root = join(remote.root, "registry"); + mkdirSync(join(root, "locks"), { recursive: true }); + writeFileSync(join(root, "locks", "repository.lock"), JSON.stringify({ pid: 999_999_999 })); + const registry = new WorkspaceRegistry(config(root, remote.remote)); + + await expect(registry.bootstrap()).resolves.toMatchObject({ + head: remote.initialCommit, + degraded: false, + }); +}); + +test.each(["manifest", "blob", "workspace", "document"])( + "rejects a corrupted %s snapshot component instead of reporting it active", + async (component) => { + const remote = await fixture(); + const root = join(remote.root, "registry"); + const registry = new WorkspaceRegistry(config(root, remote.remote)); + await registry.bootstrap(); + const snapshot = join(root, "snapshots", remote.initialCommit); + + if (component === "manifest") { + const file = join(snapshot, "snapshot.json"); + chmodSync(file, 0o600); + writeFileSync(file, "{"); + } + if (component === "blob") { + const activePath = join(root, "state", "active.json"); + const active = JSON.parse(readFileSync(activePath, "utf8")); + active.revisions[0].blob = "not-a-git-blob"; + writeFileSync(activePath, JSON.stringify(active)); + } + if (component === "workspace") { + const file = join(snapshot, "psd-clinical.yaml"); + chmodSync(file, 0o600); + writeFileSync(file, "truncated"); + } + if (component === "document") rmSync(join(snapshot, "psd-clinical.md")); + + await expect(registry.list()).rejects.toMatchObject({ code: "workspace_invalid" }); + await expect(registry.read("psd-clinical")).rejects.toMatchObject({ code: "workspace_invalid" }); + }, +); + +test("rejects a corrupt fallback snapshot instead of returning degraded active state", async () => { + const remote = await fixture(); + const root = join(remote.root, "registry"); + const registry = new WorkspaceRegistry(config(root, remote.remote)); + await registry.bootstrap(); + const document = join(root, "snapshots", remote.initialCommit, "psd-clinical.md"); + chmodSync(document, 0o600); + writeFileSync(document, "corrupt"); + rmSync(remote.remote, { recursive: true, force: true }); + + await expect(registry.pull()).rejects.toMatchObject({ code: "workspace_invalid" }); +}); diff --git a/backend/test/workspaces-git-repository.test.ts b/backend/test/workspaces-git-repository.test.ts index ef0460d6..cdf95f5a 100644 --- a/backend/test/workspaces-git-repository.test.ts +++ b/backend/test/workspaces-git-repository.test.ts @@ -4,7 +4,7 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import { promisify } from "node:util"; import { afterEach, expect, test } from "vitest"; -import { GitWorkspaceRepository } from "../src/workspaces/git-repository.js"; +import { GitWorkspaceRepository, WorkspaceRepositoryLock } from "../src/workspaces/git-repository.js"; import type { WorkspaceRegistryConfig } from "../src/workspaces/types.js"; const validYaml = `workspace: @@ -105,3 +105,35 @@ test("redacts failed Git checkout details behind a stable error code", async () }); expect((error as Error).message).not.toContain(remote); }); + +test("maps registry-layout failures to a stable redacted error", async () => { + const root = mkdtempSync(join(tmpdir(), "thoth-workspace-git-layout-")); + temporaryRoots.push(root); + const file = join(root, "not-a-directory"); + writeFileSync(file, "occupied"); + const repository = new GitWorkspaceRepository(config(file, join(root, "remote.git"))); + + const error = await repository.bootstrap().catch((error: unknown) => error); + + expect(error).toMatchObject({ + code: "git_unavailable", + message: "Workspace registry storage is unavailable", + }); + expect((error as Error).message).not.toContain(file); +}); + +test("maps lock filesystem failures to a stable redacted error", async () => { + const root = mkdtempSync(join(tmpdir(), "thoth-workspace-git-lock-")); + temporaryRoots.push(root); + const file = join(root, "not-a-directory"); + writeFileSync(file, "occupied"); + const lock = new WorkspaceRepositoryLock(file); + + const error = await lock.run(async () => undefined).catch((error: unknown) => error); + + expect(error).toMatchObject({ + code: "git_unavailable", + message: "Workspace registry lock is unavailable", + }); + expect((error as Error).message).not.toContain(file); +});