From 03730a491525d472f2a325b994253e87ddf73770 Mon Sep 17 00:00:00 2001 From: KingDavid9999 Date: Sun, 26 Jul 2026 10:35:47 -0700 Subject: [PATCH] feat(indexer): add structured batch start/completion logs Add per-batch structured logging to processIndexerChainEvents so operators can see indexer throughput and batch size without querying the database. - Emit an info-level "indexer_batch_started" log before processing, with from_ledger, to_ledger (min/max ledger across the batch), and batch_size (raw event count before dedup). - Emit a debug-level "indexer_batch_completed" log after processing, with from_ledger, to_ledger, events_processed (unique count after dedup), and duration_ms (wall-clock time via the existing monotonic clock helper). - Both logs fire exactly once per batch call, not per individual ledger/event; the existing per-event log is untouched. Closes #648 --- docs/indexer/EVENT_PROCESSING.md | 20 ++ .../indexer-event-processor.utils.test.ts | 244 ++++++++++++------ src/utils/indexer-event-processor.utils.ts | 115 +++++++-- 3 files changed, 279 insertions(+), 100 deletions(-) diff --git a/docs/indexer/EVENT_PROCESSING.md b/docs/indexer/EVENT_PROCESSING.md index d0e9e94..a3e9ebf 100644 --- a/docs/indexer/EVENT_PROCESSING.md +++ b/docs/indexer/EVENT_PROCESSING.md @@ -4,6 +4,8 @@ The indexer processes events from the blockchain to update the read models and a Before changing this pipeline, review [Indexer Contributor Expectations](./CONTRIBUTOR_EXPECTATIONS.md) for the invariants, testing expectations, and deployment notes that apply to indexer work. +For the end-to-end architecture including the polling actor and write path, see [Indexer Architecture](./ARCHITECTURE.md). + ## 1. Deduplication Before processing a batch of events, they should be deduped based on their unique identifier on the chain: `transactionHash` and `eventIndex`. @@ -54,6 +56,24 @@ The log includes: Use `processIndexerChainEvents` to dedupe a batch and log once per unique event. +## 5. Structured batch logging + +In addition to the per-event log above, `processIndexerChainEvents` emits +exactly two logs per batch — never per individual ledger: + +| Log | Level | Fields | +| :------------------------ | :---- | :---------------------------------------------------------------------------------- | +| `indexer_batch_started` | info | `from_ledger`, `to_ledger`, `batch_size` (raw batch size, before dedup) | +| `indexer_batch_completed` | debug | `from_ledger`, `to_ledger`, `events_processed` (unique, after dedup), `duration_ms` | + +`from_ledger`/`to_ledger` are the min/max `ledger` values across the events +in the batch (`undefined` if no event in the batch carries a `ledger`). +`duration_ms` is measured with the same monotonic clock used for per-event +timing, from the start of the batch to the completion of the last event. + +These logs give operators batch throughput and size at a glance without +requiring per-ledger database queries. + ## 3. Error Handling If an event fails to process after multiple retries, it is moved to the [Dead-Letter Queue (DLQ)](./DLQ_WORKFLOW.md) for manual investigation. diff --git a/src/utils/indexer-event-processor.utils.test.ts b/src/utils/indexer-event-processor.utils.test.ts index 94ae3be..c555dc9 100644 --- a/src/utils/indexer-event-processor.utils.test.ts +++ b/src/utils/indexer-event-processor.utils.test.ts @@ -1,87 +1,183 @@ import { logger } from './logger.utils'; import { - getChainEventId, - processIndexerChainEvent, - processIndexerChainEvents, - IndexerChainEvent, + getChainEventId, + processIndexerChainEvent, + processIndexerChainEvents, + IndexerChainEvent, } from './indexer-event-processor.utils'; jest.mock('./logger.utils', () => ({ - logger: { - info: jest.fn(), - }, + logger: { + info: jest.fn(), + debug: jest.fn(), + }, })); const infoMock = logger.info as jest.Mock; +const debugMock = logger.debug as jest.Mock; -function makeEvent(overrides: Partial = {}): IndexerChainEvent { - return { - txHash: '0xabc123', - eventIndex: 0, - eventType: 'CREATOR_REGISTERED', - ...overrides, - }; +function makeEvent( + overrides: Partial = {} +): IndexerChainEvent { + return { + txHash: '0xabc123', + eventIndex: 0, + eventType: 'CREATOR_REGISTERED', + ...overrides, + }; } describe('indexer-event-processor.utils', () => { - beforeEach(() => { - infoMock.mockClear(); - }); - - describe('getChainEventId', () => { - it('combines txHash and eventIndex', () => { - expect(getChainEventId(makeEvent({ txHash: '0xdead', eventIndex: 3 }))).toBe( - '0xdead:3' - ); - }); - }); - - describe('processIndexerChainEvent', () => { - it('emits one structured log after the handler completes', async () => { - const event = makeEvent(); - const handler = jest.fn().mockResolvedValue(undefined); - - await processIndexerChainEvent(event, handler); - - expect(handler).toHaveBeenCalledWith(event); - expect(infoMock).toHaveBeenCalledTimes(1); - expect(infoMock).toHaveBeenCalledWith( - expect.objectContaining({ - type: 'indexer_event_processed', - eventType: 'CREATOR_REGISTERED', - eventId: '0xabc123:0', - txHash: '0xabc123', - eventIndex: 0, - elapsedMs: expect.any(Number), - }), - 'Indexer chain event processed' - ); - }); - - it('does not emit a log when the handler throws', async () => { - const handler = jest.fn().mockRejectedValue(new Error('handler failed')); - - await expect( - processIndexerChainEvent(makeEvent(), handler) - ).rejects.toThrow('handler failed'); - - expect(infoMock).not.toHaveBeenCalled(); - }); - }); - - describe('processIndexerChainEvents', () => { - it('dedupes events and logs once per unique event', async () => { - const events: IndexerChainEvent[] = [ - makeEvent({ txHash: '0x1', eventIndex: 0, eventType: 'KEY_BOUGHT' }), - makeEvent({ txHash: '0x1', eventIndex: 0, eventType: 'KEY_BOUGHT' }), - makeEvent({ txHash: '0x1', eventIndex: 1, eventType: 'KEY_SOLD' }), - ]; - const handler = jest.fn().mockResolvedValue(undefined); - - await processIndexerChainEvents(events, handler); - - expect(handler).toHaveBeenCalledTimes(2); - expect(infoMock).toHaveBeenCalledTimes(2); - }); - }); + beforeEach(() => { + infoMock.mockClear(); + debugMock.mockClear(); + }); + + describe('getChainEventId', () => { + it('combines txHash and eventIndex', () => { + expect( + getChainEventId(makeEvent({ txHash: '0xdead', eventIndex: 3 })) + ).toBe('0xdead:3'); + }); + }); + + describe('processIndexerChainEvent', () => { + it('emits one structured log after the handler completes', async () => { + const event = makeEvent(); + const handler = jest.fn().mockResolvedValue(undefined); + + await processIndexerChainEvent(event, handler); + + expect(handler).toHaveBeenCalledWith(event); + expect(infoMock).toHaveBeenCalledTimes(1); + expect(infoMock).toHaveBeenCalledWith( + expect.objectContaining({ + type: 'indexer_event_processed', + eventType: 'CREATOR_REGISTERED', + eventId: '0xabc123:0', + txHash: '0xabc123', + eventIndex: 0, + elapsedMs: expect.any(Number), + }), + 'Indexer chain event processed' + ); + }); + + it('does not emit a log when the handler throws', async () => { + const handler = jest + .fn() + .mockRejectedValue(new Error('handler failed')); + + await expect( + processIndexerChainEvent(makeEvent(), handler) + ).rejects.toThrow('handler failed'); + + expect(infoMock).not.toHaveBeenCalled(); + }); + }); + + describe('processIndexerChainEvents', () => { + it('dedupes events and logs once per unique event', async () => { + const events: IndexerChainEvent[] = [ + makeEvent({ + txHash: '0x1', + eventIndex: 0, + eventType: 'KEY_BOUGHT', + }), + makeEvent({ + txHash: '0x1', + eventIndex: 0, + eventType: 'KEY_BOUGHT', + }), + makeEvent({ txHash: '0x1', eventIndex: 1, eventType: 'KEY_SOLD' }), + ]; + const handler = jest.fn().mockResolvedValue(undefined); + + await processIndexerChainEvents(events, handler); + + expect(handler).toHaveBeenCalledTimes(2); + // 1 batch-start log + 2 per-event logs + expect(infoMock).toHaveBeenCalledTimes(3); + }); + + it('emits an info-level batch-start log with the ledger range and raw batch size', async () => { + const events: IndexerChainEvent[] = [ + makeEvent({ txHash: '0x1', eventIndex: 0, ledger: 105 }), + makeEvent({ txHash: '0x1', eventIndex: 0, ledger: 105 }), // duplicate + makeEvent({ txHash: '0x2', eventIndex: 0, ledger: 110 }), + ]; + const handler = jest.fn().mockResolvedValue(undefined); + + await processIndexerChainEvents(events, handler); + + expect(infoMock).toHaveBeenNthCalledWith( + 1, + expect.objectContaining({ + type: 'indexer_batch_started', + from_ledger: 105, + to_ledger: 110, + batch_size: 3, // size before dedup + }), + 'Indexer batch started' + ); + }); + + it('emits a debug-level batch-completion log with events processed and duration', async () => { + const events: IndexerChainEvent[] = [ + makeEvent({ txHash: '0x1', eventIndex: 0, ledger: 105 }), + makeEvent({ txHash: '0x1', eventIndex: 0, ledger: 105 }), // duplicate + makeEvent({ txHash: '0x2', eventIndex: 0, ledger: 110 }), + ]; + const handler = jest.fn().mockResolvedValue(undefined); + + await processIndexerChainEvents(events, handler); + + expect(debugMock).toHaveBeenCalledTimes(1); + expect(debugMock).toHaveBeenCalledWith( + expect.objectContaining({ + type: 'indexer_batch_completed', + from_ledger: 105, + to_ledger: 110, + events_processed: 2, // unique events after dedup + duration_ms: expect.any(Number), + }), + 'Indexer batch completed' + ); + }); + + it('emits exactly one batch-start and one batch-completion log regardless of batch size', async () => { + const events: IndexerChainEvent[] = [ + makeEvent({ txHash: '0x1', eventIndex: 0, ledger: 1 }), + makeEvent({ txHash: '0x2', eventIndex: 0, ledger: 2 }), + makeEvent({ txHash: '0x3', eventIndex: 0, ledger: 3 }), + ]; + const handler = jest.fn().mockResolvedValue(undefined); + + await processIndexerChainEvents(events, handler); + + // 1 batch-start + 3 per-event info logs + expect(infoMock).toHaveBeenCalledTimes(4); + expect(debugMock).toHaveBeenCalledTimes(1); + }); + + it('handles a batch with no ledger values by omitting the range', async () => { + const events: IndexerChainEvent[] = [ + makeEvent({ txHash: '0x1', eventIndex: 0, ledger: undefined }), + ]; + const handler = jest.fn().mockResolvedValue(undefined); + + await processIndexerChainEvents(events, handler); + + expect(infoMock).toHaveBeenNthCalledWith( + 1, + expect.objectContaining({ + type: 'indexer_batch_started', + from_ledger: undefined, + to_ledger: undefined, + batch_size: 1, + }), + 'Indexer batch started' + ); + }); + }); }); diff --git a/src/utils/indexer-event-processor.utils.ts b/src/utils/indexer-event-processor.utils.ts index 125810b..45e390f 100644 --- a/src/utils/indexer-event-processor.utils.ts +++ b/src/utils/indexer-event-processor.utils.ts @@ -6,15 +6,15 @@ import { elapsedMs, startTimer } from './monotonic-clock.utils'; * Minimal chain event shape required for indexer processing and logging. */ export interface IndexerChainEvent extends ChainEvent { - /** Domain event type (e.g. CREATOR_REGISTERED, KEY_BOUGHT). */ - eventType: string; + /** Domain event type (e.g. CREATOR_REGISTERED, KEY_BOUGHT). */ + eventType: string; } /** * Stable identifier for a chain event, used for deduplication and log correlation. */ export function getChainEventId(event: ChainEvent): string { - return `${event.txHash}:${event.eventIndex}`; + return `${event.txHash}:${event.eventIndex}`; } /** @@ -25,40 +25,103 @@ export function getChainEventId(event: ChainEvent): string { * monotonic clock. */ export async function processIndexerChainEvent( - event: T, - handler: (event: T) => Promise + event: T, + handler: (event: T) => Promise ): Promise { - const timer = startTimer(); - const eventId = getChainEventId(event); + const timer = startTimer(); + const eventId = getChainEventId(event); - await handler(event); + await handler(event); - logger.info( - { - type: 'indexer_event_processed', - eventType: event.eventType, - eventId, - txHash: event.txHash, - eventIndex: event.eventIndex, - ledger: event.ledger, - elapsedMs: elapsedMs(timer), - }, - 'Indexer chain event processed' - ); + logger.info( + { + type: 'indexer_event_processed', + eventType: event.eventType, + eventId, + txHash: event.txHash, + eventIndex: event.eventIndex, + ledger: event.ledger, + elapsedMs: elapsedMs(timer), + }, + 'Indexer chain event processed' + ); +} + +/** + * Ledger range covered by a batch of chain events, derived from the + * `ledger` field present on each event. + */ +interface BatchLedgerRange { + fromLedger: number | undefined; + toLedger: number | undefined; +} + +/** + * Derives the ledger range (min/max) covered by a batch of chain events. + * + * Events without a `ledger` value are ignored. If no event in the batch + * has a `ledger`, both bounds are `undefined`. + */ +function getBatchLedgerRange( + events: T[] +): BatchLedgerRange { + const ledgers = events + .map(event => event.ledger) + .filter((ledger): ledger is number => typeof ledger === 'number'); + + if (ledgers.length === 0) { + return { fromLedger: undefined, toLedger: undefined }; + } + + return { + fromLedger: Math.min(...ledgers), + toLedger: Math.max(...ledgers), + }; } /** * Dedupes a batch of chain events and processes each unique event sequentially. * * Each event emits one structured log entry via {@link processIndexerChainEvent}. + * + * The batch itself also emits exactly two structured logs: + * - An info-level log when the batch starts, with the ledger range and + * the size of the incoming batch (before deduplication). + * - A debug-level log when the batch completes, with the ledger range, + * the number of unique events actually processed, and the wall-clock + * duration of the whole batch. */ export async function processIndexerChainEvents( - events: T[], - handler: (event: T) => Promise + events: T[], + handler: (event: T) => Promise ): Promise { - const uniqueEvents = dedupeChainEvents(events); + const { fromLedger, toLedger } = getBatchLedgerRange(events); + const batchTimer = startTimer(); + + logger.info( + { + type: 'indexer_batch_started', + from_ledger: fromLedger, + to_ledger: toLedger, + batch_size: events.length, + }, + 'Indexer batch started' + ); + + const uniqueEvents = dedupeChainEvents(events); + + for (const event of uniqueEvents) { + await processIndexerChainEvent(event, handler); + } - for (const event of uniqueEvents) { - await processIndexerChainEvent(event, handler); - } + logger.debug( + { + type: 'indexer_batch_completed', + from_ledger: fromLedger, + to_ledger: toLedger, + events_processed: uniqueEvents.length, + duration_ms: elapsedMs(batchTimer), + }, + 'Indexer batch completed' + ); }