diff --git a/listener/.env.example b/listener/.env.example index 3a156366..39167b05 100644 --- a/listener/.env.example +++ b/listener/.env.example @@ -216,6 +216,10 @@ RETRY_SCHEDULER_PROCESSOR_ID= # Number of retry jobs to process per scheduler tick. RETRY_SCHEDULER_BATCH_SIZE=10 +# Request timeout for outbound webhook delivery (ms). Timed-out requests are +# treated as failures and are subject to normal retry backoff. +# WEBHOOK_DELIVERY_TIMEOUT_MS=10000 + # ----------------------------------------------------------------------------- # Retry Policy # ----------------------------------------------------------------------------- diff --git a/listener/src/api/events-server.ts b/listener/src/api/events-server.ts index d7277a86..d5ed4f66 100644 --- a/listener/src/api/events-server.ts +++ b/listener/src/api/events-server.ts @@ -14,8 +14,10 @@ import { handleTemplateRoutes } from './template-routes'; import { sendOk, sendErr, sendJson, ErrorCode } from '../utils/response'; import { normalizePaginationParams } from '../utils/pagination'; import { handleApiError, ApiError } from './error-handler'; -import { applyRequestContext } from '../utils/request-id'; import { applyRequestIdMiddleware } from '../middleware/request-id'; +import { TemplateService } from '../services/template-service'; +import { handleTemplateRoutes } from './template-routes'; +import { sendOk, sendErr, sendJson, ErrorCode } from '../utils/response'; import { validateContentType, getMimeType } from '../middleware/content-type'; import { NotificationHistoryService } from '../services/notification-history'; import { SearchSuggestionService } from '../services/search-suggestion'; diff --git a/listener/src/config.ts b/listener/src/config.ts index 2720252c..d755ba4f 100644 --- a/listener/src/config.ts +++ b/listener/src/config.ts @@ -262,6 +262,7 @@ function loadAnalyticsConfig(): AnalyticsConfig { function loadRetrySchedulerConfig(policy: RetryPolicyOptions): RetrySchedulerOptions { return { enabled: trimEnv('RETRY_SCHEDULER_ENABLED') !== 'false', + webhookTimeoutMs: parseIntegerEnv('WEBHOOK_DELIVERY_TIMEOUT_MS', '10000'), pollIntervalMs: parseIntegerEnv('RETRY_SCHEDULER_POLL_INTERVAL_MS', '15000'), lockTimeoutMs: parseIntegerEnv('RETRY_SCHEDULER_LOCK_TIMEOUT_MS', '60000'), processorId: trimEnv('RETRY_SCHEDULER_PROCESSOR_ID'), diff --git a/listener/src/database/migration-system.ts b/listener/src/database/migration-system.ts index 8cc12e06..fdf4375b 100644 --- a/listener/src/database/migration-system.ts +++ b/listener/src/database/migration-system.ts @@ -51,7 +51,7 @@ export class MigrationRunner { const rows = await this.db.all<{ id: string }>( 'SELECT id FROM migrations ORDER BY applied_at' ); - return rows.map((row) => row.id); + return (rows as unknown as { id: string }[]).map((row) => row.id); } async applyMigration(migration: Migration): Promise { @@ -67,7 +67,7 @@ export class MigrationRunner { 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); + logger.error(`Migration ${migration.id} failed, rolling back: ${(error as Error)?.message ?? String(error)}`); throw error; } }); diff --git a/listener/src/database/schema.sql b/listener/src/database/schema.sql index 9b67ac22..d5e65843 100644 --- a/listener/src/database/schema.sql +++ b/listener/src/database/schema.sql @@ -36,7 +36,8 @@ CREATE TABLE IF NOT EXISTS scheduled_notifications ( contract_address TEXT, -- Stellar contract address (if applicable) priority INTEGER NOT NULL DEFAULT 5, -- 1-10, lower = higher priority metadata TEXT, -- Additional JSON metadata - next_retry_at DATETIME -- When the next retry should be attempted + 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 ); -- Indexes for performance optimization @@ -62,6 +63,10 @@ CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_created_at CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_event_id ON scheduled_notifications(event_id); +CREATE UNIQUE INDEX IF NOT EXISTS idx_scheduled_notifications_dedup_key + ON scheduled_notifications(deduplication_key) + WHERE deduplication_key IS NOT NULL; + CREATE INDEX IF NOT EXISTS idx_scheduled_notifications_target ON scheduled_notifications(target_recipient, status); diff --git a/listener/src/middleware/security-headers.ts b/listener/src/middleware/security-headers.ts index 5de7d64e..d1880c17 100644 --- a/listener/src/middleware/security-headers.ts +++ b/listener/src/middleware/security-headers.ts @@ -10,6 +10,7 @@ * See: https://owasp.org/www-project-secure-headers/ */ +import type { ServerResponse } from 'http'; import type http from 'http'; import type { ServerResponse } from 'http'; diff --git a/listener/src/migrations/003-notification-deduplication-key.ts b/listener/src/migrations/003-notification-deduplication-key.ts new file mode 100644 index 00000000..8d107df3 --- /dev/null +++ b/listener/src/migrations/003-notification-deduplication-key.ts @@ -0,0 +1,34 @@ +import * as sqlite3 from 'sqlite3'; + +const migration = { + id: '003', + name: 'notification-deduplication-key', + up: async (db: sqlite3.Database) => { + await db.run(` + ALTER TABLE scheduled_notifications + ADD COLUMN deduplication_key TEXT + `); + await db.run(` + CREATE UNIQUE INDEX IF NOT EXISTS idx_scheduled_notifications_dedup_key + ON scheduled_notifications(deduplication_key) + WHERE deduplication_key IS NOT NULL + `); + }, + down: async (db: sqlite3.Database) => { + await db.run('DROP INDEX IF EXISTS idx_scheduled_notifications_dedup_key'); + // SQLite does not support DROP COLUMN before 3.35; recreate without the column. + await db.run(` + CREATE TABLE scheduled_notifications_backup AS + SELECT id, payload, payload_hash, notification_type, target_recipient, + execute_at, created_at, updated_at, status, retry_count, max_retries, + processing_started_at, processing_completed_at, processor_id, + lock_expires_at, last_error, error_details, event_id, + contract_address, priority, metadata, next_retry_at + FROM scheduled_notifications + `); + await db.run('DROP TABLE scheduled_notifications'); + await db.run('ALTER TABLE scheduled_notifications_backup RENAME TO scheduled_notifications'); + }, +}; + +export default migration; diff --git a/listener/src/schema/compatibility-check.ts b/listener/src/schema/compatibility-check.ts index 3f9df1eb..47c89f2d 100644 --- a/listener/src/schema/compatibility-check.ts +++ b/listener/src/schema/compatibility-check.ts @@ -172,6 +172,14 @@ export function checkSchemaCompatibility( breakingChanges.push('Consumer event schema is empty or malformed'); } + // Check 4: Topic structure changes + // Topics are appended as trailing topics - existing consumers ignore them + // Breaking change only if the core topic (event name) changes position + if (consumer.expectedTopics && contractEvent.topics.length < consumer.expectedTopics.length) { + const missingTopics = consumer.expectedTopics.filter( + (t) => !contractEvent.topics.includes(t) + ); + if (missingTopics.length > 0) { for (const consumerEvent of consumerSchema.events) { const contractEvent = contractEvents.get(consumerEvent.eventName); if (!contractEvent) { @@ -185,6 +193,20 @@ export function checkSchemaCompatibility( safeAdditions.push(...result.safeAdditions); } + // Check 5: New data fields are safe (backward compatible) + // Any new data fields in the contract that weren't expected by the consumer + // are simply ignored - this is the Soroban trailing-topic pattern + const newDataFields = contractEvent.dataFields.filter( + (f) => !consumer.expectedFields.includes(f.name) + ); + safeAdditions.push( + ...newDataFields.map( + (f) => `Safe addition: New data field '${f.name}' in '${contractEvent.name}' (ignored by existing consumers)` + ) + ); + + // Determine compatibility + const compatible = breakingChanges.length === 0; for (const eventName of contractEvents.keys()) { if (!consumerNames.has(eventName)) { safeAdditions.push(`Safe addition: New event '${eventName}' is ignored by existing consumers`); diff --git a/listener/src/services/discord-notification.ts b/listener/src/services/discord-notification.ts index 02ec64d3..7ce1e316 100644 --- a/listener/src/services/discord-notification.ts +++ b/listener/src/services/discord-notification.ts @@ -407,6 +407,7 @@ export class DiscordNotificationService implements NotificationProvider { const originalLength = embed.title?.length ?? title.length; title = title.slice(0, MAX_DISCORD_EMBED_TITLE_LENGTH - 3) + '...'; logger.warn('Discord embed title truncated', { + originalLength: title.length + 3, originalLength, maxLength: MAX_DISCORD_EMBED_TITLE_LENGTH, }); @@ -427,8 +428,10 @@ export class DiscordNotificationService implements NotificationProvider { let footer = embed.footer; if (footer?.text && footer.text.length > MAX_DISCORD_FOOTER_TEXT_LENGTH) { + const originalLength = footer.text.length; footer = { text: footer.text.slice(0, MAX_DISCORD_FOOTER_TEXT_LENGTH - 3) + '...' }; logger.warn('Discord footer text truncated', { + originalLength, originalLength: footer.text.length, maxLength: MAX_DISCORD_FOOTER_TEXT_LENGTH, }); diff --git a/listener/src/services/event-subscriber.ts b/listener/src/services/event-subscriber.ts index 3deae703..38f91278 100644 --- a/listener/src/services/event-subscriber.ts +++ b/listener/src/services/event-subscriber.ts @@ -31,6 +31,7 @@ export class EventSubscriber { private eventQueue: EventProcessingQueue | null = null; private expirationService: NotificationExpirationService | null = null; private lastSuccessfulPollAt: number | null = null; + private backfillStartLedger: number | null = null; private circuitBreaker: CircuitBreaker | null = null; private backfillStartLedger: number | null = null; /** Cold-start ledger resolved once per session by resolveBackfillStartLedger(). */ @@ -389,6 +390,9 @@ export class EventSubscriber { if (lastCursor) { // Normal real-time polling: continue from the last known cursor. request = { + filters: [{ contractIds: [contractConfig.address], type: 'contract' }], + cursor: lastCursor, + limit: 100, filters: [ { contractIds: [contractConfig.address], @@ -405,6 +409,9 @@ export class EventSubscriber { // Cold start: apply the backfill safety limit. const startLedger = await this.resolveBackfillStartLedger(); request = { + filters: [{ contractIds: [contractConfig.address], type: 'contract' }], + startLedger, + limit: 100, filters: [ { contractIds: [contractConfig.address], diff --git a/listener/src/services/notification-retry-queue.ts b/listener/src/services/notification-retry-queue.ts index f6432983..8b39a285 100644 --- a/listener/src/services/notification-retry-queue.ts +++ b/listener/src/services/notification-retry-queue.ts @@ -367,6 +367,6 @@ function buildRetryFingerprint( contractAddress: string ): string { const eventName = - getEventName(event.topic) ?? event.topic.map((entry) => entry.toString()).join('|'); + getEventName(event.topic) ?? event.topic.map((entry: { toString(): string }) => entry.toString()).join('|'); return `${contractAddress}:${event.id}:${eventName}:${event.txHash ?? ''}`; } diff --git a/listener/src/services/notification-stats-cache.ts b/listener/src/services/notification-stats-cache.ts index 74df1e86..9625eb47 100644 --- a/listener/src/services/notification-stats-cache.ts +++ b/listener/src/services/notification-stats-cache.ts @@ -161,6 +161,7 @@ export class NotificationStatsCache { * @notice Check if stats are currently cached */ has(): boolean { + return this.cache.get(this.CACHE_KEY) !== undefined; return (this.cache as any).has(this.CACHE_KEY); } diff --git a/listener/src/services/retry-scheduler.ts b/listener/src/services/retry-scheduler.ts index 517814cb..0c0ae8ad 100644 --- a/listener/src/services/retry-scheduler.ts +++ b/listener/src/services/retry-scheduler.ts @@ -36,6 +36,8 @@ export interface RetrySchedulerConfig { maxDelayMs: number; /** Add ±25 % random jitter to prevent thundering herd. Default: true. */ jitter: boolean; + /** Request timeout for outbound webhook delivery (ms). Default: 10 000. */ + webhookTimeoutMs: number; /** * Hard ceiling on total delivery attempts, including the first one. * `undefined` (default) leaves each notification's own `maxRetries` in @@ -59,6 +61,7 @@ export const RETRY_SCHEDULER_DEFAULTS: RetrySchedulerConfig = { multiplier: 2, maxDelayMs: 60 * 60 * 1_000, jitter: true, + webhookTimeoutMs: 10_000, }; /** @@ -125,6 +128,8 @@ export class RetryScheduler { this.processorId = this.config.processorId ?? `retry-${uuidv4()}`; this.repository = repository; this.discordService = discordService ?? null; + this.webhookDeliveryService = + webhookDeliveryService ?? new WebhookDeliveryService({ timeoutMs: this.config.webhookTimeoutMs }); this.webhookDeliveryService = webhookDeliveryService ?? new WebhookDeliveryService(); this.deliveryReceiptRepository = deliveryReceiptRepository; } diff --git a/listener/src/services/scheduled-notification-repository.ts b/listener/src/services/scheduled-notification-repository.ts index 05feef4e..d19a26b2 100644 --- a/listener/src/services/scheduled-notification-repository.ts +++ b/listener/src/services/scheduled-notification-repository.ts @@ -32,7 +32,9 @@ export class ScheduledNotificationRepository { } /** - * Create a new scheduled notification + * Create a new scheduled notification. + * If a deduplicationKey is provided and a notification with that key already + * exists, the existing notification's id is returned without creating a duplicate. */ async create(input: CreateScheduledNotificationInput, requestId?: string): Promise { const payloadJson = JSON.stringify(input.payload); @@ -42,8 +44,8 @@ export class ScheduledNotificationRepository { const sql = ` INSERT INTO scheduled_notifications ( payload, payload_hash, notification_type, target_recipient, execute_at, - max_retries, event_id, contract_address, priority, metadata - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + max_retries, event_id, contract_address, priority, metadata, deduplication_key + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) `; const serializedPayload = compressPayload(input.payload); @@ -59,8 +61,20 @@ export class ScheduledNotificationRepository { input.contractAddress ?? null, input.priority ?? 5, input.metadata ? JSON.stringify(input.metadata) : null, + input.deduplicationKey ?? null, ]; + try { + const result = await this.db.run(sql, params); + + this.statsCache.invalidate(); + + logger.info('Scheduled notification created', { + requestId, + id: result.lastID, + executeAt: input.executeAt, + type: input.notificationType, + }); const result = await this.db.run(sql, params); // Invalidate stats cache after creation @@ -73,7 +87,29 @@ export class ScheduledNotificationRepository { type: input.notificationType, }); - return result.lastID; + return result.lastID; + } catch (err) { + if ( + input.deduplicationKey && + (err as any)?.message?.includes('UNIQUE constraint failed') + ) { + const existing = await this.db.get<{ id: number }>( + 'SELECT id FROM scheduled_notifications WHERE deduplication_key = ?', + [input.deduplicationKey], + ); + + if (existing) { + logger.info('Duplicate notification skipped — deduplication key already exists', { + requestId, + deduplicationKey: input.deduplicationKey, + existingId: existing.id, + }); + return existing.id; + } + } + + throw err; + } } /** @@ -902,6 +938,7 @@ export class ScheduledNotificationRepository { priority: row.priority, metadata: row.metadata, nextRetryAt: row.next_retry_at ? new Date(row.next_retry_at) : null, + deduplicationKey: row.deduplication_key ?? null, }; } } diff --git a/listener/src/types/index.ts b/listener/src/types/index.ts index 7d9b57cc..ae708e8c 100644 --- a/listener/src/types/index.ts +++ b/listener/src/types/index.ts @@ -191,6 +191,8 @@ export interface RetrySchedulerOptions { multiplier: number; maxDelayMs: number; jitter: boolean; + /** Request timeout for outbound webhook delivery (ms). Default: 10 000. */ + webhookTimeoutMs: number; /** * Retry-policy ceiling on total attempts. Mirrors `RetrySchedulerConfig`; * `undefined` leaves each notification's own `maxRetries` in control. diff --git a/listener/src/types/scheduled-notification.ts b/listener/src/types/scheduled-notification.ts index e2123a90..bfb5c854 100644 --- a/listener/src/types/scheduled-notification.ts +++ b/listener/src/types/scheduled-notification.ts @@ -43,6 +43,8 @@ export interface ScheduledNotification { cancellationReason?: string | null; /** When the next retry should be attempted (null if not retrying). */ nextRetryAt?: Date | null; + /** Stable logical key used to prevent duplicate delivery of the same notification. */ + deduplicationKey?: string | null; } export interface CreateScheduledNotificationInput { @@ -51,6 +53,7 @@ export interface CreateScheduledNotificationInput { targetRecipient: string; executeAt: Date; maxRetries?: number; + deduplicationKey?: string; eventId?: string; contractAddress?: string; priority?: number; @@ -80,6 +83,7 @@ export interface ScheduledNotificationRow { priority: number; metadata: string | null; next_retry_at: string | null; + deduplication_key: string | null; } export interface NotificationExecutionLog { diff --git a/listener/src/utils/benchmark-utils.ts b/listener/src/utils/benchmark-utils.ts index 8bc10816..6c220272 100644 --- a/listener/src/utils/benchmark-utils.ts +++ b/listener/src/utils/benchmark-utils.ts @@ -235,7 +235,7 @@ export function exportMetricsToJSON(metrics: BenchmarkMetrics[], outputPath: str } export class ResourceMonitor { - private intervalId?: NodeJS.Timer; + private intervalId?: ReturnType; private samples: Array<{ timestamp: number; memory: NodeJS.MemoryUsage; cpu: NodeJS.CpuUsage }> = []; start(intervalMs: number = 1000): void {