feat: complete catalog fleet management workflow

This commit is contained in:
Codex
2026-08-31 15:58:43 +02:00
parent 866aee4249
commit 9c697dc062
56 changed files with 6692 additions and 611 deletions
+255
View File
@@ -11,6 +11,7 @@ import {
type Transaction,
} from "kysely";
import { Pool } from "pg";
import { createCatalogMetrics } from "./metrics.js";
import {
CatalogConflictError,
CatalogConnectorError,
@@ -21,6 +22,7 @@ import {
type CatalogDescriptionTarget,
type CatalogDatabaseMetadataDeleteTarget,
type CatalogMetadataDeleteCounts,
type CatalogMetrics,
type CatalogRelationship,
type CatalogSchemaDiff,
type CatalogSyncCounts,
@@ -40,6 +42,10 @@ import {
type DescriptionGenerationScope,
type ObservedCatalogTable,
type ObservedSchemaSnapshot,
type SensitiveDataSuggestionEvent,
type SensitiveDataSuggestionRun,
type SensitiveDataSuggestionRunUpdate,
type SensitiveDataSuggestionScope,
type TableSyncRepositoryResult,
type WorkspaceDatabase,
} from "./types.js";
@@ -166,6 +172,30 @@ interface DescriptionGenerationEventTable {
createdAt: Timestamp;
}
interface SensitiveDataSuggestionRunTable {
id: string;
databaseId: string;
scope: SensitiveDataSuggestionScope;
modelId: string;
status: SensitiveDataSuggestionRun["status"];
total: number;
suggestedSensitive: number;
suggestedNonSensitive: number;
createdAt: Timestamp;
startedAt: Timestamp;
updatedAt: Timestamp;
finishedAt: Timestamp | null;
errorSummary: string | null;
}
interface SensitiveDataSuggestionEventTable {
runId: string;
sequence: number;
level: SensitiveDataSuggestionEvent["level"];
message: string;
createdAt: Timestamp;
}
interface CatalogSyncRunTable {
id: string;
databaseId: string;
@@ -213,6 +243,8 @@ export interface CatalogDatabase {
catalogRelationshipColumns: CatalogRelationshipColumnTable;
descriptionGenerationRuns: DescriptionGenerationRunTable;
descriptionGenerationEvents: DescriptionGenerationEventTable;
sensitiveDataSuggestionRuns: SensitiveDataSuggestionRunTable;
sensitiveDataSuggestionEvents: SensitiveDataSuggestionEventTable;
catalogSyncRuns: CatalogSyncRunTable;
catalogSyncEvents: CatalogSyncEventTable;
}
@@ -220,6 +252,17 @@ export interface CatalogDatabase {
type DbOrTransaction = Kysely<CatalogDatabase> | Transaction<CatalogDatabase>;
type JoinedRow = Selectable<WorkspaceDatabaseTable> & Selectable<DatabaseBindingTable>;
interface CatalogMetricsRow {
databaseCount: number;
tables: number;
columns: number;
sensitiveColumns: number;
relationships: number;
describedTables: number;
describedColumns: number;
updatedAt: Date | string | null;
}
function present<T>(value: T | null): T | undefined {
return value === null ? undefined : value;
}
@@ -338,6 +381,25 @@ function serializeDescriptionGenerationEvent(
return { ...row, createdAt: new Date(row.createdAt).toISOString() };
}
function serializeSensitiveDataSuggestionRun(
row: Selectable<SensitiveDataSuggestionRunTable>,
): SensitiveDataSuggestionRun {
const stamp = (value: Date | string | null) => value === null ? null : new Date(value).toISOString();
return {
...row,
createdAt: new Date(row.createdAt).toISOString(),
startedAt: new Date(row.startedAt).toISOString(),
updatedAt: new Date(row.updatedAt).toISOString(),
finishedAt: stamp(row.finishedAt),
};
}
function serializeSensitiveDataSuggestionEvent(
row: Selectable<SensitiveDataSuggestionEventTable>,
): SensitiveDataSuggestionEvent {
return { ...row, createdAt: new Date(row.createdAt).toISOString() };
}
function bindingValues(databaseId: string, binding: DatabaseBinding) {
return {
databaseId,
@@ -388,6 +450,83 @@ export class KyselyCatalogRepository implements CatalogRepository {
return id ? await this.get(id.id) : undefined;
}
async getCatalogMetrics(databaseId?: string): Promise<CatalogMetrics | undefined> {
const result = await sql<CatalogMetricsRow>`
WITH requested_database AS (
SELECT ${databaseId ?? null}::uuid AS id
),
selected_databases AS (
SELECT workspace_databases.id, workspace_databases.schema_synced_at
FROM workspace_databases
CROSS JOIN requested_database
WHERE requested_database.id IS NULL
OR workspace_databases.id = requested_database.id
),
table_metrics AS (
SELECT
count(*)::int AS tables,
count(*) FILTER (
WHERE nullif(btrim(catalog_tables.description), '') IS NOT NULL
OR nullif(btrim(catalog_tables.generated_description), '') IS NOT NULL
)::int AS "describedTables",
max(catalog_tables.updated_at) AS updated_at
FROM catalog_tables
INNER JOIN selected_databases
ON selected_databases.id = catalog_tables.database_id
),
column_metrics AS (
SELECT
count(*)::int AS columns,
count(*) FILTER (WHERE catalog_columns.sensitive)::int AS "sensitiveColumns",
count(*) FILTER (
WHERE nullif(btrim(catalog_columns.description), '') IS NOT NULL
OR nullif(btrim(catalog_columns.generated_description), '') IS NOT NULL
)::int AS "describedColumns",
max(catalog_columns.updated_at) AS updated_at
FROM catalog_columns
INNER JOIN catalog_tables ON catalog_tables.id = catalog_columns.table_id
INNER JOIN selected_databases
ON selected_databases.id = catalog_tables.database_id
),
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
)
SELECT
(SELECT count(*)::int FROM selected_databases) AS "databaseCount",
table_metrics.tables,
column_metrics.columns,
column_metrics."sensitiveColumns",
relationship_metrics.relationships,
table_metrics."describedTables",
column_metrics."describedColumns",
greatest(
(SELECT max(schema_synced_at) FROM selected_databases),
table_metrics.updated_at,
column_metrics.updated_at,
relationship_metrics.updated_at
) AS "updatedAt"
FROM table_metrics
CROSS JOIN column_metrics
CROSS JOIN relationship_metrics
`.execute(this.db);
const row = result.rows[0];
if (!row || (databaseId !== undefined && Number(row.databaseCount) === 0)) return undefined;
return createCatalogMetrics(databaseId, {
tables: Number(row.tables),
columns: Number(row.columns),
sensitiveColumns: Number(row.sensitiveColumns),
relationships: Number(row.relationships),
describedTables: Number(row.describedTables),
describedColumns: Number(row.describedColumns),
}, row.updatedAt === null ? null : new Date(row.updatedAt).toISOString());
}
async create(input: DatabaseConfigurationInput): Promise<WorkspaceDatabase> {
try {
return await this.db.transaction().execute(async (trx) => {
@@ -775,6 +914,114 @@ export class KyselyCatalogRepository implements CatalogRepository {
return rows.map(serializeDescriptionGenerationEvent);
}
async createSensitiveDataSuggestionRun(
databaseId: string,
scope: SensitiveDataSuggestionScope,
modelId: string,
): Promise<SensitiveDataSuggestionRun> {
const row = await this.db.insertInto("sensitiveDataSuggestionRuns").values({
id: randomUUID(),
databaseId,
scope,
modelId,
status: "running",
total: 0,
suggestedSensitive: 0,
suggestedNonSensitive: 0,
finishedAt: null,
errorSummary: null,
}).returningAll().executeTakeFirstOrThrow();
return serializeSensitiveDataSuggestionRun(row);
}
async getSensitiveDataSuggestionRun(
runId: string,
): Promise<SensitiveDataSuggestionRun | undefined> {
const row = await this.db.selectFrom("sensitiveDataSuggestionRuns")
.selectAll()
.where("id", "=", runId)
.executeTakeFirst();
return row ? serializeSensitiveDataSuggestionRun(row) : undefined;
}
async listSensitiveDataSuggestionRuns(limit = 50): Promise<SensitiveDataSuggestionRun[]> {
const rows = await this.db.selectFrom("sensitiveDataSuggestionRuns")
.selectAll()
.orderBy("createdAt", "desc")
.orderBy("id", "desc")
.limit(limit)
.execute();
return rows.map(serializeSensitiveDataSuggestionRun);
}
async interruptActiveSensitiveDataSuggestionRuns(
errorSummary: string,
): Promise<SensitiveDataSuggestionRun[]> {
const rows = await this.db.updateTable("sensitiveDataSuggestionRuns")
.set({
status: "interrupted",
finishedAt: sql`now()`,
updatedAt: sql`now()`,
errorSummary,
})
.where("status", "=", "running")
.returningAll()
.execute();
return rows.map(serializeSensitiveDataSuggestionRun);
}
async updateSensitiveDataSuggestionRun(
runId: string,
update: SensitiveDataSuggestionRunUpdate,
): Promise<SensitiveDataSuggestionRun | undefined> {
const values: any = { ...update, updatedAt: sql`now()` };
const row = await this.db.updateTable("sensitiveDataSuggestionRuns")
.set(values)
.where("id", "=", runId)
.returningAll()
.executeTakeFirst();
return row ? serializeSensitiveDataSuggestionRun(row) : undefined;
}
async appendSensitiveDataSuggestionEvent(
runId: string,
level: SensitiveDataSuggestionEvent["level"],
message: string,
): Promise<SensitiveDataSuggestionEvent> {
return await this.db.transaction().execute(async (trx) => {
const run = await trx.selectFrom("sensitiveDataSuggestionRuns")
.select("id")
.where("id", "=", runId)
.forUpdate()
.executeTakeFirst();
if (!run) throw new CatalogConflictError("Sensitive Data Suggestion Run does not exist");
const current = await trx.selectFrom("sensitiveDataSuggestionEvents")
.select(sql<number>`coalesce(max(sequence), 0)::int`.as("sequence"))
.where("runId", "=", runId)
.executeTakeFirst();
const row = await trx.insertInto("sensitiveDataSuggestionEvents").values({
runId,
sequence: Number(current?.sequence ?? 0) + 1,
level,
message,
}).returningAll().executeTakeFirstOrThrow();
return serializeSensitiveDataSuggestionEvent(row);
});
}
async listSensitiveDataSuggestionEvents(
runId: string,
afterSequence = 0,
): Promise<SensitiveDataSuggestionEvent[]> {
const rows = await this.db.selectFrom("sensitiveDataSuggestionEvents")
.selectAll()
.where("runId", "=", runId)
.where("sequence", ">", afterSequence)
.orderBy("sequence")
.execute();
return rows.map(serializeSensitiveDataSuggestionEvent);
}
async listRelationships(databaseId: string): Promise<CatalogRelationship[]> {
const rows = await this.db.selectFrom("catalogRelationships as relationship")
.innerJoin("catalogTables as sourceTable", "sourceTable.id", "relationship.sourceTableId")
@@ -1345,6 +1592,7 @@ export class UnavailableCatalogRepository implements CatalogRepository {
async list(): Promise<WorkspaceDatabase[]> { return this.fail(); }
async get(): Promise<WorkspaceDatabase | undefined> { return this.fail(); }
async getByWorkspace(): Promise<WorkspaceDatabase | undefined> { return this.fail(); }
async getCatalogMetrics(): Promise<CatalogMetrics | undefined> { return this.fail(); }
async create(): Promise<WorkspaceDatabase> { return this.fail(); }
async update(): Promise<WorkspaceDatabase | undefined> { return this.fail(); }
async recordTest(): Promise<WorkspaceDatabase | undefined> { return this.fail(); }
@@ -1366,6 +1614,13 @@ export class UnavailableCatalogRepository implements CatalogRepository {
async updateDescriptionGenerationRun(): Promise<DescriptionGenerationRun | undefined> { return this.fail(); }
async appendDescriptionGenerationEvent(): Promise<DescriptionGenerationEvent> { return this.fail(); }
async listDescriptionGenerationEvents(): Promise<DescriptionGenerationEvent[]> { return this.fail(); }
async createSensitiveDataSuggestionRun(): Promise<SensitiveDataSuggestionRun> { return this.fail(); }
async getSensitiveDataSuggestionRun(): Promise<SensitiveDataSuggestionRun | undefined> { return this.fail(); }
async listSensitiveDataSuggestionRuns(): Promise<SensitiveDataSuggestionRun[]> { return this.fail(); }
async interruptActiveSensitiveDataSuggestionRuns(): Promise<SensitiveDataSuggestionRun[]> { return this.fail(); }
async updateSensitiveDataSuggestionRun(): Promise<SensitiveDataSuggestionRun | undefined> { return this.fail(); }
async appendSensitiveDataSuggestionEvent(): Promise<SensitiveDataSuggestionEvent> { return this.fail(); }
async listSensitiveDataSuggestionEvents(): Promise<SensitiveDataSuggestionEvent[]> { return this.fail(); }
async listRelationships(): Promise<CatalogRelationship[]> { return this.fail(); }
async deleteDatabaseMetadata(): Promise<CatalogMetadataDeleteCounts | undefined> { return this.fail(); }
async deleteTableMetadata(): Promise<CatalogMetadataDeleteCounts | undefined> { return this.fail(); }