From c0383c17c4d1e30a166a9a283932eb70dafad2e0 Mon Sep 17 00:00:00 2001 From: Hamda-gbade Date: Sun, 30 Aug 2026 14:05:49 +0000 Subject: [PATCH] Fix indexer freeze when a single event fails to process (#1214) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A failing event used to freeze lastCursor for every subsequent event in the batch (the !hasError guard), so a single always-failing event reprocessed all later events on every poll indefinitely. The cursor now always advances past successfully processed events, and failed events are recorded in a new IndexerDeadLetterEvent table (raw payload, error, attempt count) for manual triage, with a retry cap (INDEXER_DEAD_LETTER_MAX_RETRIES, default 5) after which the event is abandoned and the cursor advances past it. 🤖 Generated with Codebuff Co-Authored-By: Codebuff --- backend/docs/SSE_ARCHITECTURE.md | 3 + .../migration.sql | 31 +++ backend/prisma/schema.prisma | 20 ++ backend/src/workers/soroban-event-worker.ts | 76 +++++- backend/tests/soroban-event-worker.test.ts | 240 +++++++++++++++++- 5 files changed, 360 insertions(+), 10 deletions(-) create mode 100644 backend/prisma/migrations/20260830000000_add_indexer_dead_letter_event/migration.sql diff --git a/backend/docs/SSE_ARCHITECTURE.md b/backend/docs/SSE_ARCHITECTURE.md index 3bf547c3..cff662a3 100644 --- a/backend/docs/SSE_ARCHITECTURE.md +++ b/backend/docs/SSE_ARCHITECTURE.md @@ -294,6 +294,9 @@ The indexer worker behavior is controlled by environment variables configured in | `SOROBAN_RPC_URL` | Endpoint URL for the Soroban RPC node. | `"https://soroban-testnet.stellar.org"` | Uses default public testnet RPC URL. | | `INDEXER_POLL_INTERVAL_MS` | Polling interval in milliseconds between event fetch cycles. | `"5000"` (5 seconds) | Uses default 5000 ms interval. | | `INDEXER_START_LEDGER` | Starting Stellar ledger sequence number for cold starts when no `IndexerState` record exists in the database. | `"0"` | Starts indexing from ledger 0 on initial setup. | +| `INDEXER_DEAD_LETTER_MAX_RETRIES` | Max failed processing attempts before an event is abandoned to the `IndexerDeadLetterEvent` table and the cursor advances past it. | `"5"` | Uses default of 5 attempts. | + +Failed events never freeze the indexer: the cursor always advances past successfully processed events even when an earlier event in the batch failed, and each failing event is recorded (with its raw payload) in the `IndexerDeadLetterEvent` table for manual triage. After `INDEXER_DEAD_LETTER_MAX_RETRIES` attempts the event is abandoned and the cursor advances past it. --- diff --git a/backend/prisma/migrations/20260830000000_add_indexer_dead_letter_event/migration.sql b/backend/prisma/migrations/20260830000000_add_indexer_dead_letter_event/migration.sql new file mode 100644 index 00000000..81f4d4d0 --- /dev/null +++ b/backend/prisma/migrations/20260830000000_add_indexer_dead_letter_event/migration.sql @@ -0,0 +1,31 @@ +-- Dead-letter table for Soroban events that failed to process. A single +-- malformed event must not freeze the indexer: after N failed attempts the +-- worker abandons the event (recording it here with its raw payload for +-- manual triage) and advances the cursor past it. + +-- CreateTable +CREATE TABLE "IndexerDeadLetterEvent" ( + "id" TEXT NOT NULL, + "eventId" TEXT NOT NULL, + "ledger" INTEGER NOT NULL, + "transactionHash" TEXT NOT NULL, + "rawPayload" TEXT NOT NULL, + "errorMessage" TEXT NOT NULL, + "attempts" INTEGER NOT NULL DEFAULT 1, + "lastAttemptAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT "IndexerDeadLetterEvent_pkey" PRIMARY KEY ("id") +); + +-- CreateIndex +CREATE UNIQUE INDEX "IndexerDeadLetterEvent_eventId_key" ON "IndexerDeadLetterEvent"("eventId"); + +-- CreateIndex +CREATE INDEX "IndexerDeadLetterEvent_ledger_idx" ON "IndexerDeadLetterEvent"("ledger"); + +-- CreateIndex +CREATE INDEX "IndexerDeadLetterEvent_transactionHash_idx" ON "IndexerDeadLetterEvent"("transactionHash"); + +-- CreateIndex +CREATE INDEX "IndexerDeadLetterEvent_createdAt_idx" ON "IndexerDeadLetterEvent"("createdAt"); diff --git a/backend/prisma/schema.prisma b/backend/prisma/schema.prisma index 320c1306..49979057 100644 --- a/backend/prisma/schema.prisma +++ b/backend/prisma/schema.prisma @@ -65,6 +65,26 @@ model IndexerState { updatedAt DateTime @updatedAt } +// IndexerDeadLetterEvent model - Soroban events that failed to process, kept +// with their raw payload for manual triage. The worker retries an event at +// most INDEXER_DEAD_LETTER_MAX_RETRIES times, then abandons it and advances +// the cursor so a single malformed event can never freeze the indexer. +model IndexerDeadLetterEvent { + id String @id @default(uuid()) + eventId String @unique // RPC paging-token id of the event + ledger Int // Ledger sequence the event was emitted in + transactionHash String // Stellar transaction hash + rawPayload String // Full raw event JSON for manual triage/replay + errorMessage String // Last error thrown while processing + attempts Int @default(1) // Number of failed processing attempts + lastAttemptAt DateTime @default(now()) + createdAt DateTime @default(now()) + + @@index([ledger]) + @@index([transactionHash]) + @@index([createdAt]) +} + // StreamEvent model - indexer events for tracking all on-chain stream activities model StreamEvent { id String @id @default(uuid()) diff --git a/backend/src/workers/soroban-event-worker.ts b/backend/src/workers/soroban-event-worker.ts index 495e2539..2023caf1 100644 --- a/backend/src/workers/soroban-event-worker.ts +++ b/backend/src/workers/soroban-event-worker.ts @@ -10,6 +10,9 @@ import { rpcPool } from "../lib/rpc-pool.js"; // ─── Config ────────────────────────────────────────────────────────────────── +/** Default max failed attempts before an event is abandoned to the dead-letter table. */ +const DEAD_LETTER_MAX_RETRIES_DEFAULT = 5; + // ─── XDR Decoding Helpers ──────────────────────────────────────────────────── /** Decode an ScVal symbol to a string. */ @@ -94,6 +97,8 @@ export class SorobanEventWorker { private readonly server: rpc.Server; private readonly pollIntervalMs: number; private readonly startLedger: number; + /** Max failed processing attempts before an event is abandoned (dead-lettered). */ + private readonly deadLetterMaxRetries: number; private isRunning = false; private pollTimer: NodeJS.Timeout | undefined; @@ -123,6 +128,11 @@ export class SorobanEventWorker { 10, ); this.startLedger = parseInt(process.env.INDEXER_START_LEDGER ?? "0", 10); + this.deadLetterMaxRetries = parseInt( + process.env.INDEXER_DEAD_LETTER_MAX_RETRIES ?? + String(DEAD_LETTER_MAX_RETRIES_DEFAULT), + 10, + ); } /** @@ -363,11 +373,13 @@ export class SorobanEventWorker { await this.processEvent(event); this.eventsProcessed += 1; this.recordOutcome(true); - if (!hasError) { - // Use the event ID as the cursor if pagingToken is not available - lastCursor = event.id; - lastLedger = event.ledger; - } + // Always advance the cursor past a successfully processed event, + // even when an earlier event in this batch failed. The previous + // `!hasError` guard froze the cursor after the first failure, so a + // single always-failing event (e.g. malformed body) reprocessed + // every later event on every poll indefinitely. + lastCursor = event.id; + lastLedger = event.ledger; } catch (err) { hasError = true; this.eventsFailed += 1; @@ -377,6 +389,18 @@ export class SorobanEventWorker { `[SorobanWorker] Failed to process event ${event.id}:`, err, ); + + // Record the event in the dead-letter table (raw payload preserved + // for manual triage). Once it has failed `deadLetterMaxRetries` + // times, abandon it and advance the cursor past it so the indexer + // is never frozen by a permanently-bad event. + if (await this.deadLetterEvent(event, err)) { + logger.warn( + `[SorobanWorker] Event ${event.id} exceeded ${this.deadLetterMaxRetries} attempts — abandoning (see IndexerDeadLetterEvent for triage).`, + ); + lastCursor = event.id; + lastLedger = event.ledger; + } // Continue processing subsequent events rather than halting. } } @@ -401,6 +425,48 @@ export class SorobanEventWorker { ); } + /** + * Record a failed event in the dead-letter table (with its raw payload for + * manual triage), incrementing its attempt counter. + * + * @returns `true` when the event has reached the retry cap and should be + * abandoned (cursor advanced past it); `false` to leave it for a retry + * on a future poll. Never throws — a dead-letter write failure must not + * abort the batch; in that case the event is simply left for the next + * poll. + */ + private async deadLetterEvent( + event: rpc.Api.EventResponse, + err: unknown, + ): Promise { + try { + const row = await prisma.indexerDeadLetterEvent.upsert({ + where: { eventId: event.id }, + create: { + eventId: event.id, + ledger: event.ledger, + transactionHash: event.txHash, + rawPayload: JSON.stringify(event), + errorMessage: err instanceof Error ? err.message : String(err), + attempts: 1, + lastAttemptAt: new Date(), + }, + update: { + errorMessage: err instanceof Error ? err.message : String(err), + attempts: { increment: 1 }, + lastAttemptAt: new Date(), + }, + }); + return row.attempts >= this.deadLetterMaxRetries; + } catch (dlErr) { + logger.error( + `[SorobanWorker] Failed to write dead-letter entry for event ${event.id}:`, + dlErr, + ); + return false; + } + } + /** * Dispatch a single contract event to the appropriate handler based on the * first topic symbol. diff --git a/backend/tests/soroban-event-worker.test.ts b/backend/tests/soroban-event-worker.test.ts index 8d7cf386..98ee7aa6 100644 --- a/backend/tests/soroban-event-worker.test.ts +++ b/backend/tests/soroban-event-worker.test.ts @@ -19,6 +19,9 @@ const mockPrismaObj = vi.hoisted(() => ({ upsert: vi.fn(), create: vi.fn(), }, + indexerDeadLetterEvent: { + upsert: vi.fn(), + }, $transaction: vi.fn((cb) => cb({ streamEvent: { findUnique: vi.fn(), upsert: vi.fn() }, user: { upsert: vi.fn() }, stream: { upsert: vi.fn(), update: vi.fn() } })), $disconnect: vi.fn(), })); @@ -55,6 +58,59 @@ import { SorobanEventWorker } from '../src/workers/soroban-event-worker.js'; import { prisma } from '../src/lib/prisma.js'; import logger from '../src/logger.js'; +/** Build a valid admin_transferred event (processes without throwing). */ +function makeAdminTransferredEvent( + id: string, + ledger: number, + txHash: string, +): rpc.Api.EventResponse { + return { + id, + type: 'contract', + ledger, + ledgerClosedAt: '2024-01-01T00:00:00Z', + txHash, + transactionIndex: 0, + operationIndex: 0, + inSuccessfulContractCall: true, + topic: [ + { switch: () => ({ value: 0 }), sym: () => 'admin_transferred' } as any, + ], + value: { + switch: () => ({ value: 4 }), + map: () => [ + { key: () => ({ sym: () => 'previous_admin' }), val: () => ({ address: () => ({ switch: () => ({ value: 0 }), accountId: () => ({ ed25519: () => Buffer.alloc(32) }) }) }) }, + { key: () => ({ sym: () => 'new_admin' }), val: () => ({ address: () => ({ switch: () => ({ value: 0 }), accountId: () => ({ ed25519: () => Buffer.alloc(32) }) }) }) }, + ] as any, + } as any, + }; +} + +/** Build a malformed fee_config_updated event (missing body fields → throws). */ +function makeMalformedEvent( + id: string, + ledger: number, + txHash: string, +): rpc.Api.EventResponse { + return { + id, + type: 'contract', + ledger, + ledgerClosedAt: '2024-01-01T00:00:00Z', + txHash, + transactionIndex: 0, + operationIndex: 0, + inSuccessfulContractCall: true, + topic: [ + { switch: () => ({ value: 0 }), sym: () => 'fee_config_updated' } as any, + ], + value: { + switch: () => ({ value: 4 }), + map: () => [] as any, + } as any, + }; +} + describe('SorobanEventWorker', () => { let worker: SorobanEventWorker; @@ -687,7 +743,7 @@ describe('SorobanEventWorker', () => { expect(typeof capturedEventUpsert?.create?.streamId).toBe('bigint'); }); - it('cursor_does_not_advance_past_failed_event_in_mixed_batch', async () => { + it('advances cursor past a failed event when later events in the batch succeed', async () => { // Setup initial state: lastCursor is 'cursor-initial' (prisma.indexerState.findUnique as ReturnType).mockResolvedValue({ id: 'singleton', @@ -702,6 +758,13 @@ describe('SorobanEventWorker', () => { updatedAt: new Date(), }); + // Dead-letter upsert: first (and only) failure stays below the retry cap. + (prisma.indexerDeadLetterEvent.upsert as ReturnType).mockResolvedValue({ + id: 'dl-1', + eventId: 'cursor-event-1', + attempts: 1, + }); + // Event 1: Missing required body fields for fee_config_updated -> handleFeeConfigUpdated throws const event1: rpc.Api.EventResponse = { id: 'cursor-event-1', @@ -804,13 +867,180 @@ describe('SorobanEventWorker', () => { expect(event2Writes.length).toBe(1); expect(event3Writes.length).toBe(1); - // Assert: persisted IndexerState.lastCursor is NOT advanced past the failed event's position - // (i.e. it must not be set to 'cursor-event-2' or 'cursor-event-3' after a failure in event 1) + // Assert: the failed event was dead-lettered with its raw payload so it + // can be triaged manually. + expect(prisma.indexerDeadLetterEvent.upsert).toHaveBeenCalledWith( + expect.objectContaining({ + where: { eventId: 'cursor-event-1' }, + create: expect.objectContaining({ + eventId: 'cursor-event-1', + ledger: 101, + transactionHash: 'tx-failed-1', + attempts: 1, + errorMessage: expect.any(String), + }), + }), + ); + + // Assert: the persisted IndexerState.lastCursor DID advance past the + // failed event's position, because events 2 and 3 processed + // successfully. A single bad event must not freeze the indexer. const indexerUpsertCalls = (prisma.indexerState.upsert as ReturnType).mock.calls; const lastSaveCall = indexerUpsertCalls[indexerUpsertCalls.length - 1]![0]; - expect(lastSaveCall.update.lastCursor).not.toBe('cursor-event-2'); - expect(lastSaveCall.update.lastCursor).not.toBe('cursor-event-3'); + expect(lastSaveCall.update.lastCursor).toBe('cursor-event-3'); + expect(lastSaveCall.update.lastLedger).toBe(103); + }); + + it('processes the other four events and advances the cursor when one event in a batch of five is malformed', async () => { + (prisma.indexerState.findUnique as ReturnType).mockResolvedValue({ + id: 'singleton', + lastLedger: 100, + lastCursor: 'batch-cursor-0', + updatedAt: new Date(), + }); + (prisma.indexerState.upsert as ReturnType).mockResolvedValue({ + id: 'singleton', + lastLedger: 100, + lastCursor: 'batch-cursor-0', + updatedAt: new Date(), + }); + (prisma.indexerDeadLetterEvent.upsert as ReturnType).mockResolvedValue({ + id: 'dl-batch', + eventId: 'batch-bad-3', + attempts: 1, + }); + + // One deliberately malformed event (position 3) in a batch of five. + const events = [ + makeAdminTransferredEvent('batch-ok-1', 101, 'tx-batch-1'), + makeAdminTransferredEvent('batch-ok-2', 102, 'tx-batch-2'), + makeMalformedEvent('batch-bad-3', 103, 'tx-batch-3'), + makeAdminTransferredEvent('batch-ok-4', 104, 'tx-batch-4'), + makeAdminTransferredEvent('batch-ok-5', 105, 'tx-batch-5'), + ]; + + vi.spyOn((worker as any).server, 'getEvents').mockResolvedValue({ events }); + + const upsertedStreamEvents: any[] = []; + const mockTx = { + user: { upsert: vi.fn().mockResolvedValue({}) }, + stream: { upsert: vi.fn().mockResolvedValue({ streamId: 0n, isActive: false }) }, + streamEvent: { + findUnique: vi.fn().mockResolvedValue(null), + upsert: vi.fn().mockImplementation((args) => { + upsertedStreamEvents.push(args); + return Promise.resolve({ id: 'event-id' }); + }), + }, + }; + (prisma.$transaction as ReturnType).mockImplementation((cb) => cb(mockTx)); + + await (worker as any).fetchAndProcessEvents(); + + // All four valid events processed exactly once; malformed one never written. + for (const txHash of ['tx-batch-1', 'tx-batch-2', 'tx-batch-4', 'tx-batch-5']) { + expect( + upsertedStreamEvents.filter((e) => e.create?.transactionHash === txHash).length, + ).toBe(1); + } + expect( + upsertedStreamEvents.filter((e) => e.create?.transactionHash === 'tx-batch-3').length, + ).toBe(0); + + // Malformed event dead-lettered with its raw payload for manual triage. + expect(prisma.indexerDeadLetterEvent.upsert).toHaveBeenCalledWith( + expect.objectContaining({ + where: { eventId: 'batch-bad-3' }, + create: expect.objectContaining({ + transactionHash: 'tx-batch-3', + rawPayload: expect.stringContaining('batch-bad-3'), + attempts: 1, + }), + }), + ); + + // Cursor advances past the malformed event to the last processed event. + const lastSaveCall = + (prisma.indexerState.upsert as ReturnType).mock.calls.at(-1)![0]; + expect(lastSaveCall.update.lastCursor).toBe('batch-ok-5'); + expect(lastSaveCall.update.lastLedger).toBe(105); + + // Counters reflect 4 processed, 1 failed. + const counters = worker.getEventCounters(); + expect(counters.eventsProcessed).toBe(4); + expect(counters.eventsFailed).toBe(1); + expect(counters.lastErrorAt).not.toBeNull(); + }); + + it('abandons an always-failing event after the retry cap and advances past it', async () => { + process.env.INDEXER_DEAD_LETTER_MAX_RETRIES = '2'; + worker = new SorobanEventWorker(); + + (prisma.indexerState.findUnique as ReturnType).mockResolvedValue({ + id: 'singleton', + lastLedger: 100, + lastCursor: 'cap-cursor-0', + updatedAt: new Date(), + }); + (prisma.indexerState.upsert as ReturnType).mockResolvedValue({ + id: 'singleton', + lastLedger: 100, + lastCursor: 'cap-cursor-0', + updatedAt: new Date(), + }); + + const okEvent = makeAdminTransferredEvent('cap-ok-1', 101, 'tx-cap-1'); + const badEvent = makeMalformedEvent('cap-bad-2', 102, 'tx-cap-2'); + + // Poll 1: both events. Poll 2: only the bad tail event is re-fetched. + // Poll 3: cursor moved past it — nothing left to fetch. + const getEvents = vi + .spyOn((worker as any).server, 'getEvents') + .mockResolvedValueOnce({ events: [okEvent, badEvent] }) + .mockResolvedValueOnce({ events: [badEvent] }) + .mockResolvedValueOnce({ events: [] }); + + // Dead-letter attempts: 1 (below cap) then 2 (cap reached → abandon). + (prisma.indexerDeadLetterEvent.upsert as ReturnType) + .mockResolvedValueOnce({ id: 'dl-1', eventId: 'cap-bad-2', attempts: 1 }) + .mockResolvedValueOnce({ id: 'dl-2', eventId: 'cap-bad-2', attempts: 2 }); + + const mockTx = { + user: { upsert: vi.fn().mockResolvedValue({}) }, + stream: { upsert: vi.fn().mockResolvedValue({ streamId: 0n, isActive: false }) }, + streamEvent: { + findUnique: vi.fn().mockResolvedValue(null), + upsert: vi.fn().mockResolvedValue({ id: 'event-id' }), + }, + }; + (prisma.$transaction as ReturnType).mockImplementation((cb) => cb(mockTx)); + + // Poll 1: ok event processed, bad event fails but stays below the cap. + await (worker as any).fetchAndProcessEvents(); + let lastSaveCall = + (prisma.indexerState.upsert as ReturnType).mock.calls.at(-1)![0]; + expect(lastSaveCall.update.lastCursor).toBe('cap-ok-1'); + + // Poll 2: bad event retried, hits the cap → abandoned, cursor advances past it. + await (worker as any).fetchAndProcessEvents(); + lastSaveCall = + (prisma.indexerState.upsert as ReturnType).mock.calls.at(-1)![0]; + expect(lastSaveCall.update.lastCursor).toBe('cap-bad-2'); + expect(logger.warn).toHaveBeenCalledWith( + expect.stringContaining('exceeded 2 attempts — abandoning'), + ); + + // Poll 3: cursor is past the bad event — no more re-processing. + await (worker as any).fetchAndProcessEvents(); + expect(getEvents).toHaveBeenCalledTimes(3); + expect((prisma.indexerState.upsert as ReturnType).mock.calls.length).toBe(2); + + const counters = worker.getEventCounters(); + expect(counters.eventsProcessed).toBe(1); + expect(counters.eventsFailed).toBe(2); + + delete process.env.INDEXER_DEAD_LETTER_MAX_RETRIES; }); });