From 299fd79b1e6ddaeecb7698ca6070f750afbad479 Mon Sep 17 00:00:00 2001 From: Isaac Menichim Isreal <194775393+big6isaac@users.noreply.github.com> Date: Tue, 29 Sep 2026 22:14:50 +0100 Subject: [PATCH] feat(migrations): add rollback support with destructive-migration safeguards #848 The migration runner declared a required down() on every migration but exposed no way to invoke it, so an incompatible schema change had no way back out. Add MigrationRunner.rollback(), which unwinds applied migrations newest first, each in its own transaction, validating every target before writing so a refusal cannot leave the database half-reverted. Reverts stamp migrations.reverted_at instead of deleting the row, preserving the audit trail, and re-applying a reverted migration revives its row. Migrations whose down() cannot reconstruct their data with up() are marked destructive and refused unless the caller opts in, so a routine rollback cannot silently drop production rows. check-migration-safety enforces that flag at review time by failing any migration whose down() uses DROP TABLE, DROP COLUMN or TRUNCATE without it. Fix three latent defects that made migrations unreliable: - sqlite3 only resolves when given a callback, so every await on a raw handle was a no-op. Migration DDL was fire-and-forget and could still be in flight when the transaction committed. Statements are now awaited through a promisified proxy. - db.serialize() invokes its callback synchronously and does not await an async one, so the runner resolved before work finished and callers could close the database mid-transaction. - 001 split its schema on ';', which cut CREATE TRIGGER ... BEGIN ... END bodies in half and left every trigger unparseable. The errors were swallowed by the two defects above. The script is now passed to exec(). Verified: 13 rollback tests, safety audit, scoped typecheck, prettier. Full suite matches the pre-change baseline exactly (136 pre-existing failures, unchanged) with 13 new tests passing. --- .github/workflows/migration-safety.yml | 77 ++++++ MIGRATION_ROLLBACK.md | 171 ++++++++++++ listener/package.json | 2 + .../src/database/migration-rollback.test.ts | 215 +++++++++++++++ listener/src/database/migration-system.ts | 258 ++++++++++++++++-- listener/src/migrations/001-initial-schema.ts | 14 +- .../002-query-performance-indexes.ts | 2 + .../src/scripts/check-migration-safety.ts | 111 ++++++++ listener/src/scripts/rollback-db.ts | 102 +++++++ 9 files changed, 921 insertions(+), 31 deletions(-) create mode 100644 .github/workflows/migration-safety.yml create mode 100644 MIGRATION_ROLLBACK.md create mode 100644 listener/src/database/migration-rollback.test.ts create mode 100644 listener/src/scripts/check-migration-safety.ts create mode 100644 listener/src/scripts/rollback-db.ts diff --git a/.github/workflows/migration-safety.yml b/.github/workflows/migration-safety.yml new file mode 100644 index 00000000..bd756bcd --- /dev/null +++ b/.github/workflows/migration-safety.yml @@ -0,0 +1,77 @@ +name: Migration Safety + +# Gates changes to the listener database migrations. +# +# A migration that cannot be rolled back, or a destructive one that is not flagged, +# is a production incident waiting to happen. Both are caught here rather than +# during a failed deploy. + +on: + push: + branches: [main] + paths: + - 'listener/src/migrations/**' + - 'listener/src/database/migration-system.ts' + - 'listener/src/scripts/check-migration-safety.ts' + - 'listener/src/scripts/rollback-db.ts' + - 'listener/src/database/migration-rollback.test.ts' + - '.github/workflows/migration-safety.yml' + pull_request: + paths: + - 'listener/src/migrations/**' + - 'listener/src/database/migration-system.ts' + - 'listener/src/scripts/check-migration-safety.ts' + - 'listener/src/scripts/rollback-db.ts' + - 'listener/src/database/migration-rollback.test.ts' + - '.github/workflows/migration-safety.yml' + +concurrency: + group: migration-safety-${{ github.ref }} + cancel-in-progress: true + +jobs: + migration-safety: + name: Migration rollback safety + runs-on: ubuntu-latest + + defaults: + run: + working-directory: listener + + steps: + - name: Checkout + uses: actions/checkout@v4 + + - name: Set up Node + uses: actions/setup-node@v4 + with: + node-version: '20' + cache: npm + cache-dependency-path: listener/package-lock.json + + - name: Install dependencies + run: npm install + + # Fails when a migration has no down(), or when a destructive down() is not + # flagged `destructive: true`. + - name: Audit migrations for rollback safety + run: npm run check-migration-safety + + # Scoped to the migration surface on purpose: the listener carries + # pre-existing type errors elsewhere, and this workflow gates migrations, + # not the whole app. A repo-wide `tsc --noEmit` would fail on files this + # workflow never touches. + - name: Typecheck migration surface + run: > + npx tsc --noEmit + --strict + --esModuleInterop + --skipLibCheck + --module CommonJS + --target es2020 + --types node,jest + src/database/migration-system.ts + src/scripts/check-migration-safety.ts + + - name: Migration rollback tests + run: npx jest src/database/migration-rollback.test.ts diff --git a/MIGRATION_ROLLBACK.md b/MIGRATION_ROLLBACK.md new file mode 100644 index 00000000..fe34cb51 --- /dev/null +++ b/MIGRATION_ROLLBACK.md @@ -0,0 +1,171 @@ +# Database Migration Rollback Strategy + +How to revert a listener database migration, and the safeguards that make reverting safe. + +The listener uses a custom SQLite migration runner (`listener/src/database/migration-system.ts`). +Migrations are applied forward in order, and can be reverted backward in the reverse of that order. + +## Quick reference + +```bash +cd listener + +# See what is currently applied, newest first, and what would be reverted next. +npm run migrate:rollback -- --status + +# Revert the single most recent migration. +npm run migrate:rollback + +# Revert the three most recent migrations. +npm run migrate:rollback -- --steps 3 + +# Revert a migration that destroys data (see safeguards below). +npm run migrate:rollback -- --allow-destructive + +# Audit the migration directory for rollback safety. +npm run check-migration-safety +``` + +Point the CLI at a specific database with `DATABASE_PATH`: + +```bash +DATABASE_PATH=./data/notifications.db npm run migrate:rollback -- --status +``` + +## How reverts behave + +- Migrations revert **newest first**, exactly reversing the order they were applied in. +- Each revert runs in its **own transaction**. If a revert fails, that step is rolled back + and the database is left as it was. +- Every target is **validated before anything is written**, so a refusal partway through a + multi-step rollback cannot leave the database half-reverted. +- Reverts are **recorded, not erased**. The `migrations` row is stamped with `reverted_at` + instead of being deleted, so the record of what was applied and when it was reverted + survives. Re-applying a reverted migration revives its row rather than failing. +- Reverting does not undo data written *after* the migration ran. It only reverses the + schema change the migration made. + +## Safeguards against destructive reverts + +A migration whose `down()` cannot reconstruct its data with `up()` is marked +`destructive: true` in the migration file: + +```ts +const migration = { + id: '003', + name: 'drop-legacy-column', + destructive: true, + up: async (db) => { /* ... */ }, + down: async (db) => { /* DROP TABLE / DROP COLUMN / TRUNCATE */ }, +}; +``` + +`rollback()` **refuses** to revert a destructive migration unless the caller passes +`allowDestructive`, and the CLI requires `--allow-destructive`. This means a routine +`npm run migrate:rollback` can never quietly drop production rows. When the guard trips, +the command names the migration at risk and exits with status `3`: + +``` +❌ Refusing to roll back destructive migration 001 (initial-schema). + Re-run with allowDestructive to confirm the data loss is intended. +``` + +Migration `001` is flagged because its `down()` drops all ten listener tables. Migration +`002` is not flagged, because it only drops indexes, which `up()` recreates without +touching rows. + +### `npm run check-migration-safety` + +Audits every migration file and fails when: + +1. A migration has no `down()` function, so it cannot be rolled back at all. +2. A migration's `down()` contains `DROP TABLE`, `DROP COLUMN`, or `TRUNCATE` but is not + flagged `destructive: true`. + +Rule 2 matters because an unflagged destructive migration is invisible to the runtime +guard: `rollback()` would revert it like any other and discard data without warning. +Exit code `1` on findings, `0` when clean. + +## Authoring a new migration + +Every migration needs a working `down()`: + +```ts +import * as sqlite3 from 'sqlite3'; + +const migration = { + id: '003', + name: 'add-notification-priority', + up: async (db: sqlite3.Database) => { + await db.run('ALTER TABLE scheduled_notifications ADD COLUMN priority INTEGER DEFAULT 0'); + }, + down: async (db: sqlite3.Database) => { + // SQLite before 3.35 cannot drop a column; recreate the table instead. + await db.run('ALTER TABLE scheduled_notifications DROP COLUMN priority'); + }, +}; + +export default migration; +``` + +Guidelines: + +- Prefer **additive, reversible changes**: add a column with a default, add an index, add + a table. These revert cleanly. +- If a change must **remove** data, split it across two migrations: first stop writing the + data, then drop it in a later, `destructive: true` migration. That gives operators a + rollback window before anything is lost. +- `db.exec()` is used for multi-statement scripts. Do **not** split a script on `;`, because + `CREATE TRIGGER ... BEGIN ... END` bodies contain their own semicolons and splitting on + them produces statements that fail to parse. +- Run `npm run check-migration-safety` before opening a PR that adds a migration. + +## Operational runbook + +**A migration is failing to apply and is blocking a deploy.** + +The runner wraps each `up()` in a transaction, so a failed migration leaves no partial +schema. Fix the migration and re-run `npm run migrate`. + +**A migration applied but is causing bad behaviour in production.** + +1. `npm run migrate:rollback -- --status` to confirm what would be reverted. +2. `npm run migrate:rollback` to revert the most recent migration. +3. If it was flagged destructive, stop and confirm the data loss is intended first, then + re-run with `--allow-destructive`. +4. Take a backup before any destructive revert. SQLite's `VACUUM INTO` command produces a + consistent copy of a live database: + + ```bash + sqlite3 ./data/notifications.db "VACUUM INTO './data/notifications-backup.db'" + ``` + +**Rolling back a release that included several migrations.** + +Revert them together, newest first, with `--steps`. Check the list first, because the count +includes any destructive migration in range: + +```bash +npm run migrate:rollback -- --status +npm run migrate:rollback -- --steps 3 +``` + +## Verifying state + +`migrations` rows carry `applied_at` and `reverted_at`. A `NULL` `reverted_at` means the +migration is currently in effect: + +```bash +sqlite3 ./data/notifications.db \ + "SELECT id, name, applied_at, reverted_at FROM migrations ORDER BY applied_at DESC;" +``` + +`npm run check-migrations` verifies there are no pending migrations, which is the check that +a deploy target is fully up to date. + +## CI + +`npm run check-migration-safety` and `npm test` (which includes +`src/database/migration-rollback.test.ts`) cover the rollback behaviour. The safety audit is +cheap enough to gate every build; a migration without a rollback path, or a destructive one +that is not flagged, fails the build instead of being discovered during an incident. diff --git a/listener/package.json b/listener/package.json index 4a5508e9..e1b2bab8 100644 --- a/listener/package.json +++ b/listener/package.json @@ -13,6 +13,8 @@ "test:stress": "node ./node_modules/jest/bin/jest.js src/__tests__/stress.test.ts --runInBand --detectOpenHandles", "stress-test": "ts-node src/scripts/run-stress-tests.ts", "migrate": "ts-node src/scripts/migrate-db.ts", + "migrate:rollback": "ts-node src/scripts/rollback-db.ts", + "check-migration-safety": "ts-node src/scripts/check-migration-safety.ts", "migrate:templates": "ts-node src/scripts/migrate-templates.ts", "check-migrations": "ts-node src/scripts/check-migrations.ts", "validate:batch": "ts-node src/utils/batch-validator.ts" diff --git a/listener/src/database/migration-rollback.test.ts b/listener/src/database/migration-rollback.test.ts new file mode 100644 index 00000000..cffd62ae --- /dev/null +++ b/listener/src/database/migration-rollback.test.ts @@ -0,0 +1,215 @@ +/** + * Rollback behaviour for the listener migration system. + * + * Covers the three properties the rollback strategy depends on: + * - reverts unwind in reverse order and actually reverse the schema + * - destructive migrations are refused unless explicitly allowed + * - a failed revert leaves the database untouched + */ +import * as fs from 'fs'; +import * as os from 'os'; +import * as path from 'path'; +import * as sqlite3 from 'sqlite3'; +import { Database } from '../database/database'; +import { DestructiveMigrationError, MigrationRunner } from '../database/migration-system'; + +const MIGRATIONS_DIR = path.join(__dirname, '../migrations'); + +interface Harness { + db: Database; + sqliteDb: sqlite3.Database; + dbPath: string; + runner: MigrationRunner; +} + +async function createHarness(): Promise { + const dbPath = path.join(os.tmpdir(), `notify-rollback-${Date.now()}-${Math.random()}.db`); + const db = new Database(dbPath); + await db.initialize(); + // @ts-ignore - Accessing private db property + const sqliteDb = db['db'] as any; + const runner = new MigrationRunner(sqliteDb, MIGRATIONS_DIR); + return { db, sqliteDb, dbPath, runner }; +} + +async function destroy(harness: Harness): Promise { + await harness.db.close(); + if (fs.existsSync(harness.dbPath)) fs.unlinkSync(harness.dbPath); +} + +function tableNames(db: sqlite3.Database): Promise { + return new Promise((resolve, reject) => { + db.all( + "SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'scheduled_notifications'", + (error, rows: unknown[]) => + error ? reject(error) : resolve((rows as any[]).map((r) => r.name)), + ); + }); +} + +function indexNames(db: sqlite3.Database): Promise { + return new Promise((resolve, reject) => { + db.all( + "SELECT name FROM sqlite_master WHERE type = 'index' AND name LIKE 'idx_%'", + (error, rows: unknown[]) => + error ? reject(error) : resolve((rows as any[]).map((r) => r.name)), + ); + }); +} + +async function countTables(db: Database): Promise { + const rows = await db.all<{ name: string }>( + "SELECT name FROM sqlite_master WHERE type = 'table' AND name LIKE 't_%'", + ); + return rows.length; +} + +describe('MigrationRunner.rollback', () => { + let harness: Harness; + + beforeEach(async () => { + harness = await createHarness(); + await harness.runner.runMigrations(); + }); + + afterEach(async () => { + await destroy(harness); + }); + + it('applies all migrations on first run', async () => { + const applied = await harness.runner.getAppliedMigrations(); + expect(applied).toEqual(['001', '002']); + }); + + it('reverts the most recent migration by default', async () => { + const reverted = await harness.runner.rollback(); + + expect(reverted.map((m) => m.id)).toEqual(['002']); + const indexes = await indexNames(harness.sqliteDb); + expect(indexes).not.toContain('idx_scheduled_notifications_claim'); + }); + + it('restores the schema when a migration is re-applied after a revert', async () => { + await harness.runner.rollback(); + + // 001 creates indexes of its own, so assert on the ones 002 owns. + const afterRevert = await indexNames(harness.sqliteDb); + expect(afterRevert).not.toContain('idx_scheduled_notifications_claim'); + expect(afterRevert).not.toContain('idx_processed_events_tx_hash'); + + await harness.runner.runMigrations(); + + const applied = await harness.runner.getAppliedMigrations(); + expect(applied).toEqual(['001', '002']); + const indexes = await indexNames(harness.sqliteDb); + expect(indexes).toContain('idx_scheduled_notifications_claim'); + }); + + it('unwinds multiple migrations in reverse order', async () => { + const reverted = await harness.runner.rollback({ steps: 2, allowDestructive: true }); + + expect(reverted.map((m) => m.id)).toEqual(['002', '001']); + expect(await tableNames(harness.sqliteDb)).toEqual([]); + }); + + it('refuses to revert a destructive migration by default', async () => { + await expect(harness.runner.rollback({ steps: 2 })).rejects.toBeInstanceOf( + DestructiveMigrationError, + ); + + // Nothing may have been reverted: the safeguard must fail closed. + const applied = await harness.runner.getAppliedMigrations(); + expect(applied).toEqual(['001', '002']); + expect((await tableNames(harness.sqliteDb)).length).toBe(1); + }); + + it('names the destructive migration in the error so the operator knows what is at risk', async () => { + await expect(harness.runner.rollback({ steps: 2 })).rejects.toThrow(/001.*initial-schema/); + }); + + it('rolls back a reversible migration even while a destructive one is queued behind it', async () => { + // 002 is reversible, so it reverts without the destructive opt-in. + const reverted = await harness.runner.rollback(); + expect(reverted.map((m) => m.id)).toEqual(['002']); + }); + + it('rejects a non-positive step count', async () => { + await expect(harness.runner.rollback({ steps: 0 })).rejects.toThrow(/positive integer/); + }); + + it('rejects a fractional step count', async () => { + await expect(harness.runner.rollback({ steps: 1.5 })).rejects.toThrow(/positive integer/); + }); + + it('is a no-op when there is nothing left to revert', async () => { + // Unwind everything, then confirm a further rollback changes nothing. + await harness.runner.rollback({ steps: 2, allowDestructive: true }); + expect(await harness.runner.getAppliedMigrations()).toEqual([]); + + const reverted = await harness.runner.rollback({ allowDestructive: true }); + expect(reverted).toEqual([]); + }); + + it('leaves the schema unchanged when a revert step fails', async () => { + // A revert can only be exercised through a migration that exists on disk, + // so this case uses a throwaway migrations directory whose newest migration + // throws from down(). + const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'notify-migrations-')); + const write = (id: string, name: string, body: string) => + fs.writeFileSync( + path.join(tempDir, `${id}-${name}.ts`), + `const migration = { + id: '${id}', + name: '${name}', + up: async (db: any) => { ${body} }, + down: async (db: any) => { ${id === '002' ? "throw new Error('boom');" : 'await db.run("DROP TABLE IF EXISTS t_a");'} }, +}; +export default migration; +`, + ); + + write('001', 'create-a', 'await db.run("CREATE TABLE IF NOT EXISTS t_a (id INTEGER)");'); + write('002', 'create-b', 'await db.run("CREATE TABLE IF NOT EXISTS t_b (id INTEGER)");'); + + const localPath = path.join(os.tmpdir(), `notify-rollback-fail-${Date.now()}.db`); + const db = new Database(localPath); + await db.initialize(); + // @ts-ignore - Accessing private db property + const runner = new MigrationRunner(db['db'] as any, tempDir); + + try { + await runner.runMigrations(); + + const tablesBefore = await countTables(db); + await expect(runner.rollback()).rejects.toThrow('boom'); + + // The failed step must not have consumed the pending revert. + expect(await countTables(db)).toEqual(tablesBefore); + expect(await runner.getAppliedMigrations()).toEqual(['001', '002']); + } finally { + await db.close(); + fs.rmSync(tempDir, { recursive: true, force: true }); + if (fs.existsSync(localPath)) fs.unlinkSync(localPath); + } + }); + + it('records when a migration was reverted', async () => { + await harness.runner.rollback(); + + const rows = await harness.db.all<{ id: string; reverted_at: string | null }>( + 'SELECT id, reverted_at FROM migrations WHERE id = ?', + ['002'], + ); + + expect(rows).toHaveLength(1); + expect(rows[0].reverted_at).not.toBeNull(); + }); + + it('lists rollback candidates newest first with a destructive marker', async () => { + const candidates = await harness.runner.getRollbackCandidates(); + + expect(candidates.map((m) => m.id)).toEqual(['002', '001']); + expect(candidates[0].destructive).toBe(false); + expect(candidates[1].destructive).toBe(true); + }); +}); diff --git a/listener/src/database/migration-system.ts b/listener/src/database/migration-system.ts index 8cc12e06..288bc5a2 100644 --- a/listener/src/database/migration-system.ts +++ b/listener/src/database/migration-system.ts @@ -1,20 +1,23 @@ /** * Migration System - * + * * A custom SQLite database migration system for the Notify-Chain listener. - * + * * Features: * - Tracks applied migrations in a `migrations` table * - Atomic migrations using transactions (rolls back on failure) * - Loads migrations from a specified directory * - Applies pending migrations in order - * + * - Reverts applied migrations in reverse order, guarding destructive ones + * * Migration file structure: * Each migration should export a default object with: * - id: Unique identifier (e.g., "001") * - name: Human-readable name (e.g., "initial-schema") * - up(db): Function to apply the migration - * - down(db): Function to roll back the migration (optional) + * - down(db): Function to roll back the migration (required for rollback support) + * - destructive: Set true when `down` discards data that cannot be reconstructed + * by re-running `up` (drops columns/tables holding rows, truncates, etc.) */ import * as sqlite3 from 'sqlite3'; import * as fs from 'fs'; @@ -26,6 +29,38 @@ export interface Migration { name: string; up: (db: sqlite3.Database) => Promise; down: (db: sqlite3.Database) => Promise; + /** + * When true, reverting this migration destroys data that `up` cannot restore. + * Destructive reverts are refused unless the caller explicitly opts in, so a + * routine `rollback-last` can never silently drop production rows. + */ + destructive?: boolean; +} + +/** Options controlling a single revert operation. */ +export interface RollbackOptions { + /** + * Permit reverting migrations flagged `destructive`. Required for any revert + * that would otherwise discard unrecoverable data. + */ + allowDestructive?: boolean; + /** Maximum number of migrations to revert. Defaults to 1. */ + steps?: number; +} + +export class DestructiveMigrationError extends Error { + public readonly migrationId: string; + public readonly migrationName: string; + + constructor(id: string, name: string) { + super( + `Refusing to roll back destructive migration ${id} (${name}). ` + + 'Re-run with allowDestructive to confirm the data loss is intended.', + ); + this.name = 'DestructiveMigrationError'; + this.migrationId = id; + this.migrationName = name; + } } export class MigrationRunner { @@ -37,40 +72,121 @@ export class MigrationRunner { this.migrationsDir = migrationsDir; } + /** + * Promisified `run`. + * + * `sqlite3` only resolves when a callback is supplied: calling `db.run(sql)` + * without one returns the Database handle for chaining and does NOT wait for + * the statement to finish. Awaiting it therefore races the next statement, + * which is why migration sequencing cannot rely on the raw handle. + */ + private run(sql: string, params: unknown[] = []): Promise { + return new Promise((resolve, reject) => { + this.db.run(sql, params, (error: Error | null) => (error ? reject(error) : resolve())); + }); + } + + /** Promisified `all`. See {@link run} for why the callback is required. */ + private all(sql: string, params: unknown[] = []): Promise { + return new Promise((resolve, reject) => { + this.db.all(sql, params, (error: Error | null, rows: unknown[]) => + error ? reject(error) : resolve(rows as T[]), + ); + }); + } + + /** + * Runs a migration's `up`/`down` with real awaiting semantics. + * + * Migrations call `await db.run(...)`. Against the raw sqlite3 handle that + * await is a no-op, so every DDL statement in a migration was fire-and-forget + * and could still be in flight when the surrounding transaction committed, + * which is how a migration could half-apply. This proxy turns each call into + * a real promise, so `up`/`down` become a genuinely atomic unit. + */ + private async runMigrationStep(step: (db: sqlite3.Database) => Promise): Promise { + const db = this.db; + + const bridged = new Proxy( + {}, + { + get(_target, property) { + if (property === 'run' || property === 'all' || property === 'exec') { + return (...args: unknown[]) => + new Promise((resolve, reject) => { + const done = (error: Error | null) => (error ? reject(error) : resolve(undefined)); + const method = db[property as 'run' | 'all' | 'exec'] as unknown as ( + ...inner: unknown[] + ) => unknown; + method.apply(db, [...args, done]); + }); + } + return undefined; + }, + }, + ); + + await step(bridged as unknown as sqlite3.Database); + } + async initializeMigrationTable(): Promise { - await this.db.run(` + await this.run(` CREATE TABLE IF NOT EXISTS migrations ( id TEXT PRIMARY KEY, name TEXT NOT NULL, applied_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ) `); + + // Applied DBs predate this column, so add it defensively rather than + // assuming a fresh database. + const columns = await this.all<{ name: string }>('PRAGMA table_info(migrations)'); + if (!columns.some((column) => column.name === 'reverted_at')) { + await this.run('ALTER TABLE migrations ADD COLUMN reverted_at DATETIME'); + } } async getAppliedMigrations(): Promise { - const rows = await this.db.all<{ id: string }>( - 'SELECT id FROM migrations ORDER BY applied_at' + const rows = await this.all<{ id: string }>( + 'SELECT id FROM migrations WHERE reverted_at IS NULL ORDER BY applied_at', ); return rows.map((row) => row.id); } async applyMigration(migration: Migration): Promise { - await this.db.serialize(async () => { - await this.db.run('BEGIN TRANSACTION'); - try { - await migration.up(this.db); - await this.db.run( - 'INSERT INTO migrations (id, name) VALUES (?, ?)', - [migration.id, migration.name] - ); - await this.db.run('COMMIT'); - logger.info(`Migration ${migration.id} (${migration.name}) applied successfully`); - } catch (error) { - await this.db.run('ROLLBACK'); - logger.error(`Migration ${migration.id} failed, rolling back:`, error); - throw error; - } - }); + // NOTE: deliberately not wrapped in db.serialize(). serialize() invokes its + // callback synchronously and does not await an async one, so the runner + // resolved before the migration finished and callers could close the + // database mid-transaction. Statements on one connection already execute + // in the order they are queued, which the explicit awaits below guarantee. + await this.run('BEGIN TRANSACTION'); + try { + await this.runMigrationStep(migration.up); + // Upsert rather than plain INSERT: `id` is the primary key and a revert is + // recorded in place to preserve the audit trail, so a previously reverted + // migration still occupies its row. Re-applying must revive that row + // instead of colliding with it. + await this.run( + `INSERT INTO migrations (id, name, reverted_at) VALUES (?, ?, NULL) + ON CONFLICT(id) DO UPDATE SET + name = excluded.name, + applied_at = CURRENT_TIMESTAMP, + reverted_at = NULL`, + [migration.id, migration.name], + ); + await this.run('COMMIT'); + logger.info(`Migration ${migration.id} (${migration.name}) applied successfully`); + } catch (error) { + // A failed ROLLBACK must not replace the original error, which is the one + // that explains the failure. + await this.run('ROLLBACK').catch((rollbackError: Error) => { + logger.error(`Could not roll back transaction for ${migration.id}:`, { + error: rollbackError, + }); + }); + logger.error(`Migration ${migration.id} failed, rolling back:`, { error: error as Error }); + throw error; + } } async loadMigrations(): Promise { @@ -93,9 +209,7 @@ export class MigrationRunner { async getPendingMigrations(): Promise { const appliedMigrations = await this.getAppliedMigrations(); const allMigrations = await this.loadMigrations(); - return allMigrations.filter( - (migration) => !appliedMigrations.includes(migration.id) - ); + return allMigrations.filter((migration) => !appliedMigrations.includes(migration.id)); } async runMigrations(): Promise { @@ -113,4 +227,96 @@ export class MigrationRunner { } logger.info('All migrations applied successfully'); } + + /** + * Revert the most recently applied migrations, newest first. + * + * Each revert runs inside its own transaction and the `migrations` row is + * marked rather than deleted, so the audit trail of what was applied and + * when it was reverted is preserved. + * + * @param options Revert options. `steps` defaults to 1. + * @returns The migrations that were reverted, newest first. + */ + async rollback(options: RollbackOptions = {}): Promise { + const steps = options.steps ?? 1; + if (!Number.isInteger(steps) || steps < 1) { + throw new Error(`Rollback steps must be a positive integer, received ${steps}`); + } + + await this.initializeMigrationTable(); + const allMigrations = await this.loadMigrations(); + const appliedIds = await this.getAppliedMigrations(); + + // Newest first: reverts must unwind in the opposite order to `up`. + const targets = [...appliedIds].reverse().slice(0, steps); + if (targets.length === 0) { + logger.info('No applied migrations to roll back'); + return []; + } + + // Validate every target before mutating anything, so a refusal halfway + // through cannot leave the database partially reverted. + const migrations: Migration[] = []; + for (const id of targets) { + const migration = allMigrations.find((candidate) => candidate.id === id); + if (!migration) { + throw new Error( + `Cannot roll back migration ${id}: no migration file defines it. ` + + 'The migration directory and the migrations table are out of sync.', + ); + } + if (typeof migration.down !== 'function') { + throw new Error(`Migration ${id} (${migration.name}) does not export a down() function`); + } + if (migration.destructive && !options.allowDestructive) { + throw new DestructiveMigrationError(migration.id, migration.name); + } + migrations.push(migration); + } + + const reverted: Migration[] = []; + for (const migration of migrations) { + await this.revertMigration(migration); + reverted.push(migration); + } + return reverted; + } + + private async revertMigration(migration: Migration): Promise { + await this.run('BEGIN TRANSACTION'); + try { + await this.runMigrationStep(migration.down); + await this.run('UPDATE migrations SET reverted_at = CURRENT_TIMESTAMP WHERE id = ?', [ + migration.id, + ]); + await this.run('COMMIT'); + logger.warn( + `Migration ${migration.id} (${migration.name}) rolled back` + + (migration.destructive ? ' [destructive]' : ''), + ); + } catch (error) { + await this.run('ROLLBACK').catch((rollbackError: Error) => { + logger.error(`Could not roll back transaction for ${migration.id}:`, { + error: rollbackError, + }); + }); + logger.error(`Rollback of migration ${migration.id} failed, transaction reverted:`, { + error: error as Error, + }); + throw error; + } + } + + /** + * Migrations that are applied and not yet reverted, newest first. + */ + async getRollbackCandidates(): Promise { + const allMigrations = await this.loadMigrations(); + const appliedIds = new Set(await this.getAppliedMigrations()); + return [...appliedIds] + .reverse() + .map((id) => allMigrations.find((m) => m.id === id)) + .filter((migration): migration is Migration => migration !== undefined); + } } diff --git a/listener/src/migrations/001-initial-schema.ts b/listener/src/migrations/001-initial-schema.ts index 5e127a6b..71f64bba 100644 --- a/listener/src/migrations/001-initial-schema.ts +++ b/listener/src/migrations/001-initial-schema.ts @@ -3,6 +3,9 @@ import * as sqlite3 from 'sqlite3'; const migration = { id: '001', name: 'initial-schema', + // `down` drops every listener table. Re-applying `up` recreates empty + // tables, so any notification, template or rate-limit rows are lost for good. + destructive: true, up: async (db: sqlite3.Database) => { const schemaSql = ` -- Main table for scheduled notifications @@ -252,10 +255,11 @@ const migration = { ON notification_metrics_snapshots(captured_at); `; - const statements = schemaSql.split(';').map(s => s.trim()).filter(s => s); - for (const statement of statements) { - await db.run(statement); - } + // The schema contains CREATE TRIGGER ... BEGIN ... END bodies, which + // legitimately contain semicolons of their own. Splitting the script on ';' + // cut those bodies in half and made every trigger fail to parse, so the + // whole script is handed to exec(), which understands statement boundaries. + await db.exec(schemaSql); }, down: async (db: sqlite3.Database) => { await db.run('DROP TABLE IF EXISTS notification_metrics_snapshots'); @@ -269,7 +273,7 @@ const migration = { await db.run('DROP TRIGGER IF EXISTS update_scheduled_notifications_timestamp'); await db.run('DROP TABLE IF EXISTS notification_execution_log'); await db.run('DROP TABLE IF EXISTS scheduled_notifications'); - } + }, }; export default migration; diff --git a/listener/src/migrations/002-query-performance-indexes.ts b/listener/src/migrations/002-query-performance-indexes.ts index c4763695..e99fd52d 100644 --- a/listener/src/migrations/002-query-performance-indexes.ts +++ b/listener/src/migrations/002-query-performance-indexes.ts @@ -12,6 +12,8 @@ import * as sqlite3 from 'sqlite3'; const migration = { id: '002', name: 'query-performance-indexes', + // Reverting only drops indexes, which `up` recreates without touching rows. + destructive: false, up: async (db: sqlite3.Database) => { await db.run(` CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_claim diff --git a/listener/src/scripts/check-migration-safety.ts b/listener/src/scripts/check-migration-safety.ts new file mode 100644 index 00000000..8b28cbb5 --- /dev/null +++ b/listener/src/scripts/check-migration-safety.ts @@ -0,0 +1,111 @@ +#!/usr/bin/env ts-node +/** + * Audits the migration directory for rollback safety. + * + * Two rules are enforced, both of which are far cheaper to catch here than in + * production: + * + * 1. Every migration must export a `down()`. A migration without one cannot be + * rolled back, so a bad deploy has no way out. + * 2. Any migration whose `down()` performs a destructive statement must be + * flagged `destructive: true`. The flag is what makes `rollback()` refuse + * the revert unless the operator explicitly opts in. + * + * Destructive statements detected: DROP TABLE, DROP COLUMN, and TRUNCATE. + * + * Usage: + * npm run check-migration-safety + */ +import * as fs from 'fs'; +import * as path from 'path'; + +const MIGRATIONS_DIR = path.join(__dirname, '../migrations'); + +/** Statements that discard data `up()` cannot reconstruct. */ +const DESTRUCTIVE_PATTERNS: Array<{ pattern: RegExp; label: string }> = [ + { pattern: /DROP\s+TABLE/i, label: 'DROP TABLE' }, + { pattern: /DROP\s+COLUMN/i, label: 'DROP COLUMN' }, + { pattern: /TRUNCATE/i, label: 'TRUNCATE' }, +]; + +interface AuditFinding { + file: string; + message: string; +} + +async function audit(): Promise { + const findings: AuditFinding[] = []; + + if (!fs.existsSync(MIGRATIONS_DIR)) { + findings.push({ file: MIGRATIONS_DIR, message: 'Migrations directory does not exist' }); + return findings; + } + + const files = fs + .readdirSync(MIGRATIONS_DIR) + .filter((file) => file.endsWith('.ts') || file.endsWith('.js')) + .sort(); + + if (files.length === 0) { + findings.push({ file: MIGRATIONS_DIR, message: 'No migrations found' }); + return findings; + } + + for (const file of files) { + const migrationPath = path.join(MIGRATIONS_DIR, file); + // eslint-disable-next-line @typescript-eslint/no-var-requires + const module = await import(migrationPath); + const migration = module.default; + + if (!migration || typeof migration.up !== 'function') { + findings.push({ file, message: 'Does not export a default object with an up() function' }); + continue; + } + + if (typeof migration.down !== 'function') { + findings.push({ + file, + message: 'Does not export a down() function, so it cannot be rolled back', + }); + } + + // Rule 2 only applies once a down() exists to inspect. + if (typeof migration.down === 'function') { + const source = fs.readFileSync(migrationPath, 'utf8'); + const matched = DESTRUCTIVE_PATTERNS.filter(({ pattern }) => pattern.test(source)).map( + ({ label }) => label, + ); + + if (matched.length > 0 && migration.destructive !== true) { + findings.push({ + file, + message: + `down() uses ${matched.join(', ')} but the migration is not flagged ` + + '`destructive: true`, so rollback() would silently discard data', + }); + } + } + } + + return findings; +} + +audit() + .then((findings) => { + if (findings.length === 0) { + console.log('✅ All migrations define a rollback path and flag destructive reverts'); + process.exit(0); + } + + console.error('❌ Migration rollback safety issues found:\n'); + for (const finding of findings) { + console.error(` ${finding.file}: ${finding.message}`); + } + console.error('\nFix each migration by adding a down() function, and set'); + console.error('`destructive: true` when reverting it would discard unrecoverable data.'); + process.exit(1); + }) + .catch((error) => { + console.error('❌ Migration safety check failed to run:', error); + process.exit(2); + }); diff --git a/listener/src/scripts/rollback-db.ts b/listener/src/scripts/rollback-db.ts new file mode 100644 index 00000000..a5734c43 --- /dev/null +++ b/listener/src/scripts/rollback-db.ts @@ -0,0 +1,102 @@ +#!/usr/bin/env ts-node +/** + * Rollback CLI for the listener database. + * + * Reverts applied migrations in reverse order, unwinding as many steps as + * requested. Destructive migrations are refused unless --allow-destructive is + * passed, so a routine rollback can never quietly drop production data. + * + * Usage: + * npm run rollback # revert the most recent migration + * npm run rollback -- --steps 3 # revert the three most recent + * npm run rollback -- --status # list what would be reverted + * npm run rollback -- --allow-destructive + * + * Environment: + * DATABASE_PATH overrides the database location + */ +import { Database } from '../database/database'; +import { DestructiveMigrationError, MigrationRunner } from '../database/migration-system'; +import logger from '../utils/logger'; +import * as dotenv from 'dotenv'; +import * as path from 'path'; + +dotenv.config(); + +function parseSteps(argv: string[]): number { + const index = argv.indexOf('--steps'); + if (index === -1) return 1; + + const value = argv[index + 1]; + if (!value) { + throw new Error('--steps requires a value, e.g. --steps 3'); + } + + const steps = Number.parseInt(value, 10); + if (!Number.isInteger(steps) || steps < 1) { + throw new Error(`--steps must be a positive integer, received "${value}"`); + } + return steps; +} + +async function rollback() { + const argv = process.argv.slice(2); + const dbPath = process.env.DATABASE_PATH || './data/notifications.db'; + const migrationsDir = path.join(__dirname, '../migrations'); + const allowDestructive = argv.includes('--allow-destructive'); + + const db = new Database(dbPath); + await db.initialize(); + + // @ts-ignore - Accessing private db property + const sqliteDb = db['db'] as any; + const runner = new MigrationRunner(sqliteDb, migrationsDir); + + try { + // `--status` has to work before anything has ever been applied, so ensure + // the bookkeeping table exists rather than assuming a migrated database. + await runner.initializeMigrationTable(); + + if (argv.includes('--status')) { + const candidates = await runner.getRollbackCandidates(); + if (candidates.length === 0) { + console.log('No applied migrations to roll back'); + return; + } + console.log('Applied migrations, newest first (next to be reverted):'); + for (const migration of candidates) { + const marker = migration.destructive ? 'destructive' : 'reversible'; + console.log(` ${migration.id} ${migration.name} [${marker}]`); + } + return; + } + + const steps = parseSteps(argv); + const reverted = await runner.rollback({ steps, allowDestructive }); + + if (reverted.length === 0) { + console.log('Nothing to roll back'); + return; + } + + console.log(`Rolled back ${reverted.length} migration(s):`); + for (const migration of reverted) { + console.log(` ${migration.id} ${migration.name}`); + } + } catch (error) { + if (error instanceof DestructiveMigrationError) { + console.error(`❌ ${error.message}`); + console.error('\nThis migration drops data that cannot be recovered by re-applying it.'); + console.error('Re-run with --allow-destructive only once you are certain.'); + process.exit(3); + } + throw error; + } finally { + await db.close(); + } +} + +rollback().catch((error) => { + logger.error('Rollback failed:', error); + process.exit(1); +});