From 905c932a63adb92b00d31c8c0245d91ebd52615a Mon Sep 17 00:00:00 2001 From: Abdullahi Abubakar Sadiq Date: Fri, 28 Aug 2026 16:53:21 +0000 Subject: [PATCH] feat(listener): make blockchain event batch size configurable (#627) --- ENVIRONMENT_VARIABLES_AND_SECRETS.md | 1 + listener/.env.example | 3 ++ listener/src/config.test.ts | 32 ++++++++++++++++++- listener/src/config.ts | 5 +++ .../services/event-subscriber-reorg.test.ts | 1 + .../src/services/event-subscriber.test.ts | 15 +++++++++ listener/src/services/event-subscriber.ts | 4 +-- listener/src/types/index.ts | 2 ++ 8 files changed, 60 insertions(+), 3 deletions(-) diff --git a/ENVIRONMENT_VARIABLES_AND_SECRETS.md b/ENVIRONMENT_VARIABLES_AND_SECRETS.md index 4f527416..cd765f7c 100644 --- a/ENVIRONMENT_VARIABLES_AND_SECRETS.md +++ b/ENVIRONMENT_VARIABLES_AND_SECRETS.md @@ -126,6 +126,7 @@ Both variables must be provided together or neither. | Variable | Default | Required | Description | |---|---|---|---| | `POLL_INTERVAL_MS` | `30000` | No | How often the listener polls Stellar for new contract events (ms). | +| `EVENT_BATCH_SIZE` | `100` | No | Maximum number of blockchain events fetched in each polling cycle. Must be at least `1`. | | `MAX_RECONNECT_ATTEMPTS` | `5` | No | Maximum number of reconnect attempts when the RPC endpoint fails. | | `RECONNECT_DELAY_MS` | `5000` | No | Delay between reconnect attempts (ms). | diff --git a/listener/.env.example b/listener/.env.example index 1f6a7641..c95a212b 100644 --- a/listener/.env.example +++ b/listener/.env.example @@ -114,6 +114,9 @@ WEBHOOK_SECRETS=[{"id":"default","secret":"whsec_your_secret_here"}] # How often to poll Stellar for new contract events (ms). POLL_INTERVAL_MS=30000 +# Maximum number of blockchain events fetched in each polling cycle. +EVENT_BATCH_SIZE=100 + # Maximum number of reconnect attempts when the RPC endpoint fails. MAX_RECONNECT_ATTEMPTS=5 diff --git a/listener/src/config.test.ts b/listener/src/config.test.ts index 8ebb90ac..2b3c1b1f 100644 --- a/listener/src/config.test.ts +++ b/listener/src/config.test.ts @@ -1,4 +1,4 @@ -import { ConfigError, loadConfig } from './config'; +import { ConfigError, loadConfig, validateConfig } from './config'; describe('Config validation', () => { const originalEnv = process.env; @@ -55,6 +55,36 @@ describe('Config validation', () => { expect(() => loadConfig()).toThrow('EVENTS_API_PORT must be a valid integer, got "eighty"'); }); + it('loads the default blockchain event batch size', () => { + delete process.env.EVENT_BATCH_SIZE; + + expect(loadConfig().eventBatchSize).toBe(100); + }); + + it('loads a configured blockchain event batch size', () => { + process.env.EVENT_BATCH_SIZE = '250'; + + expect(loadConfig().eventBatchSize).toBe(250); + }); + + it('rejects a non-integer blockchain event batch size', () => { + process.env.EVENT_BATCH_SIZE = 'many'; + + expect(() => loadConfig()).toThrow( + 'EVENT_BATCH_SIZE must be a valid integer, got "many"' + ); + }); + + it('rejects a non-positive blockchain event batch size', () => { + process.env.EVENT_BATCH_SIZE = '0'; + + const config = loadConfig(); + + expect(() => validateConfig(config)).toThrow( + 'EVENT_BATCH_SIZE must be >= 1 (received: 0).' + ); + }); + it('loads default values when optional environment variables are omitted', () => { process.env.CONTRACT_ADDRESSES = JSON.stringify([{ address: 'CTEST', events: ['*'] }]); delete process.env.STELLAR_NETWORK; diff --git a/listener/src/config.ts b/listener/src/config.ts index 2d4e9e62..d1bcaf82 100644 --- a/listener/src/config.ts +++ b/listener/src/config.ts @@ -234,6 +234,7 @@ export function loadConfig(): Config { stellarNetworkPassphrase: trimEnv('STELLAR_NETWORK_PASSPHRASE') || 'Test SDF Network ; September 2015', contractAddresses: validateContractAddresses(rawContractAddresses), pollIntervalMs: parseIntegerEnv('POLL_INTERVAL_MS', '30000'), + eventBatchSize: parseIntegerEnv('EVENT_BATCH_SIZE', '100'), maxReconnectAttempts: parseIntegerEnv('MAX_RECONNECT_ATTEMPTS', '5'), reconnectDelayMs: parseIntegerEnv('RECONNECT_DELAY_MS', '5000'), eventsApiPort: parseIntegerEnv('EVENTS_API_PORT', '8787'), @@ -307,6 +308,10 @@ export function validateConfig(config: Config): void { ); } + if (config.eventBatchSize < 1) { + errors.push(`EVENT_BATCH_SIZE must be >= 1 (received: ${config.eventBatchSize}).`); + } + if (config.maxReconnectAttempts < 1) { errors.push( `MAX_RECONNECT_ATTEMPTS must be >= 1 (received: ${config.maxReconnectAttempts}).`, diff --git a/listener/src/services/event-subscriber-reorg.test.ts b/listener/src/services/event-subscriber-reorg.test.ts index c2038ec0..d8ba5ffd 100644 --- a/listener/src/services/event-subscriber-reorg.test.ts +++ b/listener/src/services/event-subscriber-reorg.test.ts @@ -54,6 +54,7 @@ const testConfig: Config = { stellarRpcUrl: 'https://soroban-testnet.stellar.org:443', contractAddresses: [contractConfig], pollIntervalMs: 30000, + eventBatchSize: 100, maxReconnectAttempts: 5, reconnectDelayMs: 100, eventsApiPort: 8787, diff --git a/listener/src/services/event-subscriber.test.ts b/listener/src/services/event-subscriber.test.ts index 2fd5bf23..7ef88c96 100644 --- a/listener/src/services/event-subscriber.test.ts +++ b/listener/src/services/event-subscriber.test.ts @@ -54,6 +54,7 @@ const testConfig: Config = { stellarRpcUrl: 'https://soroban-testnet.stellar.org:443', contractAddresses: [contractConfig], pollIntervalMs: 30000, + eventBatchSize: 100, maxReconnectAttempts: 5, reconnectDelayMs: 100, eventsApiPort: 8787, @@ -199,6 +200,20 @@ describe('EventSubscriber', () => { expect(mockGetEvents.mock.calls[1][0]).toMatchObject({ cursor: 'cursor-next' }); }); + it('uses the configured event batch size for RPC requests', async () => { + const configuredBatchSize = 25; + const subscriber = new EventSubscriber({ + ...testConfig, + eventBatchSize: configuredBatchSize, + }); + + await (subscriber as any).checkForEvents(); + + expect(mockGetEvents.mock.calls[0][0]).toMatchObject({ + limit: configuredBatchSize, + }); + }); + it('tracks cursors independently per contract', async () => { const secondContract: ContractConfig = { address: 'CBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB', diff --git a/listener/src/services/event-subscriber.ts b/listener/src/services/event-subscriber.ts index 8c3d9f2a..d31f1e5c 100644 --- a/listener/src/services/event-subscriber.ts +++ b/listener/src/services/event-subscriber.ts @@ -230,7 +230,7 @@ export class EventSubscriber { }, ], cursor: lastCursor, - limit: 100, + limit: this.config.eventBatchSize, } : { filters: [ @@ -240,7 +240,7 @@ export class EventSubscriber { }, ], startLedger: 1, - limit: 100, + limit: this.config.eventBatchSize, }; return await this.server.getEvents(request); diff --git a/listener/src/types/index.ts b/listener/src/types/index.ts index 5ff0c023..d9d0d837 100644 --- a/listener/src/types/index.ts +++ b/listener/src/types/index.ts @@ -46,6 +46,8 @@ export interface Config { stellarNetworkPassphrase: string; contractAddresses: ContractConfig[]; pollIntervalMs: number; + /** Maximum number of blockchain events fetched per polling cycle (default: 100). */ + eventBatchSize: number; maxReconnectAttempts: number; reconnectDelayMs: number; eventsApiPort: number;