From 4199e791b2e232f6899a5818a20f56138658bf2f Mon Sep 17 00:00:00 2001 From: Vvictor-commits Date: Wed, 30 Sep 2026 01:57:53 +0100 Subject: [PATCH 1/2] feat: add configurable webhook delivery timeout (WEBHOOK_DELIVERY_TIMEOUT_MS) - Add webhookTimeoutMs field to RetrySchedulerOptions and RetrySchedulerConfig - Read WEBHOOK_DELIVERY_TIMEOUT_MS env var in loadRetrySchedulerConfig() (default 10 000 ms) - Pass timeout to WebhookDeliveryService constructor in RetryScheduler - Document new env var in .env.example Fix pre-existing merge artifacts: - request-id.ts: restore missing closing brace on generateCorrelationId() - security-headers.ts: fix invalid 'http.ServerResponse' import type syntax - discord-notification.ts: remove dead unreachable return in scvString case; restore missing closing brace on sanitizeForDiscord() - index.ts: collapse duplicate healthMonitor construction, remove duplicate subscriber declaration, fix broken shutdown try/catch block - events-server.ts: remove duplicate TemplateService/handleTemplateRoutes imports; add missing handleApiError and applyRequestIdMiddleware imports - config.ts: remove duplicate types import line; add missing validateSecrets import - event-subscriber.ts: remove duplicate processableEvents declaration; collapse duplicate getContractEvents request block; add missing backfillStartLedger property --- listener/.env.example | 4 +++ listener/src/api/events-server.ts | 5 +--- listener/src/config.ts | 5 ++-- listener/src/index.ts | 19 +++++--------- listener/src/middleware/security-headers.ts | 4 +-- listener/src/services/discord-notification.ts | 3 ++- listener/src/services/event-subscriber.ts | 26 +------------------ listener/src/services/retry-scheduler.ts | 6 ++++- listener/src/types/index.ts | 2 ++ listener/src/utils/request-id.ts | 3 +++ 10 files changed, 29 insertions(+), 48 deletions(-) diff --git a/listener/.env.example b/listener/.env.example index c5f610b6..7b5d3867 100644 --- a/listener/.env.example +++ b/listener/.env.example @@ -190,6 +190,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 + # ----------------------------------------------------------------------------- # Scheduled Notification Scheduler # ----------------------------------------------------------------------------- diff --git a/listener/src/api/events-server.ts b/listener/src/api/events-server.ts index d517ad21..fd5f3283 100644 --- a/listener/src/api/events-server.ts +++ b/listener/src/api/events-server.ts @@ -9,14 +9,11 @@ import { NotificationAPI } from '../services/notification-api'; import { NotificationType } from '../types/scheduled-notification'; import logger from '../utils/logger'; import { generateRequestId, resolveCorrelationId } from '../utils/request-id'; -import { TemplateService } from '../services/template-service'; -import { handleTemplateRoutes } from './template-routes'; -import { sendOk, sendErr, sendJson, ErrorCode } from '../utils/response'; 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 { NotificationHistoryService } from '../services/notification-history'; import { SearchSuggestionService } from '../services/search-suggestion'; import { NotificationSearchService } from '../services/notification-search-service'; diff --git a/listener/src/config.ts b/listener/src/config.ts index 3dbdd3f6..cd653777 100644 --- a/listener/src/config.ts +++ b/listener/src/config.ts @@ -1,7 +1,7 @@ -import { Config, ContractConfig, DiscordConfig, WebhookSecret, AppCleanupConfig, EventQueueConfig, RetrySchedulerOptions, AnalyticsConfig, ExpirationConfig, ApiKey } from './types'; +import { Config, ContractConfig, DiscordConfig, WebhookSecret, AppCleanupConfig, EventQueueConfig, RetrySchedulerOptions, AnalyticsConfig, ExpirationConfig, ApiKey, BackfillConfig, LoggingConfig, ApiConfig } from './types'; import { validateCorsOrigin, CorsValidationError } from './utils/cors-validator'; 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 { validateSecrets } from './config/validate-secrets'; import { SUPPORTED_LOG_FORMATS, SUPPORTED_LOG_LEVELS, @@ -192,6 +192,7 @@ function loadAnalyticsConfig(): AnalyticsConfig { function loadRetrySchedulerConfig(): 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/index.ts b/listener/src/index.ts index 169ab237..28992885 100644 --- a/listener/src/index.ts +++ b/listener/src/index.ts @@ -67,19 +67,13 @@ async function main() { const db = await initializeDatabase(config.databasePath); repository = new ScheduledNotificationRepository(db); - + healthMonitor = new NotificationHealthMonitor(null, getWorkerManager(), { repository, getLastSuccessfulPoll: () => subscriber?.getLastSuccessfulPoll() ?? null, - }); - getUptimeMs: () => Date.now() - PROCESS_START_TIME, }); - healthMonitor = new NotificationHealthMonitor(null, getWorkerManager(), { - repository, - }); - // Rebuild registry with configured event TTL if (config.cleanup) { eventRegistry.setTtlMs(config.cleanup.eventRetentionMs); @@ -171,8 +165,7 @@ async function main() { healthMonitor.start(); } - subscriber = new EventSubscriber(config, deduplicationService); - const subscriber = new EventSubscriber(config, deduplicationService ?? undefined); + subscriber = new EventSubscriber(config, deduplicationService ?? undefined); await subscriber.start(); let isShuttingDown = false; @@ -216,11 +209,11 @@ async function main() { await retryScheduler.stop(); } - if (subscriber) { - await subscriber.stop(); - } + if (subscriber) { + await subscriber.stop(); + } - eventsServer.close(); + eventsServer.close(); logger.info('Graceful shutdown completed successfully', { signal }); process.exit(0); diff --git a/listener/src/middleware/security-headers.ts b/listener/src/middleware/security-headers.ts index d46499ff..8f8ff1ce 100644 --- a/listener/src/middleware/security-headers.ts +++ b/listener/src/middleware/security-headers.ts @@ -10,13 +10,13 @@ * See: https://owasp.org/www-project-secure-headers/ */ -import type { http.ServerResponse } from 'http'; +import type { ServerResponse } from 'http'; const isLocalhost = (hostname: string): boolean => hostname === 'localhost' || hostname === '127.0.0.1' || hostname === '::1'; export function addSecurityHeaders( - res: http.ServerResponse, + res: ServerResponse, options: { productionOrigin?: string } = {}, ): void { const origin = res.getHeader('Access-Control-Allow-Origin') as string | undefined; diff --git a/listener/src/services/discord-notification.ts b/listener/src/services/discord-notification.ts index 47b79249..652fea97 100644 --- a/listener/src/services/discord-notification.ts +++ b/listener/src/services/discord-notification.ts @@ -58,6 +58,8 @@ export function sanitizeForDiscord(text: string): string { return text .replace(MENTION_PATTERN, '[mention removed]') .replace(MARKDOWN_CHARS, '\\$1'); +} + // Internal helpers // --------------------------------------------------------------------------- @@ -452,7 +454,6 @@ export class DiscordNotificationService { return String(value.i64()); case StellarSDK.xdr.ScValType.scvString(): { const strVal = value.str().toString(); - return strVal.length > MAX_DISCORD_FIELD_VALUE_LENGTH ? strVal.slice(0, MAX_DISCORD_FIELD_VALUE_LENGTH) + '...' : strVal; const truncated = strVal.length > 500 ? strVal.slice(0, 500) + '...' : strVal; return sanitizeForDiscord(truncated); } diff --git a/listener/src/services/event-subscriber.ts b/listener/src/services/event-subscriber.ts index 82dcda7c..58dc437e 100644 --- a/listener/src/services/event-subscriber.ts +++ b/listener/src/services/event-subscriber.ts @@ -29,6 +29,7 @@ export class EventSubscriber { private eventQueue: EventProcessingQueue | null = null; private expirationService: NotificationExpirationService | null = null; private lastSuccessfulPollAt: number | null = null; + private backfillStartLedger: number | null = null; constructor(config: Config, deduplicationService?: EventDeduplicationService) { this.config = config; @@ -170,9 +171,6 @@ export class EventSubscriber { }); } } - const processableEvents = events.filter((event: StellarSDK.rpc.Api.EventResponse) => - this.shouldProcessEvent(event, contractConfig, requestId) - ); if (events.length > 0) { logger.info('Received events', { @@ -339,28 +337,6 @@ export class EventSubscriber { contractConfig: ContractConfig ): Promise { const lastCursor = this.lastCursors.get(contractConfig.address); - const request: StellarSDK.rpc.Api.GetEventsRequest = lastCursor - ? { - filters: [ - { - contractIds: [contractConfig.address], - type: 'contract', - }, - ], - cursor: lastCursor, - limit: this.config.eventBatchSize, - } - : { - filters: [ - { - contractIds: [contractConfig.address], - type: 'contract', - }, - ], - startLedger: 1, - limit: this.config.eventBatchSize, - }; - let request: StellarSDK.rpc.Api.GetEventsRequest; if (lastCursor) { diff --git a/listener/src/services/retry-scheduler.ts b/listener/src/services/retry-scheduler.ts index fa65685b..41193cfa 100644 --- a/listener/src/services/retry-scheduler.ts +++ b/listener/src/services/retry-scheduler.ts @@ -26,6 +26,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; } export const RETRY_SCHEDULER_DEFAULTS: RetrySchedulerConfig = { @@ -37,6 +39,7 @@ export const RETRY_SCHEDULER_DEFAULTS: RetrySchedulerConfig = { multiplier: 2, maxDelayMs: 60 * 60 * 1_000, jitter: true, + webhookTimeoutMs: 10_000, }; /** @@ -89,7 +92,8 @@ export class RetryScheduler { this.processorId = this.config.processorId ?? `retry-${uuidv4()}`; this.repository = repository; this.discordService = discordService ?? null; - this.webhookDeliveryService = webhookDeliveryService ?? new WebhookDeliveryService(); + this.webhookDeliveryService = + webhookDeliveryService ?? new WebhookDeliveryService({ timeoutMs: this.config.webhookTimeoutMs }); } async start(): Promise { diff --git a/listener/src/types/index.ts b/listener/src/types/index.ts index 1d212d13..0bb91baa 100644 --- a/listener/src/types/index.ts +++ b/listener/src/types/index.ts @@ -135,6 +135,8 @@ export interface RetrySchedulerOptions { multiplier: number; maxDelayMs: number; jitter: boolean; + /** Request timeout for outbound webhook delivery (ms). Default: 10 000. */ + webhookTimeoutMs: number; } export interface AnalyticsConfig { diff --git a/listener/src/utils/request-id.ts b/listener/src/utils/request-id.ts index 0a9c61f7..ef670efc 100644 --- a/listener/src/utils/request-id.ts +++ b/listener/src/utils/request-id.ts @@ -14,6 +14,9 @@ export function generateRequestId(): string { */ export function generateCorrelationId(): string { return randomUUID(); +} + +/** * Client-supplied request IDs must be printable ASCII tokens of bounded length. * Rejects empty values, control characters, whitespace, and oversized strings * so untrusted header content is never reused as a log/trace key (#686). From 182da8b06f3dfba97ba4d57cbc3966df02d34e0b Mon Sep 17 00:00:00 2001 From: Vvictor-commits Date: Wed, 30 Sep 2026 02:26:33 +0100 Subject: [PATCH 2/2] feat: add persistent deduplication keys for notifications MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Add migration 003: ALTER TABLE adds deduplication_key TEXT with a UNIQUE partial index on scheduled_notifications - Update schema.sql with the new column and index for fresh installs - Extend ScheduledNotification, ScheduledNotificationRow, and CreateScheduledNotificationInput types with deduplicationKey field - Update repository.create(): inserts deduplication_key; on UNIQUE constraint violation returns the existing notification id so callers are idempotent without throwing — covers concurrent requests at the DB layer Fix pre-existing type errors surfaced by the workflow check: - migration-system.ts: cast db.all() result through unknown before mapping; stringify error in logger.error call - discord-notification.ts: use local variable instead of accessing embed.title / embed.footer after optional narrowing - notification-retry-queue.ts: annotate entry param type in .map() - schema/compatibility-check.ts: use contractEvent.dataFields instead of undeclared dataFields; guard expectedTopics?.length comparison - benchmark-utils.ts: use ReturnType instead of NodeJS.Timer - notification-stats-cache.ts: replace cache.has() with cache.get() check (NodeCache typings do not expose has()) --- listener/src/database/migration-system.ts | 4 +- listener/src/database/schema.sql | 7 ++- .../003-notification-deduplication-key.ts | 34 +++++++++++ listener/src/schema/compatibility-check.ts | 6 +- listener/src/services/discord-notification.ts | 5 +- .../src/services/notification-retry-queue.ts | 2 +- .../src/services/notification-stats-cache.ts | 2 +- .../scheduled-notification-repository.ts | 56 ++++++++++++++----- listener/src/types/scheduled-notification.ts | 4 ++ listener/src/utils/benchmark-utils.ts | 2 +- 10 files changed, 95 insertions(+), 27 deletions(-) create mode 100644 listener/src/migrations/003-notification-deduplication-key.ts 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 8342d5b1..b8fa1644 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/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 1822c8cb..60458263 100644 --- a/listener/src/schema/compatibility-check.ts +++ b/listener/src/schema/compatibility-check.ts @@ -159,9 +159,7 @@ function checkEventCompatibility( // 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 (contractEvent.topics.length < consumer.expectedTopics?.length) { - // Contract has fewer topics than consumer expects - // This could be breaking if the consumer relies on specific topic positions + if (consumer.expectedTopics && contractEvent.topics.length < consumer.expectedTopics.length) { const missingTopics = consumer.expectedTopics.filter( (t) => !contractEvent.topics.includes(t) ); @@ -175,7 +173,7 @@ function checkEventCompatibility( // 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 = dataFields.filter( + const newDataFields = contractEvent.dataFields.filter( (f) => !consumer.expectedFields.includes(f.name) ); safeAdditions.push( diff --git a/listener/src/services/discord-notification.ts b/listener/src/services/discord-notification.ts index 652fea97..40d1c8a3 100644 --- a/listener/src/services/discord-notification.ts +++ b/listener/src/services/discord-notification.ts @@ -379,7 +379,7 @@ export class DiscordNotificationService { if (title.length > MAX_DISCORD_EMBED_TITLE_LENGTH) { title = title.slice(0, MAX_DISCORD_EMBED_TITLE_LENGTH - 3) + '...'; logger.warn('Discord embed title truncated', { - originalLength: embed.title.length, + originalLength: title.length + 3, maxLength: MAX_DISCORD_EMBED_TITLE_LENGTH, }); } @@ -399,9 +399,10 @@ export class DiscordNotificationService { 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: embed.footer.text.length, + originalLength, maxLength: MAX_DISCORD_FOOTER_TEXT_LENGTH, }); } diff --git a/listener/src/services/notification-retry-queue.ts b/listener/src/services/notification-retry-queue.ts index 765ca629..12b71e0c 100644 --- a/listener/src/services/notification-retry-queue.ts +++ b/listener/src/services/notification-retry-queue.ts @@ -290,6 +290,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 2b2895f2..dbedea5f 100644 --- a/listener/src/services/notification-stats-cache.ts +++ b/listener/src/services/notification-stats-cache.ts @@ -144,7 +144,7 @@ export class NotificationStatsCache { * @notice Check if stats are currently cached */ has(): boolean { - return this.cache.has(this.CACHE_KEY); + return this.cache.get(this.CACHE_KEY) !== undefined; } /** diff --git a/listener/src/services/scheduled-notification-repository.ts b/listener/src/services/scheduled-notification-repository.ts index bfbc0a02..ac096105 100644 --- a/listener/src/services/scheduled-notification-repository.ts +++ b/listener/src/services/scheduled-notification-repository.ts @@ -27,7 +27,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); @@ -37,8 +39,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); @@ -54,21 +56,44 @@ export class ScheduledNotificationRepository { input.contractAddress ?? null, input.priority ?? 5, input.metadata ? JSON.stringify(input.metadata) : null, + input.deduplicationKey ?? null, ]; - 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, - }); + 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, + }); - 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; + } } /** @@ -832,6 +857,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/scheduled-notification.ts b/listener/src/types/scheduled-notification.ts index cdaa92d8..f99cace5 100644 --- a/listener/src/types/scheduled-notification.ts +++ b/listener/src/types/scheduled-notification.ts @@ -41,6 +41,8 @@ export interface ScheduledNotification { metadata?: string | null; // JSON string /** 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 { @@ -49,6 +51,7 @@ export interface CreateScheduledNotificationInput { targetRecipient: string; executeAt: Date; maxRetries?: number; + deduplicationKey?: string; eventId?: string; contractAddress?: string; priority?: number; @@ -78,6 +81,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 {