From fa30d81b781cba80c5110904f8ccfd9dfb2280cc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E2=80=9CEnchanterme=E2=80=9D?= <“encountermehere@gmail.com”> Date: Wed, 30 Sep 2026 06:05:16 +0100 Subject: [PATCH] feat(indexer): add bounded write queue with fetch backpressure Add BoundedEventQueue between event fetch and Postgres writes. When the buffer is full, fetch pauses and drains a write batch first so bursts of on-chain events cannot pile up unbounded memory. Configurable via INDEXER_MAX_QUEUE_SIZE and INDEXER_WRITE_BATCH_SIZE. --- src/workers/indexer/BoundedEventQueue.ts | 74 +++++++++++++++++ src/workers/indexer/IndexerWorker.ts | 101 +++++++++++++++++++---- 2 files changed, 160 insertions(+), 15 deletions(-) create mode 100644 src/workers/indexer/BoundedEventQueue.ts diff --git a/src/workers/indexer/BoundedEventQueue.ts b/src/workers/indexer/BoundedEventQueue.ts new file mode 100644 index 0000000..b49507f --- /dev/null +++ b/src/workers/indexer/BoundedEventQueue.ts @@ -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 { + 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 { + while (this.isFull) { + await new Promise((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(); + } +} diff --git a/src/workers/indexer/IndexerWorker.ts b/src/workers/indexer/IndexerWorker.ts index c072ffb..bbbc620 100644 --- a/src/workers/indexer/IndexerWorker.ts +++ b/src/workers/indexer/IndexerWorker.ts @@ -23,6 +23,8 @@ 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); @@ -30,11 +32,19 @@ 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 { @@ -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( @@ -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()); @@ -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, @@ -199,6 +216,7 @@ export class IndexerWorker { let errors = 0; let cursor: string | null = null; const seenCursors = new Set(); + const queue = new BoundedEventQueue(this.maxQueueSize); let pageCursor = startCursor; while (pageCursor && !this.stopped) { @@ -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; @@ -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, + ): 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, ): Promise<{ processed: number; skipped: number; errors: number }> { let processed = 0; let skipped = 0; let errors = 0; + const queue = sharedQueue ?? new BoundedEventQueue(this.maxQueueSize); for (const event of page.events) { const dto = contract.decoder.decode(event); @@ -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(