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 { validateWorkspaceDescriptor } from "../src/workspaces/schema.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 workspace = validateWorkspaceDescriptor({ workspace: { schema_version: 4, id: "catalog-workspace", name: "Catalog only", language: "en", }, }); function runtime() { return { workspace, workspaceId: "catalog-workspace", workspaceRevision: "a".repeat(40), descriptorBlob: "b".repeat(40), catalogBlob: "c".repeat(40), configLease: { path: `/data/sessions/catalog-workspace/preprocessing/runtime-config/${"a".repeat(40)}.yaml`, workspaceId: "catalog-workspace", workspaceRevision: "a".repeat(40), descriptorBlob: "b".repeat(40), catalogBlob: "c".repeat(40), configDigest: "sha256:config", bindingDigest: "sha256:bindings", semanticQdrantUrl: "http://qdrant:6333", effectiveConfig: { schemaVersion: 3, dwh: { engine: "postgres", database: "warehouse", schema: "analytics", transport: "postgres_direct", host: "db.internal", port: 5432, user: "reader", }, vector: { collections: { reference: "catalog-workspace-reference", memory: "catalog-workspace-memory", }, 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://catalog-workspace@v1:" + "d".repeat(64), configFingerprint: "sha256:" + "e".repeat(64), inputFingerprint: "sha256:" + "f".repeat(64), release: () => undefined, }, }; } const database = { id: "database-1", workspaceId: "catalog-workspace", engine: "postgres", databaseName: "warehouse", schema: "analytics", version: 4, createdAt: "2026-01-01T00:00:00Z", updatedAt: "2026-01-01T00:00:00Z", binding: { transport: "postgres_direct", host: "db", port: 5432, username: "reader" }, connectionStatus: "reachable", schemaSyncedVersion: 4, metadataContentRevision: 9, preprocessingStatus: "running", } as const; function repository(overrides: Record = {}) { const finishPreprocessing = vi.fn(async () => ({ ...database, preprocessingStatus: "succeeded" as const, preprocessedMetadataRevision: 9, })); return { beginPreprocessing: vi.fn(async () => ({ kind: "started" as const, database })), finishPreprocessing, clearPreprocessing: vi.fn(async () => ({ kind: "cleared" as const, database })), getByWorkspace: vi.fn(async () => database), listTables: vi.fn(async () => [{ id: "table-1", databaseId: database.id, name: "orders", description: "Curated orders", generatedDescription: "Generated orders", sourceComment: "PostgreSQL orders", version: 1, createdAt: database.createdAt, updatedAt: database.updatedAt, lastSyncedDatabaseVersion: 4, lastSyncedAt: database.updatedAt, }]), listColumns: vi.fn(async () => [{ id: "column-1", tableId: "table-1", name: "customer_id", ordinalPosition: 1, dataType: "uuid", isNullable: false, defaultExpression: null, primaryKeyPosition: null, isPrimaryKey: false, isForeignKey: true, foreignKeyCount: 1, sourceComment: "PostgreSQL customer", description: null, generatedDescription: "Generated customer", sensitive: true, sensitivityReason: "identifier", lastSyncedDatabaseVersion: 4, lastSyncedAt: database.updatedAt, version: 1, createdAt: database.createdAt, updatedAt: database.updatedAt, }]), listRelationships: vi.fn(async () => []), listLogicalRelationships: vi.fn(async () => []), ...overrides, } as any; } function service(options: { repository?: any; runChild?: (request: ChildProcessRequest) => Promise<{ exitCode: number; stdout: string; stderr: string; }>; semanticPreflight?: () => Promise< { ok: true } | { ok: false; code: "semantic_index_incompatible" } >; } = {}) { const dataRoot = mkdtempSync(join(tmpdir(), "tht-catalog-preprocessing-")); roots.push(dataRoot); return new WorkspacePreprocessingService({ dataRoot, acquireActiveRuntime: async () => runtime(), runChild: options.runChild ?? (async () => ({ exitCode: 0, stdout: JSON.stringify({ counts: { tables: 1, columns: 1, relationships: 0 } }), stderr: "", })), semanticPreflight: options.semanticPreflight ?? (async () => ({ ok: true })), evidencePreflight: async () => ({ ok: true }), catalogRepository: options.repository, }); } test("complete preprocessing snapshots Catalog metadata and commits its PostgreSQL state", async () => { const catalog = repository(); const runChild = vi.fn(async (request: ChildProcessRequest) => { const snapshotPath = request.argv[request.argv.indexOf("--catalog-metadata") + 1]!; const snapshot = JSON.parse(readFileSync(snapshotPath, "utf8")); expect(snapshot).toMatchObject({ workspaceId: "catalog-workspace", databaseName: "warehouse", metadataContentRevision: 9, tables: [{ name: "orders", description: "Curated orders", descriptionSource: "curated", columns: [{ name: "customer_id", description: "Generated customer", descriptionSource: "generated", sensitive: true, }], }], }); return { exitCode: 0, stdout: JSON.stringify({ counts: { tables: 1, columns: 1, relationships: 0 } }), stderr: "", }; }); const result = await service({ repository: catalog, runChild }).run({ workspaceId: "catalog-workspace", }); expect(result).toMatchObject({ status: "succeeded", code: "ok", completedStages: ["catalog_snapshot", "catalog_metadata", "lsh", "schema_index"], }); expect(runChild).toHaveBeenCalledWith(expect.objectContaining({ argv: [ "preprocess", "catalog", "--catalog-metadata", expect.stringMatching(/\/catalog-metadata\.json$/), "--json", "-c", "/dev/fd/3", ], })); expect(catalog.finishPreprocessing).toHaveBeenCalledWith( "catalog-workspace", 9, "sha256:" + "f".repeat(64), { status: "succeeded" }, ); }); test("preprocessing fails closed when the PostgreSQL Catalog is unavailable", async () => { const result = await service().run({ workspaceId: "catalog-workspace" }); expect(result).toMatchObject({ status: "failed", code: "catalog_not_ready" }); }); test("preprocessing refuses a concurrent Catalog run without touching derived data", async () => { const catalog = repository({ beginPreprocessing: vi.fn(async () => ({ kind: "already_running" as const })), }); const runChild = vi.fn(); const result = await service({ repository: catalog, runChild }).run({ workspaceId: "catalog-workspace", }); expect(result).toMatchObject({ status: "failed", code: "preprocessing_conflict" }); expect(runChild).not.toHaveBeenCalled(); expect(catalog.finishPreprocessing).not.toHaveBeenCalled(); }); test("preprocessing records a failed PostgreSQL state when its hidden worker fails", async () => { const catalog = repository(); const instance = service({ repository: catalog, runChild: async () => ({ exitCode: 1, stdout: "", stderr: "private failure" }), }); await expect(instance.run({ workspaceId: "catalog-workspace" })).rejects.toThrow( "workspace child failed", ); expect(catalog.finishPreprocessing).toHaveBeenCalledWith( "catalog-workspace", 9, "sha256:" + "f".repeat(64), { status: "failed", errorCode: "schema_index_failed" }, ); }); test("semantic preflight failure is persisted before returning", async () => { const catalog = repository(); const result = await service({ repository: catalog, semanticPreflight: async () => ({ ok: false, code: "semantic_index_incompatible" }), }).run({ workspaceId: "catalog-workspace" }); expect(result).toMatchObject({ status: "failed", code: "semantic_index_incompatible" }); expect(catalog.finishPreprocessing).toHaveBeenCalledWith( "catalog-workspace", 9, "sha256:" + "f".repeat(64), { status: "failed", errorCode: "semantic_index_incompatible" }, ); }); test("clear invalidates Catalog readiness before clearing only derived worker data", async () => { const catalog = repository(); const runChild = vi.fn(async () => ({ exitCode: 0, stdout: JSON.stringify({ counts: { referenceCollections: 1, derivedPaths: 3 } }), stderr: "", })); const result = await service({ repository: catalog, runChild }).clear({ workspaceId: "catalog-workspace", }); expect(catalog.clearPreprocessing).toHaveBeenCalledWith("catalog-workspace"); expect(runChild).toHaveBeenCalledWith({ argv: ["preprocess", "clear", "--json", "-c", "/dev/fd/3"], configPath: expect.any(String), }); expect(result).toMatchObject({ status: "succeeded", code: "ok", operation: "preprocess clear", counts: { referenceCollections: 1, derivedPaths: 3 }, }); });