203 lines
7.4 KiB
TypeScript
203 lines
7.4 KiB
TypeScript
import { expect, test, vi } from "vitest";
|
|
import { mkdtempSync, rmSync, writeFileSync } from "node:fs";
|
|
import { tmpdir } from "node:os";
|
|
import { join } from "node:path";
|
|
import {
|
|
ConcreteDescriptionSourceSampler,
|
|
type DescriptionSourceSamplingTarget,
|
|
} from "../src/catalog/description-source-sampler.js";
|
|
import type {
|
|
CatalogDatabaseClient,
|
|
CatalogPostgresAccess,
|
|
} from "../src/catalog/postgres-access.js";
|
|
import type { WorkspaceDatabase } from "../src/catalog/types.js";
|
|
import type { WorkspaceSecretStore } from "../src/workspaces/secret-store.js";
|
|
import { CATALOG_SECRET_IDS } from "../src/catalog/secrets.js";
|
|
|
|
const database: WorkspaceDatabase = {
|
|
id: "11111111-1111-4111-8111-111111111111",
|
|
workspaceId: "psd-clinical",
|
|
engine: "postgres",
|
|
databaseName: "warehouse",
|
|
schema: 'clinical"data',
|
|
version: 1,
|
|
createdAt: "2026-08-28T08:00:00Z",
|
|
updatedAt: "2026-08-28T08:00:00Z",
|
|
connectionStatus: "reachable",
|
|
binding: {
|
|
transport: "postgres_direct",
|
|
host: "db.internal",
|
|
port: 5432,
|
|
username: "reader",
|
|
},
|
|
};
|
|
|
|
const target: DescriptionSourceSamplingTarget = {
|
|
targetId: "22222222-2222-4222-8222-222222222222",
|
|
tableName: 'patient"facts',
|
|
columnNames: ['status"code', "ward"],
|
|
};
|
|
|
|
test("samples at most five source rows and five distinct non-null examples in a read-only transaction", async () => {
|
|
const query = vi.fn(async (sql: string) => {
|
|
if (!sql.startsWith("SELECT")) return { rows: [] };
|
|
return {
|
|
rows: [
|
|
{ 'status"code': "active", ward: null },
|
|
{ 'status"code': "pending", ward: "A" },
|
|
{ 'status"code': "closed", ward: "A" },
|
|
{ 'status"code': "transferred", ward: "B" },
|
|
{ 'status"code': "unknown", ward: "C" },
|
|
{ 'status"code': "must-not-be-sampled", ward: "D" },
|
|
],
|
|
};
|
|
});
|
|
const end = vi.fn(async () => undefined);
|
|
const access: CatalogPostgresAccess = {
|
|
connect: vi.fn(async () => ({ query, end }) as CatalogDatabaseClient),
|
|
};
|
|
const sampler = new ConcreteDescriptionSourceSampler(access);
|
|
const controller = new AbortController();
|
|
|
|
const samples = await sampler.sample(database, [target], controller.signal);
|
|
|
|
expect(samples).toEqual([{
|
|
targetId: target.targetId,
|
|
tableName: target.tableName,
|
|
rows: [
|
|
{ fields: [{ name: 'status"code', value: "active" }, { name: "ward", value: null }] },
|
|
{ fields: [{ name: 'status"code', value: "pending" }, { name: "ward", value: "A" }] },
|
|
{ fields: [{ name: 'status"code', value: "closed" }, { name: "ward", value: "A" }] },
|
|
{ fields: [{ name: 'status"code', value: "transferred" }, { name: "ward", value: "B" }] },
|
|
{ fields: [{ name: 'status"code', value: "unknown" }, { name: "ward", value: "C" }] },
|
|
],
|
|
representativeValues: [
|
|
{
|
|
column: 'status"code',
|
|
values: ["active", "pending", "closed", "transferred"],
|
|
},
|
|
{ column: "ward", values: ["A"] },
|
|
],
|
|
}]);
|
|
expect(access.connect).toHaveBeenCalledWith(database, controller.signal);
|
|
expect(samples[0]!.representativeValues.flatMap((entry) => entry.values)).toHaveLength(5);
|
|
expect(query.mock.calls).toEqual([
|
|
["BEGIN TRANSACTION READ ONLY", []],
|
|
[
|
|
'SELECT LEFT(("status""code")::text, $1) AS "status""code", LEFT(("ward")::text, $1) AS "ward" FROM "clinical""data"."patient""facts" LIMIT $2',
|
|
[256, 5],
|
|
],
|
|
["ROLLBACK", []],
|
|
]);
|
|
expect(query.mock.calls.map(([sql]) => String(sql).split(" ")[0])).toEqual([
|
|
"BEGIN",
|
|
"SELECT",
|
|
"ROLLBACK",
|
|
]);
|
|
expect(end).toHaveBeenCalledOnce();
|
|
});
|
|
|
|
test("samples source rows through the configured REST run_query binding", async () => {
|
|
const root = mkdtempSync(join(tmpdir(), "tht-source-rest-"));
|
|
const credentialFile = join(root, "api-key");
|
|
writeFileSync(credentialFile, "test-api-key\n", { mode: 0o600 });
|
|
const release = vi.fn();
|
|
const secretStore = {
|
|
materialize: vi.fn(() => ({
|
|
files: new Map([[CATALOG_SECRET_IDS.apiKey, credentialFile]]),
|
|
release,
|
|
})),
|
|
} as unknown as WorkspaceSecretStore;
|
|
const fetchMock = vi.fn(async () => new Response(JSON.stringify([
|
|
{ 'status"code': "active", ward: null },
|
|
{ 'status"code': "pending", ward: "A" },
|
|
]), { status: 200, headers: { "content-type": "application/json" } }));
|
|
vi.stubGlobal("fetch", fetchMock);
|
|
const access: CatalogPostgresAccess = {
|
|
connect: vi.fn(async () => { throw new Error("PostgreSQL access must not be used"); }),
|
|
};
|
|
const sampler = new ConcreteDescriptionSourceSampler(access, secretStore);
|
|
const restDatabase: WorkspaceDatabase = {
|
|
...database,
|
|
binding: {
|
|
transport: "rest_api",
|
|
baseUrl: "https://dwh.example.test/root/",
|
|
restPath: "/health",
|
|
restAuth: "x-api-key",
|
|
},
|
|
};
|
|
|
|
try {
|
|
await expect(sampler.sample(restDatabase, [target], new AbortController().signal)).resolves.toEqual([{
|
|
targetId: target.targetId,
|
|
tableName: target.tableName,
|
|
rows: [
|
|
{ fields: [{ name: 'status"code', value: "active" }, { name: "ward", value: null }] },
|
|
{ fields: [{ name: 'status"code', value: "pending" }, { name: "ward", value: "A" }] },
|
|
],
|
|
representativeValues: [
|
|
{ column: 'status"code', values: ["active", "pending"] },
|
|
{ column: "ward", values: ["A"] },
|
|
],
|
|
}]);
|
|
expect(access.connect).not.toHaveBeenCalled();
|
|
expect(fetchMock).toHaveBeenCalledWith("https://dwh.example.test/root/rpc/run_query", expect.objectContaining({
|
|
method: "POST",
|
|
headers: { "content-type": "application/json", "x-api-key": "test-api-key" },
|
|
body: JSON.stringify({
|
|
query_text: 'SELECT LEFT(("status""code")::text, 256) AS "status""code", LEFT(("ward")::text, 256) AS "ward" FROM "clinical""data"."patient""facts" LIMIT 5',
|
|
}),
|
|
}));
|
|
expect(release).toHaveBeenCalledOnce();
|
|
} finally {
|
|
vi.unstubAllGlobals();
|
|
rmSync(root, { recursive: true, force: true });
|
|
}
|
|
});
|
|
|
|
test("does not issue a SELECT when a protected target has no source columns", async () => {
|
|
const query = vi.fn(async () => ({ rows: [] }));
|
|
const end = vi.fn(async () => undefined);
|
|
const access: CatalogPostgresAccess = {
|
|
connect: vi.fn(async () => ({ query, end }) as CatalogDatabaseClient),
|
|
};
|
|
const sampler = new ConcreteDescriptionSourceSampler(access);
|
|
|
|
const samples = await sampler.sample(database, [{
|
|
targetId: target.targetId,
|
|
tableName: target.tableName,
|
|
columnNames: [],
|
|
}], new AbortController().signal);
|
|
|
|
expect(samples).toEqual([{
|
|
targetId: target.targetId,
|
|
tableName: target.tableName,
|
|
rows: [],
|
|
representativeValues: [],
|
|
}]);
|
|
expect(query.mock.calls).toEqual([
|
|
["BEGIN TRANSACTION READ ONLY", []],
|
|
["ROLLBACK", []],
|
|
]);
|
|
expect(query.mock.calls.some(([sql]) => String(sql).startsWith("SELECT"))).toBe(false);
|
|
expect(end).toHaveBeenCalledOnce();
|
|
});
|
|
|
|
test("rolls back and closes the source connection when sampling fails", async () => {
|
|
const query = vi.fn(async (sql: string) => {
|
|
if (sql.startsWith("SELECT")) throw new Error("distinctive-source-secret");
|
|
return { rows: [] };
|
|
});
|
|
const end = vi.fn(async () => undefined);
|
|
const access: CatalogPostgresAccess = {
|
|
connect: vi.fn(async () => ({ query, end }) as CatalogDatabaseClient),
|
|
};
|
|
const sampler = new ConcreteDescriptionSourceSampler(access);
|
|
const controller = new AbortController();
|
|
|
|
await expect(sampler.sample(database, [target], controller.signal)).rejects.toThrow();
|
|
|
|
expect(query).toHaveBeenCalledWith("ROLLBACK", []);
|
|
expect(end).toHaveBeenCalledOnce();
|
|
});
|