diff --git a/listener/.env.example b/listener/.env.example index de4548d9..e3eb452e 100644 --- a/listener/.env.example +++ b/listener/.env.example @@ -387,3 +387,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 ce85f0f7..ecd9c38a 100644 --- a/listener/src/config.ts +++ b/listener/src/config.ts @@ -1,31 +1,7 @@ -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 { Config, ContractConfig, DiscordConfig, WebhookSecret, AppCleanupConfig, EventQueueConfig, RetrySchedulerOptions, AnalyticsConfig, ExpirationConfig, ApiKey, BackfillConfig, LoggingConfig, ApiConfig } from './types'; 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, diff --git a/listener/src/services/event-subscriber.ts b/listener/src/services/event-subscriber.ts index b9e9cea7..cc81f1d0 100644 --- a/listener/src/services/event-subscriber.ts +++ b/listener/src/services/event-subscriber.ts @@ -423,7 +423,6 @@ export class EventSubscriber { } const lastCursor = this.lastCursors.get(contractConfig.address); - let request: StellarSDK.rpc.Api.GetEventsRequest; if (lastCursor) { @@ -476,32 +475,6 @@ export class EventSubscriber { } return await this.server.getEvents(request); - 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: await this.resolveBackfillStartLedger(), - limit: this.config.eventBatchSize, - }; - - return await this.rpcManager.executeWithFallback( - (server) => server.getEvents(request), - { operationName: `getEvents(${contractConfig.address})` } - ); } private async processEvent( @@ -513,22 +486,13 @@ export class EventSubscriber { correlationId = correlationId || requestId || generateCorrelationId(); const eventStart = Date.now(); const eventName = getEventName(event.topic); + const isDryRun = this.config.dryRun === true; - // Atomically claim the event before doing any work. Only one concurrent - // processor (poll cycle, backfill, queue worker or another listener - // instance sharing the database) wins the claim; everyone else skips. - // A separate isDuplicate() check followed by a later write would leave a - // window in which two processors both send the notification. - if (this.deduplicationService) { - const claim = await this.deduplicationService.claimEvent( - event.id, - contractConfig.address, - event.ledger, - event.txHash, - event.type, - ); - if (!claim.claimed) { - logger.warn('Skipping event: already processed or in progress (persistent deduplication)', { + // Check persistent deduplication first (to catch reorg duplicates) + 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)', { requestId: correlationId, correlationId, eventId: event.id, @@ -567,12 +531,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', { @@ -629,10 +594,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 56855813..24d41547 100644 --- a/listener/src/types/index.ts +++ b/listener/src/types/index.ts @@ -92,35 +92,8 @@ export interface Config { backfill?: BackfillConfig; logging?: LoggingConfig; api?: ApiConfig; - circuitBreaker?: CircuitBreakerConfig; -} - -/** Configurable fallback RPC endpoint settings */ -export interface RpcFallbackConfig { - /** Array of fallback RPC URLs to try when primary fails */ - fallbackUrls: string[]; - /** Number of consecutive failures before marking endpoint unhealthy and failing over (default: 3) */ - failureThreshold: number; - /** Cooldown duration in ms before attempting to reuse a failed endpoint (default: 60000) */ - cooldownMs: number; - /** Timeout for RPC requests in milliseconds (default: 10000) */ - requestTimeoutMs: number; - /** Maximum number of endpoint retries for a single operation across pool (default: pool size) */ - maxRetries?: number; -} - -/** Operational status metrics for an RPC endpoint */ -export interface RpcEndpointStatus { - url: string; - isPrimary: boolean; - status: 'healthy' | 'degraded' | 'unhealthy'; - consecutiveFailures: number; - totalRequests: number; - totalSuccesses: number; - totalFailures: number; - lastFailureTime: number | null; - lastSuccessTime: number | null; - lastError: string | null; + /** Dry-run mode: parse and validate events without persisting or delivering notifications. */ + dryRun?: boolean; } /** Observability settings, sourced from LOG_LEVEL / LOG_FORMAT. */