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
3 changes: 3 additions & 0 deletions backend/docs/SSE_ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

---

Expand Down
Original file line number Diff line number Diff line change
@@ -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");
20 changes: 20 additions & 0 deletions backend/prisma/schema.prisma
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,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())
Expand Down
76 changes: 71 additions & 5 deletions backend/src/workers/soroban-event-worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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. */
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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,
);
}

/**
Expand Down Expand Up @@ -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;
Expand All @@ -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.
}
}
Expand All @@ -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<boolean> {
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.
Expand Down
Loading