feat: atomic revision-qualified annotations sync with ownership manifest (P5)
This commit is contained in:
@@ -251,6 +251,7 @@ export function loadConfig(env: Record<string, string | undefined>): 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(
|
||||
|
||||
@@ -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<string, unknown>;
|
||||
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 };
|
||||
}
|
||||
@@ -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);
|
||||
|
||||
@@ -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 =
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
@@ -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 });
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user