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 { syncAnnotations } from "../src/workspaces/annotations-sync.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", internalEmbeddingId: "ollama/qwen3-embedding:0.6b", internalEmbeddingModel: "qwen3-embedding:0.6b", internalEmbeddingDimensions: 1024, }; const baseWorkspace = parseWorkspaceYaml(`workspace: schema_version: 4 id: psd-clinical name: Runtime Lease language: en dwh: engine: postgres database: analytics schema: mart supported_transports: [postgres_direct] `); const filesystemWorkspace = parseWorkspaceYaml(`${baseWorkspace ? '' : ''}workspace: schema_version: 4 id: fs-workspace name: Filesystem language: en dwh: engine: postgres database: analytics schema: mart supported_transports: [postgres_direct] evidence: source: type: filesystem uri: fs-workspace/evidence `); const privateHttpWorkspace = parseWorkspaceYaml(`workspace: schema_version: 4 id: http-workspace name: Http language: en dwh: engine: postgres database: analytics schema: mart supported_transports: [postgres_direct] 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)}-identitysuffix.yaml`, workspaceId, workspaceRevision: "a".repeat(40), descriptorBlob: "b".repeat(40), catalogBlob: "c".repeat(40), configDigest: "sha256:config", bindingDigest: "sha256:bindings", semanticQdrantUrl: "http://qdrant:6333", effectiveConfig: { schemaVersion: 2, dwh: { engine: "postgres", database: "analytics", schema: "mart", transport: "postgres_direct", host: "dwh.internal", port: 5432, user: "reader", }, vector: { collection: workspaceId, dimensions: 1024, distance: "cosine" }, embedding: { id: "ollama/qwen3-embedding:0.6b", model: "qwen3-embedding:0.6b", dimensions: 1024 }, roots: { artifacts: "/data/artifacts", indexes: "/data/indexes" }, }, effectiveConfigIdentity: "workspace://psd-clinical@v1:" + "d".repeat(64), configFingerprint: "sha256:" + "e".repeat(64), inputFingerprint: "sha256:" + "f".repeat(64), 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 }), evidencePreflight: 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 and stages annotations via a temp file without recording a 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!)).toBeUndefined(); }); 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" }), evidencePreflight: async () => ({ ok: true }), }); const result = await service.indexSchema({ workspaceId: "psd-clinical" }); expect(result).toMatchObject({ status: "failed", code: "semantic_index_incompatible" }); expect(runChild).not.toHaveBeenCalled(); }); test("filesystem Evidence proceeds after materialization and private HTTP hosts outside the allowlist are refused", async () => { const filesystem = fixture(filesystemWorkspace); filesystem.runChild.mockResolvedValueOnce({ exitCode: 0, stdout: JSON.stringify({ status: "succeeded", counts: { added: 1 } }), stderr: "", }); const materialized = await filesystem.service.preprocessEvidence({ workspaceId: "fs-workspace" }); expect(materialized).toMatchObject({ status: "succeeded", code: "ok" }); expect(filesystem.runChild).toHaveBeenCalledTimes(1); 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 }), evidencePreflight: 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", ]); }); test("schema accept validates the synced Git blob and records 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" }); const runId = suggest.runId!; const candidateDigest = suggest.artifactIdentities![0]!.digest; const synced = syncAnnotations({ dataRoot: f.dataRoot, workspaceId: "psd-clinical", commit: "a".repeat(40), blobId: "b".repeat(40), contents: Buffer.from("tables: {}\n"), }); f.runChild.mockResolvedValueOnce({ exitCode: 0, stdout: JSON.stringify({ status: "succeeded", orphan_count: 0, annotations_digest: synced.contentDigest, reviewed_candidates_digest: candidateDigest, }), stderr: "", }); const result = await f.service.acceptSchema({ workspaceId: "psd-clinical", runId, yes: true }); expect(result).toMatchObject({ status: "succeeded", code: "ok", operation: "schema accept" }); const state = new PreprocessingStateStore({ dataRoot: f.dataRoot, workspaceId: "psd-clinical" }); expect(state.readFkReview(runId)).toMatchObject({ reviewedCandidatesDigest: candidateDigest, annotationsDigest: synced.contentDigest, workspaceRevision: "a".repeat(40), blobId: "b".repeat(40), }); }); test("schema accept fails closed without --yes, for an unknown run, or with no synced annotations", 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" }); const runId = suggest.runId!; await expect(f.service.acceptSchema({ workspaceId: "psd-clinical", runId, yes: false })) .resolves.toMatchObject({ status: "failed", code: "annotation_invalid" }); await expect(f.service.acceptSchema({ workspaceId: "psd-clinical", runId: "e".repeat(32), yes: true })) .resolves.toMatchObject({ status: "failed", code: "annotation_invalid" }); await expect(f.service.acceptSchema({ workspaceId: "psd-clinical", runId, yes: true })) .resolves.toMatchObject({ status: "failed", code: "annotation_invalid" }); expect(f.runChild).toHaveBeenCalledTimes(1); }); test("full runs continue after schema accept only when the accepted blob matches the current revision", async () => { const f = fixture(); const revision = "a".repeat(40); 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 blocked = await f.service.run({ workspaceId: "psd-clinical" }); expect(blocked).toMatchObject({ status: "blocked", code: "manual_review_required" }); const runId = blocked.runId!; const state = new PreprocessingStateStore({ dataRoot: f.dataRoot, workspaceId: "psd-clinical" }); const candidateDigest = state.readFkCandidates(runId)!.digest; const synced = syncAnnotations({ dataRoot: f.dataRoot, workspaceId: "psd-clinical", commit: revision, blobId: "b".repeat(40), contents: Buffer.from("tables: {}\n"), }); // Accept writes the review; the resume then passes the gate and reaches schema indexing. f.runChild.mockResolvedValueOnce({ exitCode: 0, stdout: JSON.stringify({ status: "succeeded", orphan_count: 0, annotations_digest: synced.contentDigest, reviewed_candidates_digest: candidateDigest }), stderr: "", }); await f.service.acceptSchema({ workspaceId: "psd-clinical", runId, yes: true }); f.runChild.mockResolvedValueOnce({ exitCode: 0, stdout: JSON.stringify({ status: "succeeded", counts: { added: 1, updated: 0, deleted: 0, unchanged: 0 } }), stderr: "", }); const resumed = await f.service.run({ workspaceId: "psd-clinical", resumeRunId: runId }); expect(resumed).toMatchObject({ status: "succeeded", code: "ok" }); // A review whose accepted blob digest no longer matches the current revision stays blocked. const second = fixture(); second.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 secondBlocked = await second.service.run({ workspaceId: "psd-clinical" }); const secondRunId = secondBlocked.runId!; const secondState = new PreprocessingStateStore({ dataRoot: second.dataRoot, workspaceId: "psd-clinical" }); const secondCandidate = secondState.readFkCandidates(secondRunId)!.digest; secondState.writeFkReview(secondRunId, { reviewedCandidatesDigest: secondCandidate, annotationsDigest: "sha256:" + "0".repeat(64), workspaceRevision: revision, blobId: "b".repeat(40), }); const stillBlocked = await second.service.run({ workspaceId: "psd-clinical", resumeRunId: secondRunId }); expect(stillBlocked).toMatchObject({ status: "blocked", code: "manual_review_required" }); }); test("vector rebuild recreates the full collection contract including keyword indexes", async () => { const f = fixture(); // mock fetch: DELETE ok, then reconcileCollection self-heals create + indexes (real fetch in deps) const calls: string[] = []; const fakeFetch = async (url: string, init?: any) => { calls.push(`${init?.method ?? "GET"} ${url}`); if ((init?.method ?? "GET") === "DELETE") return new Response("", { status: 200 }); if (url.endsWith("/collections/psd-clinical") && init?.method === "PUT") return new Response("", { status: 200 }); if (url.endsWith("/collections/psd-clinical") && init?.method === "GET") { return new Response(JSON.stringify({ result: { config: { params: { vectors: { size: 1024, distance: "Cosine" } } }, payload_schema: { content_hash: { data_type: "keyword" }, document_id: { data_type: "keyword" }, kind: { data_type: "keyword" }, record_key: { data_type: "keyword" }, record_kind: { data_type: "keyword" }, vector_generation: { data_type: "keyword" }, workspace_id: { data_type: "keyword" }, workspace_revision: { data_type: "keyword" } } } }), { status: 200 }); } if (url.endsWith("/collections/psd-clinical/index") && init?.method === "PUT") return new Response("", { status: 200 }); return new Response(JSON.stringify({ result: {} }), { status: 200 }); }; const service = new WorkspacePreprocessingService({ dataRoot: f.dataRoot, acquireActiveRuntime: async () => runtime(baseWorkspace), runChild: vi.fn(), listSessions: async () => [], semanticPreflight: async () => ({ ok: true }), evidencePreflight: async () => ({ ok: true }), }); // replace global fetch used by vectorRebuild/reconcileCollection const original = globalThis.fetch; globalThis.fetch = fakeFetch as any; try { const result = await service.vectorRebuild({ workspaceId: "psd-clinical", collection: "psd-clinical", confirm: "psd-clinical", destroy: true }); expect(result).toMatchObject({ status: "succeeded", code: "ok" }); } finally { globalThis.fetch = original; } expect(calls.some((c) => c.startsWith("DELETE "))).toBe(true); });