From f639b2e0ce5f9236b9ee4f9370e52b9ae113ac37 Mon Sep 17 00:00:00 2001 From: aurorabini <319766468+aurorabini@users.noreply.github.com> Date: Sat, 29 Aug 2026 08:30:35 +0200 Subject: [PATCH] feat(indexer): handle unknown and future event versions without stalling Unknown event types are already persisted raw to EventLog; this surfaces them with a warn log, counts them in the /metrics indexer section (so a silent drop is impossible), and treats a future version marker under a known topic as a still-known projection. Covers unknown type, future version, malformed payload, and dedupe in tests. Closes #1261 --- backend/src/common/observability/metrics.ts | 22 ++++++ backend/src/indexer/projections.ts | 13 ++++ backend/tests/unknown-events.test.ts | 77 +++++++++++++++++++++ 3 files changed, 112 insertions(+) create mode 100644 backend/tests/unknown-events.test.ts diff --git a/backend/src/common/observability/metrics.ts b/backend/src/common/observability/metrics.ts index 353148aa..1686895c 100644 --- a/backend/src/common/observability/metrics.ts +++ b/backend/src/common/observability/metrics.ts @@ -41,6 +41,12 @@ export interface MetricsData { retention: { rows_pruned_total: Record; }; + indexer?: { + /** Events whose topic/version is not yet understood by the indexer. */ + unknown_events_total: number; + /** Last ledger successfully processed by the indexer (for lag = chainHead - this). */ + last_processed_ledger: number | null; + }; circuitBreaker?: Record; timeouts?: { request_timeout_ms: number; @@ -58,6 +64,8 @@ let latencyCount = 0; let slowQueryCount = 0; let poolSaturationCount = 0; const retentionPrunedCounts: Record = {}; +let unknownEventCount = 0; +let lastProcessedLedger: number | null = null; export function recordRequest(duration: number) { requestCount++; @@ -83,6 +91,16 @@ export function recordRetentionPruned(model: string, count: number): void { retentionPrunedCounts[model] = (retentionPrunedCounts[model] ?? 0) + count; } +/** Records an indexer event the indexer does not yet understand (issue #1261). */ +export function recordUnknownEvent(): void { + unknownEventCount++; +} + +/** Records the last ledger successfully processed by the indexer (issue #1258 / #1261). */ +export function recordIndexerLedgerProcessed(ledger: number): void { + lastProcessedLedger = ledger; +} + export async function getMetrics(): Promise { logger.debug('Collecting metrics'); @@ -145,6 +163,10 @@ export async function getMetrics(): Promise { retention: { rows_pruned_total: { ...retentionPrunedCounts }, }, + indexer: { + unknown_events_total: unknownEventCount, + last_processed_ledger: lastProcessedLedger, + }, circuitBreaker, timeouts: { request_timeout_ms: env.REQUEST_TIMEOUT_MS, diff --git a/backend/src/indexer/projections.ts b/backend/src/indexer/projections.ts index faf6d6e0..a22efcf8 100644 --- a/backend/src/indexer/projections.ts +++ b/backend/src/indexer/projections.ts @@ -4,6 +4,7 @@ import { logger } from '../common/utils/logger.js'; import type { DecodedEvent } from './sorobanClient.js'; import { publishProjection } from './realtime-publisher.js'; import * as notificationsService from '../modules/notifications/notifications.service.js'; +import { recordUnknownEvent, recordIndexerLedgerProcessed } from '../common/observability/metrics.js'; /** Event topics that represent an on-chain tip. */ const TIP_TOPICS = new Set(['tip', 'tip_sent']); @@ -45,12 +46,22 @@ export async function projectEvent(event: DecodedEvent): Promise { if (isNewEvent) { await publishProjection(event); } + recordIndexerLedgerProcessed(event.ledger); return; } const handler = PROJECTIONS[event.topic]; if (handler) { await handler(event, isNewEvent); + } else if (!REFUND_TOPICS.has(event.topic)) { + // Unknown event type or version (issue #1261). The raw event was already + // persisted to EventLog above, so it can be replayed once a decoder ships. + // Surface it loudly and count it — never crash or stall the pipeline. + logger.warn( + { txHash: event.txHash, topic: event.topic, ledger: event.ledger }, + 'Indexer encountered an unknown event type/version; raw event persisted to EventLog', + ); + recordUnknownEvent(); } if (REFUND_TOPICS.has(event.topic)) { await projectRefund(event); @@ -59,6 +70,7 @@ export async function projectEvent(event: DecodedEvent): Promise { if (isNewEvent) { await publishProjection(event); } + recordIndexerLedgerProcessed(event.ledger); } /** @@ -606,4 +618,5 @@ function addDays(from: Date, days: number): Date { function warnUnparseable(event: DecodedEvent, topic: string): void { logger.warn({ txHash: event.txHash, topic }, 'Skipping event with unparseable payload'); + recordUnknownEvent(); } \ No newline at end of file diff --git a/backend/tests/unknown-events.test.ts b/backend/tests/unknown-events.test.ts new file mode 100644 index 00000000..adb1e6b9 --- /dev/null +++ b/backend/tests/unknown-events.test.ts @@ -0,0 +1,77 @@ +import { beforeEach, describe, expect, it } from 'vitest'; +import { prisma } from './helpers/db.js'; +import { resetDb } from './helpers/db.js'; +import { getMetrics } from '../src/common/observability/metrics.js'; +import { projectEvent } from '../src/indexer/projections.js'; +import { ADDR_A, ADDR_B } from '../src/indexer/fixtures/events.js'; +import type { DecodedEvent } from '../src/indexer/sorobanClient.js'; + +beforeEach(() => resetDb()); + +const baseEvent: DecodedEvent = { + ledger: 200, + txHash: 'unknown-tx-0001', + pagingToken: '200-0', + topic: 'some_future_topic', + value: { + from: ADDR_A, + to: ADDR_B, + amount: '1000000', + }, +}; + +async function unknownTotal(): Promise { + const metrics = await getMetrics(); + return metrics.indexer?.unknown_events_total ?? 0; +} + +describe('indexer unknown & future event handling (issue #1261)', () => { + it('persists an unknown event type raw to EventLog and warns', async () => { + await projectEvent(baseEvent); + + const log = await prisma.eventLog.findFirst({ where: { txHash: baseEvent.txHash } }); + expect(log).not.toBeNull(); + expect(log?.topic).toBe('some_future_topic'); + expect(log?.ledger).toBe(200); + expect(await unknownTotal()).toBe(1); + }); + + it('does not crash or create projection rows for an unknown event', async () => { + await projectEvent(baseEvent); + + const tips = await prisma.tip.count(); + expect(tips).toBe(0); + }); + + it('counts unknown events — visible in metrics, not silent', async () => { + await projectEvent(baseEvent); + await projectEvent({ ...baseEvent, txHash: 'unknown-tx-0002' }); + expect(await unknownTotal()).toBe(2); + }); + + it('persists unknown future-version events raw for later replay', async () => { + const futEvent: DecodedEvent = { + ...baseEvent, + txHash: 'future-version-tx', + topic: 'tip_sent', + value: { from: ADDR_A, to: ADDR_B, amount: '5000000', version: 99 }, + }; + // tip_sent is known; a future "version" marker doesn't change handling, and + // an unknown subtype under a known topic is still a known projection. + await projectEvent(futEvent); + const tip = await prisma.tip.findFirst({ where: { txHash: 'future-version-tx' } }); + expect(tip).not.toBeNull(); + }); + + it('counts malformed known-event payloads as unknown events', async () => { + await projectEvent({ ...baseEvent, topic: 'tip_sent', value: { invalid: true } }); + expect(await unknownTotal()).toBe(1); + }); + + it('dedupes unknown events by EventLog unique key (txHash, topic, ledger)', async () => { + await projectEvent(baseEvent); + await projectEvent(baseEvent); + const logs = await prisma.eventLog.findMany({ where: { txHash: baseEvent.txHash } }); + expect(logs).toHaveLength(1); + }); +}); \ No newline at end of file