diff --git a/backend/src/workspaces/runtime-config-lease.ts b/backend/src/workspaces/runtime-config-lease.ts index 488acc6e..5dd6fab4 100644 --- a/backend/src/workspaces/runtime-config-lease.ts +++ b/backend/src/workspaces/runtime-config-lease.ts @@ -96,7 +96,9 @@ export function normalizePrivateHostAllowlist(value: string | readonly string[] if (result.includes(host)) throw new Error("HTTP private host allowlist contains a duplicate hostname"); result.push(host); } - return result; + // The policy is a set. Canonical ordering prevents semantically identical + // installation overlays from producing different durable config bytes. + return result.sort(); } diff --git a/backend/test/workspace-runtime-config-lease.test.ts b/backend/test/workspace-runtime-config-lease.test.ts index c56d10b6..3c6578af 100644 --- a/backend/test/workspace-runtime-config-lease.test.ts +++ b/backend/test/workspace-runtime-config-lease.test.ts @@ -5,6 +5,7 @@ import { dirname, join } from "node:path"; import { execFileSync } from "node:child_process"; import { createHash } from "node:crypto"; import { WorkspaceRuntimeConfigLeaseFactory } from "../src/workspaces/runtime-config-lease.js"; +import { parse } from "yaml"; import { parseWorkspaceYaml, serializeWorkspaceYaml } from "../src/workspaces/schema.js"; const workspace = "abc"; @@ -68,7 +69,7 @@ function fixture(extraEnv: Record = {}) { writeFileSync(secret, "secret", { mode: 0o600 }); const configPath = join(harness, "config.yaml"); writeFileSync(configPath, "profile: workstation\n"); - const factory = new WorkspaceRuntimeConfigLeaseFactory({ + const factoryInput = { dataRoot, runtimeSnapshotRoot: snapshots, harnessDir: harness, configPath, env: { THT_WS_ABC_DWH_TRANSPORT: "postgres_direct", THT_WS_ABC_DWH_HOST: "dwh", @@ -78,8 +79,9 @@ function fixture(extraEnv: Record = {}) { internalQdrantUrl: "http://qdrant:6333", internalEmbeddingUrl: "http://embedding:11434", internalEmbeddingModel: "qwen3-embedding:0.6b", internalEmbeddingDimensions: 1024, }, - }); - return { root, repo, snapshotPath, factory, canonicalDescriptor, snapshotManifest: join(snapshotsDir, "snapshot.json") }; + }; + const factory = new WorkspaceRuntimeConfigLeaseFactory(factoryInput); + return { root, repo, snapshotPath, factory, factoryInput, canonicalDescriptor, snapshotManifest: join(snapshotsDir, "snapshot.json") }; } test("session and maintenance share deterministic bytes and path", () => { @@ -251,3 +253,310 @@ test("runtime config symlink replacement is refused", () => { expect(() => f.factory.acquireSession(f.snapshotPath)).toThrow(/trusted|changed|configuration|symbolic/i); } finally { rmSync(f.root, { recursive: true, force: true }); } }); + + +function realHarnessBinding(config: string): Record { + const helper = join(process.cwd(), "..", "harness", "tht", "runtime_config_lease_io.py"); + const python = join(process.cwd(), "..", "harness", ".venv", "bin", "python"); + return JSON.parse(execFileSync(python, [helper], { + cwd: join(process.cwd(), "..", "harness"), encoding: "utf8", + input: JSON.stringify({ action: "binding", config_hex: Buffer.from(config).toString("hex") }), + })); +} + +test("explicit installation overlay is canonical and has one real harness binding", () => { + const f = fixture(); + const overlay = { + profile: "workstation", + session_storage: { type: "postgres_direct", connection: { + host: "session-db", port: 5432, database: "sessions", schema: "public", + user: "runtime", password: "secret", sslmode: "verify-full", + } }, + egress: { http_private_host_allowlist: ["zeta.example", "alpha.example", "warehouse.example"] }, + }; + const sessionFactory = new WorkspaceRuntimeConfigLeaseFactory({ ...f.factoryInput, installationOverlay: overlay }); + const maintenanceFactory = new WorkspaceRuntimeConfigLeaseFactory({ ...f.factoryInput, installationOverlay: overlay }); + try { + const session = sessionFactory.acquireSession(f.snapshotPath); + const maintenance = maintenanceFactory.acquireMaintenance({ workspaceConfigPath: f.snapshotPath }); + const sessionYaml = readFileSync(session.path, "utf8"); + const maintenanceYaml = readFileSync(maintenance.path, "utf8"); + expect(session.path).toBe(maintenance.path); + expect(sessionYaml).toBe(maintenanceYaml); + expect(parse(sessionYaml)).toMatchObject({ + profile: overlay.profile, + session_storage: overlay.session_storage, + egress: { http_private_host_allowlist: ["alpha.example", "warehouse.example", "zeta.example"] }, + }); + const firstBinding = realHarnessBinding(sessionYaml); + const secondBinding = realHarnessBinding(maintenanceYaml); + expect(firstBinding).toEqual(secondBinding); + expect(JSON.parse(readFileSync(session.manifestPath, "utf8")).config_dwh_binding).toEqual(firstBinding); + } finally { rmSync(f.root, { recursive: true, force: true }); } +}); + +test.each([ + ["config-file", "config-file"], ["config-parent", "config-parent"], + ["manifest-file", "manifest-file"], ["manifest-parent", "manifest-parent"], +] as const)("fsync fault at %s fails closed and retries to the same durable pair", (_label, stage) => { + const f = fixture({ THT_RUNTIME_CONFIG_FSYNC_FAIL: stage }); + try { + expect(() => f.factory.acquireSession(f.snapshotPath)).toThrow(/fsync|failed/i); + const runtime = join(f.root, "data", "sessions", workspace, "preprocessing"); + for (const dir of ["runtime-config", "runtime-config-manifests"]) { + if (existsSync(join(runtime, dir))) { + expect(readdirSync(join(runtime, dir)).filter((name) => name.includes("staging")).length).toBe(0); + } + } + const recovered = new WorkspaceRuntimeConfigLeaseFactory({ ...f.factoryInput, env: { + ...f.factoryInput.env, THT_RUNTIME_CONFIG_FSYNC_FAIL: undefined, + } }); + const lease = recovered.acquireSession(f.snapshotPath); + expect(existsSync(lease.path)).toBe(true); + expect(existsSync(lease.manifestPath)).toBe(true); + expect(JSON.parse(readFileSync(lease.manifestPath, "utf8")).config_sha256) + .toBe(createHash("sha256").update(readFileSync(lease.path)).digest("hex")); + } finally { rmSync(f.root, { recursive: true, force: true }); } +}); + +for (const [label, mutate] of [ + ["wrong commit", (f: ReturnType) => f.snapshotPath.replace(/\/[0-9a-f]{40}\//, "/" + "0".repeat(40) + "/")], + ["outside path", (f: ReturnType) => join(f.root, "outside.yaml")], + ["wrong workspace id", (f: ReturnType) => f.snapshotPath.replace("abc.yaml", "abd.yaml")], +] as const) { + test(`rejects ${label} before publication`, () => { + const f = fixture(); + try { + expect(() => f.factory.acquireSession(mutate(f))).toThrow(/trusted|snapshot|identity|path|Git|unavailable|ENOENT|No such/i); + } finally { rmSync(f.root, { recursive: true, force: true }); } + }); +} + +for (const [label, replace] of [ + ["descriptor symlink", (path: string, root: string) => { const target = `${path}.target`; writeFileSync(target, readFileSync(path), { mode: 0o400 }); rmSync(path); execFileSync("ln", ["-s", target, path]); }], + ["descriptor hardlink", (path: string, root: string) => { const target = `${path}.target`; execFileSync("ln", [path, target]); rmSync(path); execFileSync("ln", [target, path]); }], +] as const) { + test(`rejects ${label}`, () => { + const f = fixture(); + try { + replace(f.snapshotPath, f.root); + expect(() => f.factory.acquireSession(f.snapshotPath)).toThrow(/trusted|integrity|identity|link/i); + } finally { rmSync(f.root, { recursive: true, force: true }); } + }); +} + +for (const [label, target] of [ + ["config symlink", "config"], ["config hardlink", "config"], + ["manifest symlink", "manifest"], ["manifest hardlink", "manifest"], +] as const) { + test(`rejects destination ${label} and recovers safely`, () => { + const f = fixture(); + try { + const first = f.factory.acquireSession(f.snapshotPath); + const path = target === "config" ? first.path : first.manifestPath; + const backup = `${path}.target`; + writeFileSync(backup, readFileSync(path), { mode: target === "config" ? 0o400 : 0o600 }); + rmSync(path); + if (label.includes("symlink")) execFileSync("ln", ["-s", backup, path]); + else execFileSync("ln", [backup, path]); + expect(() => f.factory.acquireSession(f.snapshotPath)).toThrow(/trusted|configuration|manifest|link|changed/i); + rmSync(path); + // A replaced inode can never be trusted again. Remove the paired durable + // publication and let a fresh no-replace publication recover the layout. + rmSync(first.path, { force: true }); + rmSync(first.manifestPath, { force: true }); + const recovered = f.factory.acquireSession(f.snapshotPath); + expect(recovered.path).toBe(first.path); + } finally { rmSync(f.root, { recursive: true, force: true }); } + }); +} + +for (const [label, target, mode] of [ + ["config", "config", 0o600], ["manifest", "manifest", 0o400], +] as const) { + test(`refuses wrong ${label} mode then recovers after restoring mode`, () => { + const f = fixture(); + try { + const first = f.factory.acquireSession(f.snapshotPath); + const path = target === "config" ? first.path : first.manifestPath; + chmodSync(path, mode); + expect(() => f.factory.acquireSession(f.snapshotPath)).toThrow(/trusted|mode|configuration|manifest/i); + chmodSync(path, target === "config" ? 0o400 : 0o600); + expect(() => f.factory.acquireSession(f.snapshotPath)).not.toThrow(); + } finally { rmSync(f.root, { recursive: true, force: true }); } + }); +} + + test("release retains durable state while changed binding is refused by a new factory", () => { + const f = fixture(); + try { + const first = f.factory.acquireSession(f.snapshotPath); + first.release(); + const changed = new WorkspaceRuntimeConfigLeaseFactory({ ...f.factoryInput, env: { + ...f.factoryInput.env, THT_WS_ABC_DWH_HOST: "other-dwh", + } }); + expect(() => changed.acquireSession(f.snapshotPath)).toThrow(/changed|mismatch|configuration/i); + expect(existsSync(first.path)).toBe(true); + } finally { rmSync(f.root, { recursive: true, force: true }); } +}); + +test("strict manifest rejects an extra field and recovery preserves exact bytes", () => { + const f = fixture(); + try { + const first = f.factory.acquireSession(f.snapshotPath); + const original = readFileSync(first.manifestPath, "utf8"); + const manifest = JSON.parse(original); + manifest.extra = "reject"; + chmodSync(first.manifestPath, 0o600); + writeFileSync(first.manifestPath, JSON.stringify(manifest), { mode: 0o600 }); + expect(() => f.factory.acquireSession(f.snapshotPath)).toThrow(/manifest|invalid|changed/i); + writeFileSync(first.manifestPath, original, { mode: 0o600 }); + expect(() => f.factory.acquireSession(f.snapshotPath)).not.toThrow(); + } finally { rmSync(f.root, { recursive: true, force: true }); } +}); + + +test.each([ + ["duplicate", "alpha.example,alpha.example"], + ["uppercase", "Alpha.example"], + ["ip address", "127.0.0.1"], +] as const)("rejects %s private-host policy", (_label, allowlist) => { + expect(() => fixture({ THT_HTTP_PRIVATE_HOST_ALLOWLIST: allowlist })) + .toThrow(/allowlist|hostname|duplicate|invalid/i); +}); + +test.each([ + ["runtime-config", "config"], ["runtime-config-manifests", "manifest"], +] as const)("rejects a replaced %s destination ancestor", (directory, _kind) => { + const f = fixture(); + try { + const lease = f.factory.acquireSession(f.snapshotPath); + const parent = join(f.root, "data", "sessions", workspace, "preprocessing", directory); + const moved = `${parent}.moved`; + renameSync(parent, moved); + execFileSync("ln", ["-s", moved, parent]); + expect(() => f.factory.acquireSession(f.snapshotPath)).toThrow(/trusted|symbolic|changed|directory/i); + rmSync(parent); + renameSync(moved, parent); + expect(() => f.factory.acquireSession(f.snapshotPath)).not.toThrow(); + lease.release(); + } finally { rmSync(f.root, { recursive: true, force: true }); } +}); + +test("release is idempotent and does not remove either durable publication", () => { + const f = fixture(); + try { + const lease = f.factory.acquireSession(f.snapshotPath); + lease.release(); lease.release(); + expect(existsSync(lease.path)).toBe(true); + expect(existsSync(lease.manifestPath)).toBe(true); + } finally { rmSync(f.root, { recursive: true, force: true }); } +}); + +test("partial os.write calls are completed by the real Python publication helper", () => { + const python = join(process.cwd(), "..", "harness", ".venv", "bin", "python"); + const helper = join(process.cwd(), "..", "harness", "tht", "runtime_config_lease_io.py"); + const code = `import os, sys; sys.path.insert(0, ${JSON.stringify(dirname(helper))}); import runtime_config_lease_io as m; real=os.write; os.write=lambda fd,b: real(fd,b[:3]); m.write_all(1, b'partial-write-ok\\n')`; + const output = execFileSync(python, ["-c", code], { encoding: "utf8" }); + expect(output).toBe("partial-write-ok\n"); +}); + +test("a clean existing equal publication is reconciled by a new factory", () => { + const f = fixture(); + try { + const first = f.factory.acquireSession(f.snapshotPath); + const second = new WorkspaceRuntimeConfigLeaseFactory(f.factoryInput).acquireMaintenance({ snapshotPath: f.snapshotPath }); + expect(readFileSync(second.path)).toEqual(readFileSync(first.path)); + expect(readFileSync(second.manifestPath)).toEqual(readFileSync(first.manifestPath)); + } finally { rmSync(f.root, { recursive: true, force: true }); } +}); + +test.each([["config"], ["manifest"]] as const)("rename failure is scoped to the %s branch and leaves no staging", (kind) => { + const f = fixture({ THT_RUNTIME_CONFIG_RENAME_FAIL: kind }); + try { + expect(() => f.factory.acquireSession(f.snapshotPath)).toThrow(/rename|failed/i); + const runtime = join(f.root, "data", "sessions", workspace, "preprocessing"); + for (const dir of ["runtime-config", "runtime-config-manifests"]) { + if (existsSync(join(runtime, dir))) expect(readdirSync(join(runtime, dir)).some((name) => name.includes("staging"))).toBe(false); + } + } finally { rmSync(f.root, { recursive: true, force: true }); } +}); + +test("snapshot manifest rejects an undeclared extra immutable file", () => { + const f = fixture(); + try { + const snapshot = JSON.parse(readFileSync(f.snapshotManifest, "utf8")); + writeFileSync(join(dirname(f.snapshotPath), "smuggled.txt"), "smuggled", { mode: 0o400 }); + snapshot.files["smuggled.txt"] = createHash("sha256").update("smuggled").digest("hex"); + chmodSync(f.snapshotManifest, 0o600); + writeFileSync(f.snapshotManifest, JSON.stringify(snapshot), { mode: 0o400 }); + expect(() => f.factory.acquireSession(f.snapshotPath)).toThrow(/integrity|snapshot|trusted/i); + } finally { rmSync(f.root, { recursive: true, force: true }); } +}); + + +test("independent OS publishers converge when equal and elect one winner when unequal", () => { + const f = fixture(); + try { + const seed = f.factory.acquireSession(f.snapshotPath); + const full = JSON.parse(readFileSync(seed.manifestPath, "utf8")); + const base = Object.fromEntries(["workspace_id", "workspace_revision", "descriptor_git_blob", + "descriptor_sha256", "descriptor_dev", "descriptor_ino", "config_dwh_binding"] + .map((key) => [key, full[key]])); + const python = join(process.cwd(), "..", "harness", ".venv", "bin", "python"); + const helper = join(process.cwd(), "..", "harness", "tht", "runtime_config_lease_io.py"); + const configBytes = readFileSync(seed.path); + // Leave the workspace-owned destination directories in place, but remove + // both durable leaves: the following OS processes race on a clean layout. + rmSync(seed.path); rmSync(seed.manifestPath); + const payload = JSON.stringify({ action: "publish", data_root: f.factoryInput.dataRoot, + workspace_id: workspace, workspace_revision: full.workspace_revision, + config_hex: Buffer.from(configBytes).toString("hex"), manifest_base: base }); + const code = `import json,multiprocessing,sys +multiprocessing.set_start_method("fork") +sys.path.insert(0, ${JSON.stringify(dirname(helper))}) +import runtime_config_lease_io as m +import contextlib,os +def run(x): + try: + with open(os.devnull,"w") as error, contextlib.redirect_stderr(error): m.publish(x) + except Exception: raise SystemExit(1) +x=json.loads(sys.argv[1]); y=json.loads(sys.argv[1]) +if len(sys.argv)>2: y["config_hex"]="646966666572656e742d72756e74696d652d636f6e666967" +p=[multiprocessing.Process(target=run,args=(x,)),multiprocessing.Process(target=run,args=(y,))] +[q.start() for q in p]; [q.join() for q in p] +print(json.dumps([q.exitcode for q in p]))` + const equal = JSON.parse(execFileSync(python, ["-c", code, payload], { encoding: "utf8" }).trim()); + expect(equal.sort()).toEqual([0, 0]); + rmSync(join(f.factoryInput.dataRoot, "sessions", workspace, "preprocessing", "runtime-config", `${full.workspace_revision}.yaml`)); + rmSync(join(f.factoryInput.dataRoot, "sessions", workspace, "preprocessing", "runtime-config-manifests", `${full.workspace_revision}.json`)); + const unequalPayload = JSON.stringify({ ...JSON.parse(payload), config_hex: Buffer.from("different-runtime-config").toString("hex") }); + const unequal = JSON.parse(execFileSync(python, ["-c", code, payload, "unequal"], { encoding: "utf8" }).trim()); + expect(unequal.filter((exit: number) => exit === 0)).toHaveLength(1); + expect(unequal.filter((exit: number) => exit !== 0)).toHaveLength(1); + } finally { rmSync(f.root, { recursive: true, force: true }); } +}); + + +test.each(["leaf", "ancestor"] as const)("actual harness rejects canonical %s swap", (kind) => { + const f = fixture(); + try { + const lease = f.factory.acquireSession(f.snapshotPath); + const harnessBin = join(process.cwd(), "..", "harness", ".venv", "bin", "tht"); + const harnessCwd = join(process.cwd(), "..", "harness"); + const env = { ...process.env, THT_RUNTIME_CONFIG_MANIFEST_SHA256: lease.manifestSha256 }; + if (kind === "leaf") { + const moved = `${lease.path}.moved`; + renameSync(lease.path, moved); + execFileSync("ln", ["-s", moved, lease.path]); + } else { + const parent = dirname(lease.path); + const moved = `${parent}.moved`; + renameSync(parent, moved); + execFileSync("ln", ["-s", moved, parent]); + } + expect(() => execFileSync(harnessBin, ["config", "check", "-c", lease.path], { + cwd: harnessCwd, env, stdio: "pipe", + })).toThrow(); + } finally { rmSync(f.root, { recursive: true, force: true }); } +}); diff --git a/harness/tht/runtime_config_lease_io.py b/harness/tht/runtime_config_lease_io.py index 9ef4076e..9fc9a2b5 100644 --- a/harness/tht/runtime_config_lease_io.py +++ b/harness/tht/runtime_config_lease_io.py @@ -236,6 +236,18 @@ def write_all(fd: int, data: bytes) -> None: pos += n +def publication_fsync(fd: int, stage: str) -> None: + """Fsync one publication boundary, with an explicit test-only fault seam. + + The seam is inert unless the caller opts in with the exact stage name. It is + intentionally kept here (rather than in TypeScript) so the real helper's + durable ordering can be exercised without weakening production behaviour. + """ + if os.environ.get("THT_RUNTIME_CONFIG_FSYNC_FAIL") == stage: + fail(f"runtime config fsync failed at {stage}") + os.fsync(fd) + + def read_all(fd: int, limit: int = 16 * 1024 * 1024) -> bytes: os.lseek(fd, 0, os.SEEK_SET) chunks: list[bytes] = [] @@ -361,7 +373,7 @@ def publish(inp: dict) -> dict: try: write_all(fd, content) os.fchmod(fd, 0o400) - os.fsync(fd) + publication_fsync(fd, "config-file") try: _rename_noreplace(stage, name, cfgdir, "config") except FileExistsError: @@ -384,7 +396,7 @@ def publish(inp: dict) -> dict: fail("same-revision runtime configuration changed") finally: os.close(fd) - os.fsync(cfgdir) + publication_fsync(cfgdir, "config-parent") # Identity is deliberately recorded after final no-replace publication. got = current(cfgdir, name, 0o400) assert got @@ -423,7 +435,7 @@ def publish(inp: dict) -> dict: try: write_all(fd, mb) os.fchmod(fd, 0o600) - os.fsync(fd) + publication_fsync(fd, "manifest-file") try: _rename_noreplace(stage, mname, mandir, "manifest") except FileExistsError: @@ -449,7 +461,7 @@ def publish(inp: dict) -> dict: os.unlink(stage, dir_fd=mandir) except FileNotFoundError: pass - os.fsync(mandir) + publication_fsync(mandir, "manifest-parent") return { "path": f"{root}/sessions/{wid}/preprocessing/runtime-config/{name}", "manifestPath": f"{root}/sessions/{wid}/preprocessing/runtime-config-manifests/{mname}",