diff --git a/backend/src/config.ts b/backend/src/config.ts index ee7193d7..87160468 100644 --- a/backend/src/config.ts +++ b/backend/src/config.ts @@ -251,6 +251,7 @@ export function loadConfig(env: Record): AppConfig { secretRoots, maxImportBytes: positiveImportLimit(env.THT_WORKSPACE_MAX_IMPORT_BYTES, 10 * 1024 * 1024), maxImportEntries: positiveImportLimit(env.THT_WORKSPACE_MAX_IMPORT_ENTRIES, 32), + dataRoot: env.THT_DATA_ROOT, }; const settingsFile = env.SETTINGS_FILE ?? "data/settings.json"; const internalQdrantUrl = internalServiceUrl( diff --git a/backend/src/workspaces/annotations-sync.ts b/backend/src/workspaces/annotations-sync.ts new file mode 100644 index 00000000..1592308e --- /dev/null +++ b/backend/src/workspaces/annotations-sync.ts @@ -0,0 +1,174 @@ +import { createHash, randomBytes } from "node:crypto"; +import { + closeSync, + constants as fsConstants, + fchmodSync, + fstatSync, + fsyncSync, + lstatSync, + mkdirSync, + openSync, + readFileSync, + renameSync, + unlinkSync, + writeFileSync, +} from "node:fs"; +import { dirname, isAbsolute, join } from "node:path"; +import { parseAnnotationsYaml } from "./annotations.js"; + +export interface AnnotationsSyncInput { + dataRoot: string; + workspaceId: string; + commit: string; + blobId: string; + contents: Buffer; +} + +export interface AnnotationsSyncResult { + path: string; + manifestPath: string; + contentDigest: string; +} + +interface AnnotationsOwnershipManifest { + workspace: string; + commit: string; + blobId: string; + contentDigest: string; + destination: string; +} + +export function annotationsSyncRoot(dataRoot: string, workspaceId: string, commit: string): string { + return join(dataRoot, "sessions", workspaceId, "revisions", commit, "artifacts"); +} + +function sha256(value: Buffer | string): string { + return `sha256:${createHash("sha256").update(value).digest("hex")}`; +} + +function ensureDirectory(path: string): void { + mkdirSync(path, { recursive: true, mode: 0o700 }); + const entry = lstatSync(path); + if (!entry.isDirectory() || entry.isSymbolicLink()) { + throw new Error("annotations sync directory is unavailable"); + } +} + +function syncDirectory(directory: string): void { + if (process.platform === "win32") return; + const fd = openSync(directory, "r"); + try { fsyncSync(fd); } finally { closeSync(fd); } +} + +function writeAtomicFile(path: string, contents: string | Buffer, mode: number): void { + ensureDirectory(dirname(path)); + const staging = `${path}.tmp-${process.pid}-${Date.now()}-${randomBytes(6).toString("hex")}`; + const fd = openSync( + staging, + fsConstants.O_WRONLY | fsConstants.O_CREAT | fsConstants.O_EXCL | fsConstants.O_NOFOLLOW, + 0o600, + ); + let closed = false; + try { + writeFileSync(fd, contents); + fsyncSync(fd); + fchmodSync(fd, mode); + closeSync(fd); + closed = true; + renameSync(staging, path); + syncDirectory(dirname(path)); + } catch (error) { + if (!closed) try { closeSync(fd); } catch { /* preserve original failure */ } + try { unlinkSync(staging); } catch { /* best effort */ } + throw error; + } +} + +function readTrustedFile(path: string): Buffer { + const entry = lstatSync(path); + if (!entry.isFile() || entry.isSymbolicLink()) throw new Error("annotations sync file is invalid"); + const fd = openSync(path, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW); + try { + const before = fstatSync(fd); + if (!before.isFile() || before.nlink !== 1) throw new Error("annotations sync file is invalid"); + const contents = readFileSync(fd); + const after = fstatSync(fd); + if (before.dev !== after.dev || before.ino !== after.ino || before.size !== after.size || before.nlink !== after.nlink) { + throw new Error("annotations sync file changed while reading"); + } + return contents; + } finally { + closeSync(fd); + } +} + +function parseManifest(source: string): AnnotationsOwnershipManifest { + const parsed = JSON.parse(source) as Record; + if ( + typeof parsed.workspace !== "string" + || typeof parsed.commit !== "string" + || typeof parsed.blobId !== "string" + || typeof parsed.contentDigest !== "string" + || typeof parsed.destination !== "string" + ) { + throw new Error("annotations ownership manifest is invalid"); + } + return parsed as unknown as AnnotationsOwnershipManifest; +} + +/** + * Atomically synchronize a curated annotation blob to its immutable revision-qualified runtime + * root and write an adjacent ownership manifest. Idempotent: an existing destination is re-verified + * against the exact blob/digest and fails closed on any mismatch. Never follows symlinks. + */ +export function syncAnnotations(input: AnnotationsSyncInput): AnnotationsSyncResult { + if (!isAbsolute(input.dataRoot)) throw new Error("annotations sync data root must be absolute"); + if (!/^[a-z][a-z0-9-]{2,62}$/.test(input.workspaceId) || input.workspaceId === "workspace-docs") { + throw new Error("annotations sync workspace id is invalid"); + } + if (!/^[0-9a-f]{40}$/.test(input.commit)) throw new Error("annotations sync commit is invalid"); + if (!/^[0-9a-f]{40}$/.test(input.blobId)) throw new Error("annotations sync blob id is invalid"); + parseAnnotationsYaml(input.contents.toString("utf8")); + + const root = annotationsSyncRoot(input.dataRoot, input.workspaceId, input.commit); + const directory = join(root, "mschema"); + const path = join(directory, "annotations.yaml"); + const manifestPath = join(directory, "annotations.ownership.json"); + const contentDigest = sha256(input.contents); + const manifest: AnnotationsOwnershipManifest = { + workspace: input.workspaceId, + commit: input.commit, + blobId: input.blobId, + contentDigest, + destination: path, + }; + + let existingContents: Buffer | undefined; + try { + existingContents = readTrustedFile(path); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; + existingContents = undefined; + } + + if (existingContents !== undefined) { + if (sha256(existingContents) !== contentDigest) { + throw new Error("annotations sync destination does not match the pinned revision"); + } + const existingManifest = parseManifest(readTrustedFile(manifestPath).toString("utf8")); + if ( + existingManifest.workspace !== manifest.workspace + || existingManifest.commit !== manifest.commit + || existingManifest.blobId !== manifest.blobId + || existingManifest.contentDigest !== manifest.contentDigest + || existingManifest.destination !== manifest.destination + ) { + throw new Error("annotations ownership manifest does not match the pinned revision"); + } + return { path, manifestPath, contentDigest }; + } + + writeAtomicFile(path, input.contents, 0o400); + writeAtomicFile(manifestPath, `${JSON.stringify(manifest)}\n`, 0o600); + return { path, manifestPath, contentDigest }; +} diff --git a/backend/src/workspaces/registry.ts b/backend/src/workspaces/registry.ts index 2f05e04a..82e71e64 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -4,6 +4,7 @@ import { mkdir, readdir, readFile, rename, rm, writeFile } from "node:fs/promise import { isAbsolute, join } from "node:path"; import { buildInstallationContract, renderWorkspaceDocs } from "./contracts.js"; import { parseAnnotationsYaml } from "./annotations.js"; +import { syncAnnotations } from "./annotations-sync.js"; import { assertCatalogMatchesDescriptor, parseWorkspaceCatalogYaml, type WorkspaceCatalog, type WorkspaceCatalogEntry } from "./catalog.js"; import { GitWorkspaceRepository, @@ -585,6 +586,15 @@ export class WorkspaceRegistry { const annotations = await this.repository.annotationsObject(safeHead, id); if (annotations !== undefined) { parseAnnotationsYaml(annotations.contents.toString("utf8")); + if (this.config.dataRoot !== undefined) { + syncAnnotations({ + dataRoot: this.config.dataRoot, + workspaceId: id, + commit: safeHead, + blobId: annotations.blobId, + contents: annotations.contents, + }); + } } const collection = workspace.semantic_index.vector_store.collection; const owner = collectionOwners.get(collection); diff --git a/backend/src/workspaces/types.ts b/backend/src/workspaces/types.ts index e10909bb..b703ae8b 100644 --- a/backend/src/workspaces/types.ts +++ b/backend/src/workspaces/types.ts @@ -8,6 +8,8 @@ export interface WorkspaceRegistryConfig { secretRoots: readonly string[]; maxImportBytes: number; maxImportEntries: number; + /** Absolute runtime data root; when set, activation also syncs curated annotations per revision. */ + dataRoot?: string; } export type WorkspaceErrorCode = diff --git a/backend/test/annotations-sync.test.ts b/backend/test/annotations-sync.test.ts new file mode 100644 index 00000000..e93a4652 --- /dev/null +++ b/backend/test/annotations-sync.test.ts @@ -0,0 +1,116 @@ +import { chmodSync, lstatSync, mkdirSync, mkdtempSync, readFileSync, rmSync, symlinkSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { afterEach, expect, test } from "vitest"; +import { syncAnnotations } from "../src/workspaces/annotations-sync.js"; + +const temporaryRoots: string[] = []; + +afterEach(() => { + temporaryRoots.splice(0).forEach((root) => rmSync(root, { recursive: true, force: true })); +}); + +function input(overrides: Partial<{ + dataRoot: string; workspaceId: string; commit: string; blobId: string; contents: Buffer; +}> = {}) { + return { + dataRoot: overrides.dataRoot ?? "", + workspaceId: overrides.workspaceId ?? "research", + commit: overrides.commit ?? "a".repeat(40), + blobId: overrides.blobId ?? "b".repeat(40), + contents: overrides.contents ?? Buffer.from("tables: {}\n"), + }; +} + +test("writes the revision-qualified annotations file and ownership manifest", () => { + const dataRoot = mkdtempSync(join(tmpdir(), "thoth-annotations-sync-")); + temporaryRoots.push(dataRoot); + const target = input({ dataRoot }); + + const result = syncAnnotations(target); + + expect(readFileSync(result.path, "utf8")).toBe("tables: {}\n"); + const manifest = JSON.parse(readFileSync(result.manifestPath, "utf8")); + expect(manifest).toMatchObject({ + workspace: "research", + commit: "a".repeat(40), + blobId: "b".repeat(40), + contentDigest: result.contentDigest, + destination: result.path, + }); + expect(result.path).toContain("/revisions/" + "a".repeat(40) + "/artifacts/mschema/annotations.yaml"); + expect(lstatSync(result.path).mode & 0o777).toBe(0o400); +}); + +test("is idempotent for the exact same blob and manifest", () => { + const dataRoot = mkdtempSync(join(tmpdir(), "thoth-annotations-sync-")); + temporaryRoots.push(dataRoot); + const target = input({ dataRoot }); + + expect(syncAnnotations(target)).toEqual(syncAnnotations(target)); +}); + +test("fails closed when the destination was tampered", () => { + const dataRoot = mkdtempSync(join(tmpdir(), "thoth-annotations-sync-")); + temporaryRoots.push(dataRoot); + const target = input({ dataRoot }); + syncAnnotations(target); + + const dest = join(dataRoot, "sessions", "research", "revisions", "a".repeat(40), "artifacts", "mschema", "annotations.yaml"); + chmodSync(dest, 0o600); + writeFileSync(dest, "tampered\n"); + + expect(() => syncAnnotations(target)).toThrow(); +}); + +test("fails closed when the ownership manifest does not match", () => { + const dataRoot = mkdtempSync(join(tmpdir(), "thoth-annotations-sync-")); + temporaryRoots.push(dataRoot); + const target = input({ dataRoot }); + syncAnnotations(target); + + const manifestPath = join(dataRoot, "sessions", "research", "revisions", "a".repeat(40), "artifacts", "mschema", "annotations.ownership.json"); + writeFileSync(manifestPath, JSON.stringify({ workspace: "other", commit: "c".repeat(40), blobId: "d".repeat(40), contentDigest: "sha256:x", destination: "/other" })); + + expect(() => syncAnnotations(target)).toThrow(); +}); + +test("writes different revisions to different directories", () => { + const dataRoot = mkdtempSync(join(tmpdir(), "thoth-annotations-sync-")); + temporaryRoots.push(dataRoot); + + const first = syncAnnotations(input({ dataRoot, commit: "a".repeat(40) })); + const second = syncAnnotations(input({ dataRoot, commit: "b".repeat(40) })); + + expect(first.path).not.toBe(second.path); + expect(readFileSync(second.path, "utf8")).toBe("tables: {}\n"); +}); + +test("rejects a symlink destination and malformed inputs before writing", () => { + const dataRoot = mkdtempSync(join(tmpdir(), "thoth-annotations-sync-")); + temporaryRoots.push(dataRoot); + const root = join(dataRoot, "sessions", "research", "revisions", "a".repeat(40), "artifacts"); + mkdirSync(root, { recursive: true }); + symlinkSync(join(root, "mschema-target"), join(root, "mschema")); + mkdirSync(join(root, "mschema-target")); + + expect(() => syncAnnotations(input({ dataRoot }))).toThrow(); + expect(() => syncAnnotations(input({ dataRoot: "relative" }))).toThrow(); + expect(() => syncAnnotations(input({ workspaceId: "../x" }))).toThrow(); + expect(() => syncAnnotations(input({ commit: "HEAD" }))).toThrow(); + expect(() => syncAnnotations(input({ blobId: "not-hex" }))).toThrow(); + expect(() => syncAnnotations(input({ contents: Buffer.from("tables: [bad]\n") }))).toThrow(); +}); + +test("does not follow a symlinked destination when verifying", () => { + const dataRoot = mkdtempSync(join(tmpdir(), "thoth-annotations-sync-")); + temporaryRoots.push(dataRoot); + const target = input({ dataRoot }); + syncAnnotations(target); + + const dest = join(dataRoot, "sessions", "research", "revisions", "a".repeat(40), "artifacts", "mschema", "annotations.yaml"); + rmSync(dest, { force: true }); + symlinkSync("/etc/hosts", dest); + + expect(() => syncAnnotations(target)).toThrow(); +}); diff --git a/backend/test/registry-annotations.test.ts b/backend/test/registry-annotations.test.ts index 1f612b64..094d3152 100644 --- a/backend/test/registry-annotations.test.ts +++ b/backend/test/registry-annotations.test.ts @@ -1,5 +1,5 @@ import { execFile } from "node:child_process"; -import { mkdtempSync, mkdirSync, rmSync, writeFileSync } from "node:fs"; +import { mkdtempSync, mkdirSync, readFileSync, rmSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { promisify } from "node:util"; @@ -59,7 +59,7 @@ function config(root: string, remoteUrl: string): WorkspaceRegistryConfig { type AnnotationsLayout = "absent" | "valid" | "malformed" | "dir"; -async function fixture(layout: AnnotationsLayout): Promise<{ root: string; remote: string }> { +async function fixture(layout: AnnotationsLayout): Promise<{ root: string; remote: string; commit: string }> { const root = mkdtempSync(join(tmpdir(), "thoth-registry-annotations-")); temporaryRoots.push(root); const remote = join(root, "remote.git"); @@ -84,7 +84,8 @@ async function fixture(layout: AnnotationsLayout): Promise<{ root: string; remot await git(source, ["commit", "-m", "initial"]); await git(source, ["remote", "add", "origin", remote]); await git(source, ["push", "origin", "main"]); - return { root, remote }; + const commit = await git(source, ["rev-parse", "HEAD"]); + return { root, remote, commit }; } test("activation accepts a valid curated annotation blob", async () => { @@ -114,3 +115,21 @@ test("activation rejects a tree at the annotations path", async () => { await expect(registry.bootstrap()).rejects.toMatchObject({ code: "workspace_invalid" }); }); + +test("activation syncs the curated annotations to the revision root when a data root is set", async () => { + const fixtureValue = await fixture("valid"); + const dataRoot = mkdtempSync(join(tmpdir(), "thoth-registry-annotations-data-")); + temporaryRoots.push(dataRoot); + const registry = new WorkspaceRegistry({ + ...config(join(fixtureValue.root, "registry"), fixtureValue.remote), + dataRoot, + }); + + await registry.bootstrap(); + + const annotationsPath = join(dataRoot, "sessions", "psd-clinical", "revisions", fixtureValue.commit, "artifacts", "mschema", "annotations.yaml"); + const manifestPath = join(dataRoot, "sessions", "psd-clinical", "revisions", fixtureValue.commit, "artifacts", "mschema", "annotations.ownership.json"); + expect(readFileSync(annotationsPath, "utf8")).toBe("tables: {}\n"); + const manifest = JSON.parse(readFileSync(manifestPath, "utf8")); + expect(manifest).toMatchObject({ workspace: "psd-clinical", commit: fixtureValue.commit }); +});