195 lines
9.2 KiB
TypeScript
195 lines
9.2 KiB
TypeScript
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<CatalogSchemaCapabilityUnavailableError>({
|
|
capability: "schema_snapshot",
|
|
});
|
|
});
|