diff --git a/NOTIFICATION_LIFECYCLE.md b/NOTIFICATION_LIFECYCLE.md index 0913b337..53589c57 100644 --- a/NOTIFICATION_LIFECYCLE.md +++ b/NOTIFICATION_LIFECYCLE.md @@ -310,7 +310,7 @@ Service-level callers can query through `DeliveryReceiptRepository` using ### 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). @@ -411,6 +411,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/dashboard/src/pages/NotificationSearchPage.tsx b/dashboard/src/pages/NotificationSearchPage.tsx index 09c533e2..20b952ad 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 ff87a2a6..7d3f7137 100644 --- a/listener/.env.example +++ b/listener/.env.example @@ -70,6 +70,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 901a1f1e..8b157a6b 100644 --- a/listener/src/api/events-server.ts +++ b/listener/src/api/events-server.ts @@ -1189,6 +1189,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 afba57bf..57a785da 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, @@ -48,6 +72,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 parseStrictIntegerEnv(name: string, defaultValue: string): number { const rawValue = trimEnv(name); if (rawValue !== undefined && !/^-?\d+$/.test(rawValue)) { @@ -358,6 +394,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(); @@ -439,6 +493,7 @@ export function loadConfig(): Config { cleanup: loadCleanupConfig(), analytics: loadAnalyticsConfig(), expiration: loadExpirationConfig(), + notificationDefaultTtlSeconds: loadNotificationDefaultTtlSeconds(), backfill: loadBackfillConfig(), rpcRateLimit: loadRpcRateLimitConfig(), logging: loadLoggingConfig(), diff --git a/listener/src/database/archive-schema.sql b/listener/src/database/archive-schema.sql index 546e1095..c5f016c1 100644 --- a/listener/src/database/archive-schema.sql +++ b/listener/src/database/archive-schema.sql @@ -11,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/database.ts b/listener/src/database/database.ts index b754cf01..d3f284c5 100644 --- a/listener/src/database/database.ts +++ b/listener/src/database/database.ts @@ -167,6 +167,7 @@ export class Database { } reject(err); } else { + this.db = handle; logger.info('Connected to SQLite database', { path: this.dbPath }); // Wait (instead of failing with SQLITE_BUSY) when another connection // — e.g. a second listener instance on the same file — holds the @@ -178,7 +179,6 @@ export class Database { if (pragmaErr) { reject(pragmaErr); } else { - this.db = handle; resolve(); } }); diff --git a/listener/src/database/migration-system.ts b/listener/src/database/migration-system.ts index fdf4375b..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' - ); - return (rows as unknown as { id: string }[]).map((row) => row.id); + 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 as Error)?.message ?? String(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 d5e65843..e6e969d3 100644 --- a/listener/src/database/schema.sql +++ b/listener/src/database/schema.sql @@ -8,18 +8,19 @@ 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 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', -- 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', 'EXPIRED', 'DEAD_LETTERED')), + 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 +35,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 deduplication_key TEXT -- Caller-supplied key; duplicate inserts with the same key are silently skipped @@ -76,14 +77,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 ); @@ -97,12 +98,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 ); @@ -152,12 +153,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 @@ -196,7 +197,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, @@ -234,19 +235,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 ); @@ -276,12 +277,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 @@ -307,7 +308,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 ); @@ -326,12 +327,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 @@ -353,12 +354,40 @@ 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, + expires_at DATETIME, + created_at DATETIME NOT NULL, + processing_completed_at DATETIME, + status VARCHAR(20) NOT NULL CHECK (status IN ('COMPLETED', 'FAILED', 'CANCELLED', 'EXPIRED', 'DEAD_LETTERED')), + 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) @@ -425,5 +454,3 @@ CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_type_lower_created -- Processed-event search: case-insensitive type filter + processed_at sort CREATE INDEX IF NOT EXISTS idx_processed_events_type_lower_processed ON processed_events(LOWER(event_type), processed_at); - - diff --git a/listener/src/index.ts b/listener/src/index.ts index 090e7b12..89497178 100644 --- a/listener/src/index.ts +++ b/listener/src/index.ts @@ -70,7 +70,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, + ); deliveryReceiptRepository = new DeliveryReceiptRepository(db); 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..b210aeda --- /dev/null +++ b/listener/src/migrations/004-database-data-integrity.test.ts @@ -0,0 +1,314 @@ +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, + deduplication_key TEXT + ); + 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, 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, 'dedupe-key-1')`, + ); + 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`); + } + 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, ''); + 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 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, + 'SELECT deduplication_key FROM scheduled_notifications WHERE id = 1', + ), + ).toEqual([{ deduplication_key: 'dedupe-key-1' }]); + + 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 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 run( + db, + `INSERT INTO notification_archive + (original_id, payload, notification_type, target_recipient, execute_at, created_at, status) + VALUES (2, '{}', 'discord', 'user-archive-dead-lettered', '2026-01-01', '2026-01-01', 'DEAD_LETTERED')`, + ); + 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 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 freshDb.run(`INSERT INTO notification_archive + (original_id, payload, notification_type, target_recipient, execute_at, created_at, status) + VALUES (1, '{}', 'discord', 'user-archive-dead-lettered', '2026-01-01', '2026-01-01', 'DEAD_LETTERED')`); + 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..b4f59fd1 --- /dev/null +++ b/listener/src/migrations/004-database-data-integrity.ts @@ -0,0 +1,447 @@ +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', '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, + 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, + deduplication_key TEXT + )`, + }, + { + 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', 'DEAD_LETTERED')), + 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', '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`, + }, + { + 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', 'DEAD_LETTERED') + 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; 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..3c9031e8 --- /dev/null +++ b/listener/src/migrations/005-notification-expiration.test.ts @@ -0,0 +1,318 @@ +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'; + +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, + deduplication_key TEXT + ); + 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, 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, 'dedupe-key-1')`, + ); + 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, 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'), + 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); +} + +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; + + 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<{ 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, + '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 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 run( + db, + `INSERT INTO notification_archive + (original_id, payload, notification_type, target_recipient, execute_at, created_at, status) + VALUES (2, '{}', 'discord', 'user-archive-dead-lettered', '2026-01-01', '2026-01-01', 'DEAD_LETTERED')`, + ); + 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..a76f7ccf --- /dev/null +++ b/listener/src/migrations/005-notification-expiration.ts @@ -0,0 +1,322 @@ +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', 'DEAD_LETTERED') + 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', 'DEAD_LETTERED') + 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', 'DEAD_LETTERED', '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, + deduplication_key TEXT + )`, + ); + 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, 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, + deduplication_key + 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', 'DEAD_LETTERED')), + 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; 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 2be182de..8f69bc9e 100644 --- a/listener/src/services/notification-api.ts +++ b/listener/src/services/notification-api.ts @@ -26,6 +26,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 @@ -113,6 +132,11 @@ export class NotificationAPI { '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 5af8c260..1a224e65 100644 --- a/listener/src/services/notification-scheduler.ts +++ b/listener/src/services/notification-scheduler.ts @@ -145,8 +145,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) { @@ -156,7 +163,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('; ')}`), @@ -169,7 +176,7 @@ export class NotificationScheduler { logger.info('Processing batch of scheduled notifications', { requestId, - count: notifications.length, + count: activeNotifications.length, processorId: this.processorId, }); @@ -178,10 +185,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'), @@ -198,7 +205,7 @@ export class NotificationScheduler { // the original serial behavior exactly. const jobMonitor = getJobMonitor(); const concurrency = Math.max(1, Math.floor(this.config.concurrency ?? 1)); - const queue = [...notifications]; + const queue = [...activeNotifications]; const processOne = async (notification: ScheduledNotification): Promise => { const jobId = `notification-${notification.id}`; if (!workerManager.startJob(jobId)) { @@ -244,7 +251,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) { @@ -303,6 +310,8 @@ export class NotificationScheduler { let receiptRecorded = false; try { + if (await this.expireIfPastDeadline(notification, requestId, jobId)) return; + logger.info('Processing scheduled notification', { requestId, id: notification.id, @@ -468,6 +477,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. * @@ -601,4 +637,4 @@ export class NotificationScheduler { errorMessage: result.success ? null : result.errorMessage ?? 'Provider delivery failed', }); } -} +} \ No newline at end of file diff --git a/listener/src/services/retry-scheduler.ts b/listener/src/services/retry-scheduler.ts index 87f7864e..52b74a60 100644 --- a/listener/src/services/retry-scheduler.ts +++ b/listener/src/services/retry-scheduler.ts @@ -286,6 +286,24 @@ export class RetryScheduler { const startMs = Date.now(); let receiptRecorded = false; + 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, @@ -333,20 +351,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 +435,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 +449,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 +466,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..dbf586d5 100644 --- a/listener/src/services/scheduled-notification-repository.ts +++ b/listener/src/services/scheduled-notification-repository.ts @@ -27,6 +27,7 @@ export class ScheduledNotificationRepository { constructor( private db: Database, statsCache?: NotificationStatsCache, + private defaultTtlSeconds: number = 0, ) { this.statsCache = statsCache ?? getStatsCache(); } @@ -40,12 +41,15 @@ export class ScheduledNotificationRepository { const payloadJson = JSON.stringify(input.payload); const secret = process.env.PAYLOAD_INTEGRITY_SECRET; const payloadHash = secret ? hashPayload(payloadJson, secret) : null; + const now = new Date().toISOString(); const sql = ` INSERT INTO scheduled_notifications ( - payload, payload_hash, notification_type, target_recipient, execute_at, - max_retries, event_id, contract_address, priority, metadata, deduplication_key - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + 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, deduplication_key + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) `; const serializedPayload = compressPayload(input.payload); @@ -56,11 +60,29 @@ 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, + now, + now, + NotificationStatus.PENDING, + 0, input.maxRetries ?? 3, + null, + null, + null, + null, + null, + null, input.eventId ?? null, input.contractAddress ?? null, input.priority ?? 5, input.metadata ? JSON.stringify(input.metadata) : null, + null, input.deduplicationKey ?? null, ]; @@ -75,17 +97,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) { @@ -112,6 +123,22 @@ export class ScheduledNotificationRepository { } } + 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(); + } + /** * Dequeue pending notifications due for execution with distributed locking. * Uses a single atomic UPDATE ... WHERE id IN (SELECT ...) so two workers @@ -333,6 +360,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) */ @@ -640,14 +689,16 @@ export class ScheduledNotificationRepository { const sql = ` DELETE FROM scheduled_notifications - WHERE status IN (?, ?, ?) + WHERE status IN (?, ?, ?, ?, ?) AND updated_at < ? `; const result = await this.db.run(sql, [ NotificationStatus.COMPLETED, NotificationStatus.FAILED, + NotificationStatus.DEAD_LETTERED, NotificationStatus.CANCELLED, + NotificationStatus.EXPIRED, cutoff.toISOString(), ]); @@ -963,6 +1014,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 b2b87288..b7d8b899 100644 --- a/listener/src/tests/notification-scheduler.test.ts +++ b/listener/src/tests/notification-scheduler.test.ts @@ -1,4 +1,4 @@ -jest.mock('../utils/request-id', () => ({ generateRequestId: () => 'scheduler-test-request-id' })); +jest.mock('../utils/request-id', () => ({ generateRequestId: () => 'test-request-id' })); import { Database } from '../database/database'; import { ScheduledNotificationRepository } from '../services/scheduled-notification-repository'; @@ -166,6 +166,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 @@ -187,7 +309,7 @@ describe('NotificationScheduler', () => { const notifications = await repository.fetchAndLockPendingNotifications( processorId, 30000, - 10 + 10, ); expect(notifications.length).toBe(2); @@ -225,7 +347,7 @@ describe('NotificationScheduler', () => { timingBufferMs: 1000, retryDelayMs: 2000, }, - discordService + discordService, ); await (scheduler as any).processPendingNotifications(); @@ -255,7 +377,7 @@ describe('NotificationScheduler', () => { maxDelayMs: 1000, jitter: false, }, - discordService + discordService, ); await retryScheduler.runOnce(); @@ -279,12 +401,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); @@ -365,8 +487,8 @@ describe('NotificationScheduler', () => { await repository.markAsFailedOrRetry(id, error, 2, 3); const notification = await repository.getById(id); - expect(notification!.status).toBe(NotificationStatus.FAILED); - expect(notification!.retryCount).toBe(3); + expect(notification!.status).toBe(NotificationStatus.DEAD_LETTERED); + expect(notification!.retryCount).toBe(2); }); test('should cancel pending notification', async () => { @@ -448,8 +570,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 () => { @@ -462,7 +586,7 @@ describe('NotificationScheduler', () => { notificationType: NotificationType.DISCORD, targetRecipient: 'test-webhook', executeAt: now, - }) + }), ).rejects.toThrow('executeAt must be a future timestamp'); }); @@ -473,7 +597,7 @@ describe('NotificationScheduler', () => { notificationType: NotificationType.DISCORD, targetRecipient: 'test-webhook', executeAt: new Date('not-a-date'), - }) + }), ).rejects.toThrow('executeAt must be a valid date'); }); @@ -497,7 +621,7 @@ describe('NotificationScheduler', () => { 'https://discord.com/webhook/test', { content: 'Hello World' }, executeAt, - { priority: 1, maxRetries: 5 } + { priority: 1, maxRetries: 5 }, ); expect(id).toBeGreaterThan(0); @@ -600,7 +724,7 @@ describe('Stale cache regression tests', () => { await repository.markAsFailedOrRetry(id, new Error('Max retries exceeded'), 2, 2); const notification = await repository.getById(id); - expect(notification?.status).toBe(NotificationStatus.FAILED); + expect(notification?.status).toBe(NotificationStatus.DEAD_LETTERED); expect(notification?.updatedAt).toBeDefined(); expect(notification?.updatedAt).toBeInstanceOf(Date); expect(isNaN(notification?.updatedAt?.getTime() ?? NaN)).toBe(false); diff --git a/listener/src/types/index.ts b/listener/src/types/index.ts index cb1638c4..b4a5e6e9 100644 --- a/listener/src/types/index.ts +++ b/listener/src/types/index.ts @@ -1,4 +1,5 @@ import type { RetryFailureType } from '../services/retry-policy'; +import type { CircuitBreakerConfig } from '../services/circuit-breaker'; import * as StellarSDK from '@stellar/stellar-sdk'; export interface NotificationProvider { @@ -82,6 +83,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; @@ -143,6 +146,7 @@ export interface SchedulerConfig { lockTimeoutMs: number; processorId?: string; batchSize: number; + concurrency: number; timingBufferMs: number; } @@ -277,4 +281,3 @@ export interface RpcRateLimitConfig { /** Delay in ms to apply when throttled (default: 1000). */ throttleDelayMs: number; } - diff --git a/listener/src/types/scheduled-notification.ts b/listener/src/types/scheduled-notification.ts index 10d4c111..64b8bc54 100644 --- a/listener/src/types/scheduled-notification.ts +++ b/listener/src/types/scheduled-notification.ts @@ -9,6 +9,7 @@ export enum NotificationStatus { FAILED = 'FAILED', DEAD_LETTERED = 'DEAD_LETTERED', CANCELLED = 'CANCELLED', + EXPIRED = 'EXPIRED', } export enum NotificationType { @@ -25,6 +26,7 @@ export interface ScheduledNotification { notificationType: NotificationType; targetRecipient: string; executeAt: Date; + expiresAt?: Date | null; createdAt?: Date; updatedAt?: Date; status: NotificationStatus; @@ -53,6 +55,7 @@ export interface CreateScheduledNotificationInput { notificationType: NotificationType; targetRecipient: string; executeAt: Date; + expiresAt?: Date | string | number | null; maxRetries?: number; deduplicationKey?: string; eventId?: string; @@ -68,6 +71,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;