diff --git a/NOTIFICATION_LIFECYCLE.md b/NOTIFICATION_LIFECYCLE.md index a11e2367..b72fc9fe 100644 --- a/NOTIFICATION_LIFECYCLE.md +++ b/NOTIFICATION_LIFECYCLE.md @@ -31,7 +31,8 @@ This document covers: 9. [Completion and Archival](#completion-and-archival) 10. [Dashboard Visibility](#dashboard-visibility) 11. [Developer Notes](#developer-notes) -12. [Troubleshooting](#troubleshooting) +12. [Database Cleanup](#database-cleanup) +13. [Troubleshooting](#troubleshooting) --- @@ -449,6 +450,45 @@ So the dashboard sits **after** off-chain ingestion: contract → listener → A - Batch validation: `POST /api/notifications/validate-batch` plus scheduler pre-process batch checks. +## Database Cleanup + +`DatabaseCleanupJob` runs independently from the notification archiver. It +removes expired idempotency keys, processed-event fingerprints, old dead-letter +records, execution-log rows not associated with pending/processing +notifications, expired rate-limit windows, and old `DEACTIVATED` backpressure +events. `ACTIVATED` backpressure records and all `PENDING`/`PROCESSING` +notifications are retained. Scheduled notifications are never directly +deleted by this job; `ArchiveService` owns their terminal-state archival and +the age-based, status-agnostic purge of `notification_archive`. Metrics +snapshots are retained by `NotificationMetricsRunner`. + +| Setting | Default | Minimum | Purpose | +|---------|---------|---------|---------| +| `CLEANUP_ENABLED` | `true` | `true` / `false` | Enable the scheduled cleanup job | +| `CLEANUP_INTERVAL_MS` | `3600000` | `60000` | Run interval in milliseconds | +| `CLEANUP_RETENTION_DAYS` | `30` | `1` | Global age threshold for cleanup-managed tables | + +Explicit legacy overrides remain available for processed events +(`PROCESSED_EVENT_RETENTION_MS`), execution logs +(`EXECUTION_LOG_RETENTION_MS`), and rate-limit audit records +(`RATE_LIMIT_EVENT_RETENTION_MS`). Idempotency keys are removed when +`expires_at` is past, or when status is `EXPIRED` and `created_at` is older +than retention. A future-dated `PROCESSED` key is never removed. + +Each run logs a correlation `runId`, `perTableDeleted` counts, skipped tables, +failed tables, configured interval and retention, and duration. Missing required +timestamp columns cause that table to be warned and skipped; a table failure is +logged and does not prevent remaining tables from being cleaned. Deletes run in +batches of at most 1,000 rows, each in its own transaction. + +Useful checks: + +```sql +SELECT status, COUNT(*) FROM scheduled_notifications GROUP BY status; +SELECT status, COUNT(*) FROM notification_archive GROUP BY status; +SELECT COUNT(*) FROM idempotency_keys WHERE datetime(expires_at) < datetime('now'); +``` + --- ## Troubleshooting diff --git a/listener/.env.example b/listener/.env.example index 8958248f..56967fa0 100644 --- a/listener/.env.example +++ b/listener/.env.example @@ -288,6 +288,12 @@ EXPIRATION_DEFAULT_MS=86400000 # How often the cleanup job runs (ms). Default: 1 hour. # CLEANUP_INTERVAL_MS=3600000 +# Enable the scheduled database cleanup job. Default: true. +# CLEANUP_ENABLED=true + +# Global database cleanup retention. Default: 30 days. +# CLEANUP_RETENTION_DAYS=30 + # How long processed notifications are retained before deletion (ms). Default: 7 days. # NOTIFICATION_RETENTION_MS=604800000 diff --git a/listener/src/config-schema.ts b/listener/src/config-schema.ts index eae05734..637ed93f 100644 --- a/listener/src/config-schema.ts +++ b/listener/src/config-schema.ts @@ -237,7 +237,14 @@ export const APP_CONFIG_SCHEMA: ConfigSchema = { snapshotRetentionDays: { type: 'number', min: 1 }, }, cleanup: { + enabled: { type: 'boolean' }, intervalMs: { type: 'number', min: 60000 }, + retentionDays: { type: 'number', min: 1 }, + retentionOverridesMs: { + processedEvents: { type: 'number', min: 60000 }, + executionLogs: { type: 'number', min: 60000 }, + rateLimitEvents: { type: 'number', min: 60000 }, + }, notificationRetentionMs: { type: 'number', min: 60000 }, rateLimitEventRetentionMs: { type: 'number', min: 60000 }, eventRetentionMs: { type: 'number', min: 60000 }, diff --git a/listener/src/config.ts b/listener/src/config.ts index 9496b660..d979266e 100644 --- a/listener/src/config.ts +++ b/listener/src/config.ts @@ -2,6 +2,7 @@ import { Config, ContractConfig, DiscordConfig, WebhookSecret, AppCleanupConfig, import { validateCorsOrigin, CorsValidationError } from './utils/cors-validator'; import { validateSecrets } from './config/validate-secrets'; import { ConfigurationSchemaValidator, APP_CONFIG_SCHEMA } from './config-schema'; +import { Config, ContractConfig, DiscordConfig, WebhookSecret, AppCleanupConfig, EventQueueConfig, RetrySchedulerOptions, AnalyticsConfig, ExpirationConfig, ApiKey, BackfillConfig, LoggingConfig, ApiConfig } from './types'; import { Config, ContractConfig, DiscordConfig, WebhookSecret, AppCleanupConfig, EventQueueConfig, RetrySchedulerOptions, AnalyticsConfig, ExpirationConfig, ApiKey, BackfillConfig, LoggingConfig, ApiConfig, RetryPolicyOptions } from './types'; import { DEFAULT_RETRYABLE_FAILURE_TYPES, @@ -59,6 +60,27 @@ function parseIntegerEnv(name: string, defaultValue: string): number { return parsed; } +function parseStrictIntegerEnv(name: string, defaultValue: string): number { + const rawValue = trimEnv(name); + if (rawValue !== undefined && !/^-?\d+$/.test(rawValue)) { + throw new ConfigError(`${name} must be a valid integer, got "${rawValue}"`); + } + return parseIntegerEnv(name, defaultValue); +} + +function parseBooleanEnv(name: string, defaultValue: boolean): boolean { + const rawValue = trimEnv(name); + if (rawValue === undefined) return defaultValue; + if (rawValue === 'true') return true; + if (rawValue === 'false') return false; + throw new ConfigError(`${name} must be either "true" or "false", got "${rawValue}"`); +} + +function parseOptionalIntegerEnv(name: string): number | undefined { + const rawValue = trimEnv(name); + return rawValue ? parseStrictIntegerEnv(name, rawValue) : undefined; +} + function parseJsonEnv(name: string, defaultValue: string): T { const rawValue = trimEnv(name) ?? defaultValue; try { @@ -196,19 +218,32 @@ function validateApiKeys(value: unknown): ApiKey[] { } function loadCleanupConfig(): AppCleanupConfig { + const processedEventRetentionMs = parseIntegerEnv( + 'PROCESSED_EVENT_RETENTION_MS', + String(30 * 24 * 60 * 60 * 1000), + ); + const executionLogRetentionMs = parseIntegerEnv( + 'EXECUTION_LOG_RETENTION_MS', + String(90 * 24 * 60 * 60 * 1000), + ); + const rateLimitEventRetentionMs = parseIntegerEnv( + 'RATE_LIMIT_EVENT_RETENTION_MS', + String(24 * 60 * 60 * 1000), + ); return { - intervalMs: parseIntegerEnv('CLEANUP_INTERVAL_MS', String(60 * 60 * 1000)), + enabled: parseBooleanEnv('CLEANUP_ENABLED', true), + intervalMs: parseStrictIntegerEnv('CLEANUP_INTERVAL_MS', String(60 * 60 * 1000)), + retentionDays: parseStrictIntegerEnv('CLEANUP_RETENTION_DAYS', '30'), + retentionOverridesMs: { + processedEvents: parseOptionalIntegerEnv('PROCESSED_EVENT_RETENTION_MS'), + executionLogs: parseOptionalIntegerEnv('EXECUTION_LOG_RETENTION_MS'), + rateLimitEvents: parseOptionalIntegerEnv('RATE_LIMIT_EVENT_RETENTION_MS'), + }, notificationRetentionMs: parseIntegerEnv('NOTIFICATION_RETENTION_MS', String(7 * 24 * 60 * 60 * 1000)), - rateLimitEventRetentionMs: parseIntegerEnv('RATE_LIMIT_EVENT_RETENTION_MS', String(24 * 60 * 60 * 1000)), + rateLimitEventRetentionMs, eventRetentionMs: parseIntegerEnv('EVENT_RETENTION_MS', String(24 * 60 * 60 * 1000)), - processedEventRetentionMs: parseIntegerEnv( - 'PROCESSED_EVENT_RETENTION_MS', - String(30 * 24 * 60 * 60 * 1000), - ), - executionLogRetentionMs: parseIntegerEnv( - 'EXECUTION_LOG_RETENTION_MS', - String(90 * 24 * 60 * 60 * 1000), - ), + processedEventRetentionMs, + executionLogRetentionMs, }; } @@ -804,6 +839,11 @@ export function validateConfig(config: Config): void { `(received: ${config.cleanup.processedEventRetentionMs}).`, ); } + if (config.cleanup.retentionDays < 1) { + errors.push( + `CLEANUP_RETENTION_DAYS must be >= 1 (received: ${config.cleanup.retentionDays}).`, + ); + } } // ── Backfill ─────────────────────────────────────────────────────────────── diff --git a/listener/src/index.ts b/listener/src/index.ts index 73512031..f3d74469 100644 --- a/listener/src/index.ts +++ b/listener/src/index.ts @@ -9,7 +9,7 @@ import { NotificationTemplateService } from './services/notification-template-se import { TemplateAuditTrail } from './services/template-audit-trail'; import { getTemplateCache } from './services/notification-template-cache'; import { NotificationAPI } from './services/notification-api'; -import { CleanupService } from './services/cleanup-service'; +import { DatabaseCleanupJob } from './services/database-cleanup-job'; import { ArchiveService } from './services/archive-service'; import { ArchiveStore } from './services/archive-store'; import { loadArchiveConfig } from './services/archive-config'; @@ -49,7 +49,7 @@ async function main() { let templateService: NotificationTemplateService | null = null; let legacyTemplateService: TemplateService | null = null; - let cleanupService: CleanupService | null = null; + let databaseCleanupJob: DatabaseCleanupJob | null = null; let repository: ScheduledNotificationRepository | null = null; let reconciliationEngine: IndexingReconciliationEngine | null = null; let archiveService: ArchiveService | null = null; @@ -79,8 +79,10 @@ async function main() { eventRegistry.setTtlMs(config.cleanup.eventRetentionMs); } - cleanupService = new CleanupService(db, eventRegistry, config.cleanup); - cleanupService.start(); + if (config.cleanup) { + databaseCleanupJob = new DatabaseCleanupJob(db, config.cleanup, eventRegistry); + databaseCleanupJob.start(); + } reconciliationEngine = new IndexingReconciliationEngine({ db, @@ -185,8 +187,8 @@ async function main() { healthMonitor.stop(); } - if (cleanupService) { - await cleanupService.stop(); + if (databaseCleanupJob) { + await databaseCleanupJob.stop(); } if (reconciliationEngine) { diff --git a/listener/src/services/archive.test.ts b/listener/src/services/archive.test.ts index 7e6b3b9d..4dfe79c3 100644 --- a/listener/src/services/archive.test.ts +++ b/listener/src/services/archive.test.ts @@ -383,16 +383,21 @@ describe('ArchiveService', () => { expect(result.archived).toBe(0); }); - it('purges archive rows older than deleteAfterMs', async () => { - // Manually plant an "old" archive row - (db as any).tables.notification_archive.push({ - id: 1, - original_id: 100, - archived_at: new Date(Date.now() - 91 * 24 * 60 * 60 * 1000).toISOString(), - }); + it('purges archive rows by age regardless of status', async () => { + const oldArchivedAt = new Date(Date.now() - 91 * 24 * 60 * 60 * 1000).toISOString(); + (db as any).tables.notification_archive.push( + { id: 1, original_id: 100, status: 'COMPLETED', archived_at: oldArchivedAt }, + { id: 2, original_id: 101, status: 'EXPIRED', archived_at: oldArchivedAt }, + { + id: 3, + original_id: 102, + status: 'FAILED', + archived_at: new Date(Date.now() - 1 * 24 * 60 * 60 * 1000).toISOString(), + }, + ); const result = await service.runCycle(); - expect(result.purged).toBe(1); - expect(db.archiveCount()).toBe(0); + expect(result.purged).toBe(2); + expect(db.archiveCount()).toBe(1); }); it('skips purge when deleteAfterMs is 0', async () => { diff --git a/listener/src/services/database-cleanup-job.test.ts b/listener/src/services/database-cleanup-job.test.ts new file mode 100644 index 00000000..033dc052 --- /dev/null +++ b/listener/src/services/database-cleanup-job.test.ts @@ -0,0 +1,363 @@ +import { Database } from '../database/database'; +import { loadConfig, validateConfig, ConfigError } from '../config'; +import logger from '../utils/logger'; +import { resetWorkerManager } from './worker-manager'; +import { DatabaseCleanupJob } from './database-cleanup-job'; + +jest.mock('../utils/logger', () => ({ + __esModule: true, + parseLogFormat: jest.fn(() => 'pretty'), + parseLogLevel: jest.fn(() => 'info'), + SUPPORTED_LOG_FORMATS: ['json', 'pretty'], + SUPPORTED_LOG_LEVELS: ['error', 'warn', 'info', 'http', 'verbose', 'debug', 'silly'], + default: { + debug: jest.fn(), + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + }, +})); + +const DAY_MS = 24 * 60 * 60 * 1000; +const OLD_DATE = new Date(Date.now() - 40 * DAY_MS).toISOString(); +const RECENT_DATE = new Date(Date.now() - DAY_MS).toISOString(); + +function cleanupConfig( + overrides: Partial<{ + enabled: boolean; + intervalMs: number; + retentionDays: number; + retentionOverridesMs: { + processedEvents?: number; + executionLogs?: number; + rateLimitEvents?: number; + }; + }> = {}, +) { + return { + enabled: true, + intervalMs: 60_000, + retentionDays: 30, + notificationRetentionMs: DAY_MS * 7, + rateLimitEventRetentionMs: DAY_MS, + eventRetentionMs: DAY_MS, + processedEventRetentionMs: DAY_MS * 30, + executionLogRetentionMs: DAY_MS * 90, + ...overrides, + }; +} + +describe('DatabaseCleanupJob', () => { + let db: Database; + const originalEnv = { ...process.env }; + + beforeEach(async () => { + jest.clearAllMocks(); + resetWorkerManager(); + db = new Database(':memory:'); + await db.initialize(); + }); + + afterEach(async () => { + resetWorkerManager(); + await db.close(); + process.env = originalEnv; + jest.useRealTimers(); + }); + + async function addNotification(status: string): Promise { + const result = await db.run( + `INSERT INTO scheduled_notifications ( + payload, notification_type, target_recipient, execute_at, status + ) VALUES ('{}', 'discord', 'test-recipient', CURRENT_TIMESTAMP, ?)`, + [status], + ); + return result.lastID; + } + + async function addProcessedEvent(eventId: string, processedAt: string): Promise { + await db.run( + `INSERT INTO processed_events ( + event_id, contract_address, fingerprint, ledger_number, event_type, processed_at + ) VALUES (?, 'CABC', ?, 1, 'contract', ?)`, + [eventId, `fingerprint-${eventId}`, processedAt], + ); + } + + async function addIdempotencyKey( + id: string, + notificationId: number, + expiresAt: string, + createdAt: string, + status = 'PROCESSED', + ): Promise { + await db.run( + `INSERT INTO idempotency_keys ( + idempotency_key, request_hash, response_notification_id, response_data, + created_at, expires_at, status + ) VALUES (?, 'hash', ?, '{}', ?, ?, ?)`, + [id, notificationId, createdAt, expiresAt, status], + ); + } + + it('purges stale temporary records but preserves active rows and logs table counts', async () => { + const pendingId = await addNotification('PENDING'); + const processingId = await addNotification('PROCESSING'); + const terminalId = await addNotification('FAILED'); + const secondTerminalId = await addNotification('FAILED'); + + await addIdempotencyKey('expired-key', terminalId, OLD_DATE, OLD_DATE); + await addIdempotencyKey( + 'future-key', + terminalId, + new Date(Date.now() + DAY_MS).toISOString(), + OLD_DATE, + ); + await addIdempotencyKey( + 'expired-status-key', + terminalId, + new Date(Date.now() + DAY_MS).toISOString(), + OLD_DATE, + 'EXPIRED', + ); + + await addProcessedEvent('old-event', OLD_DATE); + await addProcessedEvent('recent-event', RECENT_DATE); + + await db.run( + `INSERT INTO dead_letter_queue ( + scheduled_notification_id, notification_type, target_recipient, payload, failure_reason, created_at + ) VALUES (?, 'discord', 'recipient', '{}', 'failed', ?)`, + [terminalId, OLD_DATE], + ); + await db.run( + `INSERT INTO dead_letter_queue ( + scheduled_notification_id, notification_type, target_recipient, payload, failure_reason, created_at + ) VALUES (?, 'discord', 'recipient', '{}', 'failed', ?)`, + [secondTerminalId, RECENT_DATE], + ); + + for (const [notificationId, executionTime] of [ + [pendingId, OLD_DATE], + [processingId, OLD_DATE], + [terminalId, OLD_DATE], + [secondTerminalId, RECENT_DATE], + ] as Array<[number, string]>) { + await db.run( + `INSERT INTO notification_execution_log ( + scheduled_notification_id, execution_attempt, execution_time, status + ) VALUES (?, 1, ?, 'FAILED')`, + [notificationId, executionTime], + ); + } + + await db.run( + `INSERT INTO rate_limit_events ( + client_id, client_type, endpoint, method, timestamp, limit_threshold, window_ms + ) VALUES ('old', 'IP', '/', 'GET', ?, 10, 60000)`, + [OLD_DATE], + ); + await db.run( + `INSERT INTO rate_limit_events ( + client_id, client_type, endpoint, method, timestamp, limit_threshold, window_ms + ) VALUES ('recent', 'IP', '/', 'GET', ?, 10, 60000)`, + [RECENT_DATE], + ); + + await db.run( + `INSERT INTO backpressure_events (event_type, queue_size, target_throughput_per_sec, timestamp) + VALUES ('ACTIVATED', 100, 10, ?)`, + [OLD_DATE], + ); + await db.run( + `INSERT INTO backpressure_events (event_type, queue_size, target_throughput_per_sec, timestamp) + VALUES ('DEACTIVATED', 0, 10, ?)`, + [OLD_DATE], + ); + await db.run( + `INSERT INTO backpressure_events (event_type, queue_size, target_throughput_per_sec, timestamp) + VALUES ('DEACTIVATED', 0, 10, ?)`, + [RECENT_DATE], + ); + + const job = new DatabaseCleanupJob(db, cleanupConfig()); + const result = await job.runOnce(); + + expect( + await db.get('SELECT id FROM idempotency_keys WHERE idempotency_key = ?', ['expired-key']), + ).toBeUndefined(); + expect( + await db.get('SELECT id FROM idempotency_keys WHERE idempotency_key = ?', ['future-key']), + ).toBeDefined(); + expect( + await db.get('SELECT id FROM idempotency_keys WHERE idempotency_key = ?', [ + 'expired-status-key', + ]), + ).toBeUndefined(); + expect( + await db.get('SELECT id FROM processed_events WHERE event_id = ?', ['old-event']), + ).toBeUndefined(); + expect( + await db.get('SELECT id FROM processed_events WHERE event_id = ?', ['recent-event']), + ).toBeDefined(); + expect( + await db.get('SELECT id FROM dead_letter_queue WHERE scheduled_notification_id = ?', [ + terminalId, + ]), + ).toBeUndefined(); + expect( + await db.get('SELECT id FROM dead_letter_queue WHERE scheduled_notification_id = ?', [ + secondTerminalId, + ]), + ).toBeDefined(); + + const logRows = await db.all<{ scheduled_notification_id: number }>( + 'SELECT scheduled_notification_id FROM notification_execution_log ORDER BY scheduled_notification_id', + ); + expect(logRows.map((row) => row.scheduled_notification_id)).toEqual([ + pendingId, + processingId, + secondTerminalId, + ]); + expect( + await db.get('SELECT id FROM rate_limit_events WHERE client_id = ?', ['old']), + ).toBeUndefined(); + expect( + await db.get('SELECT id FROM rate_limit_events WHERE client_id = ?', ['recent']), + ).toBeDefined(); + + const backpressureRows = await db.all<{ event_type: string }>( + 'SELECT event_type FROM backpressure_events ORDER BY id', + ); + expect(backpressureRows.map((row) => row.event_type)).toEqual(['ACTIVATED', 'DEACTIVATED']); + expect( + await db.get('SELECT status FROM scheduled_notifications WHERE id = ?', [pendingId]), + ).toBeDefined(); + expect( + await db.get('SELECT status FROM scheduled_notifications WHERE id = ?', [processingId]), + ).toBeDefined(); + + expect(result?.deletedCounts).toMatchObject({ + idempotency_keys: 2, + processed_events: 1, + dead_letter_queue: 1, + notification_execution_log: 1, + rate_limit_events: 1, + backpressure_events: 1, + }); + expect(result?.skippedTables.notification_archive).toContain('ArchiveService'); + expect(logger.info).toHaveBeenCalledWith( + 'Database cleanup run completed', + expect.objectContaining({ + perTableDeleted: result?.deletedCounts, + retentionDays: 30, + intervalMs: 60_000, + }), + ); + }); + + it('rejects cleanup configuration below the supported interval and retention', () => { + process.env.CONTRACT_ADDRESSES = JSON.stringify([ + { address: 'CXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXX', events: ['*'] }, + ]); + process.env.CLEANUP_INTERVAL_MS = '59999'; + process.env.CLEANUP_RETENTION_DAYS = '0'; + + const config = loadConfig(); + expect(() => validateConfig(config)).toThrow(ConfigError); + expect(() => validateConfig(config)).toThrow('CLEANUP_INTERVAL_MS must be >= 60000 ms'); + expect(() => validateConfig(config)).toThrow('CLEANUP_RETENTION_DAYS must be >= 1'); + }); + + it('limits each delete transaction to 1000 rows', async () => { + for (let index = 0; index < 1001; index += 1) { + await addProcessedEvent(`batch-event-${index}`, OLD_DATE); + } + const runSpy = jest.spyOn(db, 'run'); + + const result = await new DatabaseCleanupJob(db, cleanupConfig()).runOnce(); + + const processedEventDeletes = runSpy.mock.calls.filter(([sql]) => + sql.includes('DELETE FROM processed_events'), + ); + expect(result?.deletedCounts.processed_events).toBe(1001); + expect(processedEventDeletes).toHaveLength(2); + expect(processedEventDeletes.every(([, params]) => params?.[params.length - 1] === 1000)).toBe( + true, + ); + }); + + it('continues with other tables when one table purge fails', async () => { + await addProcessedEvent('event-survives-other-table-failure', OLD_DATE); + const terminalId = await addNotification('FAILED'); + await db.run( + `INSERT INTO dead_letter_queue ( + scheduled_notification_id, notification_type, target_recipient, payload, failure_reason, created_at + ) VALUES (?, 'discord', 'recipient', '{}', 'failed', ?)`, + [terminalId, OLD_DATE], + ); + const originalRun = db.run.bind(db); + jest.spyOn(db, 'run').mockImplementation((sql, params = []) => { + if (sql.includes('DELETE FROM dead_letter_queue')) { + return Promise.reject(new Error('injected table failure')); + } + return originalRun(sql, params); + }); + + const result = await new DatabaseCleanupJob(db, cleanupConfig()).runOnce(); + + expect(result?.failedTables).toContain('dead_letter_queue'); + expect(result?.deletedCounts.processed_events).toBe(1); + expect( + await db.get('SELECT id FROM processed_events WHERE event_id = ?', [ + 'event-survives-other-table-failure', + ]), + ).toBeUndefined(); + expect( + await db.get('SELECT id FROM dead_letter_queue WHERE scheduled_notification_id = ?', [ + terminalId, + ]), + ).toBeDefined(); + expect(logger.error).toHaveBeenCalledWith( + 'Database cleanup table failed', + expect.objectContaining({ table: 'dead_letter_queue' }), + ); + }); + + it('rejects invalid cleanup enabled values', () => { + process.env.CLEANUP_ENABLED = 'yes'; + expect(() => loadConfig()).toThrow('CLEANUP_ENABLED must be either "true" or "false"'); + }); + + it('defaults cleanup to enabled and honors explicit disablement', () => { + process.env.CONTRACT_ADDRESSES = JSON.stringify([ + { address: 'CXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXXX', events: ['*'] }, + ]); + delete process.env.CLEANUP_ENABLED; + delete process.env.CLEANUP_INTERVAL_MS; + delete process.env.CLEANUP_RETENTION_DAYS; + + expect(loadConfig().cleanup).toMatchObject({ + enabled: true, + intervalMs: 3_600_000, + retentionDays: 30, + }); + + process.env.CLEANUP_ENABLED = 'false'; + expect(loadConfig().cleanup?.enabled).toBe(false); + }); + + it('stops the interval and waits for its active cleanup run', async () => { + jest.useFakeTimers(); + const intervalSpy = jest.spyOn(global, 'setInterval'); + const clearIntervalSpy = jest.spyOn(global, 'clearInterval'); + const job = new DatabaseCleanupJob(db, cleanupConfig()); + job.start(); + const jobTimer = intervalSpy.mock.results[0].value; + await job.stop(); + + expect(clearIntervalSpy).toHaveBeenCalledWith(jobTimer); + intervalSpy.mockRestore(); + clearIntervalSpy.mockRestore(); + }); +}); diff --git a/listener/src/services/database-cleanup-job.ts b/listener/src/services/database-cleanup-job.ts new file mode 100644 index 00000000..fb12e37c --- /dev/null +++ b/listener/src/services/database-cleanup-job.ts @@ -0,0 +1,253 @@ +import { Database } from '../database/database'; +import { AppCleanupConfig } from '../types'; +import { EventRegistry } from '../store/event-registry'; +import logger from '../utils/logger'; +import { getWorkerManager } from './worker-manager'; + +const DELETE_BATCH_SIZE = 1000; +const DAY_MS = 24 * 60 * 60 * 1000; + +type RetentionOverride = keyof NonNullable; + +interface CleanupTarget { + table: string; + requiredColumns: string[]; + where: string; + params: (now: string, cutoff: string) => unknown[]; + retentionOverride?: RetentionOverride; +} + +export interface DatabaseCleanupResult { + runId: string; + deletedCounts: Record; + skippedTables: Record; + failedTables: string[]; + durationMs: number; +} + +const CLEANUP_TARGETS: CleanupTarget[] = [ + { + table: 'idempotency_keys', + requiredColumns: ['expires_at', 'created_at', 'status'], + where: `( + (datetime(expires_at) IS NOT NULL AND datetime(expires_at) < datetime(?)) + OR (status = 'EXPIRED' AND datetime(created_at) < datetime(?)) + )`, + params: (now, cutoff) => [now, cutoff], + }, + { + table: 'processed_events', + requiredColumns: ['processed_at'], + where: 'datetime(processed_at) < datetime(?)', + params: (_now, cutoff) => [cutoff], + retentionOverride: 'processedEvents', + }, + { + table: 'dead_letter_queue', + requiredColumns: ['created_at'], + where: 'datetime(created_at) < datetime(?)', + params: (_now, cutoff) => [cutoff], + }, + { + table: 'notification_execution_log', + requiredColumns: ['execution_time', 'scheduled_notification_id'], + where: `datetime(execution_time) < datetime(?) + AND NOT EXISTS ( + SELECT 1 FROM scheduled_notifications n + WHERE n.id = notification_execution_log.scheduled_notification_id + AND n.status IN ('PENDING', 'PROCESSING') + )`, + params: (_now, cutoff) => [cutoff], + retentionOverride: 'executionLogs', + }, + { + table: 'rate_limit_events', + requiredColumns: ['timestamp', 'window_ms'], + where: `datetime(timestamp) < datetime(?) + AND julianday(timestamp) + (window_ms / 86400000.0) < julianday(?)`, + params: (now, cutoff) => [cutoff, now], + retentionOverride: 'rateLimitEvents', + }, + { + table: 'backpressure_events', + requiredColumns: ['event_type', 'timestamp'], + where: `event_type = 'DEACTIVATED' AND datetime(timestamp) < datetime(?)`, + params: (_now, cutoff) => [cutoff], + }, +]; + +/** Periodically removes expired temporary and derived records in small transactions. */ +export class DatabaseCleanupJob { + private timer: ReturnType | null = null; + private activeRun: Promise | null = null; + private runSequence = 0; + + constructor( + private readonly db: Database, + private readonly config: AppCleanupConfig, + private readonly registry?: EventRegistry, + ) {} + + start(): void { + if (this.timer) return; + if (!this.config.enabled) { + logger.info('Database cleanup job is disabled', { + intervalMs: this.config.intervalMs, + retentionDays: this.config.retentionDays, + }); + return; + } + + this.registry?.startCleanup(this.config.intervalMs); + logger.info('Database cleanup job started', { + intervalMs: this.config.intervalMs, + retentionDays: this.config.retentionDays, + }); + void this.runOnce(); + this.timer = setInterval(() => void this.runOnce(), this.config.intervalMs); + this.timer.unref?.(); + } + + async stop(): Promise { + if (this.timer) { + clearInterval(this.timer); + this.timer = null; + } + this.registry?.stopCleanup(); + if (this.activeRun) await this.activeRun; + logger.info('Database cleanup job stopped'); + } + + async runOnce(): Promise { + if (this.activeRun) return this.activeRun; + const run = this.executeRun(); + this.activeRun = run; + try { + return await run; + } finally { + if (this.activeRun === run) this.activeRun = null; + } + } + + private async executeRun(): Promise { + const runId = `database-cleanup-${Date.now()}-${++this.runSequence}`; + const workerManager = getWorkerManager(); + if (!workerManager.startJob(runId)) { + logger.warn('Database cleanup skipped during shutdown', { runId }); + return null; + } + + const startedAt = Date.now(); + const now = new Date(startedAt).toISOString(); + const globalRetentionMs = this.config.retentionDays * DAY_MS; + const deletedCounts: Record = {}; + const skippedTables: Record = { + scheduled_notifications: + 'Rows are retained or moved by ArchiveService; stale processing locks are recovered by the repository.', + notification_archive: 'ArchiveService owns archived_at-based retention and purging.', + notification_metrics_snapshots: 'NotificationMetricsRunner owns snapshot retention.', + polling_cursors: 'Cursor and reorg state is required to resume event polling safely.', + notification_templates: 'Template definitions remain active application configuration.', + notification_template_audit_log: + 'Template audit history is immutable and retained for compliance.', + }; + const failedTables: string[] = []; + + try { + logger.warn('Database cleanup delegated tables are skipped', { runId, skippedTables }); + for (const target of CLEANUP_TARGETS) { + deletedCounts[target.table] = 0; + try { + const columns = await this.db.all<{ name: string }>(`PRAGMA table_info(${target.table})`); + const availableColumns = new Set(columns.map((column) => column.name)); + const missingColumns = target.requiredColumns.filter( + (column) => !availableColumns.has(column), + ); + if (missingColumns.length > 0) { + skippedTables[target.table] = + `Missing required timestamp/schema column(s): ${missingColumns.join(', ')}`; + logger.warn('Database cleanup table skipped', { + runId, + table: target.table, + missingColumns, + }); + continue; + } + + const retentionMs = target.retentionOverride + ? (this.config.retentionOverridesMs?.[target.retentionOverride] ?? globalRetentionMs) + : globalRetentionMs; + const cutoff = new Date(startedAt - retentionMs).toISOString(); + const deletion = await this.deleteInBatches(target, now, cutoff); + deletedCounts[target.table] = deletion.deleted; + if (deletion.error) { + failedTables.push(target.table); + logger.error('Database cleanup table failed', { + runId, + table: target.table, + error: deletion.error, + }); + } + } catch (error) { + failedTables.push(target.table); + logger.error('Database cleanup table failed', { + runId, + table: target.table, + error, + }); + } + } + } catch (error) { + logger.error('Database cleanup run failed', { runId, error }); + } finally { + workerManager.completeJob(runId); + } + + const result: DatabaseCleanupResult = { + runId, + deletedCounts, + skippedTables, + failedTables, + durationMs: Date.now() - startedAt, + }; + logger.info('Database cleanup run completed', { + runId, + perTableDeleted: deletedCounts, + skippedTables, + failedTables, + retentionDays: this.config.retentionDays, + intervalMs: this.config.intervalMs, + durationMs: result.durationMs, + }); + return result; + } + + private async deleteInBatches( + target: CleanupTarget, + now: string, + cutoff: string, + ): Promise<{ deleted: number; error?: unknown }> { + let totalDeleted = 0; + while (true) { + let deleted = 0; + try { + await this.db.transaction(async () => { + const result = await this.db.run( + `DELETE FROM ${target.table} + WHERE rowid IN ( + SELECT rowid FROM ${target.table} + WHERE ${target.where} + LIMIT ? + )`, + [...target.params(now, cutoff), DELETE_BATCH_SIZE], + ); + deleted = result.changes; + }); + } catch (error) { + return { deleted: totalDeleted, error }; + } + totalDeleted += deleted; + if (deleted < DELETE_BATCH_SIZE) return { deleted: totalDeleted }; + } + } +} diff --git a/listener/src/types/index.ts b/listener/src/types/index.ts index 3f5c3492..2f80a22c 100644 --- a/listener/src/types/index.ts +++ b/listener/src/types/index.ts @@ -156,8 +156,18 @@ export interface EventQueueConfig { } export interface AppCleanupConfig { + /** Whether scheduled database cleanup is enabled. */ + enabled: boolean; /** How often to run cleanup jobs (ms). */ intervalMs: number; + /** Global retention period for database cleanup (days). */ + retentionDays: number; + /** Explicit legacy per-table overrides, when supplied. */ + retentionOverridesMs?: { + processedEvents?: number; + executionLogs?: number; + rateLimitEvents?: number; + }; /** Retain completed/failed/cancelled notifications for this long (ms). */ notificationRetentionMs: number; /** Retain rate-limit audit rows for this long (ms). */