Files
ThothII/backend/src/catalog/repository.ts
T

1416 lines
59 KiB
TypeScript

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 {
CatalogConflictError,
CatalogConnectorError,
CatalogUnavailableError,
DescriptionGenerationRunActiveError,
type CatalogColumn,
type CatalogDescriptionConsolidationCounts,
type CatalogDescriptionTarget,
type CatalogDatabaseMetadataDeleteTarget,
type CatalogMetadataDeleteCounts,
type CatalogRelationship,
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 TableSyncRepositoryResult,
type WorkspaceDatabase,
} from "./types.js";
type Timestamp = ColumnType<Date, Date | string | undefined, Date | string>;
interface WorkspaceDatabaseTable {
id: string;
workspaceId: string;
engine: "postgres";
databaseName: string;
schemaName: string;
version: Generated<number>;
createdAt: Timestamp;
updatedAt: Timestamp;
schemaSyncedVersion: number | null;
schemaSyncedAt: Timestamp | null;
}
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<number>;
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;
lastSyncedDatabaseVersion: number | null;
lastSyncedAt: Timestamp | null;
version: Generated<number>;
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 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;
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 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<string[], string, never>;
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<number>;
runId: string;
sequence: number;
level: CatalogSyncEvent["level"];
eventType: string;
message: string;
data: Record<string, unknown>;
createdAt: Timestamp;
}
export interface CatalogDatabase {
workspaceDatabases: WorkspaceDatabaseTable;
databaseBindings: DatabaseBindingTable;
catalogTables: CatalogTableTable;
catalogColumns: CatalogColumnTable;
catalogRelationships: CatalogRelationshipTable;
catalogRelationshipColumns: CatalogRelationshipColumnTable;
descriptionGenerationRuns: DescriptionGenerationRunTable;
descriptionGenerationEvents: DescriptionGenerationEventTable;
catalogSyncRuns: CatalogSyncRunTable;
catalogSyncEvents: CatalogSyncEventTable;
}
type DbOrTransaction = Kysely<CatalogDatabase> | Transaction<CatalogDatabase>;
type JoinedRow = Selectable<WorkspaceDatabaseTable> & Selectable<DatabaseBindingTable>;
function present<T>(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(),
};
}
function serializeTable(row: Selectable<CatalogTableTable>): 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<CatalogColumnTable>, 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,
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<CatalogSyncRunTable>): 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<CatalogSyncEventTable>): CatalogSyncEvent {
return { ...row, id: Number(row.id), data: row.data ?? {}, createdAt: new Date(row.createdAt).toISOString() };
}
function serializeDescriptionGenerationRun(
row: Selectable<DescriptionGenerationRunTable>,
): 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<DescriptionGenerationEventTable>,
): DescriptionGenerationEvent {
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<WorkspaceDatabase | undefined> {
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<CatalogDatabase>) {}
async list(): Promise<WorkspaceDatabase[]> {
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<WorkspaceDatabase | undefined> {
return await selectOne(this.db, id);
}
async getByWorkspace(workspaceId: string): Promise<WorkspaceDatabase | undefined> {
const id = await this.db.selectFrom("workspaceDatabases").select("id")
.where("workspaceId", "=", workspaceId).executeTakeFirst();
return id ? await this.get(id.id) : undefined;
}
async create(input: DatabaseConfigurationInput): Promise<WorkspaceDatabase> {
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<WorkspaceDatabase | undefined> {
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<WorkspaceDatabase | undefined> {
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<WorkspaceDatabase | undefined> {
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<boolean> {
const result = await this.db.deleteFrom("workspaceDatabases")
.where("id", "=", id).where("version", "=", expectedVersion).executeTakeFirst();
return result.numDeletedRows === 1n;
}
async listTables(databaseId: string): Promise<CatalogTable[]> {
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<CatalogTable | undefined> {
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<CatalogTable | undefined> {
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<CatalogTable | undefined> {
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<CatalogColumn[]> {
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<number>`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<CatalogColumn | undefined> {
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<number>`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,
): Promise<CatalogColumn | undefined> {
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,
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<CatalogDescriptionConsolidationCounts | undefined> {
const selectedTargetIds = [...new Set(targetIds)];
if (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")
.where("databaseId", "=", databaseId).execute();
const rows = tableRows.length === 0 ? [] : await trx.selectFrom("catalogColumns")
.select(["id", "generatedDescription"])
.where("tableId", "in", tableRows.map((table) => table.id))
.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("catalogColumns").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,
};
});
}
async createDescriptionGenerationRun(
databaseId: string,
scope: DescriptionGenerationScope,
modelId: string,
language: DescriptionGenerationRun["language"],
total: number,
): Promise<DescriptionGenerationRun> {
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,
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<DescriptionGenerationRun | undefined> {
const row = await this.db.selectFrom("descriptionGenerationRuns")
.selectAll()
.where("id", "=", runId)
.executeTakeFirst();
return row ? serializeDescriptionGenerationRun(row) : undefined;
}
async listDescriptionGenerationRuns(limit = 50): Promise<DescriptionGenerationRun[]> {
const rows = await this.db.selectFrom("descriptionGenerationRuns")
.selectAll()
.orderBy("createdAt", "desc")
.orderBy("id", "desc")
.limit(limit)
.execute();
return rows.map(serializeDescriptionGenerationRun);
}
async getActiveDescriptionGenerationRun(): Promise<DescriptionGenerationRun | undefined> {
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<DescriptionGenerationRun[]> {
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<DescriptionGenerationRun | undefined> {
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<DescriptionGenerationEvent> {
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<number>`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<DescriptionGenerationEvent[]> {
const rows = await this.db.selectFrom("descriptionGenerationEvents")
.selectAll()
.where("runId", "=", runId)
.where("sequence", ">", afterSequence)
.orderBy("sequence")
.execute();
return rows.map(serializeDescriptionGenerationEvent);
}
async listRelationships(databaseId: string): Promise<CatalogRelationship[]> {
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<string, CatalogRelationship["columns"]>();
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(),
}));
}
async deleteDatabaseMetadata(
databaseIds: readonly string[],
target: CatalogDatabaseMetadataDeleteTarget,
): Promise<CatalogMetadataDeleteCounts | undefined> {
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 count = await trx.selectFrom("catalogRelationships")
.select(sql<number>`count(*)::int`.as("count"))
.where("databaseId", "in", selectedDatabaseIds).executeTakeFirst();
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) };
}
const tableCount = await trx.selectFrom("catalogTables")
.select(sql<number>`count(*)::int`.as("count"))
.where("databaseId", "in", selectedDatabaseIds).executeTakeFirst();
const columnCount = await trx.selectFrom("catalogColumns")
.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")
.select(sql<number>`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(relationshipCount?.count ?? 0),
};
});
}
async deleteTableMetadata(
databaseId: string,
tableIds: readonly string[],
target: CatalogTableMetadataDeleteTarget,
): Promise<CatalogMetadataDeleteCounts | undefined> {
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<number>`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 count = 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();
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(count?.count ?? 0) };
});
}
async planSchemaSync(
databaseId: string,
scope: CatalogSyncScope,
tableIds: readonly string[],
snapshot: ObservedSchemaSnapshot,
): Promise<CatalogSchemaDiff> {
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,
): Promise<CatalogSyncCounts | undefined> {
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();
}
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<CatalogSyncRun> {
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<CatalogSyncRun | undefined> {
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<CatalogSyncRun | undefined> {
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<CatalogSyncRun[]> {
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<CatalogSyncRun | undefined> {
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<CatalogSyncRun | undefined> {
return await this.updateSyncRun(runId, { cancelRequested: true });
}
async appendSyncEvent(
runId: string,
level: CatalogSyncEvent["level"],
eventType: string,
message: string,
data: Record<string, unknown> = {},
): Promise<CatalogSyncEvent> {
return await this.db.transaction().execute(async (trx) => {
const current = await trx.selectFrom("catalogSyncEvents")
.select(sql<number>`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<CatalogSyncEvent[]> {
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<void> {
await this.db.deleteFrom("catalogSyncEvents").where("createdAt", "<", new Date(before)).execute();
}
async interruptActiveSyncRuns(): Promise<void> {
await this.db.updateTable("catalogSyncRuns").set({
state: "interrupted", phase: "completed", 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<TableSyncRepositoryResult | undefined> {
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<boolean> {
try {
await sql`select 1`.execute(this.db);
return true;
} catch {
return false;
}
}
async close(): Promise<void> { await this.db.destroy(); }
}
export class UnavailableCatalogRepository implements CatalogRepository {
private fail(): never { throw new CatalogUnavailableError("Database catalog is unavailable"); }
async list(): Promise<WorkspaceDatabase[]> { return this.fail(); }
async get(): Promise<WorkspaceDatabase | undefined> { return this.fail(); }
async getByWorkspace(): Promise<WorkspaceDatabase | 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(); }
async touch(): Promise<WorkspaceDatabase | undefined> { return this.fail(); }
async delete(): Promise<boolean> { return this.fail(); }
async listTables(): Promise<CatalogTable[]> { return this.fail(); }
async getTable(): Promise<CatalogTable | undefined> { return this.fail(); }
async updateTableDescription(): Promise<CatalogTable | undefined> { return this.fail(); }
async updateTableMetadata(): Promise<CatalogTable | undefined> { return this.fail(); }
async listColumns(): Promise<CatalogColumn[]> { return this.fail(); }
async getColumn(): Promise<CatalogColumn | undefined> { return this.fail(); }
async updateColumnMetadata(): Promise<CatalogColumn | undefined> { return this.fail(); }
async consolidateGeneratedDescriptions(): Promise<CatalogDescriptionConsolidationCounts | undefined> { return this.fail(); }
async createDescriptionGenerationRun(): Promise<DescriptionGenerationRun> { return this.fail(); }
async getDescriptionGenerationRun(): Promise<DescriptionGenerationRun | undefined> { return this.fail(); }
async listDescriptionGenerationRuns(): Promise<DescriptionGenerationRun[]> { return this.fail(); }
async getActiveDescriptionGenerationRun(): Promise<DescriptionGenerationRun | undefined> { return this.fail(); }
async interruptActiveDescriptionGenerationRuns(): Promise<DescriptionGenerationRun[]> { return this.fail(); }
async updateDescriptionGenerationRun(): Promise<DescriptionGenerationRun | undefined> { return this.fail(); }
async appendDescriptionGenerationEvent(): Promise<DescriptionGenerationEvent> { return this.fail(); }
async listDescriptionGenerationEvents(): Promise<DescriptionGenerationEvent[]> { 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(); }
async planSchemaSync(): Promise<CatalogSchemaDiff> { return this.fail(); }
async applySchemaSync(): Promise<CatalogSyncCounts | undefined> { return this.fail(); }
async createSyncRun(): Promise<CatalogSyncRun> { return this.fail(); }
async getSyncRun(): Promise<CatalogSyncRun | undefined> { return this.fail(); }
async claimSyncRun(): Promise<CatalogSyncRun | undefined> { return this.fail(); }
async listSyncRuns(): Promise<CatalogSyncRun[]> { return this.fail(); }
async updateSyncRun(): Promise<CatalogSyncRun | undefined> { return this.fail(); }
async requestSyncRunCancellation(): Promise<CatalogSyncRun | undefined> { return this.fail(); }
async appendSyncEvent(): Promise<CatalogSyncEvent> { return this.fail(); }
async listSyncEvents(): Promise<CatalogSyncEvent[]> { return this.fail(); }
async pruneSyncEvents(): Promise<void> { return this.fail(); }
async interruptActiveSyncRuns(): Promise<void> { return this.fail(); }
async reconcileTables(): Promise<TableSyncRepositoryResult | undefined> { return this.fail(); }
async available(): Promise<boolean> { 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<CatalogDatabase>({
dialect: new PostgresDialect({ pool }),
plugins: [new CamelCasePlugin()],
});
return new KyselyCatalogRepository(db);
}