2064 lines
85 KiB
TypeScript
2064 lines
85 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 { 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<Date, Date | string | undefined, Date | string>;
|
|
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<number>;
|
|
createdAt: Timestamp;
|
|
updatedAt: Timestamp;
|
|
schemaSyncedVersion: number | null;
|
|
schemaSyncedAt: Timestamp | null;
|
|
metadataContentRevision: Generated<number>;
|
|
preprocessingStatus: Generated<WorkspaceDatabase["preprocessingStatus"]>;
|
|
preprocessingInputFingerprint: Generated<string | null>;
|
|
preprocessedMetadataRevision: Generated<number | null>;
|
|
preprocessingStartedAt: NullableTimestamp;
|
|
preprocessingFinishedAt: NullableTimestamp;
|
|
preprocessingErrorCode: Generated<string | 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;
|
|
sensitive: Generated<boolean>;
|
|
sensitivityReason: Generated<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 CatalogLogicalRelationshipTable {
|
|
id: string;
|
|
databaseId: string;
|
|
sourceColumnId: string;
|
|
targetColumnId: string;
|
|
generated: Generated<boolean>;
|
|
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<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;
|
|
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<CatalogDatabase> | Transaction<CatalogDatabase>;
|
|
type JoinedRow = Selectable<WorkspaceDatabaseTable> & Selectable<DatabaseBindingTable>;
|
|
|
|
interface CatalogMetricsRow {
|
|
databaseCount: number;
|
|
tables: number;
|
|
columns: number;
|
|
sensitiveColumns: number;
|
|
relationships: number;
|
|
describedTables: number;
|
|
describedColumns: number;
|
|
updatedAt: Date | string | null;
|
|
}
|
|
|
|
function present<T>(value: T | null): T | undefined {
|
|
return value === null ? undefined : value;
|
|
}
|
|
|
|
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<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,
|
|
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<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 serializeSensitivityAnalysisRun(
|
|
row: Selectable<SensitivityAnalysisRunTable>,
|
|
): 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<SensitivityAnalysisEventTable>,
|
|
): 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<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 beginPreprocessing(
|
|
workspaceId: string,
|
|
inputFingerprint: string,
|
|
): Promise<CatalogPreprocessingStartResult> {
|
|
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<WorkspaceDatabase | undefined> {
|
|
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<CatalogPreprocessingClearResult> {
|
|
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<CatalogMetrics | undefined> {
|
|
const result = await sql<CatalogMetricsRow>`
|
|
WITH requested_database AS (
|
|
SELECT ${databaseId ?? null}::uuid AS id
|
|
),
|
|
selected_databases AS (
|
|
SELECT workspace_databases.id, workspace_databases.schema_synced_at
|
|
FROM workspace_databases
|
|
CROSS JOIN requested_database
|
|
WHERE requested_database.id IS NULL
|
|
OR workspace_databases.id = requested_database.id
|
|
),
|
|
table_metrics AS (
|
|
SELECT
|
|
count(*)::int AS tables,
|
|
count(*) FILTER (
|
|
WHERE nullif(btrim(catalog_tables.description), '') IS NOT NULL
|
|
OR nullif(btrim(catalog_tables.generated_description), '') IS NOT NULL
|
|
)::int AS "describedTables",
|
|
max(catalog_tables.updated_at) AS updated_at
|
|
FROM catalog_tables
|
|
INNER JOIN selected_databases
|
|
ON selected_databases.id = catalog_tables.database_id
|
|
),
|
|
column_metrics AS (
|
|
SELECT
|
|
count(*)::int AS columns,
|
|
count(*) FILTER (WHERE catalog_columns.sensitive)::int AS "sensitiveColumns",
|
|
count(*) FILTER (
|
|
WHERE nullif(btrim(catalog_columns.description), '') IS NOT NULL
|
|
OR nullif(btrim(catalog_columns.generated_description), '') IS NOT NULL
|
|
)::int AS "describedColumns",
|
|
max(catalog_columns.updated_at) AS updated_at
|
|
FROM catalog_columns
|
|
INNER JOIN catalog_tables ON catalog_tables.id = catalog_columns.table_id
|
|
INNER JOIN selected_databases
|
|
ON selected_databases.id = catalog_tables.database_id
|
|
),
|
|
relationship_metrics AS (
|
|
SELECT
|
|
count(*)::int AS relationships,
|
|
max(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<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,
|
|
sensitive?: boolean,
|
|
sensitivityReason?: 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,
|
|
...(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<CatalogDescriptionConsolidationCounts | undefined> {
|
|
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<boolean>`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<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,
|
|
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<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 createSensitivityAnalysisRun(
|
|
databaseId: string,
|
|
scope: SensitivityAnalysisScope,
|
|
origin: { engine: "llm"; modelId: string } | { engine: "local"; policyVersion: string },
|
|
): Promise<SensitivityAnalysisRun> {
|
|
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<SensitivityAnalysisRun | undefined> {
|
|
const row = await this.db.selectFrom("sensitiveDataSuggestionRuns")
|
|
.selectAll()
|
|
.where("id", "=", runId)
|
|
.executeTakeFirst();
|
|
return row ? serializeSensitivityAnalysisRun(row) : undefined;
|
|
}
|
|
|
|
async listSensitivityAnalysisRuns(limit = 50): Promise<SensitivityAnalysisRun[]> {
|
|
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<SensitivityAnalysisRun[]> {
|
|
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<SensitivityAnalysisRun | undefined> {
|
|
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<SensitivityAnalysisEvent> {
|
|
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<number>`coalesce(max(sequence), 0)::int`.as("sequence"))
|
|
.where("runId", "=", runId)
|
|
.executeTakeFirst();
|
|
const row = await trx.insertInto("sensitiveDataSuggestionEvents").values({
|
|
runId,
|
|
sequence: Number(current?.sequence ?? 0) + 1,
|
|
level,
|
|
message,
|
|
}).returningAll().executeTakeFirstOrThrow();
|
|
return serializeSensitivityAnalysisEvent(row);
|
|
});
|
|
}
|
|
|
|
async listSensitivityAnalysisEvents(
|
|
runId: string,
|
|
afterSequence = 0,
|
|
): Promise<SensitivityAnalysisEvent[]> {
|
|
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<CatalogPhysicalRelationship[]> {
|
|
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, CatalogPhysicalRelationship["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(),
|
|
origin: "physical" as const,
|
|
status: "active" as const,
|
|
}));
|
|
}
|
|
|
|
async listLogicalRelationships(databaseId: string): Promise<CatalogLogicalRelationship[]> {
|
|
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<CatalogLogicalRelationshipContext | undefined> {
|
|
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<string, number>();
|
|
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<CatalogLogicalRelationship | undefined> {
|
|
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<number> {
|
|
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<CatalogLogicalRelationship | undefined> {
|
|
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<boolean> {
|
|
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<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 physicalCount = await trx.selectFrom("catalogRelationships")
|
|
.select(sql<number>`count(*)::int`.as("count"))
|
|
.where("databaseId", "in", selectedDatabaseIds).executeTakeFirst();
|
|
const logicalCount = await trx.selectFrom("catalogLogicalRelationships")
|
|
.select(sql<number>`count(*)::int`.as("count"))
|
|
.where("databaseId", "in", selectedDatabaseIds).executeTakeFirst();
|
|
await trx.deleteFrom("catalogLogicalRelationships")
|
|
.where("databaseId", "in", selectedDatabaseIds).execute();
|
|
await trx.deleteFrom("catalogRelationships")
|
|
.where("databaseId", "in", selectedDatabaseIds).execute();
|
|
await trx.updateTable("workspaceDatabases").set({
|
|
schemaSyncedVersion: null,
|
|
schemaSyncedAt: null,
|
|
}).where("id", "in", selectedDatabaseIds).execute();
|
|
return {
|
|
tables: 0,
|
|
columns: 0,
|
|
relationships: Number(physicalCount?.count ?? 0) + Number(logicalCount?.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 physicalRelationshipCount = await trx.selectFrom("catalogRelationships")
|
|
.select(sql<number>`count(*)::int`.as("count"))
|
|
.where("databaseId", "in", selectedDatabaseIds).executeTakeFirst();
|
|
const logicalRelationshipCount = await trx.selectFrom("catalogLogicalRelationships")
|
|
.select(sql<number>`count(*)::int`.as("count"))
|
|
.where("databaseId", "in", selectedDatabaseIds).executeTakeFirst();
|
|
await trx.deleteFrom("catalogTables")
|
|
.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<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 physicalCount = await trx.selectFrom("catalogRelationships")
|
|
.select(sql<number>`count(*)::int`.as("count"))
|
|
.where("databaseId", "=", databaseId)
|
|
.where((eb) => eb.or([
|
|
eb("sourceTableId", "in", selectedTableIds),
|
|
eb("targetTableId", "in", selectedTableIds),
|
|
])).executeTakeFirst();
|
|
const selectedColumnIds = (await trx.selectFrom("catalogColumns")
|
|
.select("id")
|
|
.where("tableId", "in", selectedTableIds)
|
|
.execute()).map((column) => column.id);
|
|
const logicalCount = selectedColumnIds.length === 0
|
|
? undefined
|
|
: await trx.selectFrom("catalogLogicalRelationships")
|
|
.select(sql<number>`count(*)::int`.as("count"))
|
|
.where("databaseId", "=", databaseId)
|
|
.where((eb) => eb.or([
|
|
eb("sourceColumnId", "in", selectedColumnIds),
|
|
eb("targetColumnId", "in", selectedColumnIds),
|
|
])).executeTakeFirst();
|
|
if (selectedColumnIds.length > 0) {
|
|
await trx.deleteFrom("catalogLogicalRelationships")
|
|
.where("databaseId", "=", databaseId)
|
|
.where((eb) => eb.or([
|
|
eb("sourceColumnId", "in", selectedColumnIds),
|
|
eb("targetColumnId", "in", selectedColumnIds),
|
|
])).execute();
|
|
}
|
|
await trx.deleteFrom("catalogRelationships")
|
|
.where("databaseId", "=", databaseId)
|
|
.where((eb) => eb.or([
|
|
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<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 beginPreprocessing(): Promise<CatalogPreprocessingStartResult> { return this.fail(); }
|
|
async finishPreprocessing(): Promise<WorkspaceDatabase | undefined> { return this.fail(); }
|
|
async clearPreprocessing(): Promise<CatalogPreprocessingClearResult> { return this.fail(); }
|
|
async getCatalogMetrics(): Promise<CatalogMetrics | undefined> { return this.fail(); }
|
|
async create(): Promise<WorkspaceDatabase> { return this.fail(); }
|
|
async update(): Promise<WorkspaceDatabase | undefined> { return this.fail(); }
|
|
async recordTest(): Promise<WorkspaceDatabase | undefined> { return this.fail(); }
|
|
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 createSensitivityAnalysisRun(): Promise<SensitivityAnalysisRun> { return this.fail(); }
|
|
async getSensitivityAnalysisRun(): Promise<SensitivityAnalysisRun | undefined> { return this.fail(); }
|
|
async listSensitivityAnalysisRuns(): Promise<SensitivityAnalysisRun[]> { return this.fail(); }
|
|
async interruptActiveSensitivityAnalysisRuns(): Promise<SensitivityAnalysisRun[]> { return this.fail(); }
|
|
async updateSensitivityAnalysisRun(): Promise<SensitivityAnalysisRun | undefined> { return this.fail(); }
|
|
async appendSensitivityAnalysisEvent(): Promise<SensitivityAnalysisEvent> { return this.fail(); }
|
|
async listSensitivityAnalysisEvents(): Promise<SensitivityAnalysisEvent[]> { return this.fail(); }
|
|
async listRelationships(): Promise<CatalogPhysicalRelationship[]> { return this.fail(); }
|
|
async listLogicalRelationships(): Promise<CatalogLogicalRelationship[]> { return this.fail(); }
|
|
async getLogicalRelationshipContext(): Promise<CatalogLogicalRelationshipContext | undefined> { return this.fail(); }
|
|
async insertLogicalRelationship(): Promise<CatalogLogicalRelationship | undefined> { return this.fail(); }
|
|
async insertGeneratedLogicalRelationships(): Promise<number> { return this.fail(); }
|
|
async setLogicalRelationshipStatus(): Promise<CatalogLogicalRelationship | undefined> { return this.fail(); }
|
|
async deleteLogicalRelationship(): Promise<boolean> { 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);
|
|
}
|