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
74 changes: 74 additions & 0 deletions src/workers/indexer/BoundedEventQueue.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
/**
* Bounded in-memory queue between event fetch and DB write.
*
* When Postgres write throughput cannot keep up with a burst of on-chain
* events, the queue fills up and producers pause (await `enqueue`) instead
* of piling up unbounded memory. Consumers drain via `dequeueBatch`.
*/
export class BoundedEventQueue<T> {
private items: T[] = [];
private waiters: Array<() => void> = [];

constructor(private readonly maxSize: number) {
if (!Number.isInteger(maxSize) || maxSize <= 0) {
throw new Error('BoundedEventQueue maxSize must be a positive integer');
}
}

get capacity(): number {
return this.maxSize;
}

get size(): number {
return this.items.length;
}

get isFull(): boolean {
return this.items.length >= this.maxSize;
}

get isEmpty(): boolean {
return this.items.length === 0;
}

/** Non-blocking enqueue. Returns false when the queue is full. */
tryEnqueue(item: T): boolean {
if (this.isFull) return false;
this.items.push(item);
return true;
}

/**
* Enqueue, waiting while the queue is full. This is the backpressure
* mechanism: fetch pauses here until the DB writer drains space.
*/
async enqueue(item: T): Promise<void> {
while (this.isFull) {
await new Promise<void>((resolve) => this.waiters.push(resolve));
}
this.items.push(item);
}

/** Remove up to `batchSize` items (FIFO). Returns [] when empty. */
dequeueBatch(batchSize: number): T[] {
if (batchSize <= 0 || this.items.length === 0) return [];
const batch = this.items.splice(0, Math.min(batchSize, this.items.length));
this.notifySpaceAvailable();
return batch;
}

/** Drain everything currently buffered. */
drain(): T[] {
return this.dequeueBatch(this.items.length);
}

clear(): void {
this.items = [];
this.notifySpaceAvailable();
}

private notifySpaceAvailable(): void {
const waiters = this.waiters.splice(0);
for (const w of waiters) w();
}
}
101 changes: 86 additions & 15 deletions src/workers/indexer/IndexerWorker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,18 +23,28 @@ import {
} from '../../services/Soroban/IndexerContractRegistry';
import { getSorobanNetwork } from '../../config/soroban';
import logger from '../../config/logger';
import { BoundedEventQueue } from './BoundedEventQueue';
import { InsertIndexedEventDTO } from '../../services/IndexedEventService';

const DEFAULT_POLL_INTERVAL_MS = parseInt(process.env.INDEXER_POLL_INTERVAL_MS || '15000', 10);
const EVENTS_PER_PAGE = parseInt(process.env.INDEXER_PAGE_SIZE || '200', 10);
const LAG_MONITOR_INTERVAL_MS = parseInt(
process.env.INDEXER_LAG_MONITOR_INTERVAL_MS || '15000',
10,
);
// Bounded buffer between fetch and DB write. When the queue is full the
// fetch loop pauses until the writer drains space (backpressure).
const MAX_QUEUE_SIZE = parseInt(process.env.INDEXER_MAX_QUEUE_SIZE || '1000', 10);
const WRITE_BATCH_SIZE = parseInt(process.env.INDEXER_WRITE_BATCH_SIZE || '100', 10);

export interface PollConfig {
pollIntervalMs?: number;
pageSize?: number;
lagMonitorIntervalMs?: number;
/** Max buffered decoded events between fetch and DB write. */
maxQueueSize?: number;
/** Max DB writes flushed per batch. */
writeBatchSize?: number;
}

export interface PollOutcome {
Expand All @@ -57,6 +67,9 @@ export class IndexerWorker {
private pollIntervalMs: number;
private pageSize: number;
private lagMonitorIntervalMs: number;
private maxQueueSize: number;
private writeBatchSize: number;
private queueDepth = 0;
private stopped = false;

constructor(
Expand All @@ -76,6 +89,8 @@ export class IndexerWorker {
this.pollIntervalMs = options.pollConfig?.pollIntervalMs ?? DEFAULT_POLL_INTERVAL_MS;
this.pageSize = options.pollConfig?.pageSize ?? EVENTS_PER_PAGE;
this.lagMonitorIntervalMs = options.pollConfig?.lagMonitorIntervalMs ?? LAG_MONITOR_INTERVAL_MS;
this.maxQueueSize = options.pollConfig?.maxQueueSize ?? MAX_QUEUE_SIZE;
this.writeBatchSize = options.pollConfig?.writeBatchSize ?? WRITE_BATCH_SIZE;

for (const contract of this.contracts) {
this.guards.set(contract.contractType, new IndexerReorgGuard());
Expand Down Expand Up @@ -187,6 +202,8 @@ export class IndexerWorker {

/**
* Drain remaining events via cursor pagination until caught up.
* Uses a bounded queue: when the buffer is full, the next fetch pauses
* until the DB writer drains space (backpressure).
*/
private async drainPages(
contract: IndexerContract,
Expand All @@ -199,6 +216,7 @@ export class IndexerWorker {
let errors = 0;
let cursor: string | null = null;
const seenCursors = new Set<string>();
const queue = new BoundedEventQueue<InsertIndexedEventDTO>(this.maxQueueSize);

let pageCursor = startCursor;
while (pageCursor && !this.stopped) {
Expand All @@ -211,12 +229,25 @@ export class IndexerWorker {
}
seenCursors.add(pageCursor);

// Backpressure: pause fetch while the queue is full, draining one
// write batch first so memory stays bounded.
while (queue.isFull && !this.stopped) {
logger.warn(
{ contract: contract.contractType, queueDepth: queue.size },
'Indexer write queue full; pausing fetch until DB drains',
);
const drained = await this.flushQueue(contract, queue);
processed += drained.processed;
errors += drained.errors;
}
if (this.stopped) break;

const page = await this.reader.fetchEvents({
contractIds: [contract.contractId],
cursor: pageCursor,
limit: this.pageSize,
});
const outcome = await this.persistProgress(contract, page, prevLedger);
const outcome = await this.persistProgress(contract, page, prevLedger, queue);
fetched += page.events.length;
processed += outcome.processed;
skipped += outcome.skipped;
Expand All @@ -225,17 +256,54 @@ export class IndexerWorker {
pageCursor = page.cursor || undefined;
}

// Flush any buffered writes before returning.
if (!queue.isEmpty) {
const drained = await this.flushQueue(contract, queue);
processed += drained.processed;
errors += drained.errors;
}
this.queueDepth = queue.size;

return { fetched, processed, skipped, errors, cursor };
}

/** Current buffered (unflushed) write depth. Useful for metrics/tests. */
getQueueDepth(): number {
return this.queueDepth;
}

private async flushQueue(
contract: IndexerContract,
queue: BoundedEventQueue<InsertIndexedEventDTO>,
): Promise<{ processed: number; errors: number }> {
let processed = 0;
let errors = 0;
while (!queue.isEmpty) {
const batch = queue.dequeueBatch(this.writeBatchSize);
for (const dto of batch) {
try {
await this.eventService.upsertEvent(dto);
processed += 1;
} catch (err) {
errors += 1;
logger.error({ contract: contract.contractType, err }, 'Failed to persist indexed event');
}
}
this.queueDepth = queue.size;
}
return { processed, errors };
}

private async persistProgress(
contract: IndexerContract,
page: EventsPage,
prevLedger: number,
sharedQueue?: BoundedEventQueue<InsertIndexedEventDTO>,
): Promise<{ processed: number; skipped: number; errors: number }> {
let processed = 0;
let skipped = 0;
let errors = 0;
const queue = sharedQueue ?? new BoundedEventQueue<InsertIndexedEventDTO>(this.maxQueueSize);

for (const event of page.events) {
const dto = contract.decoder.decode(event);
Expand All @@ -252,24 +320,27 @@ export class IndexerWorker {
);
continue;
}
try {
await this.eventService.upsertEvent(dto);
processed += 1;
} catch (err) {
errors += 1;
logger.error(
{
contract: contract.contractType,
eventId: event.id,
ledger: event.ledger,
txHash: event.txHash,
err,
},
'Failed to persist indexed event',
// Enqueue with backpressure: if the buffer is full, flush a batch
// first (pausing fetch) instead of growing memory unbounded.
if (queue.isFull) {
logger.warn(
{ contract: contract.contractType, queueDepth: queue.size },
'Indexer write queue full; pausing fetch until DB drains',
);
const drained = await this.flushQueue(contract, queue);
processed += drained.processed;
errors += drained.errors;
}
queue.tryEnqueue(dto);
this.queueDepth = queue.size;
}

// Flush what this page buffered so progress/cursor only advances for
// durable writes.
const flushed = await this.flushQueue(contract, queue);
processed += flushed.processed;
errors += flushed.errors;

const highestLedger = page.events.reduce((max, ev) => Math.max(max, ev.ledger), prevLedger);
if (processed > 0 || highestLedger > prevLedger) {
await this.indexerService.recordProgress(
Expand Down