feat: complete catalog sensitivity enhancements

This commit is contained in:
Codex
2026-09-04 15:11:18 +02:00
parent b891246664
commit 7b1d69a65b
72 changed files with 33866 additions and 358 deletions
+11 -1
View File
@@ -33,7 +33,11 @@ import { ReadinessManager } from "./runtime/readiness-manager.js";
import { MaintenanceBarrier } from "./runtime/maintenance-gate.js";
import { WorkspaceRegistry } from "./workspaces/registry.js";
import { createProductionWorkspaceDiagnoser } from "./workspaces/diagnostics.js";
import { workspaceRoutes, type WorkspaceDiagnoser } from "./routes/workspaces.js";
import {
workspaceRoutes,
type WorkspaceDatabaseTester,
type WorkspaceDiagnoser,
} from "./routes/workspaces.js";
import { piManagementRoutes } from "./routes/pi-management.js";
import { supportsSessionRuntime } from "./workspaces/bindings.js";
import { resolveRuntimeBindingsWithWorkspaceSecrets } from "./workspaces/secret-requirements.js";
@@ -83,6 +87,7 @@ export interface BuildAppDeps {
hub?: SseHub;
workspaceRegistry?: WorkspaceRegistry;
workspaceDiagnoser?: WorkspaceDiagnoser;
workspaceDatabaseTester?: WorkspaceDatabaseTester;
workspaceSecretStore?: WorkspaceSecretStore;
catalogRepository?: CatalogRepository;
catalogService?: CatalogService;
@@ -230,6 +235,10 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
catalogPostgresAccess,
catalogOperationCoordinator,
);
const workspaceDatabaseTester = deps?.workspaceDatabaseTester ?? (async (workspaceId: string) => {
const database = await catalogRepository.getByWorkspace(workspaceId);
return database ? catalogService.test(database) : undefined;
});
const catalogTableService = deps?.catalogTableService ?? new CatalogTableService(catalogRepository);
const catalogLogicalRelationshipService = deps?.catalogLogicalRelationshipService
?? new CatalogLogicalRelationshipService(catalogRepository);
@@ -485,6 +494,7 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
diagnose: workspaceDiagnoser,
authDiagnoser,
secretStore: workspaceSecretStore,
testDatabaseConnection: workspaceDatabaseTester,
});
catalogDatabaseRoutes(app, { repository: catalogRepository, service: catalogService, operations: catalogOperationCoordinator });
catalogTableRoutes(app, {
+46 -10
View File
@@ -83,8 +83,10 @@ export class MemoryCatalogRepository implements CatalogRepository {
const tableIds = new Set(tables.map((table) => table.id));
const columns = [...this.columns.values()]
.filter((column) => tableIds.has(column.tableId));
const relationships = [...this.relationships.values()]
.filter((relationship) => selectedDatabaseIds.has(relationship.databaseId));
const relationships = [
...this.relationships.values(),
...this.logicalRelationships.values(),
].filter((relationship) => selectedDatabaseIds.has(relationship.databaseId));
return createCatalogMetrics(databaseId, {
tables: tables.length,
@@ -255,6 +257,7 @@ export class MemoryCatalogRepository implements CatalogRepository {
description: string | null,
generatedDescription: string | null,
sensitive?: boolean,
sensitivityReason?: string | null,
): Promise<CatalogColumn | undefined> {
const current = await this.getColumn(databaseId, tableId, columnId);
if (!current || current.version !== expectedVersion) return undefined;
@@ -263,6 +266,9 @@ export class MemoryCatalogRepository implements CatalogRepository {
description,
generatedDescription,
sensitive: sensitive ?? current.sensitive,
sensitivityReason: sensitive === false
? null
: sensitivityReason === undefined ? current.sensitivityReason : sensitivityReason,
version: current.version + 1,
updatedAt: new Date().toISOString(),
};
@@ -276,10 +282,26 @@ export class MemoryCatalogRepository implements CatalogRepository {
targetIds: readonly string[],
): Promise<CatalogDescriptionConsolidationCounts | undefined> {
const selectedTargetIds = [...new Set(targetIds)];
if (!this.records.has(databaseId) || selectedTargetIds.length === 0) {
if (!this.records.has(databaseId)) {
return undefined;
}
const now = new Date().toISOString();
if (target === "database_columns") {
const targets = [...this.columns.values()].filter((column) => (
this.tables.get(column.tableId)?.databaseId === databaseId
));
const copied = targets.filter((column) => Boolean(column.generatedDescription?.trim()));
for (const column of copied) {
this.columns.set(column.id, {
...column,
description: column.generatedDescription,
version: column.version + 1,
updatedAt: now,
});
}
return { copied: copied.length, skipped: targets.length - copied.length };
}
if (selectedTargetIds.length === 0) return undefined;
if (target === "tables") {
const targets = selectedTargetIds.map((id) => this.tables.get(id));
if (targets.some((table) => !table || table.databaseId !== databaseId)) return undefined;
@@ -687,18 +709,22 @@ export class MemoryCatalogRepository implements CatalogRepository {
const tables = [...this.tables.values()].filter((table) => selected.has(table.databaseId));
const tableIds = new Set(tables.map((table) => table.id));
const columns = [...this.columns.values()].filter((column) => tableIds.has(column.tableId));
const relationships = [...this.relationships.values()]
const physicalRelationships = [...this.relationships.values()]
.filter((relationship) => selected.has(relationship.databaseId));
const logicalRelationships = [...this.logicalRelationships.values()]
.filter((relationship) => selected.has(relationship.databaseId));
const relationshipCount = physicalRelationships.length + logicalRelationships.length;
if (target === "tables") {
for (const table of tables) this.deleteTable(table.id);
this.markCatalogIncomplete(selectedDatabaseIds);
return { tables: tables.length, columns: columns.length, relationships: relationships.length };
return { tables: tables.length, columns: columns.length, relationships: relationshipCount };
}
for (const relationship of relationships) this.relationships.delete(relationship.id);
for (const relationship of physicalRelationships) this.relationships.delete(relationship.id);
for (const relationship of logicalRelationships) this.logicalRelationships.delete(relationship.id);
for (const databaseId of selectedDatabaseIds) this.refreshForeignKeyFlags(databaseId);
this.markCatalogIncomplete(selectedDatabaseIds);
return { tables: 0, columns: 0, relationships: relationships.length };
return { tables: 0, columns: 0, relationships: relationshipCount };
}
async deleteTableMetadata(
@@ -736,14 +762,23 @@ export class MemoryCatalogRepository implements CatalogRepository {
return { tables: 0, columns: columns.length, relationships: 0 };
}
const relationships = [...this.relationships.values()].filter((relationship) => (
const physicalRelationships = [...this.relationships.values()].filter((relationship) => (
relationship.databaseId === databaseId
&& (selected.has(relationship.sourceTableId) || selected.has(relationship.targetTableId))
));
for (const relationship of relationships) this.relationships.delete(relationship.id);
const logicalRelationships = [...this.logicalRelationships.values()].filter((relationship) => (
relationship.databaseId === databaseId
&& (selected.has(relationship.sourceTableId) || selected.has(relationship.targetTableId))
));
for (const relationship of physicalRelationships) this.relationships.delete(relationship.id);
for (const relationship of logicalRelationships) this.logicalRelationships.delete(relationship.id);
this.refreshForeignKeyFlags(databaseId);
this.markCatalogIncomplete([databaseId]);
return { tables: 0, columns: 0, relationships: relationships.length };
return {
tables: 0,
columns: 0,
relationships: physicalRelationships.length + logicalRelationships.length,
};
}
async planSchemaSync(
@@ -927,6 +962,7 @@ export class MemoryCatalogRepository implements CatalogRepository {
description: null,
generatedDescription: null,
sensitive: false,
sensitivityReason: null,
lastSyncedDatabaseVersion: expectedDatabaseVersion,
lastSyncedAt: now,
version: 1,
+2
View File
@@ -14,6 +14,7 @@ import * as catalogLogicalRelationshipsMigration from "./migrations/008_catalog_
import * as aiTokenUsageMigration from "./migrations/009_ai_token_usage.js";
import * as canonicalModelIdsMigration from "./migrations/010_canonical_model_ids.js";
import * as localSensitivityAnalysisMigration from "./migrations/011_local_sensitivity_analysis.js";
import * as sensitivityReasonMigration from "./migrations/012_sensitivity_reason.js";
const connectionString = process.env.THT_CATALOG_MIGRATOR_DATABASE_URL;
const host = process.env.THT_CATALOG_DB_HOST;
@@ -50,6 +51,7 @@ const provider: MigrationProvider = {
"009_ai_token_usage": aiTokenUsageMigration,
"010_canonical_model_ids": canonicalModelIdsMigration,
"011_local_sensitivity_analysis": localSensitivityAnalysisMigration,
"012_sensitivity_reason": sensitivityReasonMigration,
};
},
};
@@ -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("sensitivity_reason", "text")
.execute();
}
export async function down(db: Kysely<CatalogDatabase>): Promise<void> {
await db.schema.alterTable("catalog_columns").dropColumn("sensitivity_reason").execute();
}
+74 -16
View File
@@ -118,6 +118,7 @@ interface CatalogColumnTable {
description: string | null;
generatedDescription: string | null;
sensitive: Generated<boolean>;
sensitivityReason: Generated<string | null>;
lastSyncedDatabaseVersion: number | null;
lastSyncedAt: Timestamp | null;
version: Generated<number>;
@@ -360,6 +361,7 @@ function serializeColumn(row: Selectable<CatalogColumnTable>, foreignKeyCount =
description: row.description,
generatedDescription: row.generatedDescription,
sensitive: row.sensitive,
sensitivityReason: row.sensitivityReason,
lastSyncedDatabaseVersion: row.lastSyncedDatabaseVersion,
lastSyncedAt: row.lastSyncedAt === null ? null : new Date(row.lastSyncedAt).toISOString(),
version: row.version,
@@ -516,10 +518,18 @@ export class KyselyCatalogRepository implements CatalogRepository {
relationship_metrics AS (
SELECT
count(*)::int AS relationships,
max(catalog_relationships.updated_at) AS updated_at
FROM catalog_relationships
INNER JOIN selected_databases
ON selected_databases.id = catalog_relationships.database_id
max(relationship.updated_at) AS updated_at
FROM (
SELECT catalog_relationships.updated_at
FROM catalog_relationships
INNER JOIN selected_databases
ON selected_databases.id = catalog_relationships.database_id
UNION ALL
SELECT catalog_logical_relationships.updated_at
FROM catalog_logical_relationships
INNER JOIN selected_databases
ON selected_databases.id = catalog_logical_relationships.database_id
) AS relationship
)
SELECT
(SELECT count(*)::int FROM selected_databases) AS "databaseCount",
@@ -728,6 +738,7 @@ export class KyselyCatalogRepository implements CatalogRepository {
description: string | null,
generatedDescription: string | null,
sensitive?: boolean,
sensitivityReason?: string | null,
): Promise<CatalogColumn | undefined> {
const belongs = await this.db.selectFrom("catalogTables").select("id")
.where("id", "=", tableId).where("databaseId", "=", databaseId).executeTakeFirst();
@@ -736,6 +747,9 @@ export class KyselyCatalogRepository implements CatalogRepository {
description,
generatedDescription,
...(sensitive === undefined ? {} : { sensitive }),
...(sensitive === false
? { sensitivityReason: null }
: sensitivityReason === undefined ? {} : { sensitivityReason }),
version: sql`version + 1`,
updatedAt: sql`now()`,
}).where("id", "=", columnId).where("tableId", "=", tableId)
@@ -749,7 +763,7 @@ export class KyselyCatalogRepository implements CatalogRepository {
targetIds: readonly string[],
): Promise<CatalogDescriptionConsolidationCounts | undefined> {
const selectedTargetIds = [...new Set(targetIds)];
if (selectedTargetIds.length === 0) return undefined;
if (target !== "database_columns" && selectedTargetIds.length === 0) return undefined;
return await this.db.transaction().execute(async (trx) => {
const database = await trx.selectFrom("workspaceDatabases").select("id")
.where("id", "=", databaseId).forUpdate().executeTakeFirst();
@@ -784,24 +798,30 @@ export class KyselyCatalogRepository implements CatalogRepository {
const rows = tableRows.length === 0 ? [] : await trx.selectFrom("catalogColumns")
.select(["id", "generatedDescription"])
.where("tableId", "in", tableRows.map((table) => table.id))
.where("id", "in", selectedTargetIds)
.$if(target === "columns", (query) => query.where("id", "in", selectedTargetIds))
.orderBy("id")
.forUpdate()
.execute();
if (rows.length !== selectedTargetIds.length) return undefined;
if (target === "columns" && rows.length !== selectedTargetIds.length) return undefined;
const copiedIds = rows
.filter((row) => Boolean(row.generatedDescription?.trim()))
.map((row) => row.id);
if (copiedIds.length > 0) {
await trx.updateTable("catalogColumns").set({
let update = trx.updateTable("catalogColumns").set({
description: sql`generated_description`,
version: sql`version + 1`,
updatedAt: sql`now()`,
}).where("id", "in", copiedIds).execute();
});
update = target === "database_columns"
? update
.where("tableId", "in", tableRows.map((table) => table.id))
.where(sql<boolean>`nullif(btrim(generated_description), '') is not null`)
: update.where("id", "in", copiedIds);
await update.execute();
}
return {
copied: copiedIds.length,
skipped: selectedTargetIds.length - copiedIds.length,
skipped: rows.length - copiedIds.length,
};
});
}
@@ -1261,16 +1281,25 @@ export class KyselyCatalogRepository implements CatalogRepository {
if (databases.length !== selectedDatabaseIds.length) return undefined;
if (target === "relationships") {
const count = await trx.selectFrom("catalogRelationships")
const physicalCount = await trx.selectFrom("catalogRelationships")
.select(sql<number>`count(*)::int`.as("count"))
.where("databaseId", "in", selectedDatabaseIds).executeTakeFirst();
const logicalCount = await trx.selectFrom("catalogLogicalRelationships")
.select(sql<number>`count(*)::int`.as("count"))
.where("databaseId", "in", selectedDatabaseIds).executeTakeFirst();
await trx.deleteFrom("catalogLogicalRelationships")
.where("databaseId", "in", selectedDatabaseIds).execute();
await trx.deleteFrom("catalogRelationships")
.where("databaseId", "in", selectedDatabaseIds).execute();
await trx.updateTable("workspaceDatabases").set({
schemaSyncedVersion: null,
schemaSyncedAt: null,
}).where("id", "in", selectedDatabaseIds).execute();
return { tables: 0, columns: 0, relationships: Number(count?.count ?? 0) };
return {
tables: 0,
columns: 0,
relationships: Number(physicalCount?.count ?? 0) + Number(logicalCount?.count ?? 0),
};
}
const tableCount = await trx.selectFrom("catalogTables")
@@ -1280,7 +1309,10 @@ export class KyselyCatalogRepository implements CatalogRepository {
.innerJoin("catalogTables", "catalogTables.id", "catalogColumns.tableId")
.select(sql<number>`count(*)::int`.as("count"))
.where("catalogTables.databaseId", "in", selectedDatabaseIds).executeTakeFirst();
const relationshipCount = await trx.selectFrom("catalogRelationships")
const physicalRelationshipCount = await trx.selectFrom("catalogRelationships")
.select(sql<number>`count(*)::int`.as("count"))
.where("databaseId", "in", selectedDatabaseIds).executeTakeFirst();
const logicalRelationshipCount = await trx.selectFrom("catalogLogicalRelationships")
.select(sql<number>`count(*)::int`.as("count"))
.where("databaseId", "in", selectedDatabaseIds).executeTakeFirst();
await trx.deleteFrom("catalogTables")
@@ -1292,7 +1324,8 @@ export class KyselyCatalogRepository implements CatalogRepository {
return {
tables: Number(tableCount?.count ?? 0),
columns: Number(columnCount?.count ?? 0),
relationships: Number(relationshipCount?.count ?? 0),
relationships: Number(physicalRelationshipCount?.count ?? 0)
+ Number(logicalRelationshipCount?.count ?? 0),
};
});
}
@@ -1326,13 +1359,34 @@ export class KyselyCatalogRepository implements CatalogRepository {
return { tables: 0, columns: Number(count?.count ?? 0), relationships: 0 };
}
const count = await trx.selectFrom("catalogRelationships")
const physicalCount = await trx.selectFrom("catalogRelationships")
.select(sql<number>`count(*)::int`.as("count"))
.where("databaseId", "=", databaseId)
.where((eb) => eb.or([
eb("sourceTableId", "in", selectedTableIds),
eb("targetTableId", "in", selectedTableIds),
])).executeTakeFirst();
const selectedColumnIds = (await trx.selectFrom("catalogColumns")
.select("id")
.where("tableId", "in", selectedTableIds)
.execute()).map((column) => column.id);
const logicalCount = selectedColumnIds.length === 0
? undefined
: await trx.selectFrom("catalogLogicalRelationships")
.select(sql<number>`count(*)::int`.as("count"))
.where("databaseId", "=", databaseId)
.where((eb) => eb.or([
eb("sourceColumnId", "in", selectedColumnIds),
eb("targetColumnId", "in", selectedColumnIds),
])).executeTakeFirst();
if (selectedColumnIds.length > 0) {
await trx.deleteFrom("catalogLogicalRelationships")
.where("databaseId", "=", databaseId)
.where((eb) => eb.or([
eb("sourceColumnId", "in", selectedColumnIds),
eb("targetColumnId", "in", selectedColumnIds),
])).execute();
}
await trx.deleteFrom("catalogRelationships")
.where("databaseId", "=", databaseId)
.where((eb) => eb.or([
@@ -1343,7 +1397,11 @@ export class KyselyCatalogRepository implements CatalogRepository {
schemaSyncedVersion: null,
schemaSyncedAt: null,
}).where("id", "=", databaseId).execute();
return { tables: 0, columns: 0, relationships: Number(count?.count ?? 0) };
return {
tables: 0,
columns: 0,
relationships: Number(physicalCount?.count ?? 0) + Number(logicalCount?.count ?? 0),
};
});
}
@@ -116,6 +116,15 @@ export class SensitivityAnalysisRunner {
);
ensureActive(signal);
},
async (message) => {
ensureActive(signal);
await this.repository.appendSensitivityAnalysisEvent(
started.id,
"info",
message,
);
ensureActive(signal);
},
);
ensureActive(signal);
const suggestedSensitive = suggestions.filter(
@@ -12,7 +12,7 @@ import type {
} from "./types.js";
export type { SensitivityAnalysisScope } from "./types.js";
export const SENSITIVITY_POLICY_VERSION = "sensitivity-v2";
export const SENSITIVITY_POLICY_VERSION = "sensitivity-v4";
interface SelectedColumn {
table: CatalogTable;
@@ -116,6 +116,7 @@ export class SensitivityAnalysisService {
signal: AbortSignal,
onPrepared?: (total: number) => void | Promise<void>,
onProgress?: (processed: number, suggestions: readonly SensitivityReviewItem[]) => void | Promise<void>,
onActivity?: (message: string) => void | Promise<void>,
): Promise<readonly SensitivityReviewItem[]> {
const configuredNerBudget = this.options.nerBudgetMs ?? 10_000;
const nerBudget: SensitivityNerBudget = {
@@ -144,7 +145,12 @@ export class SensitivityAnalysisService {
columns: items.map(({ column }) => column),
};
});
const assessments = await this.classifier.assess(tableTargets, signal, nerBudget);
const assessments = await this.classifier.assess(
tableTargets,
signal,
nerBudget,
onActivity,
);
ensureActive(signal);
const assessmentById = new Map(assessments.map((assessment) => [
assessment.columnId,
+52 -15
View File
@@ -128,6 +128,25 @@ function metadataEvidence(column: CatalogColumn): SensitivityEvidence | undefine
return ruleId ? { kind: "metadata", ruleId } : undefined;
}
function nonSensitiveStructuralEvidence(column: CatalogColumn): SensitivityEvidence | undefined {
if (column.dataType.trim().toLowerCase() !== "bigint") return undefined;
if (column.isPrimaryKey || column.primaryKeyPosition !== null) {
return {
kind: "type",
ruleId: "type.bigint_primary_key_non_informative",
label: "non-informative bigint primary key",
};
}
if (normalizedName(column.name) === "pk") {
return {
kind: "metadata",
ruleId: "metadata.bigint_pk_identifier_non_informative",
label: "non-informative conventional bigint primary-key identifier",
};
}
return undefined;
}
function sensitiveNameRule(value: string): string | undefined {
const name = normalizedName(value);
if (DIRECT_IDENTIFIER_NAMES.has(name)) {
@@ -299,6 +318,7 @@ function contentEvidence(value: string): SensitivityEvidence | undefined {
interface ColumnState {
column: CatalogColumn;
evidence: SensitivityEvidence[];
nonSensitiveEvidence?: SensitivityEvidence;
observedValues: number;
nerCandidates: string[];
coverage: "metadata" | "complete" | "sampled" | "no_values";
@@ -323,14 +343,16 @@ export class SensitivityClassifier {
targets: readonly SensitivityTableTarget[],
signal: AbortSignal,
sharedNerBudget?: SensitivityNerBudget,
onActivity?: (message: string) => void | Promise<void>,
): Promise<readonly SensitivityColumnAssessment[]> {
const now = this.options.now ?? Date.now;
const maxNerValuesPerColumn = boundedCount(this.options.maxNerValuesPerColumn, 8, 8);
const states = new Map<string, ColumnState>();
for (const target of targets) {
for (const column of target.columns) {
const metadataMatch = metadataEvidence(column);
const binary = UNSUPPORTED_BINARY_TYPE.test(column.dataType);
const nonSensitiveEvidence = nonSensitiveStructuralEvidence(column);
const metadataMatch = nonSensitiveEvidence ? undefined : metadataEvidence(column);
const binary = !nonSensitiveEvidence && UNSUPPORTED_BINARY_TYPE.test(column.dataType);
states.set(column.id, {
column,
evidence: metadataMatch
@@ -338,9 +360,10 @@ export class SensitivityClassifier {
: binary
? [{ kind: "type", ruleId: "type.binary_uninspectable" }]
: [],
...(nonSensitiveEvidence ? { nonSensitiveEvidence } : {}),
observedValues: 0,
nerCandidates: [],
coverage: metadataMatch || binary ? "metadata" : "no_values",
coverage: nonSensitiveEvidence || metadataMatch || binary ? "metadata" : "no_values",
sampledTarget: 0,
});
}
@@ -349,6 +372,12 @@ export class SensitivityClassifier {
const completeTables = new Set<string>();
for (const [phaseIndex, phase] of SENSITIVITY_SAMPLE_PHASES.entries()) {
for (let offset = 0; offset < targets.length; offset += MAX_CONCURRENT_TABLE_SCANS) {
signal.throwIfAborted();
const batchNumber = Math.floor(offset / MAX_CONCURRENT_TABLE_SCANS) + 1;
const batchCount = Math.ceil(targets.length / MAX_CONCURRENT_TABLE_SCANS);
await onActivity?.(
`Scanning source data: pass ${phaseIndex + 1} of ${SENSITIVITY_SAMPLE_PHASES.length}, table batch ${batchNumber} of ${batchCount}.`,
);
signal.throwIfAborted();
const peerController = new AbortController();
const scanSignal = AbortSignal.any([signal, peerController.signal]);
@@ -357,7 +386,7 @@ export class SensitivityClassifier {
if (completeTables.has(target.table.id)) return;
const columns = target.columns.filter((column) => {
const state = states.get(column.id)!;
return state.evidence.length === 0
return state.evidence.length === 0 && !state.nonSensitiveEvidence
&& (!phase.deepTextOnly || DEEP_TEXT_TYPE.test(column.dataType));
});
if (columns.length === 0) return;
@@ -411,14 +440,14 @@ export class SensitivityClassifier {
&& nerBudget.remainingMs > 0) {
const maxCandidates = boundedCount(this.options.maxNerCandidatesPerTable, 2, 1_024);
const threshold = this.options.nerConfidenceThreshold ?? 0.8;
for (const target of targets) {
for (const [targetIndex, target] of targets.entries()) {
signal.throwIfAborted();
if (nerBudget.remainingMs <= 0) break;
const candidates: LocalNerCandidate[] = [];
candidateSelection: for (let valueIndex = 0; valueIndex < maxNerValuesPerColumn; valueIndex += 1) {
for (const column of target.columns) {
const state = states.get(column.id)!;
if (state.evidence.length > 0) continue;
if (state.evidence.length > 0 || state.nonSensitiveEvidence) continue;
const text = state.nerCandidates[valueIndex];
if (text === undefined) continue;
candidates.push({ columnId: column.id, text });
@@ -426,6 +455,10 @@ export class SensitivityClassifier {
}
}
if (candidates.length === 0) continue;
await onActivity?.(
`Running local entity detection: table ${targetIndex + 1} of ${targets.length}.`,
);
signal.throwIfAborted();
const startedAt = now();
const deadline = startedAt + nerBudget.remainingMs;
try {
@@ -464,17 +497,21 @@ export class SensitivityClassifier {
return targets.flatMap((target) => target.columns.map((column) => {
const state = states.get(column.id)!;
const sensitive = state.evidence.length > 0;
const coverage = state.observedValues === 0 && !sensitive ? "no_values" : state.coverage;
const coverage = state.nonSensitiveEvidence
? "metadata"
: state.observedValues === 0 && !sensitive ? "no_values" : state.coverage;
const coverageEvidence: SensitivityEvidence[] = sensitive
? state.evidence
: [{
kind: "coverage",
ruleId: coverage === "complete"
? "coverage.complete"
: coverage === "no_values"
? "coverage.no_values"
: `coverage.sampled_${state.sampledTarget}`,
}];
: state.nonSensitiveEvidence
? [state.nonSensitiveEvidence]
: [{
kind: "coverage",
ruleId: coverage === "complete"
? "coverage.complete"
: coverage === "no_values"
? "coverage.no_values"
: `coverage.sampled_${state.sampledTarget}`,
}];
return {
columnId: column.id,
assessment: sensitive ? "sensitive" : "non_sensitive",
+4 -9
View File
@@ -39,7 +39,10 @@ function safeFailure(error: unknown): { code: string; message: string } {
};
}
if (error instanceof CatalogConnectorError) {
return { code: "schema_introspection_failed", message: "The database schema could not be read safely." };
return {
code: "schema_introspection_failed",
message: "The database schema could not be read. Check the connection and credentials, then try again.",
};
}
return { code: "schema_sync_failed", message: "Schema synchronization failed." };
}
@@ -64,7 +67,6 @@ export class CatalogSyncWorker {
}
async start(database: WorkspaceDatabase, scope: CatalogSyncScope, tableIds: readonly string[]): Promise<CatalogSyncRun> {
this.assertReady(database);
const uniqueTableIds = [...new Set(tableIds)];
if (scope === "columns") {
const tables = await Promise.all(uniqueTableIds.map((tableId) => this.repository.getTable(database.id, tableId)));
@@ -177,7 +179,6 @@ export class CatalogSyncWorker {
if (!database || database.version !== claimed.requestedDatabaseVersion) {
throw new CatalogConflictError("Database binding changed before synchronization started");
}
this.assertReady(database);
const progress: CatalogSchemaScanProgress = async (phase, counts) => {
await this.checkCancelled(runId);
await this.repository.updateSyncRun(runId, {
@@ -277,12 +278,6 @@ export class CatalogSyncWorker {
}
}
private assertReady(database: WorkspaceDatabase): void {
if (database.connectionStatus !== "reachable" || database.testedVersion !== database.version) {
throw new CatalogConflictError("Test the current database binding before synchronizing its schema");
}
}
private assertCapability(scope: CatalogSyncScope, snapshot: ObservedSchemaSnapshot): void {
const required = scope === "all" ? ["tables", "columns", "relationships"] as const : [scope] as const;
for (const name of required) {
+3 -1
View File
@@ -102,6 +102,7 @@ export interface CatalogColumn {
description: string | null;
generatedDescription: string | null;
sensitive: boolean;
sensitivityReason: string | null;
lastSyncedDatabaseVersion: number | null;
lastSyncedAt: string | null;
version: number;
@@ -195,7 +196,7 @@ export interface CatalogLogicalRelationshipCandidate {
export type CatalogDatabaseMetadataDeleteTarget = "tables" | "relationships";
export type CatalogTableMetadataDeleteTarget = "columns" | "relationships";
export type CatalogDescriptionTarget = "tables" | "columns";
export type CatalogDescriptionTarget = "tables" | "columns" | "database_columns";
export interface CatalogMetadataDeleteCounts {
tables: number;
@@ -463,6 +464,7 @@ export interface CatalogRepository {
description: string | null,
generatedDescription: string | null,
sensitive?: boolean,
sensitivityReason?: string | null,
): Promise<CatalogColumn | undefined>;
consolidateGeneratedDescriptions(
databaseId: string,
@@ -9,10 +9,13 @@ import {
} from "../catalog/types.js";
const idSchema = z.uuid();
const consolidationSchema = z.object({
target: z.enum(["tables", "columns"]),
targetIds: z.array(idSchema).min(1).max(10_000),
}).strict();
const consolidationSchema = z.discriminatedUnion("target", [
z.object({
target: z.enum(["tables", "columns"]),
targetIds: z.array(idSchema).min(1).max(10_000),
}).strict(),
z.object({ target: z.literal("database_columns") }).strict(),
]);
function manage(request: FastifyRequest, reply: FastifyReply) {
return isPrincipalContext(requirePermission(request, reply, "database.manage"));
@@ -49,7 +52,7 @@ export function catalogDescriptionConsolidationRoutes(
try {
const databaseId = idSchema.parse((request.params as { databaseId?: unknown }).databaseId);
const input = consolidationSchema.parse(request.body);
const targetIds = [...new Set(input.targetIds)];
const targetIds = "targetIds" in input ? [...new Set(input.targetIds)] : [];
const result = await deps.operations.run(
databaseId,
async () => await deps.repository.consolidateGeneratedDescriptions(
+18 -2
View File
@@ -19,8 +19,14 @@ const metadataSchema = z.object({
description: z.string().max(20_000).nullable().optional(),
generatedDescription: z.string().max(20_000).nullable().optional(),
sensitive: z.boolean().optional(),
sensitivityReason: z.string().max(2_000).nullable().optional(),
}).strict().refine((value) => (
"description" in value || "generatedDescription" in value || "sensitive" in value
"description" in value
|| "generatedDescription" in value
|| "sensitive" in value
|| "sensitivityReason" in value
)).refine((value) => (
value.sensitivityReason == null || value.sensitive === true
));
const createRunSchema = z.object({
version: z.number().int().positive(),
@@ -56,7 +62,10 @@ function safeError(reply: FastifyReply, error: unknown) {
return reply.code(409).send({ code: "schema_sync_conflict", message: error.message });
}
if (error instanceof CatalogConnectorError) {
return reply.code(502).send({ code: "schema_introspection_failed", message: "The database schema could not be read safely." });
return reply.code(502).send({
code: "schema_introspection_failed",
message: "The database schema could not be read. Check the connection and credentials, then try again.",
});
}
if (error instanceof z.ZodError) {
return reply.code(400).send({ code: "schema_request_invalid", message: "Schema request is invalid." });
@@ -113,6 +122,12 @@ export function catalogSchemaRoutes(
if (current.version !== input.version) {
return reply.code(409).send({ code: "column_stale", message: "Column metadata changed. Reload and try again." });
}
const nextSensitive = input.sensitive ?? current.sensitive;
const nextSensitivityReason = nextSensitive
? ("sensitivityReason" in input
? normalized(input.sensitivityReason ?? null)
: current.sensitivityReason)
: null;
const updated = await deps.repository.updateColumnMetadata(
databaseId,
tableId,
@@ -123,6 +138,7 @@ export function catalogSchemaRoutes(
? normalized(input.generatedDescription ?? null)
: current.generatedDescription,
input.sensitive,
nextSensitivityReason,
);
if (!updated) return reply.code(409).send({ code: "column_stale", message: "Column metadata changed. Reload and try again." });
return updated;
+48 -5
View File
@@ -16,23 +16,33 @@ import {
type WorkspaceDescriptor,
} from "../workspaces/schema.js";
import type { RuntimeBindings } from "../workspaces/runtime-renderer.js";
import type { ConnectorDiagnostics } from "../workspaces/diagnostics.js";
import type {
ConnectorDiagnostics,
Diagnostic,
WorkspaceDiagnosticOptions,
} from "../workspaces/diagnostics.js";
import { isPrincipalContext, requirePermission } from "../auth/authorization.js";
import type { AuthDiagnoser } from "../auth/diagnostics.js";
import { decodeAuthDiagnostics, type AuthDiagnostics } from "../auth/group-catalog.js";
import type { WorkspaceDatabase } from "../catalog/types.js";
export type WorkspaceDiagnoser = (
workspace: WorkspaceDescriptor,
bindings: RuntimeBindings,
options: { writeProbe: boolean },
options: WorkspaceDiagnosticOptions,
) => Promise<ConnectorDiagnostics>;
export type WorkspaceDatabaseTester = (
workspaceId: string,
) => Promise<WorkspaceDatabase | undefined>;
interface WorkspaceRoutesDeps {
registry: WorkspaceRegistry;
config: WorkspaceRegistryConfig;
diagnose: WorkspaceDiagnoser;
authDiagnoser: AuthDiagnoser;
secretStore: WorkspaceSecretStore;
testDatabaseConnection: WorkspaceDatabaseTester;
}
const workspaceId = z.string().regex(/^[a-z][a-z0-9-]{2,62}$/);
@@ -56,6 +66,20 @@ const SAFE_MESSAGES = {
semantic_index_incompatible: "Semantic index is incompatible with this workspace.",
} as const;
const catalogConnectionUnavailable = (): Diagnostic => ({
level: "error",
code: "connector_unavailable",
field: "dwh",
message: "The configured database could not be reached or authenticated.",
});
const catalogConnectionMissing = (): Diagnostic => ({
level: "error",
code: "binding_missing",
field: "dwh",
message: "Configure this workspace in Database Management before testing connections.",
});
function authenticationReport(value: unknown): AuthDiagnostics {
const report = decodeAuthDiagnostics(value);
if (!report) throw new Error("invalid authentication diagnostic report");
@@ -238,14 +262,33 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
deps.secretStore,
);
try {
const [workspaceDiagnostics, inspectedAuthentication] = await Promise.all([
deps.diagnose(operational, lease.bindings, { writeProbe: false }),
const [workspaceDiagnostics, testedDatabase, inspectedAuthentication] = await Promise.all([
deps.diagnose(operational, lease.bindings, {
writeProbe: false,
skipDwh: true,
}),
deps.testDatabaseConnection(id),
deps.authDiagnoser.inspect({ live: true }),
]);
const authentication = authenticationReport(inspectedAuthentication);
const catalogConnectionReady = testedDatabase?.connectionStatus === "reachable";
const catalogConnectionDiagnostic = !testedDatabase
? catalogConnectionMissing()
: catalogConnectionReady
? undefined
: catalogConnectionUnavailable();
const diagnostics = catalogConnectionDiagnostic
? [
...workspaceDiagnostics.diagnostics.filter(({ code }) => code !== "binding_ok"),
catalogConnectionDiagnostic,
]
: workspaceDiagnostics.diagnostics;
return {
...workspaceDiagnostics,
activatable: workspaceDiagnostics.activatable && authentication.ready,
activatable: workspaceDiagnostics.activatable
&& catalogConnectionReady
&& authentication.ready,
diagnostics,
authentication,
};
} finally {
+68 -51
View File
@@ -130,6 +130,11 @@ export interface DiagnosticAdapters {
probeEmbedding(request: EmbeddingDiagnosticRequest): Promise<EmbeddingDiagnosticResult>;
}
export interface WorkspaceDiagnosticOptions {
writeProbe: boolean;
skipDwh?: boolean;
}
export const DEFAULT_WORKSPACE_DIAGNOSTIC_TIMEOUT_MS = 5_000;
async function secretPresent(file: string): Promise<boolean> {
@@ -421,6 +426,7 @@ async function diagnoseValidatedWorkspace(
adapters: DiagnosticAdapters,
timeoutMs: number,
semanticRuntime: SemanticRuntimeConfig,
skipDwh: boolean,
): Promise<ConnectorDiagnostics> {
const evidenceField = descriptor.evidence?.source.type === "http"
? "evidence.source.authentication"
@@ -432,75 +438,79 @@ async function diagnoseValidatedWorkspace(
variable,
}));
const diagnostics = [
...[...bindings.dwh.missing].sort().map((field) => diagnosticError("binding_missing", field)),
...(skipDwh
? []
: [...bindings.dwh.missing].sort().map((field) => diagnosticError("binding_missing", field))),
...evidenceDiagnostics,
];
if (diagnostics.length > 0) return { activatable: false, diagnostics };
const dwhTimeout = boundedTimeout(descriptor.dwh.timeout_ms, timeoutMs);
let activatable = true;
const dwhValues = bindings.dwh.values;
const dwhField = (suffix: string) => bindingName(descriptor, suffix);
const dwhResource = { database: descriptor.dwh.database, schema: descriptor.dwh.schema };
let dwhRequest: ConnectorDiagnosticRequest | undefined;
if (bindings.dwh.transport === "rest_api") {
const diagnostic = descriptor.diagnostics?.dwh_rest;
const baseUrl = dwhValues[dwhField("BASE_URL")];
if (diagnostic && baseUrl) {
const credentialFile = diagnostic.auth === "none" ? undefined : dwhValues[dwhField("API_KEY_FILE")];
if (diagnostic.auth === "none" || credentialFile !== undefined) {
if (!skipDwh) {
const dwhTimeout = boundedTimeout(descriptor.dwh.timeout_ms, timeoutMs);
const dwhValues = bindings.dwh.values;
const dwhField = (suffix: string) => bindingName(descriptor, suffix);
const dwhResource = { database: descriptor.dwh.database, schema: descriptor.dwh.schema };
let dwhRequest: ConnectorDiagnosticRequest | undefined;
if (bindings.dwh.transport === "rest_api") {
const diagnostic = descriptor.diagnostics?.dwh_rest;
const baseUrl = dwhValues[dwhField("BASE_URL")];
if (diagnostic && baseUrl) {
const credentialFile = diagnostic.auth === "none" ? undefined : dwhValues[dwhField("API_KEY_FILE")];
if (diagnostic.auth === "none" || credentialFile !== undefined) {
dwhRequest = {
role: "dwh",
transport: "rest_api",
baseUrl,
credentialFile,
tlsCaFile: dwhValues[dwhField("TLS_CA_FILE")],
resource: dwhResource,
timeoutMs: dwhTimeout,
signal: new AbortController().signal,
diagnostic,
};
}
}
} else if (bindings.dwh.transport === "postgres_direct") {
const host = dwhValues[dwhField("HOST")];
const port = numericBinding(dwhValues, dwhField("PORT"));
const user = dwhValues[dwhField("USER")];
const credentialFile = dwhValues[dwhField("PASSWORD_FILE")];
if (host && port && user && credentialFile) {
dwhRequest = {
role: "dwh",
transport: "rest_api",
baseUrl,
transport: "postgres_direct",
host,
port,
user,
credentialFile,
tlsCaFile: dwhValues[dwhField("TLS_CA_FILE")],
resource: dwhResource,
timeoutMs: dwhTimeout,
signal: new AbortController().signal,
diagnostic,
};
}
}
} else if (bindings.dwh.transport === "postgres_direct") {
const host = dwhValues[dwhField("HOST")];
const port = numericBinding(dwhValues, dwhField("PORT"));
const user = dwhValues[dwhField("USER")];
const credentialFile = dwhValues[dwhField("PASSWORD_FILE")];
if (host && port && user && credentialFile) {
dwhRequest = {
role: "dwh",
transport: "postgres_direct",
host,
port,
user,
credentialFile,
tlsCaFile: dwhValues[dwhField("TLS_CA_FILE")],
resource: dwhResource,
timeoutMs: dwhTimeout,
signal: new AbortController().signal,
};
if (!dwhRequest) {
diagnostics.push(diagnosticError("workspace_not_activatable"));
return { activatable: false, diagnostics };
}
}
if (!dwhRequest) {
diagnostics.push(diagnosticError("workspace_not_activatable"));
return { activatable: false, diagnostics };
}
try {
const dwhResult = await withTimeout(dwhTimeout, (signal) => adapters.probeConnector({
...dwhRequest,
signal,
timeoutMs: dwhTimeout,
}));
if (!hasRequiredConnectorChecks(dwhResult, dwhRequest.resource)) {
try {
const dwhResult = await withTimeout(dwhTimeout, (signal) => adapters.probeConnector({
...dwhRequest,
signal,
timeoutMs: dwhTimeout,
}));
if (!hasRequiredConnectorChecks(dwhResult, dwhRequest.resource)) {
diagnostics.push(diagnosticError("connector_unavailable"));
activatable = false;
}
} catch {
diagnostics.push(diagnosticError("connector_unavailable"));
activatable = false;
}
} catch {
diagnostics.push(diagnosticError("connector_unavailable"));
activatable = false;
}
try {
@@ -571,11 +581,18 @@ export function createWorkspaceDiagnoser(
return async function diagnose(
workspace: WorkspaceDescriptor,
bindings: RuntimeBindings,
_options: { writeProbe: boolean },
diagnosticOptions: WorkspaceDiagnosticOptions,
): Promise<ConnectorDiagnostics> {
requireSupportedDescriptor(workspace);
const descriptor = validateWorkspaceDescriptor(workspace);
return await diagnoseValidatedWorkspace(descriptor, bindings, adapters, timeoutMs, semanticRuntime);
return await diagnoseValidatedWorkspace(
descriptor,
bindings,
adapters,
timeoutMs,
semanticRuntime,
diagnosticOptions.skipDwh ?? false,
);
};
}