A TypeScript library for reliable EVM indexing with PostgreSQL.
🇬🇧 English | 🇷🇺 Русский
Voryn helps you build indexers that read blocks from EVM RPC, store normalized fetched data, commit chain progress in strict order, and run your application logic on transactions and events.
The library handles the boring but critical infrastructure: block queues, retries, cursors, singleton locks, reorg protection, retention, and metrics. You write business logic on top of committed data.
Voryn is a good fit for teams building:
- backends for DeFi, NFT, payments, wallets, and on-chain analytics;
- event-driven services that react to contract logs;
- pipelines that load blocks, transactions, and events into PostgreSQL;
- application-owned custom indexers;
- multi-chain services where each chain needs isolated progress tracking.
Voryn focuses on durable PostgreSQL-backed indexing pipelines with application-owned storage and processing.
- Ingestion pipeline:
HeadWorkerenqueues blocks,FetchWorkerdownloads data, andSequencerWorkercommits only the strictN -> N+1sequence. - Reorg handling:
parentHashchecks, common ancestor lookup, and rollback of replaced fetched data. - Horizontal fetch scaling: multiple fetch workers can safely share one PostgreSQL-backed queue.
- Durable reactions:
EventReactionWorkerandTransactionReactionWorkerread committed streams by block position and maintain their own cursors. - Operational tools:
RetentionWorker,PipelineMetrics,BlockJobRecovery,ConsoleLogger, and a PostgreSQL schema helper. - Replaceable pieces: you can bring your own
BlockSource, logger, repositories, transaction manager, or leader lock.
flowchart LR
RPC["EVM RPC"] --> Head["HeadWorker"]
Head --> Jobs["block_jobs"]
Jobs --> Fetch["FetchWorker x N"]
Fetch --> Data["blocks / transactions / events"]
Data --> Sequencer["SequencerWorker"]
Jobs --> Sequencer
Sequencer --> Cursor["chain_cursor"]
Data --> Reactions["Reaction workers"]
Cursor --> Reactions
Fetchcan be scaled horizontally. It writes normalizedblocks,transactions, andevents.Sequencervalidates order through block hashes, advanceschain_cursor, and marks the matchingblock_jobsrows as committed.Head,Sequencer,Retention, and Reaction workers run as singleton processes throughLeaderLock.
npm install @drillcoder/vorynThe package is published as ESM and uses ethers v6 and pg.
Voryn expects the base PostgreSQL schema from the SQL file included in the npm package.
The physical file path after installation is:
node_modules/@drillcoder/voryn/dist/sql/postgres-schema.sqlThe same file is also exposed as a package subpath:
@drillcoder/voryn/sql/postgres-schema.sqlYou can apply the SQL with your own migration flow:
psql "$DATABASE_URL" -f node_modules/@drillcoder/voryn/dist/sql/postgres-schema.sqlOr use the built-in helper:
import { Pool } from "pg";
import { ConsoleLogger, applySqlFileToPostgresDb } from "@drillcoder/voryn";
const pool = new Pool({ connectionString: process.env.DATABASE_URL });
const logger = new ConsoleLogger({ minLevel: "info" });
await applySqlFileToPostgresDb({
pool,
sqlFilePath: "node_modules/@drillcoder/voryn/dist/sql/postgres-schema.sql",
logger,
});
await pool.end();Full example: examples/db-apply-sql.ts
The minimal ingestion pipeline consists of head, fetch, and sequencer. In production, they are usually started as separate processes or containers.
import { FetchWorker, HeadWorker, SequencerWorker } from "@drillcoder/voryn";
const dbUrl = "postgres://user:pass@localhost:5432/voryn";
const rpcUrls = ["https://rpc.example.org", "https://fallback-rpc.example.org"];
const chainId = 1;
const logLevel = "info";
const headOptions = {
sourceConfig: {
network: { chainId, rpcUrls },
requestTimeoutMs: 5_000,
operationTimeoutMs: 60_000,
},
delayBetweenTicksMs: 1_000,
depthBlocks: 65_000,
// initialBlock: 20_000_000, // first block on a new chain cursor
logLevel,
dbUrl,
};
const fetchOptions = {
sourceConfig: {
network: { chainId, rpcUrls },
requestTimeoutMs: 30_000,
operationTimeoutMs: 60_000,
},
delayBetweenTicksMs: 100,
fetchBatchSize: 10,
fetchConcurrency: 2,
fetchClaimTtlMs: 125_000,
retryMaxAttempts: 10,
retryBaseDelayMs: 1_000,
retryMaxDelayMs: 10_000,
logLevel,
dbUrl,
};
const sequencerOptions = {
sourceConfig: {
network: { chainId, rpcUrls },
requestTimeoutMs: 5_000,
operationTimeoutMs: 60_000,
},
delayBetweenTicksMs: 100,
maxBlocksPerTick: 10,
logLevel,
dbUrl,
};
const head = await HeadWorker.create(headOptions);
const fetch = await FetchWorker.create(fetchOptions);
const sequencer = await SequencerWorker.create(sequencerOptions);
const handleSingletonFailure = (error: Error): void => {
console.error(error);
process.exitCode = 1;
};
head.onFailure(handleSingletonFailure);
sequencer.onFailure(handleSingletonFailure);
await Promise.all([
head.start(),
fetch.start(),
sequencer.start(),
]);Full worker examples:
Reaction workers read only committed data. Each workerName has its own persisted cursor, so handlers can be
restarted safely.
import type { EventReactionHandler, EventReactionWorkerOptions, ReactionHandlerResult } from "@drillcoder/voryn";
import { EventReactionWorker } from "@drillcoder/voryn";
const dbUrl = "postgres://user:pass@localhost:5432/voryn";
const logLevel = "info";
const handler: EventReactionHandler = async (event): Promise<ReactionHandlerResult> => {
console.info("event_received", {
blockNumber: event.blockNumber,
transactionHash: event.transactionHash,
logIndex: event.index,
address: event.address,
});
return event.index === 10 ? "processed" : "skipped";
};
const options: EventReactionWorkerOptions = {
chainId: 1,
workerName: "contract-events",
delayBetweenTicksMs: 500,
batchSize: 1000,
skipFlushInterval: 100,
confirmations: 12,
// initialBlock: 20_000_000, // first block when this worker has no cursor
logLevel,
dbUrl,
handler,
};
const worker = await EventReactionWorker.create(options);
worker.onFailure((error) => {
console.error(error);
process.exitCode = 1;
});
await worker.start();Handlers may return "processed" or "skipped". Processed items advance the worker cursor immediately. Skipped
items are safe to advance too, but their cursor writes are batched by skipFlushInterval and flushed at
the end of the tick or before rethrowing a handler error.
A handler can be called more than once for the same item if it fails before returning a result, if the worker stops before the cursor write is persisted, or if a reorg moves the item to a new committed block. Reorg rollback rewinds affected reaction cursors, but it cannot undo external side effects. Keep handlers idempotent.
initialBlock is optional for HeadWorker and both reaction workers. It is used only when their cursor is missing;
an existing cursor always resumes from its saved position. For HeadWorker, initialBlock must be within
[max(0, latestBlock - depthBlocks + 1), latestBlock]. The first head tick initializes the cursor; the next
enqueues blocks starting with initialBlock. An out-of-range value causes tick errors without creating a cursor or
jobs. Without initialBlock, a new HeadWorker cursor starts at the latest block and enqueues only later blocks.
A new reaction worker starts with items in the current committed block and skips earlier committed blocks.
Each reaction worker chooses its own non-negative safe integer confirmations. Its read boundary is the lower of committed
progress and lastEnqueuedBlock - confirmations. Retention does not wait for reaction workers, so configure:
retentionDepthBlocks > max reaction confirmations + maximum expected processing lag + operational reserve
Examples:
PipelineMetrics returns a pipeline snapshot: current RPC head, stage lag, data freshness, block job statuses, failed blocks, and reaction worker lag.
import { PipelineMetrics } from "@drillcoder/voryn";
const dbUrl = "postgres://user:pass@localhost:5432/voryn";
const metrics = await PipelineMetrics.create({
dbUrl,
sourceConfig: {
networks: [
{
chainId: 1,
rpcUrls: ["https://mainnet-rpc.example.org", "https://mainnet-fallback-rpc.example.org"],
},
{
chainId: 56,
rpcUrls: ["https://bsc-rpc.example.org", "https://bsc-fallback-rpc.example.org"],
},
],
requestTimeoutMs: 5_000,
operationTimeoutMs: 60_000,
},
});
const snapshot = await metrics.get();
const prometheusText = await metrics.getPrometheus();
await metrics.close();get()returns one aggregate snapshot with achainsarray.getPrometheus()returns one Prometheus text document for all configured chains. Serve it from your own/metricsendpoint.
Use BlockJobRecovery to manually put failed blocks back into processing.
Examples:
The RPC branch of sourceConfig makes Voryn create and own an internal RPC-pool-backed BlockSource. The pool pins every
block-source operation to one endpoint. Eligible transport or endpoint-data failures retry the whole operation on
another endpoint. Fetch job retries remain a separate pipeline-level recovery mechanism.
The internal adapter validates hashes, addresses, data fields, transaction indexes, and block number consistency.
To use another data source, implement the public BlockSource interface. For direct RPC pool access, use
@drillcoder/ethers-rpc-pool.
Version 1.1 intentionally changes the source API. For a single-chain worker, replace top-level chainId plus
rpcUrl/fallbackRpcUrl with sourceConfig: { network: { chainId, rpcUrls: [rpcUrl, fallbackRpcUrl] } }. Move
rpcRequestTimeoutMs to sourceConfig.requestTimeoutMs; sourceConfig.operationTimeoutMs is the optional whole-operation
deadline. A custom source uses sourceConfig: { chainId, source }, without network. Replace direct
EthersBlockSource usage with one of these branches; for a custom integration, implement BlockSource.
For PipelineMetrics, use either sourceConfig: { networks, requestTimeoutMs?, operationTimeoutMs? } or
sourceConfig: { chainIds, source }. RPC metrics derive chain IDs from networks; custom metrics require chainIds.
Import the supported public API from the package root, @drillcoder/voryn. Package-internal paths are not stable.
The public API is grouped by use case:
- pipeline workers:
HeadWorker,FetchWorker,SequencerWorker, andRetentionWorker, with a corresponding*Optionstype for eachcreate()factory; - reactions:
EventReactionWorker,TransactionReactionWorker, their*Optionsand handler types, and thePipelineEventandPipelineTransactionrecords delivered to handlers; - block sources:
BlockSourceand the chain data types needed to implement a custom source; - operations:
PipelineMetrics,PipelineMetricsResult,BlockJobRecovery, and their option and result types; - logging:
Logger,noopLogger,ConsoleLogger, and its configuration types; - PostgreSQL infrastructure: repository interfaces and implementations,
LeaderLock,TransactionManager, database executor types, and schema helpers.
The *DatabaseDependencies types describe the overrides accepted by factory options. They are useful when replacing
one or more PostgreSQL-backed dependencies. Other exported data types are the inputs and results of these public
extension points.
Commands are collected in dev/Makefile. Checks that require project dependencies should be run through the tools container:
make lint
make test
make buildFor the local development environment:
cp dev/.env.example dev/.env
make init
make ingestion-up