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

1297 lines
53 KiB
TypeScript

import { randomUUID } from "node:crypto";
import {
createCatalogMetrics,
hasCatalogDescription,
latestCatalogTimestamp,
} from "./metrics.js";
import {
CatalogConflictError,
CatalogConnectorError,
DescriptionGenerationRunActiveError,
type CatalogColumn,
type CatalogDescriptionConsolidationCounts,
type CatalogDescriptionTarget,
type CatalogDatabaseMetadataDeleteTarget,
type CatalogMetadataDeleteCounts,
type CatalogMetrics,
type CatalogLogicalRelationship,
type CatalogLogicalRelationshipCandidate,
type CatalogLogicalRelationshipContext,
type CatalogPhysicalRelationship,
type CatalogSchemaDiff,
type CatalogSyncCounts,
type CatalogSyncEvent,
type CatalogSyncRun,
type CatalogSyncRunUpdate,
type CatalogSyncScope,
type CatalogTable,
type CatalogTableMetadataDeleteTarget,
type CatalogRepository,
type DatabaseConfigurationInput,
type DatabaseTestResult,
type ObservedCatalogTable,
type ObservedSchemaSnapshot,
type DescriptionGenerationEvent,
type DescriptionGenerationRun,
type DescriptionGenerationRunUpdate,
type DescriptionGenerationScope,
type SensitivityAnalysisEvent,
type SensitivityAnalysisRun,
type SensitivityAnalysisRunUpdate,
type SensitivityAnalysisScope,
type TableSyncRepositoryResult,
type WorkspaceDatabase,
} from "./types.js";
function clone(value: WorkspaceDatabase): WorkspaceDatabase {
return structuredClone(value);
}
/** Deterministic repository used by route tests and local contract consumers. */
export class MemoryCatalogRepository implements CatalogRepository {
private readonly records = new Map<string, WorkspaceDatabase>();
private readonly tables = new Map<string, CatalogTable>();
private readonly columns = new Map<string, CatalogColumn>();
private readonly relationships = new Map<string, CatalogPhysicalRelationship>();
private readonly logicalRelationships = new Map<string, CatalogLogicalRelationship>();
private readonly descriptionGenerationRuns = new Map<string, DescriptionGenerationRun>();
private readonly descriptionGenerationEvents = new Map<string, DescriptionGenerationEvent[]>();
private readonly sensitivityAnalysisRuns = new Map<string, SensitivityAnalysisRun>();
private readonly sensitivityAnalysisEvents = new Map<string, SensitivityAnalysisEvent[]>();
private readonly syncRuns = new Map<string, CatalogSyncRun>();
private readonly syncEvents = new Map<string, CatalogSyncEvent[]>();
async list(): Promise<WorkspaceDatabase[]> {
return [...this.records.values()].sort((a, b) => a.workspaceId.localeCompare(b.workspaceId)).map(clone);
}
async get(id: string): Promise<WorkspaceDatabase | undefined> {
const value = this.records.get(id);
return value ? clone(value) : undefined;
}
async getByWorkspace(workspaceId: string): Promise<WorkspaceDatabase | undefined> {
const value = [...this.records.values()].find((record) => record.workspaceId === workspaceId);
return value ? clone(value) : undefined;
}
async getCatalogMetrics(databaseId?: string): Promise<CatalogMetrics | undefined> {
if (databaseId !== undefined && !this.records.has(databaseId)) return undefined;
const selectedDatabaseIds = new Set(
databaseId === undefined ? this.records.keys() : [databaseId],
);
const tables = [...this.tables.values()]
.filter((table) => selectedDatabaseIds.has(table.databaseId));
const tableIds = new Set(tables.map((table) => table.id));
const columns = [...this.columns.values()]
.filter((column) => tableIds.has(column.tableId));
const relationships = [
...this.relationships.values(),
...this.logicalRelationships.values(),
].filter((relationship) => selectedDatabaseIds.has(relationship.databaseId));
return createCatalogMetrics(databaseId, {
tables: tables.length,
columns: columns.length,
sensitiveColumns: columns.filter((column) => column.sensitive).length,
relationships: relationships.length,
describedTables: tables.filter(hasCatalogDescription).length,
describedColumns: columns.filter(hasCatalogDescription).length,
}, latestCatalogTimestamp([
...[...selectedDatabaseIds].map((id) => this.records.get(id)?.schemaSyncedAt),
...tables.map((table) => table.updatedAt),
...columns.map((column) => column.updatedAt),
...relationships.map((relationship) => relationship.updatedAt),
]));
}
async create(input: DatabaseConfigurationInput): Promise<WorkspaceDatabase> {
if ([...this.records.values()].some((record) => record.workspaceId === input.workspaceId)) {
throw new CatalogConflictError("Workspace database already exists");
}
const now = new Date().toISOString();
const record: WorkspaceDatabase = {
id: randomUUID(),
...structuredClone(input),
version: 1,
createdAt: now,
updatedAt: now,
connectionStatus: "untested",
};
this.records.set(record.id, record);
return clone(record);
}
async update(id: string, expectedVersion: number, input: DatabaseConfigurationInput): Promise<WorkspaceDatabase | undefined> {
const current = this.records.get(id);
if (!current || current.version !== expectedVersion) return undefined;
if ([...this.records.values()].some((record) => record.id !== id && record.workspaceId === input.workspaceId)) {
throw new CatalogConflictError("Workspace database already exists");
}
const updated: WorkspaceDatabase = {
...current,
...structuredClone(input),
version: current.version + 1,
updatedAt: new Date().toISOString(),
connectionStatus: "untested",
testedVersion: undefined,
lastTestedAt: undefined,
lastErrorCode: undefined,
lastErrorMessage: undefined,
schemaSyncedVersion: undefined,
schemaSyncedAt: undefined,
};
this.records.set(id, updated);
return clone(updated);
}
async recordTest(id: string, expectedVersion: number, result: DatabaseTestResult): Promise<WorkspaceDatabase | undefined> {
const current = this.records.get(id);
if (!current || current.version !== expectedVersion) return undefined;
const updated = {
...current,
connectionStatus: result.connectionStatus,
testedVersion: result.testedVersion,
lastTestedAt: result.lastTestedAt,
lastErrorCode: result.errorCode,
lastErrorMessage: result.errorMessage,
};
this.records.set(id, updated);
return clone(updated);
}
async touch(id: string, expectedVersion: number): Promise<WorkspaceDatabase | undefined> {
const current = this.records.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 current = this.records.get(id);
if (!current || current.version !== expectedVersion) return false;
for (const [tableId, table] of this.tables) {
if (table.databaseId === id) {
this.tables.delete(tableId);
for (const [columnId, column] of this.columns) if (column.tableId === tableId) this.columns.delete(columnId);
}
}
for (const [relationshipId, relationship] of this.relationships) {
if (relationship.databaseId === id) this.relationships.delete(relationshipId);
}
for (const [relationshipId, relationship] of this.logicalRelationships) {
if (relationship.databaseId === id) this.logicalRelationships.delete(relationshipId);
}
for (const [runId, run] of this.descriptionGenerationRuns) {
if (run.databaseId !== id) continue;
this.descriptionGenerationRuns.delete(runId);
this.descriptionGenerationEvents.delete(runId);
}
for (const [runId, run] of this.sensitivityAnalysisRuns) {
if (run.databaseId !== id) continue;
this.sensitivityAnalysisRuns.delete(runId);
this.sensitivityAnalysisEvents.delete(runId);
}
return this.records.delete(id);
}
async listTables(databaseId: string): Promise<CatalogTable[]> {
return [...this.tables.values()]
.filter((table) => table.databaseId === databaseId)
.sort((a, b) => a.name.localeCompare(b.name))
.map((table) => structuredClone(table));
}
async getTable(databaseId: string, tableId: string): Promise<CatalogTable | undefined> {
const table = this.tables.get(tableId);
return table?.databaseId === databaseId ? structuredClone(table) : undefined;
}
async updateTableDescription(
databaseId: string,
tableId: string,
expectedVersion: number,
description: string | null,
): Promise<CatalogTable | undefined> {
const current = this.tables.get(tableId);
if (!current || current.databaseId !== databaseId || current.version !== expectedVersion) return undefined;
const updated = {
...current,
description,
version: current.version + 1,
updatedAt: new Date().toISOString(),
};
this.tables.set(tableId, updated);
return structuredClone(updated);
}
async updateTableMetadata(
databaseId: string,
tableId: string,
expectedVersion: number,
description: string | null,
generatedDescription: string | null,
): Promise<CatalogTable | undefined> {
const current = this.tables.get(tableId);
if (!current || current.databaseId !== databaseId || current.version !== expectedVersion) return undefined;
const updated = {
...current, description, generatedDescription, version: current.version + 1,
updatedAt: new Date().toISOString(),
};
this.tables.set(tableId, updated);
return structuredClone(updated);
}
async listColumns(databaseId: string, tableId: string): Promise<CatalogColumn[]> {
const table = this.tables.get(tableId);
if (!table || table.databaseId !== databaseId) return [];
return [...this.columns.values()].filter((column) => column.tableId === tableId)
.sort((a, b) => a.ordinalPosition - b.ordinalPosition).map((column) => structuredClone(column));
}
async getColumn(databaseId: string, tableId: string, columnId: string): Promise<CatalogColumn | undefined> {
const table = this.tables.get(tableId);
const column = this.columns.get(columnId);
return table?.databaseId === databaseId && column?.tableId === tableId ? structuredClone(column) : undefined;
}
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 current = await this.getColumn(databaseId, tableId, columnId);
if (!current || current.version !== expectedVersion) return undefined;
const updated = {
...current,
description,
generatedDescription,
sensitive: sensitive ?? current.sensitive,
sensitivityReason: sensitive === false
? null
: sensitivityReason === undefined ? current.sensitivityReason : sensitivityReason,
version: current.version + 1,
updatedAt: new Date().toISOString(),
};
this.columns.set(columnId, updated);
return structuredClone(updated);
}
async consolidateGeneratedDescriptions(
databaseId: string,
target: CatalogDescriptionTarget,
targetIds: readonly string[],
): Promise<CatalogDescriptionConsolidationCounts | undefined> {
const selectedTargetIds = [...new Set(targetIds)];
if (!this.records.has(databaseId)) {
return undefined;
}
const now = new Date().toISOString();
if (target === "database_columns") {
const targets = [...this.columns.values()].filter((column) => (
this.tables.get(column.tableId)?.databaseId === databaseId
));
const copied = targets.filter((column) => Boolean(column.generatedDescription?.trim()));
for (const column of copied) {
this.columns.set(column.id, {
...column,
description: column.generatedDescription,
version: column.version + 1,
updatedAt: now,
});
}
return { copied: copied.length, skipped: targets.length - copied.length };
}
if (selectedTargetIds.length === 0) return undefined;
if (target === "tables") {
const targets = selectedTargetIds.map((id) => this.tables.get(id));
if (targets.some((table) => !table || table.databaseId !== databaseId)) return undefined;
const copied = targets.filter((table) => Boolean(table!.generatedDescription?.trim()));
for (const table of copied) {
this.tables.set(table!.id, {
...table!,
description: table!.generatedDescription,
version: table!.version + 1,
updatedAt: now,
});
}
return { copied: copied.length, skipped: targets.length - copied.length };
}
const targets = selectedTargetIds.map((id) => this.columns.get(id));
if (targets.some((column) => (
!column || this.tables.get(column.tableId)?.databaseId !== databaseId
))) return undefined;
const copied = targets.filter((column) => Boolean(column!.generatedDescription?.trim()));
for (const column of copied) {
this.columns.set(column!.id, {
...column!,
description: column!.generatedDescription,
version: column!.version + 1,
updatedAt: now,
});
}
return { copied: copied.length, skipped: targets.length - copied.length };
}
async createDescriptionGenerationRun(
databaseId: string,
scope: DescriptionGenerationScope,
modelId: string,
language: DescriptionGenerationRun["language"],
total: number,
): Promise<DescriptionGenerationRun> {
if ([...this.descriptionGenerationRuns.values()].some((run) => (
run.status === "queued" || run.status === "running"
))) {
throw new DescriptionGenerationRunActiveError("A description generation run is already active");
}
const now = new Date().toISOString();
const run: DescriptionGenerationRun = {
id: randomUUID(),
databaseId,
scope,
modelId,
language,
status: "queued",
total,
processed: 0,
generated: 0,
nonGeneratable: 0,
failed: 0,
inputTokens: 0,
cacheReadTokens: 0,
outputTokens: 0,
createdAt: now,
startedAt: null,
updatedAt: now,
finishedAt: null,
errorSummary: null,
};
this.descriptionGenerationRuns.set(run.id, run);
return structuredClone(run);
}
async getDescriptionGenerationRun(runId: string): Promise<DescriptionGenerationRun | undefined> {
const run = this.descriptionGenerationRuns.get(runId);
return run ? structuredClone(run) : undefined;
}
async listDescriptionGenerationRuns(limit = 50): Promise<DescriptionGenerationRun[]> {
return [...this.descriptionGenerationRuns.values()]
.sort((a, b) => b.createdAt.localeCompare(a.createdAt) || b.id.localeCompare(a.id))
.slice(0, limit)
.map((run) => structuredClone(run));
}
async getActiveDescriptionGenerationRun(): Promise<DescriptionGenerationRun | undefined> {
const run = [...this.descriptionGenerationRuns.values()]
.filter((candidate) => candidate.status === "queued" || candidate.status === "running")
.sort((a, b) => b.createdAt.localeCompare(a.createdAt))[0];
return run ? structuredClone(run) : undefined;
}
async interruptActiveDescriptionGenerationRuns(
errorSummary: string,
): Promise<DescriptionGenerationRun[]> {
const interrupted: DescriptionGenerationRun[] = [];
for (const run of this.descriptionGenerationRuns.values()) {
if (run.status !== "queued" && run.status !== "running") continue;
const now = new Date().toISOString();
const updated: DescriptionGenerationRun = {
...run,
status: "interrupted",
updatedAt: now,
finishedAt: now,
errorSummary,
};
this.descriptionGenerationRuns.set(run.id, updated);
interrupted.push(structuredClone(updated));
}
return interrupted;
}
async updateDescriptionGenerationRun(
runId: string,
update: DescriptionGenerationRunUpdate,
): Promise<DescriptionGenerationRun | undefined> {
const current = this.descriptionGenerationRuns.get(runId);
if (!current) return undefined;
const updated = {
...current,
...structuredClone(update),
updatedAt: new Date().toISOString(),
};
this.descriptionGenerationRuns.set(runId, updated);
return structuredClone(updated);
}
async appendDescriptionGenerationEvent(
runId: string,
level: DescriptionGenerationEvent["level"],
message: string,
): Promise<DescriptionGenerationEvent> {
if (!this.descriptionGenerationRuns.has(runId)) {
throw new CatalogConflictError("Description Generation Run does not exist");
}
const events = this.descriptionGenerationEvents.get(runId) ?? [];
const event: DescriptionGenerationEvent = {
runId,
sequence: events.length + 1,
level,
message,
createdAt: new Date().toISOString(),
};
events.push(event);
this.descriptionGenerationEvents.set(runId, events);
return structuredClone(event);
}
async listDescriptionGenerationEvents(
runId: string,
afterSequence = 0,
): Promise<DescriptionGenerationEvent[]> {
return (this.descriptionGenerationEvents.get(runId) ?? [])
.filter((event) => event.sequence > afterSequence)
.map((event) => structuredClone(event));
}
async createSensitivityAnalysisRun(
databaseId: string,
scope: SensitivityAnalysisScope,
origin: { engine: "llm"; modelId: string } | { engine: "local"; policyVersion: string },
): Promise<SensitivityAnalysisRun> {
const now = new Date().toISOString();
const run: SensitivityAnalysisRun = {
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,
createdAt: now,
startedAt: now,
updatedAt: now,
finishedAt: null,
errorSummary: null,
};
this.sensitivityAnalysisRuns.set(run.id, run);
return structuredClone(run);
}
async getSensitivityAnalysisRun(
runId: string,
): Promise<SensitivityAnalysisRun | undefined> {
const run = this.sensitivityAnalysisRuns.get(runId);
return run ? structuredClone(run) : undefined;
}
async listSensitivityAnalysisRuns(limit = 50): Promise<SensitivityAnalysisRun[]> {
return [...this.sensitivityAnalysisRuns.values()]
.sort((a, b) => b.createdAt.localeCompare(a.createdAt) || b.id.localeCompare(a.id))
.slice(0, limit)
.map((run) => structuredClone(run));
}
async interruptActiveSensitivityAnalysisRuns(
errorSummary: string,
): Promise<SensitivityAnalysisRun[]> {
const interrupted: SensitivityAnalysisRun[] = [];
for (const run of this.sensitivityAnalysisRuns.values()) {
if (run.status !== "running") continue;
const now = new Date().toISOString();
const updated: SensitivityAnalysisRun = {
...run,
status: "interrupted",
updatedAt: now,
finishedAt: now,
errorSummary,
};
this.sensitivityAnalysisRuns.set(run.id, updated);
interrupted.push(structuredClone(updated));
}
return interrupted;
}
async updateSensitivityAnalysisRun(
runId: string,
update: SensitivityAnalysisRunUpdate,
): Promise<SensitivityAnalysisRun | undefined> {
const current = this.sensitivityAnalysisRuns.get(runId);
if (!current) return undefined;
const updated = {
...current,
...structuredClone(update),
updatedAt: new Date().toISOString(),
};
this.sensitivityAnalysisRuns.set(runId, updated);
return structuredClone(updated);
}
async appendSensitivityAnalysisEvent(
runId: string,
level: SensitivityAnalysisEvent["level"],
message: string,
): Promise<SensitivityAnalysisEvent> {
if (!this.sensitivityAnalysisRuns.has(runId)) {
throw new CatalogConflictError("Sensitivity Analysis Run does not exist");
}
const events = this.sensitivityAnalysisEvents.get(runId) ?? [];
const event: SensitivityAnalysisEvent = {
runId,
sequence: events.length + 1,
level,
message,
createdAt: new Date().toISOString(),
};
events.push(event);
this.sensitivityAnalysisEvents.set(runId, events);
return structuredClone(event);
}
async listSensitivityAnalysisEvents(
runId: string,
afterSequence = 0,
): Promise<SensitivityAnalysisEvent[]> {
return (this.sensitivityAnalysisEvents.get(runId) ?? [])
.filter((event) => event.sequence > afterSequence)
.map((event) => structuredClone(event));
}
async listRelationships(databaseId: string): Promise<CatalogPhysicalRelationship[]> {
return [...this.relationships.values()].filter((relationship) => relationship.databaseId === databaseId)
.sort((a, b) => `${a.sourceTableName}.${a.constraintName}`.localeCompare(`${b.sourceTableName}.${b.constraintName}`))
.map((relationship) => structuredClone(relationship));
}
async listLogicalRelationships(databaseId: string): Promise<CatalogLogicalRelationship[]> {
return [...this.logicalRelationships.values()]
.filter((relationship) => relationship.databaseId === databaseId)
.sort((a, b) => {
const left = `${a.sourceTableName}.${a.columns[0].sourceColumnName}.${a.targetTableName}.${a.columns[0].targetColumnName}`;
const right = `${b.sourceTableName}.${b.columns[0].sourceColumnName}.${b.targetTableName}.${b.columns[0].targetColumnName}`;
return left.localeCompare(right);
})
.map((relationship) => structuredClone(relationship));
}
async getLogicalRelationshipContext(
databaseId: string,
): Promise<CatalogLogicalRelationshipContext | undefined> {
if (!this.records.has(databaseId)) return undefined;
const tables = [...this.tables.values()].filter((table) => table.databaseId === databaseId);
const tableById = new Map(tables.map((table) => [table.id, table]));
const columns = [...this.columns.values()].filter((column) => tableById.has(column.tableId));
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);
}
}
return {
endpoints: columns.map((column) => ({
columnId: column.id,
columnName: column.name,
tableId: column.tableId,
tableName: tableById.get(column.tableId)!.name,
dataType: column.dataType,
primaryKeyPosition: column.primaryKeyPosition,
tablePrimaryKeyColumnCount: primaryKeyCounts.get(column.tableId) ?? 0,
})),
physicalPairs: [...this.relationships.values()]
.filter((relationship) => relationship.databaseId === databaseId)
.flatMap((relationship) => relationship.columns.map((column) => ({
sourceColumnId: column.sourceColumnId,
targetColumnId: column.targetColumnId,
}))),
logicalRelationships: await this.listLogicalRelationships(databaseId),
};
}
async insertLogicalRelationship(
databaseId: string,
sourceColumnId: string,
targetColumnId: string,
generated: boolean,
): Promise<CatalogLogicalRelationship | undefined> {
if ([...this.logicalRelationships.values()].some((relationship) => (
relationship.databaseId === databaseId
&& relationship.columns[0].sourceColumnId === sourceColumnId
&& relationship.columns[0].targetColumnId === targetColumnId
))) return undefined;
const sourceColumn = this.columns.get(sourceColumnId);
const targetColumn = this.columns.get(targetColumnId);
const sourceTable = sourceColumn ? this.tables.get(sourceColumn.tableId) : undefined;
const targetTable = targetColumn ? this.tables.get(targetColumn.tableId) : undefined;
if (!sourceColumn || !targetColumn || !sourceTable || !targetTable
|| sourceTable.databaseId !== databaseId || targetTable.databaseId !== databaseId
|| sourceColumnId === targetColumnId) return undefined;
const now = new Date().toISOString();
const relationship: CatalogLogicalRelationship = {
id: randomUUID(),
databaseId,
constraintName: null,
sourceTableId: sourceTable.id,
sourceTableName: sourceTable.name,
targetTableId: targetTable.id,
targetTableName: targetTable.name,
updateRule: null,
deleteRule: null,
deferrable: false,
initiallyDeferred: false,
columns: [{
position: 1,
sourceColumnId,
sourceColumnName: sourceColumn.name,
targetColumnId,
targetColumnName: targetColumn.name,
}],
lastSyncedDatabaseVersion: null,
lastSyncedAt: null,
createdAt: now,
updatedAt: now,
origin: generated ? "generated" : "manual",
status: "active",
};
this.logicalRelationships.set(relationship.id, relationship);
return structuredClone(relationship);
}
async insertGeneratedLogicalRelationships(
databaseId: string,
candidates: readonly CatalogLogicalRelationshipCandidate[],
): Promise<number> {
let added = 0;
for (const candidate of candidates) {
if (await this.insertLogicalRelationship(
databaseId,
candidate.sourceColumnId,
candidate.targetColumnId,
true,
)) added += 1;
}
return added;
}
async setLogicalRelationshipStatus(
databaseId: string,
relationshipId: string,
status: CatalogLogicalRelationship["status"],
): Promise<CatalogLogicalRelationship | undefined> {
const current = this.logicalRelationships.get(relationshipId);
if (!current || current.databaseId !== databaseId) return undefined;
const updated = { ...current, status, updatedAt: new Date().toISOString() };
this.logicalRelationships.set(relationshipId, updated);
return structuredClone(updated);
}
async deleteLogicalRelationship(databaseId: string, relationshipId: string): Promise<boolean> {
const current = this.logicalRelationships.get(relationshipId);
return Boolean(current?.databaseId === databaseId && this.logicalRelationships.delete(relationshipId));
}
async deleteDatabaseMetadata(
databaseIds: readonly string[],
target: CatalogDatabaseMetadataDeleteTarget,
): Promise<CatalogMetadataDeleteCounts | undefined> {
const selectedDatabaseIds = [...new Set(databaseIds)];
if (selectedDatabaseIds.length === 0 || selectedDatabaseIds.some((id) => !this.records.has(id))) {
return undefined;
}
const selected = new Set(selectedDatabaseIds);
const tables = [...this.tables.values()].filter((table) => selected.has(table.databaseId));
const tableIds = new Set(tables.map((table) => table.id));
const columns = [...this.columns.values()].filter((column) => tableIds.has(column.tableId));
const physicalRelationships = [...this.relationships.values()]
.filter((relationship) => selected.has(relationship.databaseId));
const logicalRelationships = [...this.logicalRelationships.values()]
.filter((relationship) => selected.has(relationship.databaseId));
const relationshipCount = physicalRelationships.length + logicalRelationships.length;
if (target === "tables") {
for (const table of tables) this.deleteTable(table.id);
this.markCatalogIncomplete(selectedDatabaseIds);
return { tables: tables.length, columns: columns.length, relationships: relationshipCount };
}
for (const relationship of physicalRelationships) this.relationships.delete(relationship.id);
for (const relationship of logicalRelationships) this.logicalRelationships.delete(relationship.id);
for (const databaseId of selectedDatabaseIds) this.refreshForeignKeyFlags(databaseId);
this.markCatalogIncomplete(selectedDatabaseIds);
return { tables: 0, columns: 0, relationships: relationshipCount };
}
async deleteTableMetadata(
databaseId: string,
tableIds: readonly string[],
target: CatalogTableMetadataDeleteTarget,
): Promise<CatalogMetadataDeleteCounts | undefined> {
const selectedTableIds = [...new Set(tableIds)];
if (!this.records.has(databaseId) || selectedTableIds.length === 0) return undefined;
const tables = selectedTableIds.map((id) => this.tables.get(id));
if (tables.some((table) => !table || table.databaseId !== databaseId)) return undefined;
const selected = new Set(selectedTableIds);
if (target === "columns") {
const columns = [...this.columns.values()].filter((column) => selected.has(column.tableId));
const deletedColumnIds = new Set(columns.map((column) => column.id));
for (const column of columns) this.columns.delete(column.id);
for (const relationship of [...this.logicalRelationships.values()]) {
const pair = relationship.columns[0];
if (deletedColumnIds.has(pair.sourceColumnId) || deletedColumnIds.has(pair.targetColumnId)) {
this.logicalRelationships.delete(relationship.id);
}
}
for (const [relationshipId, relationship] of this.relationships) {
if (relationship.databaseId !== databaseId) continue;
this.relationships.set(relationshipId, {
...relationship,
columns: relationship.columns.filter((pair) => (
!deletedColumnIds.has(pair.sourceColumnId) && !deletedColumnIds.has(pair.targetColumnId)
)),
});
}
this.refreshForeignKeyFlags(databaseId);
this.markCatalogIncomplete([databaseId]);
return { tables: 0, columns: columns.length, relationships: 0 };
}
const physicalRelationships = [...this.relationships.values()].filter((relationship) => (
relationship.databaseId === databaseId
&& (selected.has(relationship.sourceTableId) || selected.has(relationship.targetTableId))
));
const logicalRelationships = [...this.logicalRelationships.values()].filter((relationship) => (
relationship.databaseId === databaseId
&& (selected.has(relationship.sourceTableId) || selected.has(relationship.targetTableId))
));
for (const relationship of physicalRelationships) this.relationships.delete(relationship.id);
for (const relationship of logicalRelationships) this.logicalRelationships.delete(relationship.id);
this.refreshForeignKeyFlags(databaseId);
this.markCatalogIncomplete([databaseId]);
return {
tables: 0,
columns: 0,
relationships: physicalRelationships.length + logicalRelationships.length,
};
}
async planSchemaSync(
databaseId: string,
scope: CatalogSyncScope,
tableIds: readonly string[],
snapshot: ObservedSchemaSnapshot,
): Promise<CatalogSchemaDiff> {
this.assertSnapshotCapability(scope, snapshot);
const tables = await this.listTables(databaseId);
const selectedTableIds = new Set(tableIds);
const selectedTables = scope === "columns" && selectedTableIds.size > 0
? tables.filter((table) => selectedTableIds.has(table.id))
: tables;
const observedTables = new Set(snapshot.tables.map((table) => table.name));
const observedColumns = new Set(snapshot.columns.map((column) => `${column.tableName}\u0000${column.name}`));
const observedRelationships = new Set(
snapshot.relationships.map((relationship) => `${relationship.sourceTableName}\u0000${relationship.constraintName}`),
);
const deletedTableIds = new Set(tables.filter((table) => !observedTables.has(table.name)).map((table) => table.id));
return {
deletedTables: scope === "tables" || scope === "all"
? tables.filter((table) => !observedTables.has(table.name)).map((table) => table.name).sort()
: [],
deletedColumns: scope === "tables" || scope === "columns" || scope === "all"
? [...this.columns.values()]
.filter((column) => scope === "tables" ? deletedTableIds.has(column.tableId) : selectedTables.some((table) => table.id === column.tableId))
.filter((column) => {
const table = tables.find((candidate) => candidate.id === column.tableId);
return table && (scope === "tables" || !observedColumns.has(`${table.name}\u0000${column.name}`));
})
.map((column) => ({
tableName: tables.find((table) => table.id === column.tableId)?.name ?? "",
columnName: column.name,
}))
.sort((a, b) => `${a.tableName}.${a.columnName}`.localeCompare(`${b.tableName}.${b.columnName}`))
: [],
deletedRelationships: scope === "tables" || scope === "relationships" || scope === "all"
? [...this.relationships.values()]
.filter((relationship) => relationship.databaseId === databaseId)
.filter((relationship) => scope === "tables"
? deletedTableIds.has(relationship.sourceTableId) || deletedTableIds.has(relationship.targetTableId)
: !observedRelationships.has(`${relationship.sourceTableName}\u0000${relationship.constraintName}`))
.map((relationship) => ({
sourceTableName: relationship.sourceTableName,
constraintName: relationship.constraintName,
}))
.sort((a, b) => `${a.sourceTableName}.${a.constraintName}`.localeCompare(`${b.sourceTableName}.${b.constraintName}`))
: [],
};
}
async applySchemaSync(
databaseId: string,
expectedDatabaseVersion: number,
scope: CatalogSyncScope,
tableIds: readonly string[],
snapshot: ObservedSchemaSnapshot,
): Promise<CatalogSyncCounts | undefined> {
const database = this.records.get(databaseId);
if (!database || database.version !== expectedDatabaseVersion) return undefined;
this.assertSnapshotCapability(scope, snapshot);
const tableBackup = structuredClone([...this.tables.entries()]);
const columnBackup = structuredClone([...this.columns.entries()]);
const relationshipBackup = structuredClone([...this.relationships.entries()]);
const now = new Date().toISOString();
let created = 0;
let updated = 0;
let deleted = 0;
try {
if (scope === "tables" || scope === "all") {
const observedByName = new Map(snapshot.tables.map((table) => [table.name, table]));
const existing = await this.listTables(databaseId);
for (const table of existing) {
const observed = observedByName.get(table.name);
if (!observed) {
this.deleteTable(table.id);
deleted += 1;
} else {
const changed = table.sourceComment !== observed.sourceComment;
this.tables.set(table.id, {
...table,
sourceComment: observed.sourceComment,
lastSyncedDatabaseVersion: expectedDatabaseVersion,
lastSyncedAt: now,
version: changed ? table.version + 1 : table.version,
updatedAt: changed ? now : table.updatedAt,
});
if (changed) updated += 1;
observedByName.delete(table.name);
}
}
for (const observed of observedByName.values()) {
const id = randomUUID();
this.tables.set(id, {
id,
databaseId,
name: observed.name,
sourceComment: observed.sourceComment,
description: null,
generatedDescription: null,
lastSyncedDatabaseVersion: expectedDatabaseVersion,
lastSyncedAt: now,
version: 1,
createdAt: now,
updatedAt: now,
});
created += 1;
}
}
if (scope === "columns" || scope === "all") {
const currentTables = await this.listTables(databaseId);
const selectedIds = scope === "all" || tableIds.length === 0
? new Set(currentTables.map((table) => table.id))
: new Set(tableIds);
const selectedTables = currentTables.filter((table) => selectedIds.has(table.id));
if (scope === "columns" && selectedTables.length !== selectedIds.size) {
throw new CatalogConnectorError("One or more selected tables no longer exist");
}
const selectedNames = new Set(selectedTables.map((table) => table.name));
const tableByName = new Map(selectedTables.map((table) => [table.name, table]));
const observedByKey = new Map(
snapshot.columns
.filter((column) => selectedNames.has(column.tableName))
.map((column) => [`${column.tableName}\u0000${column.name}`, column]),
);
for (const column of [...this.columns.values()].filter((candidate) => selectedIds.has(candidate.tableId))) {
const table = selectedTables.find((candidate) => candidate.id === column.tableId);
if (!table) continue;
const key = `${table.name}\u0000${column.name}`;
const observed = observedByKey.get(key);
if (!observed) {
this.deleteColumn(column.id);
deleted += 1;
} else {
const changed = column.ordinalPosition !== observed.ordinalPosition
|| column.dataType !== observed.dataType
|| column.isNullable !== observed.isNullable
|| column.defaultExpression !== observed.defaultExpression
|| column.primaryKeyPosition !== observed.primaryKeyPosition
|| column.sourceComment !== observed.sourceComment;
this.columns.set(column.id, {
...column,
ordinalPosition: observed.ordinalPosition,
dataType: observed.dataType,
isNullable: observed.isNullable,
defaultExpression: observed.defaultExpression,
primaryKeyPosition: observed.primaryKeyPosition,
isPrimaryKey: observed.primaryKeyPosition !== null,
sourceComment: observed.sourceComment,
lastSyncedDatabaseVersion: expectedDatabaseVersion,
lastSyncedAt: now,
version: changed ? column.version + 1 : column.version,
updatedAt: changed ? now : column.updatedAt,
});
if (changed) updated += 1;
observedByKey.delete(key);
}
}
for (const observed of observedByKey.values()) {
const table = tableByName.get(observed.tableName);
if (!table) continue;
const id = randomUUID();
this.columns.set(id, {
id,
tableId: table.id,
name: observed.name,
ordinalPosition: observed.ordinalPosition,
dataType: observed.dataType,
isNullable: observed.isNullable,
defaultExpression: observed.defaultExpression,
primaryKeyPosition: observed.primaryKeyPosition,
isPrimaryKey: observed.primaryKeyPosition !== null,
isForeignKey: false,
foreignKeyCount: 0,
sourceComment: observed.sourceComment,
description: null,
generatedDescription: null,
sensitive: false,
sensitivityReason: null,
lastSyncedDatabaseVersion: expectedDatabaseVersion,
lastSyncedAt: now,
version: 1,
createdAt: now,
updatedAt: now,
});
created += 1;
}
}
if (scope === "relationships" || scope === "all") {
const tables = await this.listTables(databaseId);
const tableByName = new Map(tables.map((table) => [table.name, table]));
const existing = [...this.relationships.values()].filter((relationship) => relationship.databaseId === databaseId);
const existingByKey = new Map(existing.map((relationship) => [`${relationship.sourceTableName}\u0000${relationship.constraintName}`, relationship]));
const observedKeys = new Set(snapshot.relationships.map((relationship) => `${relationship.sourceTableName}\u0000${relationship.constraintName}`));
for (const relationship of existing) {
if (!observedKeys.has(`${relationship.sourceTableName}\u0000${relationship.constraintName}`)) {
this.relationships.delete(relationship.id);
deleted += 1;
}
}
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 ${observed.constraintName} refers to an unknown table`);
}
const pairs = observed.columns.map((pair) => {
const source = [...this.columns.values()].find((column) => column.tableId === sourceTable.id && column.name === pair.sourceColumnName);
const target = [...this.columns.values()].find((column) => column.tableId === targetTable.id && column.name === pair.targetColumnName);
if (!source || !target) {
throw new CatalogConnectorError(`Relationship ${observed.constraintName} refers to an unknown column`);
}
return {
position: pair.position,
sourceColumnId: source.id,
sourceColumnName: source.name,
targetColumnId: target.id,
targetColumnName: target.name,
};
}).sort((a, b) => a.position - b.position);
const key = `${observed.sourceTableName}\u0000${observed.constraintName}`;
const current = existingByKey.get(key);
const comparable = current && JSON.stringify({
target: current.targetTableName,
update: current.updateRule,
delete: current.deleteRule,
deferrable: current.deferrable,
deferred: current.initiallyDeferred,
columns: current.columns.map((pair) => [pair.position, pair.sourceColumnName, pair.targetColumnName]),
});
const nextComparable = JSON.stringify({
target: observed.targetTableName,
update: observed.updateRule,
delete: observed.deleteRule,
deferrable: observed.deferrable,
deferred: observed.initiallyDeferred,
columns: pairs.map((pair) => [pair.position, pair.sourceColumnName, pair.targetColumnName]),
});
const id = current?.id ?? randomUUID();
this.relationships.set(id, {
id,
databaseId,
constraintName: observed.constraintName,
sourceTableId: sourceTable.id,
sourceTableName: sourceTable.name,
targetTableId: targetTable.id,
targetTableName: targetTable.name,
updateRule: observed.updateRule,
deleteRule: observed.deleteRule,
deferrable: observed.deferrable,
initiallyDeferred: observed.initiallyDeferred,
columns: pairs,
lastSyncedDatabaseVersion: expectedDatabaseVersion,
lastSyncedAt: now,
createdAt: current?.createdAt ?? now,
updatedAt: comparable === nextComparable ? (current?.updatedAt ?? now) : now,
origin: "physical",
status: "active",
});
if (!current) created += 1;
else if (comparable !== nextComparable) updated += 1;
}
this.refreshForeignKeyFlags(databaseId);
}
if (scope === "all") {
this.records.set(databaseId, { ...database, schemaSyncedVersion: expectedDatabaseVersion, schemaSyncedAt: now });
}
return {
tables: (await this.listTables(databaseId)).length,
columns: [...this.columns.values()].filter((column) => this.tables.get(column.tableId)?.databaseId === databaseId).length,
relationships: (await this.listRelationships(databaseId)).length,
created,
updated,
deleted,
};
} catch (error) {
this.tables.clear();
this.columns.clear();
this.relationships.clear();
for (const [id, table] of tableBackup) this.tables.set(id, table);
for (const [id, column] of columnBackup) this.columns.set(id, column);
for (const [id, relationship] of relationshipBackup) this.relationships.set(id, relationship);
throw error;
}
}
async createSyncRun(
databaseId: string,
scope: CatalogSyncScope,
tableIds: readonly string[],
requestedDatabaseVersion: number,
): Promise<CatalogSyncRun> {
const activeStates = new Set<CatalogSyncRun["state"]>(["queued", "running", "awaiting_confirmation", "applying"]);
if ([...this.syncRuns.values()].some((run) => run.databaseId === databaseId && activeStates.has(run.state))) {
throw new CatalogConflictError("A schema synchronization is already active for this database");
}
const now = new Date().toISOString();
const run: CatalogSyncRun = {
id: randomUUID(), databaseId, scope, tableIds: [...tableIds], state: "queued", phase: "queued",
requestedDatabaseVersion, observedSnapshot: null, plannedDiff: null, confirmationToken: null,
counts: {}, errorCode: null, errorMessage: null, cancelRequested: false,
createdAt: now, startedAt: null, updatedAt: now, finishedAt: null, heartbeatAt: null,
leaseOwner: null, leaseExpiresAt: null,
};
this.syncRuns.set(run.id, run);
return structuredClone(run);
}
async getSyncRun(runId: string): Promise<CatalogSyncRun | undefined> {
const run = this.syncRuns.get(runId);
return run ? structuredClone(run) : undefined;
}
async claimSyncRun(runId: string, workerId: string, leaseExpiresAt: string): Promise<CatalogSyncRun | undefined> {
const run = this.syncRuns.get(runId);
if (!run || run.state !== "queued") return undefined;
if (run.leaseOwner && run.leaseOwner !== workerId && run.leaseExpiresAt && run.leaseExpiresAt > new Date().toISOString()) {
return undefined;
}
return await this.updateSyncRun(runId, {
state: "running",
startedAt: run.startedAt ?? new Date().toISOString(),
heartbeatAt: new Date().toISOString(),
leaseOwner: workerId,
leaseExpiresAt,
});
}
async listSyncRuns(databaseId: string, limit = 20): Promise<CatalogSyncRun[]> {
return [...this.syncRuns.values()].filter((run) => run.databaseId === databaseId)
.sort((a, b) => b.createdAt.localeCompare(a.createdAt)).slice(0, limit).map((run) => structuredClone(run));
}
async updateSyncRun(runId: string, update: CatalogSyncRunUpdate): Promise<CatalogSyncRun | undefined> {
const current = this.syncRuns.get(runId);
if (!current) return undefined;
const updated = { ...current, ...structuredClone(update), updatedAt: new Date().toISOString() };
this.syncRuns.set(runId, updated);
return structuredClone(updated);
}
async requestSyncRunCancellation(runId: string): Promise<CatalogSyncRun | undefined> {
const run = this.syncRuns.get(runId);
if (!run) return undefined;
if (!["queued", "running", "awaiting_confirmation"].includes(run.state)) return structuredClone(run);
return await this.updateSyncRun(runId, { cancelRequested: true });
}
async appendSyncEvent(
runId: string,
level: CatalogSyncEvent["level"],
eventType: string,
message: string,
data: Record<string, unknown> = {},
): Promise<CatalogSyncEvent> {
const events = this.syncEvents.get(runId) ?? [];
const event: CatalogSyncEvent = {
id: [...this.syncEvents.values()].reduce((count, values) => count + values.length, 0) + 1,
runId, sequence: events.length + 1, level, eventType, message, data: structuredClone(data),
createdAt: new Date().toISOString(),
};
events.push(event);
this.syncEvents.set(runId, events);
return structuredClone(event);
}
async listSyncEvents(runId: string, afterSequence = 0): Promise<CatalogSyncEvent[]> {
return (this.syncEvents.get(runId) ?? []).filter((event) => event.sequence > afterSequence).map((event) => structuredClone(event));
}
async pruneSyncEvents(before: string): Promise<void> {
for (const [runId, events] of this.syncEvents) {
this.syncEvents.set(runId, events.filter((event) => event.createdAt >= before));
}
}
async interruptActiveSyncRuns(): Promise<void> {
for (const run of this.syncRuns.values()) {
if (["queued", "running", "awaiting_confirmation", "applying"].includes(run.state)) {
await this.updateSyncRun(run.id, {
state: "interrupted", phase: "completed", finishedAt: new Date().toISOString(),
errorCode: "SYNC_INTERRUPTED", errorMessage: "Synchronization was interrupted by a service restart",
});
}
}
}
async reconcileTables(
databaseId: string,
expectedDatabaseVersion: number,
observed: readonly ObservedCatalogTable[],
confirmedDeletedNames: readonly string[],
): Promise<TableSyncRepositoryResult | undefined> {
const database = this.records.get(databaseId);
if (!database || database.version !== expectedDatabaseVersion) return undefined;
const existing = await this.listTables(databaseId);
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 };
}
const byName = new Map(existing.map((table) => [table.name, table]));
let createdCount = 0;
let updatedCount = 0;
for (const observedTable of observed) {
const current = byName.get(observedTable.name);
const now = new Date().toISOString();
if (!current) {
const created: CatalogTable = {
id: randomUUID(),
databaseId,
name: observedTable.name,
sourceComment: observedTable.sourceComment,
description: null,
generatedDescription: null,
lastSyncedDatabaseVersion: expectedDatabaseVersion,
lastSyncedAt: now,
version: 1,
createdAt: now,
updatedAt: now,
};
this.tables.set(created.id, created);
createdCount += 1;
} else if (current.sourceComment !== observedTable.sourceComment) {
this.tables.set(current.id, {
...current,
sourceComment: observedTable.sourceComment,
version: current.version + 1,
updatedAt: now,
});
updatedCount += 1;
}
}
for (const table of existing) {
if (deletedNames.includes(table.name)) this.deleteTable(table.id);
}
return {
kind: "applied",
createdCount,
updatedCount,
deletedCount: deletedNames.length,
tables: await this.listTables(databaseId),
};
}
private assertSnapshotCapability(scope: CatalogSyncScope, snapshot: ObservedSchemaSnapshot): void {
const required = scope === "all" ? ["tables", "columns", "relationships"] as const : [scope] as const;
for (const capability of required) {
if (snapshot.capabilities[capability] !== "available") {
throw new CatalogConnectorError(`Schema introspection capability '${capability}' is unavailable`);
}
}
}
private deleteTable(tableId: string): void {
this.tables.delete(tableId);
for (const column of [...this.columns.values()]) if (column.tableId === tableId) this.deleteColumn(column.id);
for (const relationship of [...this.relationships.values()]) {
if (relationship.sourceTableId === tableId || relationship.targetTableId === tableId) {
this.relationships.delete(relationship.id);
}
}
}
private deleteColumn(columnId: string): void {
this.columns.delete(columnId);
for (const relationship of [...this.relationships.values()]) {
if (relationship.columns.some((pair) => pair.sourceColumnId === columnId || pair.targetColumnId === columnId)) {
this.relationships.delete(relationship.id);
}
}
for (const relationship of [...this.logicalRelationships.values()]) {
const pair = relationship.columns[0];
if (pair.sourceColumnId === columnId || pair.targetColumnId === columnId) {
this.logicalRelationships.delete(relationship.id);
}
}
}
private markCatalogIncomplete(databaseIds: readonly string[]): void {
for (const databaseId of databaseIds) {
const database = this.records.get(databaseId);
if (!database) continue;
this.records.set(databaseId, {
...database,
schemaSyncedVersion: undefined,
schemaSyncedAt: undefined,
});
}
}
private refreshForeignKeyFlags(databaseId: string): void {
const counts = new Map<string, number>();
for (const relationship of this.relationships.values()) {
if (relationship.databaseId !== databaseId) continue;
for (const pair of relationship.columns) counts.set(pair.sourceColumnId, (counts.get(pair.sourceColumnId) ?? 0) + 1);
}
for (const column of this.columns.values()) {
if (this.tables.get(column.tableId)?.databaseId !== databaseId) continue;
const foreignKeyCount = counts.get(column.id) ?? 0;
this.columns.set(column.id, { ...column, isForeignKey: foreignKeyCount > 0, foreignKeyCount });
}
}
async available(): Promise<boolean> { return true; }
}