From cb7344c51eec3826c063f74786f5d6c79c6a3234 Mon Sep 17 00:00:00 2001 From: Jess Date: Sat, 26 Sep 2026 13:32:30 +0100 Subject: [PATCH] feat(listener): add event processing checkpoints (#783) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Persist the latest successfully processed ledger position so the listener can resume from a known checkpoint after interruption, rather than replaying from the beginning on every restart. Changes ------- event-subscriber.ts - Add restoreCheckpoints() private method: loops over all configured contract addresses, calls deduplicationService.getLastCursor() for each, and populates this.lastCursors with the stored cursor string. Per-contract failures are caught and logged as warnings so a single bad DB row never prevents startup. Logs per-contract and a summary. - Call await this.restoreCheckpoints() from start() before queue start and the poll loop — the only moment restoration is needed. - Add private backfillStartLedger field (was referenced by resolveBackfillStartLedger but never declared). - Remove duplicate processableEvents re-declaration and duplicate request re-declaration in getContractEvents (pre-existing merge artifacts that caused SyntaxError: Identifier already declared). index.ts - Add deduplicationService = new EventDeduplicationService(db) immediately after initializeDatabase(). The let declaration and import already existed; the missing assignment meant every if (this.deduplicationService) guard inside EventSubscriber was always false at runtime — cursor persistence and checkpoint restoration were dead code. request-id.ts - Fix missing closing } on generateCorrelationId(): the JSDoc comment for isValidRequestId had leaked inside the function body, causing a parse error that prevented the module from loading. event-subscriber-checkpoint.test.ts (new) - 16 tests across three describe blocks using in-memory SQLite. Acceptance criteria met ----------------------- Restarting the listener does not require processing the entire history restoreCheckpoints() loads the persisted cursor from polling_cursors into this.lastCursors before the first poll; getContractEvents() uses the cursor branch (not startLedger) on the first call after restart. Checkpoints are only advanced after successful processing updatePollingCursor() is only called inside the response.cursor block at the end of checkForEvents(), after all events in the batch have been processed. A poll that throws never reaches that block. Recovery behaviour is covered by tests checkpoint persistence: cursor written after success, advances per poll, absent when no cursor returned, stable when RPC throws, independent per contract. checkpoint restoration: cursor loaded into lastCursors, used in first RPC call, cold start uses startLedger when no record, multi- contract restore, partial restore, no-op with null dedup service, logging verified. end-to-end: second subscriber instance resumes from the first's cursor; checkpoint only advances on success across failure cycles. --- listener/src/index.ts | 2 + .../event-subscriber-checkpoint.test.ts | 526 ++++++++++++++++++ listener/src/services/event-subscriber.ts | 72 ++- listener/src/utils/request-id.ts | 3 + 4 files changed, 578 insertions(+), 25 deletions(-) create mode 100644 listener/src/services/event-subscriber-checkpoint.test.ts diff --git a/listener/src/index.ts b/listener/src/index.ts index 169ab237..bbe889a1 100644 --- a/listener/src/index.ts +++ b/listener/src/index.ts @@ -66,6 +66,8 @@ async function main() { logger.info('Initializing database'); const db = await initializeDatabase(config.databasePath); + deduplicationService = new EventDeduplicationService(db); + repository = new ScheduledNotificationRepository(db); healthMonitor = new NotificationHealthMonitor(null, getWorkerManager(), { diff --git a/listener/src/services/event-subscriber-checkpoint.test.ts b/listener/src/services/event-subscriber-checkpoint.test.ts new file mode 100644 index 00000000..65225a7b --- /dev/null +++ b/listener/src/services/event-subscriber-checkpoint.test.ts @@ -0,0 +1,526 @@ +/** + * Event processing checkpoint tests (#783) + * + * Verifies that: + * - Checkpoints are stored (polling_cursors) after successful event processing. + * - Restarting the subscriber restores the persisted cursor into lastCursors + * so the first poll resumes from that position rather than ledger 1. + * - Checkpoints are only advanced when processing succeeds. + * - Recovery behaviour: a subscriber with no prior checkpoint starts from + * the beginning; one with a checkpoint resumes correctly. + * + * Uses an in-memory SQLite database (no file system state required) and a + * real EventDeduplicationService so the full persistence path is exercised. + */ + +import * as StellarSDK from '@stellar/stellar-sdk'; +import { EventSubscriber } from './event-subscriber'; +import { EventDeduplicationService } from './event-deduplication-service'; +import { Database } from '../database/database'; +import { Config, ContractConfig } from '../types'; +import logger from '../utils/logger'; + +// --------------------------------------------------------------------------- +// Module mocks +// --------------------------------------------------------------------------- + +jest.mock('../utils/logger', () => ({ + __esModule: true, + default: { + info: jest.fn(), + warn: jest.fn(), + error: jest.fn(), + debug: jest.fn(), + }, +})); + +jest.mock('./discord-notification', () => ({ + DiscordNotificationService: jest.fn().mockImplementation(() => ({ + sendEventNotification: jest.fn().mockResolvedValue(true), + })), +})); + +jest.mock('../store/preference-store', () => ({ + preferenceStore: { + isCategoryEnabled: jest.fn().mockReturnValue(true), + }, +})); + +jest.mock('./polling-metrics', () => ({ + pollingMetrics: { record: jest.fn() }, +})); + +const mockGetEvents = jest.fn(); + +jest.mock('@stellar/stellar-sdk', () => { + const actual = jest.requireActual('@stellar/stellar-sdk'); + return { + ...actual, + rpc: { + Server: jest.fn().mockImplementation(() => ({ + getEvents: mockGetEvents, + getLatestLedger: jest.fn().mockResolvedValue({ sequence: 50_000 }), + })), + }, + }; +}); + +// --------------------------------------------------------------------------- +// Fixtures +// --------------------------------------------------------------------------- + +const CONTRACT_ADDRESS = 'CCEMX6Q5V5F5F5F5F5F5F5F5F5F5F5F5F5F5F5F5F5F5F5F5F5F5F5'; + +const contractConfig: ContractConfig = { + address: CONTRACT_ADDRESS, + events: ['*'], +}; + +const testConfig: Config = { + stellarNetwork: 'testnet', + stellarNetworkPassphrase: 'Test SDF Network ; September 2015', + stellarRpcUrl: 'https://soroban-testnet.stellar.org:443', + contractAddresses: [contractConfig], + pollIntervalMs: 30000, + eventBatchSize: 100, + maxReconnectAttempts: 5, + reconnectDelayMs: 100, + eventsApiPort: 8787, + eventsApiCorsOrigin: 'http://localhost:5173', +}; + +function makeEvent( + id: string, + ledger: number, + overrides: Partial = {} +): StellarSDK.rpc.Api.EventResponse { + return { + id, + type: 'contract', + ledger, + ledgerClosedAt: '2026-01-01T00:00:00Z', + transactionIndex: 0, + operationIndex: 0, + inSuccessfulContractCall: true, + txHash: `tx-${id}`, + // Use plain arrays/objects — we never decode these in checkpoint tests + // and xdr is only available after jest.mock runs, not at module scope. + topic: [] as any, + value: {} as any, + ...overrides, + }; +} + +// --------------------------------------------------------------------------- +// Setup / teardown +// --------------------------------------------------------------------------- + +async function setupDb(): Promise<{ db: Database; dedup: EventDeduplicationService }> { + const db = new Database(':memory:'); + await db.initialize(); + const dedup = new EventDeduplicationService(db); + return { db, dedup }; +} + +// --------------------------------------------------------------------------- +// Tests +// --------------------------------------------------------------------------- + +describe('EventSubscriber — checkpoint persistence', () => { + let db: Database; + let dedup: EventDeduplicationService; + + beforeEach(async () => { + jest.clearAllMocks(); + ({ db, dedup } = await setupDb()); + // Default: no events, empty cursor + mockGetEvents.mockResolvedValue({ events: [], cursor: undefined }); + }); + + afterEach(async () => { + await db.close(); + }); + + it('persists the cursor to polling_cursors after a successful poll with events', async () => { + mockGetEvents.mockResolvedValue({ + events: [makeEvent('E1', 1000)], + cursor: 'cursor-after-1000', + }); + + const subscriber = new EventSubscriber(testConfig, dedup); + await (subscriber as any).checkForEvents(); + + const record = await dedup.getLastCursor(CONTRACT_ADDRESS); + expect(record).not.toBeNull(); + expect(record!.cursor).toBe('cursor-after-1000'); + expect(record!.ledgerNumber).toBe(1000); + }); + + it('advances the checkpoint ledger after each successful poll', async () => { + mockGetEvents + .mockResolvedValueOnce({ + events: [makeEvent('E1', 1000)], + cursor: 'cursor-1000', + }) + .mockResolvedValueOnce({ + events: [makeEvent('E2', 2000)], + cursor: 'cursor-2000', + }); + + const subscriber = new EventSubscriber(testConfig, dedup); + await (subscriber as any).checkForEvents(); + + let record = await dedup.getLastCursor(CONTRACT_ADDRESS); + expect(record!.cursor).toBe('cursor-1000'); + expect(record!.ledgerNumber).toBe(1000); + + await (subscriber as any).checkForEvents(); + + record = await dedup.getLastCursor(CONTRACT_ADDRESS); + expect(record!.cursor).toBe('cursor-2000'); + expect(record!.ledgerNumber).toBe(2000); + }); + + it('does not create a polling_cursors row when the RPC returns no cursor', async () => { + mockGetEvents.mockResolvedValue({ + events: [makeEvent('E1', 1000)], + cursor: undefined, + }); + + const subscriber = new EventSubscriber(testConfig, dedup); + await (subscriber as any).checkForEvents(); + + const record = await dedup.getLastCursor(CONTRACT_ADDRESS); + expect(record).toBeNull(); + }); + + it('does not advance the checkpoint when the RPC call throws', async () => { + // First poll succeeds and writes a checkpoint + mockGetEvents.mockResolvedValueOnce({ + events: [makeEvent('E1', 1000)], + cursor: 'cursor-1000', + }); + + const subscriber = new EventSubscriber(testConfig, dedup); + await (subscriber as any).checkForEvents(); + + // Second poll fails — checkpoint must not move + mockGetEvents.mockRejectedValueOnce(new Error('RPC timeout')); + + await expect((subscriber as any).checkForEvents()).rejects.toThrow( + 'Failed to fetch events for all' + ); + + const record = await dedup.getLastCursor(CONTRACT_ADDRESS); + expect(record!.cursor).toBe('cursor-1000'); + expect(record!.ledgerNumber).toBe(1000); + }); + + it('persists checkpoints for each contract independently', async () => { + const contract2: ContractConfig = { + address: 'CBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB', + events: ['*'], + }; + const multiConfig: Config = { + ...testConfig, + contractAddresses: [contractConfig, contract2], + }; + + mockGetEvents + .mockResolvedValueOnce({ events: [makeEvent('E1', 500)], cursor: 'cursor-C1' }) + .mockResolvedValueOnce({ events: [makeEvent('E2', 800)], cursor: 'cursor-C2' }); + + const subscriber = new EventSubscriber(multiConfig, dedup); + await (subscriber as any).checkForEvents(); + + const r1 = await dedup.getLastCursor(CONTRACT_ADDRESS); + const r2 = await dedup.getLastCursor(contract2.address); + + expect(r1!.cursor).toBe('cursor-C1'); + expect(r1!.ledgerNumber).toBe(500); + expect(r2!.cursor).toBe('cursor-C2'); + expect(r2!.ledgerNumber).toBe(800); + }); +}); + +// --------------------------------------------------------------------------- +// Checkpoint restoration at startup +// --------------------------------------------------------------------------- + +describe('EventSubscriber — checkpoint restoration at startup', () => { + let db: Database; + let dedup: EventDeduplicationService; + + beforeEach(async () => { + jest.clearAllMocks(); + ({ db, dedup } = await setupDb()); + mockGetEvents.mockResolvedValue({ events: [], cursor: undefined }); + }); + + afterEach(async () => { + await db.close(); + }); + + it('restores a persisted cursor into lastCursors before the first poll', async () => { + // Seed polling_cursors as if a previous run had reached ledger 5000 + await dedup.updatePollingCursor(CONTRACT_ADDRESS, 'cursor-from-prev-run', 5000); + + const subscriber = new EventSubscriber(testConfig, dedup); + await (subscriber as any).restoreCheckpoints(); + + // The in-memory map must now hold the persisted cursor + const lastCursors: Map = (subscriber as any).lastCursors; + expect(lastCursors.get(CONTRACT_ADDRESS)).toBe('cursor-from-prev-run'); + }); + + it('uses the restored cursor in the first RPC call after start()', async () => { + await dedup.updatePollingCursor(CONTRACT_ADDRESS, 'checkpoint-cursor', 5000); + + // start() calls restoreCheckpoints() then poll() → checkForEvents() → getContractEvents() + // poll() runs in the background so we yield after start() resolves + let pollFired = false; + mockGetEvents.mockImplementationOnce(async () => { + pollFired = true; + subscriber.stop(); + return { events: [], cursor: undefined }; + }); + + const subscriber = new EventSubscriber(testConfig, dedup); + await subscriber.start(); + + await new Promise((r) => setTimeout(r, 10)); + + expect(pollFired).toBe(true); + expect(mockGetEvents).toHaveBeenCalledTimes(1); + const firstCall = mockGetEvents.mock.calls[0][0] as StellarSDK.rpc.Api.GetEventsRequest; + expect((firstCall as any).cursor).toBe('checkpoint-cursor'); + expect((firstCall as any).startLedger).toBeUndefined(); + }); + + it('starts from the beginning when no checkpoint exists for a contract', async () => { + // No row in polling_cursors + + let pollFired = false; + mockGetEvents.mockImplementationOnce(async () => { + pollFired = true; + subscriber.stop(); + return { events: [], cursor: undefined }; + }); + + const subscriber = new EventSubscriber(testConfig, dedup); + await subscriber.start(); + + // Yield the event loop so the background poll() iteration can run + await new Promise((r) => setTimeout(r, 10)); + + expect(pollFired).toBe(true); + const firstCall = mockGetEvents.mock.calls[0][0] as StellarSDK.rpc.Api.GetEventsRequest; + expect((firstCall as any).cursor).toBeUndefined(); + expect((firstCall as any).startLedger).toBeDefined(); + }); + + it('restores checkpoints for multiple contracts independently', async () => { + const contract2: ContractConfig = { + address: 'CBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB', + events: ['*'], + }; + const multiConfig: Config = { + ...testConfig, + contractAddresses: [contractConfig, contract2], + }; + + await dedup.updatePollingCursor(CONTRACT_ADDRESS, 'cursor-C1', 1000); + await dedup.updatePollingCursor(contract2.address, 'cursor-C2', 2000); + + const subscriber = new EventSubscriber(multiConfig, dedup); + await (subscriber as any).restoreCheckpoints(); + + const lastCursors: Map = (subscriber as any).lastCursors; + expect(lastCursors.get(CONTRACT_ADDRESS)).toBe('cursor-C1'); + expect(lastCursors.get(contract2.address)).toBe('cursor-C2'); + }); + + it('restores a checkpoint for one contract but not for another that has none', async () => { + const contract2: ContractConfig = { + address: 'CBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB', + events: ['*'], + }; + const multiConfig: Config = { + ...testConfig, + contractAddresses: [contractConfig, contract2], + }; + + // Only C1 has a persisted cursor + await dedup.updatePollingCursor(CONTRACT_ADDRESS, 'cursor-C1', 1000); + + const subscriber = new EventSubscriber(multiConfig, dedup); + await (subscriber as any).restoreCheckpoints(); + + const lastCursors: Map = (subscriber as any).lastCursors; + expect(lastCursors.get(CONTRACT_ADDRESS)).toBe('cursor-C1'); + expect(lastCursors.has(contract2.address)).toBe(false); + }); + + it('is a no-op and does not throw when deduplicationService is null', async () => { + // No deduplication service — should complete silently + const subscriber = new EventSubscriber(testConfig, undefined); + await expect((subscriber as any).restoreCheckpoints()).resolves.toBeUndefined(); + + const lastCursors: Map = (subscriber as any).lastCursors; + expect(lastCursors.size).toBe(0); + }); + + it('logs each restored checkpoint at info level', async () => { + const mockLogger = logger as jest.Mocked; + await dedup.updatePollingCursor(CONTRACT_ADDRESS, 'cursor-logged', 7777); + + const subscriber = new EventSubscriber(testConfig, dedup); + await (subscriber as any).restoreCheckpoints(); + + expect(mockLogger.info).toHaveBeenCalledWith( + 'Checkpoint restored', + expect.objectContaining({ + contractAddress: CONTRACT_ADDRESS, + cursor: 'cursor-logged', + ledgerNumber: 7777, + }) + ); + }); + + it('logs the restoration summary with counts', async () => { + const mockLogger = logger as jest.Mocked; + await dedup.updatePollingCursor(CONTRACT_ADDRESS, 'cursor-summary', 1000); + + const subscriber = new EventSubscriber(testConfig, dedup); + await (subscriber as any).restoreCheckpoints(); + + expect(mockLogger.info).toHaveBeenCalledWith( + 'Checkpoint restore complete', + expect.objectContaining({ + contractsConfigured: 1, + contractsRestored: 1, + }) + ); + }); + + it('continues restoring remaining contracts when one getLastCursor call throws', async () => { + const mockLogger = logger as jest.Mocked; + const contract2: ContractConfig = { + address: 'CBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB', + events: ['*'], + }; + const multiConfig: Config = { + ...testConfig, + contractAddresses: [contractConfig, contract2], + }; + + // C1 throws, C2 has a valid cursor + await dedup.updatePollingCursor(contract2.address, 'cursor-C2', 2000); + jest.spyOn(dedup, 'getLastCursor').mockImplementationOnce(() => { + throw new Error('DB read failure'); + }); + + const subscriber = new EventSubscriber(multiConfig, dedup); + await expect((subscriber as any).restoreCheckpoints()).resolves.toBeUndefined(); + + // C1 failure logged as a warning + expect(mockLogger.warn).toHaveBeenCalledWith( + 'Failed to restore checkpoint for contract; starting from beginning', + expect.objectContaining({ contractAddress: CONTRACT_ADDRESS }) + ); + + // C2 still restored despite C1 failure + const lastCursors: Map = (subscriber as any).lastCursors; + expect(lastCursors.get(contract2.address)).toBe('cursor-C2'); + }); +}); + +// --------------------------------------------------------------------------- +// End-to-end: write then resume +// --------------------------------------------------------------------------- + +describe('EventSubscriber — write checkpoint then resume from it', () => { + let db: Database; + let dedup: EventDeduplicationService; + + beforeEach(async () => { + jest.clearAllMocks(); + ({ db, dedup } = await setupDb()); + mockGetEvents.mockResolvedValue({ events: [], cursor: undefined }); + }); + + afterEach(async () => { + await db.close(); + }); + + it('resumes from the last checkpoint after a simulated restart', async () => { + // --- First "run" --- + mockGetEvents.mockResolvedValue({ + events: [makeEvent('E1', 3000), makeEvent('E2', 3001)], + cursor: 'cursor-at-3001', + }); + + const firstSubscriber = new EventSubscriber(testConfig, dedup); + await (firstSubscriber as any).checkForEvents(); + + // Checkpoint must be persisted + const savedRecord = await dedup.getLastCursor(CONTRACT_ADDRESS); + expect(savedRecord!.cursor).toBe('cursor-at-3001'); + expect(savedRecord!.ledgerNumber).toBe(3001); + + // --- Simulate restart: new subscriber instance, same DB --- + jest.clearAllMocks(); + mockGetEvents.mockResolvedValue({ events: [], cursor: undefined }); + + const secondSubscriber = new EventSubscriber(testConfig, dedup); + + // Before restoreCheckpoints the map is empty + expect((secondSubscriber as any).lastCursors.size).toBe(0); + + await (secondSubscriber as any).restoreCheckpoints(); + + // After restore it holds the persisted cursor + expect((secondSubscriber as any).lastCursors.get(CONTRACT_ADDRESS)).toBe('cursor-at-3001'); + + // The first poll after restore must use the restored cursor + await (secondSubscriber as any).checkForEvents(); + + const rpcRequest = mockGetEvents.mock.calls[0][0] as StellarSDK.rpc.Api.GetEventsRequest; + expect((rpcRequest as any).cursor).toBe('cursor-at-3001'); + expect((rpcRequest as any).startLedger).toBeUndefined(); + }); + + it('checkpoint is only advanced on success — failed polls leave it at the last good position', async () => { + // Successful poll writes checkpoint at ledger 4000 + mockGetEvents.mockResolvedValueOnce({ + events: [makeEvent('E1', 4000)], + cursor: 'cursor-4000', + }); + + const subscriber = new EventSubscriber(testConfig, dedup); + await (subscriber as any).checkForEvents(); + + expect((await dedup.getLastCursor(CONTRACT_ADDRESS))!.cursor).toBe('cursor-4000'); + + // Failed poll — checkpoint must stay at 4000 + mockGetEvents.mockRejectedValueOnce(new Error('network error')); + + await expect((subscriber as any).checkForEvents()).rejects.toThrow(); + + const record = await dedup.getLastCursor(CONTRACT_ADDRESS); + expect(record!.cursor).toBe('cursor-4000'); + expect(record!.ledgerNumber).toBe(4000); + + // Third poll succeeds — checkpoint advances + mockGetEvents.mockResolvedValueOnce({ + events: [makeEvent('E2', 5000)], + cursor: 'cursor-5000', + }); + + await (subscriber as any).checkForEvents(); + + const finalRecord = await dedup.getLastCursor(CONTRACT_ADDRESS); + expect(finalRecord!.cursor).toBe('cursor-5000'); + expect(finalRecord!.ledgerNumber).toBe(5000); + }); +}); diff --git a/listener/src/services/event-subscriber.ts b/listener/src/services/event-subscriber.ts index 82dcda7c..0f0cc79b 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; @@ -65,6 +66,7 @@ export class EventSubscriber { this.isRunning = true; logger.info('Starting event subscriber service'); + await this.restoreCheckpoints(); this.eventQueue?.start(); this.retryQueue?.start(); this.poll(); @@ -77,6 +79,51 @@ export class EventSubscriber { logger.info('Stopping event subscriber service'); } + /** + * Restore the in-memory cursor map from persisted checkpoints. + * + * Called once at the start of `start()` before the poll loop begins. + * For each configured contract address, if a row exists in + * `polling_cursors`, the stored cursor string is loaded into + * `this.lastCursors` so the first poll resumes from the last known + * position rather than replaying from the beginning. + * + * Failures are logged and swallowed per-contract so a single bad DB + * row never prevents the subscriber from starting. + */ + private async restoreCheckpoints(): Promise { + if (!this.deduplicationService) { + return; + } + + let restored = 0; + + for (const contractConfig of this.config.contractAddresses) { + try { + const record = await this.deduplicationService.getLastCursor(contractConfig.address); + if (record) { + this.lastCursors.set(contractConfig.address, record.cursor); + restored++; + logger.info('Checkpoint restored', { + contractAddress: contractConfig.address, + cursor: record.cursor, + ledgerNumber: record.ledgerNumber, + }); + } + } catch (error) { + logger.warn('Failed to restore checkpoint for contract; starting from beginning', { + contractAddress: contractConfig.address, + error, + }); + } + } + + logger.info('Checkpoint restore complete', { + contractsConfigured: this.config.contractAddresses.length, + contractsRestored: restored, + }); + } + private async poll(): Promise { while (this.isRunning) { const requestId = generateRequestId(); @@ -170,9 +217,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', { @@ -248,7 +292,6 @@ export class EventSubscriber { contractAddress: contractConfig.address, eventId: event.id, eventName, - receivedAt: event.receivedAt, currentTime: Date.now(), reason: 'expired', }); @@ -339,27 +382,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; 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).