cryo-indexer¶
The cryo-indexer is the primary execution layer indexer for the Gnosis Analytics pipeline. It extracts blockchain data from an RPC node using the Cryo binary (a high-performance Rust data extraction tool) and loads it into ClickHouse.
Purpose¶
cryo-indexer handles the full lifecycle of execution layer data acquisition:
- Real-time indexing of new blocks as they are produced
- Historical bulk loading of past block ranges
- Automatic recovery from failures and incomplete ranges
- Validation of data integrity
Architecture¶
Blockchain RPC --> Cryo Binary --> Parquet Files --> ClickHouse
|
Worker Process
|
State Manager
|
indexing_state table
The indexer orchestrates the Cryo binary to extract data from the RPC node into intermediate Parquet files, which are then loaded into ClickHouse. A single indexing_state table in ClickHouse serves as the source of truth for all processing state.
Key Design Decisions¶
- Blocks first -- Block headers are always processed before other datasets, because downstream datasets require valid block timestamps.
- 1000-block chunks -- All processing uses fixed 1000-block ranges for predictable resource usage.
- Atomic ranges -- A range is either fully completed or marked as failed. No partial writes.
- Single state table -- The
indexing_statetable tracks all datasets and ranges in one place.
Datasets¶
cryo-indexer supports 11 datasets organized into four indexing modes:
Indexing Modes¶
| Mode | Datasets Included | Use Case | Approx. Storage per 1M Blocks |
|---|---|---|---|
| minimal (default) | blocks, transactions, logs | Standard DeFi/DApp analysis | ~50 GB |
| extra | contracts, native_transfers, traces | Contract and trace analysis | ~100 GB |
| diffs | balance_diffs, code_diffs, nonce_diffs, storage_diffs | State change tracking | ~200 GB |
| full | All 11 datasets | Complete blockchain analysis | ~500 GB |
| custom | User-defined | Tailored to specific needs | Variable |
Dataset Reference¶
| Dataset | Description | ClickHouse Table |
|---|---|---|
blocks | Block headers, timestamps, gas, withdrawals root | blocks |
transactions | Transaction data including gas, value, status, input | transactions |
logs | Smart contract event log emissions | logs |
contracts | Contract creation events | contracts |
native_transfers | xDAI/ETH native token transfers | native_transfers |
traces | Internal transaction execution traces | traces |
balance_diffs | Account balance state changes per block | balance_diffs |
code_diffs | Smart contract bytecode changes | code_diffs |
nonce_diffs | Account nonce changes | nonce_diffs |
storage_diffs | Contract storage slot changes | storage_diffs |
withdrawals | Validator withdrawals (auto-populated with blocks) | withdrawals |
Note
The withdrawals dataset is automatically populated whenever blocks are processed. It does not require a separate extraction step.
Operating and recovering¶
How it runs¶
One image and no command: the entrypoint dispatches on OPERATION ∈ continuous | historical | maintain | auto-maintain | validate. DATASETS is honoured only when MODE=custom; MODE=minimal silently ignores it and gives you blocks, transactions, logs. All state is <db>.indexing_state; the pod holds nothing.
Three continuous writers, one per database: Gnosis (execution, CONFIRMATION_BLOCKS=720, MODE=custom with six datasets — blocks, transactions, logs, contracts, native_transfers, traces), Gnosis live (execution_live, CONFIRMATION_BLOCKS=6, two-day TTL) and Celo (celo_execution, MODE=minimal). Each chain also has an auto-maintain CronJob, seven times a day, that re-extracts recent failed or missing ranges. Cryo itself and its two patches live in cryo-base; odd Cryo behaviour starts there.
Single writer per database, Recreate, one replica — enforced by convention only. The target tables do not dedupe re-inserted rows, and maintain DELETEs a range before it claims it.
Health¶
Config failures are narrow: the entrypoint exits 1 if ETH_RPC_URL or CLICKHOUSE_HOST is empty. A failed secret sync shows as a container-config error, not a crash loop.
The authoritative liveness check needs only warehouse access:
SELECT dataset, max(end_block) AS highest_completed, max(created_at) AS last_write
FROM execution.indexing_state WHERE status = 'completed' GROUP BY dataset ORDER BY dataset;
last_write older than a few minutes means the writer is gone (POLL_INTERVAL=60, BATCH_SIZE=100). Repeat for execution_live and celo_execution. Head lag is ~66 minutes on execution (720 confirmations at 5 s) and ~1 minute on execution_live; neither is a stall.
A pod that is alive but wedged: restart the workload. Never above one replica and never a rolling update.
Detecting a gap¶
Never judge completeness from blocks row counts
On Celo, count() and max - min read 99.98% complete while ~993,000 blocks had a blocks row, no transactions or logs rows, and no indexing_state row at all — so nothing retried them and no alert fired. Nothing alerts on "has a blocks row but no transactions row"; only this check finds it.
Per-dataset batch coverage, run once per dataset:
WITH expected AS (SELECT toUInt32(⟨first_block⟩ + number * 100) AS s FROM numbers(⟨n_batches⟩))
SELECT s AS missing_start_block FROM expected
WHERE s NOT IN (SELECT DISTINCT start_block FROM execution.indexing_state WHERE dataset = 'transactions')
ORDER BY s;
Repeat for blocks, logs, and on execution also contracts, native_transfers, traces. Derive the bounds from the chain; do not copy numbers. A hole in one dataset does not imply a hole in the others.
Is a transaction-less batch a real hole or a quiet chain? Gas burned means transactions existed:
SELECT intDiv(block_number,100)*100 AS batch, count() AS blocks, sum(gas_used) AS gas
FROM execution.blocks WHERE block_number >= ⟨lo⟩ AND block_number < ⟨hi⟩
GROUP BY batch HAVING gas > 0 ORDER BY batch;
OPERATION=validate with a block range writes nothing and exits non-zero when gaps exist, so it can run inside the live pod. Also read the continuous log for max retries, All datasets exhausted, Moving to next range — the loop gives up on a dataset after MAX_RETRIES and moves the pointer forward, leaving a hole behind.
Repairing¶
Recent gaps: auto-maintain. Its reach is AUTO_MAINTAIN_LOOKBACK_HOURS × 720 blocks, and 720 blocks/hour is hardcoded for Gnosis' 5 s blocks — so 48 h is real on Gnosis and only ~9.6 h on Celo. Anything older is permanently out of its range. Trigger an extra run by cloning its CronJob (one-shot jobs).
auto-maintain is not read-only beside the live indexer
It DELETEs a range and then attempts the claim, so a claim conflict is reported as "Skipped" after the delete has already run. Prefer it for recent gaps, but know what it does.
Older gaps: a scoped maintain. It selects every non-completed range (processing included) and DELETEs it before re-extracting, with no claim, so the live writer for that database must be stopped first — the stop sequence and the one-shot recipe are on the one-shot jobs page. Chunk to a few hundred thousand blocks per job.
Two absolute limits, neither guarded in code
Never MODE=full — it expands to all ten datasets and starts a 41.9M-block backfill. Never 0/0 bounds — the unscoped range query self-joins ~1.3M rows and OOMs the warehouse.
Gnosis needs MODE=custom with all six live datasets — MODE=minimal repairs three of six and leaves the rest holed. Celo runs MODE=minimal live, so there minimal is the whole set; its RPC settings are 20× more aggressive than Gnosis' — keep them.
Corrupt range¶
Bookkeeping first — this is the step everyone misses
A range whose latest status is completed is invisible to every repair path. Re-open it by hand, using the exact stored bounds (a different start/end creates a phantom range):
SELECT dataset, start_block, end_block, argMax(status, created_at) AS st
FROM execution.indexing_state
WHERE start_block >= ⟨lo⟩ AND end_block <= ⟨hi⟩
GROUP BY dataset, start_block, end_block ORDER BY dataset, start_block;
INSERT INTO execution.indexing_state
(dataset, start_block, end_block, status, error_message, attempt_count)
VALUES
('blocks', ⟨lo⟩, ⟨hi⟩, 'failed', 'manual re-open: bad RPC data', 1),
('transactions', ⟨lo⟩, ⟨hi⟩, 'failed', 'manual re-open: bad RPC data', 1),
('logs', ⟨lo⟩, ⟨hi⟩, 'failed', 'manual re-open: bad RPC data', 1);
Then the scoped maintain above.
Order matters: blocks must be correct before transactions / logs. The other datasets derive block_timestamp by joining blocks and hard-fail if any block has an invalid timestamp. That failure logs as CRITICAL: Cannot add timestamps for <table>! ... Process blocks first! on transactions, which points diagnosis at the wrong dataset. Garbage timestamps, with the indexer's own predicate:
SELECT count() FROM execution.blocks FINAL
WHERE block_number >= ⟨lo⟩ AND block_number < ⟨hi⟩
AND (timestamp IS NULL OR timestamp = 0 OR toDateTime(timestamp) <= toDateTime('1971-01-01'));
Cold partitions need image 19e81f5 or later. Before that, the delete removed the rows but not the Keeper dedup block-id, so the re-insert was silently refused while the client saw written_rows=N. The only trace is system.part_log.error = 389.
Never hand-DELETE from a data table
maintain's own delete is what keeps the withdrawals side-table in step with blocks.
Verify:
SELECT count() raw, countDistinct(block_number) d, max(block_number)-min(block_number)+1 span
FROM execution.blocks WHERE block_number >= ⟨lo⟩ AND block_number < ⟨hi⟩; -- want raw = d = span
SELECT status, rows_indexed FROM execution.indexing_state
WHERE dataset='transactions' AND start_block=⟨lo⟩ AND end_block=⟨hi⟩
ORDER BY created_at DESC LIMIT 1; -- want 'completed' with rows_indexed > 0
A completed with rows_indexed = 0 is the silent-failure shape auto-maintain hunts.
Restart¶
Killing mid-range is safe for correctness, but SIGTERM exits immediately without finishing or failing the in-flight range — it stays processing until STUCK_RANGE_TIMEOUT_HOURS (2). Until then it is invisible to gap detection, and find_gaps deliberately clamps below anything still processing so a repair cannot stomp a live range. If you cannot wait two hours before a scoped repair over that range, re-open it as failed with the INSERT above.
A durable stop or start is a replica change in the stack, applied (Deployment); a live scale is [drift]. Restart resumes from the minimum across datasets of max(end_block) WHERE status='completed', aligned up to a batch boundary. On celo_execution, START_BLOCK is a floor, not a resume point.
Then dbt¶
A raw repair does not fix the decode layer — decode models are append with an embedded watermark and cannot see anything backfilled below it. Go to dbt reprocess. Do not run the daily microbatch runner: it only advances watermarks and produces a completely green run that fixes nothing.
Expected noise
- A uniform ~80 s
auto-maintainruntime means "nothing in window", not "healthy". No data found for <dataset> in blocks X-Y— legitimately empty ranges for contracts, native_transfers and traces.blocks ... not fully visible yet (0/100 valid)then succeeding on retry — read-after-write lag on SharedMergeTree, not bad data. Onlyreason='invalid'is a real fault.find_gapsskips any gap smaller than 100 blocks. "No gaps found" does not mean no missing blocks.Skipped: N (being processed by another worker)from auto-maintain — a claim conflict, benign.- The Celo chain-lag alert cannot fire during an RPC outage; the metric goes stale rather than growing. Absence of that alert is not evidence of health.
Internal runbook
runbooks/20-cryo-indexer.md — private repository; carries the cluster-specific commands for this page.
Redeploying¶
On an image roll the continuous writers restart (the auto-maintain CronJob picks the image up at its next slot). Apply outside the two hours before an auto-maintain slot, in an idle gap between ranges. The decisive proof is grid coverage over the restart window whole and no orphan range. Procedure: Redeploying a service.
Configuration¶
Required Settings¶
| Variable | Description |
|---|---|
ETH_RPC_URL | Blockchain RPC endpoint URL |
CLICKHOUSE_HOST | ClickHouse server hostname |
CLICKHOUSE_PASSWORD | ClickHouse authentication password |
Core Settings¶
| Variable | Default | Description |
|---|---|---|
NETWORK_NAME | ethereum | Network name passed to Cryo |
CLICKHOUSE_USER | default | ClickHouse username |
CLICKHOUSE_DATABASE | blockchain | Target database name (auto-created) |
CLICKHOUSE_PORT | 8443 | ClickHouse HTTP port |
CLICKHOUSE_SECURE | true | Use HTTPS for ClickHouse connection |
Operation Settings¶
| Variable | Default | Description |
|---|---|---|
OPERATION | continuous | continuous, historical, maintain, auto-maintain, validate |
MODE | minimal | minimal, extra, diffs, full, custom. Never full in production |
DATASETS | (derived from MODE) | Comma-separated dataset list — honoured only with MODE=custom |
START_BLOCK | 0 | Starting block number (a floor for continuous) |
END_BLOCK | 0 | Ending block number (0 = chain tip). Never 0/0 for maintain |
AUTO_MAINTAIN_LOOKBACK_HOURS | — | Reach of auto-maintain, converted at 720 blocks/hour regardless of chain |
STUCK_RANGE_TIMEOUT_HOURS | 2 | Age at which a processing range becomes eligible for repair |
Performance Settings¶
The production deployment overrides several code defaults; the values that explain healthy numbers are listed.
| Variable | Code default | Production | Description |
|---|---|---|---|
WORKERS | 1 | 1 | Parallel workers (use 4–16 for historical) |
BATCH_SIZE | 100 | 100 | Blocks per processing batch and per indexing_state range |
MAX_RETRIES | 3 | — | Retries before the continuous loop moves past a range |
REQUESTS_PER_SECOND | 20 | per chain | RPC request rate limit (Celo is set ~20× higher) |
MAX_CONCURRENT_REQUESTS | 2 | per chain | Maximum concurrent RPC requests |
CRYO_TIMEOUT | 600 | — | Cryo command timeout in seconds |
CONFIRMATION_BLOCKS | 12 | 720 / 6 | Blocks behind head (execution / execution_live) |
POLL_INTERVAL | 10 | 60 | Seconds between chain tip polls |
State Management¶
All indexing state is tracked in the indexing_state table:
-- Composite key: (mode, dataset, start_block, end_block)
-- Status flow: pending -> processing -> completed | failed
| Field | Description |
|---|---|
mode | Indexing mode that created this range |
dataset | Dataset name (blocks, transactions, etc.) |
start_block, end_block | Block range boundaries |
status | Current state: pending, processing, completed, failed |
worker_id | ID of the worker processing this range |
attempt_count | Number of processing attempts |
rows_indexed | Number of rows inserted |
error_message | Error details if status is failed |
The table is append-only: the current status of a range is its latest row (argMax(status, created_at)). A range left in processing by an unclean stop becomes eligible for repair only after STUCK_RANGE_TIMEOUT_HOURS; it is not reset on startup.
ClickHouse Table Schemas¶
All tables are stored in the execution database.
Table: execution.blocks
Engine: ReplacingMergeTree(insert_version) PARTITION BY: toStartOfMonth(block_timestamp) ORDER BY: (block_number)
| Column | Type | Notes |
|---|---|---|
block_number | UInt32 (Nullable) | Block height |
block_hash | String (Nullable) | Block hash |
parent_hash | String (Nullable) | Parent block hash |
author | String (Nullable) | Block author / miner |
state_root | String (Nullable) | State trie root hash |
gas_used | UInt64 (Nullable) | Total gas used in block |
gas_limit | UInt64 (Nullable) | Block gas limit |
timestamp | UInt32 (Nullable) | Unix timestamp |
size | UInt64 (Nullable) | Block size in bytes |
base_fee_per_gas | UInt64 (Nullable) | EIP-1559 base fee |
withdrawals_root | String (Nullable) | Withdrawals trie root |
chain_id | UInt64 (Nullable) | Chain identifier |
block_timestamp | DateTime | Materialized from timestamp |
insert_version | UInt64 | Materialized; deduplication version |
Table: execution.transactions
Engine: ReplacingMergeTree(insert_version) PARTITION BY: toStartOfMonth(block_timestamp) ORDER BY: (block_number, transaction_index)
| Column | Type | Notes |
|---|---|---|
block_number | UInt32 (Nullable) | Block height |
transaction_index | UInt64 (Nullable) | Position within the block |
transaction_hash | String (Nullable) | Transaction hash |
nonce | UInt64 (Nullable) | Sender nonce |
from_address | String (Nullable) | Sender address |
to_address | String (Nullable) | Recipient address |
value_string | String (Nullable) | Transfer value (string) |
value_f64 | Float64 (Nullable) | Transfer value (float) |
input | String (Nullable) | Calldata |
gas_limit | UInt64 (Nullable) | Gas limit |
gas_used | UInt64 (Nullable) | Gas consumed |
gas_price | UInt64 (Nullable) | Gas price in wei |
transaction_type | UInt32 (Nullable) | EIP-2718 transaction type |
max_priority_fee_per_gas | UInt64 (Nullable) | EIP-1559 priority fee |
max_fee_per_gas | UInt64 (Nullable) | EIP-1559 max fee |
success | UInt8 (Nullable) | 1 = success, 0 = revert |
chain_id | UInt64 (Nullable) | Chain identifier |
block_timestamp | DateTime64(0, 'UTC') | Block timestamp |
insert_version | UInt64 | Materialized; deduplication version |
Table: execution.logs
Engine: ReplacingMergeTree(insert_version) ORDER BY: (block_number, transaction_index, log_index)
| Column | Type | Notes |
|---|---|---|
block_number | UInt32 (Nullable) | Block height |
log_index | UInt32 (Nullable) | Log position within the block |
transaction_hash | String (Nullable) | Parent transaction hash |
address | String (Nullable) | Emitting contract address |
topic0 | String (Nullable) | Event signature hash |
topic1 | String (Nullable) | Indexed parameter 1 |
topic2 | String (Nullable) | Indexed parameter 2 |
topic3 | String (Nullable) | Indexed parameter 3 |
data | String (Nullable) | Non-indexed event data |
chain_id | UInt64 (Nullable) | Chain identifier |
block_timestamp | DateTime64(0, 'UTC') | Block timestamp |
insert_version | UInt64 | Materialized; deduplication version |
Table: execution.traces
Engine: ReplacingMergeTree(insert_version) ORDER BY: (block_number, transaction_index, trace_address)
| Column | Type | Notes |
|---|---|---|
action_from | String (Nullable) | Caller address |
action_to | String (Nullable) | Callee address |
action_value | String (Nullable) | Value transferred |
action_gas | String (Nullable) | Gas provided |
action_input | String (Nullable) | Input data |
action_call_type | String (Nullable) | Call type (call, delegatecall, etc.) |
action_type | String (Nullable) | Trace action type |
result_gas_used | UInt32 (Nullable) | Gas consumed by trace |
result_output | String (Nullable) | Return data |
result_code | String (Nullable) | Deployed bytecode (create traces) |
result_address | String (Nullable) | Created contract address |
trace_address | String (Nullable) | Position in the trace tree |
subtraces | UInt32 (Nullable) | Number of child traces |
transaction_index | UInt32 (Nullable) | Transaction position in block |
block_number | UInt32 (Nullable) | Block height |
error | String (Nullable) | Error message if trace reverted |
chain_id | UInt64 (Nullable) | Chain identifier |
block_timestamp | DateTime64 | Block timestamp |
insert_version | UInt64 | Materialized; deduplication version |
Table: execution.contracts
Engine: ReplacingMergeTree(insert_version) ORDER BY: (block_number, create_index)
| Column | Type | Notes |
|---|---|---|
block_number | UInt32 (Nullable) | Block height |
contract_address | String (Nullable) | Deployed contract address |
deployer | String (Nullable) | EOA that initiated the deploy |
factory | String (Nullable) | Factory contract (if created via CREATE2) |
code | String (Nullable) | Deployed bytecode |
code_hash | String (Nullable) | Keccak256 of bytecode |
chain_id | UInt64 (Nullable) | Chain identifier |
block_timestamp | DateTime64 | Block timestamp |
insert_version | UInt64 | Materialized; deduplication version |
Table: execution.native_transfers
Engine: ReplacingMergeTree(insert_version) ORDER BY: (block_number, transfer_index)
| Column | Type | Notes |
|---|---|---|
block_number | UInt32 (Nullable) | Block height |
transfer_index | UInt32 (Nullable) | Transfer position in block |
transaction_hash | String (Nullable) | Parent transaction hash |
from_address | String (Nullable) | Sender address |
to_address | String (Nullable) | Recipient address |
value_string | String (Nullable) | Transfer value (string) |
value_f64 | Float64 (Nullable) | Transfer value (float) |
chain_id | UInt64 | Chain identifier |
block_timestamp | DateTime64 | Block timestamp |
insert_version | UInt64 | Materialized; deduplication version |
Table: execution.balance_diffs
Engine: ReplacingMergeTree(insert_version) ORDER BY: (block_number, transaction_index, address)
| Column | Type | Notes |
|---|---|---|
block_number | UInt32 | Block height |
transaction_index | UInt32 | Transaction position in block |
address | String | Account address |
from_value_f64 | Float64 | Balance before the change |
to_value_f64 | Float64 | Balance after the change |
block_timestamp | DateTime64 | Block timestamp |
insert_version | UInt64 | Materialized; deduplication version |
Table: execution.withdrawals
Engine: ReplacingMergeTree(insert_version) ORDER BY: (block_number, withdrawal_index)
| Column | Type | Notes |
|---|---|---|
block_number | UInt32 | Block height |
withdrawal_index | String | Withdrawal sequence index |
validator_index | String | Validator that triggered the withdrawal |
address | String | Recipient address |
amount | String | Withdrawal amount |
block_timestamp | DateTime64 | Block timestamp |
insert_version | UInt64 | Materialized; deduplication version |
Table: execution.indexing_state
Engine: ReplacingMergeTree(insert_version) ORDER BY: (mode, dataset, start_block)
| Column | Type | Notes |
|---|---|---|
mode | String | Indexing mode (minimal, extra, etc.) |
dataset | String | Dataset name (blocks, transactions, etc.) |
start_block | UInt32 | Range start block |
end_block | UInt32 | Range end block |
status | String | pending / processing / completed / failed |
worker_id | String | ID of the processing worker |
attempt_count | UInt8 | Number of processing attempts |
created_at | DateTime | Range creation timestamp |
completed_at | DateTime (Nullable) | Completion timestamp |
rows_indexed | UInt64 (Nullable) | Number of rows inserted |
error_message | String (Nullable) | Error details if failed |