import { mkdtempSync, rmSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { afterEach, expect, test, vi } from "vitest"; import type { CatalogDatabaseClient, CatalogPostgresAccess } from "../src/catalog/postgres-access.js"; import { ConcreteCatalogSchemaIntrospector } from "../src/catalog/schema-introspector.js"; import { CatalogSchemaCapabilityUnavailableError, type WorkspaceDatabase, } from "../src/catalog/types.js"; import { WorkspaceSecretStore } from "../src/workspaces/secret-store.js"; const roots: string[] = []; afterEach(() => { vi.unstubAllGlobals(); for (const root of roots.splice(0)) rmSync(root, { recursive: true, force: true }); }); function store() { const root = mkdtempSync(join(tmpdir(), "catalog-schema-introspection-secrets-")); const runtimeRoot = mkdtempSync(join(tmpdir(), "catalog-schema-introspection-runtime-")); roots.push(root, runtimeRoot); return new WorkspaceSecretStore({ root, runtimeRoot, installationId: "test" }); } function database(binding: WorkspaceDatabase["binding"]): WorkspaceDatabase { return { id: "11111111-1111-4111-8111-111111111111", workspaceId: "psd-clinical", engine: "postgres", databaseName: "warehouse", schema: "datawarehouse", binding, version: 4, connectionStatus: "reachable", testedVersion: 4, createdAt: "2026-08-27T08:00:00Z", updatedAt: "2026-08-27T09:00:00Z", }; } test("reads columns, ordered composite keys, and physical relationships from one PostgreSQL connection", async () => { const query = vi.fn() .mockResolvedValueOnce({ rows: [{ present: true }] }) .mockResolvedValueOnce({ rows: [{ name: "visits", source_comment: "Visits" }] }) .mockResolvedValueOnce({ rows: [ { table_name: "visits", name: "tenant_id", ordinal_position: 1, data_type: "uuid", is_nullable: false, default_expression: null, primary_key_position: 1, source_comment: null }, { table_name: "visits", name: "patient_id", ordinal_position: 2, data_type: "bigint", is_nullable: false, default_expression: null, primary_key_position: 2, source_comment: "Patient" }, ] }) .mockResolvedValueOnce({ rows: [ { constraint_name: "visits_patient_fkey", source_table_name: "visits", target_table_name: "patients", update_action: "a", delete_action: "c", deferrable: true, initially_deferred: false, position: 1, source_column_name: "tenant_id", target_column_name: "tenant_id" }, { constraint_name: "visits_patient_fkey", source_table_name: "visits", target_table_name: "patients", update_action: "a", delete_action: "c", deferrable: true, initially_deferred: false, position: 2, source_column_name: "patient_id", target_column_name: "id" }, ] }); const end = vi.fn(async () => undefined); const client: CatalogDatabaseClient = { query, end }; const postgres: CatalogPostgresAccess = { connect: vi.fn(async () => client) }; const introspector = new ConcreteCatalogSchemaIntrospector(postgres, store()); const result = await introspector.scan(database({ transport: "ssh_tunnel", username: "reader", sshHost: "bastion.internal", sshPort: 22, sshUsername: "tunnel", sshTargetHost: "db.internal", sshTargetPort: 5432, }), new AbortController().signal); expect(result.capabilities).toEqual({ tables: "available", columns: "available", relationships: "available" }); expect(result.columns).toMatchObject([ { name: "tenant_id", primaryKeyPosition: 1, isNullable: false }, { name: "patient_id", primaryKeyPosition: 2, sourceComment: "Patient" }, ]); expect(result.relationships).toEqual([expect.objectContaining({ constraintName: "visits_patient_fkey", updateRule: "NO ACTION", deleteRule: "CASCADE", deferrable: true, columns: [ { position: 1, sourceColumnName: "tenant_id", targetColumnName: "tenant_id" }, { position: 2, sourceColumnName: "patient_id", targetColumnName: "id" }, ], })]); expect(query.mock.calls[2][0]).toContain("format_type"); expect(query.mock.calls[3][0]).toContain("WITH ORDINALITY"); expect(query.mock.calls.slice(1).every((call) => call[1][0] === "datawarehouse")).toBe(true); expect(end).toHaveBeenCalledOnce(); }); test("uses the typed full REST snapshot RPC and preserves explicit capability unavailability", async () => { const response = { schemaVersion: 1, capabilities: { tables: "available", columns: "unavailable", relationships: "unavailable" }, tables: [{ name: "patients", sourceComment: null }], columns: [], relationships: [], }; const fetchMock = vi.fn(async () => new Response(JSON.stringify(response), { status: 200, headers: { "content-type": "application/json" } })); vi.stubGlobal("fetch", fetchMock); const postgres: CatalogPostgresAccess = { connect: vi.fn(async () => { throw new Error("wire access must not be used"); }) }; const introspector = new ConcreteCatalogSchemaIntrospector(postgres, store()); const result = await introspector.scan(database({ transport: "rest_api", baseUrl: "https://connector.internal/api/", restPath: "/health", restAuth: "none", }), new AbortController().signal); expect(result).toEqual(response); expect(fetchMock).toHaveBeenCalledWith( "https://connector.internal/api/rpc/schema_snapshot", expect.objectContaining({ method: "POST", body: JSON.stringify({ schema_name: "datawarehouse" }) }), ); }); test("falls back to one read-only REST query when the snapshot RPC is absent", async () => { const response = { schemaVersion: 1 as const, capabilities: { tables: "available" as const, columns: "available" as const, relationships: "available" as const }, tables: [ { name: "patients", sourceComment: "Clinical patients" }, { name: "visits", sourceComment: null }, ], columns: [ { tableName: "patients", name: "tenant_id", ordinalPosition: 1, dataType: "uuid", isNullable: false, defaultExpression: null, primaryKeyPosition: 1, sourceComment: "Tenant" }, { tableName: "patients", name: "id", ordinalPosition: 2, dataType: "bigint", isNullable: false, defaultExpression: "nextval('patients_id_seq'::regclass)", primaryKeyPosition: 2, sourceComment: null }, { tableName: "visits", name: "tenant_id", ordinalPosition: 1, dataType: "uuid", isNullable: false, defaultExpression: null, primaryKeyPosition: null, sourceComment: null }, { tableName: "visits", name: "patient_id", ordinalPosition: 2, dataType: "bigint", isNullable: true, defaultExpression: null, primaryKeyPosition: null, sourceComment: "Owning patient" }, ], relationships: [{ constraintName: "visits_patient_fkey", sourceTableName: "visits", targetTableName: "patients", updateRule: "CASCADE", deleteRule: "RESTRICT", deferrable: true, initiallyDeferred: false, columns: [ { position: 1, sourceColumnName: "tenant_id", targetColumnName: "tenant_id" }, { position: 2, sourceColumnName: "patient_id", targetColumnName: "id" }, ], }], }; const fetchMock = vi.fn() .mockResolvedValueOnce(new Response(null, { status: 404 })) .mockResolvedValueOnce(new Response(JSON.stringify([response]), { status: 200, headers: { "content-type": "application/json" }, })); vi.stubGlobal("fetch", fetchMock); const postgres: CatalogPostgresAccess = { connect: vi.fn(async () => { throw new Error("wire access must not be used"); }), }; const introspector = new ConcreteCatalogSchemaIntrospector(postgres, store()); const result = await introspector.scan(database({ transport: "rest_api", baseUrl: "https://connector.internal/api/", restPath: "/health", restAuth: "none", }), new AbortController().signal); expect(result).toEqual(response); expect(fetchMock).toHaveBeenCalledTimes(2); expect(fetchMock.mock.calls[0]).toEqual([ "https://connector.internal/api/rpc/schema_snapshot", expect.objectContaining({ method: "POST", body: JSON.stringify({ schema_name: "datawarehouse" }) }), ]); expect(fetchMock.mock.calls[1][0]).toBe("https://connector.internal/api/rpc/run_query"); const fallbackRequest = fetchMock.mock.calls[1][1] as RequestInit; expect(fallbackRequest).toMatchObject({ method: "POST" }); const fallbackBody = JSON.parse(String(fallbackRequest.body)) as { query_text: string }; expect(Object.keys(fallbackBody)).toEqual(["query_text"]); expect(fallbackBody.query_text).toMatch(/^\s*WITH\b/); expect(fallbackBody.query_text).toContain("pg_catalog.pg_constraint"); expect(fallbackBody.query_text).not.toMatch(/\b(INSERT|UPDATE|DROP|ALTER|CREATE|TRUNCATE)\b/i); }); test("classifies a missing REST snapshot RPC as an explicit binding capability", async () => { vi.stubGlobal("fetch", vi.fn(async () => new Response(null, { status: 404 }))); const postgres: CatalogPostgresAccess = { connect: vi.fn(async () => { throw new Error("wire access must not be used"); }), }; const introspector = new ConcreteCatalogSchemaIntrospector(postgres, store()); const scan = introspector.scan(database({ transport: "rest_api", baseUrl: "https://connector.internal/api/", restPath: "/health", restAuth: "none", }), new AbortController().signal); await expect(scan).rejects.toMatchObject({ capability: "schema_snapshot", }); });