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
22 changes: 21 additions & 1 deletion NOTIFICATION_LIFECYCLE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).

Expand Down Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions dashboard/src/pages/NotificationSearchPage.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -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' },
];

Expand Down
4 changes: 4 additions & 0 deletions listener/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions listener/src/api/events-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
57 changes: 56 additions & 1 deletion listener/src/config.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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)) {
Expand Down Expand Up @@ -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();

Expand Down Expand Up @@ -439,6 +493,7 @@ export function loadConfig(): Config {
cleanup: loadCleanupConfig(),
analytics: loadAnalyticsConfig(),
expiration: loadExpirationConfig(),
notificationDefaultTtlSeconds: loadNotificationDefaultTtlSeconds(),
backfill: loadBackfillConfig(),
rpcRateLimit: loadRpcRateLimitConfig(),
logging: loadLoggingConfig(),
Expand Down
5 changes: 3 additions & 2 deletions listener/src/database/archive-schema.sql
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion listener/src/database/database.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -178,7 +179,6 @@ export class Database {
if (pragmaErr) {
reject(pragmaErr);
} else {
this.db = handle;
resolve();
}
});
Expand Down
77 changes: 58 additions & 19 deletions listener/src/database/migration-system.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,8 +37,49 @@ export class MigrationRunner {
this.migrationsDir = migrationsDir;
}

private run(sql: string, params: unknown[] = []): Promise<void> {
return new Promise((resolve, reject) => {
this.db.run(sql, params, (error) => (error ? reject(error) : resolve()));
});
}

private all<T>(sql: string): Promise<T[]> {
return new Promise((resolve, reject) => {
this.db.all(sql, (error, rows) => (error ? reject(error) : resolve(rows as T[])));
});
}

private async ensureForeignKeysEnabled(): Promise<void> {
const current = await new Promise<number>((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<void>((resolve, reject) => {
this.db.run('PRAGMA foreign_keys = ON', (error) => {
if (error) reject(error);
else resolve();
});
});
}

const verified = await new Promise<number>((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<void> {
await this.db.run(`
await this.ensureForeignKeysEnabled();
await this.run(`
CREATE TABLE IF NOT EXISTS migrations (
id TEXT PRIMARY KEY,
name TEXT NOT NULL,
Expand All @@ -48,29 +89,27 @@ export class MigrationRunner {
}

async getAppliedMigrations(): Promise<string[]> {
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<void> {
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<Migration[]> {
Expand Down
Loading
Loading