import { randomUUID } from "node:crypto"; import { readFile } from "node:fs/promises"; import { CamelCasePlugin, Kysely, PostgresDialect, sql, type ColumnType, type Generated, type Selectable, type Transaction, } from "kysely"; import { Pool } from "pg"; import { createCatalogMetrics } from "./metrics.js"; import { CatalogConflictError, CatalogConnectorError, CatalogUnavailableError, DescriptionGenerationRunActiveError, type CatalogColumn, type CatalogDescriptionConsolidationCounts, type CatalogDescriptionTarget, type CatalogDatabaseMetadataDeleteTarget, type CatalogMetadataDeleteCounts, type CatalogMetrics, type CatalogLogicalRelationship, type CatalogLogicalRelationshipCandidate, type CatalogLogicalRelationshipContext, type CatalogPhysicalRelationship, type CatalogPreprocessingStartResult, type CatalogPreprocessingClearResult, type CatalogSchemaDiff, type CatalogSyncCounts, type CatalogSyncEvent, type CatalogSyncRun, type CatalogSyncRunUpdate, type CatalogSyncScope, type CatalogTable, type CatalogTableMetadataDeleteTarget, type CatalogRepository, type DatabaseBinding, type DatabaseConfigurationInput, type DatabaseTestResult, type DescriptionGenerationEvent, type DescriptionGenerationRun, type DescriptionGenerationRunUpdate, type DescriptionGenerationScope, type ObservedCatalogTable, type ObservedSchemaSnapshot, type SensitivityAnalysisEvent, type SensitivityAnalysisRun, type SensitivityAnalysisRunUpdate, type SensitivityAnalysisScope, type TableSyncRepositoryResult, type WorkspaceDatabase, } from "./types.js"; type Timestamp = ColumnType; type NullableTimestamp = ColumnType< Date | null, Date | string | null | undefined, Date | string | null >; interface WorkspaceDatabaseTable { id: string; workspaceId: string; engine: "postgres"; databaseName: string; schemaName: string; version: Generated; createdAt: Timestamp; updatedAt: Timestamp; schemaSyncedVersion: number | null; schemaSyncedAt: Timestamp | null; metadataContentRevision: Generated; preprocessingStatus: Generated; preprocessingInputFingerprint: Generated; preprocessedMetadataRevision: Generated; preprocessingStartedAt: NullableTimestamp; preprocessingFinishedAt: NullableTimestamp; preprocessingErrorCode: Generated; } interface DatabaseBindingTable { databaseId: string; transport: DatabaseBinding["transport"]; host: string | null; port: number | null; username: string | null; baseUrl: string | null; restPath: string | null; restAuth: DatabaseBinding["restAuth"] | null; tlsServername: string | null; sshHost: string | null; sshPort: number | null; sshUsername: string | null; sshTargetHost: string | null; sshTargetPort: number | null; connectionStatus: WorkspaceDatabase["connectionStatus"]; testedVersion: number | null; lastTestedAt: Timestamp | null; lastErrorCode: string | null; lastErrorMessage: string | null; updatedAt: Timestamp; } interface CatalogTableTable { id: string; databaseId: string; name: string; sourceComment: string | null; description: string | null; generatedDescription: string | null; lastSyncedDatabaseVersion: number | null; lastSyncedAt: Timestamp | null; version: Generated; createdAt: Timestamp; updatedAt: Timestamp; } interface CatalogColumnTable { id: string; tableId: string; name: string; ordinalPosition: number; dataType: string; isNullable: boolean; defaultExpression: string | null; primaryKeyPosition: number | null; sourceComment: string | null; description: string | null; generatedDescription: string | null; sensitive: Generated; sensitivityReason: Generated; lastSyncedDatabaseVersion: number | null; lastSyncedAt: Timestamp | null; version: Generated; createdAt: Timestamp; updatedAt: Timestamp; } interface CatalogRelationshipTable { id: string; databaseId: string; constraintName: string; sourceTableId: string; targetTableId: string; updateRule: string; deleteRule: string; deferrable: boolean; initiallyDeferred: boolean; lastSyncedDatabaseVersion: number | null; lastSyncedAt: Timestamp | null; createdAt: Timestamp; updatedAt: Timestamp; } interface CatalogRelationshipColumnTable { relationshipId: string; position: number; sourceColumnId: string; targetColumnId: string; } interface CatalogLogicalRelationshipTable { id: string; databaseId: string; sourceColumnId: string; targetColumnId: string; generated: Generated; deletedAt: Timestamp | null; createdAt: Timestamp; updatedAt: Timestamp; } interface DescriptionGenerationRunTable { id: string; databaseId: string; scope: DescriptionGenerationScope; modelId: string; language: DescriptionGenerationRun["language"]; status: DescriptionGenerationRun["status"]; total: number; processed: number; generated: number; nonGeneratable: number; failed: number; inputTokens: number; cacheReadTokens: number; outputTokens: number; createdAt: Timestamp; startedAt: Timestamp | null; updatedAt: Timestamp; finishedAt: Timestamp | null; errorSummary: string | null; } interface DescriptionGenerationEventTable { runId: string; sequence: number; level: DescriptionGenerationEvent["level"]; message: string; createdAt: Timestamp; } interface SensitivityAnalysisRunTable { id: string; databaseId: string; scope: SensitivityAnalysisScope; engine: SensitivityAnalysisRun["engine"]; modelId: string | null; policyVersion: string | null; status: SensitivityAnalysisRun["status"]; total: number; suggestedSensitive: number; suggestedNonSensitive: number; unknown: number; inputTokens: number; cacheReadTokens: number; outputTokens: number; createdAt: Timestamp; startedAt: Timestamp; updatedAt: Timestamp; finishedAt: Timestamp | null; errorSummary: string | null; } interface SensitivityAnalysisEventTable { runId: string; sequence: number; level: SensitivityAnalysisEvent["level"]; message: string; createdAt: Timestamp; } interface CatalogSyncRunTable { id: string; databaseId: string; scope: CatalogSyncScope; // node-postgres serializes JavaScript arrays as PostgreSQL arrays. The // catalog column is JSONB, so inserts must cross the driver boundary as a // JSON string while reads are decoded back to a string array. tableIds: ColumnType; state: CatalogSyncRun["state"]; phase: CatalogSyncRun["phase"]; requestedDatabaseVersion: number; observedSnapshot: ObservedSchemaSnapshot | null; plannedDiff: CatalogSchemaDiff | null; confirmationToken: string | null; counts: CatalogSyncCounts; errorCode: string | null; errorMessage: string | null; cancelRequested: boolean; createdAt: Timestamp; startedAt: Timestamp | null; updatedAt: Timestamp; finishedAt: Timestamp | null; heartbeatAt: Timestamp | null; leaseOwner: string | null; leaseExpiresAt: Timestamp | null; } interface CatalogSyncEventTable { id: Generated; runId: string; sequence: number; level: CatalogSyncEvent["level"]; eventType: string; message: string; data: Record; createdAt: Timestamp; } export interface CatalogDatabase { workspaceDatabases: WorkspaceDatabaseTable; databaseBindings: DatabaseBindingTable; catalogTables: CatalogTableTable; catalogColumns: CatalogColumnTable; catalogRelationships: CatalogRelationshipTable; catalogRelationshipColumns: CatalogRelationshipColumnTable; catalogLogicalRelationships: CatalogLogicalRelationshipTable; descriptionGenerationRuns: DescriptionGenerationRunTable; descriptionGenerationEvents: DescriptionGenerationEventTable; // Legacy physical table names retained for migration and storage compatibility. sensitiveDataSuggestionRuns: SensitivityAnalysisRunTable; sensitiveDataSuggestionEvents: SensitivityAnalysisEventTable; catalogSyncRuns: CatalogSyncRunTable; catalogSyncEvents: CatalogSyncEventTable; } type DbOrTransaction = Kysely | Transaction; type JoinedRow = Selectable & Selectable; interface CatalogMetricsRow { databaseCount: number; tables: number; columns: number; sensitiveColumns: number; relationships: number; describedTables: number; describedColumns: number; updatedAt: Date | string | null; } function present(value: T | null): T | undefined { return value === null ? undefined : value; } function serialize(row: JoinedRow): WorkspaceDatabase { return { id: row.id, workspaceId: row.workspaceId, engine: row.engine, databaseName: row.databaseName, schema: row.schemaName, version: row.version, createdAt: new Date(row.createdAt).toISOString(), updatedAt: new Date(row.updatedAt).toISOString(), binding: { transport: row.transport, host: present(row.host), port: present(row.port), username: present(row.username), baseUrl: present(row.baseUrl), restPath: present(row.restPath), restAuth: present(row.restAuth), tlsServername: present(row.tlsServername), sshHost: present(row.sshHost), sshPort: present(row.sshPort), sshUsername: present(row.sshUsername), sshTargetHost: present(row.sshTargetHost), sshTargetPort: present(row.sshTargetPort), }, connectionStatus: row.connectionStatus, testedVersion: present(row.testedVersion), lastTestedAt: row.lastTestedAt === null ? undefined : new Date(row.lastTestedAt).toISOString(), lastErrorCode: present(row.lastErrorCode), lastErrorMessage: present(row.lastErrorMessage), schemaSyncedVersion: present(row.schemaSyncedVersion), schemaSyncedAt: row.schemaSyncedAt == null ? undefined : new Date(row.schemaSyncedAt).toISOString(), metadataContentRevision: Number(row.metadataContentRevision), preprocessingStatus: row.preprocessingStatus, preprocessingInputFingerprint: present(row.preprocessingInputFingerprint), preprocessedMetadataRevision: row.preprocessedMetadataRevision == null ? undefined : Number(row.preprocessedMetadataRevision), preprocessingStartedAt: row.preprocessingStartedAt == null ? undefined : new Date(row.preprocessingStartedAt).toISOString(), preprocessingFinishedAt: row.preprocessingFinishedAt == null ? undefined : new Date(row.preprocessingFinishedAt).toISOString(), preprocessingErrorCode: present(row.preprocessingErrorCode), }; } function serializeTable(row: Selectable): CatalogTable { return { id: row.id, databaseId: row.databaseId, name: row.name, sourceComment: row.sourceComment, description: row.description, generatedDescription: row.generatedDescription, lastSyncedDatabaseVersion: row.lastSyncedDatabaseVersion, lastSyncedAt: row.lastSyncedAt === null ? null : new Date(row.lastSyncedAt).toISOString(), version: row.version, createdAt: new Date(row.createdAt).toISOString(), updatedAt: new Date(row.updatedAt).toISOString(), }; } function serializeColumn(row: Selectable, foreignKeyCount = 0): CatalogColumn { return { id: row.id, tableId: row.tableId, name: row.name, ordinalPosition: row.ordinalPosition, dataType: row.dataType, isNullable: row.isNullable, defaultExpression: row.defaultExpression, primaryKeyPosition: row.primaryKeyPosition, isPrimaryKey: row.primaryKeyPosition !== null, isForeignKey: foreignKeyCount > 0, foreignKeyCount, sourceComment: row.sourceComment, 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, createdAt: new Date(row.createdAt).toISOString(), updatedAt: new Date(row.updatedAt).toISOString(), }; } function serializeSyncRun(row: Selectable): CatalogSyncRun { const stamp = (value: Date | string | null) => value === null ? null : new Date(value).toISOString(); return { ...row, tableIds: row.tableIds ?? [], counts: row.counts ?? {}, createdAt: new Date(row.createdAt).toISOString(), startedAt: stamp(row.startedAt), updatedAt: new Date(row.updatedAt).toISOString(), finishedAt: stamp(row.finishedAt), heartbeatAt: stamp(row.heartbeatAt), leaseExpiresAt: stamp(row.leaseExpiresAt), }; } function serializeSyncEvent(row: Selectable): CatalogSyncEvent { return { ...row, id: Number(row.id), data: row.data ?? {}, createdAt: new Date(row.createdAt).toISOString() }; } function serializeDescriptionGenerationRun( row: Selectable, ): DescriptionGenerationRun { const stamp = (value: Date | string | null) => value === null ? null : new Date(value).toISOString(); return { ...row, createdAt: new Date(row.createdAt).toISOString(), startedAt: stamp(row.startedAt), updatedAt: new Date(row.updatedAt).toISOString(), finishedAt: stamp(row.finishedAt), }; } function serializeDescriptionGenerationEvent( row: Selectable, ): DescriptionGenerationEvent { return { ...row, createdAt: new Date(row.createdAt).toISOString() }; } function serializeSensitivityAnalysisRun( row: Selectable, ): SensitivityAnalysisRun { 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 serializeSensitivityAnalysisEvent( row: Selectable, ): SensitivityAnalysisEvent { return { ...row, createdAt: new Date(row.createdAt).toISOString() }; } function bindingValues(databaseId: string, binding: DatabaseBinding) { return { databaseId, transport: binding.transport, host: binding.host ?? null, port: binding.port ?? null, username: binding.username ?? null, baseUrl: binding.baseUrl ?? null, restPath: binding.restPath ?? null, restAuth: binding.restAuth ?? null, tlsServername: binding.tlsServername ?? null, sshHost: binding.sshHost ?? null, sshPort: binding.sshPort ?? null, sshUsername: binding.sshUsername ?? null, sshTargetHost: binding.sshTargetHost ?? null, sshTargetPort: binding.sshTargetPort ?? null, }; } async function selectOne(db: DbOrTransaction, id: string): Promise { const row = await db.selectFrom("workspaceDatabases") .innerJoin("databaseBindings", "databaseBindings.databaseId", "workspaceDatabases.id") .selectAll() .where("workspaceDatabases.id", "=", id) .executeTakeFirst(); return row ? serialize(row as JoinedRow) : undefined; } export class KyselyCatalogRepository implements CatalogRepository { constructor(private readonly db: Kysely) {} async list(): Promise { const rows = await this.db.selectFrom("workspaceDatabases") .innerJoin("databaseBindings", "databaseBindings.databaseId", "workspaceDatabases.id") .selectAll() .orderBy("workspaceDatabases.workspaceId") .execute(); return rows.map((row) => serialize(row as JoinedRow)); } async get(id: string): Promise { return await selectOne(this.db, id); } async getByWorkspace(workspaceId: string): Promise { const id = await this.db.selectFrom("workspaceDatabases").select("id") .where("workspaceId", "=", workspaceId).executeTakeFirst(); return id ? await this.get(id.id) : undefined; } async beginPreprocessing( workspaceId: string, inputFingerprint: string, ): Promise { return await this.db.transaction().execute(async (trx) => { const database = await trx.selectFrom("workspaceDatabases") .selectAll() .where("workspaceId", "=", workspaceId) .forUpdate() .executeTakeFirst(); if (!database) return { kind: "not_found" }; if (database.preprocessingStatus === "running") return { kind: "already_running" }; if (database.schemaSyncedVersion !== database.version) return { kind: "schema_stale" }; const activeSync = await trx.selectFrom("catalogSyncRuns") .select("id") .where("databaseId", "=", database.id) .where("state", "in", ["queued", "running", "awaiting_confirmation", "applying"]) .executeTakeFirst(); const activeDescriptions = await trx.selectFrom("descriptionGenerationRuns") .select("id") .where("databaseId", "=", database.id) .where("status", "in", ["queued", "running"]) .executeTakeFirst(); const activeSensitivity = await trx.selectFrom("sensitiveDataSuggestionRuns") .select("id") .where("databaseId", "=", database.id) .where("status", "=", "running") .executeTakeFirst(); if (activeSync || activeDescriptions || activeSensitivity) return { kind: "catalog_busy" }; await trx.updateTable("workspaceDatabases") .set({ preprocessingStatus: "running", preprocessingInputFingerprint: inputFingerprint, preprocessedMetadataRevision: null, preprocessingStartedAt: sql`now()`, preprocessingFinishedAt: null, preprocessingErrorCode: null, updatedAt: sql`now()`, }) .where("id", "=", database.id) .execute(); return { kind: "started", database: (await selectOne(trx, database.id))! }; }); } async finishPreprocessing( workspaceId: string, metadataContentRevision: number, inputFingerprint: string, outcome: { status: "succeeded" } | { status: "failed"; errorCode: string }, ): Promise { const result = await this.db.updateTable("workspaceDatabases") .set({ preprocessingStatus: outcome.status, preprocessedMetadataRevision: outcome.status === "succeeded" ? metadataContentRevision : null, preprocessingFinishedAt: sql`now()`, preprocessingErrorCode: outcome.status === "failed" ? outcome.errorCode : null, updatedAt: sql`now()`, }) .where("workspaceId", "=", workspaceId) .where("preprocessingStatus", "=", "running") .where("metadataContentRevision", "=", metadataContentRevision) .where("preprocessingInputFingerprint", "=", inputFingerprint) .returning("id") .executeTakeFirst(); return result ? await this.get(result.id) : undefined; } async clearPreprocessing(workspaceId: string): Promise { return await this.db.transaction().execute(async (trx) => { const database = await trx.selectFrom("workspaceDatabases") .select(["id", "preprocessingStatus"]) .where("workspaceId", "=", workspaceId) .forUpdate() .executeTakeFirst(); if (!database) return { kind: "not_found" }; if (database.preprocessingStatus === "running") return { kind: "already_running" }; await trx.updateTable("workspaceDatabases").set({ preprocessingStatus: "failed", preprocessingInputFingerprint: null, preprocessedMetadataRevision: null, preprocessingStartedAt: null, preprocessingFinishedAt: sql`now()`, preprocessingErrorCode: "derived_data_cleared", updatedAt: sql`now()`, }).where("id", "=", database.id).execute(); return { kind: "cleared", database: (await selectOne(trx, database.id))! }; }); } async getCatalogMetrics(databaseId?: string): Promise { const result = await sql` 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(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", 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 { try { return await this.db.transaction().execute(async (trx) => { const id = randomUUID(); await trx.insertInto("workspaceDatabases").values({ id, workspaceId: input.workspaceId, engine: input.engine, databaseName: input.databaseName, schemaName: input.schema, }).execute(); await trx.insertInto("databaseBindings").values({ ...bindingValues(id, input.binding), connectionStatus: "untested", testedVersion: null, lastTestedAt: null, lastErrorCode: null, lastErrorMessage: null, }).execute(); return (await selectOne(trx, id))!; }); } catch (error: any) { if (error?.code === "23505") throw new CatalogConflictError("Workspace database already exists"); throw error; } } async update(id: string, expectedVersion: number, input: DatabaseConfigurationInput): Promise { return await this.db.transaction().execute(async (trx) => { const changed = await trx.updateTable("workspaceDatabases").set({ workspaceId: input.workspaceId, engine: input.engine, databaseName: input.databaseName, schemaName: input.schema, version: sql`version + 1`, updatedAt: sql`now()`, }).where("id", "=", id).where("version", "=", expectedVersion).returning("id").executeTakeFirst(); if (!changed) return undefined; await trx.updateTable("databaseBindings").set({ ...bindingValues(id, input.binding), connectionStatus: "untested", testedVersion: null, lastTestedAt: null, lastErrorCode: null, lastErrorMessage: null, updatedAt: sql`now()`, }).where("databaseId", "=", id).execute(); return await selectOne(trx, id); }); } async recordTest(id: string, expectedVersion: number, result: DatabaseTestResult): Promise { const version = await this.db.selectFrom("workspaceDatabases").select("version") .where("id", "=", id).where("version", "=", expectedVersion).executeTakeFirst(); if (!version) return undefined; await this.db.updateTable("databaseBindings").set({ connectionStatus: result.connectionStatus, testedVersion: result.testedVersion, lastTestedAt: result.lastTestedAt, lastErrorCode: result.errorCode ?? null, lastErrorMessage: result.errorMessage ?? null, updatedAt: sql`now()`, }).where("databaseId", "=", id).execute(); return await this.get(id); } async touch(id: string, expectedVersion: number): Promise { const current = await this.get(id); if (!current || current.version !== expectedVersion) return undefined; return await this.update(id, expectedVersion, { workspaceId: current.workspaceId, engine: current.engine, databaseName: current.databaseName, schema: current.schema, binding: current.binding, }); } async delete(id: string, expectedVersion: number): Promise { const result = await this.db.deleteFrom("workspaceDatabases") .where("id", "=", id).where("version", "=", expectedVersion).executeTakeFirst(); return result.numDeletedRows === 1n; } async listTables(databaseId: string): Promise { const rows = await this.db.selectFrom("catalogTables").selectAll() .where("databaseId", "=", databaseId).orderBy("name").execute(); return rows.map(serializeTable); } async getTable(databaseId: string, tableId: string): Promise { const row = await this.db.selectFrom("catalogTables").selectAll() .where("databaseId", "=", databaseId).where("id", "=", tableId).executeTakeFirst(); return row ? serializeTable(row) : undefined; } async updateTableDescription( databaseId: string, tableId: string, expectedVersion: number, description: string | null, ): Promise { const row = await this.db.updateTable("catalogTables").set({ description, version: sql`version + 1`, updatedAt: sql`now()`, }).where("databaseId", "=", databaseId) .where("id", "=", tableId) .where("version", "=", expectedVersion) .returningAll() .executeTakeFirst(); return row ? serializeTable(row) : undefined; } async updateTableMetadata( databaseId: string, tableId: string, expectedVersion: number, description: string | null, generatedDescription: string | null, ): Promise { const row = await this.db.updateTable("catalogTables").set({ description, generatedDescription, version: sql`version + 1`, updatedAt: sql`now()`, }).where("databaseId", "=", databaseId) .where("id", "=", tableId) .where("version", "=", expectedVersion) .returningAll() .executeTakeFirst(); return row ? serializeTable(row) : undefined; } async listColumns(databaseId: string, tableId: string): Promise { const rows = await this.db.selectFrom("catalogColumns") .innerJoin("catalogTables", "catalogTables.id", "catalogColumns.tableId") .selectAll("catalogColumns") .where("catalogTables.databaseId", "=", databaseId) .where("catalogColumns.tableId", "=", tableId) .orderBy("catalogColumns.ordinalPosition") .execute(); if (rows.length === 0) return []; const counts = await this.db.selectFrom("catalogRelationshipColumns") .select(["sourceColumnId", sql`count(*)::int`.as("count")]) .where("sourceColumnId", "in", rows.map((row) => row.id)) .groupBy("sourceColumnId") .execute(); const byColumn = new Map(counts.map((row) => [row.sourceColumnId, Number(row.count)])); return rows.map((row) => serializeColumn(row, byColumn.get(row.id) ?? 0)); } async getColumn(databaseId: string, tableId: string, columnId: string): Promise { const row = await this.db.selectFrom("catalogColumns") .innerJoin("catalogTables", "catalogTables.id", "catalogColumns.tableId") .selectAll("catalogColumns") .where("catalogTables.databaseId", "=", databaseId) .where("catalogColumns.tableId", "=", tableId) .where("catalogColumns.id", "=", columnId) .executeTakeFirst(); if (!row) return undefined; const count = await this.db.selectFrom("catalogRelationshipColumns") .select(sql`count(*)::int`.as("count")) .where("sourceColumnId", "=", columnId) .executeTakeFirst(); return serializeColumn(row, Number(count?.count ?? 0)); } async updateColumnMetadata( databaseId: string, tableId: string, columnId: string, expectedVersion: number, description: string | null, generatedDescription: string | null, sensitive?: boolean, sensitivityReason?: string | null, ): Promise { const belongs = await this.db.selectFrom("catalogTables").select("id") .where("id", "=", tableId).where("databaseId", "=", databaseId).executeTakeFirst(); if (!belongs) return undefined; const row = await this.db.updateTable("catalogColumns").set({ 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) .where("version", "=", expectedVersion).returningAll().executeTakeFirst(); return row ? await this.getColumn(databaseId, tableId, row.id) : undefined; } async consolidateGeneratedDescriptions( databaseId: string, target: CatalogDescriptionTarget, targetIds: readonly string[], ): Promise { const selectedTargetIds = [...new Set(targetIds)]; if (target !== "database" && 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(); if (!database) return undefined; if (target === "tables") { const rows = await trx.selectFrom("catalogTables") .select(["id", "generatedDescription"]) .where("databaseId", "=", databaseId) .where("id", "in", selectedTargetIds) .orderBy("id") .forUpdate() .execute(); if (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("catalogTables").set({ description: sql`generated_description`, version: sql`version + 1`, updatedAt: sql`now()`, }).where("id", "in", copiedIds).execute(); } return { copied: copiedIds.length, skipped: selectedTargetIds.length - copiedIds.length, }; } const tableRows = await trx.selectFrom("catalogTables").select(["id", "generatedDescription"]) .where("databaseId", "=", databaseId) .orderBy("id") .forUpdate() .execute(); const copiedTableIds = target === "database" ? tableRows .filter((row) => Boolean(row.generatedDescription?.trim())) .map((row) => row.id) : []; if (copiedTableIds.length > 0) { await trx.updateTable("catalogTables").set({ description: sql`generated_description`, version: sql`version + 1`, updatedAt: sql`now()`, }).where("id", "in", copiedTableIds).execute(); } const rows = tableRows.length === 0 ? [] : await trx.selectFrom("catalogColumns") .select(["id", "generatedDescription"]) .where("tableId", "in", tableRows.map((table) => table.id)) .$if(target === "columns", (query) => query.where("id", "in", selectedTargetIds)) .orderBy("id") .forUpdate() .execute(); 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) { let update = trx.updateTable("catalogColumns").set({ description: sql`generated_description`, version: sql`version + 1`, updatedAt: sql`now()`, }); update = target === "database" || target === "database_columns" ? update .where("tableId", "in", tableRows.map((table) => table.id)) .where(sql`nullif(btrim(generated_description), '') is not null`) : update.where("id", "in", copiedIds); await update.execute(); } return { copied: copiedTableIds.length + copiedIds.length, skipped: (target === "database" ? tableRows.length : 0) - copiedTableIds.length + rows.length - copiedIds.length, }; }); } async createDescriptionGenerationRun( databaseId: string, scope: DescriptionGenerationScope, modelId: string, language: DescriptionGenerationRun["language"], total: number, ): Promise { try { const row = await this.db.insertInto("descriptionGenerationRuns").values({ id: randomUUID(), databaseId, scope, modelId, language, status: "queued", total, processed: 0, generated: 0, nonGeneratable: 0, failed: 0, inputTokens: 0, cacheReadTokens: 0, outputTokens: 0, startedAt: null, finishedAt: null, errorSummary: null, }).returningAll().executeTakeFirstOrThrow(); return serializeDescriptionGenerationRun(row); } catch (error: any) { if (error?.code === "23505" && error?.constraint === "description_generation_runs_one_active") { throw new DescriptionGenerationRunActiveError( "A description generation run is already active", ); } throw error; } } async getDescriptionGenerationRun( runId: string, ): Promise { const row = await this.db.selectFrom("descriptionGenerationRuns") .selectAll() .where("id", "=", runId) .executeTakeFirst(); return row ? serializeDescriptionGenerationRun(row) : undefined; } async listDescriptionGenerationRuns(limit = 50): Promise { const rows = await this.db.selectFrom("descriptionGenerationRuns") .selectAll() .orderBy("createdAt", "desc") .orderBy("id", "desc") .limit(limit) .execute(); return rows.map(serializeDescriptionGenerationRun); } async getActiveDescriptionGenerationRun(): Promise { const row = await this.db.selectFrom("descriptionGenerationRuns") .selectAll() .where("status", "in", ["queued", "running"]) .orderBy("createdAt", "desc") .executeTakeFirst(); return row ? serializeDescriptionGenerationRun(row) : undefined; } async interruptActiveDescriptionGenerationRuns( errorSummary: string, ): Promise { const rows = await this.db.updateTable("descriptionGenerationRuns") .set({ status: "interrupted", finishedAt: sql`now()`, updatedAt: sql`now()`, errorSummary, }) .where("status", "in", ["queued", "running"]) .returningAll() .execute(); return rows.map(serializeDescriptionGenerationRun); } async updateDescriptionGenerationRun( runId: string, update: DescriptionGenerationRunUpdate, ): Promise { const values: any = { ...update, updatedAt: sql`now()` }; const row = await this.db.updateTable("descriptionGenerationRuns") .set(values) .where("id", "=", runId) .returningAll() .executeTakeFirst(); return row ? serializeDescriptionGenerationRun(row) : undefined; } async appendDescriptionGenerationEvent( runId: string, level: DescriptionGenerationEvent["level"], message: string, ): Promise { return await this.db.transaction().execute(async (trx) => { const run = await trx.selectFrom("descriptionGenerationRuns") .select("id") .where("id", "=", runId) .forUpdate() .executeTakeFirst(); if (!run) throw new CatalogConflictError("Description Generation Run does not exist"); const current = await trx.selectFrom("descriptionGenerationEvents") .select(sql`coalesce(max(sequence), 0)::int`.as("sequence")) .where("runId", "=", runId) .executeTakeFirst(); const row = await trx.insertInto("descriptionGenerationEvents").values({ runId, sequence: Number(current?.sequence ?? 0) + 1, level, message, }).returningAll().executeTakeFirstOrThrow(); return serializeDescriptionGenerationEvent(row); }); } async listDescriptionGenerationEvents( runId: string, afterSequence = 0, ): Promise { const rows = await this.db.selectFrom("descriptionGenerationEvents") .selectAll() .where("runId", "=", runId) .where("sequence", ">", afterSequence) .orderBy("sequence") .execute(); return rows.map(serializeDescriptionGenerationEvent); } async createSensitivityAnalysisRun( databaseId: string, scope: SensitivityAnalysisScope, origin: { engine: "llm"; modelId: string } | { engine: "local"; policyVersion: string }, ): Promise { const row = await this.db.insertInto("sensitiveDataSuggestionRuns").values({ id: randomUUID(), databaseId, scope, engine: origin.engine, modelId: origin.engine === "llm" ? origin.modelId : null, policyVersion: origin.engine === "local" ? origin.policyVersion : null, status: "running", total: 0, suggestedSensitive: 0, suggestedNonSensitive: 0, unknown: 0, inputTokens: 0, cacheReadTokens: 0, outputTokens: 0, finishedAt: null, errorSummary: null, }).returningAll().executeTakeFirstOrThrow(); return serializeSensitivityAnalysisRun(row); } async getSensitivityAnalysisRun( runId: string, ): Promise { const row = await this.db.selectFrom("sensitiveDataSuggestionRuns") .selectAll() .where("id", "=", runId) .executeTakeFirst(); return row ? serializeSensitivityAnalysisRun(row) : undefined; } async listSensitivityAnalysisRuns(limit = 50): Promise { const rows = await this.db.selectFrom("sensitiveDataSuggestionRuns") .selectAll() .orderBy("createdAt", "desc") .orderBy("id", "desc") .limit(limit) .execute(); return rows.map(serializeSensitivityAnalysisRun); } async interruptActiveSensitivityAnalysisRuns( errorSummary: string, ): Promise { 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(serializeSensitivityAnalysisRun); } async updateSensitivityAnalysisRun( runId: string, update: SensitivityAnalysisRunUpdate, ): Promise { const values: any = { ...update, updatedAt: sql`now()` }; const row = await this.db.updateTable("sensitiveDataSuggestionRuns") .set(values) .where("id", "=", runId) .returningAll() .executeTakeFirst(); return row ? serializeSensitivityAnalysisRun(row) : undefined; } async appendSensitivityAnalysisEvent( runId: string, level: SensitivityAnalysisEvent["level"], message: string, ): Promise { 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("Sensitivity Analysis Run does not exist"); const current = await trx.selectFrom("sensitiveDataSuggestionEvents") .select(sql`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 serializeSensitivityAnalysisEvent(row); }); } async listSensitivityAnalysisEvents( runId: string, afterSequence = 0, ): Promise { const rows = await this.db.selectFrom("sensitiveDataSuggestionEvents") .selectAll() .where("runId", "=", runId) .where("sequence", ">", afterSequence) .orderBy("sequence") .execute(); return rows.map(serializeSensitivityAnalysisEvent); } async listRelationships(databaseId: string): Promise { const rows = await this.db.selectFrom("catalogRelationships as relationship") .innerJoin("catalogTables as sourceTable", "sourceTable.id", "relationship.sourceTableId") .innerJoin("catalogTables as targetTable", "targetTable.id", "relationship.targetTableId") .select([ "relationship.id", "relationship.databaseId", "relationship.constraintName", "relationship.sourceTableId", "sourceTable.name as sourceTableName", "relationship.targetTableId", "targetTable.name as targetTableName", "relationship.updateRule", "relationship.deleteRule", "relationship.deferrable", "relationship.initiallyDeferred", "relationship.lastSyncedDatabaseVersion", "relationship.lastSyncedAt", "relationship.createdAt", "relationship.updatedAt", ]) .where("relationship.databaseId", "=", databaseId) .orderBy("sourceTable.name").orderBy("relationship.constraintName") .execute(); if (rows.length === 0) return []; const pairs = await this.db.selectFrom("catalogRelationshipColumns as pair") .innerJoin("catalogColumns as sourceColumn", "sourceColumn.id", "pair.sourceColumnId") .innerJoin("catalogColumns as targetColumn", "targetColumn.id", "pair.targetColumnId") .select([ "pair.relationshipId", "pair.position", "pair.sourceColumnId", "sourceColumn.name as sourceColumnName", "pair.targetColumnId", "targetColumn.name as targetColumnName", ]) .where("pair.relationshipId", "in", rows.map((row) => row.id)) .orderBy("pair.relationshipId").orderBy("pair.position") .execute(); const byRelationship = new Map(); for (const pair of pairs) { const items = byRelationship.get(pair.relationshipId) ?? []; items.push(pair); byRelationship.set(pair.relationshipId, items); } return rows.map((row) => ({ ...row, columns: byRelationship.get(row.id) ?? [], lastSyncedAt: row.lastSyncedAt === null ? null : new Date(row.lastSyncedAt).toISOString(), createdAt: new Date(row.createdAt).toISOString(), updatedAt: new Date(row.updatedAt).toISOString(), origin: "physical" as const, status: "active" as const, })); } async listLogicalRelationships(databaseId: string): Promise { const rows = await this.db.selectFrom("catalogLogicalRelationships as relationship") .innerJoin("catalogColumns as sourceColumn", "sourceColumn.id", "relationship.sourceColumnId") .innerJoin("catalogTables as sourceTable", "sourceTable.id", "sourceColumn.tableId") .innerJoin("catalogColumns as targetColumn", "targetColumn.id", "relationship.targetColumnId") .innerJoin("catalogTables as targetTable", "targetTable.id", "targetColumn.tableId") .select([ "relationship.id", "relationship.databaseId", "relationship.generated", "relationship.deletedAt", "relationship.createdAt", "relationship.updatedAt", "sourceTable.id as sourceTableId", "sourceTable.name as sourceTableName", "sourceColumn.id as sourceColumnId", "sourceColumn.name as sourceColumnName", "targetTable.id as targetTableId", "targetTable.name as targetTableName", "targetColumn.id as targetColumnId", "targetColumn.name as targetColumnName", ]) .where("relationship.databaseId", "=", databaseId) .orderBy("sourceTable.name") .orderBy("sourceColumn.name") .orderBy("targetTable.name") .orderBy("targetColumn.name") .execute(); return rows.map((row) => ({ id: row.id, databaseId: row.databaseId, constraintName: null, sourceTableId: row.sourceTableId, sourceTableName: row.sourceTableName, targetTableId: row.targetTableId, targetTableName: row.targetTableName, updateRule: null, deleteRule: null, deferrable: false, initiallyDeferred: false, columns: [{ position: 1, sourceColumnId: row.sourceColumnId, sourceColumnName: row.sourceColumnName, targetColumnId: row.targetColumnId, targetColumnName: row.targetColumnName, }], lastSyncedDatabaseVersion: null, lastSyncedAt: null, createdAt: new Date(row.createdAt).toISOString(), updatedAt: new Date(row.updatedAt).toISOString(), origin: row.generated ? "generated" : "manual", status: row.deletedAt === null ? "active" : "excluded", })); } async getLogicalRelationshipContext( databaseId: string, ): Promise { if (!(await selectOne(this.db, databaseId))) return undefined; const columns = await this.db.selectFrom("catalogColumns as column") .innerJoin("catalogTables as table", "table.id", "column.tableId") .select([ "column.id as columnId", "column.name as columnName", "column.dataType", "column.primaryKeyPosition", "table.id as tableId", "table.name as tableName", ]) .where("table.databaseId", "=", databaseId) .orderBy("table.name") .orderBy("column.ordinalPosition") .execute(); const primaryKeyCounts = new Map(); for (const column of columns) { if (column.primaryKeyPosition !== null) { primaryKeyCounts.set(column.tableId, (primaryKeyCounts.get(column.tableId) ?? 0) + 1); } } const physicalPairs = await this.db.selectFrom("catalogRelationshipColumns as pair") .innerJoin("catalogRelationships as relationship", "relationship.id", "pair.relationshipId") .select(["pair.sourceColumnId", "pair.targetColumnId"]) .where("relationship.databaseId", "=", databaseId) .execute(); return { endpoints: columns.map((column) => ({ ...column, tablePrimaryKeyColumnCount: primaryKeyCounts.get(column.tableId) ?? 0, })), physicalPairs, logicalRelationships: await this.listLogicalRelationships(databaseId), }; } async insertLogicalRelationship( databaseId: string, sourceColumnId: string, targetColumnId: string, generated: boolean, ): Promise { const id = randomUUID(); const inserted = await this.db.insertInto("catalogLogicalRelationships").values({ id, databaseId, sourceColumnId, targetColumnId, generated, }).onConflict((conflict) => conflict .columns(["databaseId", "sourceColumnId", "targetColumnId"]) .doNothing()) .returning("id") .executeTakeFirst(); if (!inserted) return undefined; return (await this.listLogicalRelationships(databaseId)).find((item) => item.id === id); } async insertGeneratedLogicalRelationships( databaseId: string, candidates: readonly CatalogLogicalRelationshipCandidate[], ): Promise { if (candidates.length === 0) return 0; return await this.db.transaction().execute(async (trx) => { let added = 0; for (const candidate of candidates) { const inserted = await trx.insertInto("catalogLogicalRelationships").values({ id: randomUUID(), databaseId, sourceColumnId: candidate.sourceColumnId, targetColumnId: candidate.targetColumnId, generated: true, }).onConflict((conflict) => conflict .columns(["databaseId", "sourceColumnId", "targetColumnId"]) .doNothing()) .returning("id") .executeTakeFirst(); if (inserted) added += 1; } return added; }); } async setLogicalRelationshipStatus( databaseId: string, relationshipId: string, status: CatalogLogicalRelationship["status"], ): Promise { const updated = await this.db.updateTable("catalogLogicalRelationships") .set({ deletedAt: status === "excluded" ? sql`now()` : null, updatedAt: sql`now()` }) .where("databaseId", "=", databaseId) .where("id", "=", relationshipId) .returning("id") .executeTakeFirst(); if (!updated) return undefined; return (await this.listLogicalRelationships(databaseId)).find((item) => item.id === relationshipId); } async deleteLogicalRelationship(databaseId: string, relationshipId: string): Promise { const result = await this.db.deleteFrom("catalogLogicalRelationships") .where("databaseId", "=", databaseId) .where("id", "=", relationshipId) .executeTakeFirst(); return result.numDeletedRows > 0n; } async deleteDatabaseMetadata( databaseIds: readonly string[], target: CatalogDatabaseMetadataDeleteTarget, ): Promise { const selectedDatabaseIds = [...new Set(databaseIds)]; if (selectedDatabaseIds.length === 0) return undefined; return await this.db.transaction().execute(async (trx) => { const databases = await trx.selectFrom("workspaceDatabases").select("id") .where("id", "in", selectedDatabaseIds).orderBy("id").forUpdate().execute(); if (databases.length !== selectedDatabaseIds.length) return undefined; if (target === "relationships") { const physicalCount = await trx.selectFrom("catalogRelationships") .select(sql`count(*)::int`.as("count")) .where("databaseId", "in", selectedDatabaseIds).executeTakeFirst(); const logicalCount = await trx.selectFrom("catalogLogicalRelationships") .select(sql`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(physicalCount?.count ?? 0) + Number(logicalCount?.count ?? 0), }; } const tableCount = await trx.selectFrom("catalogTables") .select(sql`count(*)::int`.as("count")) .where("databaseId", "in", selectedDatabaseIds).executeTakeFirst(); const columnCount = await trx.selectFrom("catalogColumns") .innerJoin("catalogTables", "catalogTables.id", "catalogColumns.tableId") .select(sql`count(*)::int`.as("count")) .where("catalogTables.databaseId", "in", selectedDatabaseIds).executeTakeFirst(); const physicalRelationshipCount = await trx.selectFrom("catalogRelationships") .select(sql`count(*)::int`.as("count")) .where("databaseId", "in", selectedDatabaseIds).executeTakeFirst(); const logicalRelationshipCount = await trx.selectFrom("catalogLogicalRelationships") .select(sql`count(*)::int`.as("count")) .where("databaseId", "in", selectedDatabaseIds).executeTakeFirst(); await trx.deleteFrom("catalogTables") .where("databaseId", "in", selectedDatabaseIds).execute(); await trx.updateTable("workspaceDatabases").set({ schemaSyncedVersion: null, schemaSyncedAt: null, }).where("id", "in", selectedDatabaseIds).execute(); return { tables: Number(tableCount?.count ?? 0), columns: Number(columnCount?.count ?? 0), relationships: Number(physicalRelationshipCount?.count ?? 0) + Number(logicalRelationshipCount?.count ?? 0), }; }); } async deleteTableMetadata( databaseId: string, tableIds: readonly string[], target: CatalogTableMetadataDeleteTarget, ): Promise { const selectedTableIds = [...new Set(tableIds)]; if (selectedTableIds.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(); if (!database) return undefined; const tables = await trx.selectFrom("catalogTables").select("id") .where("databaseId", "=", databaseId) .where("id", "in", selectedTableIds).execute(); if (tables.length !== selectedTableIds.length) return undefined; if (target === "columns") { const count = await trx.selectFrom("catalogColumns") .select(sql`count(*)::int`.as("count")) .where("tableId", "in", selectedTableIds).executeTakeFirst(); await trx.deleteFrom("catalogColumns") .where("tableId", "in", selectedTableIds).execute(); await trx.updateTable("workspaceDatabases").set({ schemaSyncedVersion: null, schemaSyncedAt: null, }).where("id", "=", databaseId).execute(); return { tables: 0, columns: Number(count?.count ?? 0), relationships: 0 }; } const physicalCount = await trx.selectFrom("catalogRelationships") .select(sql`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`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([ eb("sourceTableId", "in", selectedTableIds), eb("targetTableId", "in", selectedTableIds), ])).execute(); await trx.updateTable("workspaceDatabases").set({ schemaSyncedVersion: null, schemaSyncedAt: null, }).where("id", "=", databaseId).execute(); return { tables: 0, columns: 0, relationships: Number(physicalCount?.count ?? 0) + Number(logicalCount?.count ?? 0), }; }); } async planSchemaSync( databaseId: string, scope: CatalogSyncScope, tableIds: readonly string[], snapshot: ObservedSchemaSnapshot, ): Promise { const tables = await this.db.selectFrom("catalogTables").select(["id", "name"]) .where("databaseId", "=", databaseId).execute(); const selected = new Set(tableIds); const targetTables = scope === "columns" && selected.size > 0 ? tables.filter((table) => selected.has(table.id)) : tables; const observedTableNames = new Set(snapshot.tables.map((table) => table.name)); const deletedTables = scope === "tables" || scope === "all" ? tables.filter((table) => !observedTableNames.has(table.name)).map((table) => table.name).sort() : []; const deletedTableIds = new Set(tables.filter((table) => deletedTables.includes(table.name)).map((table) => table.id)); const deletedColumns: CatalogSchemaDiff["deletedColumns"] = []; if (scope === "tables" || scope === "columns" || scope === "all") { const columnTables = scope === "tables" ? tables.filter((table) => deletedTableIds.has(table.id)) : targetTables; const columns = columnTables.length === 0 ? [] : await this.db.selectFrom("catalogColumns") .select(["tableId", "name"]).where("tableId", "in", columnTables.map((table) => table.id)).execute(); const tableNameById = new Map(columnTables.map((table) => [table.id, table.name])); const observed = new Set(snapshot.columns.map((column) => `${column.tableName}\u0000${column.name}`)); for (const column of columns) { const tableName = tableNameById.get(column.tableId)!; if (scope === "tables" || !observed.has(`${tableName}\u0000${column.name}`)) { deletedColumns.push({ tableName, columnName: column.name }); } } deletedColumns.sort((a, b) => `${a.tableName}.${a.columnName}`.localeCompare(`${b.tableName}.${b.columnName}`)); } const deletedRelationships: CatalogSchemaDiff["deletedRelationships"] = []; if (scope === "tables" || scope === "relationships" || scope === "all") { const relationships = await this.db.selectFrom("catalogRelationships as relationship") .innerJoin("catalogTables as sourceTable", "sourceTable.id", "relationship.sourceTableId") .select(["relationship.constraintName", "relationship.sourceTableId", "relationship.targetTableId", "sourceTable.name as sourceTableName"]) .where("relationship.databaseId", "=", databaseId).execute(); const observed = new Set(snapshot.relationships.map((relationship) => `${relationship.sourceTableName}\u0000${relationship.constraintName}`)); for (const relationship of relationships) { if ( (scope === "tables" && (deletedTableIds.has(relationship.sourceTableId) || deletedTableIds.has(relationship.targetTableId))) || (scope !== "tables" && !observed.has(`${relationship.sourceTableName}\u0000${relationship.constraintName}`)) ) { deletedRelationships.push({ sourceTableName: relationship.sourceTableName, constraintName: relationship.constraintName, }); } } deletedRelationships.sort((a, b) => `${a.sourceTableName}.${a.constraintName}`.localeCompare(`${b.sourceTableName}.${b.constraintName}`)); } return { deletedTables, deletedColumns, deletedRelationships }; } async applySchemaSync( databaseId: string, expectedDatabaseVersion: number, scope: CatalogSyncScope, tableIds: readonly string[], snapshot: ObservedSchemaSnapshot, syncRunId?: string, ): Promise { return await this.db.transaction().execute(async (trx) => { const database = await trx.selectFrom("workspaceDatabases").select("version") .where("id", "=", databaseId).forUpdate().executeTakeFirst(); if (!database || database.version !== expectedDatabaseVersion) return undefined; const now = new Date().toISOString(); let created = 0; let updated = 0; let deleted = 0; let tables = await trx.selectFrom("catalogTables").selectAll() .where("databaseId", "=", databaseId).execute(); if (scope === "tables" || scope === "all") { const byName = new Map(tables.map((table) => [table.name, table])); const observedNames = new Set(snapshot.tables.map((table) => table.name)); const removed = tables.filter((table) => !observedNames.has(table.name)); if (removed.length > 0) { await trx.deleteFrom("catalogTables").where("id", "in", removed.map((table) => table.id)).execute(); deleted += removed.length; } for (const observed of snapshot.tables) { const current = byName.get(observed.name); if (!current) { await trx.insertInto("catalogTables").values({ id: randomUUID(), databaseId, name: observed.name, sourceComment: observed.sourceComment, description: null, generatedDescription: null, lastSyncedDatabaseVersion: expectedDatabaseVersion, lastSyncedAt: now, }).execute(); created += 1; } else { const changed = current.sourceComment !== observed.sourceComment; await trx.updateTable("catalogTables").set({ sourceComment: observed.sourceComment, lastSyncedDatabaseVersion: expectedDatabaseVersion, lastSyncedAt: now, ...(changed ? { version: sql`version + 1`, updatedAt: sql`now()` } : {}), }).where("id", "=", current.id).execute(); if (changed) updated += 1; } } tables = await trx.selectFrom("catalogTables").selectAll().where("databaseId", "=", databaseId).execute(); } if (scope === "columns" || scope === "all") { const selected = new Set(tableIds); const targetTables = scope === "columns" && selected.size > 0 ? tables.filter((table) => selected.has(table.id)) : tables; for (const table of targetTables) { const observedColumns = snapshot.columns.filter((column) => column.tableName === table.name); const existing = await trx.selectFrom("catalogColumns").selectAll().where("tableId", "=", table.id).execute(); const observedNames = new Set(observedColumns.map((column) => column.name)); const removed = existing.filter((column) => !observedNames.has(column.name)); if (removed.length > 0) { await trx.deleteFrom("catalogColumns").where("id", "in", removed.map((column) => column.id)).execute(); deleted += removed.length; } const byName = new Map(existing.map((column) => [column.name, column])); for (const observed of observedColumns) { const current = byName.get(observed.name); if (!current) { await trx.insertInto("catalogColumns").values({ id: randomUUID(), tableId: table.id, name: observed.name, ordinalPosition: observed.ordinalPosition, dataType: observed.dataType, isNullable: observed.isNullable, defaultExpression: observed.defaultExpression, primaryKeyPosition: observed.primaryKeyPosition, sourceComment: observed.sourceComment, description: null, generatedDescription: null, lastSyncedDatabaseVersion: expectedDatabaseVersion, lastSyncedAt: now, }).execute(); created += 1; } else { const changed = current.ordinalPosition !== observed.ordinalPosition || current.dataType !== observed.dataType || current.isNullable !== observed.isNullable || current.defaultExpression !== observed.defaultExpression || current.primaryKeyPosition !== observed.primaryKeyPosition || current.sourceComment !== observed.sourceComment; await trx.updateTable("catalogColumns").set({ ordinalPosition: observed.ordinalPosition, dataType: observed.dataType, isNullable: observed.isNullable, defaultExpression: observed.defaultExpression, primaryKeyPosition: observed.primaryKeyPosition, sourceComment: observed.sourceComment, lastSyncedDatabaseVersion: expectedDatabaseVersion, lastSyncedAt: now, ...(changed ? { version: sql`version + 1`, updatedAt: sql`now()` } : {}), }).where("id", "=", current.id).execute(); if (changed) updated += 1; } } } } if (scope === "relationships" || scope === "all") { tables = await trx.selectFrom("catalogTables").selectAll().where("databaseId", "=", databaseId).execute(); const tableByName = new Map(tables.map((table) => [table.name, table])); const existing = await trx.selectFrom("catalogRelationships").selectAll() .where("databaseId", "=", databaseId).execute(); const tableNameById = new Map(tables.map((table) => [table.id, table.name])); const observedKeys = new Set(snapshot.relationships.map((relationship) => `${relationship.sourceTableName}\u0000${relationship.constraintName}`)); const removed = existing.filter((relationship) => !observedKeys.has(`${tableNameById.get(relationship.sourceTableId)}\u0000${relationship.constraintName}`)); if (removed.length > 0) { await trx.deleteFrom("catalogRelationships").where("id", "in", removed.map((relationship) => relationship.id)).execute(); deleted += removed.length; } const byKey = new Map(existing.map((relationship) => [`${tableNameById.get(relationship.sourceTableId)}\u0000${relationship.constraintName}`, relationship])); for (const observed of snapshot.relationships) { const sourceTable = tableByName.get(observed.sourceTableName); const targetTable = tableByName.get(observed.targetTableName); if (!sourceTable || !targetTable) throw new CatalogConnectorError("Relationship endpoint is not present in the catalog"); const sourceColumns = await trx.selectFrom("catalogColumns").select(["id", "name"]) .where("tableId", "=", sourceTable.id).execute(); const targetColumns = await trx.selectFrom("catalogColumns").select(["id", "name"]) .where("tableId", "=", targetTable.id).execute(); const sourceByName = new Map(sourceColumns.map((column) => [column.name, column.id])); const targetByName = new Map(targetColumns.map((column) => [column.name, column.id])); const key = `${observed.sourceTableName}\u0000${observed.constraintName}`; const current = byKey.get(key); const relationshipId = current?.id ?? randomUUID(); const desiredPairs = observed.columns.map((pair) => { const sourceColumnId = sourceByName.get(pair.sourceColumnName); const targetColumnId = targetByName.get(pair.targetColumnName); if (!sourceColumnId || !targetColumnId) { throw new CatalogConnectorError("Relationship column is not present in the catalog"); } return { relationshipId, position: pair.position, sourceColumnId, targetColumnId }; }).sort((left, right) => left.position - right.position); const currentPairs = current ? await trx.selectFrom("catalogRelationshipColumns") .select(["position", "sourceColumnId", "targetColumnId"]) .where("relationshipId", "=", current.id).orderBy("position").execute() : []; const changed = Boolean(current && ( current.targetTableId !== targetTable.id || current.updateRule !== observed.updateRule || current.deleteRule !== observed.deleteRule || current.deferrable !== observed.deferrable || current.initiallyDeferred !== observed.initiallyDeferred || JSON.stringify(currentPairs) !== JSON.stringify(desiredPairs.map( ({ position, sourceColumnId, targetColumnId }) => ({ position, sourceColumnId, targetColumnId }), )) )); if (!current) { await trx.insertInto("catalogRelationships").values({ id: relationshipId, databaseId, constraintName: observed.constraintName, sourceTableId: sourceTable.id, targetTableId: targetTable.id, updateRule: observed.updateRule, deleteRule: observed.deleteRule, deferrable: observed.deferrable, initiallyDeferred: observed.initiallyDeferred, lastSyncedDatabaseVersion: expectedDatabaseVersion, lastSyncedAt: now, }).execute(); created += 1; } else { await trx.updateTable("catalogRelationships").set({ targetTableId: targetTable.id, updateRule: observed.updateRule, deleteRule: observed.deleteRule, deferrable: observed.deferrable, initiallyDeferred: observed.initiallyDeferred, lastSyncedDatabaseVersion: expectedDatabaseVersion, lastSyncedAt: now, ...(changed ? { updatedAt: sql`now()` } : {}), }).where("id", "=", current.id).execute(); if (changed) { await trx.deleteFrom("catalogRelationshipColumns").where("relationshipId", "=", current.id).execute(); updated += 1; } } if (!current || changed) { for (const pair of desiredPairs) { await trx.insertInto("catalogRelationshipColumns").values(pair).execute(); } } } } if (scope === "all") { await trx.updateTable("workspaceDatabases").set({ schemaSyncedVersion: expectedDatabaseVersion, schemaSyncedAt: now, }).where("id", "=", databaseId).execute(); } if (syncRunId) { await trx.updateTable("catalogSyncRuns").set({ phase: "memory_cleanup" }) .where("id", "=", syncRunId).where("databaseId", "=", databaseId).execute(); } return { tables: snapshot.tables.length, columns: scope === "tables" ? undefined : snapshot.columns.length, relationships: scope === "tables" || scope === "columns" ? undefined : snapshot.relationships.length, created, updated, deleted, }; }); } async createSyncRun( databaseId: string, scope: CatalogSyncScope, tableIds: readonly string[], requestedDatabaseVersion: number, ): Promise { try { const row = await this.db.insertInto("catalogSyncRuns").values({ id: randomUUID(), databaseId, scope, tableIds: JSON.stringify([...tableIds]), state: "queued", phase: "queued", requestedDatabaseVersion, observedSnapshot: null, plannedDiff: null, confirmationToken: null, counts: {}, errorCode: null, errorMessage: null, cancelRequested: false, startedAt: null, finishedAt: null, heartbeatAt: null, leaseOwner: null, leaseExpiresAt: null, }).returningAll().executeTakeFirstOrThrow(); return serializeSyncRun(row); } catch (error: any) { if (error?.code === "23505") throw new CatalogConflictError("A synchronization is already active for this database"); throw error; } } async getSyncRun(runId: string): Promise { const row = await this.db.selectFrom("catalogSyncRuns").selectAll().where("id", "=", runId).executeTakeFirst(); return row ? serializeSyncRun(row) : undefined; } async claimSyncRun(runId: string, workerId: string, leaseExpiresAt: string): Promise { const now = new Date(); const row = await this.db.updateTable("catalogSyncRuns").set({ state: "running", leaseOwner: workerId, leaseExpiresAt, heartbeatAt: now, startedAt: sql`coalesce(started_at, now())`, updatedAt: now, }).where("id", "=", runId) .where("state", "=", "queued") .where((eb) => eb.or([ eb("leaseOwner", "is", null), eb("leaseExpiresAt", "<", now), eb("leaseOwner", "=", workerId), ])) .returningAll().executeTakeFirst(); return row ? serializeSyncRun(row) : undefined; } async listSyncRuns(databaseId: string, limit = 50): Promise { const rows = await this.db.selectFrom("catalogSyncRuns").selectAll() .where("databaseId", "=", databaseId).orderBy("createdAt", "desc").limit(limit).execute(); return rows.map(serializeSyncRun); } async updateSyncRun(runId: string, update: CatalogSyncRunUpdate): Promise { const values: any = { ...update, updatedAt: sql`now()` }; const row = await this.db.updateTable("catalogSyncRuns").set(values) .where("id", "=", runId).returningAll().executeTakeFirst(); return row ? serializeSyncRun(row) : undefined; } async requestSyncRunCancellation(runId: string): Promise { return await this.updateSyncRun(runId, { cancelRequested: true }); } async appendSyncEvent( runId: string, level: CatalogSyncEvent["level"], eventType: string, message: string, data: Record = {}, ): Promise { return await this.db.transaction().execute(async (trx) => { const current = await trx.selectFrom("catalogSyncEvents") .select(sql`coalesce(max(sequence), 0)::int`.as("sequence")) .where("runId", "=", runId).executeTakeFirst(); const row = await trx.insertInto("catalogSyncEvents").values({ runId, sequence: Number(current?.sequence ?? 0) + 1, level, eventType, message, data, }).returningAll().executeTakeFirstOrThrow(); return serializeSyncEvent(row); }); } async listSyncEvents(runId: string, afterSequence = 0): Promise { const rows = await this.db.selectFrom("catalogSyncEvents").selectAll() .where("runId", "=", runId).where("sequence", ">", afterSequence) .orderBy("sequence").execute(); return rows.map(serializeSyncEvent); } async pruneSyncEvents(before: string): Promise { await this.db.deleteFrom("catalogSyncEvents").where("createdAt", "<", new Date(before)).execute(); } async interruptActiveSyncRuns(): Promise { await this.db.updateTable("catalogSyncRuns").set({ state: "interrupted", phase: sql`case when phase='memory_cleanup' then phase else 'completed' end`, errorCode: "worker_restarted", errorMessage: "Synchronization was interrupted by a backend restart.", finishedAt: sql`now()`, updatedAt: sql`now()`, leaseOwner: null, leaseExpiresAt: null, }).where("state", "in", ["queued", "running", "awaiting_confirmation", "applying"]).execute(); } async reconcileTables( databaseId: string, expectedDatabaseVersion: number, observed: readonly ObservedCatalogTable[], confirmedDeletedNames: readonly string[], ): Promise { return await this.db.transaction().execute(async (trx) => { const database = await trx.selectFrom("workspaceDatabases").select("version") .where("id", "=", databaseId).forUpdate().executeTakeFirst(); if (!database || database.version !== expectedDatabaseVersion) return undefined; const existing = await trx.selectFrom("catalogTables").selectAll() .where("databaseId", "=", databaseId).orderBy("name").execute(); const byName = new Map(existing.map((table) => [table.name, table])); const observedNames = new Set(observed.map((table) => table.name)); const deletedNames = existing .filter((table) => !observedNames.has(table.name)) .map((table) => table.name) .sort(); const confirmation = [...new Set(confirmedDeletedNames)].sort(); if (deletedNames.length > 0 && JSON.stringify(deletedNames) !== JSON.stringify(confirmation)) { return { kind: "confirmation_required", deletedNames }; } let createdCount = 0; let updatedCount = 0; for (const table of observed) { const current = byName.get(table.name); if (!current) { await trx.insertInto("catalogTables").values({ id: randomUUID(), databaseId, name: table.name, sourceComment: table.sourceComment, description: null, generatedDescription: null, }).execute(); createdCount += 1; } else if (current.sourceComment !== table.sourceComment) { await trx.updateTable("catalogTables").set({ sourceComment: table.sourceComment, version: sql`version + 1`, updatedAt: sql`now()`, }).where("id", "=", current.id).execute(); updatedCount += 1; } } if (deletedNames.length > 0) { await trx.deleteFrom("catalogTables") .where("databaseId", "=", databaseId) .where("name", "in", deletedNames) .execute(); } const rows = await trx.selectFrom("catalogTables").selectAll() .where("databaseId", "=", databaseId).orderBy("name").execute(); return { kind: "applied", createdCount, updatedCount, deletedCount: deletedNames.length, tables: rows.map(serializeTable), }; }); } async available(): Promise { try { await sql`select 1`.execute(this.db); return true; } catch { return false; } } async close(): Promise { await this.db.destroy(); } } export class UnavailableCatalogRepository implements CatalogRepository { private fail(): never { throw new CatalogUnavailableError("Database catalog is unavailable"); } async list(): Promise { return this.fail(); } async get(): Promise { return this.fail(); } async getByWorkspace(): Promise { return this.fail(); } async beginPreprocessing(): Promise { return this.fail(); } async finishPreprocessing(): Promise { return this.fail(); } async clearPreprocessing(): Promise { return this.fail(); } async getCatalogMetrics(): Promise { return this.fail(); } async create(): Promise { return this.fail(); } async update(): Promise { return this.fail(); } async recordTest(): Promise { return this.fail(); } async touch(): Promise { return this.fail(); } async delete(): Promise { return this.fail(); } async listTables(): Promise { return this.fail(); } async getTable(): Promise { return this.fail(); } async updateTableDescription(): Promise { return this.fail(); } async updateTableMetadata(): Promise { return this.fail(); } async listColumns(): Promise { return this.fail(); } async getColumn(): Promise { return this.fail(); } async updateColumnMetadata(): Promise { return this.fail(); } async consolidateGeneratedDescriptions(): Promise { return this.fail(); } async createDescriptionGenerationRun(): Promise { return this.fail(); } async getDescriptionGenerationRun(): Promise { return this.fail(); } async listDescriptionGenerationRuns(): Promise { return this.fail(); } async getActiveDescriptionGenerationRun(): Promise { return this.fail(); } async interruptActiveDescriptionGenerationRuns(): Promise { return this.fail(); } async updateDescriptionGenerationRun(): Promise { return this.fail(); } async appendDescriptionGenerationEvent(): Promise { return this.fail(); } async listDescriptionGenerationEvents(): Promise { return this.fail(); } async createSensitivityAnalysisRun(): Promise { return this.fail(); } async getSensitivityAnalysisRun(): Promise { return this.fail(); } async listSensitivityAnalysisRuns(): Promise { return this.fail(); } async interruptActiveSensitivityAnalysisRuns(): Promise { return this.fail(); } async updateSensitivityAnalysisRun(): Promise { return this.fail(); } async appendSensitivityAnalysisEvent(): Promise { return this.fail(); } async listSensitivityAnalysisEvents(): Promise { return this.fail(); } async listRelationships(): Promise { return this.fail(); } async listLogicalRelationships(): Promise { return this.fail(); } async getLogicalRelationshipContext(): Promise { return this.fail(); } async insertLogicalRelationship(): Promise { return this.fail(); } async insertGeneratedLogicalRelationships(): Promise { return this.fail(); } async setLogicalRelationshipStatus(): Promise { return this.fail(); } async deleteLogicalRelationship(): Promise { return this.fail(); } async deleteDatabaseMetadata(): Promise { return this.fail(); } async deleteTableMetadata(): Promise { return this.fail(); } async planSchemaSync(): Promise { return this.fail(); } async applySchemaSync(): Promise { return this.fail(); } async createSyncRun(): Promise { return this.fail(); } async getSyncRun(): Promise { return this.fail(); } async claimSyncRun(): Promise { return this.fail(); } async listSyncRuns(): Promise { return this.fail(); } async updateSyncRun(): Promise { return this.fail(); } async requestSyncRunCancellation(): Promise { return this.fail(); } async appendSyncEvent(): Promise { return this.fail(); } async listSyncEvents(): Promise { return this.fail(); } async pruneSyncEvents(): Promise { return this.fail(); } async interruptActiveSyncRuns(): Promise { return this.fail(); } async reconcileTables(): Promise { return this.fail(); } async available(): Promise { return false; } } export interface CatalogConnectionConfig { connectionString?: string; host?: string; port?: number; database?: string; user?: string; passwordFile?: string; } export function catalogPool(config: CatalogConnectionConfig, max = 5): Pool { return new Pool(config.connectionString ? { connectionString: config.connectionString, max, connectionTimeoutMillis: 3_000 } : { host: config.host, port: config.port, database: config.database, user: config.user, password: async () => (await readFile(config.passwordFile!, "utf8")).trim(), max, connectionTimeoutMillis: 3_000, }); } export function createCatalogRepository(config: CatalogConnectionConfig | undefined): CatalogRepository { if (!config) return new UnavailableCatalogRepository(); const pool = catalogPool(config); const db = new Kysely({ dialect: new PostgresDialect({ pool }), plugins: [new CamelCasePlugin()], }); return new KyselyCatalogRepository(db); }