Files
ThothII/backend/src/catalog/sync-worker.ts
T

307 lines
12 KiB
TypeScript

import { randomUUID } from "node:crypto";
import type { CatalogOperationCoordinator } from "./operation-coordinator.js";
import type { CatalogSchemaIntrospector, CatalogSchemaScanProgress } from "./schema-introspector.js";
import {
CatalogConflictError,
CatalogConnectorError,
CatalogSchemaCapabilityUnavailableError,
type CatalogRepository,
type CatalogSchemaDiff,
type CatalogSyncCounts,
type CatalogSyncRun,
type CatalogSyncScope,
type ObservedSchemaSnapshot,
type WorkspaceDatabase,
} from "./types.js";
class SyncCancelledError extends Error {}
const TERMINAL_STATES = new Set<CatalogSyncRun["state"]>([
"succeeded", "failed", "cancelled", "interrupted",
]);
function destructive(diff: CatalogSchemaDiff): boolean {
return diff.deletedTables.length > 0
|| diff.deletedColumns.length > 0
|| diff.deletedRelationships.length > 0;
}
function fingerprint(snapshot: ObservedSchemaSnapshot): string {
return JSON.stringify(snapshot);
}
function safeFailure(error: unknown): { code: string; message: string } {
if (error instanceof CatalogSchemaCapabilityUnavailableError) {
const label = error.capability === "schema_snapshot" ? "schema snapshot" : error.capability;
return {
code: "schema_capability_unavailable",
message: `This database binding does not provide the ${label} capability.`,
};
}
if (error instanceof CatalogConnectorError) {
return {
code: "schema_introspection_failed",
message: "The database schema could not be read. Check the connection and credentials, then try again.",
};
}
return { code: "schema_sync_failed", message: "Schema synchronization failed." };
}
export class CatalogSyncWorker {
private readonly workerId = randomUUID();
private readonly controllers = new Map<string, AbortController>();
private readonly reservations = new Map<string, () => void>();
private stopping = false;
constructor(
private readonly repository: CatalogRepository,
private readonly introspector: CatalogSchemaIntrospector,
private readonly operations: CatalogOperationCoordinator,
private readonly timeoutMs: number,
) {}
async initialize(): Promise<void> {
if (!(await this.repository.available())) return;
await this.repository.interruptActiveSyncRuns();
await this.repository.pruneSyncEvents(new Date(Date.now() - 30 * 24 * 60 * 60 * 1_000).toISOString());
}
async start(database: WorkspaceDatabase, scope: CatalogSyncScope, tableIds: readonly string[]): Promise<CatalogSyncRun> {
const uniqueTableIds = [...new Set(tableIds)];
if (scope === "columns") {
const tables = await Promise.all(uniqueTableIds.map((tableId) => this.repository.getTable(database.id, tableId)));
if (tables.some((table) => !table)) throw new CatalogConflictError("One or more selected tables no longer exist");
}
if (scope !== "columns" && uniqueTableIds.length > 0) {
throw new CatalogConflictError("Table selection is only valid for a column synchronization");
}
const release = this.operations.reserve(database.id);
try {
const run = await this.repository.createSyncRun(database.id, scope, uniqueTableIds, database.version);
this.reservations.set(run.id, release);
await this.repository.appendSyncEvent(run.id, "info", "queued", "Synchronization queued.", { scope });
this.launch(run.id);
return run;
} catch (error) {
release();
throw error;
}
}
async confirm(runId: string, confirmationToken: string): Promise<CatalogSyncRun | undefined> {
const run = await this.repository.getSyncRun(runId);
if (!run) return undefined;
if (run.state !== "awaiting_confirmation" || !run.observedSnapshot || run.confirmationToken !== confirmationToken) {
throw new CatalogConflictError("Synchronization confirmation is no longer valid");
}
await this.repository.appendSyncEvent(run.id, "info", "confirmation_received", "Destructive changes were confirmed.");
const queued = await this.repository.updateSyncRun(run.id, {
state: "queued",
phase: "queued",
confirmationToken: null,
leaseOwner: null,
leaseExpiresAt: null,
});
this.launch(run.id, fingerprint(run.observedSnapshot));
return queued;
}
async cancel(runId: string): Promise<CatalogSyncRun | undefined> {
const run = await this.repository.getSyncRun(runId);
if (!run) return undefined;
if (run.state === "applying" || TERMINAL_STATES.has(run.state)) return run;
await this.repository.requestSyncRunCancellation(runId);
this.controllers.get(runId)?.abort();
if (run.state === "queued" || run.state === "awaiting_confirmation") {
const cancelled = await this.repository.updateSyncRun(runId, {
state: "cancelled",
phase: "completed",
finishedAt: new Date().toISOString(),
observedSnapshot: null,
plannedDiff: null,
confirmationToken: null,
leaseOwner: null,
leaseExpiresAt: null,
});
await this.repository.appendSyncEvent(runId, "warning", "cancelled", "Synchronization cancelled.");
this.release(runId);
return cancelled;
}
return await this.repository.getSyncRun(runId);
}
async retry(runId: string): Promise<CatalogSyncRun | undefined> {
const previous = await this.repository.getSyncRun(runId);
if (!previous) return undefined;
if (!TERMINAL_STATES.has(previous.state)) {
throw new CatalogConflictError("Only a finished synchronization can be retried");
}
const database = await this.repository.get(previous.databaseId);
if (!database) return undefined;
return await this.start(database, previous.scope, previous.tableIds);
}
async stop(): Promise<void> {
this.stopping = true;
for (const controller of this.controllers.values()) controller.abort();
if (await this.repository.available()) await this.repository.interruptActiveSyncRuns();
for (const runId of [...this.reservations.keys()]) this.release(runId);
}
private launch(runId: string, confirmedFingerprint?: string): void {
queueMicrotask(() => {
void this.execute(runId, confirmedFingerprint).catch(() => undefined);
});
}
private async execute(runId: string, confirmedFingerprint?: string): Promise<void> {
if (this.stopping) return;
const leaseExpiresAt = new Date(Date.now() + 20_000).toISOString();
const claimed = await this.repository.claimSyncRun(runId, this.workerId, leaseExpiresAt);
if (!claimed) return;
const controller = new AbortController();
this.controllers.set(runId, controller);
let timedOut = false;
const timeout = setTimeout(() => {
timedOut = true;
controller.abort();
}, this.timeoutMs);
const heartbeat = setInterval(() => {
void this.repository.updateSyncRun(runId, {
heartbeatAt: new Date().toISOString(),
leaseExpiresAt: new Date(Date.now() + 20_000).toISOString(),
});
}, 5_000);
try {
await this.repository.appendSyncEvent(runId, "info", "started", "Synchronization started.");
const database = await this.repository.get(claimed.databaseId);
if (!database || database.version !== claimed.requestedDatabaseVersion) {
throw new CatalogConflictError("Database binding changed before synchronization started");
}
const progress: CatalogSchemaScanProgress = async (phase, counts) => {
await this.checkCancelled(runId);
await this.repository.updateSyncRun(runId, {
phase,
heartbeatAt: new Date().toISOString(),
...(counts ? { counts } : {}),
});
await this.repository.appendSyncEvent(runId, "info", phase, this.phaseMessage(phase), counts ?? {});
};
const snapshot = await this.introspector.scan(database, controller.signal, progress);
await this.checkCancelled(runId);
this.assertCapability(claimed.scope, snapshot);
const counts: CatalogSyncCounts = {
tables: snapshot.tables.length,
columns: snapshot.columns.length,
relationships: snapshot.relationships.length,
};
await this.repository.updateSyncRun(runId, { phase: "planning", counts, observedSnapshot: snapshot });
await this.repository.appendSyncEvent(runId, "info", "planning", "Schema changes are being planned.", { ...counts });
const diff = await this.repository.planSchemaSync(claimed.databaseId, claimed.scope, claimed.tableIds, snapshot);
await this.checkCancelled(runId);
if (destructive(diff) && fingerprint(snapshot) !== confirmedFingerprint) {
const token = randomUUID();
const waiting = await this.repository.updateSyncRun(runId, {
state: "awaiting_confirmation",
phase: "awaiting_confirmation",
observedSnapshot: snapshot,
plannedDiff: diff,
confirmationToken: token,
counts,
leaseOwner: null,
leaseExpiresAt: null,
});
await this.repository.appendSyncEvent(runId, "warning", "confirmation_required", "Confirmation is required before removing catalog objects.", {
deletedTables: diff.deletedTables.length,
deletedColumns: diff.deletedColumns.length,
deletedRelationships: diff.deletedRelationships.length,
});
if (!waiting) throw new Error("Synchronization run disappeared");
return;
}
await this.repository.updateSyncRun(runId, { state: "applying", phase: "applying", plannedDiff: diff });
await this.repository.appendSyncEvent(runId, "info", "applying", "Catalog changes are being applied atomically.");
this.controllers.delete(runId);
const applied = await this.repository.applySchemaSync(
claimed.databaseId,
claimed.requestedDatabaseVersion,
claimed.scope,
claimed.tableIds,
snapshot,
);
if (!applied) throw new CatalogConflictError("Database binding changed before schema changes were applied");
await this.repository.updateSyncRun(runId, {
state: "succeeded",
phase: "completed",
counts: applied,
finishedAt: new Date().toISOString(),
observedSnapshot: null,
plannedDiff: null,
confirmationToken: null,
heartbeatAt: new Date().toISOString(),
leaseOwner: null,
leaseExpiresAt: null,
});
await this.repository.appendSyncEvent(runId, "info", "succeeded", "Synchronization completed.", { ...applied });
this.release(runId);
} catch (error) {
const current = await this.repository.getSyncRun(runId);
const cancelled = !timedOut && (error instanceof SyncCancelledError || controller.signal.aborted || current?.cancelRequested);
const failure = timedOut
? { code: "schema_sync_timed_out", message: "Schema synchronization timed out." }
: safeFailure(error);
await this.repository.updateSyncRun(runId, {
state: cancelled ? "cancelled" : "failed",
phase: "completed",
errorCode: cancelled ? null : failure.code,
errorMessage: cancelled ? null : failure.message,
finishedAt: new Date().toISOString(),
observedSnapshot: null,
plannedDiff: null,
confirmationToken: null,
leaseOwner: null,
leaseExpiresAt: null,
});
await this.repository.appendSyncEvent(
runId,
cancelled ? "warning" : "error",
cancelled ? "cancelled" : "failed",
cancelled ? "Synchronization cancelled." : failure.message,
);
this.release(runId);
} finally {
clearTimeout(timeout);
clearInterval(heartbeat);
this.controllers.delete(runId);
}
}
private assertCapability(scope: CatalogSyncScope, snapshot: ObservedSchemaSnapshot): void {
const required = scope === "all" ? ["tables", "columns", "relationships"] as const : [scope] as const;
for (const name of required) {
if (snapshot.capabilities[name] !== "available") {
throw new CatalogSchemaCapabilityUnavailableError(name);
}
}
}
private async checkCancelled(runId: string): Promise<void> {
const run = await this.repository.getSyncRun(runId);
if (run?.cancelRequested) throw new SyncCancelledError("Synchronization cancelled");
}
private release(runId: string): void {
this.reservations.get(runId)?.();
this.reservations.delete(runId);
}
private phaseMessage(phase: Parameters<CatalogSchemaScanProgress>[0]): string {
if (phase === "connecting") return "Connecting to the database.";
if (phase === "scanning_tables") return "Reading tables.";
if (phase === "scanning_columns") return "Reading columns and primary keys.";
return "Reading foreign-key relationships.";
}
}