From ce109c6aa9566d958d708d2eba3744beef4309e9 Mon Sep 17 00:00:00 2001 From: Najib Ishiyaku Njidda Date: Tue, 29 Sep 2026 18:03:02 +0000 Subject: [PATCH 1/5] feat(listener): enforce database integrity constraints on core tables. Closes #844. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CHECK constraints on all closed status/state enums (scheduled notifications, execution attempts, processed events, idempotency, backpressure, rate-limit client types, notification archive); FK actions preserved deliberately (CASCADE for cleanup children, RESTRICT for template audit); PRAGMA foreign_keys verified at connect and in the migration runner. Migration 004 rebuilds tables in a single idempotent transaction with preflight audits that abort fail-closed on invalid legacy rows; valid legacy data passes through byte-identical (legacy-shape test asserts survival, clean foreign_key_check, and orphan/status/duplicate rejection). Constraints mirrored in schema.sql and archive-schema.sql for fresh-install parity. Also fixes a latent migration-runner defect: callback-based sqlite3 run/all were awaited as promises, undermining transaction/rollback guarantees. Deliberately NOT added: UNIQUE(notification_id, attempt) — the dead-letter retry path resets retry_count, so attempts legitimately repeat. --- NOTIFICATION_LIFECYCLE.md | 20 + listener/src/database/archive-schema.sql | 3 + listener/src/database/migration-system.ts | 75 ++- listener/src/database/schema.sql | 85 ++-- .../004-database-data-integrity.test.ts | 282 +++++++++++ .../migrations/004-database-data-integrity.ts | 446 ++++++++++++++++++ 6 files changed, 864 insertions(+), 47 deletions(-) create mode 100644 listener/src/migrations/004-database-data-integrity.test.ts create mode 100644 listener/src/migrations/004-database-data-integrity.ts diff --git a/NOTIFICATION_LIFECYCLE.md b/NOTIFICATION_LIFECYCLE.md index a11e2367..d5cd3688 100644 --- a/NOTIFICATION_LIFECYCLE.md +++ b/NOTIFICATION_LIFECYCLE.md @@ -393,6 +393,26 @@ before returning failure to the caller. Details: [NOTIFICATION_FAILURE_RECOVERY.md](NOTIFICATION_FAILURE_RECOVERY.md). +## Database Integrity Constraints + +The active SQLite schema rejects unknown notification types/statuses, event +states, idempotency states, backpressure event types, and invalid counters or +boolean values. Notification execution logs, dead-letter entries, and +idempotency records reference their scheduled notification with `ON DELETE +CASCADE`; template audit records reference their template with `ON DELETE +RESTRICT`. Processed-event fingerprints, cursor contract addresses, +idempotency keys, and dead-letter notification references retain their unique +keys. + +Fresh databases receive these constraints from +`listener/src/database/schema.sql`. Run `npm run migrate` to apply migration +004 to an existing database. The migration audits legacy rows before rebuilding +tables in a transaction. If it finds invalid state, an orphaned reference, or a +duplicate unique key, it aborts and reports the affected table and row count; +it never deletes or rewrites invalid data automatically. Repair the listed rows +and rerun the migration. SQLite foreign-key checks are enabled on application +database connections and verified by the migration runner. + --- ## Completion and Archival diff --git a/listener/src/database/archive-schema.sql b/listener/src/database/archive-schema.sql index 546e1095..e961b14c 100644 --- a/listener/src/database/archive-schema.sql +++ b/listener/src/database/archive-schema.sql @@ -1,3 +1,6 @@ + notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), + status VARCHAR(20) NOT NULL CHECK (status IN ('COMPLETED', 'FAILED', 'CANCELLED')), + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), -- Archive table for notifications moved out of active storage. -- Records here are read-only for audit purposes and are never modified. CREATE TABLE IF NOT EXISTS notification_archive ( diff --git a/listener/src/database/migration-system.ts b/listener/src/database/migration-system.ts index 8cc12e06..66a58e0f 100644 --- a/listener/src/database/migration-system.ts +++ b/listener/src/database/migration-system.ts @@ -37,8 +37,49 @@ export class MigrationRunner { this.migrationsDir = migrationsDir; } + private run(sql: string, params: unknown[] = []): Promise { + return new Promise((resolve, reject) => { + this.db.run(sql, params, (error) => (error ? reject(error) : resolve())); + }); + } + + private all(sql: string): Promise { + return new Promise((resolve, reject) => { + this.db.all(sql, (error, rows) => (error ? reject(error) : resolve(rows as T[]))); + }); + } + + private async ensureForeignKeysEnabled(): Promise { + const current = await new Promise((resolve, reject) => { + this.db.get('PRAGMA foreign_keys', (error, row: { foreign_keys: number }) => { + if (error) reject(error); + else resolve(row.foreign_keys); + }); + }); + + if (current !== 1) { + await new Promise((resolve, reject) => { + this.db.run('PRAGMA foreign_keys = ON', (error) => { + if (error) reject(error); + else resolve(); + }); + }); + } + + const verified = await new Promise((resolve, reject) => { + this.db.get('PRAGMA foreign_keys', (error, row: { foreign_keys: number }) => { + if (error) reject(error); + else resolve(row.foreign_keys); + }); + }); + if (verified !== 1) { + throw new Error('SQLite foreign key enforcement could not be enabled'); + } + } + async initializeMigrationTable(): Promise { - await this.db.run(` + await this.ensureForeignKeysEnabled(); + await this.run(` CREATE TABLE IF NOT EXISTS migrations ( id TEXT PRIMARY KEY, name TEXT NOT NULL, @@ -48,29 +89,27 @@ export class MigrationRunner { } 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 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'); + await this.ensureForeignKeysEnabled(); + await this.run('BEGIN TRANSACTION'); + try { + await migration.up(this.db); + await this.run('INSERT INTO migrations (id, name) VALUES (?, ?)', [migration.id, migration.name]); + await this.run('COMMIT'); + logger.info(`Migration ${migration.id} (${migration.name}) applied successfully`); + } catch (error) { 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; + await this.run('ROLLBACK'); + } catch (rollbackError) { + logger.error(`Migration ${migration.id} rollback failed`, { error: rollbackError }); } - }); + logger.error(`Migration ${migration.id} failed, rolling back`, { error }); + throw error; + } } async loadMigrations(): Promise { diff --git a/listener/src/database/schema.sql b/listener/src/database/schema.sql index 8342d5b1..cb867162 100644 --- a/listener/src/database/schema.sql +++ b/listener/src/database/schema.sql @@ -8,7 +8,7 @@ CREATE TABLE IF NOT EXISTS scheduled_notifications ( -- Notification content and metadata payload TEXT NOT NULL, -- JSON payload of the notification (compressed when large) payload_hash TEXT, -- HMAC hash of the raw JSON payload for integrity checks - notification_type VARCHAR(50) NOT NULL, -- Type: 'discord', 'email', 'webhook', etc. + notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), target_recipient TEXT NOT NULL, -- User ID, webhook URL, or recipient identifier -- Scheduling information @@ -17,9 +17,9 @@ CREATE TABLE IF NOT EXISTS scheduled_notifications ( updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, -- Status tracking - status VARCHAR(20) NOT NULL DEFAULT 'PENDING', -- PENDING, PROCESSING, COMPLETED, FAILED, CANCELLED - retry_count INTEGER NOT NULL DEFAULT 0, - max_retries INTEGER NOT NULL DEFAULT 3, + status VARCHAR(20) NOT NULL DEFAULT 'PENDING' CHECK (status IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'CANCELLED')), + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), + max_retries INTEGER NOT NULL DEFAULT 3 CHECK (max_retries >= 0), -- Processing metadata processing_started_at DATETIME, @@ -34,7 +34,7 @@ CREATE TABLE IF NOT EXISTS scheduled_notifications ( -- Additional metadata event_id TEXT, -- Reference to the original event (if applicable) contract_address TEXT, -- Stellar contract address (if applicable) - priority INTEGER NOT NULL DEFAULT 5, -- 1-10, lower = higher priority + priority INTEGER NOT NULL DEFAULT 5 CHECK (priority BETWEEN 1 AND 10), metadata TEXT, -- Additional JSON metadata next_retry_at DATETIME -- When the next retry should be attempted ); @@ -71,14 +71,14 @@ CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_target CREATE TABLE IF NOT EXISTS dead_letter_queue ( id INTEGER PRIMARY KEY AUTOINCREMENT, scheduled_notification_id INTEGER NOT NULL UNIQUE, - notification_type VARCHAR(50) NOT NULL, + notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), target_recipient TEXT NOT NULL, payload TEXT NOT NULL, failure_reason TEXT NOT NULL, error_details TEXT, created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, last_retried_at DATETIME, - retry_count INTEGER NOT NULL DEFAULT 0, + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), FOREIGN KEY (scheduled_notification_id) REFERENCES scheduled_notifications(id) ON DELETE CASCADE ); @@ -92,12 +92,12 @@ CREATE INDEX IF NOT EXISTS idx_dead_letter_queue_notification_type CREATE TABLE IF NOT EXISTS notification_execution_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, scheduled_notification_id INTEGER NOT NULL, - execution_attempt INTEGER NOT NULL, + execution_attempt INTEGER NOT NULL CHECK (execution_attempt > 0), execution_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, - status VARCHAR(20) NOT NULL, -- SUCCESS, FAILED, RETRY + status VARCHAR(20) NOT NULL CHECK (status IN ('SUCCESS', 'FAILED', 'RETRY')), error_message TEXT, response_data TEXT, -- JSON response from notification service - duration_ms INTEGER, + duration_ms INTEGER CHECK (duration_ms IS NULL OR duration_ms >= 0), FOREIGN KEY (scheduled_notification_id) REFERENCES scheduled_notifications(id) ON DELETE CASCADE ); @@ -127,12 +127,12 @@ END; CREATE TABLE IF NOT EXISTS rate_limit_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, client_id TEXT NOT NULL, -- IP address or API key - client_type VARCHAR(20) NOT NULL, -- 'IP' or 'API_KEY' + client_type VARCHAR(20) NOT NULL CHECK (client_type IN ('IP', 'API_KEY')), endpoint TEXT NOT NULL, -- Request path/method method VARCHAR(10) NOT NULL, -- Request method timestamp DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, - limit_threshold INTEGER NOT NULL, - window_ms INTEGER NOT NULL + limit_threshold INTEGER NOT NULL CHECK (limit_threshold > 0), + window_ms INTEGER NOT NULL CHECK (window_ms > 0) ); CREATE INDEX IF NOT EXISTS idx_rate_limit_events_timestamp @@ -171,7 +171,7 @@ CREATE TABLE IF NOT EXISTS notification_template_audit_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, template_id TEXT NOT NULL, actor TEXT NOT NULL, - action TEXT NOT NULL DEFAULT 'UPDATE', + action TEXT NOT NULL DEFAULT 'UPDATE' CHECK (action IN ('UPDATE')), changed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, previous_snapshot TEXT NOT NULL, new_snapshot TEXT NOT NULL, @@ -209,19 +209,19 @@ CREATE TABLE IF NOT EXISTS processed_events ( fingerprint TEXT NOT NULL UNIQUE, -- Composite key: contract_address:event_id (for faster lookups) -- Processing metadata - ledger_number INTEGER NOT NULL, -- Ledger in which the event occurred + ledger_number INTEGER NOT NULL CHECK (ledger_number >= 0), tx_hash TEXT, -- Transaction hash (if available) event_type VARCHAR(50) NOT NULL, -- Type from RPC (contract, system, etc) processed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, -- Reorg detection and tracking - is_reorg_duplicate BOOLEAN NOT NULL DEFAULT 0, -- Flag indicating this is a duplicate from a reorg - reorg_detection_count INTEGER NOT NULL DEFAULT 0, -- Number of times this event was redetected + is_reorg_duplicate BOOLEAN NOT NULL DEFAULT 0 CHECK (is_reorg_duplicate IN (0, 1)), + reorg_detection_count INTEGER NOT NULL DEFAULT 0 CHECK (reorg_detection_count >= 0), last_redetected_at DATETIME, -- When the event was last detected again (for reorg monitoring) -- Status and metadata - status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED', -- PROCESSED, SKIPPED, ERROR - notification_sent BOOLEAN NOT NULL DEFAULT 0, -- Whether a notification was sent for this event + status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED' CHECK (status IN ('PROCESSED', 'SKIPPED', 'ERROR')), + notification_sent BOOLEAN NOT NULL DEFAULT 0 CHECK (notification_sent IN (0, 1)), error_reason TEXT -- If status is ERROR, what went wrong ); @@ -251,12 +251,12 @@ CREATE TABLE IF NOT EXISTS polling_cursors ( -- Cursor information cursor TEXT NOT NULL, -- Last known cursor from RPC - ledger_number INTEGER NOT NULL, -- Ledger number associated with this cursor + ledger_number INTEGER NOT NULL CHECK (ledger_number >= 0), updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, -- Reorg detection - reorg_detected BOOLEAN NOT NULL DEFAULT 0, -- Whether a reorg was detected on the last poll - reorg_detection_count INTEGER NOT NULL DEFAULT 0 -- Total number of reorgs detected for this contract + reorg_detected BOOLEAN NOT NULL DEFAULT 0 CHECK (reorg_detected IN (0, 1)), + reorg_detection_count INTEGER NOT NULL DEFAULT 0 CHECK (reorg_detection_count >= 0) ); CREATE INDEX IF NOT EXISTS idx_polling_cursors_contract @@ -282,7 +282,7 @@ CREATE TABLE IF NOT EXISTS idempotency_keys ( expires_at DATETIME NOT NULL, -- When this key should be purged -- Status tracking - status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED', -- PROCESSED, EXPIRED + status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED' CHECK (status IN ('PROCESSED', 'EXPIRED')), FOREIGN KEY (response_notification_id) REFERENCES scheduled_notifications(id) ON DELETE CASCADE ); @@ -301,12 +301,12 @@ CREATE TABLE IF NOT EXISTS backpressure_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, -- Event tracking - event_type VARCHAR(20) NOT NULL, -- ACTIVATED or DEACTIVATED - queue_size INTEGER NOT NULL, -- Queue size when event occurred - target_throughput_per_sec INTEGER NOT NULL, -- Target throughput limit during this event + event_type VARCHAR(20) NOT NULL CHECK (event_type IN ('ACTIVATED', 'DEACTIVATED')), + queue_size INTEGER NOT NULL CHECK (queue_size >= 0), + target_throughput_per_sec INTEGER NOT NULL CHECK (target_throughput_per_sec >= 0), -- Duration tracking (for deactivation events) - duration_ms INTEGER, -- How long backpressure was active + duration_ms INTEGER CHECK (duration_ms IS NULL OR duration_ms >= 0), -- Additional metadata reason TEXT, -- Optional reason/context for the event @@ -328,12 +328,39 @@ CREATE TABLE IF NOT EXISTS notification_metrics_snapshots ( captured_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, window_start INTEGER NOT NULL, window_end INTEGER NOT NULL, - total_recorded INTEGER NOT NULL, - snapshot_json TEXT NOT NULL + total_recorded INTEGER NOT NULL CHECK (total_recorded >= 0), + snapshot_json TEXT NOT NULL, + CHECK (window_start <= window_end) ); CREATE INDEX IF NOT EXISTS idx_metrics_snapshots_captured_at ON notification_metrics_snapshots(captured_at); +-- Archived notification rows retain only terminal delivery states +CREATE TABLE IF NOT EXISTS notification_archive ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + original_id INTEGER NOT NULL, + payload TEXT NOT NULL, + notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), + target_recipient TEXT NOT NULL, + execute_at DATETIME NOT NULL, + created_at DATETIME NOT NULL, + processing_completed_at DATETIME, + status VARCHAR(20) NOT NULL CHECK (status IN ('COMPLETED', 'FAILED', 'CANCELLED')), + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), + last_error TEXT, + event_id TEXT, + contract_address TEXT, + metadata TEXT, + archived_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_archive_original_id ON notification_archive(original_id); +CREATE INDEX IF NOT EXISTS idx_archive_archived_at ON notification_archive(archived_at); +CREATE INDEX IF NOT EXISTS idx_archive_status ON notification_archive(status); +CREATE INDEX IF NOT EXISTS idx_archive_contract_address ON notification_archive(contract_address) + WHERE contract_address IS NOT NULL; +CREATE INDEX IF NOT EXISTS idx_archive_event_id ON notification_archive(event_id) + WHERE event_id IS NOT NULL; -- =============================================== -- QUERY PERFORMANCE INDEXES (migration 002) diff --git a/listener/src/migrations/004-database-data-integrity.test.ts b/listener/src/migrations/004-database-data-integrity.test.ts new file mode 100644 index 00000000..8acb823c --- /dev/null +++ b/listener/src/migrations/004-database-data-integrity.test.ts @@ -0,0 +1,282 @@ +import * as sqlite3 from 'sqlite3'; +import { Database } from '../database/database'; +import { MigrationRunner } from '../database/migration-system'; +import migration from './004-database-data-integrity'; + +function openDatabase(): Promise { + return new Promise((resolve, reject) => { + const db = new sqlite3.Database(':memory:', (error) => { + if (error) reject(error); + else resolve(db); + }); + }); +} + +function run(db: sqlite3.Database, sql: string, params: unknown[] = []): Promise { + return new Promise((resolve, reject) => { + db.run(sql, params, (error) => (error ? reject(error) : resolve())); + }); +} + +function exec(db: sqlite3.Database, sql: string): Promise { + return new Promise((resolve, reject) => { + db.exec(sql, (error) => (error ? reject(error) : resolve())); + }); +} + +function all(db: sqlite3.Database, sql: string): Promise { + return new Promise((resolve, reject) => { + db.all(sql, (error, rows) => (error ? reject(error) : resolve(rows as T[]))); + }); +} + +function close(db: sqlite3.Database): Promise { + return new Promise((resolve, reject) => { + db.close((error) => (error ? reject(error) : resolve())); + }); +} + +async function createLegacySchema(db: sqlite3.Database): Promise { + await exec( + db, + ` + CREATE TABLE scheduled_notifications ( + id INTEGER PRIMARY KEY AUTOINCREMENT, payload TEXT NOT NULL, payload_hash TEXT, + notification_type VARCHAR(50) NOT NULL, target_recipient TEXT NOT NULL, + execute_at DATETIME NOT NULL, created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, status VARCHAR(20) NOT NULL DEFAULT 'PENDING', + retry_count INTEGER NOT NULL DEFAULT 0, max_retries INTEGER NOT NULL DEFAULT 3, + processing_started_at DATETIME, processing_completed_at DATETIME, processor_id VARCHAR(100), + lock_expires_at DATETIME, last_error TEXT, error_details TEXT, event_id TEXT, + contract_address TEXT, priority INTEGER NOT NULL DEFAULT 5, metadata TEXT, next_retry_at DATETIME + ); + CREATE TABLE notification_execution_log ( + id INTEGER PRIMARY KEY AUTOINCREMENT, scheduled_notification_id INTEGER NOT NULL, + execution_attempt INTEGER NOT NULL, execution_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + status VARCHAR(20) NOT NULL, error_message TEXT, response_data TEXT, duration_ms INTEGER, + FOREIGN KEY (scheduled_notification_id) REFERENCES scheduled_notifications(id) ON DELETE CASCADE + ); + CREATE TABLE dead_letter_queue ( + id INTEGER PRIMARY KEY AUTOINCREMENT, scheduled_notification_id INTEGER NOT NULL UNIQUE, + notification_type VARCHAR(50) NOT NULL, target_recipient TEXT NOT NULL, payload TEXT NOT NULL, + failure_reason TEXT NOT NULL, error_details TEXT, created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + last_retried_at DATETIME, retry_count INTEGER NOT NULL DEFAULT 0, + FOREIGN KEY (scheduled_notification_id) REFERENCES scheduled_notifications(id) ON DELETE CASCADE + ); + CREATE TABLE processed_events ( + id INTEGER PRIMARY KEY AUTOINCREMENT, event_id TEXT NOT NULL, contract_address TEXT NOT NULL, + fingerprint TEXT NOT NULL UNIQUE, ledger_number INTEGER NOT NULL, tx_hash TEXT, + event_type VARCHAR(50) NOT NULL, processed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + is_reorg_duplicate BOOLEAN NOT NULL DEFAULT 0, reorg_detection_count INTEGER NOT NULL DEFAULT 0, + last_redetected_at DATETIME, status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED', + notification_sent BOOLEAN NOT NULL DEFAULT 0, error_reason TEXT + ); + CREATE TABLE polling_cursors ( + id INTEGER PRIMARY KEY AUTOINCREMENT, contract_address TEXT NOT NULL UNIQUE, + cursor TEXT NOT NULL, ledger_number INTEGER NOT NULL, updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + reorg_detected BOOLEAN NOT NULL DEFAULT 0, reorg_detection_count INTEGER NOT NULL DEFAULT 0 + ); + CREATE TABLE idempotency_keys ( + id INTEGER PRIMARY KEY AUTOINCREMENT, idempotency_key TEXT NOT NULL UNIQUE, request_hash TEXT NOT NULL, + response_notification_id INTEGER NOT NULL, response_data TEXT NOT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, expires_at DATETIME NOT NULL, + status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED', + FOREIGN KEY (response_notification_id) REFERENCES scheduled_notifications(id) ON DELETE CASCADE + ); + CREATE TABLE rate_limit_events ( + id INTEGER PRIMARY KEY AUTOINCREMENT, client_id TEXT NOT NULL, client_type VARCHAR(20) NOT NULL, + endpoint TEXT NOT NULL, method VARCHAR(10) NOT NULL, timestamp DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + limit_threshold INTEGER NOT NULL, window_ms INTEGER NOT NULL + ); + CREATE TABLE notification_templates (id TEXT PRIMARY KEY); + CREATE TABLE notification_template_audit_log ( + id INTEGER PRIMARY KEY AUTOINCREMENT, template_id TEXT NOT NULL, actor TEXT NOT NULL, + action TEXT NOT NULL DEFAULT 'UPDATE', changed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + previous_snapshot TEXT NOT NULL, new_snapshot TEXT NOT NULL, + FOREIGN KEY (template_id) REFERENCES notification_templates(id) ON DELETE RESTRICT + ); + CREATE TABLE backpressure_events ( + id INTEGER PRIMARY KEY AUTOINCREMENT, event_type VARCHAR(20) NOT NULL, queue_size INTEGER NOT NULL, + target_throughput_per_sec INTEGER NOT NULL, duration_ms INTEGER, reason TEXT, + timestamp DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP + ); + CREATE TABLE notification_metrics_snapshots ( + id INTEGER PRIMARY KEY AUTOINCREMENT, captured_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + window_start INTEGER NOT NULL, window_end INTEGER NOT NULL, total_recorded INTEGER NOT NULL, + snapshot_json TEXT NOT NULL + ); + CREATE TABLE notification_archive ( + id INTEGER PRIMARY KEY AUTOINCREMENT, original_id INTEGER NOT NULL, payload TEXT NOT NULL, + notification_type VARCHAR(50) NOT NULL, target_recipient TEXT NOT NULL, + execute_at DATETIME NOT NULL, created_at DATETIME NOT NULL, processing_completed_at DATETIME, + status VARCHAR(20) NOT NULL, retry_count INTEGER NOT NULL DEFAULT 0, last_error TEXT, + event_id TEXT, contract_address TEXT, metadata TEXT, + archived_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP + ); + `, + ); +} + +async function seedValidRows(db: sqlite3.Database): Promise { + await run( + db, + `INSERT INTO scheduled_notifications + (id, payload, payload_hash, notification_type, target_recipient, execute_at, created_at, + updated_at, status, retry_count, max_retries, processing_started_at, processing_completed_at, + processor_id, lock_expires_at, last_error, error_details, event_id, contract_address, priority, + metadata, next_retry_at) + VALUES (1, '{"n":1}', 'hash', 'discord', 'user-1', '2026-01-01', 'created', 'updated', + 'FAILED', 1, 3, NULL, 'completed', 'worker-1', NULL, 'error', NULL, 'event-1', 'contract-1', + 5, '{"meta":true}', NULL)`, + ); + await run( + db, + `INSERT INTO notification_execution_log VALUES (1, 1, 1, 'attempt-time', 'FAILED', 'failed', NULL, 3)`, + ); + await run( + db, + `INSERT INTO dead_letter_queue VALUES (1, 1, 'discord', 'user-1', '{"n":1}', 'failed', NULL, 'dlq-time', NULL, 1)`, + ); + await run( + db, + `INSERT INTO processed_events VALUES (1, 'event-1', 'contract-1', 'fingerprint-1', 42, 'tx-1', 'contract', 'processed-time', 1, 2, 'redetected-time', 'ERROR', 0, 'failed')`, + ); + await run( + db, + `INSERT INTO polling_cursors VALUES (1, 'contract-1', 'cursor-42', 42, 'cursor-time', 1, 2)`, + ); + await run( + db, + `INSERT INTO idempotency_keys VALUES (1, 'key-1', 'hash-1', 1, '{"id":1}', 'key-time', 'expiry-time', 'EXPIRED')`, + ); + await run( + db, + `INSERT INTO rate_limit_events VALUES (1, 'client-1', 'API_KEY', '/api', 'POST', 'rate-time', 10, 1000)`, + ); + await run(db, `INSERT INTO notification_templates VALUES ('template-1')`); + await run( + db, + `INSERT INTO notification_template_audit_log VALUES (1, 'template-1', 'operator', 'UPDATE', 'audit-time', '{}', '{}')`, + ); + await run( + db, + `INSERT INTO backpressure_events VALUES (1, 'ACTIVATED', 12, 100, NULL, 'load', 'pressure-time')`, + ); + await run( + db, + `INSERT INTO notification_metrics_snapshots VALUES (1, 'capture-time', 10, 20, 5, '{}')`, + ); + await run( + db, + `INSERT INTO notification_archive VALUES (1, 1, '{"n":1}', 'discord', 'user-1', 'execute-time', 'created', 'completed', 'COMPLETED', 1, 'done', 'event-1', 'contract-1', '{}', 'archive-time')`, + ); +} + +const PRESERVED_TABLES = [ + 'scheduled_notifications', + 'notification_execution_log', + 'dead_letter_queue', + 'processed_events', + 'polling_cursors', + 'idempotency_keys', + 'rate_limit_events', + 'notification_template_audit_log', + 'backpressure_events', + 'notification_metrics_snapshots', + 'notification_archive', +]; + +describe('migration 004 database data integrity', () => { + let db: sqlite3.Database; + + beforeEach(async () => { + db = await openDatabase(); + await run(db, 'PRAGMA foreign_keys = ON'); + const pragma = await all<{ foreign_keys: number }>(db, 'PRAGMA foreign_keys'); + expect(pragma[0].foreign_keys).toBe(1); + await createLegacySchema(db); + await seedValidRows(db); + }); + + afterEach(async () => { + await close(db); + }); + + it('preserves valid legacy rows and enforces FK, status, and existing unique-key constraints', async () => { + const before: Record = {}; + for (const table of PRESERVED_TABLES) { + before[table] = await all(db, `SELECT * FROM ${table} ORDER BY id`); + } + + await run(db, 'PRAGMA foreign_keys = OFF'); + const runner = new MigrationRunner(db, ''); + await runner.initializeMigrationTable(); + const runnerPragma = await all<{ foreign_keys: number }>(db, 'PRAGMA foreign_keys'); + expect(runnerPragma[0].foreign_keys).toBe(1); + await runner.applyMigration(migration); + await runner.applyMigration({ ...migration, id: '004-repeat' }); + + for (const table of PRESERVED_TABLES) { + expect(await all(db, `SELECT * FROM ${table} ORDER BY id`)).toEqual(before[table]); + } + + const foreignKeyViolations = await all(db, 'PRAGMA foreign_key_check'); + expect(foreignKeyViolations).toHaveLength(0); + + await expect( + run( + db, + `INSERT INTO notification_execution_log (scheduled_notification_id, execution_attempt, status) + VALUES (999, 1, 'SUCCESS')`, + ), + ).rejects.toThrow(/FOREIGN KEY constraint failed/); + await expect( + run( + db, + `INSERT INTO scheduled_notifications (payload, notification_type, target_recipient, execute_at, status) + VALUES ('{}', 'discord', 'user-2', '2026-01-01', 'UNKNOWN')`, + ), + ).rejects.toThrow(/CHECK constraint failed/); + await expect( + run( + db, + `INSERT INTO processed_events + (event_id, contract_address, fingerprint, ledger_number, event_type) + VALUES ('event-duplicate', 'contract-1', 'fingerprint-1', 43, 'contract')`, + ), + ).rejects.toThrow(/UNIQUE constraint failed/); + }); + + it('fails closed with a table and row count when legacy state is invalid', async () => { + await run(db, `UPDATE processed_events SET status = 'UNKNOWN' WHERE id = 1`); + const runner = new MigrationRunner(db, ''); + await runner.initializeMigrationTable(); + + await expect(runner.applyMigration(migration)).rejects.toThrow( + /Migration 004 aborted.*processed_events: 1 row\(s\)/, + ); + const rows = await all<{ status: string }>( + db, + 'SELECT status FROM processed_events WHERE id = 1', + ); + expect(rows[0].status).toBe('UNKNOWN'); + const temporaryTable = await all(db, `SELECT name FROM sqlite_master WHERE name LIKE '%_v004'`); + expect(temporaryTable).toHaveLength(0); + }); + + it('enables foreign keys and installs matching constraints on fresh bootstrap databases', async () => { + const freshDb = new Database(':memory:'); + await freshDb.initialize(); + try { + const pragma = await freshDb.get<{ foreign_keys: number }>('PRAGMA foreign_keys'); + expect(pragma?.foreign_keys).toBe(1); + await expect( + freshDb.run(`INSERT INTO scheduled_notifications + (payload, notification_type, target_recipient, execute_at, status) + VALUES ('{}', 'discord', 'user-1', '2026-01-01', 'UNKNOWN')`), + ).rejects.toThrow(/CHECK constraint failed/); + } finally { + await freshDb.close(); + } + }); +}); diff --git a/listener/src/migrations/004-database-data-integrity.ts b/listener/src/migrations/004-database-data-integrity.ts new file mode 100644 index 00000000..7ecaedf9 --- /dev/null +++ b/listener/src/migrations/004-database-data-integrity.ts @@ -0,0 +1,446 @@ +import * as sqlite3 from 'sqlite3'; + +const tables = [ + { + name: 'scheduled_notifications', + create: `CREATE TABLE scheduled_notifications_v004 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + payload TEXT NOT NULL, + payload_hash TEXT, + notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), + target_recipient TEXT NOT NULL, + execute_at DATETIME NOT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + status VARCHAR(20) NOT NULL DEFAULT 'PENDING' CHECK (status IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'CANCELLED')), + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), + max_retries INTEGER NOT NULL DEFAULT 3 CHECK (max_retries >= 0), + processing_started_at DATETIME, + processing_completed_at DATETIME, + processor_id VARCHAR(100), + lock_expires_at DATETIME, + last_error TEXT, + error_details TEXT, + event_id TEXT, + contract_address TEXT, + priority INTEGER NOT NULL DEFAULT 5 CHECK (priority BETWEEN 1 AND 10), + metadata TEXT, + next_retry_at DATETIME + )`, + }, + { + name: 'notification_execution_log', + create: `CREATE TABLE notification_execution_log_v004 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + scheduled_notification_id INTEGER NOT NULL, + execution_attempt INTEGER NOT NULL CHECK (execution_attempt > 0), + execution_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + status VARCHAR(20) NOT NULL CHECK (status IN ('SUCCESS', 'FAILED', 'RETRY')), + error_message TEXT, + response_data TEXT, + duration_ms INTEGER CHECK (duration_ms IS NULL OR duration_ms >= 0), + FOREIGN KEY (scheduled_notification_id) REFERENCES scheduled_notifications_v004(id) ON DELETE CASCADE + )`, + }, + { + name: 'dead_letter_queue', + create: `CREATE TABLE dead_letter_queue_v004 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + scheduled_notification_id INTEGER NOT NULL UNIQUE, + notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), + target_recipient TEXT NOT NULL, + payload TEXT NOT NULL, + failure_reason TEXT NOT NULL, + error_details TEXT, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + last_retried_at DATETIME, + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), + FOREIGN KEY (scheduled_notification_id) REFERENCES scheduled_notifications_v004(id) ON DELETE CASCADE + )`, + }, + { + name: 'processed_events', + create: `CREATE TABLE processed_events_v004 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + event_id TEXT NOT NULL, + contract_address TEXT NOT NULL, + fingerprint TEXT NOT NULL UNIQUE, + ledger_number INTEGER NOT NULL CHECK (ledger_number >= 0), + tx_hash TEXT, + event_type VARCHAR(50) NOT NULL, + processed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + is_reorg_duplicate BOOLEAN NOT NULL DEFAULT 0 CHECK (is_reorg_duplicate IN (0, 1)), + reorg_detection_count INTEGER NOT NULL DEFAULT 0 CHECK (reorg_detection_count >= 0), + last_redetected_at DATETIME, + status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED' CHECK (status IN ('PROCESSED', 'SKIPPED', 'ERROR')), + notification_sent BOOLEAN NOT NULL DEFAULT 0 CHECK (notification_sent IN (0, 1)), + error_reason TEXT + )`, + }, + { + name: 'polling_cursors', + create: `CREATE TABLE polling_cursors_v004 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + contract_address TEXT NOT NULL UNIQUE, + cursor TEXT NOT NULL, + ledger_number INTEGER NOT NULL CHECK (ledger_number >= 0), + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + reorg_detected BOOLEAN NOT NULL DEFAULT 0 CHECK (reorg_detected IN (0, 1)), + reorg_detection_count INTEGER NOT NULL DEFAULT 0 CHECK (reorg_detection_count >= 0) + )`, + }, + { + name: 'idempotency_keys', + create: `CREATE TABLE idempotency_keys_v004 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + idempotency_key TEXT NOT NULL UNIQUE, + request_hash TEXT NOT NULL, + response_notification_id INTEGER NOT NULL, + response_data TEXT NOT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + expires_at DATETIME NOT NULL, + status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED' CHECK (status IN ('PROCESSED', 'EXPIRED')), + FOREIGN KEY (response_notification_id) REFERENCES scheduled_notifications_v004(id) ON DELETE CASCADE + )`, + }, + { + name: 'rate_limit_events', + create: `CREATE TABLE rate_limit_events_v004 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + client_id TEXT NOT NULL, + client_type VARCHAR(20) NOT NULL CHECK (client_type IN ('IP', 'API_KEY')), + endpoint TEXT NOT NULL, + method VARCHAR(10) NOT NULL, + timestamp DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + limit_threshold INTEGER NOT NULL CHECK (limit_threshold >= 0), + window_ms INTEGER NOT NULL CHECK (window_ms >= 0) + )`, + }, + { + name: 'notification_template_audit_log', + create: `CREATE TABLE notification_template_audit_log_v004 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + template_id TEXT NOT NULL, + actor TEXT NOT NULL, + action TEXT NOT NULL DEFAULT 'UPDATE' CHECK (action IN ('UPDATE')), + changed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + previous_snapshot TEXT NOT NULL, + new_snapshot TEXT NOT NULL, + FOREIGN KEY (template_id) REFERENCES notification_templates(id) ON DELETE RESTRICT + )`, + }, + { + name: 'backpressure_events', + create: `CREATE TABLE backpressure_events_v004 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + event_type VARCHAR(20) NOT NULL CHECK (event_type IN ('ACTIVATED', 'DEACTIVATED')), + queue_size INTEGER NOT NULL CHECK (queue_size >= 0), + target_throughput_per_sec INTEGER NOT NULL CHECK (target_throughput_per_sec >= 0), + duration_ms INTEGER CHECK (duration_ms IS NULL OR duration_ms >= 0), + reason TEXT, + timestamp DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP + )`, + }, + { + name: 'notification_metrics_snapshots', + create: `CREATE TABLE notification_metrics_snapshots_v004 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + captured_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + window_start INTEGER NOT NULL, + window_end INTEGER NOT NULL, + total_recorded INTEGER NOT NULL CHECK (total_recorded >= 0), + snapshot_json TEXT NOT NULL, + CHECK (window_start <= window_end) + )`, + }, + { + name: 'notification_archive', + create: `CREATE TABLE notification_archive_v004 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + original_id INTEGER NOT NULL, + payload TEXT NOT NULL, + notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), + target_recipient TEXT NOT NULL, + execute_at DATETIME NOT NULL, + created_at DATETIME NOT NULL, + processing_completed_at DATETIME, + status VARCHAR(20) NOT NULL CHECK (status IN ('COMPLETED', 'FAILED', 'CANCELLED')), + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), + last_error TEXT, + event_id TEXT, + contract_address TEXT, + metadata TEXT, + archived_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP + )`, + }, +]; + +const audits: Array<{ table: string; sql: string }> = [ + { + table: 'scheduled_notifications', + sql: `SELECT COUNT(*) AS count FROM scheduled_notifications + WHERE payload IS NULL OR notification_type IS NULL OR notification_type NOT IN ('discord', 'email', 'webhook', 'sms') + OR target_recipient IS NULL OR execute_at IS NULL + OR status IS NULL OR status NOT IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'CANCELLED') + OR retry_count IS NULL OR retry_count < 0 OR max_retries IS NULL OR max_retries < 0 + OR priority IS NULL OR priority NOT BETWEEN 1 AND 10`, + }, + { + table: 'notification_execution_log', + sql: `SELECT COUNT(*) AS count FROM notification_execution_log l + WHERE l.scheduled_notification_id IS NULL OR l.execution_attempt IS NULL OR l.execution_attempt <= 0 + OR l.execution_time IS NULL OR l.status IS NULL OR l.status NOT IN ('SUCCESS', 'FAILED', 'RETRY') + OR (l.duration_ms IS NOT NULL AND l.duration_ms < 0) + OR NOT EXISTS (SELECT 1 FROM scheduled_notifications n WHERE n.id = l.scheduled_notification_id)`, + }, + { + table: 'dead_letter_queue', + sql: `SELECT COUNT(*) AS count FROM dead_letter_queue d + WHERE d.scheduled_notification_id IS NULL OR d.notification_type IS NULL + OR d.notification_type NOT IN ('discord', 'email', 'webhook', 'sms') + OR d.target_recipient IS NULL OR d.payload IS NULL OR d.failure_reason IS NULL + OR d.retry_count IS NULL OR d.retry_count < 0 + OR NOT EXISTS (SELECT 1 FROM scheduled_notifications n WHERE n.id = d.scheduled_notification_id)`, + }, + { + table: 'dead_letter_queue (duplicate notification references)', + sql: `SELECT COALESCE(SUM(row_count), 0) AS count FROM ( + SELECT COUNT(*) AS row_count FROM dead_letter_queue + GROUP BY scheduled_notification_id HAVING COUNT(*) > 1 + )`, + }, + { + table: 'processed_events', + sql: `SELECT COUNT(*) AS count FROM processed_events + WHERE event_id IS NULL OR contract_address IS NULL OR fingerprint IS NULL + OR ledger_number IS NULL OR ledger_number < 0 OR event_type IS NULL + OR is_reorg_duplicate IS NULL OR is_reorg_duplicate NOT IN (0, 1) + OR reorg_detection_count IS NULL OR reorg_detection_count < 0 + OR status IS NULL OR status NOT IN ('PROCESSED', 'SKIPPED', 'ERROR') + OR notification_sent IS NULL OR notification_sent NOT IN (0, 1)`, + }, + { + table: 'processed_events (duplicate fingerprints)', + sql: `SELECT COALESCE(SUM(row_count), 0) AS count FROM ( + SELECT COUNT(*) AS row_count FROM processed_events + GROUP BY fingerprint HAVING COUNT(*) > 1 + )`, + }, + { + table: 'polling_cursors', + sql: `SELECT COUNT(*) AS count FROM polling_cursors + WHERE contract_address IS NULL OR cursor IS NULL OR ledger_number IS NULL OR ledger_number < 0 + OR reorg_detected IS NULL OR reorg_detected NOT IN (0, 1) + OR reorg_detection_count IS NULL OR reorg_detection_count < 0`, + }, + { + table: 'polling_cursors (duplicate contracts)', + sql: `SELECT COALESCE(SUM(row_count), 0) AS count FROM ( + SELECT COUNT(*) AS row_count FROM polling_cursors + GROUP BY contract_address HAVING COUNT(*) > 1 + )`, + }, + { + table: 'idempotency_keys', + sql: `SELECT COUNT(*) AS count FROM idempotency_keys k + WHERE k.idempotency_key IS NULL OR k.request_hash IS NULL OR k.response_notification_id IS NULL + OR k.response_data IS NULL OR k.expires_at IS NULL OR k.status IS NULL OR k.status NOT IN ('PROCESSED', 'EXPIRED') + OR NOT EXISTS (SELECT 1 FROM scheduled_notifications n WHERE n.id = k.response_notification_id)`, + }, + { + table: 'idempotency_keys (duplicate keys)', + sql: `SELECT COALESCE(SUM(row_count), 0) AS count FROM ( + SELECT COUNT(*) AS row_count FROM idempotency_keys + GROUP BY idempotency_key HAVING COUNT(*) > 1 + )`, + }, + { + table: 'rate_limit_events', + sql: `SELECT COUNT(*) AS count FROM rate_limit_events + WHERE client_id IS NULL OR client_type IS NULL OR client_type NOT IN ('IP', 'API_KEY') + OR endpoint IS NULL OR method IS NULL OR limit_threshold IS NULL OR limit_threshold < 0 + OR window_ms IS NULL OR window_ms < 0`, + }, + { + table: 'notification_template_audit_log', + sql: `SELECT COUNT(*) AS count FROM notification_template_audit_log a + WHERE a.template_id IS NULL OR a.actor IS NULL OR a.action IS NULL OR a.action NOT IN ('UPDATE') + OR a.previous_snapshot IS NULL OR a.new_snapshot IS NULL + OR NOT EXISTS (SELECT 1 FROM notification_templates t WHERE t.id = a.template_id)`, + }, + { + table: 'backpressure_events', + sql: `SELECT COUNT(*) AS count FROM backpressure_events + WHERE event_type IS NULL OR event_type NOT IN ('ACTIVATED', 'DEACTIVATED') + OR queue_size IS NULL OR queue_size < 0 OR target_throughput_per_sec IS NULL OR target_throughput_per_sec < 0 + OR (duration_ms IS NOT NULL AND duration_ms < 0)`, + }, + { + table: 'notification_metrics_snapshots', + sql: `SELECT COUNT(*) AS count FROM notification_metrics_snapshots + WHERE window_start IS NULL OR window_end IS NULL OR window_start > window_end + OR total_recorded IS NULL OR total_recorded < 0 OR snapshot_json IS NULL`, + }, + { + table: 'notification_archive', + sql: `SELECT COUNT(*) AS count FROM notification_archive + WHERE original_id IS NULL OR payload IS NULL OR notification_type IS NULL + OR notification_type NOT IN ('discord', 'email', 'webhook', 'sms') + OR target_recipient IS NULL OR execute_at IS NULL OR created_at IS NULL + OR status IS NULL OR status NOT IN ('COMPLETED', 'FAILED', 'CANCELLED') + OR retry_count IS NULL OR retry_count < 0 OR archived_at IS NULL`, + }, +]; + +const indexes = [ + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_status ON scheduled_notifications(status)`, + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_status_execute_at ON scheduled_notifications(status, execute_at) WHERE status = 'PENDING'`, + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_lock_expires ON scheduled_notifications(lock_expires_at, status) WHERE status = 'PROCESSING'`, + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_next_retry_at ON scheduled_notifications(next_retry_at, status) WHERE status = 'PENDING'`, + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_created_at ON scheduled_notifications(created_at)`, + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_event_id ON scheduled_notifications(event_id) WHERE event_id IS NOT NULL`, + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_target ON scheduled_notifications(target_recipient, status)`, + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_claim ON scheduled_notifications(status, priority, execute_at) WHERE status = 'PENDING'`, + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_processor_lock ON scheduled_notifications(processor_id, status, lock_expires_at)`, + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_status_created ON scheduled_notifications(status, created_at)`, + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_type_status ON scheduled_notifications(notification_type, status)`, + `CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_contract ON scheduled_notifications(contract_address) WHERE contract_address IS NOT NULL`, + `CREATE INDEX IF NOT EXISTS idx_execution_log_notification_id ON notification_execution_log(scheduled_notification_id)`, + `CREATE INDEX IF NOT EXISTS idx_execution_log_execution_time ON notification_execution_log(execution_time)`, + `CREATE INDEX IF NOT EXISTS idx_execution_log_status_execution_time ON notification_execution_log(status, execution_time)`, + `CREATE INDEX IF NOT EXISTS idx_execution_log_notification_attempt ON notification_execution_log(scheduled_notification_id, execution_attempt)`, + `CREATE INDEX IF NOT EXISTS idx_dead_letter_queue_created_at ON dead_letter_queue(created_at)`, + `CREATE INDEX IF NOT EXISTS idx_dead_letter_queue_notification_type ON dead_letter_queue(notification_type)`, + `CREATE INDEX IF NOT EXISTS idx_processed_events_fingerprint ON processed_events(fingerprint)`, + `CREATE INDEX IF NOT EXISTS idx_processed_events_contract_event ON processed_events(contract_address, event_id)`, + `CREATE INDEX IF NOT EXISTS idx_processed_events_processed_at ON processed_events(processed_at)`, + `CREATE INDEX IF NOT EXISTS idx_processed_events_reorg_duplicates ON processed_events(is_reorg_duplicate, processed_at) WHERE is_reorg_duplicate = 1`, + `CREATE INDEX IF NOT EXISTS idx_processed_events_ledger_contract ON processed_events(ledger_number, contract_address)`, + `CREATE INDEX IF NOT EXISTS idx_processed_events_status_processed ON processed_events(status, processed_at)`, + `CREATE INDEX IF NOT EXISTS idx_processed_events_event_type_status ON processed_events(event_type, status)`, + `CREATE INDEX IF NOT EXISTS idx_processed_events_tx_hash ON processed_events(tx_hash) WHERE tx_hash IS NOT NULL`, + `CREATE INDEX IF NOT EXISTS idx_polling_cursors_contract ON polling_cursors(contract_address)`, + `CREATE INDEX IF NOT EXISTS idx_polling_cursors_updated_at ON polling_cursors(updated_at)`, + `CREATE INDEX IF NOT EXISTS idx_idempotency_keys_key ON idempotency_keys(idempotency_key)`, + `CREATE INDEX IF NOT EXISTS idx_idempotency_keys_expires_at ON idempotency_keys(expires_at)`, + `CREATE INDEX IF NOT EXISTS idx_idempotency_keys_created_at ON idempotency_keys(created_at)`, + `CREATE INDEX IF NOT EXISTS idx_rate_limit_events_timestamp ON rate_limit_events(timestamp)`, + `CREATE INDEX IF NOT EXISTS idx_rate_limit_events_client_id ON rate_limit_events(client_id)`, + `CREATE INDEX IF NOT EXISTS idx_rate_limit_events_client_timestamp ON rate_limit_events(client_id, timestamp)`, + `CREATE INDEX IF NOT EXISTS idx_template_audit_template_id ON notification_template_audit_log(template_id)`, + `CREATE INDEX IF NOT EXISTS idx_template_audit_changed_at ON notification_template_audit_log(changed_at)`, + `CREATE INDEX IF NOT EXISTS idx_backpressure_events_type ON backpressure_events(event_type)`, + `CREATE INDEX IF NOT EXISTS idx_backpressure_events_timestamp ON backpressure_events(timestamp)`, + `CREATE INDEX IF NOT EXISTS idx_backpressure_events_type_timestamp ON backpressure_events(event_type, timestamp)`, + `CREATE INDEX IF NOT EXISTS idx_metrics_snapshots_captured_at ON notification_metrics_snapshots(captured_at)`, + `CREATE INDEX IF NOT EXISTS idx_archive_original_id ON notification_archive(original_id)`, + `CREATE INDEX IF NOT EXISTS idx_archive_archived_at ON notification_archive(archived_at)`, + `CREATE INDEX IF NOT EXISTS idx_archive_status ON notification_archive(status)`, + `CREATE INDEX IF NOT EXISTS idx_archive_contract_address ON notification_archive(contract_address) WHERE contract_address IS NOT NULL`, + `CREATE INDEX IF NOT EXISTS idx_archive_event_id ON notification_archive(event_id) WHERE event_id IS NOT NULL`, +]; + +function run(db: sqlite3.Database, sql: string, params: unknown[] = []): Promise { + return new Promise((resolve, reject) => { + db.run(sql, params, (error) => (error ? reject(error) : resolve())); + }); +} + +function getCount(db: sqlite3.Database, sql: string): Promise { + return new Promise((resolve, reject) => { + db.get(sql, (error, row: { count: number }) => { + if (error) reject(error); + else resolve(row.count); + }); + }); +} + +function getRows(db: sqlite3.Database, sql: string): Promise { + return new Promise((resolve, reject) => { + db.all(sql, (error, rows) => (error ? reject(error) : resolve(rows))); + }); +} + +async function assertLegacyDataIsValid(db: sqlite3.Database): Promise { + const violations: string[] = []; + for (const audit of audits) { + const count = await getCount(db, audit.sql); + if (count > 0) violations.push(`${audit.table}: ${count} row(s)`); + } + if (violations.length > 0) { + throw new Error( + `Migration 004 aborted before rebuilding tables. Repair invalid legacy data first: ${violations.join('; ')}`, + ); + } +} + +const migration = { + id: '004', + name: 'database-data-integrity', + up: async (db: sqlite3.Database) => { + await assertLegacyDataIsValid(db); + + for (const table of tables) { + await run(db, table.create); + await run(db, `INSERT INTO ${table.name}_v004 SELECT * FROM ${table.name}`); + } + + for (const tableName of [ + 'notification_execution_log', + 'dead_letter_queue', + 'idempotency_keys', + 'notification_template_audit_log', + 'processed_events', + 'polling_cursors', + 'rate_limit_events', + 'backpressure_events', + 'notification_metrics_snapshots', + 'notification_archive', + 'scheduled_notifications', + ]) { + await run(db, `DROP TABLE ${tableName}`); + } + + await run(db, 'ALTER TABLE scheduled_notifications_v004 RENAME TO scheduled_notifications'); + for (const table of tables.slice(1)) { + await run(db, `ALTER TABLE ${table.name}_v004 RENAME TO ${table.name}`); + } + + for (const index of indexes) await run(db, index); + + await run( + db, + `CREATE TRIGGER IF NOT EXISTS update_scheduled_notifications_timestamp + AFTER UPDATE ON scheduled_notifications + FOR EACH ROW BEGIN + UPDATE scheduled_notifications SET updated_at = CURRENT_TIMESTAMP WHERE id = NEW.id; + END`, + ); + await run( + db, + `CREATE TRIGGER IF NOT EXISTS prevent_template_audit_update + BEFORE UPDATE ON notification_template_audit_log + FOR EACH ROW BEGIN SELECT RAISE(ABORT, 'Audit records are immutable'); END`, + ); + await run( + db, + `CREATE TRIGGER IF NOT EXISTS prevent_template_audit_delete + BEFORE DELETE ON notification_template_audit_log + FOR EACH ROW BEGIN SELECT RAISE(ABORT, 'Audit records are immutable'); END`, + ); + + const foreignKeyViolations = await getRows(db, 'PRAGMA foreign_key_check'); + if (foreignKeyViolations.length > 0) { + throw new Error( + `Migration 004 detected ${foreignKeyViolations.length} foreign-key violation(s) after rebuilding tables`, + ); + } + }, + down: async () => { + throw new Error( + 'Migration 004 cannot be safely rolled back; restore a database backup instead', + ); + }, +}; + +export default migration; From 9930d61a88087790cd1ffc1857aa9181f76f44d5 Mon Sep 17 00:00:00 2001 From: Najib Ishiyaku Njidda Date: Tue, 29 Sep 2026 19:20:11 +0000 Subject: [PATCH 2/5] feat(listener): expire scheduled notifications past a configured deadline. Closes #840. Adds expires_at (explicit override > NOTIFICATION_DEFAULT_TTL_SECONDS > never), EXPIRED terminal status persisted and enforced via extended CHECKs in migration 005 (scheduled + archive tables, preserving 004 constraints), expiry gates before provider dispatch in both scheduled and retry loops, terminal/archive/cleanup handling, API-boundary normalization of ISO/epoch expiry, focused tests, operator docs. Stacked on #892 so merge order is irrelevant. --- NOTIFICATION_LIFECYCLE.md | 2 +- .../src/pages/NotificationSearchPage.tsx | 1 + listener/.env.example | 4 + listener/src/api/events-server.ts | 1 + listener/src/config.ts | 13 + listener/src/database/archive-schema.sql | 8 +- listener/src/database/schema.sql | 6 +- listener/src/index.ts | 6 +- .../005-notification-expiration.test.ts | 280 ++++++++++++++++ .../migrations/005-notification-expiration.ts | 312 ++++++++++++++++++ listener/src/services/archive-service.ts | 16 +- listener/src/services/archive-store.ts | 11 +- listener/src/services/cleanup-service.ts | 2 +- listener/src/services/notification-api.ts | 24 ++ .../src/services/notification-scheduler.ts | 50 ++- listener/src/services/retry-scheduler.ts | 18 + .../scheduled-notification-repository.ts | 54 ++- .../src/tests/notification-scheduler.test.ts | 146 +++++++- listener/src/types/index.ts | 2 + listener/src/types/scheduled-notification.ts | 4 + 20 files changed, 921 insertions(+), 39 deletions(-) create mode 100644 listener/src/migrations/005-notification-expiration.test.ts create mode 100644 listener/src/migrations/005-notification-expiration.ts diff --git a/NOTIFICATION_LIFECYCLE.md b/NOTIFICATION_LIFECYCLE.md index d5cd3688..d77227ff 100644 --- a/NOTIFICATION_LIFECYCLE.md +++ b/NOTIFICATION_LIFECYCLE.md @@ -292,7 +292,7 @@ Declared but not implemented in the scheduler today: `webhook`, `email`, `sms`. ### 7. Archive / purge -Terminal rows (`COMPLETED`, `FAILED`, `CANCELLED`) are later moved by +Terminal rows (`COMPLETED`, `FAILED`, `CANCELLED`, `EXPIRED`) are later moved by `ArchiveService` into `notification_archive`, then optionally purged after retention. See [Completion and Archival](#completion-and-archival). diff --git a/dashboard/src/pages/NotificationSearchPage.tsx b/dashboard/src/pages/NotificationSearchPage.tsx index 2efc7999..b28918d4 100644 --- a/dashboard/src/pages/NotificationSearchPage.tsx +++ b/dashboard/src/pages/NotificationSearchPage.tsx @@ -29,6 +29,7 @@ export const NOTIFICATION_DELIVERY_STATUS_OPTIONS = [ { value: 'COMPLETED', label: 'Completed' }, { value: 'FAILED', label: 'Failed' }, { value: 'CANCELLED', label: 'Cancelled' }, + { value: 'EXPIRED', label: 'Expired' }, { value: 'PROCESSED', label: 'Processed' }, ]; diff --git a/listener/.env.example b/listener/.env.example index c5f610b6..2049e41f 100644 --- a/listener/.env.example +++ b/listener/.env.example @@ -61,6 +61,10 @@ EVENTS_API_CORS_ORIGIN=http://localhost:5173 # Ensure this path is persistent and writable in your deployment environment. DATABASE_PATH=./data/notifications.db +# Default lifetime for scheduled notifications in seconds. Zero means never expire. +# A per-notification expiresAt value overrides this default. +NOTIFICATION_DEFAULT_TTL_SECONDS=0 + # ----------------------------------------------------------------------------- # Discord Delivery (optional — both must be provided together or neither) # ⚠️ SECRETS: inject via a secrets manager, never commit actual values. diff --git a/listener/src/api/events-server.ts b/listener/src/api/events-server.ts index d517ad21..64e4e827 100644 --- a/listener/src/api/events-server.ts +++ b/listener/src/api/events-server.ts @@ -858,6 +858,7 @@ export function createEventsServer(options: EventsServerOptions): http.Server { notificationType: data.notificationType || NotificationType.DISCORD, targetRecipient: data.targetRecipient, executeAt, + expiresAt: data.expiresAt, maxRetries: data.maxRetries, priority: data.priority, eventId: data.eventId, diff --git a/listener/src/config.ts b/listener/src/config.ts index 3dbdd3f6..808c27ac 100644 --- a/listener/src/config.ts +++ b/listener/src/config.ts @@ -48,6 +48,18 @@ function parseIntegerEnv(name: string, defaultValue: string): number { return parsed; } +function loadNotificationDefaultTtlSeconds(): number { + const rawValue = trimEnv('NOTIFICATION_DEFAULT_TTL_SECONDS') ?? '0'; + if (!/^\d+$/.test(rawValue)) { + throw new ConfigError('NOTIFICATION_DEFAULT_TTL_SECONDS must be a non-negative integer'); + } + const seconds = Number(rawValue); + if (!Number.isSafeInteger(seconds) || seconds > Math.floor(8.64e15 / 1000)) { + throw new ConfigError('NOTIFICATION_DEFAULT_TTL_SECONDS must be a supported non-negative integer'); + } + return seconds; +} + function parseJsonEnv(name: string, defaultValue: string): T { const rawValue = trimEnv(name) ?? defaultValue; try { @@ -303,6 +315,7 @@ export function loadConfig(): Config { cleanup: loadCleanupConfig(), analytics: loadAnalyticsConfig(), expiration: loadExpirationConfig(), + notificationDefaultTtlSeconds: loadNotificationDefaultTtlSeconds(), backfill: loadBackfillConfig(), logging: loadLoggingConfig(), api: loadApiConfig(), diff --git a/listener/src/database/archive-schema.sql b/listener/src/database/archive-schema.sql index e961b14c..c5f016c1 100644 --- a/listener/src/database/archive-schema.sql +++ b/listener/src/database/archive-schema.sql @@ -1,6 +1,3 @@ - notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), - status VARCHAR(20) NOT NULL CHECK (status IN ('COMPLETED', 'FAILED', 'CANCELLED')), - retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), -- Archive table for notifications moved out of active storage. -- Records here are read-only for audit purposes and are never modified. CREATE TABLE IF NOT EXISTS notification_archive ( @@ -14,12 +11,13 @@ CREATE TABLE IF NOT EXISTS notification_archive ( -- Original scheduling / timing execute_at DATETIME NOT NULL, + expires_at DATETIME, created_at DATETIME NOT NULL, processing_completed_at DATETIME, -- Final status at time of archiving - status VARCHAR(20) NOT NULL, -- COMPLETED | FAILED | CANCELLED - retry_count INTEGER NOT NULL DEFAULT 0, + status VARCHAR(20) NOT NULL CHECK (status IN ('COMPLETED', 'FAILED', 'CANCELLED', 'EXPIRED')), + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), last_error TEXT, -- Optional references diff --git a/listener/src/database/schema.sql b/listener/src/database/schema.sql index cb867162..f6cd1e12 100644 --- a/listener/src/database/schema.sql +++ b/listener/src/database/schema.sql @@ -13,11 +13,12 @@ CREATE TABLE IF NOT EXISTS scheduled_notifications ( -- Scheduling information execute_at DATETIME NOT NULL, -- When the notification should be sent + expires_at DATETIME, -- Optional deadline after which delivery is skipped created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, -- Status tracking - status VARCHAR(20) NOT NULL DEFAULT 'PENDING' CHECK (status IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'CANCELLED')), + status VARCHAR(20) NOT NULL DEFAULT 'PENDING' CHECK (status IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'CANCELLED', 'EXPIRED')), retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), max_retries INTEGER NOT NULL DEFAULT 3 CHECK (max_retries >= 0), @@ -343,9 +344,10 @@ CREATE TABLE IF NOT EXISTS notification_archive ( notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), target_recipient TEXT NOT NULL, execute_at DATETIME NOT NULL, + expires_at DATETIME, created_at DATETIME NOT NULL, processing_completed_at DATETIME, - status VARCHAR(20) NOT NULL CHECK (status IN ('COMPLETED', 'FAILED', 'CANCELLED')), + status VARCHAR(20) NOT NULL CHECK (status IN ('COMPLETED', 'FAILED', 'CANCELLED', 'EXPIRED')), retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), last_error TEXT, event_id TEXT, diff --git a/listener/src/index.ts b/listener/src/index.ts index 169ab237..65893a22 100644 --- a/listener/src/index.ts +++ b/listener/src/index.ts @@ -66,7 +66,11 @@ async function main() { logger.info('Initializing database'); const db = await initializeDatabase(config.databasePath); - repository = new ScheduledNotificationRepository(db); + repository = new ScheduledNotificationRepository( + db, + undefined, + config.notificationDefaultTtlSeconds ?? 0, + ); healthMonitor = new NotificationHealthMonitor(null, getWorkerManager(), { repository, diff --git a/listener/src/migrations/005-notification-expiration.test.ts b/listener/src/migrations/005-notification-expiration.test.ts new file mode 100644 index 00000000..567ef897 --- /dev/null +++ b/listener/src/migrations/005-notification-expiration.test.ts @@ -0,0 +1,280 @@ +import * as sqlite3 from 'sqlite3'; +import { MigrationRunner } from '../database/migration-system'; +import migration from './005-notification-expiration'; + +function openDatabase(): Promise { + return new Promise((resolve, reject) => { + const db = new sqlite3.Database(':memory:', (error) => { + if (error) reject(error); + else resolve(db); + }); + }); +} + +function run(db: sqlite3.Database, sql: string, params: unknown[] = []): Promise { + return new Promise((resolve, reject) => { + db.run(sql, params, (error) => (error ? reject(error) : resolve())); + }); +} + +function exec(db: sqlite3.Database, sql: string): Promise { + return new Promise((resolve, reject) => { + db.exec(sql, (error) => (error ? reject(error) : resolve())); + }); +} + +function all(db: sqlite3.Database, sql: string): Promise { + return new Promise((resolve, reject) => { + db.all(sql, (error, rows) => (error ? reject(error) : resolve(rows as T[]))); + }); +} + +function close(db: sqlite3.Database): Promise { + return new Promise((resolve, reject) => { + db.close((error) => (error ? reject(error) : resolve())); + }); +} + +async function createLegacySchema(db: sqlite3.Database): Promise { + await exec( + db, + ` + CREATE TABLE scheduled_notifications ( + id INTEGER PRIMARY KEY AUTOINCREMENT, payload TEXT NOT NULL, payload_hash TEXT, + notification_type VARCHAR(50) NOT NULL, target_recipient TEXT NOT NULL, + execute_at DATETIME NOT NULL, created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, status VARCHAR(20) NOT NULL DEFAULT 'PENDING', + retry_count INTEGER NOT NULL DEFAULT 0, max_retries INTEGER NOT NULL DEFAULT 3, + processing_started_at DATETIME, processing_completed_at DATETIME, processor_id VARCHAR(100), + lock_expires_at DATETIME, last_error TEXT, error_details TEXT, event_id TEXT, + contract_address TEXT, priority INTEGER NOT NULL DEFAULT 5, metadata TEXT, next_retry_at DATETIME + ); + CREATE TABLE notification_execution_log ( + id INTEGER PRIMARY KEY AUTOINCREMENT, scheduled_notification_id INTEGER NOT NULL, + execution_attempt INTEGER NOT NULL, execution_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + status VARCHAR(20) NOT NULL, error_message TEXT, response_data TEXT, duration_ms INTEGER, + FOREIGN KEY (scheduled_notification_id) REFERENCES scheduled_notifications(id) ON DELETE CASCADE + ); + CREATE TABLE dead_letter_queue ( + id INTEGER PRIMARY KEY AUTOINCREMENT, scheduled_notification_id INTEGER NOT NULL UNIQUE, + notification_type VARCHAR(50) NOT NULL, target_recipient TEXT NOT NULL, payload TEXT NOT NULL, + failure_reason TEXT NOT NULL, error_details TEXT, created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + last_retried_at DATETIME, retry_count INTEGER NOT NULL DEFAULT 0, + FOREIGN KEY (scheduled_notification_id) REFERENCES scheduled_notifications(id) ON DELETE CASCADE + ); + CREATE TABLE idempotency_keys ( + id INTEGER PRIMARY KEY AUTOINCREMENT, idempotency_key TEXT NOT NULL UNIQUE, request_hash TEXT NOT NULL, + response_notification_id INTEGER NOT NULL, response_data TEXT NOT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, expires_at DATETIME NOT NULL, + status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED', + FOREIGN KEY (response_notification_id) REFERENCES scheduled_notifications(id) ON DELETE CASCADE + ); + CREATE TABLE notification_archive ( + id INTEGER PRIMARY KEY AUTOINCREMENT, original_id INTEGER NOT NULL, payload TEXT NOT NULL, + notification_type VARCHAR(50) NOT NULL, target_recipient TEXT NOT NULL, + execute_at DATETIME NOT NULL, created_at DATETIME NOT NULL, processing_completed_at DATETIME, + status VARCHAR(20) NOT NULL, retry_count INTEGER NOT NULL DEFAULT 0, last_error TEXT, + event_id TEXT, contract_address TEXT, metadata TEXT, + archived_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP + ); + `, + ); +} + +async function seedValidRows(db: sqlite3.Database): Promise { + await run( + db, + `INSERT INTO scheduled_notifications + (id, payload, payload_hash, notification_type, target_recipient, execute_at, created_at, + updated_at, status, retry_count, max_retries, processing_started_at, processing_completed_at, + processor_id, lock_expires_at, last_error, error_details, event_id, contract_address, priority, + metadata, next_retry_at) + VALUES (1, '{"n":1}', 'hash', 'discord', 'user-1', '2026-01-01', 'created', 'updated', + 'FAILED', 1, 3, NULL, 'completed', 'worker-1', NULL, 'error', NULL, 'event-1', 'contract-1', + 5, '{"meta":true}', NULL)`, + ); + await run( + db, + `INSERT INTO notification_execution_log VALUES + (1, 1, 1, 'attempt-time', 'FAILED', 'failed', NULL, 3)`, + ); + await run( + db, + `INSERT INTO dead_letter_queue VALUES + (1, 1, 'discord', 'user-1', '{"n":1}', 'failed', NULL, 'dlq-time', NULL, 1)`, + ); + await run( + db, + `INSERT INTO idempotency_keys VALUES + (1, 'key-1', 'hash-1', 1, '{"id":1}', 'key-time', 'expiry-time', 'EXPIRED')`, + ); + await run( + db, + `INSERT INTO notification_archive VALUES + (1, 1, '{"n":1}', 'discord', 'user-1', 'execute-time', 'created', 'completed', + 'COMPLETED', 1, 'done', 'event-1', 'contract-1', '{}', 'archive-time')`, + ); +} + +const preservedRows = async (db: sqlite3.Database) => ({ + scheduled_notifications: await all( + db, + `SELECT id, payload, payload_hash, notification_type, target_recipient, execute_at, created_at, + updated_at, status, retry_count, max_retries, processing_started_at, processing_completed_at, + processor_id, lock_expires_at, last_error, error_details, event_id, contract_address, priority, + metadata, next_retry_at FROM scheduled_notifications ORDER BY id`, + ), + notification_execution_log: await all(db, 'SELECT * FROM notification_execution_log ORDER BY id'), + dead_letter_queue: await all(db, 'SELECT * FROM dead_letter_queue ORDER BY id'), + idempotency_keys: await all(db, 'SELECT * FROM idempotency_keys ORDER BY id'), + notification_archive: await all( + db, + `SELECT id, original_id, payload, notification_type, target_recipient, execute_at, created_at, + processing_completed_at, status, retry_count, last_error, event_id, contract_address, metadata, + archived_at FROM notification_archive ORDER BY id`, + ), +}); + +async function applyMigration(db: sqlite3.Database, id = migration.id): Promise { + const runner = new MigrationRunner(db, ''); + await runner.initializeMigrationTable(); + await runner.applyMigration({ ...migration, id }); +} + +async function expectNoRebuildTables(db: sqlite3.Database): Promise { + const tables = await all<{ name: string }>( + db, + "SELECT name FROM sqlite_master WHERE type = 'table' AND name LIKE '%_v005'", + ); + expect(tables).toHaveLength(0); +} + +describe('migration 005 notification expiration', () => { + let db: sqlite3.Database; + + beforeEach(async () => { + db = await openDatabase(); + await run(db, 'PRAGMA foreign_keys = ON'); + await createLegacySchema(db); + await seedValidRows(db); + }); + + afterEach(async () => { + await close(db); + }); + + it('preserves legacy rows, enforces the extended checks, and is idempotent', async () => { + const before = await preservedRows(db); + + await applyMigration(db); + + expect(await preservedRows(db)).toEqual(before); + expect( + await all<{ expires_at: string | null }>( + db, + 'SELECT expires_at FROM scheduled_notifications', + ), + ).toEqual([{ expires_at: null }]); + expect( + await all<{ expires_at: string | null }>(db, 'SELECT expires_at FROM notification_archive'), + ).toEqual([{ expires_at: null }]); + + await run( + db, + `INSERT INTO scheduled_notifications (payload, notification_type, target_recipient, execute_at, status) + VALUES ('{}', 'discord', 'user-2', '2026-01-01', 'EXPIRED')`, + ); + await expect( + run( + db, + `INSERT INTO scheduled_notifications (payload, notification_type, target_recipient, execute_at, status) + VALUES ('{}', 'discord', 'user-3', '2026-01-01', 'UNKNOWN')`, + ), + ).rejects.toThrow(/CHECK constraint failed/); + + await expect( + run( + db, + `INSERT INTO notification_execution_log + (scheduled_notification_id, execution_attempt, status) VALUES (999, 1, 'SUCCESS')`, + ), + ).rejects.toThrow(/FOREIGN KEY constraint failed/); + await expect( + run( + db, + `INSERT INTO notification_execution_log + (scheduled_notification_id, execution_attempt, status) VALUES (1, 1, 'UNKNOWN')`, + ), + ).rejects.toThrow(/CHECK constraint failed/); + await expect( + run( + db, + `INSERT INTO dead_letter_queue + (scheduled_notification_id, notification_type, target_recipient, payload, failure_reason) + VALUES (999, 'discord', 'user-9', '{}', 'failed')`, + ), + ).rejects.toThrow(/FOREIGN KEY constraint failed/); + await expect( + run( + db, + `INSERT INTO dead_letter_queue + (scheduled_notification_id, notification_type, target_recipient, payload, failure_reason) + VALUES (1, 'unknown', 'user-9', '{}', 'failed')`, + ), + ).rejects.toThrow(/CHECK constraint failed/); + await expect( + run( + db, + `INSERT INTO idempotency_keys + (idempotency_key, request_hash, response_notification_id, response_data, expires_at) + VALUES ('orphan-key', 'hash', 999, '{}', 'expiry')`, + ), + ).rejects.toThrow(/FOREIGN KEY constraint failed/); + await expect( + run( + db, + `INSERT INTO idempotency_keys + (idempotency_key, request_hash, response_notification_id, response_data, expires_at, status) + VALUES ('unknown-status-key', 'hash', 1, '{}', 'expiry', 'UNKNOWN')`, + ), + ).rejects.toThrow(/CHECK constraint failed/); + await expect( + run( + db, + `INSERT INTO notification_archive + (original_id, payload, notification_type, target_recipient, execute_at, created_at, status) + VALUES (2, '{}', 'discord', 'user-2', 'execute', 'created', 'UNKNOWN')`, + ), + ).rejects.toThrow(/CHECK constraint failed/); + + const afterFirstApply = await preservedRows(db); + await applyMigration(db, '005-repeat'); + expect(await preservedRows(db)).toEqual(afterFirstApply); + }); + + it('aborts before rebuilding when legacy statuses are invalid', async () => { + await run(db, `UPDATE scheduled_notifications SET status = 'UNKNOWN' WHERE id = 1`); + + await expect(applyMigration(db)).rejects.toThrow( + /Migration 005 aborted before rebuilding tables.*scheduled_notifications: 1 row\(s\)/, + ); + await expectNoRebuildTables(db); + expect(await all(db, 'SELECT status FROM scheduled_notifications WHERE id = 1')).toEqual([ + { status: 'UNKNOWN' }, + ]); + }); + + it('aborts before rebuilding when a legacy expiration timestamp is unparseable', async () => { + await run(db, 'ALTER TABLE scheduled_notifications ADD COLUMN expires_at DATETIME'); + await run(db, 'ALTER TABLE notification_archive ADD COLUMN expires_at DATETIME'); + await run(db, `UPDATE scheduled_notifications SET expires_at = 'not-a-date' WHERE id = 1`); + + await expect(applyMigration(db)).rejects.toThrow( + /Migration 005 aborted before rebuilding tables.*scheduled_notifications: 1 row\(s\)/, + ); + await expectNoRebuildTables(db); + expect(await all(db, 'SELECT expires_at FROM scheduled_notifications WHERE id = 1')).toEqual([ + { expires_at: 'not-a-date' }, + ]); + }); +}); diff --git a/listener/src/migrations/005-notification-expiration.ts b/listener/src/migrations/005-notification-expiration.ts new file mode 100644 index 00000000..70476a6d --- /dev/null +++ b/listener/src/migrations/005-notification-expiration.ts @@ -0,0 +1,312 @@ +import * as sqlite3 from 'sqlite3'; + +function run(db: sqlite3.Database, sql: string, params: unknown[] = []): Promise { + return new Promise((resolve, reject) => { + db.run(sql, params, (error) => (error ? reject(error) : resolve())); + }); +} + +function getCount(db: sqlite3.Database, sql: string): Promise { + return new Promise((resolve, reject) => { + db.get(sql, (error, row: { count: number }) => { + if (error) reject(error); + else resolve(row.count); + }); + }); +} + +function getRows(db: sqlite3.Database, sql: string): Promise { + return new Promise((resolve, reject) => { + db.all(sql, (error, rows) => (error ? reject(error) : resolve(rows)); + }); +} + +async function hasColumn(db: sqlite3.Database, table: string, column: string): Promise { + const rows = await new Promise>((resolve, reject) => { + db.all(`PRAGMA table_info(${table})`, (error, result) => + error ? reject(error) : resolve(result as Array<{ name: string }>), + ); + }); + return rows.some((row) => row.name === column); +} + +async function auditLegacyRows(db: sqlite3.Database): Promise { + const audits = [ + { + table: 'scheduled_notifications', + sql: `SELECT COUNT(*) AS count FROM scheduled_notifications + WHERE payload IS NULL OR notification_type NOT IN ('discord', 'email', 'webhook', 'sms') + OR target_recipient IS NULL OR execute_at IS NULL + OR status NOT IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'CANCELLED', 'EXPIRED') + OR retry_count < 0 OR max_retries < 0 OR priority NOT BETWEEN 1 AND 10`, + }, + { + table: 'notification_execution_log', + sql: `SELECT COUNT(*) AS count FROM notification_execution_log l + WHERE l.scheduled_notification_id IS NULL OR l.execution_attempt <= 0 + OR l.execution_time IS NULL OR l.status NOT IN ('SUCCESS', 'FAILED', 'RETRY') + OR (l.duration_ms IS NOT NULL AND l.duration_ms < 0) + OR NOT EXISTS (SELECT 1 FROM scheduled_notifications n WHERE n.id = l.scheduled_notification_id)`, + }, + { + table: 'dead_letter_queue', + sql: `SELECT COUNT(*) AS count FROM dead_letter_queue d + WHERE d.scheduled_notification_id IS NULL + OR d.notification_type NOT IN ('discord', 'email', 'webhook', 'sms') + OR d.target_recipient IS NULL OR d.payload IS NULL OR d.failure_reason IS NULL + OR d.retry_count < 0 + OR NOT EXISTS (SELECT 1 FROM scheduled_notifications n WHERE n.id = d.scheduled_notification_id)`, + }, + { + table: 'dead_letter_queue (duplicate notification references)', + sql: `SELECT COUNT(*) AS count FROM ( + SELECT scheduled_notification_id FROM dead_letter_queue + GROUP BY scheduled_notification_id HAVING COUNT(*) > 1 + )`, + }, + { + table: 'idempotency_keys', + sql: `SELECT COUNT(*) AS count FROM idempotency_keys k + WHERE k.idempotency_key IS NULL OR k.request_hash IS NULL + OR k.response_notification_id IS NULL OR k.response_data IS NULL + OR k.expires_at IS NULL OR k.status NOT IN ('PROCESSED', 'EXPIRED') + OR NOT EXISTS (SELECT 1 FROM scheduled_notifications n WHERE n.id = k.response_notification_id)`, + }, + { + table: 'notification_archive', + sql: `SELECT COUNT(*) AS count FROM notification_archive + WHERE original_id IS NULL OR payload IS NULL + OR notification_type NOT IN ('discord', 'email', 'webhook', 'sms') + OR target_recipient IS NULL OR execute_at IS NULL OR created_at IS NULL + OR status NOT IN ('COMPLETED', 'FAILED', 'CANCELLED', 'EXPIRED') + OR retry_count < 0 OR archived_at IS NULL`, + }, + ]; + + if (await hasColumn(db, 'scheduled_notifications', 'expires_at')) { + audits[0].sql = `${audits[0].sql} + OR (expires_at IS NOT NULL AND datetime(expires_at) IS NULL)`; + } + if (await hasColumn(db, 'notification_archive', 'expires_at')) { + audits[5].sql = `${audits[5].sql} + OR (expires_at IS NOT NULL AND datetime(expires_at) IS NULL)`; + } + + const violations: string[] = []; + for (const audit of audits) { + const count = await getCount(db, audit.sql); + if (count > 0) violations.push(`${audit.table}: ${count} row(s)`); + } + if (violations.length > 0) { + throw new Error( + `Migration 005 aborted before rebuilding tables. Repair invalid legacy data first: ${violations.join('; ')}`, + ); + } +} + +const migration = { + id: '005', + name: 'notification-expiration', + up: async (db: sqlite3.Database) => { + await auditLegacyRows(db); + + const notificationExpiresAt = (await hasColumn(db, 'scheduled_notifications', 'expires_at')) + ? 'expires_at' + : 'NULL'; + const archiveExpiresAt = (await hasColumn(db, 'notification_archive', 'expires_at')) + ? 'expires_at' + : 'NULL'; + + await run( + db, + `CREATE TABLE scheduled_notifications_v005 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + payload TEXT NOT NULL, + payload_hash TEXT, + notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), + target_recipient TEXT NOT NULL, + execute_at DATETIME NOT NULL, + expires_at DATETIME, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + status VARCHAR(20) NOT NULL DEFAULT 'PENDING' CHECK (status IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'CANCELLED', 'EXPIRED')), + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), + max_retries INTEGER NOT NULL DEFAULT 3 CHECK (max_retries >= 0), + processing_started_at DATETIME, + processing_completed_at DATETIME, + processor_id VARCHAR(100), + lock_expires_at DATETIME, + last_error TEXT, + error_details TEXT, + event_id TEXT, + contract_address TEXT, + priority INTEGER NOT NULL DEFAULT 5 CHECK (priority BETWEEN 1 AND 10), + metadata TEXT, + next_retry_at DATETIME + )`, + ); + await run( + db, + `INSERT INTO scheduled_notifications_v005 ( + id, payload, payload_hash, notification_type, target_recipient, execute_at, expires_at, + created_at, updated_at, status, retry_count, max_retries, processing_started_at, + processing_completed_at, processor_id, lock_expires_at, last_error, error_details, + event_id, contract_address, priority, metadata, next_retry_at + ) SELECT id, payload, payload_hash, notification_type, target_recipient, execute_at, + ${notificationExpiresAt}, created_at, updated_at, status, retry_count, max_retries, + processing_started_at, processing_completed_at, processor_id, lock_expires_at, + last_error, error_details, event_id, contract_address, priority, metadata, next_retry_at + FROM scheduled_notifications`, + ); + + await run( + db, + `CREATE TABLE notification_execution_log_v005 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + scheduled_notification_id INTEGER NOT NULL, + execution_attempt INTEGER NOT NULL CHECK (execution_attempt > 0), + execution_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + status VARCHAR(20) NOT NULL CHECK (status IN ('SUCCESS', 'FAILED', 'RETRY')), + error_message TEXT, + response_data TEXT, + duration_ms INTEGER CHECK (duration_ms IS NULL OR duration_ms >= 0), + FOREIGN KEY (scheduled_notification_id) REFERENCES scheduled_notifications_v005(id) ON DELETE CASCADE + )`, + ); + await run(db, 'INSERT INTO notification_execution_log_v005 SELECT * FROM notification_execution_log'); + + await run( + db, + `CREATE TABLE dead_letter_queue_v005 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + scheduled_notification_id INTEGER NOT NULL UNIQUE, + notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), + target_recipient TEXT NOT NULL, + payload TEXT NOT NULL, + failure_reason TEXT NOT NULL, + error_details TEXT, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + last_retried_at DATETIME, + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), + FOREIGN KEY (scheduled_notification_id) REFERENCES scheduled_notifications_v005(id) ON DELETE CASCADE + )`, + ); + await run(db, 'INSERT INTO dead_letter_queue_v005 SELECT * FROM dead_letter_queue'); + + await run( + db, + `CREATE TABLE idempotency_keys_v005 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + idempotency_key TEXT NOT NULL UNIQUE, + request_hash TEXT NOT NULL, + response_notification_id INTEGER NOT NULL, + response_data TEXT NOT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + expires_at DATETIME NOT NULL, + status VARCHAR(20) NOT NULL DEFAULT 'PROCESSED' CHECK (status IN ('PROCESSED', 'EXPIRED')), + FOREIGN KEY (response_notification_id) REFERENCES scheduled_notifications_v005(id) ON DELETE CASCADE + )`, + ); + await run(db, 'INSERT INTO idempotency_keys_v005 SELECT * FROM idempotency_keys'); + + await run( + db, + `CREATE TABLE notification_archive_v005 ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + original_id INTEGER NOT NULL, + payload TEXT NOT NULL, + notification_type VARCHAR(50) NOT NULL CHECK (notification_type IN ('discord', 'email', 'webhook', 'sms')), + target_recipient TEXT NOT NULL, + execute_at DATETIME NOT NULL, + expires_at DATETIME, + created_at DATETIME NOT NULL, + processing_completed_at DATETIME, + status VARCHAR(20) NOT NULL CHECK (status IN ('COMPLETED', 'FAILED', 'CANCELLED', 'EXPIRED')), + retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), + last_error TEXT, + event_id TEXT, + contract_address TEXT, + metadata TEXT, + archived_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP + )`, + ); + await run( + db, + `INSERT INTO notification_archive_v005 ( + id, original_id, payload, notification_type, target_recipient, execute_at, expires_at, + created_at, processing_completed_at, status, retry_count, last_error, event_id, + contract_address, metadata, archived_at + ) SELECT id, original_id, payload, notification_type, target_recipient, execute_at, + ${archiveExpiresAt}, created_at, processing_completed_at, status, retry_count, + last_error, event_id, contract_address, metadata, archived_at + FROM notification_archive`, + ); + + for (const table of [ + 'notification_execution_log', + 'dead_letter_queue', + 'idempotency_keys', + 'notification_archive', + 'scheduled_notifications', + ]) { + await run(db, `DROP TABLE ${table}`); + } + + await run(db, 'ALTER TABLE scheduled_notifications_v005 RENAME TO scheduled_notifications'); + await run(db, 'ALTER TABLE notification_execution_log_v005 RENAME TO notification_execution_log'); + await run(db, 'ALTER TABLE dead_letter_queue_v005 RENAME TO dead_letter_queue'); + await run(db, 'ALTER TABLE idempotency_keys_v005 RENAME TO idempotency_keys'); + await run(db, 'ALTER TABLE notification_archive_v005 RENAME TO notification_archive'); + + const indexes = [ + 'CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_status ON scheduled_notifications(status)', + "CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_status_execute_at ON scheduled_notifications(status, execute_at) WHERE status = 'PENDING'", + "CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_lock_expires ON scheduled_notifications(lock_expires_at, status) WHERE status = 'PROCESSING'", + "CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_next_retry_at ON scheduled_notifications(next_retry_at, status) WHERE status = 'PENDING'", + 'CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_created_at ON scheduled_notifications(created_at)', + 'CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_event_id ON scheduled_notifications(event_id) WHERE event_id IS NOT NULL', + 'CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_target ON scheduled_notifications(target_recipient, status)', + "CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_claim ON scheduled_notifications(status, priority, execute_at) WHERE status = 'PENDING'", + 'CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_processor_lock ON scheduled_notifications(processor_id, status, lock_expires_at)', + 'CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_status_created ON scheduled_notifications(status, created_at)', + 'CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_type_status ON scheduled_notifications(notification_type, status)', + 'CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_contract ON scheduled_notifications(contract_address) WHERE contract_address IS NOT NULL', + 'CREATE INDEX IF NOT EXISTS idx_execution_log_notification_id ON notification_execution_log(scheduled_notification_id)', + 'CREATE INDEX IF NOT EXISTS idx_execution_log_execution_time ON notification_execution_log(execution_time)', + 'CREATE INDEX IF NOT EXISTS idx_execution_log_status_execution_time ON notification_execution_log(status, execution_time)', + 'CREATE INDEX IF NOT EXISTS idx_execution_log_notification_attempt ON notification_execution_log(scheduled_notification_id, execution_attempt)', + 'CREATE INDEX IF NOT EXISTS idx_dead_letter_queue_created_at ON dead_letter_queue(created_at)', + 'CREATE INDEX IF NOT EXISTS idx_dead_letter_queue_notification_type ON dead_letter_queue(notification_type)', + 'CREATE INDEX IF NOT EXISTS idx_idempotency_keys_key ON idempotency_keys(idempotency_key)', + 'CREATE INDEX IF NOT EXISTS idx_idempotency_keys_expires_at ON idempotency_keys(expires_at)', + 'CREATE INDEX IF NOT EXISTS idx_idempotency_keys_created_at ON idempotency_keys(created_at)', + 'CREATE INDEX IF NOT EXISTS idx_archive_original_id ON notification_archive(original_id)', + 'CREATE INDEX IF NOT EXISTS idx_archive_archived_at ON notification_archive(archived_at)', + 'CREATE INDEX IF NOT EXISTS idx_archive_status ON notification_archive(status)', + 'CREATE INDEX IF NOT EXISTS idx_archive_contract_address ON notification_archive(contract_address) WHERE contract_address IS NOT NULL', + 'CREATE INDEX IF NOT EXISTS idx_archive_event_id ON notification_archive(event_id) WHERE event_id IS NOT NULL', + ]; + for (const index of indexes) await run(db, index); + + await run( + db, + `CREATE TRIGGER IF NOT EXISTS update_scheduled_notifications_timestamp + AFTER UPDATE ON scheduled_notifications + FOR EACH ROW BEGIN + UPDATE scheduled_notifications SET updated_at = CURRENT_TIMESTAMP WHERE id = NEW.id; + END`, + ); + + const foreignKeyViolations = await getRows(db, 'PRAGMA foreign_key_check'); + if (foreignKeyViolations.length > 0) { + throw new Error( + `Migration 005 detected ${foreignKeyViolations.length} foreign-key violation(s) after rebuilding tables`, + ); + } + }, + down: async () => { + throw new Error('Migration 005 cannot be safely rolled back; restore a database backup instead'); + }, +}; + +export default migration; \ No newline at end of file diff --git a/listener/src/services/archive-service.ts b/listener/src/services/archive-service.ts index 14207c0c..9d500ae9 100644 --- a/listener/src/services/archive-service.ts +++ b/listener/src/services/archive-service.ts @@ -32,6 +32,7 @@ interface NotificationRow { notification_type: string; target_recipient: string; execute_at: string; + expires_at: string | null; created_at: string; processing_completed_at: string | null; status: string; @@ -142,12 +143,12 @@ export class ArchiveService { */ async archiveProcessedById(id: number): Promise { const row = await this.db.get( - `SELECT id, payload, notification_type, target_recipient, execute_at, - created_at, processing_completed_at, status, retry_count, + `SELECT id, payload, notification_type, target_recipient, execute_at, + expires_at, created_at, processing_completed_at, status, retry_count, last_error, event_id, contract_address, metadata FROM scheduled_notifications WHERE id = ? - AND status IN ('COMPLETED','FAILED','CANCELLED')`, + AND status IN ('COMPLETED','FAILED','CANCELLED','EXPIRED')`, [id], ); @@ -163,6 +164,7 @@ export class ArchiveService { notificationType: row.notification_type, targetRecipient: row.target_recipient, executeAt: row.execute_at, + expiresAt: row.expires_at, createdAt: row.created_at, processingCompletedAt: row.processing_completed_at, status: row.status, @@ -194,11 +196,11 @@ export class ArchiveService { }); const rows = await this.db.all( - `SELECT id, payload, notification_type, target_recipient, execute_at, - created_at, processing_completed_at, status, retry_count, + `SELECT id, payload, notification_type, target_recipient, execute_at, + expires_at, created_at, processing_completed_at, status, retry_count, last_error, event_id, contract_address, metadata FROM scheduled_notifications - WHERE status IN ('COMPLETED','FAILED','CANCELLED') + WHERE status IN ('COMPLETED','FAILED','CANCELLED','EXPIRED') AND processing_completed_at IS NOT NULL AND processing_completed_at < ? ORDER BY processing_completed_at ASC @@ -221,6 +223,7 @@ export class ArchiveService { notificationType: r.notification_type, targetRecipient: r.target_recipient, executeAt: r.execute_at, + expiresAt: r.expires_at, createdAt: r.created_at, processingCompletedAt: r.processing_completed_at, status: r.status, @@ -247,6 +250,7 @@ export class ArchiveService { completed: rows.filter((r) => r.status === 'COMPLETED').length, failed: rows.filter((r) => r.status === 'FAILED').length, cancelled: rows.filter((r) => r.status === 'CANCELLED').length, + expired: rows.filter((r) => r.status === 'EXPIRED').length, }, }); }); diff --git a/listener/src/services/archive-store.ts b/listener/src/services/archive-store.ts index e2d35e13..cbc28f69 100644 --- a/listener/src/services/archive-store.ts +++ b/listener/src/services/archive-store.ts @@ -17,6 +17,7 @@ export interface ArchivedNotification { notificationType: string; targetRecipient: string; executeAt: string; + expiresAt: string | null; createdAt: string; processingCompletedAt: string | null; status: string; @@ -36,6 +37,7 @@ interface ArchiveRow { notification_type: string; target_recipient: string; execute_at: string; + expires_at: string | null; created_at: string; processing_completed_at: string | null; status: string; @@ -73,6 +75,7 @@ function mapRow(row: ArchiveRow): ArchivedNotification { notificationType: row.notification_type, targetRecipient: row.target_recipient, executeAt: row.execute_at, + expiresAt: row.expires_at, createdAt: row.created_at, processingCompletedAt: row.processing_completed_at, status: row.status, @@ -89,7 +92,7 @@ export class ArchiveStore { constructor(private readonly db: Database) {} /** - * Insert a batch of completed/failed/cancelled notifications into the + * Insert a batch of terminal notifications into the * archive. Returns the number of rows inserted. */ async insertBatch( @@ -99,6 +102,7 @@ export class ArchiveStore { notificationType: string; targetRecipient: string; executeAt: string; + expiresAt: string | null; createdAt: string; processingCompletedAt: string | null; status: string; @@ -116,15 +120,16 @@ export class ArchiveStore { await this.db.run( `INSERT INTO notification_archive (original_id, payload, notification_type, target_recipient, - execute_at, created_at, processing_completed_at, + execute_at, expires_at, created_at, processing_completed_at, status, retry_count, last_error, event_id, contract_address, metadata) - VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)`, + VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)`, [ r.originalId, r.payload, r.notificationType, r.targetRecipient, r.executeAt, + r.expiresAt, r.createdAt, r.processingCompletedAt, r.status, diff --git a/listener/src/services/cleanup-service.ts b/listener/src/services/cleanup-service.ts index 3679617a..5742b085 100644 --- a/listener/src/services/cleanup-service.ts +++ b/listener/src/services/cleanup-service.ts @@ -65,7 +65,7 @@ export class CleanupService { const [notifResult, rateLimitResult, executionLogResult, processedEventResult] = await Promise.all([ this.db.run( `DELETE FROM scheduled_notifications - WHERE status IN ('COMPLETED','FAILED','CANCELLED') + WHERE status IN ('COMPLETED','FAILED','CANCELLED','EXPIRED') AND processing_completed_at < ?`, [notificationCutoff], ), diff --git a/listener/src/services/notification-api.ts b/listener/src/services/notification-api.ts index 82b4304e..4eddb05f 100644 --- a/listener/src/services/notification-api.ts +++ b/listener/src/services/notification-api.ts @@ -22,6 +22,25 @@ const PRIORITY_MIN = 1; const PRIORITY_MAX = 10; import { buildRetryStatisticsPayload } from './retry-statistics'; +function parseExpiresAt(value: Date | string | number): Date | null { + if (value instanceof Date) { + return Number.isNaN(value.getTime()) ? null : value; + } + + if (typeof value === 'number') { + if (!Number.isFinite(value)) return null; + const milliseconds = Math.abs(value) < 100_000_000_000 ? value * 1000 : value; + const date = new Date(milliseconds); + return Number.isNaN(date.getTime()) ? null : date; + } + + if (!/^\d{4}-\d{2}-\d{2}(?:T\d{2}:\d{2}(?::\d{2}(?:\.\d{1,9})?)?(?:Z|[+-]\d{2}:\d{2})?)?$/.test(value)) { + return null; + } + const date = new Date(value); + return Number.isNaN(date.getTime()) ? null : date; +} + /** * High-level API for scheduling notifications * This is the main interface that application code should use @@ -95,6 +114,11 @@ export class NotificationAPI { if (input.contractAddress !== undefined) { v.check(isNonEmptyString(input.contractAddress), 'contractAddress', 'must be a non-empty string'); } + if (input.expiresAt !== undefined && input.expiresAt !== null) { + const expiresAt = parseExpiresAt(input.expiresAt); + v.check(expiresAt !== null, 'expiresAt', 'must be a valid ISO-8601 timestamp or epoch'); + if (expiresAt) input = { ...input, expiresAt }; + } if (input.metadata !== undefined) { v.check(isPlainObject(input.metadata), 'metadata', 'must be an object'); } diff --git a/listener/src/services/notification-scheduler.ts b/listener/src/services/notification-scheduler.ts index 8d0cda9b..be591207 100644 --- a/listener/src/services/notification-scheduler.ts +++ b/listener/src/services/notification-scheduler.ts @@ -139,8 +139,15 @@ export class NotificationScheduler { return; } + const activeNotifications: ScheduledNotification[] = []; + for (const notification of notifications) { + if (await this.expireIfPastDeadline(notification, requestId)) continue; + activeNotifications.push(notification); + } + if (activeNotifications.length === 0) return; + const batchRejection = this.batchValidator.rejectIfInvalid( - this.toValidationBatch(notifications) + this.toValidationBatch(activeNotifications) ); if (batchRejection) { @@ -150,7 +157,7 @@ export class NotificationScheduler { errors: batchRejection.errors, }); - for (const notification of notifications) { + for (const notification of activeNotifications) { await this.repository.markAsFailedOrRetry( notification.id!, new Error(`Batch validation failed: ${batchRejection.errors.map((e) => e.message).join('; ')}`), @@ -163,7 +170,7 @@ export class NotificationScheduler { logger.info('Processing batch of scheduled notifications', { requestId, - count: notifications.length, + count: activeNotifications.length, processorId: this.processorId, }); @@ -172,10 +179,10 @@ export class NotificationScheduler { if (workerManager.isShutdownInProgress()) { logger.info('Shutdown in progress - releasing unprocessed notifications', { requestId, - count: notifications.length, + count: activeNotifications.length, }); // Release locks on unprocessed notifications - for (const notification of notifications) { + for (const notification of activeNotifications) { await this.repository.markAsFailedOrRetry( notification.id!, new Error('Scheduler shutting down'), @@ -188,7 +195,7 @@ export class NotificationScheduler { // Process each notification with job tracking + monitoring const jobMonitor = getJobMonitor(); - for (const notification of notifications) { + for (const notification of activeNotifications) { const jobId = `notification-${notification.id}`; if (!workerManager.startJob(jobId)) { // Shutdown is in progress, don't process new jobs @@ -218,7 +225,7 @@ export class NotificationScheduler { logger.info('Scheduler batch complete', { requestId, processorId: this.processorId, - count: notifications.length, + count: activeNotifications.length, durationMs: Date.now() - batchStart, }); } catch (error) { @@ -244,6 +251,8 @@ export class NotificationScheduler { const jobMonitor = getJobMonitor(); try { + if (await this.expireIfPastDeadline(notification, requestId, jobId)) return; + logger.info('Processing scheduled notification', { requestId, id: notification.id, @@ -389,6 +398,33 @@ export class NotificationScheduler { } } + private async expireIfPastDeadline( + notification: ScheduledNotification, + requestId: string, + jobId?: string, + ): Promise { + if (!notification.expiresAt || notification.expiresAt.getTime() > Date.now()) return false; + + const errorMessage = 'Notification expired before delivery'; + await this.repository.markAsExpired(notification.id!); + await this.repository.logExecution({ + scheduledNotificationId: notification.id!, + executionAttempt: notification.retryCount + 1, + executionTime: new Date(), + status: 'FAILED', + errorMessage, + durationMs: 0, + }); + if (jobId) { + getJobMonitor().failJob(jobId, errorMessage, { notificationId: notification.id }); + } + logger.info('Expired notification skipped before delivery', { + requestId, + id: notification.id, + }); + return true; + } + /** * Execute notification delivery based on type. * diff --git a/listener/src/services/retry-scheduler.ts b/listener/src/services/retry-scheduler.ts index fa65685b..bddfd6aa 100644 --- a/listener/src/services/retry-scheduler.ts +++ b/listener/src/services/retry-scheduler.ts @@ -217,6 +217,24 @@ export class RetryScheduler { const executionAttempt = priorFailures + 1; const startMs = Date.now(); + if (notification.expiresAt && notification.expiresAt.getTime() <= startMs) { + const errorMessage = 'Notification expired before delivery'; + await this.repository.markAsExpired(notification.id!); + await this.repository.logExecution({ + scheduledNotificationId: notification.id!, + executionAttempt, + executionTime: new Date(startMs), + status: 'FAILED', + errorMessage, + durationMs: 0, + }); + logger.info('Expired retry skipped before delivery', { + requestId, + id: notification.id, + }); + return; + } + logger.info('Retrying notification', { requestId, id: notification.id, diff --git a/listener/src/services/scheduled-notification-repository.ts b/listener/src/services/scheduled-notification-repository.ts index bfbc0a02..bd376e8f 100644 --- a/listener/src/services/scheduled-notification-repository.ts +++ b/listener/src/services/scheduled-notification-repository.ts @@ -22,6 +22,7 @@ export class ScheduledNotificationRepository { constructor( private db: Database, statsCache?: NotificationStatsCache, + private defaultTtlSeconds: number = 0, ) { this.statsCache = statsCache ?? getStatsCache(); } @@ -36,9 +37,9 @@ export class ScheduledNotificationRepository { const sql = ` INSERT INTO scheduled_notifications ( - payload, payload_hash, notification_type, target_recipient, execute_at, + payload, payload_hash, notification_type, target_recipient, execute_at, expires_at, max_retries, event_id, contract_address, priority, metadata - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) `; const serializedPayload = compressPayload(input.payload); @@ -49,6 +50,13 @@ export class ScheduledNotificationRepository { input.notificationType, input.targetRecipient, input.executeAt.toISOString(), + input.expiresAt !== undefined + ? input.expiresAt === null + ? null + : this.normalizeExpiration(input.expiresAt) + : this.defaultTtlSeconds > 0 + ? new Date(Date.now() + this.defaultTtlSeconds * 1000).toISOString() + : null, input.maxRetries ?? 3, input.eventId ?? null, input.contractAddress ?? null, @@ -71,6 +79,22 @@ export class ScheduledNotificationRepository { return result.lastID; } + private normalizeExpiration(value: Date | string | number): string { + const milliseconds = + typeof value === 'number' && Math.abs(value) < 100_000_000_000 + ? value * 1000 + : value instanceof Date + ? value.getTime() + : typeof value === 'number' + ? value + : Date.parse(value); + const date = new Date(milliseconds); + if (Number.isNaN(date.getTime())) { + throw new Error('expiresAt must be a valid ISO-8601 timestamp or epoch'); + } + return date.toISOString(); + } + /** * Fetch pending notifications due for execution with distributed locking * Uses atomic update to prevent race conditions @@ -244,6 +268,28 @@ export class ScheduledNotificationRepository { logger.info('Notification marked as completed', { requestId, id }); } + /** Mark a notification expired and release its processing lock. */ + async markAsExpired(id: number): Promise { + const now = new Date().toISOString(); + const errorMessage = 'Notification expired before delivery'; + await this.db.run( + `UPDATE scheduled_notifications + SET status = ?, last_error = ?, error_details = ?, processing_completed_at = ?, + updated_at = ?, processor_id = NULL, lock_expires_at = NULL, next_retry_at = NULL + WHERE id = ? AND status = ?`, + [ + NotificationStatus.EXPIRED, + errorMessage, + JSON.stringify({ message: errorMessage, timestamp: now }), + now, + now, + id, + NotificationStatus.PROCESSING, + ], + ); + this.statsCache.invalidate(); + } + /** * Mark notification as failed or retry (sets next_retry_at for backoff scheduling) */ @@ -541,7 +587,7 @@ export class ScheduledNotificationRepository { const sql = ` DELETE FROM scheduled_notifications - WHERE status IN (?, ?, ?) + WHERE status IN (?, ?, ?, ?) AND updated_at < ? `; @@ -549,6 +595,7 @@ export class ScheduledNotificationRepository { NotificationStatus.COMPLETED, NotificationStatus.FAILED, NotificationStatus.CANCELLED, + NotificationStatus.EXPIRED, cutoff.toISOString(), ]); @@ -816,6 +863,7 @@ export class ScheduledNotificationRepository { notificationType: row.notification_type as any, targetRecipient: row.target_recipient, executeAt: parseUtc(row.execute_at) as Date, + expiresAt: parseUtc(row.expires_at) ?? null, createdAt: parseUtc(row.created_at), updatedAt: parseUtc(row.updated_at), status: row.status as NotificationStatus, diff --git a/listener/src/tests/notification-scheduler.test.ts b/listener/src/tests/notification-scheduler.test.ts index 9c4cb794..ee5797cf 100644 --- a/listener/src/tests/notification-scheduler.test.ts +++ b/listener/src/tests/notification-scheduler.test.ts @@ -1,3 +1,5 @@ +jest.mock('../utils/request-id', () => ({ generateRequestId: () => 'test-request-id' })); + import { Database } from '../database/database'; import { ScheduledNotificationRepository } from '../services/scheduled-notification-repository'; import { NotificationScheduler } from '../services/notification-scheduler'; @@ -68,6 +70,128 @@ describe('NotificationScheduler', () => { expect(notification!.retryCount).toBe(0); }); + test('does not deliver an expired scheduled notification', async () => { + const id = await repository.create({ + payload: { event: {}, contractConfig: {}, message: 'Expired notification' }, + notificationType: NotificationType.DISCORD, + targetRecipient: 'webhook-b', + executeAt: new Date(Date.now() - 1000), + expiresAt: new Date(Date.now() - 1000), + }); + const discordService = { sendEventNotification: jest.fn().mockResolvedValue(true) } as any; + + scheduler = new NotificationScheduler( + repository, + { + enabled: true, + pollIntervalMs: 1000, + lockTimeoutMs: 30000, + processorId: 'expired-dispatch-test', + batchSize: 10, + timingBufferMs: 1000, + }, + discordService, + ); + + await (scheduler as any).processPendingNotifications(); + + expect(discordService.sendEventNotification).toHaveBeenCalledTimes(0); + expect((await repository.getById(id))!.status).toBe(NotificationStatus.EXPIRED); + }); + + test('does not deliver a notification that expires between retries', async () => { + const id = await repository.create({ + payload: { event: {}, contractConfig: {}, message: 'Retry expiration' }, + notificationType: NotificationType.DISCORD, + targetRecipient: 'test-webhook', + executeAt: new Date(Date.now() - 60_000), + expiresAt: new Date(Date.now() + 60_000), + maxRetries: 3, + }); + const discordService = { + sendEventNotification: jest.fn().mockRejectedValueOnce(new Error('temporary outage')), + } as any; + scheduler = new NotificationScheduler( + repository, + { + enabled: true, + pollIntervalMs: 1000, + lockTimeoutMs: 30000, + processorId: 'retry-expiration-initial-test', + batchSize: 10, + timingBufferMs: 1000, + retryDelayMs: 100, + }, + discordService, + ); + + await (scheduler as any).processPendingNotifications(); + expect(discordService.sendEventNotification).toHaveBeenCalledTimes(1); + + await db.run( + 'UPDATE scheduled_notifications SET expires_at = ?, next_retry_at = ? WHERE id = ?', + [new Date(Date.now() - 1000).toISOString(), new Date(Date.now() - 1000).toISOString(), id], + ); + + const { RetryScheduler } = await import('../services/retry-scheduler'); + const retryScheduler = new RetryScheduler( + repository, + { + enabled: true, + pollIntervalMs: 1000, + lockTimeoutMs: 30000, + processorId: 'retry-expiration-test', + batchSize: 10, + baseDelayMs: 100, + multiplier: 2, + maxDelayMs: 1000, + jitter: false, + }, + discordService, + ); + + await retryScheduler.runOnce(); + + expect(discordService.sendEventNotification).toHaveBeenCalledTimes(1); + expect((await repository.getById(id))!.status).toBe(NotificationStatus.EXPIRED); + }); + + test('delivers non-expired and NULL-expiration notifications normally', async () => { + const activeId = await repository.create({ + payload: { event: {}, contractConfig: {}, message: 'Still valid' }, + notificationType: NotificationType.DISCORD, + targetRecipient: 'webhook-a', + executeAt: new Date(Date.now() - 1000), + expiresAt: new Date(Date.now() + 60_000), + }); + const neverExpiresId = await repository.create({ + payload: { event: {}, contractConfig: {}, message: 'Never expires' }, + notificationType: NotificationType.DISCORD, + targetRecipient: 'webhook-b', + executeAt: new Date(Date.now() - 1000), + expiresAt: null, + }); + const discordService = { sendEventNotification: jest.fn().mockResolvedValue(true) } as any; + scheduler = new NotificationScheduler( + repository, + { + enabled: true, + pollIntervalMs: 1000, + lockTimeoutMs: 30000, + processorId: 'valid-expiration-test', + batchSize: 10, + timingBufferMs: 1000, + }, + discordService, + ); + + await (scheduler as any).processPendingNotifications(); + + expect(discordService.sendEventNotification).toHaveBeenCalledTimes(2); + expect((await repository.getById(activeId))!.status).toBe(NotificationStatus.COMPLETED); + expect((await repository.getById(neverExpiresId))!.status).toBe(NotificationStatus.COMPLETED); + }); + test('should fetch and lock pending notifications', async () => { const executeAt = new Date(Date.now() - 1000); // Past time @@ -89,7 +213,7 @@ describe('NotificationScheduler', () => { const notifications = await repository.fetchAndLockPendingNotifications( processorId, 30000, - 10 + 10, ); expect(notifications.length).toBe(2); @@ -127,7 +251,7 @@ describe('NotificationScheduler', () => { timingBufferMs: 1000, retryDelayMs: 2000, }, - discordService + discordService, ); await (scheduler as any).processPendingNotifications(); @@ -157,7 +281,7 @@ describe('NotificationScheduler', () => { maxDelayMs: 1000, jitter: false, }, - discordService + discordService, ); await retryScheduler.runOnce(); @@ -181,12 +305,12 @@ describe('NotificationScheduler', () => { const processor1 = await repository.fetchAndLockPendingNotifications( 'processor-1', 30000, - 10 + 10, ); const processor2 = await repository.fetchAndLockPendingNotifications( 'processor-2', 30000, - 10 + 10, ); expect(processor1.length).toBe(1); @@ -334,8 +458,10 @@ describe('NotificationScheduler', () => { notificationType: NotificationType.DISCORD, targetRecipient: 'test-webhook', executeAt: pastDate, - }) - ).rejects.toThrow('executeAt must be a future timestamp — the provided date has already expired'); + }), + ).rejects.toThrow( + 'executeAt must be a future timestamp — the provided date has already expired', + ); }); test('should reject execution time equal to now', async () => { @@ -348,7 +474,7 @@ describe('NotificationScheduler', () => { notificationType: NotificationType.DISCORD, targetRecipient: 'test-webhook', executeAt: now, - }) + }), ).rejects.toThrow('executeAt must be a future timestamp'); }); @@ -359,7 +485,7 @@ describe('NotificationScheduler', () => { notificationType: NotificationType.DISCORD, targetRecipient: 'test-webhook', executeAt: new Date('not-a-date'), - }) + }), ).rejects.toThrow('executeAt must be a valid date'); }); @@ -383,7 +509,7 @@ describe('NotificationScheduler', () => { 'https://discord.com/webhook/test', { content: 'Hello World' }, executeAt, - { priority: 1, maxRetries: 5 } + { priority: 1, maxRetries: 5 }, ); expect(id).toBeGreaterThan(0); diff --git a/listener/src/types/index.ts b/listener/src/types/index.ts index 1d212d13..d8f81ace 100644 --- a/listener/src/types/index.ts +++ b/listener/src/types/index.ts @@ -64,6 +64,8 @@ export interface Config { cleanup?: AppCleanupConfig; analytics?: AnalyticsConfig; expiration?: ExpirationConfig; + /** Default scheduled-notification lifetime in seconds; zero disables expiry. */ + notificationDefaultTtlSeconds?: number; backfill?: BackfillConfig; logging?: LoggingConfig; api?: ApiConfig; diff --git a/listener/src/types/scheduled-notification.ts b/listener/src/types/scheduled-notification.ts index cdaa92d8..10e738b3 100644 --- a/listener/src/types/scheduled-notification.ts +++ b/listener/src/types/scheduled-notification.ts @@ -8,6 +8,7 @@ export enum NotificationStatus { COMPLETED = 'COMPLETED', FAILED = 'FAILED', CANCELLED = 'CANCELLED', + EXPIRED = 'EXPIRED', } export enum NotificationType { @@ -24,6 +25,7 @@ export interface ScheduledNotification { notificationType: NotificationType; targetRecipient: string; executeAt: Date; + expiresAt?: Date | null; createdAt?: Date; updatedAt?: Date; status: NotificationStatus; @@ -48,6 +50,7 @@ export interface CreateScheduledNotificationInput { notificationType: NotificationType; targetRecipient: string; executeAt: Date; + expiresAt?: Date | string | number | null; maxRetries?: number; eventId?: string; contractAddress?: string; @@ -62,6 +65,7 @@ export interface ScheduledNotificationRow { notification_type: string; target_recipient: string; execute_at: string; + expires_at: string | null; created_at: string; updated_at: string; status: string; From 93e25ebb4e1593474c3357b1afdee66923c97e3b Mon Sep 17 00:00:00 2001 From: Najib Ishiyaku Njidda Date: Sat, 3 Oct 2026 08:35:02 +0000 Subject: [PATCH 3/5] fix(migrations): carry deduplication_key through 005 scheduled_notifications rebuild. The bootstrap schema gained deduplication_key via merged 003-notification-deduplication-key after 005 was written; the rebuild omitted it, which would have dropped the column on existing DBs. Column now copied verbatim into the _v005 CREATE and the INSERT/SELECT lists, with a legacy-survival assertion and a bootstrap-vs-rebuild column-parity test to catch future drift. Also fixed a pre-existing missing parenthesis in getRows callback that blocked Prettier formatting. --- .../005-notification-expiration.test.ts | 35 ++++++++++++++++--- .../migrations/005-notification-expiration.ts | 26 +++++++++----- 2 files changed, 49 insertions(+), 12 deletions(-) diff --git a/listener/src/migrations/005-notification-expiration.test.ts b/listener/src/migrations/005-notification-expiration.test.ts index 567ef897..16d55753 100644 --- a/listener/src/migrations/005-notification-expiration.test.ts +++ b/listener/src/migrations/005-notification-expiration.test.ts @@ -1,4 +1,6 @@ import * as sqlite3 from 'sqlite3'; +import { readFileSync } from 'fs'; +import { join } from 'path'; import { MigrationRunner } from '../database/migration-system'; import migration from './005-notification-expiration'; @@ -47,7 +49,8 @@ async function createLegacySchema(db: sqlite3.Database): Promise { retry_count INTEGER NOT NULL DEFAULT 0, max_retries INTEGER NOT NULL DEFAULT 3, processing_started_at DATETIME, processing_completed_at DATETIME, processor_id VARCHAR(100), lock_expires_at DATETIME, last_error TEXT, error_details TEXT, event_id TEXT, - contract_address TEXT, priority INTEGER NOT NULL DEFAULT 5, metadata TEXT, next_retry_at DATETIME + contract_address TEXT, priority INTEGER NOT NULL DEFAULT 5, metadata TEXT, next_retry_at DATETIME, + deduplication_key TEXT ); CREATE TABLE notification_execution_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, scheduled_notification_id INTEGER NOT NULL, @@ -88,10 +91,10 @@ async function seedValidRows(db: sqlite3.Database): Promise { (id, payload, payload_hash, notification_type, target_recipient, execute_at, created_at, updated_at, status, retry_count, max_retries, processing_started_at, processing_completed_at, processor_id, lock_expires_at, last_error, error_details, event_id, contract_address, priority, - metadata, next_retry_at) + metadata, next_retry_at, deduplication_key) VALUES (1, '{"n":1}', 'hash', 'discord', 'user-1', '2026-01-01', 'created', 'updated', 'FAILED', 1, 3, NULL, 'completed', 'worker-1', NULL, 'error', NULL, 'event-1', 'contract-1', - 5, '{"meta":true}', NULL)`, + 5, '{"meta":true}', NULL, 'dedupe-key-1')`, ); await run( db, @@ -122,7 +125,7 @@ const preservedRows = async (db: sqlite3.Database) => ({ `SELECT id, payload, payload_hash, notification_type, target_recipient, execute_at, created_at, updated_at, status, retry_count, max_retries, processing_started_at, processing_completed_at, processor_id, lock_expires_at, last_error, error_details, event_id, contract_address, priority, - metadata, next_retry_at FROM scheduled_notifications ORDER BY id`, + metadata, next_retry_at, deduplication_key FROM scheduled_notifications ORDER BY id`, ), notification_execution_log: await all(db, 'SELECT * FROM notification_execution_log ORDER BY id'), dead_letter_queue: await all(db, 'SELECT * FROM dead_letter_queue ORDER BY id'), @@ -149,6 +152,19 @@ async function expectNoRebuildTables(db: sqlite3.Database): Promise { expect(tables).toHaveLength(0); } +function bootstrapScheduledNotificationColumns(): string[] { + const schema = readFileSync(join(__dirname, '../database/schema.sql'), 'utf8'); + const tableDefinition = schema.match( + /CREATE TABLE IF NOT EXISTS scheduled_notifications \(([\s\S]*?)\n\);/i, + ); + if (!tableDefinition) throw new Error('scheduled_notifications missing from bootstrap schema'); + + return Array.from( + tableDefinition[1].matchAll(/^\s{2}([a-z][a-z0-9_]*)\s+/gim), + (match) => match[1], + ); +} + describe('migration 005 notification expiration', () => { let db: sqlite3.Database; @@ -169,6 +185,17 @@ describe('migration 005 notification expiration', () => { await applyMigration(db); expect(await preservedRows(db)).toEqual(before); + expect( + (await all<{ name: string }>(db, 'PRAGMA table_info(scheduled_notifications)')).map( + ({ name }) => name, + ), + ).toEqual(bootstrapScheduledNotificationColumns()); + expect( + await all<{ deduplication_key: string }>( + db, + 'SELECT deduplication_key FROM scheduled_notifications WHERE id = 1', + ), + ).toEqual([{ deduplication_key: 'dedupe-key-1' }]); expect( await all<{ expires_at: string | null }>( db, diff --git a/listener/src/migrations/005-notification-expiration.ts b/listener/src/migrations/005-notification-expiration.ts index 70476a6d..d8d4047d 100644 --- a/listener/src/migrations/005-notification-expiration.ts +++ b/listener/src/migrations/005-notification-expiration.ts @@ -17,7 +17,7 @@ function getCount(db: sqlite3.Database, sql: string): Promise { function getRows(db: sqlite3.Database, sql: string): Promise { return new Promise((resolve, reject) => { - db.all(sql, (error, rows) => (error ? reject(error) : resolve(rows)); + db.all(sql, (error, rows) => (error ? reject(error) : resolve(rows))); }); } @@ -142,7 +142,8 @@ const migration = { contract_address TEXT, priority INTEGER NOT NULL DEFAULT 5 CHECK (priority BETWEEN 1 AND 10), metadata TEXT, - next_retry_at DATETIME + next_retry_at DATETIME, + deduplication_key TEXT )`, ); await run( @@ -151,11 +152,12 @@ const migration = { id, payload, payload_hash, notification_type, target_recipient, execute_at, expires_at, created_at, updated_at, status, retry_count, max_retries, processing_started_at, processing_completed_at, processor_id, lock_expires_at, last_error, error_details, - event_id, contract_address, priority, metadata, next_retry_at + event_id, contract_address, priority, metadata, next_retry_at, deduplication_key ) SELECT id, payload, payload_hash, notification_type, target_recipient, execute_at, ${notificationExpiresAt}, created_at, updated_at, status, retry_count, max_retries, processing_started_at, processing_completed_at, processor_id, lock_expires_at, - last_error, error_details, event_id, contract_address, priority, metadata, next_retry_at + last_error, error_details, event_id, contract_address, priority, metadata, next_retry_at, + deduplication_key FROM scheduled_notifications`, ); @@ -173,7 +175,10 @@ const migration = { FOREIGN KEY (scheduled_notification_id) REFERENCES scheduled_notifications_v005(id) ON DELETE CASCADE )`, ); - await run(db, 'INSERT INTO notification_execution_log_v005 SELECT * FROM notification_execution_log'); + await run( + db, + 'INSERT INTO notification_execution_log_v005 SELECT * FROM notification_execution_log', + ); await run( db, @@ -253,7 +258,10 @@ const migration = { } await run(db, 'ALTER TABLE scheduled_notifications_v005 RENAME TO scheduled_notifications'); - await run(db, 'ALTER TABLE notification_execution_log_v005 RENAME TO notification_execution_log'); + await run( + db, + 'ALTER TABLE notification_execution_log_v005 RENAME TO notification_execution_log', + ); await run(db, 'ALTER TABLE dead_letter_queue_v005 RENAME TO dead_letter_queue'); await run(db, 'ALTER TABLE idempotency_keys_v005 RENAME TO idempotency_keys'); await run(db, 'ALTER TABLE notification_archive_v005 RENAME TO notification_archive'); @@ -305,8 +313,10 @@ const migration = { } }, down: async () => { - throw new Error('Migration 005 cannot be safely rolled back; restore a database backup instead'); + throw new Error( + 'Migration 005 cannot be safely rolled back; restore a database backup instead', + ); }, }; -export default migration; \ No newline at end of file +export default migration; From 94a4231a8343f154fa19ee4d40ac7f387e610bf8 Mon Sep 17 00:00:00 2001 From: Najib Ishiyaku Njidda Date: Sat, 3 Oct 2026 09:43:23 +0000 Subject: [PATCH 4/5] test(migrations): compare 004 rebuild column parity against pre-004 source, not final bootstrap. Migration 005 owns expires_at, so the intermediate post-004 state legitimately lacks it; asserting against bootstrap would go red on main once #893 merges. 004's invariant is exact preservation of its source column set. --- .../004-database-data-integrity.test.ts | 28 ++++++------------- 1 file changed, 8 insertions(+), 20 deletions(-) diff --git a/listener/src/migrations/004-database-data-integrity.test.ts b/listener/src/migrations/004-database-data-integrity.test.ts index 89802885..70208a0f 100644 --- a/listener/src/migrations/004-database-data-integrity.test.ts +++ b/listener/src/migrations/004-database-data-integrity.test.ts @@ -1,6 +1,4 @@ import * as sqlite3 from 'sqlite3'; -import { readFileSync } from 'fs'; -import { join } from 'path'; import { Database } from '../database/database'; import { MigrationRunner } from '../database/migration-system'; import migration from './004-database-data-integrity'; @@ -189,19 +187,6 @@ const PRESERVED_TABLES = [ 'notification_archive', ]; -function bootstrapScheduledNotificationColumns(): string[] { - const schema = readFileSync(join(__dirname, '../database/schema.sql'), 'utf8'); - const tableDefinition = schema.match( - /CREATE TABLE IF NOT EXISTS scheduled_notifications \(([\s\S]*?)\n\);/i, - ); - if (!tableDefinition) throw new Error('scheduled_notifications missing from bootstrap schema'); - - return Array.from( - tableDefinition[1].matchAll(/^\s{2}([a-z][a-z0-9_]*)\s+/gim), - (match) => match[1], - ); -} - describe('migration 004 database data integrity', () => { let db: sqlite3.Database; @@ -223,6 +208,9 @@ describe('migration 004 database data integrity', () => { for (const table of PRESERVED_TABLES) { before[table] = await all(db, `SELECT * FROM ${table} ORDER BY id`); } + const sourceColumns = ( + await all<{ name: string }>(db, 'PRAGMA table_info(scheduled_notifications)') + ).map(({ name }) => name); await run(db, 'PRAGMA foreign_keys = OFF'); const runner = new MigrationRunner(db, ''); @@ -235,11 +223,11 @@ describe('migration 004 database data integrity', () => { for (const table of PRESERVED_TABLES) { expect(await all(db, `SELECT * FROM ${table} ORDER BY id`)).toEqual(before[table]); } - expect( - (await all<{ name: string }>(db, 'PRAGMA table_info(scheduled_notifications)')).map( - ({ name }) => name, - ), - ).toEqual(bootstrapScheduledNotificationColumns()); + const rebuiltColumns = ( + await all<{ name: string }>(db, 'PRAGMA table_info(scheduled_notifications)') + ).map(({ name }) => name); + // Compare to pre-004 columns; later migration 005 owns expires_at. + expect(rebuiltColumns.sort()).toEqual(sourceColumns.sort()); expect( await all<{ deduplication_key: string }>( db, From 8368708f1d2cfc6b25652564a68a83eecd0941e4 Mon Sep 17 00:00:00 2001 From: Najib Ishiyaku Njidda Date: Sat, 3 Oct 2026 10:48:01 +0000 Subject: [PATCH 5/5] chore: merge upstream/main into feat/844 (union dead-letter status, RPC/circuit-breaker config, retry rework; DEAD_LETTERED added to 004 CHECKs for bootstrap/rebuild parity) --- listener/src/config.ts | 44 +++++++++++++++++- listener/src/database/migration-system.ts | 21 --------- listener/src/database/schema.sql | 3 +- .../004-database-data-integrity.test.ts | 8 ++++ .../migrations/004-database-data-integrity.ts | 4 +- listener/src/services/retry-scheduler.ts | 45 ++++--------------- .../scheduled-notification-repository.ts | 11 ----- listener/src/types/index.ts | 2 +- 8 files changed, 63 insertions(+), 75 deletions(-) diff --git a/listener/src/config.ts b/listener/src/config.ts index afba57bf..30c97892 100644 --- a/listener/src/config.ts +++ b/listener/src/config.ts @@ -1,7 +1,31 @@ -import { Config, ContractConfig, DiscordConfig, WebhookSecret, AppCleanupConfig, EventQueueConfig, RetrySchedulerOptions, AnalyticsConfig, ExpirationConfig, ApiKey, BackfillConfig, LoggingConfig, ApiConfig, RpcRateLimitConfig } from './types'; +import { + Config, + ContractConfig, + DiscordConfig, + WebhookSecret, + AppCleanupConfig, + EventQueueConfig, + RetrySchedulerOptions, + RetryPolicyOptions, + AnalyticsConfig, + ExpirationConfig, + ApiKey, + BackfillConfig, + LoggingConfig, + ApiConfig, + RpcFallbackConfig, + RpcRateLimitConfig, +} from './types'; +import { CircuitBreakerConfig } from './services/circuit-breaker'; import { validateCorsOrigin, CorsValidationError } from './utils/cors-validator'; import { validateSecrets } from './config/validate-secrets'; import { ConfigurationSchemaValidator, APP_CONFIG_SCHEMA } from './config-schema'; +import { + DEFAULT_RETRYABLE_FAILURE_TYPES, + RETRY_FAILURE_TYPES, + RetryFailureType, + parseRetryableFailureTypes, +} from './services/retry-policy'; import { SUPPORTED_LOG_FORMATS, SUPPORTED_LOG_LEVELS, @@ -358,6 +382,24 @@ function loadRpcRateLimitConfig(): RpcRateLimitConfig { }; } +function loadRpcFallbackConfig(fallbackUrls: string[]): RpcFallbackConfig { + return { + fallbackUrls, + failureThreshold: parseIntegerEnv('RPC_FAILURE_THRESHOLD', '3'), + cooldownMs: parseIntegerEnv('RPC_COOLDOWN_MS', '60000'), + requestTimeoutMs: parseIntegerEnv('RPC_REQUEST_TIMEOUT_MS', '10000'), + maxRetries: parseOptionalIntegerEnv('RPC_MAX_RETRIES'), + }; +} + +function loadCircuitBreakerConfig(): CircuitBreakerConfig { + return { + failureThreshold: parseIntegerEnv('CIRCUIT_BREAKER_FAILURE_THRESHOLD', '5'), + recoveryTimeoutMs: parseIntegerEnv('CIRCUIT_BREAKER_RECOVERY_TIMEOUT_MS', '60000'), + successThreshold: parseIntegerEnv('CIRCUIT_BREAKER_SUCCESS_THRESHOLD', '2'), + }; +} + export function loadConfig(): Config { validateRequiredEnvVars(); diff --git a/listener/src/database/migration-system.ts b/listener/src/database/migration-system.ts index c1443d01..66a58e0f 100644 --- a/listener/src/database/migration-system.ts +++ b/listener/src/database/migration-system.ts @@ -89,15 +89,8 @@ export class MigrationRunner { } async getAppliedMigrations(): Promise { -<<<<<<< HEAD const rows = await this.all<{ id: string }>('SELECT id FROM migrations ORDER BY applied_at'); return rows.map((row) => row.id); -======= - const rows = await this.db.all<{ id: string }>( - 'SELECT id FROM migrations ORDER BY applied_at' - ); - return (rows as unknown as { id: string }[]).map((row) => row.id); ->>>>>>> upstream/main } async applyMigration(migration: Migration): Promise { @@ -110,23 +103,9 @@ export class MigrationRunner { logger.info(`Migration ${migration.id} (${migration.name}) applied successfully`); } catch (error) { try { -<<<<<<< HEAD await this.run('ROLLBACK'); } catch (rollbackError) { logger.error(`Migration ${migration.id} rollback failed`, { error: rollbackError }); -======= - 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 as Error)?.message ?? String(error)}`); - throw error; ->>>>>>> upstream/main } logger.error(`Migration ${migration.id} failed, rolling back`, { error }); throw error; diff --git a/listener/src/database/schema.sql b/listener/src/database/schema.sql index 8737108f..7998e38f 100644 --- a/listener/src/database/schema.sql +++ b/listener/src/database/schema.sql @@ -17,7 +17,7 @@ CREATE TABLE IF NOT EXISTS scheduled_notifications ( updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, -- Status tracking - status VARCHAR(20) NOT NULL DEFAULT 'PENDING' CHECK (status IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'CANCELLED')), + status VARCHAR(20) NOT NULL DEFAULT 'PENDING' CHECK (status IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'DEAD_LETTERED', 'CANCELLED')), retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), max_retries INTEGER NOT NULL DEFAULT 3 CHECK (max_retries >= 0), @@ -453,4 +453,3 @@ CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_type_lower_created CREATE INDEX IF NOT EXISTS idx_processed_events_type_lower_processed ON processed_events(LOWER(event_type), processed_at); - diff --git a/listener/src/migrations/004-database-data-integrity.test.ts b/listener/src/migrations/004-database-data-integrity.test.ts index 70208a0f..5f03a59d 100644 --- a/listener/src/migrations/004-database-data-integrity.test.ts +++ b/listener/src/migrations/004-database-data-integrity.test.ts @@ -252,6 +252,11 @@ describe('migration 004 database data integrity', () => { VALUES ('{}', 'discord', 'user-2', '2026-01-01', 'UNKNOWN')`, ), ).rejects.toThrow(/CHECK constraint failed/); + await run( + db, + `INSERT INTO scheduled_notifications (payload, notification_type, target_recipient, execute_at, status) + VALUES ('{}', 'discord', 'user-dead-lettered', '2026-01-01', 'DEAD_LETTERED')`, + ); await expect( run( db, @@ -285,6 +290,9 @@ describe('migration 004 database data integrity', () => { try { const pragma = await freshDb.get<{ foreign_keys: number }>('PRAGMA foreign_keys'); expect(pragma?.foreign_keys).toBe(1); + await freshDb.run(`INSERT INTO scheduled_notifications + (payload, notification_type, target_recipient, execute_at, status) + VALUES ('{}', 'discord', 'user-dead-lettered', '2026-01-01', 'DEAD_LETTERED')`); await expect( freshDb.run(`INSERT INTO scheduled_notifications (payload, notification_type, target_recipient, execute_at, status) diff --git a/listener/src/migrations/004-database-data-integrity.ts b/listener/src/migrations/004-database-data-integrity.ts index 83291a18..a4aacc2f 100644 --- a/listener/src/migrations/004-database-data-integrity.ts +++ b/listener/src/migrations/004-database-data-integrity.ts @@ -12,7 +12,7 @@ const tables = [ execute_at DATETIME NOT NULL, created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, - status VARCHAR(20) NOT NULL DEFAULT 'PENDING' CHECK (status IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'CANCELLED')), + status VARCHAR(20) NOT NULL DEFAULT 'PENDING' CHECK (status IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'DEAD_LETTERED', 'CANCELLED')), retry_count INTEGER NOT NULL DEFAULT 0 CHECK (retry_count >= 0), max_retries INTEGER NOT NULL DEFAULT 3 CHECK (max_retries >= 0), processing_started_at DATETIME, @@ -182,7 +182,7 @@ const audits: Array<{ table: string; sql: string }> = [ sql: `SELECT COUNT(*) AS count FROM scheduled_notifications WHERE payload IS NULL OR notification_type IS NULL OR notification_type NOT IN ('discord', 'email', 'webhook', 'sms') OR target_recipient IS NULL OR execute_at IS NULL - OR status IS NULL OR status NOT IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'CANCELLED') + OR status IS NULL OR status NOT IN ('PENDING', 'PROCESSING', 'COMPLETED', 'FAILED', 'DEAD_LETTERED', 'CANCELLED') OR retry_count IS NULL OR retry_count < 0 OR max_retries IS NULL OR max_retries < 0 OR priority IS NULL OR priority NOT BETWEEN 1 AND 10`, }, diff --git a/listener/src/services/retry-scheduler.ts b/listener/src/services/retry-scheduler.ts index 87f7864e..72fd5d70 100644 --- a/listener/src/services/retry-scheduler.ts +++ b/listener/src/services/retry-scheduler.ts @@ -333,20 +333,6 @@ export class RetryScheduler { }); }); } - const isFinalAttempt = priorFailures + 1 >= notification.maxRetries; - - const nextRetryAt = isFinalAttempt - ? undefined - : new Date( - Date.now() + - calculateBackoffDelay( - priorFailures, - this.config.baseDelayMs, - this.config.multiplier, - this.config.maxDelayMs, - this.config.jitter - ) - ); const failureType = classifyError(err); const decision = this.policy.evaluate( @@ -431,17 +417,6 @@ export class RetryScheduler { ), degradedCapabilities: [], }; - if (!this.discordService) { - throw new DeliveryError( - 'Discord service not configured', - RetryFailureType.ConfigurationError, - ); - } - return this.discordService.sendEventNotification( - payload.event, - payload.contractConfig, - `retry-${notification.id}-${requestId}` - ); case 'webhook': { const targetUrl: string = notification.targetRecipient; @@ -456,6 +431,14 @@ export class RetryScheduler { payload, `retry-${notification.id}-${requestId}`, ); + if (!result.success) { + const failureType = classifyHttpStatus(result.statusCode); + throw new DeliveryError( + result.errorReason ?? `Webhook delivery failed (HTTP ${result.statusCode ?? 'unknown'})`, + failureType, + { statusCode: result.statusCode }, + ); + } return { success: result.success, degradedCapabilities: [], @@ -465,18 +448,6 @@ export class RetryScheduler { errorCode: result.errorCode, errorMessage: result.errorReason, }; - if (!result.success) { - // Surface the specific reason so it lands in markAsFailedOrRetry's - // error details, and tag it with a failure type so the retry policy - // can tell permanent rejections (4xx) from transient ones (5xx). - const failureType = classifyHttpStatus(result.statusCode); - throw new DeliveryError( - result.errorReason ?? `Webhook delivery failed (HTTP ${result.statusCode ?? 'unknown'})`, - failureType, - { statusCode: result.statusCode }, - ); - } - return true; } default: diff --git a/listener/src/services/scheduled-notification-repository.ts b/listener/src/services/scheduled-notification-repository.ts index 161044a8..5fcac875 100644 --- a/listener/src/services/scheduled-notification-repository.ts +++ b/listener/src/services/scheduled-notification-repository.ts @@ -75,17 +75,6 @@ export class ScheduledNotificationRepository { executeAt: input.executeAt, type: input.notificationType, }); - const result = await this.db.run(sql, params); - - // Invalidate stats cache after creation - this.statsCache.invalidate(); - - logger.info('Scheduled notification created', { - requestId, - id: result.lastID, - executeAt: input.executeAt, - type: input.notificationType, - }); return result.lastID; } catch (err) { diff --git a/listener/src/types/index.ts b/listener/src/types/index.ts index cb1638c4..03003e67 100644 --- a/listener/src/types/index.ts +++ b/listener/src/types/index.ts @@ -143,6 +143,7 @@ export interface SchedulerConfig { lockTimeoutMs: number; processorId?: string; batchSize: number; + concurrency: number; timingBufferMs: number; } @@ -277,4 +278,3 @@ export interface RpcRateLimitConfig { /** Delay in ms to apply when throttled (default: 1000). */ throttleDelayMs: number; } -