Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 41 additions & 1 deletion NOTIFICATION_LIFECYCLE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)

---

Expand Down Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions listener/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
7 changes: 7 additions & 0 deletions listener/src/config-schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 },
Expand Down
60 changes: 50 additions & 10 deletions listener/src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<T>(name: string, defaultValue: string): T {
const rawValue = trimEnv(name) ?? defaultValue;
try {
Expand Down Expand Up @@ -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,
};
}

Expand Down Expand Up @@ -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 ───────────────────────────────────────────────────────────────
Expand Down
14 changes: 8 additions & 6 deletions listener/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -185,8 +187,8 @@ async function main() {
healthMonitor.stop();
}

if (cleanupService) {
await cleanupService.stop();
if (databaseCleanupJob) {
await databaseCleanupJob.stop();
}

if (reconciliationEngine) {
Expand Down
23 changes: 14 additions & 9 deletions listener/src/services/archive.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand Down
Loading
Loading