Skip to content
Merged
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
18 changes: 18 additions & 0 deletions docs/indexer/EVENT_PROCESSING.md
Original file line number Diff line number Diff line change
Expand Up @@ -56,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.
86 changes: 85 additions & 1 deletion src/utils/indexer-event-processor.utils.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,10 +9,12 @@ import {
jest.mock('./logger.utils', () => ({
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> = {}
Expand All @@ -28,6 +30,7 @@ function makeEvent(
describe('indexer-event-processor.utils', () => {
beforeEach(() => {
infoMock.mockClear();
debugMock.mockClear();
});

describe('getChainEventId', () => {
Expand Down Expand Up @@ -93,7 +96,88 @@ describe('indexer-event-processor.utils', () => {
await processIndexerChainEvents(events, handler);

expect(handler).toHaveBeenCalledTimes(2);
expect(infoMock).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'
);
});
});
});
63 changes: 63 additions & 0 deletions src/utils/indexer-event-processor.utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -47,18 +47,81 @@ export async function processIndexerChainEvent<T extends IndexerChainEvent>(
);
}

/**
* 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<T extends ChainEvent>(
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<T extends IndexerChainEvent>(
events: T[],
handler: (event: T) => Promise<void>
): Promise<void> {
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);
}

logger.debug(
{
type: 'indexer_batch_completed',
from_ledger: fromLedger,
to_ledger: toLedger,
events_processed: uniqueEvents.length,
duration_ms: elapsedMs(batchTimer),
},
'Indexer batch completed'
);
}
Loading