feat: protect sensitive catalog samples

This commit is contained in:
Codex
2026-08-30 12:14:23 +02:00
parent 6278ee9d81
commit 0736983bc5
28 changed files with 1162 additions and 128 deletions
+7
View File
@@ -56,6 +56,7 @@ import { metadataGenerationModelRoutes } from "./routes/metadata-generation-mode
import { catalogDescriptionConsolidationRoutes } from "./routes/catalog-description-consolidation.js";
import { PythonModelCompleter, type ModelCompleter } from "./catalog/model-completer.js";
import { DescriptionGenerationWorker } from "./catalog/description-generation-worker.js";
import { SensitiveDataSuggester } from "./catalog/sensitive-data-suggester.js";
import {
PostgresDescriptionSourceSampler,
type DescriptionSourceSampler,
@@ -179,6 +180,11 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
catalogOperationCoordinator,
descriptionSourceSampler,
);
const sensitiveDataSuggester = new SensitiveDataSuggester(
catalogRepository,
metadataGenerationModels,
modelCompleter,
);
const catalogService = deps?.catalogService ?? new CatalogService(
catalogRepository,
workspaceRegistry,
@@ -438,6 +444,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
catalogDescriptionGenerationRoutes(app, {
repository: catalogRepository,
worker: descriptionGenerationWorker,
sensitiveDataSuggester,
});
settingsRoutes(app, { cfg: config, listModels, getSettings });
piManagementRoutes(app, { service: piManagement });
@@ -12,6 +12,7 @@ import type {
DescriptionSourceSampleValue,
DescriptionTargetSourceSample,
} from "./description-source-sampler.js";
import { syntheticSampleValue } from "./synthetic-sample-value.js";
import {
DescriptionGenerationRunActiveError,
type CatalogColumn,
@@ -306,35 +307,55 @@ function sourceSampleFor(
): PromptSourceSample | undefined {
const targetId = target.kind === "column" ? target.column.id : target.table.id;
const sample = samples.find((candidate) => candidate.targetId === targetId);
if (!sample) return undefined;
const relevantColumns = new Set(
target.kind === "column" ? [target.column.name] : target.columns.map((column) => column.name),
);
const rows = sample.rows.slice(0, budget.rows).map((row) => {
const columns = target.kind === "column" ? [target.column] : target.columns;
const sourceRows = sample?.rows ?? [];
const hasSensitiveColumns = columns.some((column) => column.sensitive);
const hasNonSensitiveColumns = columns.some((column) => !column.sensitive);
const realRowCount = hasNonSensitiveColumns
? Math.min(sourceRows.length, budget.rows)
: 0;
const rowCount = hasSensitiveColumns ? MAX_SAMPLE_ROWS_PER_REQUEST : realRowCount;
const rows = Array.from({ length: rowCount }, (_, rowIndex) => {
const row = rowIndex < realRowCount ? sourceRows[rowIndex] : undefined;
const sourceFields = new Map((row?.fields ?? []).map((field) => [field.name, field.value]));
const fields: Array<{ name: string; value: PromptSampleValue }> = [];
const seen = new Set<string>();
for (const field of row.fields) {
if (!relevantColumns.has(field.name) || seen.has(field.name)) continue;
seen.add(field.name);
for (const column of columns) {
const value = column.sensitive
? syntheticSampleValue(column, rowIndex + 1)
: sourceFields.get(column.name);
if (value === undefined) continue;
fields.push({
name: boundedJsonText(field.name, MAX_IDENTIFIER_JSON_BYTES),
value: promptSampleValue(field.value),
name: boundedJsonText(column.name, MAX_IDENTIFIER_JSON_BYTES),
value: promptSampleValue(value),
});
if (fields.length === MAX_SAMPLE_FIELDS_PER_ROW) break;
}
return { fields };
});
budget.rows -= rows.length;
budget.rows -= realRowCount;
const representativeValues: PromptSourceSample["representativeValues"] = [];
const seenColumns = new Set<string>();
let remainingRepresentativeValues = budget.representativeValues;
for (const examples of sample.representativeValues) {
if (remainingRepresentativeValues === 0) break;
if (!relevantColumns.has(examples.column) || seenColumns.has(examples.column)) continue;
seenColumns.add(examples.column);
let remainingRealRepresentativeValues = budget.representativeValues;
let remainingSyntheticRepresentativeValues = MAX_REPRESENTATIVE_VALUES_PER_REQUEST;
for (const column of columns) {
if (seenColumns.has(column.name)) continue;
const remainingRepresentativeValues = column.sensitive
? remainingSyntheticRepresentativeValues
: remainingRealRepresentativeValues;
if (remainingRepresentativeValues === 0) continue;
const sourceExamples = sample?.representativeValues.find(
(examples) => examples.column === column.name,
);
const exampleValues = column.sensitive
? Array.from(
{ length: Math.min(Math.max(rowCount, 1), remainingRepresentativeValues) },
(_, index) => syntheticSampleValue(column, index + 1),
)
: sourceExamples?.values ?? [];
seenColumns.add(column.name);
const values: Array<Exclude<PromptSampleValue, null>> = [];
const seenValues = new Set<string>();
for (const value of examples.values) {
for (const value of exampleValues) {
const normalized = promptSampleValue(value);
if (normalized === null) continue;
const key = JSON.stringify([typeof normalized, normalized]);
@@ -345,11 +366,15 @@ function sourceSampleFor(
}
if (values.length > 0) {
representativeValues.push({
column: boundedJsonText(examples.column, MAX_IDENTIFIER_JSON_BYTES),
column: boundedJsonText(column.name, MAX_IDENTIFIER_JSON_BYTES),
values,
});
remainingRepresentativeValues -= values.length;
budget.representativeValues -= values.length;
if (column.sensitive) {
remainingSyntheticRepresentativeValues -= values.length;
} else {
remainingRealRepresentativeValues -= values.length;
budget.representativeValues -= values.length;
}
}
if (representativeValues.length === MAX_SAMPLE_COLUMNS) break;
}
@@ -933,8 +958,9 @@ export class DescriptionGenerationWorker {
targetId: target.kind === "column" ? target.column.id : target.table.id,
tableName: target.table.name,
columnNames: target.kind === "column"
? [target.column.name]
: target.columns.map((column) => column.name),
? target.column.sensitive ? [] : [target.column.name]
: target.columns.filter((column) => !column.sensitive)
.map((column) => column.name),
})),
signal,
);
+10 -1
View File
@@ -202,10 +202,18 @@ export class MemoryCatalogRepository implements CatalogRepository {
expectedVersion: number,
description: string | null,
generatedDescription: string | null,
sensitive?: boolean,
): Promise<CatalogColumn | undefined> {
const current = await this.getColumn(databaseId, tableId, columnId);
if (!current || current.version !== expectedVersion) return undefined;
const updated = { ...current, description, generatedDescription, version: current.version + 1, updatedAt: new Date().toISOString() };
const updated = {
...current,
description,
generatedDescription,
sensitive: sensitive ?? current.sensitive,
version: current.version + 1,
updatedAt: new Date().toISOString(),
};
this.columns.set(columnId, updated);
return structuredClone(updated);
}
@@ -621,6 +629,7 @@ export class MemoryCatalogRepository implements CatalogRepository {
sourceComment: observed.sourceComment,
description: null,
generatedDescription: null,
sensitive: false,
lastSyncedDatabaseVersion: expectedDatabaseVersion,
lastSyncedAt: now,
version: 1,
+2
View File
@@ -8,6 +8,7 @@ import * as catalogTablesMigration from "./migrations/002_catalog_tables.js";
import * as catalogSchemaSyncMigration from "./migrations/003_catalog_schema_sync.js";
import * as catalogRuntimeSequencePrivilegesMigration from "./migrations/004_catalog_runtime_sequence_privileges.js";
import * as descriptionGenerationRunsMigration from "./migrations/005_description_generation_runs.js";
import * as sensitiveDataFlagMigration from "./migrations/006_sensitive_data_flag.js";
const connectionString = process.env.THT_CATALOG_MIGRATOR_DATABASE_URL;
const host = process.env.THT_CATALOG_DB_HOST;
@@ -38,6 +39,7 @@ const provider: MigrationProvider = {
"003_catalog_schema_sync": catalogSchemaSyncMigration,
"004_catalog_runtime_sequence_privileges": catalogRuntimeSequencePrivilegesMigration,
"005_description_generation_runs": descriptionGenerationRunsMigration,
"006_sensitive_data_flag": sensitiveDataFlagMigration,
};
},
};
@@ -0,0 +1,12 @@
import type { Kysely } from "kysely";
import type { CatalogDatabase } from "../repository.js";
export async function up(db: Kysely<CatalogDatabase>): Promise<void> {
await db.schema.alterTable("catalog_columns")
.addColumn("sensitive", "boolean", (column) => column.notNull().defaultTo(false))
.execute();
}
export async function down(db: Kysely<CatalogDatabase>): Promise<void> {
await db.schema.alterTable("catalog_columns").dropColumn("sensitive").execute();
}
+4
View File
@@ -108,6 +108,7 @@ interface CatalogColumnTable {
sourceComment: string | null;
description: string | null;
generatedDescription: string | null;
sensitive: Generated<boolean>;
lastSyncedDatabaseVersion: number | null;
lastSyncedAt: Timestamp | null;
version: Generated<number>;
@@ -290,6 +291,7 @@ function serializeColumn(row: Selectable<CatalogColumnTable>, foreignKeyCount =
sourceComment: row.sourceComment,
description: row.description,
generatedDescription: row.generatedDescription,
sensitive: row.sensitive,
lastSyncedDatabaseVersion: row.lastSyncedDatabaseVersion,
lastSyncedAt: row.lastSyncedAt === null ? null : new Date(row.lastSyncedAt).toISOString(),
version: row.version,
@@ -561,6 +563,7 @@ export class KyselyCatalogRepository implements CatalogRepository {
expectedVersion: number,
description: string | null,
generatedDescription: string | null,
sensitive?: boolean,
): Promise<CatalogColumn | undefined> {
const belongs = await this.db.selectFrom("catalogTables").select("id")
.where("id", "=", tableId).where("databaseId", "=", databaseId).executeTakeFirst();
@@ -568,6 +571,7 @@ export class KyselyCatalogRepository implements CatalogRepository {
const row = await this.db.updateTable("catalogColumns").set({
description,
generatedDescription,
...(sensitive === undefined ? {} : { sensitive }),
version: sql`version + 1`,
updatedAt: sql`now()`,
}).where("id", "=", columnId).where("tableId", "=", tableId)
@@ -0,0 +1,103 @@
import { z } from "zod";
import type { MetadataGenerationModels } from "./metadata-generation-models.js";
import type { ModelCompleter } from "./model-completer.js";
import type { CatalogRepository } from "./types.js";
const MAX_COLUMNS = 10_000;
const responseSchema = z.object({
suggestions: z.array(z.object({
columnId: z.uuid(),
sensitive: z.boolean(),
}).strict()).max(MAX_COLUMNS),
}).strict();
export interface SensitiveDataSuggestion {
columnId: string;
sensitive: boolean;
}
export class SensitiveDataSuggestionTargetNotFoundError extends Error {
constructor() {
super("database not found");
this.name = "SensitiveDataSuggestionTargetNotFoundError";
}
}
export class SensitiveDataSuggestionInvalidResponseError extends Error {
constructor() {
super("sensitive-data suggestion response is invalid");
this.name = "SensitiveDataSuggestionInvalidResponseError";
}
}
export class SensitiveDataSuggester {
constructor(
private readonly repository: CatalogRepository,
private readonly models: MetadataGenerationModels,
private readonly completer: ModelCompleter,
) {}
async suggest(
databaseId: string,
modelId: string,
signal: AbortSignal,
): Promise<readonly SensitiveDataSuggestion[]> {
const database = await this.repository.get(databaseId);
if (!database) throw new SensitiveDataSuggestionTargetNotFoundError();
const tables = await this.repository.listTables(databaseId);
const columns = (await Promise.all(tables.map(async (table) => ({
table,
columns: await this.repository.listColumns(databaseId, table.id),
})))).flatMap(({ table, columns: tableColumns }) => tableColumns.map((column) => ({
columnId: column.id,
table: table.name,
column: column.name,
dataType: column.dataType,
nullable: column.isNullable,
primaryKey: column.isPrimaryKey,
foreignKey: column.isForeignKey,
})));
if (columns.length === 0) return [];
if (columns.length > MAX_COLUMNS) throw new SensitiveDataSuggestionInvalidResponseError();
const content = await this.completer.complete({
model: this.models.resolve(modelId),
signal,
messages: [
{
role: "system",
content: [
"Classify whether each database column is likely to contain sensitive source values.",
"Use only the supplied structural metadata. Return strict JSON with this exact shape:",
'{"suggestions":[{"columnId":"uuid","sensitive":true}]}',
"Return every supplied column exactly once. Do not add explanations or markdown.",
].join("\n"),
},
{
role: "user",
content: JSON.stringify({
database: database.databaseName,
schema: database.schema,
columns,
}),
},
],
});
try {
const parsed = responseSchema.parse(JSON.parse(content));
const expected = new Set(columns.map((column) => column.columnId));
const received = new Set(parsed.suggestions.map((suggestion) => suggestion.columnId));
if (received.size !== parsed.suggestions.length
|| received.size !== expected.size
|| [...received].some((columnId) => !expected.has(columnId))) {
throw new SensitiveDataSuggestionInvalidResponseError();
}
return parsed.suggestions;
} catch (error) {
if (error instanceof SensitiveDataSuggestionInvalidResponseError) throw error;
throw new SensitiveDataSuggestionInvalidResponseError();
}
}
}
@@ -0,0 +1,67 @@
import type { DescriptionSourceSampleValue } from "./description-source-sampler.js";
const FIRST_NAMES = ["marta", "luca", "elena", "paolo", "giulia"] as const;
const LAST_NAMES = ["rossi", "bianchi", "conti", "romano", "ferrari"] as const;
const NUMERIC_TYPE = /(int|numeric|decimal|real|double|float|money)/;
const BOOLEAN_TYPE = /(bool)/;
function nameAt(index: number): string {
const offset = Math.max(0, index - 1);
return `${FIRST_NAMES[offset % FIRST_NAMES.length]} ${LAST_NAMES[offset % LAST_NAMES.length]}`;
}
export function syntheticSampleValue(
column: { name: string; dataType: string },
index: number,
): DescriptionSourceSampleValue {
const ordinal = Math.max(1, index);
const name = column.name.toLocaleLowerCase("en-US");
const type = column.dataType.toLocaleLowerCase("en-US");
const person = nameAt(ordinal).split(" ");
const safeName = name.replace(/[^a-z0-9]+/g, "_").replace(/^_+|_+$/g, "") || "value";
if (type.endsWith("[]") || type.startsWith("_") || /\barray\b/.test(type)) {
if (NUMERIC_TYPE.test(type)) {
return `{${1000 + ordinal},${1001 + ordinal}}`;
}
if (BOOLEAN_TYPE.test(type)) return `{${ordinal % 2 === 1},${ordinal % 2 !== 1}}`;
return `{${safeName}_${String(ordinal).padStart(3, "0")},${safeName}_${String(ordinal + 1).padStart(3, "0")}}`;
}
if (/^jsonb?$/.test(type)) {
return JSON.stringify({ example: `${safeName}_${String(ordinal).padStart(3, "0")}`, sequence: ordinal });
}
if (/(uuid|uniqueidentifier)/.test(type)) {
return `00000000-0000-4000-8000-${String(ordinal).padStart(12, "0")}`;
}
if (/\bcidr\b/.test(type)) return "192.0.2.0/24";
if (/\binet\b/.test(type)) return `192.0.2.${((ordinal - 1) % 254) + 1}`;
if (/(timestamp|datetime)/.test(type)) {
return `2024-01-${String(Math.min(ordinal, 28)).padStart(2, "0")}T10:30:00.000Z`;
}
if (/\bdate\b/.test(type)) {
return `198${ordinal % 10}-01-${String(Math.min(ordinal, 28)).padStart(2, "0")}`;
}
if (/\btime\b/.test(type)) {
return `10:30:${String(ordinal % 60).padStart(2, "0")}`;
}
if (/\binterval\b/.test(type)) {
return `${ordinal} days ${String(ordinal % 24).padStart(2, "0")}:00:00`;
}
if (/\bbytea\b/.test(type)) return `\\x${ordinal.toString(16).padStart(8, "0")}`;
if (BOOLEAN_TYPE.test(type)) return ordinal % 2 === 1;
if (NUMERIC_TYPE.test(type)) return 1000 + ordinal;
if (/e[-_]?mail/.test(name)) {
return `${person[0]}.${person[1]}@example.com`;
}
if (/(phone|mobile|cell|telefono|telefono_mobile|tel_)/.test(name)) {
return `+39 02 5550 ${String(1000 + ordinal).padStart(4, "0")}`;
}
if (/(first_?name|given_?name|nome)/.test(name)) return person[0]!;
if (/(last_?name|family_?name|surname|cognome)/.test(name)) return person[1]!;
if (/(full_?name|patient_?name|person_?name)/.test(name)) return nameAt(ordinal);
if (/(birth|dob|data_nascita)/.test(name)) {
return `198${ordinal % 10}-01-${String(Math.min(ordinal, 28)).padStart(2, "0")}`;
}
return `${safeName}_${String(ordinal).padStart(3, "0")}`;
}
+2
View File
@@ -88,6 +88,7 @@ export interface CatalogColumn {
sourceComment: string | null;
description: string | null;
generatedDescription: string | null;
sensitive: boolean;
lastSyncedDatabaseVersion: number | null;
lastSyncedAt: string | null;
version: number;
@@ -349,6 +350,7 @@ export interface CatalogRepository {
expectedVersion: number,
description: string | null,
generatedDescription: string | null,
sensitive?: boolean,
): Promise<CatalogColumn | undefined>;
consolidateGeneratedDescriptions(
databaseId: string,
@@ -11,6 +11,12 @@ import {
type DescriptionGenerationWorker,
} from "../catalog/description-generation-worker.js";
import { MetadataGenerationModelUnavailableError } from "../catalog/metadata-generation-models.js";
import { ModelCompletionProviderError } from "../catalog/model-completer.js";
import {
SensitiveDataSuggester,
SensitiveDataSuggestionInvalidResponseError,
SensitiveDataSuggestionTargetNotFoundError,
} from "../catalog/sensitive-data-suggester.js";
import {
CatalogOperationInProgressError,
CatalogUnavailableError,
@@ -22,6 +28,7 @@ import {
const idSchema = z.uuid();
const modelIdSchema = z.string().regex(/^[a-z][a-z0-9._-]{0,63}$/);
const suggestionSchema = z.object({ modelId: modelIdSchema }).strict();
const selectedTargetIdsSchema = z.array(idSchema).min(1);
const startSchema = z.discriminatedUnion("scope", [
z.object({
@@ -116,6 +123,19 @@ function safeError(reply: FastifyReply, error: unknown) {
message: "The selected metadata-generation model is unavailable.",
});
}
if (error instanceof SensitiveDataSuggestionTargetNotFoundError) {
return reply.code(404).send({
code: "database_not_found",
message: "Database configuration was not found.",
});
}
if (error instanceof SensitiveDataSuggestionInvalidResponseError
|| error instanceof ModelCompletionProviderError) {
return reply.code(502).send({
code: "sensitive_data_suggestion_failed",
message: "Sensitive-data suggestions could not be prepared.",
});
}
if (error instanceof DescriptionGenerationDuplicateTargetIdsError) {
return reply.code(400).send({
code: "description_generation_target_ids_duplicate",
@@ -172,8 +192,28 @@ function safeError(reply: FastifyReply, error: unknown) {
export function catalogDescriptionGenerationRoutes(
app: FastifyInstance,
deps: { repository: CatalogRepository; worker: DescriptionGenerationWorker },
deps: {
repository: CatalogRepository;
worker: DescriptionGenerationWorker;
sensitiveDataSuggester: SensitiveDataSuggester;
},
): void {
app.post("/catalog/databases/:databaseId/sensitive-data-suggestions", async (request, reply) => {
if (!manage(request, reply)) return reply;
try {
const databaseId = idSchema.parse((request.params as { databaseId?: unknown }).databaseId);
const input = suggestionSchema.parse(request.body);
const suggestions = await deps.sensitiveDataSuggester.suggest(
databaseId,
input.modelId,
new AbortController().signal,
);
return { suggestions };
} catch (error) {
return safeError(reply, error);
}
});
app.post("/catalog/databases/:databaseId/description-generation-runs", async (request, reply) => {
if (!manage(request, reply)) return reply;
try {
+2
View File
@@ -18,6 +18,7 @@ const metadataSchema = z.object({
version: z.number().int().positive(),
description: z.string().max(20_000).nullable(),
generatedDescription: z.string().max(20_000).nullable(),
sensitive: z.boolean().optional(),
}).strict();
const createRunSchema = z.object({
version: z.number().int().positive(),
@@ -117,6 +118,7 @@ export function catalogSchemaRoutes(
input.version,
normalized(input.description),
normalized(input.generatedDescription),
input.sensitive,
);
if (!updated) return reply.code(409).send({ code: "column_stale", message: "Column metadata changed. Reload and try again." });
return updated;
@@ -141,6 +141,83 @@ async function waitForTerminalRun(app: ReturnType<typeof buildApp>, runId: strin
throw new Error(`Description Generation Run ${runId} did not finish`);
}
test("suggests sensitive flags from structural metadata without persisting them", async () => {
const modelCompleter = {
complete: vi.fn(async () => JSON.stringify({
suggestions: [{ columnId: expect.any(String), sensitive: true }],
})),
};
const { app, repository, database, column } = await setup(modelCompleter);
modelCompleter.complete.mockResolvedValueOnce(JSON.stringify({
suggestions: [{ columnId: column.id, sensitive: true }],
}));
try {
const response = await app.inject({
method: "POST",
url: `/catalog/databases/${database.id}/sensitive-data-suggestions`,
payload: { modelId: configuredModel.id },
});
expect(response.statusCode).toBe(200);
expect(response.json()).toEqual({
suggestions: [{ columnId: column.id, sensitive: true }],
});
expect(await repository.getColumn(database.id, column.tableId, column.id))
.toMatchObject({ sensitive: false });
const request = modelCompleter.complete.mock.calls[0]![0] as ModelCompletionRequest;
const prompt = request.messages.map((message) => message.content).join("\n");
expect(prompt).toContain("patients");
expect(prompt).toContain("birth_date");
expect(prompt).toContain("date");
expect(prompt).not.toContain("Patient date of birth");
expect(prompt).not.toContain("test-provider-secret");
} finally {
await app.close();
}
});
test.each(["malformed", "incomplete", "duplicate"] as const)(
"fails safely when sensitive-data suggestions are %s",
async (kind) => {
const modelCompleter: ModelCompleter = {
complete: vi.fn(async () => "unused"),
};
const { app, repository, database, column } = await setup(modelCompleter);
const rawResponse = kind === "malformed"
? "RAW_PROVIDER_RESPONSE_DO_NOT_EXPOSE_{"
: kind === "incomplete"
? JSON.stringify({ suggestions: [] })
: JSON.stringify({
suggestions: [
{ columnId: column.id, sensitive: true },
{ columnId: column.id, sensitive: true },
],
});
vi.mocked(modelCompleter.complete).mockResolvedValueOnce(rawResponse);
try {
const response = await app.inject({
method: "POST",
url: `/catalog/databases/${database.id}/sensitive-data-suggestions`,
payload: { modelId: configuredModel.id },
});
expect(response.statusCode).toBe(502);
expect(response.json()).toEqual({
code: "sensitive_data_suggestion_failed",
message: "Sensitive-data suggestions could not be prepared.",
});
expect(response.body).not.toContain(rawResponse);
expect(await repository.getColumn(database.id, column.tableId, column.id))
.toMatchObject({ sensitive: false });
} finally {
await app.close();
}
},
);
interface SseFrame {
id?: string;
event?: string;
@@ -398,6 +475,102 @@ test("keeps real source samples transient across the Fastify API and application
}
});
test("never exposes a protected source value to the model, persistence, logs, or browser APIs", async () => {
let selectedColumnId = "";
const protectedValue = "PROTECTED_SOURCE_VALUE_8f4c2a";
const modelCompleter: ModelCompleter = {
complete: vi.fn(async () => JSON.stringify({
results: [{
targetId: selectedColumnId,
outcome: "generated",
description: "Data di nascita del paziente.",
}],
})),
};
const descriptionSourceSampler: DescriptionSourceSampler = {
sample: vi.fn(async (_database, targets) => [{
targetId: targets[0]!.targetId,
tableName: targets[0]!.tableName,
rows: [{ fields: [{ name: "birth_date", value: protectedValue }] }],
representativeValues: [{ column: "birth_date", values: [protectedValue] }],
}]),
};
const { app, repository, database, table, column } = await setup(
modelCompleter,
{},
"it",
descriptionSourceSampler,
);
selectedColumnId = column.id;
await repository.updateColumnMetadata(
database.id,
table.id,
column.id,
column.version,
column.description,
column.generatedDescription,
true,
);
const logSpies = [
vi.spyOn(app.log, "info"),
vi.spyOn(app.log, "warn"),
vi.spyOn(app.log, "error"),
];
try {
const start = await app.inject({
method: "POST",
url: `/catalog/databases/${database.id}/description-generation-runs`,
payload: {
modelId: configuredModel.id,
scope: "selected_columns",
targetIds: [column.id],
},
});
expect(start.statusCode).toBe(202);
await waitForTerminalRun(app, start.json().id);
expect(descriptionSourceSampler.sample).toHaveBeenCalledWith(
expect.anything(),
[expect.objectContaining({ targetId: column.id, columnNames: [] })],
expect.any(AbortSignal),
);
const completionRequest = vi.mocked(modelCompleter.complete).mock.calls[0]![0];
const providerPayload = JSON.stringify(completionRequest.messages);
expect(providerPayload).not.toContain(protectedValue);
expect(providerPayload).toContain("1981-01-01");
const apiResponses = await Promise.all([
app.inject({ method: "GET", url: `/catalog/description-generation-runs/${start.json().id}` }),
app.inject({
method: "GET",
url: `/catalog/description-generation-runs/${start.json().id}/events-list`,
}),
app.inject({ method: "GET", url: `/catalog/databases/${database.id}` }),
app.inject({ method: "GET", url: `/catalog/databases/${database.id}/tables` }),
app.inject({
method: "GET",
url: `/catalog/databases/${database.id}/tables/${table.id}/columns`,
}),
]);
expect(apiResponses.every((response) => response.statusCode === 200)).toBe(true);
expect(apiResponses.map((response) => response.body).join("\n")).not.toContain(protectedValue);
const persisted = JSON.stringify({
run: await repository.getDescriptionGenerationRun(start.json().id),
events: await repository.listDescriptionGenerationEvents(start.json().id),
database: await repository.get(database.id),
table: await repository.getTable(database.id, table.id),
column: await repository.getColumn(database.id, table.id, column.id),
});
expect(persisted).not.toContain(protectedValue);
expect(JSON.stringify(logSpies.flatMap((spy) => spy.mock.calls))).not.toContain(protectedValue);
} finally {
for (const spy of logSpies) spy.mockRestore();
await app.close();
}
});
test("wires the production sampler to the same injected CatalogPostgresAccess instance", async () => {
let selectedColumnId = "";
const sampleSecret = "PRODUCTION_WIRING_SAMPLE_f2986a";
@@ -222,7 +222,7 @@ test("adds only bounded transient source samples to the model request", async ()
tables: [{ name: "patients", sourceComment: null }],
columns: [{
tableName: "patients",
name: "status",
name: "patient_email",
ordinalPosition: 1,
dataType: "text",
isNullable: true,
@@ -243,9 +243,18 @@ test("adds only bounded transient source samples to the model request", async ()
});
const table = (await repository.listTables(database.id))[0]!;
const columns = await repository.listColumns(database.id, table.id);
const column = columns.find((candidate) => candidate.name === "status")!;
const column = columns.find((candidate) => candidate.name === "patient_email")!;
const ward = columns.find((candidate) => candidate.name === "ward")!;
const sampleSecret = "ONLY_IN_TRANSIENT_SAMPLE_7f29c8";
await repository.updateColumnMetadata(
database.id,
table.id,
column.id,
column.version,
column.description,
column.generatedDescription,
true,
);
const sampleSecret = "real.patient@hospital.invalid";
const sourceSampler: DescriptionSourceSampler = {
sample: vi.fn(async () => [{
targetId: column.id,
@@ -265,11 +274,14 @@ test("adds only bounded transient source samples to the model request", async ()
rows: [
{ fields: [{ name: ward.name, value: "row-4" }] },
{ fields: [{ name: ward.name, value: "row-5" }] },
{ fields: [{ name: ward.name, value: "row-6-must-be-omitted" }] },
{ fields: [{ name: ward.name, value: "row-6" }] },
{ fields: [{ name: ward.name, value: "row-7" }] },
{ fields: [{ name: ward.name, value: "row-8" }] },
{ fields: [{ name: ward.name, value: "row-9-must-be-omitted" }] },
],
representativeValues: [{
column: ward.name,
values: ["ward-1", "ward-2", "ward-3-must-be-omitted"],
values: ["ward-1", "ward-2", "ward-3", "ward-4", "ward-5", "ward-6-must-be-omitted"],
}],
}]),
};
@@ -321,7 +333,7 @@ test("adds only bounded transient source samples to the model request", async ()
expect(sourceSampler.sample).toHaveBeenCalledWith(
expect.objectContaining({ id: database.id, binding: database.binding }),
[
{ targetId: column.id, tableName: table.name, columnNames: [column.name] },
{ targetId: column.id, tableName: table.name, columnNames: [] },
{ targetId: ward.id, tableName: table.name, columnNames: [ward.name] },
],
expect.any(AbortSignal),
@@ -338,21 +350,29 @@ test("adds only bounded transient source samples to the model request", async ()
targetContext.sourceSample?.representativeValues.flatMap((entry) => entry.values) ?? []
),
);
expect(sampledRows).toHaveLength(5);
expect(representativeValues).toHaveLength(5);
expect(context.targets[0].sourceSample.rows).toHaveLength(3);
expect(context.targets[1].sourceSample.rows).toHaveLength(2);
expect(sampledRows).toHaveLength(10);
expect(representativeValues).toHaveLength(10);
expect(context.targets[0].sourceSample.rows).toHaveLength(5);
expect(context.targets[1].sourceSample.rows).toHaveLength(5);
expect(context.targets[0].sourceSample.representativeValues).toEqual([{
column: column.name,
values: [sampleSecret, "two", "three"],
values: [
"marta.rossi@example.com",
"luca.bianchi@example.com",
"elena.conti@example.com",
"paolo.romano@example.com",
"giulia.ferrari@example.com",
],
}]);
expect(context.targets[1].sourceSample.representativeValues).toEqual([{
column: ward.name,
values: ["ward-1", "ward-2"],
values: ["ward-1", "ward-2", "ward-3", "ward-4", "ward-5"],
}]);
expect(userMessage).toContain(sampleSecret);
expect(userMessage).not.toContain(sampleSecret);
expect(userMessage).not.toMatch(/synthetic|fake|fittizi/i);
expect(userMessage).toContain("marta.rossi@example.com");
expect(userMessage).not.toMatch(
/row-6-must-be-omitted|ward-3-must-be-omitted/,
/row-9-must-be-omitted|ward-6-must-be-omitted/,
);
const persisted = JSON.stringify({
@@ -366,6 +386,173 @@ test("adds only bounded transient source samples to the model request", async ()
expect(persisted).not.toContain(sampleSecret);
});
test("gives every sensitive column synthetic context without consuming the real sample budget", async () => {
const repository = new MemoryCatalogRepository();
const database = await repository.create({
workspaceId: "psd-clinical",
engine: "postgres",
databaseName: "warehouse",
schema: "datawarehouse",
binding: { transport: "postgres_direct", host: "db.internal", port: 5432, username: "reader" },
});
await repository.applySchemaSync(database.id, database.version, "all", [], {
schemaVersion: 1,
capabilities: { tables: "available", columns: "available", relationships: "available" },
tables: [{ name: "patients", sourceComment: null }],
columns: [{
tableName: "patients",
name: "patient_email",
ordinalPosition: 1,
dataType: "text",
isNullable: true,
defaultExpression: null,
primaryKeyPosition: null,
sourceComment: null,
}, {
tableName: "patients",
name: "patient_phone",
ordinalPosition: 2,
dataType: "text",
isNullable: true,
defaultExpression: null,
primaryKeyPosition: null,
sourceComment: null,
}, {
tableName: "patients",
name: "ward",
ordinalPosition: 3,
dataType: "text",
isNullable: true,
defaultExpression: null,
primaryKeyPosition: null,
sourceComment: null,
}],
relationships: [],
});
const table = (await repository.listTables(database.id))[0]!;
const columns = await repository.listColumns(database.id, table.id);
const email = columns.find((column) => column.name === "patient_email")!;
const phone = columns.find((column) => column.name === "patient_phone")!;
const ward = columns.find((column) => column.name === "ward")!;
for (const column of [email, phone]) {
await repository.updateColumnMetadata(
database.id,
table.id,
column.id,
column.version,
column.description,
column.generatedDescription,
true,
);
}
const wardValues = ["ward-a", "ward-b", "ward-c", "ward-d", "ward-e"];
const sourceSampler: DescriptionSourceSampler = {
sample: vi.fn(async (_database, targets) => targets.map((target) => {
if (target.columnNames.length === 0) {
return {
targetId: target.targetId,
tableName: target.tableName,
rows: [],
representativeValues: [],
};
}
const columnName = target.columnNames[0]!;
return {
targetId: target.targetId,
tableName: target.tableName,
rows: wardValues.map((value) => ({ fields: [{ name: columnName, value }] })),
representativeValues: [{ column: columnName, values: wardValues }],
};
})),
};
const completer: ModelCompleter = {
complete: vi.fn(async (request) => {
const context = JSON.parse(request.messages[1]!.content.split("\n").slice(1).join("\n"));
return JSON.stringify({
results: context.targets.map((target: { targetId: string }) => ({
targetId: target.targetId,
outcome: "generated",
description: "Descrizione generata.",
})),
});
}),
};
const models: MetadataGenerationModels = {
catalog: () => ({ models: [{ id: "openai-mini", label: "OpenAI Mini" }], default: "openai-mini" }),
resolve: () => ({
id: "openai-mini",
provider: "openai",
model: "gpt-4.1-mini",
apiKeyEnv: "OPENAI_API_KEY",
apiKey: "test-provider-secret",
}),
};
const worker = new DescriptionGenerationWorker(
repository,
{
read: vi.fn(async () => ({
workspace: { workspace: { language: "it" } },
revision: {},
})),
} as unknown as WorkspaceRegistry,
models,
completer,
new CatalogOperationCoordinator(),
sourceSampler,
);
const run = await worker.start(
database.id,
"openai-mini",
"selected_columns",
[email.id, phone.id, ward.id],
);
await worker.waitForRun(run.id);
expect(sourceSampler.sample).toHaveBeenCalledWith(
expect.objectContaining({ id: database.id }),
[
{ targetId: email.id, tableName: table.name, columnNames: [] },
{ targetId: phone.id, tableName: table.name, columnNames: [] },
{ targetId: ward.id, tableName: table.name, columnNames: [ward.name] },
],
expect.any(AbortSignal),
);
const request = vi.mocked(completer.complete).mock.calls[0]![0] as ModelCompletionRequest;
const context = JSON.parse(request.messages[1]!.content.split("\n").slice(1).join("\n"));
const targets = new Map(
context.targets.map((target: { targetId: string }) => [target.targetId, target]),
);
expect(targets.get(email.id)).toMatchObject({
sourceSample: {
rows: expect.arrayContaining([
{ fields: [{ name: email.name, value: "marta.rossi@example.com" }] },
]),
representativeValues: [{
column: email.name,
values: expect.arrayContaining(["marta.rossi@example.com"]),
}],
},
});
expect(targets.get(phone.id)).toMatchObject({
sourceSample: {
rows: expect.arrayContaining([
{ fields: [{ name: phone.name, value: "+39 02 5550 1001" }] },
]),
representativeValues: [{
column: phone.name,
values: expect.arrayContaining(["+39 02 5550 1001"]),
}],
},
});
expect(targets.get(ward.id)).toMatchObject({
sourceSample: {
rows: wardValues.map((value) => ({ fields: [{ name: ward.name, value }] })),
representativeValues: [{ column: ward.name, values: wardValues }],
},
});
});
test("continues metadata-only with one safe warning when source sampling is unavailable", async () => {
const repository = new MemoryCatalogRepository();
const database = await repository.create({
@@ -11,6 +11,7 @@ import { up as upDatabases } from "../src/catalog/migrations/001_workspace_datab
import { up as upTables } from "../src/catalog/migrations/002_catalog_tables.js";
import { up as upSchemaSync } from "../src/catalog/migrations/003_catalog_schema_sync.js";
import { up as upDescriptionGeneration } from "../src/catalog/migrations/005_description_generation_runs.js";
import { up as upSensitiveDataFlag } from "../src/catalog/migrations/006_sensitive_data_flag.js";
import { KyselyCatalogRepository, type CatalogDatabase } from "../src/catalog/repository.js";
import { loadConfig } from "../src/config.js";
import type { WorkspaceRegistry } from "../src/workspaces/registry.js";
@@ -43,6 +44,7 @@ test.skipIf(!dockerAvailable)("Fastify persists Description Generation success a
await upDatabases(db);
await upTables(db);
await upSchemaSync(db);
await upSensitiveDataFlag(db);
await upDescriptionGeneration(db);
const repository = new KyselyCatalogRepository(db);
const database = await repository.create({
@@ -92,6 +92,34 @@ test("samples at most five source rows and five distinct non-null examples in a
expect(end).toHaveBeenCalledOnce();
});
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 PostgresDescriptionSourceSampler(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");
@@ -10,6 +10,7 @@ import { up as upTables } from "../src/catalog/migrations/002_catalog_tables.js"
import { up as upSchemaSync } from "../src/catalog/migrations/003_catalog_schema_sync.js";
import { up as upRuntimeSequencePrivileges } from "../src/catalog/migrations/004_catalog_runtime_sequence_privileges.js";
import { up as upDescriptionGeneration } from "../src/catalog/migrations/005_description_generation_runs.js";
import { up as upSensitiveDataFlag } from "../src/catalog/migrations/006_sensitive_data_flag.js";
const dockerAvailable = spawnSync("docker", ["info"], { stdio: "ignore" }).status === 0;
@@ -23,6 +24,7 @@ test.skipIf(!dockerAvailable)("PostgreSQL migration enforces one database per wo
await upDatabases(db);
await upTables(db);
await upSchemaSync(db);
await upSensitiveDataFlag(db);
await sql`CREATE ROLE thothii_catalog_runtime`.execute(db);
await upRuntimeSequencePrivileges(db);
const sequencePrivilege = await sql<{ allowed: boolean }>`
@@ -75,11 +77,47 @@ test.skipIf(!dockerAvailable)("PostgreSQL migration enforces one database per wo
};
expect(await repository.applySchemaSync(created.id, 1, "columns", [], fullColumnsSnapshot))
.toMatchObject({ created: 4, deleted: 0 });
expect((await repository.listColumns(created.id, patients.id)).map((column) => column.name))
.toEqual(["id", "name"]);
expect((await repository.listColumns(created.id, patients.id)).map((column) => ({
name: column.name,
sensitive: column.sensitive,
}))).toEqual([
{ name: "id", sensitive: false },
{ name: "name", sensitive: false },
]);
expect((await repository.listColumns(created.id, visits.id)).map((column) => column.name))
.toEqual(["id", "patient_id"]);
const patientName = (await repository.listColumns(created.id, patients.id))
.find((column) => column.name === "name")!;
expect(await repository.updateColumnMetadata(
created.id,
patients.id,
patientName.id,
patientName.version,
patientName.description,
patientName.generatedDescription,
true,
)).toMatchObject({ sensitive: true });
const refreshedColumnsSnapshot: ObservedSchemaSnapshot = {
...fullColumnsSnapshot,
schemaVersion: 2,
columns: fullColumnsSnapshot.columns.map((column) => column.tableName === "patients"
&& column.name === "name"
? { ...column, sourceComment: "Sensitive patient name" }
: column),
};
expect(await repository.applySchemaSync(
created.id,
1,
"columns",
[patients.id],
refreshedColumnsSnapshot,
)).toMatchObject({ updated: 1 });
expect(await repository.getColumn(created.id, patients.id, patientName.id)).toMatchObject({
sensitive: true,
sourceComment: "Sensitive patient name",
});
const reducedColumnsSnapshot: ObservedSchemaSnapshot = {
...fullColumnsSnapshot,
columns: fullColumnsSnapshot.columns.filter((column) => column.name === "id"),
@@ -135,6 +173,7 @@ test.skipIf(!dockerAvailable)("PostgreSQL repository performs scoped metadata cl
await upDatabases(db);
await upTables(db);
await upSchemaSync(db);
await upSensitiveDataFlag(db);
const repository = new KyselyCatalogRepository(db);
const database = await repository.create({
workspaceId: "cleanup-test",
@@ -212,6 +251,7 @@ test.skipIf(!dockerAvailable)("PostgreSQL repository atomically consolidates sel
await upDatabases(db);
await upTables(db);
await upSchemaSync(db);
await upSensitiveDataFlag(db);
const repository = new KyselyCatalogRepository(db);
const database = await repository.create({
workspaceId: "consolidation-test",
@@ -310,6 +350,7 @@ test.skipIf(!dockerAvailable)("PostgreSQL repository persists globally exclusive
await upDatabases(db);
await upTables(db);
await upSchemaSync(db);
await upSensitiveDataFlag(db);
await upDescriptionGeneration(db);
const repository = new KyselyCatalogRepository(db);
const firstDatabase = await repository.create({
+17 -3
View File
@@ -234,16 +234,30 @@ test("keeps generated descriptions editable and preserves them across synchroniz
});
expect(editedTable.json()).toMatchObject({ description: null, generatedDescription: "Generated table draft" });
const idColumn = (await repository.listColumns(database.id, patients.id))[0];
expect(idColumn.sensitive).toBe(false);
const editedColumn = await app.inject({
method: "PATCH", url: `/catalog/databases/${database.id}/tables/${patients.id}/columns/${idColumn.id}`,
payload: { version: idColumn.version, description: "Reviewed key", generatedDescription: "Generated key draft" },
payload: {
version: idColumn.version,
description: "Reviewed key",
generatedDescription: "Generated key draft",
sensitive: true,
},
});
expect(editedColumn.json()).toMatchObject({
description: "Reviewed key",
generatedDescription: "Generated key draft",
sensitive: true,
});
expect(editedColumn.json()).toMatchObject({ description: "Reviewed key", generatedDescription: "Generated key draft" });
const second = await app.inject({ method: "POST", url: `/catalog/databases/${database.id}/sync-runs`, payload: { version: database.version, scope: "all", tableIds: [] } });
await waitFor(repository, second.json().id, "succeeded");
expect(await repository.getTable(database.id, patients.id)).toMatchObject({ generatedDescription: "Generated table draft" });
expect(await repository.getColumn(database.id, patients.id, idColumn.id)).toMatchObject({ description: "Reviewed key", generatedDescription: "Generated key draft" });
expect(await repository.getColumn(database.id, patients.id, idColumn.id)).toMatchObject({
description: "Reviewed key",
generatedDescription: "Generated key draft",
sensitive: true,
});
});
test("consolidates non-empty generated table descriptions and reports skipped selections", async () => {
@@ -0,0 +1,29 @@
import { expect, test } from "vitest";
import { syntheticSampleValue } from "../src/catalog/synthetic-sample-value.js";
test.each([
["text[]", "tags", 1, "{tags_001,tags_002}"],
["json", "payload", 2, '{"example":"payload_002","sequence":2}'],
["jsonb", "attributes", 3, '{"example":"attributes_003","sequence":3}'],
["inet", "client_ip", 4, "192.0.2.4"],
["cidr", "network", 5, "192.0.2.0/24"],
["time without time zone", "opening_time", 6, "10:30:06"],
["interval", "duration", 7, "7 days 07:00:00"],
["bytea", "digest", 8, "\\x00000008"],
] as const)(
"creates a deterministic PostgreSQL-shaped value for %s",
(dataType, name, index, expected) => {
const column = { name, dataType };
expect(syntheticSampleValue(column, index)).toBe(expected);
expect(syntheticSampleValue(column, index)).toBe(expected);
},
);
test("uses the PostgreSQL type before birth-related name hints", () => {
expect(syntheticSampleValue({ name: "birth_year", dataType: "integer" }, 2)).toBe(1002);
expect(syntheticSampleValue({ name: "birth_date", dataType: "date" }, 2))
.toBe("1982-01-02");
expect(syntheticSampleValue({ name: "birth_recorded_at", dataType: "timestamp" }, 2))
.toBe("2024-01-02T10:30:00.000Z");
});