From f7c2b69837c9782a3dda35bcbdb88ce9773125af Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 11 Aug 2026 18:40:11 +0200 Subject: [PATCH] feat: P2 operator, preprocessing state/service, and runtime config lease --- backend/src/tht/tht-runner.ts | 33 +- backend/src/workspace-maintenance.ts | 285 ++++++++++ .../src/workspaces/preprocessing-service.ts | 505 +++++++++++++++++ backend/src/workspaces/preprocessing-state.ts | 383 +++++++++++++ .../src/workspaces/runtime-config-lease.ts | 522 ++++++++++++++++++ backend/test/tht-runner.test.ts | 28 +- backend/test/workspace-maintenance.test.ts | 85 +++ .../workspace-preprocessing-service.test.ts | 333 +++++++++++ .../workspace-preprocessing-state.test.ts | 119 ++++ .../workspace-runtime-config-lease.test.ts | 239 ++++++++ 10 files changed, 2501 insertions(+), 31 deletions(-) create mode 100644 backend/src/workspace-maintenance.ts create mode 100644 backend/src/workspaces/preprocessing-service.ts create mode 100644 backend/src/workspaces/preprocessing-state.ts create mode 100644 backend/src/workspaces/runtime-config-lease.ts create mode 100644 backend/test/workspace-maintenance.test.ts create mode 100644 backend/test/workspace-preprocessing-service.test.ts create mode 100644 backend/test/workspace-preprocessing-state.test.ts create mode 100644 backend/test/workspace-runtime-config-lease.test.ts diff --git a/backend/src/tht/tht-runner.ts b/backend/src/tht/tht-runner.ts index db46c031..8aa57c8b 100644 --- a/backend/src/tht/tht-runner.ts +++ b/backend/src/tht/tht-runner.ts @@ -8,9 +8,8 @@ import { dirname, isAbsolute, join, relative, resolve } from "node:path"; import { parseAllDocuments } from "yaml"; import { clearPrincipalEnvironment, principalEnvironment, type PrincipalContext } from "../auth/principal.js"; import { secretValue, type SecretBundleConfig } from "../config/secret-bundle.js"; -import { resolveRuntimeBindings } from "../workspaces/bindings.js"; +import { renderWorkspaceRuntimeFromSnapshotPath } from "../workspaces/runtime-config-lease.js"; import { - renderRuntimeConfig, type RuntimeInstallationOverlay, type RuntimePaths, type SemanticRuntimeConfig, @@ -208,26 +207,22 @@ export class ThtRunner { /** Render one immutable canonical registry revision into a backend-owned harness config. */ acquireWorkspaceRuntime(workspaceConfigPath: string): RuntimeConfigLease { - const canonical = this.readCanonicalWorkspaceSnapshot(workspaceConfigPath); - const bindings = resolveRuntimeBindings( - canonical.workspace, - process.env, - this.cfg.secretRoots ?? [], - ); - const config = renderRuntimeConfig( - canonical.workspace, - bindings, - this.runtimePaths(canonical.workspaceId), - canonical, - this.installationOverlay(), - this.cfg.semanticRuntime, - ); - const path = this.createRuntimeSnapshot(config); + const rendered = renderWorkspaceRuntimeFromSnapshotPath({ + snapshotPath: workspaceConfigPath, + harnessDir: this.cfg.harnessDir, + configPath: this.cfg.configPath, + dataRoot: this.cfg.dataRoot ?? (() => { + throw new Error("registry workspace runtime requires an absolute data root"); + })(), + secretRoots: this.cfg.secretRoots ?? [], + semanticRuntime: this.cfg.semanticRuntime, + }); + const path = this.createRuntimeSnapshot(rendered.renderedConfig); let released = false; return { path, - workspaceId: canonical.workspaceId, - workspaceRevision: canonical.workspaceRevision, + workspaceId: rendered.workspaceId, + workspaceRevision: rendered.workspaceRevision, release: () => { if (released) return; released = true; diff --git a/backend/src/workspace-maintenance.ts b/backend/src/workspace-maintenance.ts new file mode 100644 index 00000000..e506eeaa --- /dev/null +++ b/backend/src/workspace-maintenance.ts @@ -0,0 +1,285 @@ +import { spawn } from "node:child_process"; +import { closeSync, constants as fsConstants, openSync } from "node:fs"; +import { readdir, readFile } from "node:fs/promises"; +import { join } from "node:path"; +import { parse } from "yaml"; +import { loadConfig } from "./config.js"; +import { ThtRunner } from "./tht/tht-runner.js"; +import { WorkspaceRegistry } from "./workspaces/registry.js"; +import { publishDeterministicRuntimeConfigLease, renderActiveWorkspaceRuntime } from "./workspaces/runtime-config-lease.js"; +import { WorkspacePreprocessingService, type WorkspaceOperationResult } from "./workspaces/preprocessing-service.js"; +import type { SessionInventoryRow } from "./workspaces/preprocessing-state.js"; + +export interface WorkspaceMaintenanceIo { + stdin: string; + stdout: string[]; + stderr: string[]; + writeStdout(value: string): void; + writeStderr(value: string): void; +} + +type Command = "inspect" | "preprocess-dwh" | "schema-suggest-fks" | "schema-check" | "index-schema" | "preprocess-evidence" | "preprocess-run"; + +function failureResult(operation: string, workspaceId = ""): WorkspaceOperationResult { + return { + schemaVersion: 1, + status: "failed", + code: "workspace_not_activatable", + workspaceId, + workspaceRevision: "", + descriptorBlob: "", + operation, + completedStages: [], + }; +} + +function boundedJson(result: WorkspaceOperationResult): string { + const encoded = JSON.stringify(result); + if (Buffer.byteLength(encoded, "utf8") > 1024 * 1024) { + return JSON.stringify(failureResult(result.operation || "unknown", result.workspaceId)); + } + return encoded; +} + +function sanitizeStderr(_error: unknown): string { + return "workspace maintenance failed\n"; +} + +function parseRequest(command: string, stdin: string): Record { + const parsed = JSON.parse(stdin) as Record; + if (!parsed || typeof parsed !== "object" || Array.isArray(parsed) || parsed.schemaVersion !== 1) { + throw new Error("invalid request"); + } + const allowedByCommand: Record = { + inspect: ["schemaVersion", "workspaceId"], + "preprocess-dwh": ["schemaVersion", "workspaceId", "resumeRunId"], + "schema-suggest-fks": ["schemaVersion", "workspaceId", "fromSql", "assume", "resumeRunId"], + "schema-check": ["schemaVersion", "workspaceId", "annotationsYaml", "reviewedCandidatesDigest"], + "index-schema": ["schemaVersion", "workspaceId", "resumeRunId"], + "preprocess-evidence": ["schemaVersion", "workspaceId", "dryRun", "resumeRunId"], + "preprocess-run": ["schemaVersion", "workspaceId", "resumeRunId"], + }; + const allowed = allowedByCommand[command]; + if (!allowed) throw new Error("unknown command"); + if (typeof parsed.workspaceId !== "string") throw new Error("invalid workspace id"); + for (const key of Object.keys(parsed)) if (!allowed.includes(key)) throw new Error("unexpected request field"); + return parsed; +} + +function exitCodeFor(result: WorkspaceOperationResult): number { + if (["succeeded", "unchanged", "dry_run"].includes(result.status)) return 0; + if (result.status === "blocked") return 3; + return 1; +} + +async function dispatch(command: Command, service: WorkspacePreprocessingService, request: Record): Promise { + switch (command) { + case "inspect": + return await service.inspect({ workspaceId: request.workspaceId as string }); + case "preprocess-dwh": + return await service.preprocessDwh({ + workspaceId: request.workspaceId as string, + resumeRunId: request.resumeRunId as string | undefined, + }); + case "schema-suggest-fks": + return await service.suggestFks({ + workspaceId: request.workspaceId as string, + fromSql: request.fromSql as any, + assume: request.assume as any, + resumeRunId: request.resumeRunId as string | undefined, + }); + case "schema-check": + return await service.checkSchema({ + workspaceId: request.workspaceId as string, + annotationsYaml: request.annotationsYaml as string | undefined, + reviewedCandidatesDigest: request.reviewedCandidatesDigest as string | undefined, + }); + case "index-schema": + return await service.indexSchema({ + workspaceId: request.workspaceId as string, + resumeRunId: request.resumeRunId as string | undefined, + }); + case "preprocess-evidence": + return await service.preprocessEvidence({ + workspaceId: request.workspaceId as string, + dryRun: request.dryRun as boolean | undefined, + resumeRunId: request.resumeRunId as string | undefined, + }); + case "preprocess-run": + return await service.run({ + workspaceId: request.workspaceId as string, + resumeRunId: request.resumeRunId as string | undefined, + }); + } +} + +export async function runWorkspaceMaintenanceCli( + argv: readonly string[], + service: WorkspacePreprocessingService, + io: WorkspaceMaintenanceIo, +): Promise { + const command = argv[2]; + if (!command) { + const result = failureResult("unknown"); + io.writeStdout(boundedJson(result)); + return 2; + } + try { + const request = parseRequest(command, io.stdin); + const result = await dispatch(command as Command, service, request); + io.writeStdout(boundedJson(result)); + return exitCodeFor(result); + } catch (error) { + const result = failureResult(command, (() => { + try { return JSON.parse(io.stdin).workspaceId ?? ""; } catch { return ""; } + })()); + io.writeStdout(boundedJson(result)); + io.writeStderr(sanitizeStderr(error)); + const message = String((error as Error).message ?? ""); + const requestError = error instanceof SyntaxError + || message === "invalid request" + || message === "unknown command" + || message === "unexpected request field" + || message === "invalid workspace id"; + return command in { + inspect: true, "preprocess-dwh": true, "schema-suggest-fks": true, "schema-check": true, + "index-schema": true, "preprocess-evidence": true, "preprocess-run": true, + } ? (requestError ? 2 : 1) : 2; + } +} + +async function readSessionInventory(dataRoot: string, workspaceId: string): Promise { + const directory = join(dataRoot, "sessions", workspaceId, "sessions"); + try { + const entries = await readdir(directory, { withFileTypes: true }); + const rows: SessionInventoryRow[] = []; + for (const entry of entries) { + if (!entry.isDirectory() || entry.isSymbolicLink()) continue; + try { + const source = await readFile(join(directory, entry.name, "session_manifest.yaml"), "utf8"); + const manifest = parse(source) as Record; + rows.push({ + id: entry.name, + status: typeof manifest.status === "string" ? manifest.status : "open", + archived: manifest.archived === true, + workspaceRevision: typeof manifest.workspace_revision === "string" ? manifest.workspace_revision : null, + }); + } catch { + // fail closed at mutation time by ignoring unreadable manifests from the resumable scan + } + } + return rows; + } catch { + return []; + } +} + +function createProductionService(): WorkspacePreprocessingService { + const config = loadConfig(process.env); + const registry = new WorkspaceRegistry(config.workspaceRegistry); + const runner = new ThtRunner({ + thtBin: config.thtBin, + harnessDir: config.harnessDir, + configPath: process.env.THT_CONFIG ?? "config/tht.yaml", + dataRoot: config.dataRoot, + runtimeSnapshotRoot: join(config.workspaceRegistry.root, "snapshots", "runtime"), + secretRoots: config.workspaceRegistry.secretRoots, + secretsFile: config.secretsFile, + secretFiles: config.secretFiles, + semanticRuntime: { + internalQdrantUrl: config.internalQdrantUrl, + internalEmbeddingUrl: config.internalEmbeddingUrl, + internalEmbeddingModel: config.internalEmbeddingModel, + internalEmbeddingDimensions: config.internalEmbeddingDimensions, + }, + }); + return new WorkspacePreprocessingService({ + dataRoot: config.dataRoot ?? "/data", + acquireActiveRuntime: async (workspaceId) => { + const active = await renderActiveWorkspaceRuntime({ + workspaceId, + registry, + registryConfig: config.workspaceRegistry, + harnessDir: config.harnessDir, + configPath: process.env.THT_CONFIG ?? "config/tht.yaml", + dataRoot: config.dataRoot ?? "/data", + secretRoots: config.workspaceRegistry.secretRoots, + semanticRuntime: { + internalQdrantUrl: config.internalQdrantUrl, + internalEmbeddingUrl: config.internalEmbeddingUrl, + internalEmbeddingModel: config.internalEmbeddingModel, + internalEmbeddingDimensions: config.internalEmbeddingDimensions, + }, + }); + const configLease = await publishDeterministicRuntimeConfigLease({ + workspaceId, + registry, + registryConfig: config.workspaceRegistry, + harnessDir: config.harnessDir, + configPath: process.env.THT_CONFIG ?? "config/tht.yaml", + dataRoot: config.dataRoot ?? "/data", + secretRoots: config.workspaceRegistry.secretRoots, + semanticRuntime: { + internalQdrantUrl: config.internalQdrantUrl, + internalEmbeddingUrl: config.internalEmbeddingUrl, + internalEmbeddingModel: config.internalEmbeddingModel, + internalEmbeddingDimensions: config.internalEmbeddingDimensions, + }, + }); + return { + workspace: active.workspace, + workspaceId: active.workspaceId, + workspaceRevision: active.workspaceRevision, + descriptorBlob: active.descriptorBlob, + catalogBlob: active.catalogBlob, + configLease, + }; + }, + runChild: async ({ argv, configPath }) => { + const configFd = openSync(configPath, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW); + try { + return await new Promise((resolve) => { + const child = spawn(config.thtBin, argv, { + cwd: config.harnessDir, + env: { ...process.env, ...(config.dataRoot ? { THT_DATA_ROOT: config.dataRoot } : {}) }, + stdio: ["ignore", "pipe", "pipe", configFd], + }); + let stdout = ""; + let stderr = ""; + child.stdout?.on("data", (chunk: Buffer) => { stdout += chunk.toString("utf8"); }); + child.stderr?.on("data", (chunk: Buffer) => { stderr += chunk.toString("utf8"); }); + child.on("close", (code) => resolve({ exitCode: code ?? 0, stdout, stderr })); + child.on("error", (error) => resolve({ exitCode: 1, stdout, stderr: String(error.message) })); + }); + } finally { + closeSync(configFd); + } + }, + listSessions: async (workspaceId) => await readSessionInventory(config.dataRoot ?? "/data", workspaceId), + semanticPreflight: async (workspace) => { + const result = await runner.qdrantEnsure(workspace, 30); + return result.ok ? { ok: true as const } : { ok: false as const, code: result.code ?? "workspace_not_activatable" }; + }, + }); +} + +if (process.argv[1] && import.meta.url === new URL(`file://${process.argv[1]}`).href) { + const stdout: string[] = []; + const stderr: string[] = []; + const io: WorkspaceMaintenanceIo = { + stdin: await new Promise((resolve) => { + let input = ""; + process.stdin.setEncoding("utf8"); + process.stdin.on("data", (chunk) => { input += chunk; }); + process.stdin.on("end", () => resolve(input)); + }), + stdout, + stderr, + writeStdout: (value) => { stdout.push(value); }, + writeStderr: (value) => { stderr.push(value); }, + }; + const exitCode = await runWorkspaceMaintenanceCli(process.argv, createProductionService(), io); + process.stdout.write(stdout.join("")); + if (stderr.length > 0) process.stderr.write(stderr.join("").slice(0, 64 * 1024)); + process.exit(exitCode); +} diff --git a/backend/src/workspaces/preprocessing-service.ts b/backend/src/workspaces/preprocessing-service.ts new file mode 100644 index 00000000..da5d7bc6 --- /dev/null +++ b/backend/src/workspaces/preprocessing-service.ts @@ -0,0 +1,505 @@ +import { createHash, randomBytes } from "node:crypto"; +import { readdirSync, readFileSync, rmSync, writeFileSync, mkdirSync } from "node:fs"; +import { isIP } from "node:net"; +import { join } from "node:path"; +import type { WorkspaceDescriptor } from "./schema.js"; +import { + PreprocessingStateStore, + type FkReviewRecord, + type PreprocessingJobState, + type SessionInventoryRow, +} from "./preprocessing-state.js"; +import type { DeterministicRuntimeConfigLease } from "./runtime-config-lease.js"; + +export interface WorkspaceOperationResult { + schemaVersion: 1; + status: "succeeded" | "unchanged" | "dry_run" | "blocked" | "failed"; + code: + | "ok" | "workspace_not_found" | "workspace_not_activatable" + | "binding_missing" | "preprocessing_conflict" + | "preprocessing_resume_mismatch" | "manual_review_required" + | "evidence_materialization_required" | "effective_config_mismatch" + | "semantic_index_incompatible" | "annotation_invalid" + | "egress_policy_refused"; + workspaceId: string; + workspaceRevision: string; + descriptorBlob: string; + operation: string; + runId?: string; + childRuns?: Record; + completedStages: string[]; + counts?: Record; + artifactIdentities?: Array<{ kind: string; digest: string }>; + /** Suggested FK annotations YAML for the operator to write to --output (schema suggest-fks). */ + suggestedFksYaml?: string; + warnings?: string[]; +} + +export interface ChildProcessRequest { + argv: string[]; + configPath: string; +} + +interface ActiveRuntime { + workspace: WorkspaceDescriptor; + workspaceId: string; + workspaceRevision: string; + descriptorBlob: string; + catalogBlob: string; + configLease: DeterministicRuntimeConfigLease; +} + +interface ChildProcessResult { + exitCode: number; + stdout: string; + stderr: string; +} + +export interface WorkspacePreprocessingServiceDeps { + dataRoot: string; + acquireActiveRuntime(workspaceId: string): Promise; + runChild(request: ChildProcessRequest): Promise; + listSessions(workspaceId: string): Promise; + semanticPreflight(workspace: WorkspaceDescriptor): Promise< + { ok: true } | { ok: false; code: "workspace_not_activatable" | "semantic_index_incompatible" } + >; + httpPrivateHostAllowlist?: readonly string[]; +} + +interface RunScope { + runtime: ActiveRuntime; + state: PreprocessingStateStore; + job: PreprocessingJobState; +} + +function digest(value: string | Buffer): string { + return `sha256:${createHash("sha256").update(value).digest("hex")}`; +} + +function baseResult( + runtime: ActiveRuntime, + operation: string, + status: WorkspaceOperationResult["status"], + code: WorkspaceOperationResult["code"], + extra: Omit, "schemaVersion" | "status" | "code" | "workspaceId" | "workspaceRevision" | "descriptorBlob" | "operation"> = {}, +): WorkspaceOperationResult { + return { + schemaVersion: 1, + status, + code, + workspaceId: runtime.workspaceId, + workspaceRevision: runtime.workspaceRevision, + descriptorBlob: runtime.descriptorBlob, + operation, + completedStages: [], + ...extra, + }; +} + +function isPrivateHost(hostname: string): boolean { + if (hostname === "localhost" || hostname === "metadata.google.internal") return true; + const address = isIP(hostname); + if (address === 4) { + if (/^127\./.test(hostname) || /^10\./.test(hostname) || /^192\.168\./.test(hostname)) return true; + if (/^169\.254\./.test(hostname) || /^0\./.test(hostname)) return true; + const match = /^172\.(\d+)\./.exec(hostname); + return Boolean(match && Number(match[1]) >= 16 && Number(match[1]) <= 31); + } + if (address === 6) { + const normalized = hostname.toLowerCase(); + return normalized === "::1" || normalized.startsWith("fe80:") || normalized.startsWith("fd") || normalized.startsWith("fc"); + } + return hostname.endsWith(".internal"); +} + +function noEvidenceWarning(workspace: WorkspaceDescriptor): string[] { + return workspace.evidence === undefined ? ["workspace has no Evidence source"] : []; +} + +export class WorkspacePreprocessingService { + constructor(private readonly deps: WorkspacePreprocessingServiceDeps) {} + + async inspect(options: { workspaceId: string }): Promise { + try { + const runtime = await this.deps.acquireActiveRuntime(options.workspaceId); + return baseResult(runtime, "inspect", "succeeded", "ok", { + artifactIdentities: [ + { kind: "descriptor", digest: runtime.descriptorBlob }, + { kind: "catalog", digest: runtime.catalogBlob }, + { kind: "runtime_config", digest: runtime.configLease.configDigest }, + ], + }); + } catch { + return { + schemaVersion: 1, + status: "failed", + code: "workspace_not_activatable", + workspaceId: options.workspaceId, + workspaceRevision: "", + descriptorBlob: "", + operation: "inspect", + completedStages: [], + }; + } + } + + async preprocessDwh(options: { workspaceId: string; resumeRunId?: string }): Promise { + const scope = await this.startRun(options.workspaceId, "preprocess dwh", options.resumeRunId); + if (scope.job.completedStages.includes("dwh")) { + return baseResult(scope.runtime, "preprocess dwh", "unchanged", "ok", { + runId: scope.job.runId, + childRuns: scope.job.childRuns, + completedStages: [...scope.job.completedStages], + }); + } + const payload = await this.runJsonStage(scope.runtime, [ + "preprocess", "dwh", "--steps", "introspect,lsh", + ...(scope.job.childRuns.dwh ? ["--resume", scope.job.childRuns.dwh] : []), + "--json", "-c", "/dev/fd/3", + ]); + const childRun = this.requireRunId(payload.run_id); + scope.job.childRuns.dwh = childRun; + if (!scope.job.completedStages.includes("dwh")) scope.job.completedStages.push("dwh"); + this.state(scope.runtime.workspaceId).writeJob(scope.job); + return baseResult(scope.runtime, "preprocess dwh", "succeeded", "ok", { + runId: scope.job.runId, + childRuns: { ...scope.job.childRuns }, + completedStages: [...scope.job.completedStages], + }); + } + + async suggestFks(options: { + workspaceId: string; + fromSql?: ReadonlyArray<{ name: string; sql: string }>; + assume?: readonly string[]; + resumeRunId?: string; + }): Promise { + const scope = await this.startRun(options.workspaceId, "schema suggest-fks", options.resumeRunId); + return await this.runSuggestStage(scope, options.fromSql ?? [], options.assume ?? []); + } + + async checkSchema(options: { + workspaceId: string; + annotationsYaml?: string; + reviewedCandidatesDigest?: string; + }): Promise { + const runtime = await this.deps.acquireActiveRuntime(options.workspaceId); + const state = this.state(runtime.workspaceId); + if ((options.annotationsYaml === undefined) !== (options.reviewedCandidatesDigest === undefined)) { + return baseResult(runtime, "schema check", "failed", "annotation_invalid"); + } + if (options.reviewedCandidatesDigest === undefined) { + const payload = await this.runJsonStage(runtime, ["schema", "check", "--json", "-c", "/dev/fd/3"]); + return baseResult(runtime, "schema check", Number(payload.orphan_count ?? 0) === 0 ? "succeeded" : "failed", Number(payload.orphan_count ?? 0) === 0 ? "ok" : "annotation_invalid"); + } + const reviewedCandidatesDigest = options.reviewedCandidatesDigest; + const runId = this.findRunIdByCandidateDigest(state, reviewedCandidatesDigest); + if (!runId) return baseResult(runtime, "schema check", "failed", "annotation_invalid"); + const request = await this.withStagedInputs(runtime.workspaceId, [ + { flag: "--annotations", name: "annotations.yaml", contents: options.annotationsYaml! }, + ], async (argv) => await this.runJsonStage(runtime, [ + "schema", "check", ...argv, + "--reviewed-candidates", reviewedCandidatesDigest, + "--json", "-c", "/dev/fd/3", + ])); + if (request.reviewed_candidates_digest !== reviewedCandidatesDigest || typeof request.annotations_digest !== "string") { + return baseResult(runtime, "schema check", "failed", "annotation_invalid", { runId }); + } + const annotationsDigest = request.annotations_digest; + const review = state.writeFkReview(runId, { + reviewedCandidatesDigest, + annotationsDigest, + workspaceRevision: runtime.workspaceRevision, + }); + const job = state.readJob(runId); + if (!job.completedStages.includes("fk_review")) { + job.reviewDigest = review.digest; + job.completedStages.push("fk_review"); + state.writeJob(job); + } + return baseResult(runtime, "schema check", "succeeded", "ok", { + runId, + completedStages: [...job.completedStages], + artifactIdentities: [{ kind: "fk_review", digest: review.digest }], + }); + } + + async indexSchema(options: { workspaceId: string; resumeRunId?: string }): Promise { + const scope = await this.startRun(options.workspaceId, "index-schema", options.resumeRunId); + const semantic = await this.deps.semanticPreflight(scope.runtime.workspace); + if (!semantic.ok) return baseResult(scope.runtime, "index-schema", "failed", semantic.code, { runId: scope.job.runId }); + if (scope.job.completedStages.includes("schema_index")) { + return baseResult(scope.runtime, "index-schema", "unchanged", "ok", { + runId: scope.job.runId, + completedStages: [...scope.job.completedStages], + }); + } + const payload = await this.runJsonStage(scope.runtime, ["vector", "index-schema", "--json", "-c", "/dev/fd/3"]); + const counts = this.numberRecord(payload.counts); + scope.job.completedStages.push("schema_index"); + this.state(scope.runtime.workspaceId).writeJob(scope.job); + return baseResult(scope.runtime, "index-schema", "succeeded", "ok", { + runId: scope.job.runId, + completedStages: [...scope.job.completedStages], + counts, + }); + } + + async preprocessEvidence(options: { workspaceId: string; dryRun?: boolean; resumeRunId?: string }): Promise { + const scope = await this.startRun(options.workspaceId, "preprocess evidence", options.resumeRunId); + if (scope.runtime.workspace.evidence === undefined) { + return baseResult(scope.runtime, "preprocess evidence", "unchanged", "ok", { + warnings: noEvidenceWarning(scope.runtime.workspace), + }); + } + const policy = this.evidencePolicy(scope.runtime.workspace); + if (policy !== undefined) return baseResult(scope.runtime, "preprocess evidence", policy.status, policy.code, { warnings: policy.warnings }); + const semantic = await this.deps.semanticPreflight(scope.runtime.workspace); + if (!semantic.ok) return baseResult(scope.runtime, "preprocess evidence", "failed", semantic.code, { runId: scope.job.runId }); + if (scope.job.completedStages.includes("evidence") && !options.dryRun) { + return baseResult(scope.runtime, "preprocess evidence", "unchanged", "ok", { + runId: scope.job.runId, + completedStages: [...scope.job.completedStages], + }); + } + const payload = await this.runJsonStage(scope.runtime, [ + "preprocess", "evidence", + ...(options.dryRun ? ["--dry-run"] : []), + ...(scope.job.childRuns.evidence ? ["--resume", scope.job.childRuns.evidence] : []), + "--json", "-c", "/dev/fd/3", + ]); + if (typeof payload.run_id === "string") scope.job.childRuns.evidence = this.requireRunId(payload.run_id); + if (!options.dryRun && !scope.job.completedStages.includes("evidence")) scope.job.completedStages.push("evidence"); + this.state(scope.runtime.workspaceId).writeJob(scope.job); + return baseResult(scope.runtime, "preprocess evidence", options.dryRun ? "dry_run" : "succeeded", "ok", { + runId: scope.job.runId, + childRuns: { ...scope.job.childRuns }, + completedStages: [...scope.job.completedStages], + counts: this.numberRecord(payload.counts), + }); + } + + async run(options: { workspaceId: string; resumeRunId?: string }): Promise { + const scope = await this.startRun(options.workspaceId, "preprocess run", options.resumeRunId); + if (!scope.job.completedStages.includes("dwh")) { + const payload = await this.runJsonStage(scope.runtime, [ + "preprocess", "dwh", "--steps", "introspect,lsh", + ...(scope.job.childRuns.dwh ? ["--resume", scope.job.childRuns.dwh] : []), + "--json", "-c", "/dev/fd/3", + ]); + scope.job.childRuns.dwh = this.requireRunId(payload.run_id); + scope.job.completedStages.push("dwh"); + this.state(scope.runtime.workspaceId).writeJob(scope.job); + } + if (!scope.job.completedStages.includes("fk_suggest")) { + const suggest = await this.runSuggestStage(scope, [], []); + if (suggest.code === "manual_review_required") return suggest; + } + const candidate = this.state(scope.runtime.workspaceId).readFkCandidates(scope.job.runId); + if (candidate && !scope.job.completedStages.includes("fk_review")) { + const review = this.state(scope.runtime.workspaceId).readFkReview(scope.job.runId); + if (!review || review.reviewedCandidatesDigest !== candidate.digest) { + return baseResult(scope.runtime, "preprocess run", "blocked", "manual_review_required", { + runId: scope.job.runId, + childRuns: { ...scope.job.childRuns }, + completedStages: [...scope.job.completedStages], + artifactIdentities: [{ kind: "fk_candidates", digest: candidate.digest }], + }); + } + scope.job.completedStages.push("fk_review"); + this.state(scope.runtime.workspaceId).writeJob(scope.job); + } + const semantic = await this.deps.semanticPreflight(scope.runtime.workspace); + if (!semantic.ok) return baseResult(scope.runtime, "preprocess run", "failed", semantic.code, { runId: scope.job.runId }); + if (!scope.job.completedStages.includes("schema_index")) { + const payload = await this.runJsonStage(scope.runtime, ["vector", "index-schema", "--json", "-c", "/dev/fd/3"]); + scope.job.completedStages.push("schema_index"); + this.state(scope.runtime.workspaceId).writeJob(scope.job); + const warnings = noEvidenceWarning(scope.runtime.workspace); + if (scope.runtime.workspace.evidence === undefined) { + return baseResult(scope.runtime, "preprocess run", "succeeded", "ok", { + runId: scope.job.runId, + childRuns: { ...scope.job.childRuns }, + completedStages: [...scope.job.completedStages], + counts: this.numberRecord(payload.counts), + warnings, + }); + } + } + if (scope.runtime.workspace.evidence === undefined) { + return baseResult(scope.runtime, "preprocess run", "succeeded", "ok", { + runId: scope.job.runId, + childRuns: { ...scope.job.childRuns }, + completedStages: [...scope.job.completedStages], + warnings: noEvidenceWarning(scope.runtime.workspace), + }); + } + const policy = this.evidencePolicy(scope.runtime.workspace); + if (policy !== undefined) { + return baseResult(scope.runtime, "preprocess run", policy.status, policy.code, { + runId: scope.job.runId, + childRuns: { ...scope.job.childRuns }, + completedStages: [...scope.job.completedStages], + warnings: policy.warnings, + }); + } + if (!scope.job.completedStages.includes("evidence")) { + const payload = await this.runJsonStage(scope.runtime, [ + "preprocess", "evidence", + ...(scope.job.childRuns.evidence ? ["--resume", scope.job.childRuns.evidence] : []), + "--json", "-c", "/dev/fd/3", + ]); + if (typeof payload.run_id === "string") scope.job.childRuns.evidence = this.requireRunId(payload.run_id); + scope.job.completedStages.push("evidence"); + this.state(scope.runtime.workspaceId).writeJob(scope.job); + return baseResult(scope.runtime, "preprocess run", "succeeded", "ok", { + runId: scope.job.runId, + childRuns: { ...scope.job.childRuns }, + completedStages: [...scope.job.completedStages], + counts: this.numberRecord(payload.counts), + }); + } + return baseResult(scope.runtime, "preprocess run", "unchanged", "ok", { + runId: scope.job.runId, + childRuns: { ...scope.job.childRuns }, + completedStages: [...scope.job.completedStages], + }); + } + + private async startRun(workspaceId: string, operation: string, resumeRunId?: string): Promise { + const runtime = await this.deps.acquireActiveRuntime(workspaceId); + const state = this.state(runtime.workspaceId); + await state.assertSessionInventoryCompatible(runtime.workspaceRevision, await this.deps.listSessions(runtime.workspaceId)); + const job = await state.beginJob({ + operation, + runId: resumeRunId, + workspaceRevision: runtime.workspaceRevision, + descriptorBlob: runtime.descriptorBlob, + catalogBlob: runtime.catalogBlob, + configDigest: runtime.configLease.configDigest, + bindingDigest: runtime.configLease.bindingDigest, + }); + return { runtime, state, job }; + } + + private state(workspaceId: string): PreprocessingStateStore { + return new PreprocessingStateStore({ dataRoot: this.deps.dataRoot, workspaceId }); + } + + private async runSuggestStage( + scope: RunScope, + fromSql: ReadonlyArray<{ name: string; sql: string }>, + assume: readonly string[], + ): Promise { + const payload = await this.withStagedInputs(scope.runtime.workspaceId, fromSql.map((entry) => ({ + flag: "--from-sql", + name: entry.name, + contents: entry.sql, + })), async (stagedArgv) => await this.runJsonStage(scope.runtime, [ + "schema", "suggest-fks", ...stagedArgv, + ...assume.flatMap((value) => ["--assume", value]), + "--json", "-c", "/dev/fd/3", + ])); + const candidateCount = Number(payload.candidate_count ?? 0); + const candidateYaml = typeof payload.candidate_yaml === "string" ? payload.candidate_yaml : ""; + let artifactIdentities: Array<{ kind: string; digest: string }> | undefined; + if (candidateCount > 0) { + const persisted = this.state(scope.runtime.workspaceId).writeFkCandidates(scope.job.runId, candidateYaml); + scope.job.candidateDigest = persisted.digest; + artifactIdentities = [{ kind: "fk_candidates", digest: persisted.digest }]; + } + if (!scope.job.completedStages.includes("fk_suggest")) scope.job.completedStages.push("fk_suggest"); + this.state(scope.runtime.workspaceId).writeJob(scope.job); + const resultExtra = { + runId: scope.job.runId, + completedStages: [...scope.job.completedStages], + ...(artifactIdentities ? { artifactIdentities } : {}), + ...(candidateYaml.length > 0 ? { suggestedFksYaml: candidateYaml } : {}), + }; + if (candidateCount > 0) { + return baseResult(scope.runtime, scope.job.operation === "preprocess run" ? "preprocess run" : "schema suggest-fks", "blocked", "manual_review_required", { + ...resultExtra, + childRuns: { ...scope.job.childRuns }, + }); + } + return baseResult(scope.runtime, scope.job.operation === "preprocess run" ? "preprocess run" : "schema suggest-fks", "succeeded", "ok", resultExtra); + } + + private async runJsonStage(runtime: ActiveRuntime, argv: string[]): Promise> { + const result = await this.deps.runChild({ argv, configPath: runtime.configLease.path }); + if (result.exitCode !== 0) throw new Error("workspace child failed"); + return JSON.parse(result.stdout) as Record; + } + + private requireRunId(value: unknown): string { + if (typeof value !== "string" || !/^[0-9a-f]{32}$/.test(value)) throw new Error("child run id is invalid"); + return value; + } + + private numberRecord(value: unknown): Record | undefined { + if (!value || typeof value !== "object" || Array.isArray(value)) return undefined; + return Object.fromEntries(Object.entries(value as Record).map(([key, nested]) => [key, Number(nested)])); + } + + private async withStagedInputs( + workspaceId: string, + inputs: ReadonlyArray<{ flag: string; name: string; contents: string }>, + fn: (argv: string[]) => Promise, + ): Promise { + if (inputs.length === 0) return await fn([]); + const root = join(this.deps.dataRoot, "sessions", workspaceId, "preprocessing", `.stage-${randomBytes(6).toString("hex")}`); + mkdirSync(root, { recursive: true, mode: 0o700 }); + const argv: string[] = []; + const paths: string[] = []; + try { + for (const input of inputs) { + const path = join(root, input.name); + writeFileSync(path, input.contents, { encoding: "utf8", flag: "wx", mode: 0o600 }); + paths.push(path); + argv.push(input.flag, path); + } + return await fn(argv); + } finally { + rmSync(root, { recursive: true, force: true }); + } + } + + private findRunIdByCandidateDigest(state: PreprocessingStateStore, digestValue: string): string | undefined { + for (const entry of readdirSync(state.fkCandidatesDirectory(), { withFileTypes: true })) { + if (!entry.isFile() || entry.isSymbolicLink() || !/^[0-9a-f]{32}\.yaml$/.test(entry.name)) continue; + const runId = entry.name.slice(0, -".yaml".length); + if (state.readFkCandidates(runId)?.digest === digestValue) return runId; + } + return undefined; + } + + private evidencePolicy(workspace: WorkspaceDescriptor): { + status: WorkspaceOperationResult["status"]; + code: WorkspaceOperationResult["code"]; + warnings?: string[]; + } | undefined { + const evidence = workspace.evidence; + if (!evidence) return undefined; + if (evidence.source.type === "filesystem") { + return { status: "blocked", code: "evidence_materialization_required" }; + } + if (evidence.source.type === "http") { + for (const value of evidence.source.uris) { + const host = new URL(value).hostname; + if (isPrivateHost(host) && !(evidence.source.allow_private_hosts && this.deps.httpPrivateHostAllowlist?.includes(host))) { + return { status: "failed", code: "egress_policy_refused" }; + } + } + return undefined; + } + if ( + evidence.source.endpoint_url !== undefined + || evidence.source.credentials === "ambient" + || evidence.source.allow_private_endpoint + || evidence.source.allow_insecure_endpoint + ) { + return { status: "failed", code: "egress_policy_refused" }; + } + return undefined; + } +} diff --git a/backend/src/workspaces/preprocessing-state.ts b/backend/src/workspaces/preprocessing-state.ts new file mode 100644 index 00000000..95ee6cd3 --- /dev/null +++ b/backend/src/workspaces/preprocessing-state.ts @@ -0,0 +1,383 @@ +import { createHash, randomBytes } from "node:crypto"; +import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process"; +import { + closeSync, + constants as fsConstants, + fchmodSync, + fstatSync, + fsyncSync, + lstatSync, + mkdirSync, + openSync, + readFileSync, + renameSync, + unlinkSync, + writeFileSync, +} from "node:fs"; +import { dirname, join } from "node:path"; + +export type PreprocessingConflictCode = "preprocessing_conflict" | "preprocessing_resume_mismatch"; + +export class PreprocessingStateError extends Error { + constructor(readonly code: PreprocessingConflictCode, message: string) { + super(message); + this.name = "PreprocessingStateError"; + } +} + +export interface SessionInventoryRow { + id: string; + status: string; + archived?: boolean; + workspaceRevision?: string | null; +} + +export interface WriterLockLease { + holderPid: number; + release(): Promise; +} + +export interface BeginPreprocessingJobOptions { + operation: string; + workspaceRevision: string; + descriptorBlob: string; + catalogBlob: string; + configDigest: string; + bindingDigest: string; + runId?: string; +} + +export interface PreprocessingJobState { + schemaVersion: 1; + runId: string; + operation: string; + workspaceId: string; + workspaceRevision: string; + descriptorBlob: string; + catalogBlob: string; + configDigest: string; + bindingDigest: string; + completedStages: string[]; + childRuns: Record; + status: "active" | "succeeded" | "blocked" | "failed"; + candidateDigest?: string; + reviewDigest?: string; +} + +export interface FkReviewRecord { + reviewedCandidatesDigest: string; + annotationsDigest: string; + workspaceRevision: string; +} + +function sha256(value: string | Buffer): string { + return `sha256:${createHash("sha256").update(value).digest("hex")}`; +} + +function syncDirectory(directory: string): void { + if (process.platform === "win32") return; + const fd = openSync(directory, "r"); + try { fsyncSync(fd); } finally { closeSync(fd); } +} + +function ensureDirectory(path: string): string { + mkdirSync(path, { recursive: true, mode: 0o700 }); + const entry = lstatSync(path); + if (!entry.isDirectory() || entry.isSymbolicLink()) { + throw new Error("preprocessing state directory is unavailable"); + } + return path; +} + +function validateRunId(runId: string): string { + if (!/^[0-9a-f]{32}$/.test(runId)) throw new Error("preprocessing run id is invalid"); + return runId; +} + +function writeAtomicFile(path: string, contents: string, mode: number): void { + ensureDirectory(join(path, "..")); + const directory = path.slice(0, path.lastIndexOf("/")); + ensureDirectory(directory); + 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, "utf8"); + fsyncSync(fd); + fchmodSync(fd, mode); + closeSync(fd); + closed = true; + renameSync(staging, path); + syncDirectory(directory); + } catch (error) { + if (!closed) try { closeSync(fd); } catch { /* preserve original failure */ } + try { unlinkSync(staging); } catch { /* best effort */ } + throw error; + } +} + +function readTrustedFile(path: string): string { + const entry = lstatSync(path); + if (!entry.isFile() || entry.isSymbolicLink()) throw new Error("preprocessing state 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("preprocessing state file is invalid"); + const contents = readFileSync(fd, "utf8"); + const after = fstatSync(fd); + if ( + before.dev !== after.dev || before.ino !== after.ino || before.size !== after.size || before.nlink !== after.nlink + ) throw new Error("preprocessing state file changed while reading"); + return contents; + } finally { + closeSync(fd); + } +} + +function decodeJob(value: unknown): PreprocessingJobState { + if (!value || typeof value !== "object" || Array.isArray(value)) { + throw new Error("preprocessing job state is invalid"); + } + const record = value as Record; + if ( + record.schemaVersion !== 1 + || typeof record.runId !== "string" + || typeof record.operation !== "string" + || typeof record.workspaceId !== "string" + || typeof record.workspaceRevision !== "string" + || typeof record.descriptorBlob !== "string" + || typeof record.catalogBlob !== "string" + || typeof record.configDigest !== "string" + || typeof record.bindingDigest !== "string" + || !Array.isArray(record.completedStages) + || typeof record.childRuns !== "object" || record.childRuns === null || Array.isArray(record.childRuns) + || !["active", "succeeded", "blocked", "failed"].includes(String(record.status)) + ) { + throw new Error("preprocessing job state is invalid"); + } + return record as unknown as PreprocessingJobState; +} + +function decodeReview(value: unknown): FkReviewRecord { + if (!value || typeof value !== "object" || Array.isArray(value)) { + throw new Error("preprocessing review state is invalid"); + } + const record = value as Record; + if ( + typeof record.reviewedCandidatesDigest !== "string" + || typeof record.annotationsDigest !== "string" + || typeof record.workspaceRevision !== "string" + ) throw new Error("preprocessing review state is invalid"); + return record as unknown as FkReviewRecord; +} + +export class PreprocessingStateStore { + private readonly root: string; + + constructor(private readonly options: { dataRoot: string; workspaceId: string }) { + if (!/^[a-z][a-z0-9-]{2,62}$/.test(options.workspaceId)) { + throw new Error("workspace id is invalid"); + } + this.root = join(options.dataRoot, "sessions", options.workspaceId, "preprocessing"); + } + + writerLockPath(): string { return join(this.root, "writer.lock"); } + runtimeConfigDirectory(): string { return join(this.root, "runtime-config"); } + runtimeConfigManifestDirectory(): string { return join(this.root, "runtime-config-manifests"); } + jobsDirectory(): string { return join(this.root, "jobs"); } + fkCandidatesDirectory(): string { return join(this.root, "fk-candidates"); } + fkReviewsDirectory(): string { return join(this.root, "fk-reviews"); } + jobPath(runId: string): string { return join(this.jobsDirectory(), `${validateRunId(runId)}.json`); } + fkCandidatesPath(runId: string): string { return join(this.fkCandidatesDirectory(), `${validateRunId(runId)}.yaml`); } + fkReviewPath(runId: string): string { return join(this.fkReviewsDirectory(), `${validateRunId(runId)}.json`); } + + private ensureLayout(): void { + ensureDirectory(this.root); + for (const directory of [ + this.runtimeConfigDirectory(), + this.runtimeConfigManifestDirectory(), + this.jobsDirectory(), + this.fkCandidatesDirectory(), + this.fkReviewsDirectory(), + ]) ensureDirectory(directory); + } + + async acquireWriterLock(): Promise { + this.ensureLayout(); + const lockPath = this.writerLockPath(); + try { + const entry = lstatSync(lockPath); + if (!entry.isFile() || entry.isSymbolicLink()) throw new Error("invalid writer lock path"); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") { + throw new PreprocessingStateError("preprocessing_conflict", "Workspace preprocessing lock is unavailable"); + } + } + const holder = spawn("python3", ["-c", PreprocessingStateStore.HOLDER_PROGRAM, lockPath], { + stdio: ["pipe", "pipe", "pipe"], + }); + await new Promise((resolve, reject) => { + let output = ""; + const fail = (error: Error) => { + holder.stdout.removeAllListeners("data"); + reject(error); + }; + holder.once("error", () => fail(new PreprocessingStateError( + "preprocessing_conflict", + "Workspace preprocessing lock is unavailable", + ))); + holder.once("exit", (code) => { + fail(new PreprocessingStateError( + "preprocessing_conflict", + code === 73 ? "Workspace preprocessing is busy" : "Workspace preprocessing lock is unavailable", + )); + }); + holder.stdout.on("data", (chunk: Buffer) => { + output += chunk.toString("utf8"); + if (output === "locked\n") { + holder.stdout.removeAllListeners("data"); + resolve(); + } + }); + }); + let released = false; + return { + holderPid: holder.pid ?? 0, + release: async () => { + if (released) return; + released = true; + if (!holder.stdin.destroyed) holder.stdin.end(); + await new Promise((resolve) => holder.once("exit", () => resolve())); + }, + }; + } + + async beginJob(options: BeginPreprocessingJobOptions): Promise { + this.ensureLayout(); + if (!/^[0-9a-f]{40}$/.test(options.workspaceRevision)) { + throw new Error("workspace revision is invalid"); + } + const runId = options.runId ?? randomBytes(16).toString("hex"); + const path = this.jobPath(runId); + try { + const existing = this.readJob(runId); + if ( + existing.operation !== options.operation + || existing.workspaceRevision !== options.workspaceRevision + || existing.descriptorBlob !== options.descriptorBlob + || existing.catalogBlob !== options.catalogBlob + || existing.configDigest !== options.configDigest + || existing.bindingDigest !== options.bindingDigest + ) { + throw new PreprocessingStateError( + "preprocessing_resume_mismatch", + "Workspace preprocessing resume no longer matches the pinned revision", + ); + } + return existing; + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") { + if (error instanceof PreprocessingStateError) throw error; + if (options.runId) throw error; + } + } + const job: PreprocessingJobState = { + schemaVersion: 1, + runId, + operation: options.operation, + workspaceId: this.options.workspaceId, + workspaceRevision: options.workspaceRevision, + descriptorBlob: options.descriptorBlob, + catalogBlob: options.catalogBlob, + configDigest: options.configDigest, + bindingDigest: options.bindingDigest, + completedStages: [], + childRuns: {}, + status: "active", + }; + writeAtomicFile(path, `${JSON.stringify(job)} +`, 0o600); + return job; + } + + readJob(runId: string): PreprocessingJobState { + return decodeJob(JSON.parse(readTrustedFile(this.jobPath(runId)))); + } + + writeJob(job: PreprocessingJobState): PreprocessingJobState { + this.ensureLayout(); + writeAtomicFile(this.jobPath(job.runId), `${JSON.stringify(job)} +`, 0o600); + return job; + } + + writeFkCandidates(runId: string, yaml: string): { path: string; digest: string } { + this.ensureLayout(); + const path = this.fkCandidatesPath(runId); + writeAtomicFile(path, yaml, 0o600); + return { path, digest: sha256(yaml) }; + } + + readFkCandidates(runId: string): { path: string; yaml: string; digest: string } | undefined { + const path = this.fkCandidatesPath(runId); + try { + const yaml = readTrustedFile(path); + return { path, yaml, digest: sha256(yaml) }; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined; + throw error; + } + } + + writeFkReview(runId: string, review: FkReviewRecord): { path: string; digest: string } { + this.ensureLayout(); + const path = this.fkReviewPath(runId); + const json = `${JSON.stringify(review)} +`; + writeAtomicFile(path, json, 0o600); + return { path, digest: sha256(json) }; + } + + readFkReview(runId: string): FkReviewRecord | undefined { + try { + return decodeReview(JSON.parse(readTrustedFile(this.fkReviewPath(runId)))); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined; + throw error; + } + } + + async assertSessionInventoryCompatible( + workspaceRevision: string, + sessions: readonly SessionInventoryRow[], + ): Promise { + const conflicting = sessions.find((session) => ( + session.workspaceRevision + && session.workspaceRevision !== workspaceRevision + && session.status !== "finalized" + && !session.archived + )); + if (conflicting) { + throw new PreprocessingStateError( + "preprocessing_conflict", + "A resumable session is pinned to a different workspace revision", + ); + } + } + + private static readonly HOLDER_PROGRAM = [ + "import fcntl, os, sys", + "fd = os.open(sys.argv[1], os.O_RDWR | os.O_CREAT | getattr(os, 'O_NOFOLLOW', 0), 0o600)", + "try:", + " fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)", + "except BlockingIOError:", + " sys.exit(73)", + "sys.stdout.write('locked\\n')", + "sys.stdout.flush()", + "sys.stdin.buffer.read()", + ].join("\n"); +} diff --git a/backend/src/workspaces/runtime-config-lease.ts b/backend/src/workspaces/runtime-config-lease.ts new file mode 100644 index 00000000..6bb79776 --- /dev/null +++ b/backend/src/workspaces/runtime-config-lease.ts @@ -0,0 +1,522 @@ +import { createHash, randomUUID } from "node:crypto"; +import { + closeSync, + constants as fsConstants, + fchmodSync, + fstatSync, + fsyncSync, + lstatSync, + mkdirSync, + openSync, + readFileSync, + readSync, + renameSync, + statSync, + unlinkSync, + writeFileSync, +} from "node:fs"; +import { dirname, isAbsolute, join, relative, resolve } from "node:path"; +import { parse, parseAllDocuments, stringify } from "yaml"; +import { resolveRuntimeBindings, type RuntimeBindings } from "./bindings.js"; +import { GitWorkspaceRepository } from "./git-repository.js"; +import { WorkspaceRegistry } from "./registry.js"; +import { + renderRuntimeConfig, + type RuntimeIdentity, + type RuntimeInstallationOverlay, + type RuntimePaths, + type RuntimeRenderContext, + type SemanticRuntimeConfig, +} from "./runtime-renderer.js"; +import { parseWorkspaceYaml, validateOperationalWorkspace, type WorkspaceDescriptor } from "./schema.js"; +import type { WorkspaceRegistryConfig } from "./types.js"; + +export interface RuntimeConfigLease { + path: string; + workspaceId: string; + workspaceRevision: string; + release(): void; +} + +export interface RenderedWorkspaceRuntime { + workspace: WorkspaceDescriptor; + workspaceId: string; + workspaceRevision: string; + revisionContentRoot: string; + runtimePaths: RuntimePaths; + installationOverlay: RuntimeInstallationOverlay; + bindings: RuntimeBindings; + bindingDigest: string; + renderedConfig: string; +} + +export interface ActiveRenderedWorkspaceRuntime extends RenderedWorkspaceRuntime { + snapshotPath: string; + descriptorBlob: string; + catalogBlob: string; +} + +export interface DeterministicRuntimeConfigLease extends RuntimeConfigLease { + manifestPath: string; + descriptorBlob: string; + catalogBlob: string; + configDigest: string; + bindingDigest: string; +} + +export class RuntimeConfigLeaseError extends Error { + constructor(readonly code: "effective_config_mismatch", message: string) { + super(message); + this.name = "RuntimeConfigLeaseError"; + } +} + +interface SnapshotIdentity extends RuntimeIdentity { + snapshotPath: string; +} + +interface PublishedRuntimeConfigManifest { + schemaVersion: 1; + workspaceId: string; + workspaceRevision: string; + descriptorBlob: string; + catalogBlob: string; + configDigest: string; + bindingDigest: string; + path: string; + file: { + dev: number; + ino: number; + size: number; + mode: number; + nlink: number; + }; +} + +function sha256(value: string | Buffer): string { + return `sha256:${createHash("sha256").update(value).digest("hex")}`; +} + +function stableBindingDigest(bindings: RuntimeBindings): string { + const encodeRecord = (value: Record) => Object.entries(value).sort(([left], [right]) => ( + left.localeCompare(right) + )); + return sha256(JSON.stringify({ + dwh: { + transport: bindings.dwh.transport, + values: encodeRecord(bindings.dwh.values), + missing: [...bindings.dwh.missing].sort(), + }, + evidence: { + values: encodeRecord(bindings.evidence.values), + missing: [...bindings.evidence.missing].sort(), + }, + })); +} + +function syncDirectory(directory: string): void { + if (process.platform === "win32") return; + const fd = openSync(directory, "r"); + try { fsyncSync(fd); } finally { closeSync(fd); } +} + +function ensureTrustedDirectory(directory: string): string { + mkdirSync(directory, { recursive: true, mode: 0o700 }); + const entry = lstatSync(directory); + if (!entry.isDirectory() || entry.isSymbolicLink()) { + throw new Error("runtime configuration directory is unavailable"); + } + return directory; +} + +function writeAtomicFile(path: string, contents: string, mode: number): void { + ensureTrustedDirectory(dirname(path)); + const staging = `${path}.tmp-${process.pid}-${Date.now()}-${randomUUID()}`; + const fd = openSync( + staging, + fsConstants.O_WRONLY | fsConstants.O_CREAT | fsConstants.O_EXCL | fsConstants.O_NOFOLLOW, + 0o600, + ); + let closed = false; + try { + writeFileSync(fd, contents, "utf8"); + 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 error */ } + try { unlinkSync(staging); } catch { /* best effort */ } + throw error; + } +} + +function readTrustedFile(path: string): { contents: string; stat: ReturnType } { + const entry = lstatSync(path); + if (!entry.isFile() || entry.isSymbolicLink()) { + throw new Error("runtime configuration file is untrusted"); + } + const fd = openSync(path, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW); + try { + const before = fstatSync(fd); + if (!before.isFile() || before.nlink !== 1) { + throw new Error("runtime configuration file is untrusted"); + } + const contents = readFileSync(fd, "utf8"); + const after = fstatSync(fd); + if ( + before.dev !== after.dev || before.ino !== after.ino || before.size !== after.size || before.nlink !== after.nlink + ) { + throw new Error("runtime configuration file changed while reading"); + } + return { contents, stat: statSync(path) }; + } finally { + closeSync(fd); + } +} + +function readStrictJson(path: string): unknown { + return JSON.parse(readTrustedFile(path).contents); +} + +function trustedSnapshotIdentity(snapshotPath: string): SnapshotIdentity { + if (!isAbsolute(snapshotPath)) throw new Error("config path is not a trusted runtime snapshot"); + const commitDirectory = dirname(snapshotPath); + const snapshotsRoot = dirname(commitDirectory); + const pathRelative = relative(snapshotsRoot, snapshotPath); + const match = /^([0-9a-f]{40})\/([a-z][a-z0-9-]{2,62})\.yaml$/.exec(pathRelative); + if (pathRelative.startsWith("..") || isAbsolute(pathRelative) || !match) { + throw new Error("config path is not a trusted runtime snapshot"); + } + const entry = lstatSync(snapshotPath); + if (!entry.isFile() || entry.isSymbolicLink()) { + throw new Error("config path is not a trusted runtime snapshot"); + } + return { snapshotPath, workspaceRevision: match[1], workspaceId: match[2] }; +} + +function readSnapshotWorkspace(snapshotPath: string): { + workspace: WorkspaceDescriptor; + identity: SnapshotIdentity; + revisionContentRoot: string; +} { + const identity = trustedSnapshotIdentity(snapshotPath); + const fd = openSync(snapshotPath, fsConstants.O_RDONLY | fsConstants.O_NOFOLLOW); + try { + const before = fstatSync(fd); + if (!before.isFile() || before.nlink !== 1) throw new Error("workspace snapshot is not a file"); + const source = readFileSync(fd, "utf8"); + const after = fstatSync(fd); + if (before.dev !== after.dev || before.ino !== after.ino || before.size !== after.size) { + throw new Error("workspace snapshot changed while reading"); + } + const workspace = validateOperationalWorkspace(parseWorkspaceYaml(source)); + if (workspace.workspace.id !== identity.workspaceId) { + throw new Error("workspace snapshot identity does not match its path"); + } + return { workspace, identity, revisionContentRoot: dirname(snapshotPath) }; + } finally { + closeSync(fd); + } +} + +function runtimePaths(dataRoot: string, workspaceId: string): RuntimePaths { + if (!isAbsolute(dataRoot)) throw new Error("registry workspace runtime requires an absolute data root"); + const root = join(dataRoot, "sessions", workspaceId); + return { + sessions: join(root, "sessions"), + artifacts: join(root, "artifacts"), + indexes: join(root, "indexes"), + }; +} + +function installationOverlay(harnessDir: string, configPath: string): RuntimeInstallationOverlay { + const path = isAbsolute(configPath) ? configPath : resolve(harnessDir, configPath); + try { + const documents = parseAllDocuments(readFileSync(path, "utf8"), { uniqueKeys: true }); + if (documents.length !== 1) throw new Error("installation config must contain one YAML document"); + const document = documents[0]; + if (document.errors.length > 0 || document.warnings.length > 0) { + throw new Error("installation config contains invalid YAML"); + } + const parsed = document.toJSON(); + if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) { + throw new Error("installation config must be a YAML mapping"); + } + const source = parsed as Record; + return { + ...(source.session_storage === undefined ? {} : { session_storage: source.session_storage }), + ...(source.profile === undefined ? {} : { profile: source.profile }), + }; + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return {}; + throw error; + } +} + +function applyCollectionLifecycle(config: string, lifecycle: "require_existing"): string { + const parsed = parse(config) as Record; + if (!parsed.resources || typeof parsed.resources !== "object") { + throw new Error("runtime configuration is missing resources"); + } + parsed.resources = { ...parsed.resources }; + parsed.resources.vector = { ...(parsed.resources.vector ?? {}), collection_lifecycle: lifecycle }; + return stringify(parsed, { lineWidth: 0, sortMapEntries: false }); +} + +function renderWorkspaceRuntimeFromWorkspace(options: { + workspace: WorkspaceDescriptor; + workspaceId: string; + workspaceRevision: string; + revisionContentRoot: string; + harnessDir: string; + configPath: string; + dataRoot: string; + secretRoots: readonly string[]; + semanticRuntime: SemanticRuntimeConfig; +}): RenderedWorkspaceRuntime { + const bindings = resolveRuntimeBindings(options.workspace, process.env, options.secretRoots); + const overlay = installationOverlay(options.harnessDir, options.configPath); + const context: RuntimeRenderContext = { + workspaceId: options.workspaceId, + workspaceRevision: options.workspaceRevision, + revisionContentRoot: options.revisionContentRoot, + }; + return { + workspace: options.workspace, + workspaceId: options.workspaceId, + workspaceRevision: options.workspaceRevision, + revisionContentRoot: options.revisionContentRoot, + runtimePaths: runtimePaths(options.dataRoot, options.workspaceId), + installationOverlay: overlay, + bindings, + bindingDigest: stableBindingDigest(bindings), + renderedConfig: renderRuntimeConfig( + options.workspace, + bindings, + runtimePaths(options.dataRoot, options.workspaceId), + context, + overlay, + options.semanticRuntime, + ), + }; +} + +export function renderWorkspaceRuntimeFromSnapshotPath(options: { + snapshotPath: string; + harnessDir: string; + configPath: string; + dataRoot: string; + secretRoots: readonly string[]; + semanticRuntime: SemanticRuntimeConfig; +}): RenderedWorkspaceRuntime { + const snapshot = readSnapshotWorkspace(options.snapshotPath); + return renderWorkspaceRuntimeFromWorkspace({ + workspace: snapshot.workspace, + workspaceId: snapshot.identity.workspaceId, + workspaceRevision: snapshot.identity.workspaceRevision, + revisionContentRoot: snapshot.revisionContentRoot, + harnessDir: options.harnessDir, + configPath: options.configPath, + dataRoot: options.dataRoot, + secretRoots: options.secretRoots, + semanticRuntime: options.semanticRuntime, + }); +} + +export async function renderActiveWorkspaceRuntime(options: { + workspaceId: string; + registry: WorkspaceRegistry; + registryConfig: WorkspaceRegistryConfig; + harnessDir: string; + configPath: string; + dataRoot: string; + secretRoots: readonly string[]; + semanticRuntime: SemanticRuntimeConfig; +}): Promise { + const { workspace, revision } = await options.registry.read(options.workspaceId); + const repository = new GitWorkspaceRepository(options.registryConfig); + await repository.ensureLayout(); + const rendered = renderWorkspaceRuntimeFromWorkspace({ + workspace, + workspaceId: revision.id, + workspaceRevision: revision.commit, + revisionContentRoot: dirname(revision.snapshotPath), + harnessDir: options.harnessDir, + configPath: options.configPath, + dataRoot: options.dataRoot, + secretRoots: options.secretRoots, + semanticRuntime: options.semanticRuntime, + }); + return { + ...rendered, + snapshotPath: revision.snapshotPath, + descriptorBlob: revision.blob, + catalogBlob: (await repository.catalogBlob(revision.commit)).trim(), + }; +} + +function decodePublishedRuntimeConfigManifest(value: unknown): PublishedRuntimeConfigManifest { + if (!value || typeof value !== "object" || Array.isArray(value)) { + throw new RuntimeConfigLeaseError("effective_config_mismatch", "runtime configuration manifest is invalid"); + } + const manifest = value as Record; + const file = manifest.file as Record | undefined; + if ( + manifest.schemaVersion !== 1 + || typeof manifest.workspaceId !== "string" + || typeof manifest.workspaceRevision !== "string" + || typeof manifest.descriptorBlob !== "string" + || typeof manifest.catalogBlob !== "string" + || typeof manifest.configDigest !== "string" + || typeof manifest.bindingDigest !== "string" + || typeof manifest.path !== "string" + || !file + || typeof file.dev !== "number" + || typeof file.ino !== "number" + || typeof file.size !== "number" + || typeof file.mode !== "number" + || typeof file.nlink !== "number" + ) { + throw new RuntimeConfigLeaseError("effective_config_mismatch", "runtime configuration manifest is invalid"); + } + return manifest as unknown as PublishedRuntimeConfigManifest; +} + +export async function publishDeterministicRuntimeConfigLease(options: { + workspaceId: string; + registry: WorkspaceRegistry; + registryConfig: WorkspaceRegistryConfig; + harnessDir: string; + configPath: string; + dataRoot: string; + secretRoots: readonly string[]; + semanticRuntime: SemanticRuntimeConfig; +}): Promise { + const rendered = await renderActiveWorkspaceRuntime(options); + const publishedConfig = applyCollectionLifecycle(rendered.renderedConfig, "require_existing"); + const preprocessingRoot = ensureTrustedDirectory(join( + options.dataRoot, + "sessions", + rendered.workspaceId, + "preprocessing", + )); + const configDirectory = ensureTrustedDirectory(join(preprocessingRoot, "runtime-config")); + const manifestDirectory = ensureTrustedDirectory(join(preprocessingRoot, "runtime-config-manifests")); + const path = join(configDirectory, `${rendered.workspaceRevision}.yaml`); + const manifestPath = join(manifestDirectory, `${rendered.workspaceRevision}.json`); + const configDigest = sha256(publishedConfig); + + const verifyPublished = (): PublishedRuntimeConfigManifest | undefined => { + let manifest: PublishedRuntimeConfigManifest | undefined; + try { + manifest = decodePublishedRuntimeConfigManifest(readStrictJson(manifestPath)); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; + } + let configStat: ReturnType | undefined; + let contents: string | undefined; + try { + const file = readTrustedFile(path); + contents = file.contents; + configStat = file.stat; + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; + } + if (contents === undefined) return manifest; + if ((Number(configStat!.mode) & 0o777) !== 0o400 || configStat!.nlink !== 1) { + throw new RuntimeConfigLeaseError("effective_config_mismatch", "runtime configuration lease changed"); + } + if (sha256(contents) !== configDigest) { + throw new RuntimeConfigLeaseError("effective_config_mismatch", "runtime configuration lease changed"); + } + if (!manifest) return undefined; + if ( + manifest.workspaceId !== rendered.workspaceId + || manifest.workspaceRevision !== rendered.workspaceRevision + || manifest.descriptorBlob !== rendered.descriptorBlob + || manifest.catalogBlob !== rendered.catalogBlob + || manifest.configDigest !== configDigest + || manifest.bindingDigest !== rendered.bindingDigest + || manifest.path !== path + || manifest.file.dev !== configStat!.dev + || manifest.file.ino !== configStat!.ino + || manifest.file.size !== configStat!.size + || manifest.file.mode !== (Number(configStat!.mode) & 0o777) + || manifest.file.nlink !== configStat!.nlink + ) { + throw new RuntimeConfigLeaseError("effective_config_mismatch", "runtime configuration lease changed"); + } + return manifest; + }; + + const existing = verifyPublished(); + if (existing) { + return { + path, + manifestPath, + workspaceId: rendered.workspaceId, + workspaceRevision: rendered.workspaceRevision, + descriptorBlob: rendered.descriptorBlob, + catalogBlob: rendered.catalogBlob, + configDigest, + bindingDigest: rendered.bindingDigest, + release: () => undefined, + }; + } + + try { + readTrustedFile(path); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") { + writeAtomicFile(path, publishedConfig, 0o400); + } else { + throw error; + } + } + const published = readTrustedFile(path); + const publishedStat = published.stat!; + if (sha256(published.contents) !== configDigest) { + throw new RuntimeConfigLeaseError("effective_config_mismatch", "runtime configuration lease changed"); + } + const manifest: PublishedRuntimeConfigManifest = { + schemaVersion: 1, + workspaceId: rendered.workspaceId, + workspaceRevision: rendered.workspaceRevision, + descriptorBlob: rendered.descriptorBlob, + catalogBlob: rendered.catalogBlob, + configDigest, + bindingDigest: rendered.bindingDigest, + path, + file: { + dev: Number(publishedStat.dev), + ino: Number(publishedStat.ino), + size: Number(publishedStat.size), + mode: Number(publishedStat.mode) & 0o777, + nlink: Number(publishedStat.nlink), + }, + }; + try { + readStrictJson(manifestPath); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") { + writeAtomicFile(manifestPath, `${JSON.stringify(manifest)}\n`, 0o600); + } else { + throw error; + } + } + verifyPublished(); + return { + path, + manifestPath, + workspaceId: rendered.workspaceId, + workspaceRevision: rendered.workspaceRevision, + descriptorBlob: rendered.descriptorBlob, + catalogBlob: rendered.catalogBlob, + configDigest, + bindingDigest: rendered.bindingDigest, + release: () => undefined, + }; +} diff --git a/backend/test/tht-runner.test.ts b/backend/test/tht-runner.test.ts index 62c58ea2..848743b9 100644 --- a/backend/test/tht-runner.test.ts +++ b/backend/test/tht-runner.test.ts @@ -10,18 +10,22 @@ import { ThtRunner } from "../src/tht/tht-runner.js"; // Spy on child_process.spawn so we can capture the resolved argv (incl. -c config) // that ThtRunner.run() builds, without launching a real process. -vi.mock("node:child_process", () => ({ - spawn: vi.fn(() => { - const ch: any = new EventEmitter(); - ch.stdout = new EventEmitter(); - ch.stderr = new EventEmitter(); - queueMicrotask(() => { - ch.stdout.emit("data", Buffer.from('{"id":"x"}')); - ch.emit("close", 0); - }); - return ch; - }), -})); +vi.mock("node:child_process", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + spawn: vi.fn(() => { + const ch: any = new EventEmitter(); + ch.stdout = new EventEmitter(); + ch.stderr = new EventEmitter(); + queueMicrotask(() => { + ch.stdout.emit("data", Buffer.from('{"id":"x"}')); + ch.emit("close", 0); + }); + return ch; + }), + }; +}); import { spawn } from "node:child_process"; test("sessionNew parses id from JSON", async () => { diff --git a/backend/test/workspace-maintenance.test.ts b/backend/test/workspace-maintenance.test.ts new file mode 100644 index 00000000..6c8ecfee --- /dev/null +++ b/backend/test/workspace-maintenance.test.ts @@ -0,0 +1,85 @@ +import { expect, test, vi } from "vitest"; +import { + runWorkspaceMaintenanceCli, + type WorkspaceMaintenanceIo, +} from "../src/workspace-maintenance.js"; +import type { WorkspaceOperationResult } from "../src/workspaces/preprocessing-service.js"; + +function io(stdin: string): WorkspaceMaintenanceIo & { stdout: string[]; stderr: string[] } { + const stdout: string[] = []; + const stderr: string[] = []; + return { + stdin, + stdout, + stderr, + writeStdout: (value) => void stdout.push(value), + writeStderr: (value) => void stderr.push(value), + }; +} + +function ok(operation: string): WorkspaceOperationResult { + return { + schemaVersion: 1, + status: "succeeded", + code: "ok", + workspaceId: "psd-clinical", + workspaceRevision: "a".repeat(40), + descriptorBlob: "b".repeat(40), + operation, + completedStages: [], + }; +} + +test("entrypoint emits exactly one pristine JSON document and maps success/block/failure exits", async () => { + const service = { + inspect: vi.fn(async () => ok("inspect")), + preprocessDwh: vi.fn(async () => ({ ...ok("preprocess dwh"), status: "blocked", code: "manual_review_required" as const })), + run: vi.fn(async () => ({ ...ok("preprocess run"), status: "failed", code: "semantic_index_incompatible" as const })), + } as any; + + const inspectIo = io(JSON.stringify({ schemaVersion: 1, workspaceId: "psd-clinical" })); + expect(await runWorkspaceMaintenanceCli(["node", "workspace-maintenance", "inspect"], service, inspectIo)).toBe(0); + expect(JSON.parse(inspectIo.stdout.join(""))).toMatchObject({ operation: "inspect", code: "ok" }); + expect(inspectIo.stderr.join("")).toBe(""); + + const blockedIo = io(JSON.stringify({ schemaVersion: 1, workspaceId: "psd-clinical" })); + expect(await runWorkspaceMaintenanceCli(["node", "workspace-maintenance", "preprocess-dwh"], service, blockedIo)).toBe(3); + expect(JSON.parse(blockedIo.stdout.join(""))).toMatchObject({ code: "manual_review_required" }); + + const failedIo = io(JSON.stringify({ schemaVersion: 1, workspaceId: "psd-clinical" })); + expect(await runWorkspaceMaintenanceCli(["node", "workspace-maintenance", "preprocess-run"], service, failedIo)).toBe(1); + expect(JSON.parse(failedIo.stdout.join(""))).toMatchObject({ code: "semantic_index_incompatible" }); +}); + +test("malformed stdin, unknown commands, and extra fields fail with exit 2 but still return bounded JSON", async () => { + const service = {} as any; + + const malformedIo = io("not-json"); + expect(await runWorkspaceMaintenanceCli(["node", "workspace-maintenance", "inspect"], service, malformedIo)).toBe(2); + expect(JSON.parse(malformedIo.stdout.join(""))).toMatchObject({ status: "failed", operation: "inspect" }); + + const extraFieldIo = io(JSON.stringify({ schemaVersion: 1, workspaceId: "psd-clinical", unexpected: true })); + expect(await runWorkspaceMaintenanceCli(["node", "workspace-maintenance", "inspect"], service, extraFieldIo)).toBe(2); + expect(JSON.parse(extraFieldIo.stdout.join(""))).toMatchObject({ status: "failed", operation: "inspect" }); + + const unknownIo = io(JSON.stringify({ schemaVersion: 1, workspaceId: "psd-clinical" })); + expect(await runWorkspaceMaintenanceCli(["node", "workspace-maintenance", "explode"], service, unknownIo)).toBe(2); + expect(JSON.parse(unknownIo.stdout.join(""))).toMatchObject({ status: "failed", operation: "explode" }); +}); + +test("raw exception text is redacted from stderr and stdout remains within the public result contract", async () => { + const service = { + inspect: vi.fn(async () => { + throw new Error("https://secret.example.invalid?q=token SELECT * FROM sensitive_table"); + }), + } as any; + const captured = io(JSON.stringify({ schemaVersion: 1, workspaceId: "psd-clinical" })); + + expect(await runWorkspaceMaintenanceCli(["node", "workspace-maintenance", "inspect"], service, captured)).toBe(1); + expect(JSON.parse(captured.stdout.join(""))).toMatchObject({ + status: "failed", + operation: "inspect", + }); + expect(captured.stderr.join("")).not.toContain("secret.example.invalid"); + expect(captured.stderr.join("")).not.toContain("SELECT *"); +}); diff --git a/backend/test/workspace-preprocessing-service.test.ts b/backend/test/workspace-preprocessing-service.test.ts new file mode 100644 index 00000000..66471a93 --- /dev/null +++ b/backend/test/workspace-preprocessing-service.test.ts @@ -0,0 +1,333 @@ +import { mkdtempSync, readFileSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { afterEach, expect, test, vi } from "vitest"; +import { parseWorkspaceYaml } from "../src/workspaces/schema.js"; +import { PreprocessingStateStore } from "../src/workspaces/preprocessing-state.js"; +import { + WorkspacePreprocessingService, + type ChildProcessRequest, +} from "../src/workspaces/preprocessing-service.js"; + +const roots: string[] = []; + +afterEach(() => { + roots.splice(0).forEach((root) => rmSync(root, { recursive: true, force: true })); +}); + +const semanticRuntime = { + internalQdrantUrl: "http://qdrant:6333", + internalEmbeddingUrl: "http://embedding:11434", + internalEmbeddingModel: "qwen3-embedding:0.6b", + internalEmbeddingDimensions: 1024, +}; + +const baseWorkspace = parseWorkspaceYaml(`workspace: + schema_version: 3 + id: psd-clinical + name: Runtime Lease + language: en +dwh: + engine: postgres + database: analytics + schema: mart + supported_transports: [postgres_direct] +semantic_index: + vector_store: + engine: qdrant + collection: psd-clinical + dimensions: 1024 + distance: cosine + embedding: + provider: ollama_internal + model: qwen3-embedding:0.6b + dimensions: 1024 +llm_policy: + allowed: [zai/glm-5.2] +`); +const filesystemWorkspace = parseWorkspaceYaml(`${baseWorkspace ? '' : ''}workspace: + schema_version: 3 + id: fs-workspace + name: Filesystem + language: en +dwh: + engine: postgres + database: analytics + schema: mart + supported_transports: [postgres_direct] +semantic_index: + vector_store: + engine: qdrant + collection: fs-workspace + dimensions: 1024 + distance: cosine + embedding: + provider: ollama_internal + model: qwen3-embedding:0.6b + dimensions: 1024 +llm_policy: + allowed: [zai/glm-5.2] +evidence: + source: + type: filesystem + uri: fs-workspace/evidence +`); +const privateHttpWorkspace = parseWorkspaceYaml(`workspace: + schema_version: 3 + id: http-workspace + name: Http + language: en +dwh: + engine: postgres + database: analytics + schema: mart + supported_transports: [postgres_direct] +semantic_index: + vector_store: + engine: qdrant + collection: http-workspace + dimensions: 1024 + distance: cosine + embedding: + provider: ollama_internal + model: qwen3-embedding:0.6b + dimensions: 1024 +llm_policy: + allowed: [zai/glm-5.2] +evidence: + source: + type: http + uris: [http://127.0.0.1/private.md] + authentication: none + connect_timeout_ms: 1000 + read_timeout_ms: 2000 + max_bytes: 100 + max_redirects: 0 + allow_private_hosts: true + max_cache_bytes: 100 +`); + +function runtime(workspace = baseWorkspace, workspaceId = workspace.workspace.id) { + return { + workspace, + workspaceId, + workspaceRevision: "a".repeat(40), + descriptorBlob: "b".repeat(40), + catalogBlob: "c".repeat(40), + configLease: { + path: `/data/sessions/${workspaceId}/preprocessing/runtime-config/${"a".repeat(40)}.yaml`, + workspaceId, + workspaceRevision: "a".repeat(40), + descriptorBlob: "b".repeat(40), + catalogBlob: "c".repeat(40), + configDigest: "sha256:config", + bindingDigest: "sha256:bindings", + release: () => undefined, + }, + }; +} + +function fixture(workspace = baseWorkspace) { + const dataRoot = mkdtempSync(join(tmpdir(), "tht-preprocessing-service-")); + roots.push(dataRoot); + const requests: ChildProcessRequest[] = []; + const runChild = vi.fn(async (request: ChildProcessRequest) => { + requests.push(request); + return { exitCode: 0, stdout: JSON.stringify({ status: "succeeded" }), stderr: "" }; + }); + const service = new WorkspacePreprocessingService({ + dataRoot, + acquireActiveRuntime: async () => runtime(workspace), + runChild, + listSessions: async () => [], + semanticPreflight: async () => ({ ok: true }), + }); + return { dataRoot, runChild, requests, service }; +} + +test("preprocess dwh uses fixed argv and resumes outer state without rerunning a completed stage", async () => { + const f = fixture(); + f.runChild.mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ status: "succeeded", run_id: "d".repeat(32) }), + stderr: "", + }); + + const first = await f.service.preprocessDwh({ workspaceId: "psd-clinical" }); + expect(first).toMatchObject({ + status: "succeeded", + code: "ok", + operation: "preprocess dwh", + completedStages: ["dwh"], + childRuns: { dwh: "d".repeat(32) }, + }); + expect((f.runChild.mock.calls[0]![0] as ChildProcessRequest).argv).toEqual([ + "preprocess", "dwh", "--steps", "introspect,lsh", "--json", "-c", "/dev/fd/3", + ]); + + const second = await f.service.preprocessDwh({ + workspaceId: "psd-clinical", + resumeRunId: first.runId!, + }); + expect(second.status).toBe("unchanged"); + expect(f.runChild).toHaveBeenCalledTimes(1); +}); + +test("schema suggest-fks publishes a candidate artifact and blocks full runs for manual review", async () => { + const f = fixture(); + f.runChild + .mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ status: "succeeded", run_id: "d".repeat(32) }), + stderr: "", + }) + .mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ + status: "succeeded", + candidate_count: 1, + candidate_digest: "sha256:" + "e".repeat(64), + candidate_yaml: "tables: []\n", + }), + stderr: "", + }); + + const result = await f.service.run({ workspaceId: "psd-clinical" }); + expect(result).toMatchObject({ + status: "blocked", + code: "manual_review_required", + completedStages: ["dwh", "fk_suggest"], + }); + expect(f.runChild.mock.calls.map(([request]) => (request as ChildProcessRequest).argv[0])).toEqual(["preprocess", "schema"]); +}); + +test("schema check requires the exact candidate digest, stages annotations via temp file, and persists the review", async () => { + const f = fixture(); + f.runChild.mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ + status: "succeeded", + candidate_count: 1, + candidate_digest: "sha256:" + "e".repeat(64), + candidate_yaml: "tables: []\n", + }), + stderr: "", + }); + const suggest = await f.service.suggestFks({ workspaceId: "psd-clinical" }); + await expect(f.service.checkSchema({ + workspaceId: "psd-clinical", + annotationsYaml: "tables: {}\n", + reviewedCandidatesDigest: "sha256:" + "f".repeat(64), + })).resolves.toMatchObject({ status: "failed", code: "annotation_invalid" }); + + const reviewedDigest = suggest.artifactIdentities![0]!.digest; + let stagedPath = ""; + f.runChild.mockImplementationOnce(async (request: ChildProcessRequest) => { + stagedPath = request.argv[request.argv.indexOf("--annotations") + 1]!; + expect(readFileSync(stagedPath, "utf8")).toBe("tables: {}\n"); + expect(request.argv).toEqual([ + "schema", "check", "--annotations", stagedPath, + "--reviewed-candidates", reviewedDigest, + "--json", "-c", "/dev/fd/3", + ]); + return { + exitCode: 0, + stdout: JSON.stringify({ + status: "succeeded", + orphan_count: 0, + annotations_digest: "sha256:annotations", + reviewed_candidates_digest: reviewedDigest, + }), + stderr: "", + }; + }); + + const checked = await f.service.checkSchema({ + workspaceId: "psd-clinical", + annotationsYaml: "tables: {}\n", + reviewedCandidatesDigest: reviewedDigest, + }); + expect(checked).toMatchObject({ status: "succeeded", code: "ok" }); + expect(() => readFileSync(stagedPath, "utf8")).toThrow(); + const state = new PreprocessingStateStore({ dataRoot: f.dataRoot, workspaceId: "psd-clinical" }); + expect(state.readFkReview(suggest.runId!)?.reviewedCandidatesDigest).toBe(reviewedDigest); +}); + +test("index schema fails closed when semantic preflight refuses the collection", async () => { + const dataRoot = mkdtempSync(join(tmpdir(), "tht-preprocessing-service-")); + roots.push(dataRoot); + const runChild = vi.fn(); + const service = new WorkspacePreprocessingService({ + dataRoot, + acquireActiveRuntime: async () => runtime(baseWorkspace), + runChild, + listSessions: async () => [], + semanticPreflight: async () => ({ ok: false, code: "semantic_index_incompatible" }), + }); + + const result = await service.indexSchema({ workspaceId: "psd-clinical" }); + expect(result).toMatchObject({ status: "failed", code: "semantic_index_incompatible" }); + expect(runChild).not.toHaveBeenCalled(); +}); + +test("evidence stops before child execution for filesystem sources and refuses private HTTP hosts outside the allowlist", async () => { + const filesystem = fixture(filesystemWorkspace); + const blocked = await filesystem.service.preprocessEvidence({ workspaceId: "fs-workspace" }); + expect(blocked).toMatchObject({ status: "blocked", code: "evidence_materialization_required" }); + expect(filesystem.runChild).not.toHaveBeenCalled(); + + const httpDataRoot = mkdtempSync(join(tmpdir(), "tht-preprocessing-service-")); + roots.push(httpDataRoot); + const httpService = new WorkspacePreprocessingService({ + dataRoot: httpDataRoot, + acquireActiveRuntime: async () => runtime(privateHttpWorkspace, "http-workspace"), + runChild: vi.fn(), + listSessions: async () => [], + semanticPreflight: async () => ({ ok: true }), + httpPrivateHostAllowlist: ["metadata.internal"], + }); + + const refused = await httpService.preprocessEvidence({ workspaceId: "http-workspace" }); + expect(refused).toMatchObject({ status: "failed", code: "egress_policy_refused" }); +}); + +test("full runs follow the explicit order and finish unchanged when no Evidence source exists", async () => { + const f = fixture(); + f.runChild + .mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ status: "succeeded", run_id: "d".repeat(32) }), + stderr: "", + }) + .mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ + status: "succeeded", + candidate_count: 0, + candidate_digest: "sha256:" + "0".repeat(64), + candidate_yaml: "tables: []\n", + }), + stderr: "", + }) + .mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ + status: "succeeded", + counts: { added: 1, updated: 0, deleted: 0, unchanged: 0 }, + }), + stderr: "", + }); + + const result = await f.service.run({ workspaceId: "psd-clinical" }); + expect(result).toMatchObject({ + status: "succeeded", + code: "ok", + completedStages: ["dwh", "fk_suggest", "schema_index"], + warnings: ["workspace has no Evidence source"], + }); + expect(f.runChild.mock.calls.map(([request]) => (request as ChildProcessRequest).argv.slice(0, 2).join(" "))).toEqual([ + "preprocess dwh", + "schema suggest-fks", + "vector index-schema", + ]); +}); diff --git a/backend/test/workspace-preprocessing-state.test.ts b/backend/test/workspace-preprocessing-state.test.ts new file mode 100644 index 00000000..b0e69788 --- /dev/null +++ b/backend/test/workspace-preprocessing-state.test.ts @@ -0,0 +1,119 @@ +import { existsSync, mkdtempSync, readFileSync, rmSync, statSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { afterEach, expect, test } from "vitest"; +import { PreprocessingStateStore } from "../src/workspaces/preprocessing-state.js"; + +const roots: string[] = []; + +afterEach(() => { + roots.splice(0).forEach((root) => rmSync(root, { recursive: true, force: true })); +}); + +function fixture() { + const dataRoot = mkdtempSync(join(tmpdir(), "tht-preprocessing-state-")); + roots.push(dataRoot); + return { + dataRoot, + store: new PreprocessingStateStore({ dataRoot, workspaceId: "psd-clinical" }), + }; +} + +test("job state creates durable 0600 JSON and enforces same-revision resume", async () => { + const { store } = fixture(); + const job = await store.beginJob({ + operation: "preprocess dwh", + workspaceRevision: "a".repeat(40), + descriptorBlob: "b".repeat(40), + catalogBlob: "c".repeat(40), + configDigest: "sha256:config", + bindingDigest: "sha256:bindings", + }); + + const path = store.jobPath(job.runId); + expect(existsSync(path)).toBe(true); + expect(statSync(path).mode & 0o777).toBe(0o600); + expect(JSON.parse(readFileSync(path, "utf8"))).toMatchObject({ + schemaVersion: 1, + operation: "preprocess dwh", + workspaceId: "psd-clinical", + workspaceRevision: "a".repeat(40), + descriptorBlob: "b".repeat(40), + catalogBlob: "c".repeat(40), + configDigest: "sha256:config", + bindingDigest: "sha256:bindings", + }); + + await expect(store.beginJob({ + operation: "preprocess dwh", + runId: job.runId, + workspaceRevision: "d".repeat(40), + descriptorBlob: "b".repeat(40), + catalogBlob: "c".repeat(40), + configDigest: "sha256:config", + bindingDigest: "sha256:bindings", + })).rejects.toMatchObject({ code: "preprocessing_resume_mismatch" }); + + const resumed = await store.beginJob({ + operation: "preprocess dwh", + runId: job.runId, + workspaceRevision: "a".repeat(40), + descriptorBlob: "b".repeat(40), + catalogBlob: "c".repeat(40), + configDigest: "sha256:config", + bindingDigest: "sha256:bindings", + }); + expect(resumed.runId).toBe(job.runId); +}); + +test("writer lock rejects a concurrent contender and the kernel releases it after holder death", async () => { + const f = fixture(); + const other = new PreprocessingStateStore({ dataRoot: f.dataRoot, workspaceId: "psd-clinical" }); + const first = await f.store.acquireWriterLock(); + await expect(other.acquireWriterLock()).rejects.toMatchObject({ code: "preprocessing_conflict" }); + + process.kill(first.holderPid, "SIGKILL"); + const deadline = Date.now() + 5_000; + while (Date.now() < deadline) { + try { + const recovered = await other.acquireWriterLock(); + await recovered.release(); + return; + } catch (error) { + if ((error as { code?: string }).code !== "preprocessing_conflict") throw error; + await new Promise((resolve) => setTimeout(resolve, 50)); + } + } + throw new Error("writer lock was not released after holder death"); +}); + +test("session inventory guard blocks resumable sessions pinned to a different revision", async () => { + const { store } = fixture(); + + await expect(store.assertSessionInventoryCompatible("a".repeat(40), [ + { id: "open-other", status: "closed", archived: false, workspaceRevision: "b".repeat(40) }, + ])).rejects.toMatchObject({ code: "preprocessing_conflict" }); + + await expect(store.assertSessionInventoryCompatible("a".repeat(40), [ + { id: "current", status: "open", archived: false, workspaceRevision: "a".repeat(40) }, + { id: "finalized", status: "finalized", archived: false, workspaceRevision: "b".repeat(40) }, + { id: "archived", status: "closed", archived: true, workspaceRevision: "c".repeat(40) }, + ])).resolves.toBeUndefined(); +}); + +test("candidate and review artifacts are digest-bound durable files", async () => { + const { store } = fixture(); + const runId = "1".repeat(32); + const candidate = await store.writeFkCandidates(runId, `tables: [] +`); + const review = await store.writeFkReview(runId, { + reviewedCandidatesDigest: candidate.digest, + annotationsDigest: "sha256:annotations", + workspaceRevision: "a".repeat(40), + }); + + expect(candidate.digest).toMatch(/^sha256:[0-9a-f]{64}$/); + expect(review.digest).toMatch(/^sha256:[0-9a-f]{64}$/); + expect(store.readFkCandidates(runId)?.digest).toBe(candidate.digest); + expect(store.readFkReview(runId)?.reviewedCandidatesDigest).toBe(candidate.digest); +}); diff --git a/backend/test/workspace-runtime-config-lease.test.ts b/backend/test/workspace-runtime-config-lease.test.ts new file mode 100644 index 00000000..b91e3e27 --- /dev/null +++ b/backend/test/workspace-runtime-config-lease.test.ts @@ -0,0 +1,239 @@ +import { execFile } from "node:child_process"; +import { + chmodSync, + existsSync, + mkdtempSync, + mkdirSync, + readFileSync, + rmSync, + statSync, + symlinkSync, + writeFileSync, +} from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { promisify } from "node:util"; +import { afterEach, expect, test, vi } from "vitest"; +import { WorkspaceRegistry } from "../src/workspaces/registry.js"; +import type { WorkspaceRegistryConfig } from "../src/workspaces/types.js"; +import { + publishDeterministicRuntimeConfigLease, + renderActiveWorkspaceRuntime, + renderWorkspaceRuntimeFromSnapshotPath, +} from "../src/workspaces/runtime-config-lease.js"; + +const runFile = promisify(execFile); +const roots: string[] = []; + +afterEach(() => { + vi.unstubAllEnvs(); + roots.splice(0).forEach((root) => rmSync(root, { recursive: true, force: true })); +}); + +async function git(cwd: string, args: string[]): Promise { + return (await runFile("git", args, { cwd })).stdout.trim(); +} + +async function fixture() { + const root = mkdtempSync(join(tmpdir(), "tht-runtime-lease-")); + roots.push(root); + const remote = join(root, "remote.git"); + const source = join(root, "source"); + const registryRoot = join(root, "registry"); + const secretRoot = join(root, "secrets"); + const dataRoot = join(root, "data"); + const harnessDir = join(root, "harness"); + mkdirSync(harnessDir, { recursive: true }); + mkdirSync(join(harnessDir, "config"), { recursive: true }); + writeFileSync(join(harnessDir, "config", "tht.yaml"), `session_storage: + mode: local +profile: server +`); + + await git(root, ["init", "--bare", "--initial-branch=main", remote]); + mkdirSync(source); + await git(source, ["init", "--initial-branch=main"]); + await git(source, ["config", "user.name", "Runtime Lease Test"]); + await git(source, ["config", "user.email", "runtime-lease@example.invalid"]); + writeFileSync(join(source, "thoth-workspaces.yaml"), `schema_version: 1 +workspaces: [{id: psd-clinical, name: Runtime Lease}] +`); + mkdirSync(join(source, "psd-clinical", "evidence"), { recursive: true }); + writeFileSync(join(source, "psd-clinical", "workspace.yaml"), `workspace: + schema_version: 3 + id: psd-clinical + name: Runtime Lease + language: en +dwh: + engine: postgres + database: analytics + schema: mart + supported_transports: [postgres_direct] +semantic_index: + vector_store: + engine: qdrant + collection: psd-clinical + dimensions: 1024 + distance: cosine + embedding: + provider: ollama_internal + model: qwen3-embedding:0.6b + dimensions: 1024 +llm_policy: + allowed: [zai/glm-5.2] +evidence: + source: + type: filesystem + uri: psd-clinical/evidence +`); + writeFileSync(join(source, "psd-clinical", "evidence", "guide.md"), `# hello +`); + await git(source, ["add", "."]); + await git(source, ["commit", "-m", "Canonical workspace"]); + await git(source, ["remote", "add", "origin", remote]); + await git(source, ["push", "origin", "main"]); + + mkdirSync(secretRoot); + const passwordFile = join(secretRoot, "dwh-password"); + writeFileSync(passwordFile, "secret", { mode: 0o600 }); + chmodSync(passwordFile, 0o600); + mkdirSync(dataRoot); + + const registryConfig: WorkspaceRegistryConfig = { + root: registryRoot, + remoteUrl: remote, + branch: "main", + gitAuthorName: "Runtime Lease Test", + gitAuthorEmail: "runtime-lease@example.invalid", + installationId: "test", + secretRoots: [secretRoot], + maxImportBytes: 1024 * 1024, + maxImportEntries: 16, + }; + const registry = new WorkspaceRegistry(registryConfig); + await registry.bootstrap(); + const revision = (await registry.list())[0]; + + vi.stubEnv("THT_WS_PSD_CLINICAL_DWH_TRANSPORT", "postgres_direct"); + vi.stubEnv("THT_WS_PSD_CLINICAL_DWH_HOST", "warehouse.internal"); + vi.stubEnv("THT_WS_PSD_CLINICAL_DWH_PORT", "5432"); + vi.stubEnv("THT_WS_PSD_CLINICAL_DWH_USER", "reader"); + vi.stubEnv("THT_WS_PSD_CLINICAL_DWH_PASSWORD_FILE", passwordFile); + + return { + dataRoot, + harnessDir, + registry, + registryConfig, + revision, + }; +} + +const semanticRuntime = { + internalQdrantUrl: "http://qdrant:6333", + internalEmbeddingUrl: "http://embedding:11434", + internalEmbeddingModel: "qwen3-embedding:0.6b", + internalEmbeddingDimensions: 1024, +}; + +test("active workspace rendering is byte-identical to direct snapshot rendering", async () => { + const f = await fixture(); + const direct = renderWorkspaceRuntimeFromSnapshotPath({ + snapshotPath: f.revision.snapshotPath, + harnessDir: f.harnessDir, + configPath: "config/tht.yaml", + dataRoot: f.dataRoot, + secretRoots: f.registryConfig.secretRoots, + semanticRuntime, + }); + const active = await renderActiveWorkspaceRuntime({ + workspaceId: "psd-clinical", + registry: f.registry, + registryConfig: f.registryConfig, + harnessDir: f.harnessDir, + configPath: "config/tht.yaml", + dataRoot: f.dataRoot, + secretRoots: f.registryConfig.secretRoots, + semanticRuntime, + }); + + expect(active.renderedConfig).toBe(direct.renderedConfig); + expect(active.workspaceRevision).toBe(f.revision.commit); + expect(active.descriptorBlob).toBe(f.revision.blob); + expect(active.catalogBlob).toMatch(/^[0-9a-f]{40}$/); +}); + +test("deterministic operator leases publish one revision-bound protected config and refuse changed same-revision bytes", async () => { + const f = await fixture(); + const first = await publishDeterministicRuntimeConfigLease({ + workspaceId: "psd-clinical", + registry: f.registry, + registryConfig: f.registryConfig, + harnessDir: f.harnessDir, + configPath: "config/tht.yaml", + dataRoot: f.dataRoot, + secretRoots: f.registryConfig.secretRoots, + semanticRuntime, + }); + const second = await publishDeterministicRuntimeConfigLease({ + workspaceId: "psd-clinical", + registry: f.registry, + registryConfig: f.registryConfig, + harnessDir: f.harnessDir, + configPath: "config/tht.yaml", + dataRoot: f.dataRoot, + secretRoots: f.registryConfig.secretRoots, + semanticRuntime, + }); + + expect(second.path).toBe(first.path); + expect(first.path).toBe(join( + f.dataRoot, + "sessions", + "psd-clinical", + "preprocessing", + "runtime-config", + `${f.revision.commit}.yaml`, + )); + expect(statSync(first.path).mode & 0o777).toBe(0o400); + expect(statSync(first.manifestPath).mode & 0o777).toBe(0o600); + expect(readFileSync(first.path, "utf8")).toContain("collection_lifecycle: require_existing"); + expect(existsSync(first.manifestPath)).toBe(true); + + vi.stubEnv("THT_WS_PSD_CLINICAL_DWH_HOST", "warehouse-two.internal"); + await expect(publishDeterministicRuntimeConfigLease({ + workspaceId: "psd-clinical", + registry: f.registry, + registryConfig: f.registryConfig, + harnessDir: f.harnessDir, + configPath: "config/tht.yaml", + dataRoot: f.dataRoot, + secretRoots: f.registryConfig.secretRoots, + semanticRuntime, + })).rejects.toMatchObject({ code: "effective_config_mismatch" }); +}); + +test("runtime rendering rejects untrusted snapshot paths and symlinks", async () => { + const f = await fixture(); + const outside = join(f.dataRoot, "outside.yaml"); + writeFileSync(outside, readFileSync(f.revision.snapshotPath, "utf8")); + const symlink = join(f.dataRoot, "alias.yaml"); + symlinkSync(f.revision.snapshotPath, symlink); + + expect(() => renderWorkspaceRuntimeFromSnapshotPath({ + snapshotPath: outside, + harnessDir: f.harnessDir, + configPath: "config/tht.yaml", + dataRoot: f.dataRoot, + secretRoots: f.registryConfig.secretRoots, + semanticRuntime, + })).toThrow(/trusted runtime snapshot/i); + expect(() => renderWorkspaceRuntimeFromSnapshotPath({ + snapshotPath: symlink, + harnessDir: f.harnessDir, + configPath: "config/tht.yaml", + dataRoot: f.dataRoot, + secretRoots: f.registryConfig.secretRoots, + semanticRuntime, + })).toThrow(/trusted runtime snapshot/i); +});