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
22 changes: 22 additions & 0 deletions backend/src/common/observability/metrics.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,12 @@ export interface MetricsData {
retention: {
rows_pruned_total: Record<string, number>;
};
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<string, { state: string; failures: number; opens: number }>;
timeouts?: {
request_timeout_ms: number;
Expand All @@ -58,6 +64,8 @@ let latencyCount = 0;
let slowQueryCount = 0;
let poolSaturationCount = 0;
const retentionPrunedCounts: Record<string, number> = {};
let unknownEventCount = 0;
let lastProcessedLedger: number | null = null;

export function recordRequest(duration: number) {
requestCount++;
Expand All @@ -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<MetricsData> {
logger.debug('Collecting metrics');

Expand Down Expand Up @@ -145,6 +163,10 @@ export async function getMetrics(): Promise<MetricsData> {
retention: {
rows_pruned_total: { ...retentionPrunedCounts },
},
indexer: {
unknown_events_total: unknownEventCount,
last_processed_ledger: lastProcessedLedger,
},
circuitBreaker,
timeouts: {
request_timeout_ms: env.REQUEST_TIMEOUT_MS,
Expand Down
13 changes: 13 additions & 0 deletions backend/src/indexer/projections.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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']);
Expand Down Expand Up @@ -45,12 +46,22 @@ export async function projectEvent(event: DecodedEvent): Promise<void> {
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);
Expand All @@ -59,6 +70,7 @@ export async function projectEvent(event: DecodedEvent): Promise<void> {
if (isNewEvent) {
await publishProjection(event);
}
recordIndexerLedgerProcessed(event.ledger);
}

/**
Expand Down Expand Up @@ -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();
}
77 changes: 77 additions & 0 deletions backend/tests/unknown-events.test.ts
Original file line number Diff line number Diff line change
@@ -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<number> {
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);
});
});