From b6206f208cf07a6c5843c5b077489f3b769199af Mon Sep 17 00:00:00 2001 From: DavisVT Date: Tue, 29 Sep 2026 17:00:46 +0100 Subject: [PATCH] Add event processing dry-run mode (#834) - Add dry-run configuration option to Config interface - Load DRY_RUN environment variable in config loader - Modify event-subscriber to skip persistence and delivery in dry-run mode - Add dry-run logging output for validation results - Add tests for dry-run mode functionality - Document DRY_RUN environment variable in .env.example Acceptance Criteria: - Dry-run execution produces useful output - Database state remains unchanged - Notification delivery is disabled --- listener/.env.example | 8 ++ listener/src/config.ts | 4 +- .../src/services/event-subscriber.test.ts | 130 +++++++++++++++++- listener/src/services/event-subscriber.ts | 49 +++---- listener/src/types/index.ts | 2 + 5 files changed, 159 insertions(+), 34 deletions(-) diff --git a/listener/.env.example b/listener/.env.example index c5f610b6..9b7681b9 100644 --- a/listener/.env.example +++ b/listener/.env.example @@ -320,3 +320,11 @@ EXPIRATION_DEFAULT_MS=86400000 # Number of days to retain analytics snapshots in the database. # ANALYTICS_SNAPSHOT_RETENTION_DAYS=30 + +# ----------------------------------------------------------------------------- +# Dry-Run Mode +# ----------------------------------------------------------------------------- + +# Enable dry-run mode to parse and validate events without persisting or delivering notifications. +# Set to 'true' to enable dry-run mode (default: false). +# DRY_RUN=false diff --git a/listener/src/config.ts b/listener/src/config.ts index 3dbdd3f6..1298500f 100644 --- a/listener/src/config.ts +++ b/listener/src/config.ts @@ -1,7 +1,6 @@ -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 { SUPPORTED_LOG_FORMATS, SUPPORTED_LOG_LEVELS, @@ -306,6 +305,7 @@ export function loadConfig(): Config { backfill: loadBackfillConfig(), logging: loadLoggingConfig(), api: loadApiConfig(), + dryRun: trimEnv('DRY_RUN') === 'true', }; } diff --git a/listener/src/services/event-subscriber.test.ts b/listener/src/services/event-subscriber.test.ts index 1abd452a..b498d960 100644 --- a/listener/src/services/event-subscriber.test.ts +++ b/listener/src/services/event-subscriber.test.ts @@ -59,7 +59,6 @@ const testConfig: Config = { reconnectDelayMs: 100, eventsApiPort: 8787, eventsApiCorsOrigin: 'http://localhost:5173', - maxPayloadSizeBytes: 64 * 1024, }; function createMockEvent( @@ -728,7 +727,134 @@ describe('EventSubscriber', () => { expect(preferenceStore.isCategoryEnabled).toHaveBeenCalledWith('global', 'discord'); }); }); -}); + + describe('dry-run mode', () => { + it('skips Discord notification delivery when dry-run is enabled', async () => { + const discordConfig = { + webhookUrl: 'https://discord.com/api/webhooks/test/webhook', + webhookId: 'test', + }; + const configWithDryRun: Config = { + ...testConfig, + discord: discordConfig, + dryRun: true, + }; + + mockGetEvents.mockResolvedValue({ + events: [createMockEvent({ id: 'event-dryrun' })], + cursor: 'cursor-dryrun', + }); + + const subscriber = new EventSubscriber(configWithDryRun); + await (subscriber as any).checkForEvents(); + + expect(mockDiscordService.sendEventNotification).not.toHaveBeenCalled(); + expect(mockLogger.info).toHaveBeenCalledWith( + 'Dry-run: event validated successfully (no persistence or delivery)', + expect.objectContaining({ + eventId: 'event-dryrun', + eventName: 'TaskCreated', + contractAddress: contractConfig.address, + }) + ); + }); + + it('logs dry-run status in processing event log', async () => { + const configWithDryRun: Config = { + ...testConfig, + dryRun: true, + }; + + mockGetEvents.mockResolvedValue({ + events: [createMockEvent({ id: 'event-dryrun-status' })], + cursor: 'cursor-dryrun-status', + }); + + const subscriber = new EventSubscriber(configWithDryRun); + await (subscriber as any).checkForEvents(); + + expect(mockLogger.info).toHaveBeenCalledWith( + 'Processing event', + expect.objectContaining({ + eventId: 'event-dryrun-status', + dryRun: true, + }) + ); + }); + + it('logs dry-run outcome in event processing complete log', async () => { + const configWithDryRun: Config = { + ...testConfig, + dryRun: true, + }; + + mockGetEvents.mockResolvedValue({ + events: [createMockEvent({ id: 'event-dryrun-outcome' })], + cursor: 'cursor-dryrun-outcome', + }); + + const subscriber = new EventSubscriber(configWithDryRun); + await (subscriber as any).checkForEvents(); + + expect(mockLogger.info).toHaveBeenCalledWith( + 'Event processing complete', + expect.objectContaining({ + eventId: 'event-dryrun-outcome', + outcome: 'dry_run', + }) + ); + }); + + it('processes events normally when dry-run is disabled', async () => { + const discordConfig = { + webhookUrl: 'https://discord.com/api/webhooks/test/webhook', + webhookId: 'test', + }; + const configWithoutDryRun: Config = { + ...testConfig, + discord: discordConfig, + dryRun: false, + }; + + mockGetEvents.mockResolvedValue({ + events: [createMockEvent({ id: 'event-normal' })], + cursor: 'cursor-normal', + }); + + const subscriber = new EventSubscriber(configWithoutDryRun); + await (subscriber as any).checkForEvents(); + + expect(mockDiscordService.sendEventNotification).toHaveBeenCalled(); + expect(mockLogger.info).not.toHaveBeenCalledWith( + 'Dry-run: event validated successfully (no persistence or delivery)', + expect.any(Object) + ); + }); + + it('skips persistent deduplication when dry-run is enabled', async () => { + const configWithDryRun: Config = { + ...testConfig, + dryRun: true, + }; + + mockGetEvents.mockResolvedValue({ + events: [createMockEvent({ id: 'event-dedup-dryrun' })], + cursor: 'cursor-dedup-dryrun', + }); + + const subscriber = new EventSubscriber(configWithDryRun); + await (subscriber as any).checkForEvents(); + + // In dry-run mode, deduplication service is not called + // This is verified by the fact that no errors are thrown and processing completes + expect(mockLogger.info).toHaveBeenCalledWith( + 'Dry-run: event validated successfully (no persistence or delivery)', + expect.objectContaining({ + eventId: 'event-dedup-dryrun', + }) + ); + }); + }); describe('notification expiration (Task 3: Requirements 2.1, 2.2, 2.3)', () => { const DEFAULT_EXPIRATION_MS = 24 * 60 * 60 * 1000; // 24 hours diff --git a/listener/src/services/event-subscriber.ts b/listener/src/services/event-subscriber.ts index 82dcda7c..cfe74d3f 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) { @@ -391,9 +367,10 @@ export class EventSubscriber { ): Promise { const eventStart = Date.now(); const eventName = getEventName(event.topic); + const isDryRun = this.config.dryRun === true; // Check persistent deduplication first (to catch reorg duplicates) - if (this.deduplicationService) { + if (this.deduplicationService && !isDryRun) { const duplicate = await this.deduplicationService.isDuplicate(event.id, contractConfig.address); if (duplicate.isDuplicate) { logger.warn('Skipping event: already processed (persistent deduplication)', { @@ -440,12 +417,13 @@ export class EventSubscriber { type: displayEvent.type, topic: displayEvent.topic, value: displayEvent.value, + dryRun: isDryRun, }); let notificationSent = false; let processingError: string | undefined; - if (this.discordService) { + if (this.discordService && !isDryRun) { const userId = contractConfig.userId ?? 'global'; if (!preferenceStore.isCategoryEnabled(userId, 'discord')) { logger.info('Skipping Discord notification: category disabled by user preferences', { @@ -483,8 +461,8 @@ export class EventSubscriber { } } - // Record the processed event for persistent deduplication - if (this.deduplicationService) { + // Record the processed event for persistent deduplication (skip in dry-run) + if (this.deduplicationService && !isDryRun) { await this.deduplicationService.recordProcessedEvent( event.id, contractConfig.address, @@ -502,10 +480,21 @@ export class EventSubscriber { correlationId, eventId: event.id, notificationSent, - outcome: !this.discordService || notificationSent ? 'success' : 'failure', + outcome: isDryRun ? 'dry_run' : (!this.discordService || notificationSent ? 'success' : 'failure'), durationMs: Date.now() - eventStart, }); + if (isDryRun) { + logger.info('Dry-run: event validated successfully (no persistence or delivery)', { + requestId: correlationId, + correlationId, + eventId: event.id, + eventName, + contractAddress: contractConfig.address, + }); + return true; + } + if (!this.discordService) return true; if (notificationSent) return true; if (processingError && this.retryQueue) return true; diff --git a/listener/src/types/index.ts b/listener/src/types/index.ts index 1d212d13..b63dba62 100644 --- a/listener/src/types/index.ts +++ b/listener/src/types/index.ts @@ -67,6 +67,8 @@ export interface Config { backfill?: BackfillConfig; logging?: LoggingConfig; api?: ApiConfig; + /** Dry-run mode: parse and validate events without persisting or delivering notifications. */ + dryRun?: boolean; } /** Observability settings, sourced from LOG_LEVEL / LOG_FORMAT. */