From d5d73dd7608e5472ef3631023b90a244ab9de996 Mon Sep 17 00:00:00 2001 From: eliassmthl-collab Date: Tue, 25 Aug 2026 14:53:48 +0000 Subject: [PATCH] feat(indexer): exponential backoff retry on event_type_filter RPC calls Introduce src/indexer/event_type_filter.ts as a dedicated module that wraps Soroban RPC getEvents with exponential-backoff retry logic for transient connection failures. Changes: - Add EVENT_TYPES constant (single source of truth, removed from poller) - Add isConnectionError() to classify retryable network errors (ECONNREFUSED, ETIMEDOUT, ECONNRESET, fetch failed, socket hang up, etc.) - Add nextBackoff() helper: doubles delay each attempt, caps at 300 000 ms - Add fetchEventsWithRetry(): retries up to MAX_ATTEMPTS (5) on connection errors with 1 s initial backoff doubling each failure; non-connection errors propagate immediately without retrying; sleep is injectable for deterministic testing - Update poller.ts to delegate getEvents to fetchEventsWithRetry so every poll cycle recovers gracefully from RPC connection dropouts - Add __tests__/event-type-filter.test.ts with 53 tests covering: - EVENT_TYPES content - isConnectionError classification (14 error patterns) - nextBackoff doubling and cap behaviour - Retry call counts match maxAttempts exactly on total failure - Sleep delays grow strictly: [1000, 2000, 4000, 8000] for 5 attempts - Immediate abort on non-connection errors - Mixed connection/non-connection error sequences All 730 existing tests continue to pass (42 suites). Closes #276 --- __tests__/event-type-filter.test.ts | 667 ++++++++++++++++++++++++++++ src/indexer/event_type_filter.ts | 202 +++++++++ src/indexer/poller.ts | 24 +- 3 files changed, 872 insertions(+), 21 deletions(-) create mode 100644 __tests__/event-type-filter.test.ts create mode 100644 src/indexer/event_type_filter.ts diff --git a/__tests__/event-type-filter.test.ts b/__tests__/event-type-filter.test.ts new file mode 100644 index 0000000..e86064d --- /dev/null +++ b/__tests__/event-type-filter.test.ts @@ -0,0 +1,667 @@ +import { jest } from "@jest/globals"; +import { + EVENT_TYPES, + DEFAULT_BACKOFF_MS, + MAX_BACKOFF_MS, + BACKOFF_MULTIPLIER, + MAX_ATTEMPTS, + isConnectionError, + nextBackoff, + fetchEventsWithRetry, + type GetEventsParams, + type RpcGetEventsResult, +} from "../src/indexer/event_type_filter.js"; + +// --------------------------------------------------------------------------- +// Helpers +// --------------------------------------------------------------------------- + +/** Instant sleep (no real delays in tests). */ +const noopSleep = (): Promise => Promise.resolve(); + +/** Build a minimal success result matching RpcGetEventsResult. */ +function makeSuccessResult( + count = 0 +): RpcGetEventsResult { + return { + events: Array.from({ length: count }, (_, i) => ({ + contractId: { contractId: () => `CONTRACT-${i}` }, + topic: [`event_${i}`], + value: { data: i }, + ledger: 100 + i, + ledgerClosedAt: new Date().toISOString(), + })), + }; +} + +/** Build a mock `Server.getEvents` that returns `result`. */ +function makeSuccessServer(result: RpcGetEventsResult = makeSuccessResult()) { + return { + getEvents: jest.fn().mockResolvedValue(result), + }; +} + +/** Build a mock `Server.getEvents` that always throws `err`. */ +function makeFailingServer(err: Error) { + return { + getEvents: jest.fn().mockRejectedValue(err), + }; +} + +/** + * Build a mock server that fails `failCount` times with `err` then + * succeeds with `successResult`. + */ +function makePartialFailServer( + failCount: number, + err: Error, + successResult: RpcGetEventsResult = makeSuccessResult(1) +) { + let calls = 0; + return { + getEvents: jest.fn().mockImplementation(() => { + calls += 1; + if (calls <= failCount) return Promise.reject(err); + return Promise.resolve(successResult); + }), + }; +} + +const BASE_PARAMS: GetEventsParams = { + startLedger: 42, + contractIds: ["CONTRACT-ALPHA"], + limit: 100, +}; + +// --------------------------------------------------------------------------- +// EVENT_TYPES +// --------------------------------------------------------------------------- + +describe("EVENT_TYPES", () => { + it("contains the expected event type strings", () => { + expect(EVENT_TYPES).toContain("initialized"); + expect(EVENT_TYPES).toContain("funded"); + expect(EVENT_TYPES).toContain("delivered"); + expect(EVENT_TYPES).toContain("approved"); + expect(EVENT_TYPES).toContain("dispute_raised"); + expect(EVENT_TYPES).toContain("dispute_resolved"); + expect(EVENT_TYPES).toContain("partial_release"); + expect(EVENT_TYPES).toContain("auto_release_claimed"); + expect(EVENT_TYPES).toContain("token_whitelisted"); + expect(EVENT_TYPES).toContain("token_removed"); + }); + + it("contains exactly 10 event types", () => { + expect(EVENT_TYPES.length).toBe(10); + }); +}); + +// --------------------------------------------------------------------------- +// isConnectionError +// --------------------------------------------------------------------------- + +describe("isConnectionError", () => { + const connectionErrors = [ + "connect ECONNREFUSED 127.0.0.1:26657", + "read ECONNRESET", + "connection timeout after 30000ms", + "ETIMEDOUT", + "getaddrinfo ENOTFOUND soroban-testnet.stellar.org", + "EHOSTUNREACH", + "Network Error occurred", + "socket hang up", + "fetch failed", + "request timeout", + "getaddrinfo ENOTFOUND host", + "connect ECONNREFUSED", + ]; + + const nonConnectionErrors = [ + "start is after end", + "Contract not found", + "Invalid request body", + "Internal Server Error", + "400 Bad Request", + "rate limit exceeded", + "Unknown event type", + ]; + + connectionErrors.forEach((msg) => { + it(`classifies "${msg}" as a connection error`, () => { + expect(isConnectionError(new Error(msg))).toBe(true); + }); + }); + + nonConnectionErrors.forEach((msg) => { + it(`does NOT classify "${msg}" as a connection error`, () => { + expect(isConnectionError(new Error(msg))).toBe(false); + }); + }); + + it("returns false for non-Error values", () => { + expect(isConnectionError("string error")).toBe(false); + expect(isConnectionError(42)).toBe(false); + expect(isConnectionError(null)).toBe(false); + expect(isConnectionError(undefined)).toBe(false); + expect(isConnectionError({ message: "ECONNREFUSED" })).toBe(false); + }); +}); + +// --------------------------------------------------------------------------- +// nextBackoff +// --------------------------------------------------------------------------- + +describe("nextBackoff", () => { + it("doubles the current backoff", () => { + expect(nextBackoff(1_000)).toBe(2_000); + expect(nextBackoff(2_000)).toBe(4_000); + expect(nextBackoff(16_000)).toBe(32_000); + }); + + it("never exceeds MAX_BACKOFF_MS", () => { + expect(nextBackoff(MAX_BACKOFF_MS)).toBe(MAX_BACKOFF_MS); + expect(nextBackoff(MAX_BACKOFF_MS * 10)).toBe(MAX_BACKOFF_MS); + }); + + it("applies BACKOFF_MULTIPLIER correctly", () => { + expect(nextBackoff(1_000)).toBe(1_000 * BACKOFF_MULTIPLIER); + }); + + it("respects the DEFAULT_BACKOFF_MS starting point", () => { + expect(nextBackoff(DEFAULT_BACKOFF_MS)).toBe( + DEFAULT_BACKOFF_MS * BACKOFF_MULTIPLIER + ); + }); +}); + +// --------------------------------------------------------------------------- +// fetchEventsWithRetry – success path +// --------------------------------------------------------------------------- + +describe("fetchEventsWithRetry – success path", () => { + it("returns the result on the first attempt", async () => { + const expected = makeSuccessResult(3); + const server = makeSuccessServer(expected); + + const result = await fetchEventsWithRetry(server, BASE_PARAMS, { + sleep: noopSleep, + }); + + expect(result).toBe(expected); + expect(server.getEvents).toHaveBeenCalledTimes(1); + }); + + it("passes EVENT_TYPES as the topic filter", async () => { + const server = makeSuccessServer(); + + await fetchEventsWithRetry(server, BASE_PARAMS, { sleep: noopSleep }); + + const callArgs = server.getEvents.mock.calls[0][0] as { + filters: Array<{ topics: string[][] }>; + }; + expect(callArgs.filters[0].topics[0]).toEqual([...EVENT_TYPES]); + }); + + it("passes contractIds and startLedger from params", async () => { + const server = makeSuccessServer(); + const params: GetEventsParams = { + startLedger: 999, + contractIds: ["C1", "C2"], + limit: 50, + }; + + await fetchEventsWithRetry(server, params, { sleep: noopSleep }); + + const callArgs = server.getEvents.mock.calls[0][0] as { + startLedger: number; + limit: number; + filters: Array<{ contractIds: string[] }>; + }; + expect(callArgs.startLedger).toBe(999); + expect(callArgs.limit).toBe(50); + expect(callArgs.filters[0].contractIds).toEqual(["C1", "C2"]); + }); + + it("defaults limit to 100 when not specified", async () => { + const server = makeSuccessServer(); + + await fetchEventsWithRetry( + server, + { startLedger: 1, contractIds: ["C"] }, + { sleep: noopSleep } + ); + + const callArgs = server.getEvents.mock.calls[0][0] as { limit: number }; + expect(callArgs.limit).toBe(100); + }); + + it("returns empty events array when RPC returns no events", async () => { + const server = makeSuccessServer(makeSuccessResult(0)); + + const result = await fetchEventsWithRetry(server, BASE_PARAMS, { + sleep: noopSleep, + }); + + expect(result.events).toHaveLength(0); + }); +}); + +// --------------------------------------------------------------------------- +// fetchEventsWithRetry – non-connection errors propagate immediately +// --------------------------------------------------------------------------- + +describe("fetchEventsWithRetry – non-connection error propagation", () => { + it("throws immediately on a non-connection error without retrying", async () => { + const err = new Error("start is after end"); + const server = makeFailingServer(err); + + await expect( + fetchEventsWithRetry(server, BASE_PARAMS, { sleep: noopSleep }) + ).rejects.toThrow("start is after end"); + + // Only one call – no retries for non-connection errors + expect(server.getEvents).toHaveBeenCalledTimes(1); + }); + + it("propagates the original error object", async () => { + const original = new Error("Invalid contract ID"); + const server = makeFailingServer(original); + + await expect( + fetchEventsWithRetry(server, BASE_PARAMS, { sleep: noopSleep }) + ).rejects.toBe(original); + }); +}); + +// --------------------------------------------------------------------------- +// fetchEventsWithRetry – retry on connection errors +// --------------------------------------------------------------------------- + +describe("fetchEventsWithRetry – retry on connection errors", () => { + it("retries on ECONNREFUSED and succeeds on second attempt", async () => { + const err = new Error("connect ECONNREFUSED 127.0.0.1:26657"); + const expected = makeSuccessResult(2); + const server = makePartialFailServer(1, err, expected); + + const result = await fetchEventsWithRetry(server, BASE_PARAMS, { + sleep: noopSleep, + }); + + expect(result).toBe(expected); + expect(server.getEvents).toHaveBeenCalledTimes(2); + }); + + it("retries up to maxAttempts and then throws the last error", async () => { + const err = new Error("connect ECONNREFUSED 127.0.0.1:26657"); + const server = makeFailingServer(err); + + await expect( + fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts: 3, + sleep: noopSleep, + }) + ).rejects.toThrow("connect ECONNREFUSED"); + + expect(server.getEvents).toHaveBeenCalledTimes(3); + }); + + it("does not exceed MAX_ATTEMPTS by default", async () => { + const err = new Error("ETIMEDOUT"); + const server = makeFailingServer(err); + + await expect( + fetchEventsWithRetry(server, BASE_PARAMS, { sleep: noopSleep }) + ).rejects.toThrow(); + + expect(server.getEvents).toHaveBeenCalledTimes(MAX_ATTEMPTS); + }); + + it("succeeds on the last possible attempt", async () => { + const err = new Error("ETIMEDOUT"); + const expected = makeSuccessResult(1); + const server = makePartialFailServer(MAX_ATTEMPTS - 1, err, expected); + + const result = await fetchEventsWithRetry(server, BASE_PARAMS, { + sleep: noopSleep, + }); + + expect(result).toBe(expected); + expect(server.getEvents).toHaveBeenCalledTimes(MAX_ATTEMPTS); + }); + + it("retries on ECONNRESET", async () => { + const server = makePartialFailServer( + 2, + new Error("read ECONNRESET"), + makeSuccessResult(1) + ); + + await fetchEventsWithRetry(server, BASE_PARAMS, { sleep: noopSleep }); + + expect(server.getEvents).toHaveBeenCalledTimes(3); + }); + + it("retries on fetch failed", async () => { + const server = makePartialFailServer( + 1, + new Error("fetch failed"), + makeSuccessResult(1) + ); + + await fetchEventsWithRetry(server, BASE_PARAMS, { sleep: noopSleep }); + + expect(server.getEvents).toHaveBeenCalledTimes(2); + }); +}); + +// --------------------------------------------------------------------------- +// fetchEventsWithRetry – exponential backoff timing +// --------------------------------------------------------------------------- + +describe("fetchEventsWithRetry – exponential backoff timing", () => { + it("waits DEFAULT_BACKOFF_MS before the first retry", async () => { + const sleepDelays: number[] = []; + const err = new Error("connect ECONNREFUSED"); + const server = makePartialFailServer(1, err, makeSuccessResult()); + + await fetchEventsWithRetry(server, BASE_PARAMS, { + sleep: (ms) => { + sleepDelays.push(ms); + return Promise.resolve(); + }, + }); + + expect(sleepDelays).toHaveLength(1); + expect(sleepDelays[0]).toBe(DEFAULT_BACKOFF_MS); + }); + + it("doubles the backoff on each consecutive failure", async () => { + const sleepDelays: number[] = []; + const err = new Error("connect ECONNREFUSED"); + + // Fail 4 times, then succeed on attempt 5 + const server = makePartialFailServer(4, err, makeSuccessResult()); + + await fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts: 5, + initialBackoffMs: 1_000, + sleep: (ms) => { + sleepDelays.push(ms); + return Promise.resolve(); + }, + }); + + // After attempt 1 → sleep 1000, attempt 2 → sleep 2000, attempt 3 → sleep 4000, attempt 4 → sleep 8000 + expect(sleepDelays).toHaveLength(4); + expect(sleepDelays[0]).toBe(1_000); + expect(sleepDelays[1]).toBe(2_000); + expect(sleepDelays[2]).toBe(4_000); + expect(sleepDelays[3]).toBe(8_000); + }); + + it("backoff increases by BACKOFF_MULTIPLIER each step", async () => { + const sleepDelays: number[] = []; + const err = new Error("ETIMEDOUT"); + + // Always fail so we collect all backoffs up to maxAttempts + const server = makeFailingServer(err); + + await fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts: 5, + initialBackoffMs: 500, + sleep: (ms) => { + sleepDelays.push(ms); + return Promise.resolve(); + }, + }).catch(() => { + // Expected to throw after exhausting retries + }); + + // 4 sleeps for attempts 1–4 (attempt 5 is the last, no sleep after it) + expect(sleepDelays).toHaveLength(4); + for (let i = 1; i < sleepDelays.length; i++) { + expect(sleepDelays[i]).toBe( + Math.min(sleepDelays[i - 1] * BACKOFF_MULTIPLIER, MAX_BACKOFF_MS) + ); + } + }); + + it("caps backoff at MAX_BACKOFF_MS when growth would exceed it", async () => { + const sleepDelays: number[] = []; + const err = new Error("ETIMEDOUT"); + const server = makeFailingServer(err); + + // Start near the cap so it hits it quickly + const startingBackoff = MAX_BACKOFF_MS / 2; + + await fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts: 4, + initialBackoffMs: startingBackoff, + sleep: (ms) => { + sleepDelays.push(ms); + return Promise.resolve(); + }, + }).catch(() => {}); + + // sleepDelays[0] = MAX_BACKOFF_MS / 2 (first delay = initialBackoffMs) + // sleepDelays[1] = MAX_BACKOFF_MS (doubled, but capped) + // sleepDelays[2] = MAX_BACKOFF_MS (already at cap) + expect(sleepDelays[0]).toBe(startingBackoff); + expect(sleepDelays[1]).toBe(MAX_BACKOFF_MS); + expect(sleepDelays[2]).toBe(MAX_BACKOFF_MS); + }); + + it("backoff sequence grows strictly until the cap, then stays flat", async () => { + const sleepDelays: number[] = []; + const err = new Error("connect ECONNREFUSED"); + const server = makeFailingServer(err); + + await fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts: 6, + initialBackoffMs: DEFAULT_BACKOFF_MS, + sleep: (ms) => { + sleepDelays.push(ms); + return Promise.resolve(); + }, + }).catch(() => {}); + + // Verify strictly increasing or plateau at MAX_BACKOFF_MS + for (let i = 1; i < sleepDelays.length; i++) { + expect(sleepDelays[i]).toBeGreaterThanOrEqual(sleepDelays[i - 1]); + expect(sleepDelays[i]).toBeLessThanOrEqual(MAX_BACKOFF_MS); + } + }); +}); + +// --------------------------------------------------------------------------- +// fetchEventsWithRetry – custom maxAttempts +// --------------------------------------------------------------------------- + +describe("fetchEventsWithRetry – maxAttempts option", () => { + it("respects maxAttempts = 1 (no retries)", async () => { + const err = new Error("ETIMEDOUT"); + const server = makeFailingServer(err); + + await expect( + fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts: 1, + sleep: noopSleep, + }) + ).rejects.toThrow(); + + expect(server.getEvents).toHaveBeenCalledTimes(1); + }); + + it("respects maxAttempts = 2 (exactly one retry)", async () => { + const err = new Error("ETIMEDOUT"); + const server = makeFailingServer(err); + + await fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts: 2, + sleep: noopSleep, + }).catch(() => {}); + + expect(server.getEvents).toHaveBeenCalledTimes(2); + }); + + it("respects a custom high maxAttempts value", async () => { + const err = new Error("connect ECONNREFUSED"); + const server = makePartialFailServer(7, err, makeSuccessResult(1)); + + const result = await fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts: 10, + sleep: noopSleep, + }); + + expect(result.events).toHaveLength(1); + expect(server.getEvents).toHaveBeenCalledTimes(8); // 7 failures + 1 success + }); +}); + +// --------------------------------------------------------------------------- +// fetchEventsWithRetry – retry frequency increases test (acceptance criteria) +// --------------------------------------------------------------------------- + +describe("fetchEventsWithRetry – retry frequency increases up to max attempts", () => { + it("verifies each successive wait is longer than the previous (growing backoff)", async () => { + const delays: number[] = []; + const err = new Error("connect ECONNREFUSED 127.0.0.1:26657"); + const server = makeFailingServer(err); + + await fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts: 5, + initialBackoffMs: 100, + sleep: (ms) => { + delays.push(ms); + return Promise.resolve(); + }, + }).catch(() => {}); + + // Must have exactly 4 delays (before attempts 2, 3, 4, 5) + expect(delays).toHaveLength(4); + + // Each delay must be strictly larger than the previous (until cap) + for (let i = 1; i < delays.length; i++) { + expect(delays[i]).toBeGreaterThan(delays[i - 1]); + } + }); + + it("total sleep time increases exponentially across failed attempts", async () => { + const delays: number[] = []; + const err = new Error("ETIMEDOUT"); + const server = makeFailingServer(err); + + await fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts: 5, + initialBackoffMs: 1_000, + sleep: (ms) => { + delays.push(ms); + return Promise.resolve(); + }, + }).catch(() => {}); + + // Expected: 1000, 2000, 4000, 8000 + expect(delays[0]).toBe(1_000); + expect(delays[1]).toBe(2_000); + expect(delays[2]).toBe(4_000); + expect(delays[3]).toBe(8_000); + + const totalSleep = delays.reduce((a, b) => a + b, 0); + expect(totalSleep).toBe(15_000); + }); + + it("call count matches maxAttempts exactly on total failure", async () => { + const err = new Error("connect ECONNREFUSED 127.0.0.1:26657"); + const callTimes: number[] = []; + let fakeNow = 0; + + const server = { + getEvents: jest.fn().mockImplementation(() => { + callTimes.push(fakeNow); + fakeNow += 1; + return Promise.reject(err); + }), + }; + + await fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts: MAX_ATTEMPTS, + sleep: (ms) => { + fakeNow += ms; + return Promise.resolve(); + }, + }).catch(() => {}); + + expect(callTimes).toHaveLength(MAX_ATTEMPTS); + }); + + it("stops retrying and throws after maxAttempts connection errors", async () => { + let callCount = 0; + const err = new Error("socket hang up"); + + const server = { + getEvents: jest.fn().mockImplementation(() => { + callCount++; + return Promise.reject(err); + }), + }; + + const maxAttempts = 4; + + await expect( + fetchEventsWithRetry(server, BASE_PARAMS, { + maxAttempts, + sleep: noopSleep, + }) + ).rejects.toThrow("socket hang up"); + + expect(callCount).toBe(maxAttempts); + }); +}); + +// --------------------------------------------------------------------------- +// fetchEventsWithRetry – mixed error scenarios +// --------------------------------------------------------------------------- + +describe("fetchEventsWithRetry – mixed error scenarios", () => { + it("stops retrying immediately when a non-connection error follows connection errors", async () => { + let callCount = 0; + const connectionErr = new Error("ETIMEDOUT"); + const nonConnectionErr = new Error("Contract not found"); + + const server = { + getEvents: jest.fn().mockImplementation(() => { + callCount++; + if (callCount <= 2) return Promise.reject(connectionErr); + return Promise.reject(nonConnectionErr); + }), + }; + + await expect( + fetchEventsWithRetry(server, BASE_PARAMS, { sleep: noopSleep }) + ).rejects.toThrow("Contract not found"); + + // 2 connection retries + 1 non-connection call = 3 total + expect(callCount).toBe(3); + }); + + it("handles alternating errors by only stopping on non-connection errors", async () => { + let callCount = 0; + const connErr = new Error("fetch failed"); + const badErr = new Error("400 Bad Request"); + + const server = { + getEvents: jest.fn().mockImplementation(() => { + callCount++; + // Fail with connection error first, then bad request + if (callCount === 1) return Promise.reject(connErr); + return Promise.reject(badErr); + }), + }; + + await expect( + fetchEventsWithRetry(server, BASE_PARAMS, { sleep: noopSleep }) + ).rejects.toThrow("400 Bad Request"); + + expect(callCount).toBe(2); + }); +}); diff --git a/src/indexer/event_type_filter.ts b/src/indexer/event_type_filter.ts new file mode 100644 index 0000000..651b00d --- /dev/null +++ b/src/indexer/event_type_filter.ts @@ -0,0 +1,202 @@ +import type { Server } from "@stellar/stellar-sdk/rpc"; +import logger from "../utils/logger.js"; + +/** + * The set of Soroban contract event types the indexer cares about. + * Only events whose first topic matches one of these strings will be fetched. + */ +export const EVENT_TYPES = [ + "initialized", + "funded", + "delivered", + "approved", + "dispute_raised", + "dispute_resolved", + "partial_release", + "auto_release_claimed", + "token_whitelisted", + "token_removed", +] as const; + +export type EventType = (typeof EVENT_TYPES)[number]; + +/** Starting delay (ms) applied after the first failed attempt. */ +export const DEFAULT_BACKOFF_MS = 1_000; + +/** Hard ceiling so exponential growth never exceeds five minutes. */ +export const MAX_BACKOFF_MS = 300_000; + +/** Multiplier applied to the current backoff on every consecutive failure. */ +export const BACKOFF_MULTIPLIER = 2; + +/** Maximum number of attempts (1 initial + N retries) before giving up. */ +export const MAX_ATTEMPTS = 5; + +/** Errors whose messages indicate an RPC transport / connection problem. */ +const CONNECTION_ERROR_PATTERNS = [ + /ECONNREFUSED/i, + /ECONNRESET/i, + /ETIMEDOUT/i, + /ENOTFOUND/i, + /EHOSTUNREACH/i, + /network.*error/i, + /connection.*timeout/i, + /socket.*hang.*up/i, + /fetch.*failed/i, + /request.*timeout/i, + /getaddrinfo/i, + /connect.*ECONNREFUSED/i, +] as const; + +/** + * Determine whether `err` represents a transient RPC connection / network + * failure that should trigger a retry rather than an immediate abort. + */ +export function isConnectionError(err: unknown): boolean { + if (!(err instanceof Error)) return false; + const msg = err.message; + return CONNECTION_ERROR_PATTERNS.some((pattern) => pattern.test(msg)); +} + +/** + * Compute the next backoff duration using exponential growth capped at + * `MAX_BACKOFF_MS`. + * + * @param currentBackoffMs The backoff used for the most-recent delay. + * @returns The next backoff in milliseconds. + */ +export function nextBackoff(currentBackoffMs: number): number { + return Math.min(currentBackoffMs * BACKOFF_MULTIPLIER, MAX_BACKOFF_MS); +} + +/** Parameters forwarded directly to `Server.getEvents`. */ +export interface GetEventsParams { + startLedger: number; + contractIds: string[]; + limit?: number; +} + +/** Shape of `Server.getEvents` that this module depends on (subset). */ +export interface RpcGetEventsResult { + events: Array<{ + contractId: { contractId(): string } | undefined; + topic: unknown[]; + value: unknown; + ledger: number; + ledgerClosedAt: string | undefined; + }>; +} + +/** + * Options for `fetchEventsWithRetry`. + */ +export interface FetchEventsOptions { + /** Override the default maximum attempts (default: MAX_ATTEMPTS). */ + maxAttempts?: number; + /** Override the initial backoff delay in ms (default: DEFAULT_BACKOFF_MS). */ + initialBackoffMs?: number; + /** + * Inject a custom sleep function so tests can skip real delays. + * Defaults to a Promise-based setTimeout wrapper. + */ + sleep?: (ms: number) => Promise; +} + +const defaultSleep = (ms: number): Promise => + new Promise((resolve) => setTimeout(resolve, ms)); + +/** + * Fetch filtered contract events from the Soroban RPC with exponential-backoff + * retry on transient connection errors. + * + * Only the event types listed in `EVENT_TYPES` are requested, so irrelevant + * contract events are never returned to the caller. + * + * On each connection failure the wait before the next attempt doubles, starting + * at `initialBackoffMs` and never exceeding `MAX_BACKOFF_MS`. Non-connection + * errors (e.g. bad request, contract-not-found) are propagated immediately + * without retrying. + * + * @param server Soroban RPC `Server` instance (or compatible mock). + * @param params Event query parameters. + * @param options Optional overrides for retry tuning and sleep injection. + * @returns The raw `getEvents` response on success. + * @throws Re-throws the last error after `maxAttempts` are exhausted, + * or the first non-connection error encountered. + */ +export async function fetchEventsWithRetry( + server: Pick, + params: GetEventsParams, + options: FetchEventsOptions = {} +): Promise { + const maxAttempts = options.maxAttempts ?? MAX_ATTEMPTS; + const sleep = options.sleep ?? defaultSleep; + + let backoffMs = options.initialBackoffMs ?? DEFAULT_BACKOFF_MS; + let lastError: unknown; + + for (let attempt = 1; attempt <= maxAttempts; attempt++) { + try { + const result = await server.getEvents({ + startLedger: params.startLedger, + filters: [ + { + type: "contract", + contractIds: params.contractIds, + topics: [[...EVENT_TYPES]], + }, + ], + limit: params.limit ?? 100, + }); + + if (attempt > 1) { + logger.info("event_type_filter: RPC getEvents succeeded after retry", { + attempt, + startLedger: params.startLedger, + }); + } + + return result as unknown as RpcGetEventsResult; + } catch (err) { + lastError = err; + + // Non-connection errors are unrecoverable – propagate immediately. + if (!isConnectionError(err)) { + logger.error("event_type_filter: non-connection error from getEvents", { + attempt, + error: err instanceof Error ? err.message : String(err), + }); + throw err; + } + + const attemptsLeft = maxAttempts - attempt; + + if (attemptsLeft === 0) { + logger.error( + "event_type_filter: all retry attempts exhausted for getEvents", + { + maxAttempts, + startLedger: params.startLedger, + error: err instanceof Error ? err.message : String(err), + } + ); + break; + } + + logger.warn( + "event_type_filter: RPC connection error – retrying with backoff", + { + attempt, + attemptsLeft, + backoffMs, + error: err instanceof Error ? err.message : String(err), + } + ); + + await sleep(backoffMs); + backoffMs = nextBackoff(backoffMs); + } + } + + throw lastError; +} diff --git a/src/indexer/poller.ts b/src/indexer/poller.ts index 275006e..97128ff 100644 --- a/src/indexer/poller.ts +++ b/src/indexer/poller.ts @@ -8,6 +8,7 @@ import { type EventRow, } from "./db.js"; import { deliverWebhooks } from "./webhook-delivery.js"; +import { fetchEventsWithRetry } from "./event_type_filter.js"; import logger from "../utils/logger.js"; const RPC_URL = @@ -15,19 +16,6 @@ const RPC_URL = const server = new Server(RPC_URL); const POLL_INTERVAL_MS = parseInt(process.env.POLL_INTERVAL_MS || "15000", 10); -const EVENT_TYPES = [ - "initialized", - "funded", - "delivered", - "approved", - "dispute_raised", - "dispute_resolved", - "partial_release", - "auto_release_claimed", - "token_whitelisted", - "token_removed", -]; - /** * Poll events for all active contract IDs stored in monitored_contracts (#85). * All events fetched in a single poll are written atomically together with the @@ -59,15 +47,9 @@ export async function pollEvents() { logger.info("Polling events", { startLedger, currentLedger }); - const events = await server.getEvents({ + const events = await fetchEventsWithRetry(server, { startLedger, - filters: [ - { - type: "contract", - contractIds, - topics: [[...EVENT_TYPES]], - }, - ], + contractIds, limit: 100, });