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
72 changes: 71 additions & 1 deletion src/indexer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ import {
pruneOldTransfers,
} from "./db";
import { resolveLpPoolIds, loadKnownLpPools } from "./indexer/lp-shares";
import { pollParallel } from "./indexer/parallel";
import { pollParallel, type ParallelIo } from "./indexer/parallel";
import { processEventBatch, batchTotal } from "./indexer/batch";
import { tombstoneExpiredContracts } from "./indexer/tombstones";
import { createSourceSwitcherWithConfig, type SourceSwitcher } from "./indexer/sources";
Expand Down Expand Up @@ -216,6 +216,76 @@ export function _resetIndexerLoops(): void {
loops.clear();
}

/**
* Test-only: build one loop's isolated state with an injectable event source.
*
* The harness (#172) must drive the real `pollOnce`/`pollParallel` code — not
* raw `db` upserts — so a missing `network` tag fails the same way production
* would. `createLoopState` hard-wires a live RPC/Horizon switcher, so this
* builds the same state then swaps in the stub source and registers the loop
* (so `/status` sees it via `runningNetworks()`/`getAllIndexerStats()`).
*/
export function _createLoopForTesting(
network: Network,
overrides?: Partial<LoopState>,
): LoopState {
const loop = createLoopState(network);
if (overrides) Object.assign(loop, overrides);
loop.allContractIds = [...new Set([...loop.sacContractIds, ...loop.nftContractIds])];
loops.set(network, loop);
return loop;
}

/** Test-only: drive one real indexer poll step for `loop`. */
export function _pollOnceForTesting(
loop: LoopState,
fromLedger: number,
latestLedger: number,
): Promise<number> {
return pollOnce(loop, fromLedger, latestLedger);
}

/**
* Test-only: mirror the `ingestWindow` dispatch so the harness covers both the
* single-worker (`pollOnce`) and `INGEST_WORKERS > 1` (`pollParallel`) paths
* through the same contracts and batch pipeline production uses.
*
* `fetchEventsOverride` replaces only the fetch half of the parallel `io`
* seam; `processBatch` always stays the real `processEventBatch`, so the
* sharded path cannot drift from the single path (#219).
*/
export async function _pollWindowForTesting(
loop: LoopState,
fromLedger: number,
targetLedger: number,
workerCount: number = INGEST_WORKERS,
fetchEventsOverride?: ParallelIo["fetchEvents"],
): Promise<number> {
const net = loop.network;
if (workerCount > 1 && loop.allContractIds.length > 1) {
const { totalInserted, highestLedger } = await pollParallel(
loop.allContractIds,
fromLedger,
targetLedger,
BATCH_SIZE,
workerCount,
net,
{
fetchEvents:
fetchEventsOverride ??
((from, to, contractIds, limit) =>
loop.sourceSwitcher.fetchEvents(from, to, contractIds, limit)),
processBatch: (events) =>
processEventBatch(events, { network: net, knownLpPools: loop.knownLpPools }),
},
);
loop.totalIndexed += totalInserted;
recordLedgerProgress(net, fromLedger, highestLedger);
return highestLedger;
}
return pollOnce(loop, fromLedger, targetLedger);
}

/**
* Record one cursor advance for `network`.
*
Expand Down
Loading
Loading