fix: harden workspace registry refresh and snapshots

This commit is contained in:
2026-08-03 22:33:39 +02:00
parent 2087fbb0c9
commit 553bb41138
4 changed files with 312 additions and 42 deletions
+89 -14
View File
@@ -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<void> {
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<void> {
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<string> {
@@ -207,18 +228,21 @@ export class WorkspaceRepositoryLock {
this.queue = new Promise<void>((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");
}
}
+108 -26
View File
@@ -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<string, string>;
}
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<string, string> = {};
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<string>();
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<void> {
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;
}
}
}
+82 -1
View File
@@ -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" });
});
+33 -1
View File
@@ -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);
});