Skip to content
Open
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
8 changes: 8 additions & 0 deletions listener/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -365,3 +365,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
13 changes: 0 additions & 13 deletions listener/src/config.ts
Original file line number Diff line number Diff line change
@@ -1,20 +1,7 @@
import { Config, ContractConfig, DiscordConfig, WebhookSecret, AppCleanupConfig, EventQueueConfig, RetrySchedulerOptions, AnalyticsConfig, ExpirationConfig, ApiKey, CircuitBreakerConfig, BackfillConfig, LoggingConfig, ApiConfig } 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 { validateSecrets } from './config/validate-secrets';
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 { Config, ContractConfig, DiscordConfig, WebhookSecret, AppCleanupConfig, EventQueueConfig, RetrySchedulerOptions, AnalyticsConfig, ExpirationConfig, ApiKey, BackfillConfig, LoggingConfig, ApiConfig, RetryPolicyOptions } from './types';
import {
DEFAULT_RETRYABLE_FAILURE_TYPES,
RETRY_FAILURE_TYPES,
RetryFailureType,
parseRetryableFailureTypes,
} from './services/retry-policy';
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, RpcFallbackConfig } from './types';
import { validateSecrets } from './config/validate-secrets';
import {
SUPPORTED_LOG_FORMATS,
SUPPORTED_LOG_LEVELS,
Expand Down
84 changes: 20 additions & 64 deletions listener/src/services/event-subscriber.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,26 +31,7 @@ export class EventSubscriber {
private eventQueue: EventProcessingQueue | null = null;
private expirationService: NotificationExpirationService | null = null;
private lastSuccessfulPollAt: number | null = null;
private circuitBreaker: CircuitBreaker | null = null;
private backfillStartLedger: number | null = null;
/** Cold-start ledger resolved once per session by resolveBackfillStartLedger(). */
private backfillStartLedger: number | null = null;
private backfillStartLedger: number | null = null;

public get server(): StellarSDK.rpc.Server {
return this.rpcManager.getActiveServer();
}

public set server(val: StellarSDK.rpc.Server) {
const activeIndex = (this.rpcManager as any).activeIndex;
if ((this.rpcManager as any).endpoints && (this.rpcManager as any).endpoints[activeIndex]) {
(this.rpcManager as any).endpoints[activeIndex].server = val;
}
}

public getRpcManager(): StellarRpcManager {
return this.rpcManager;
}

constructor(config: Config, deduplicationService?: EventDeduplicationService) {
this.config = config;
Expand Down Expand Up @@ -382,8 +363,6 @@ export class EventSubscriber {
contractConfig: ContractConfig
): Promise<StellarSDK.rpc.Api.GetEventsResponse> {
const lastCursor = this.lastCursors.get(contractConfig.address);
const limit = this.config.eventBatchSize ?? 100;

let request: StellarSDK.rpc.Api.GetEventsRequest;

if (lastCursor) {
Expand Down Expand Up @@ -430,32 +409,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(
Expand All @@ -467,22 +420,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,
Expand Down Expand Up @@ -521,12 +465,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', {
Expand Down Expand Up @@ -583,10 +528,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;
Expand Down
31 changes: 2 additions & 29 deletions listener/src/types/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -84,35 +84,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. */
Expand Down
Loading