/** * Migration runner with failure diagnostic — [E00-S03-T05]. * * The migration runner applies pending migrations to the database exactly * once and, when a migration fails, produces a structured diagnostic that * identifies the failing migration. It is built on the boundary that came * before it: the migration ledger (E00-S03-T03) records applied migrations, * so a rerun never double-applies. The advisory lock (E00-S03-T04) serializes * concurrent runners, but wiring the lock into the runner is out of scope for * this story — the runner performs no locking itself; a runner that wants to * serialize takes the lock (E00-S03-T04) around `run()`. * * Failure model: * - migrations run in the order given, oldest first; a migration whose * version is already recorded in the ledger is skipped; * - when a migration's `up` throws, the runner wraps the failure into a * `MigrationFailedError` whose `diagnostic` is a structured object that * identifies the failing migration (`migration` — its version), where the * run failed (`phase`: 'apply' when the migration's `up` threw, 'record' * when the ledger insert threw after a successful `up`), the underlying * cause, and the ledger state at failure time (`applied`/`pending` — * `applied` + `pending` cover the runner's migrations exactly); * - the error is serializable: `toJSON()` returns a plain structured object * (including a structured cause — for pg errors the `code`, e.g. `42P01`), * so operators can log/parse the diagnostic without string-matching. * * The runner lives in `database-postgres` — the single workspace package * allowed to import the PostgreSQL driver (E00-S03-T02) — and talks to the * database exclusively through the package-owned `pg` Pool and the * `MigrationLedger`, so no other package needs the driver to run migrations. * * Rollback note from the issue: revert the diagnostic/error handling changes. */ import type { Pool } from 'pg'; import type { MigrationLedger } from './ledger.js'; /** * A single migration step: an identifier (recorded in the ledger once the * step has been applied) and the apply function. `up` receives the * package-owned pool, so a migration can run any SQL (and multi-statement * work) through the same driver boundary the runner itself uses. */ export interface Migration { version: string; up(pool: Pool): Promise | void; } /** * Where a migration run failed: 'apply' when the migration's `up` threw, or * 'record' when the ledger insert threw after a successful `up`. */ export type MigrationFailurePhase = 'apply' | 'record'; /** * The structured diagnostic produced when a migration fails. `migration` * identifies the failing migration; `applied` and `pending` are disjoint and * together cover the runner's migration list in run order. */ export interface MigrationDiagnostic { /** Version of the migration that failed — identifies the failing migration. */ migration: string; /** Where the run failed: 'apply' (the migration's `up` threw) or 'record' (the ledger insert threw). */ phase: MigrationFailurePhase; /** The underlying failure (e.g. the pg error), preserved for inspection. */ cause: unknown; /** Versions recorded in the ledger when the failure happened, in run order. */ applied: string[]; /** Versions not yet recorded when the failure happened, in run order — includes the failing migration. */ pending: string[]; } /** Result of a successful migration run. */ export interface MigrationRunResult { /** Versions applied by this run, in run order (oldest first). */ applied: string[]; /** Versions skipped because they were already recorded in the ledger. */ skipped: string[]; } /** * The error thrown when a migration fails. Carries the structured diagnostic * (`.diagnostic`) and is serializable (`toJSON()`), so callers and operators * can inspect and parse the failure without string-matching the message. */ export class MigrationFailedError extends Error { readonly diagnostic: MigrationDiagnostic; constructor(diagnostic: MigrationDiagnostic) { super(`migration "${diagnostic.migration}" failed during ${diagnostic.phase}`); this.name = 'MigrationFailedError'; this.diagnostic = diagnostic; } /** Serializable form of the error and its structured diagnostic. */ toJSON(): Record { return { name: this.name, message: this.message, diagnostic: { migration: this.diagnostic.migration, phase: this.diagnostic.phase, applied: [...this.diagnostic.applied], pending: [...this.diagnostic.pending], cause: structuredCause(this.diagnostic.cause), }, }; } } /** * Reduces the underlying cause to a structured, serializable shape: for an * `Error` the name/message (plus the pg error `code` when present, e.g. * `42P01`); for anything else a `{ value }` wrapper, so `toJSON()` never * stringifies to an empty object. */ function structuredCause(cause: unknown): Record { if (cause instanceof Error) { const structured: Record = { name: cause.name, message: cause.message }; const code = (cause as Error & { code?: unknown }).code; if (code !== undefined) structured.code = code; return structured; } return { value: cause }; } /** * The migration runner: applies pending migrations in order, exactly once, * through the migration ledger. Instances are cheap and share the caller's * pool and ledger; the runner performs no locking (advisory lock is * E00-S03-T04) and no ready gating (E00-S03-T06). */ export class MigrationRunner { private readonly pool: Pool; private readonly ledger: MigrationLedger; private readonly migrations: readonly Migration[]; /** * @param pool The package-owned PostgreSQL pool (`pg.Pool`), handed to each * migration's `up`. * @param migrations The migrations to run, in apply order (oldest first). * @param ledger The migration ledger (E00-S03-T03) the runner reads applied * versions from and records applied migrations into; constructed by the * caller from the same pool (`new MigrationLedger(pool)`). */ constructor(pool: Pool, migrations: readonly Migration[], ledger: MigrationLedger) { this.pool = pool; this.migrations = migrations; this.ledger = ledger; } /** * Runs the pending migrations: creates the ledger if missing, then applies * every migration whose version is not yet recorded, recording each one as * it completes. Idempotent: a rerun skips everything already recorded, so a * second run never double-applies. When a migration fails — its `up` throws * ('apply') or the ledger insert throws ('record') — the runner throws a * `MigrationFailedError` whose `diagnostic` identifies the failing migration * and the ledger state at failure time. */ async run(): Promise { await this.ledger.ensure(); const recorded = new Set(await this.ledger.applied()); const applied: string[] = []; const skipped: string[] = []; for (const migration of this.migrations) { if (recorded.has(migration.version)) { skipped.push(migration.version); continue; } try { await migration.up(this.pool); } catch (error) { throw this.failure(migration.version, 'apply', error, recorded); } try { await this.ledger.record(migration.version); } catch (error) { throw this.failure(migration.version, 'record', error, recorded); } recorded.add(migration.version); applied.push(migration.version); } return { applied, skipped }; } /** * Builds the structured diagnostic for a failing migration: `migration` is * the failing version, `applied`/`pending` are the ledger state at failure * time scoped to this runner's migrations (disjoint, in run order — the * failing migration is still pending, since it was never recorded). */ private failure( version: string, phase: MigrationFailurePhase, cause: unknown, recorded: ReadonlySet, ): MigrationFailedError { const applied: string[] = []; const pending: string[] = []; for (const migration of this.migrations) { if (recorded.has(migration.version)) applied.push(migration.version); else pending.push(migration.version); } return new MigrationFailedError({ migration: version, phase, cause, applied, pending }); } }