Explorer indexer: the observer's EVM-side sibling for Devnet 3 (live_txs, live_accounts, explorer_blocks, explorer_state under the dn3_ prefix; receipts, logs, balances, nonces, code; the three-deep reorg re-read; 7-day retention; igneum_getNodeInfo kept for the class floor)
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
parent
fe25f566df
commit
5a20afd47b
2 changed files with 313 additions and 0 deletions
275
tools/observer/explorer-indexer.mjs
Normal file
275
tools/observer/explorer-indexer.mjs
Normal file
|
|
@ -0,0 +1,275 @@
|
|||
// Igneum explorer indexer: the sibling of tools/observer/observer.mjs for the EVM side (8 October 2026, the explorer lane).
|
||||
// The observer writes every DAG block, the shard plans, the checkpoints and the certificates. This process reads the
|
||||
// execution layer of the same node (eth_getBlockByNumber with the transactions, eth_getBlockReceipts, eth_getBalance,
|
||||
// eth_getTransactionCount, eth_getCode, igneum_getNodeInfo) and writes what /tx/<hash>, /address/<addr> and the
|
||||
// transaction lists need. It never touches the observer's tables; it creates its own three next to them under the same
|
||||
// prefix, so one LIVE_TABLE_PREFIX names one chain on both sides. Node 22 or newer, no dependencies.
|
||||
//
|
||||
// node tools/observer/explorer-indexer.mjs run (env below)
|
||||
// node tools/observer/explorer-indexer.mjs --once one tick, then exit (the smoke run)
|
||||
//
|
||||
// Environment, every value optional:
|
||||
// DATABASE_URL Neon connection string (read from ~/.config/igneum/env when unset; never in the repo)
|
||||
// LIVE_TABLE_PREFIX table prefix, default dn3_ (Devnet 3, the same prefix the observer unit on the build box writes)
|
||||
// IGNEUM_EVM_RPC the node's execution JSON-RPC, default http://127.0.0.1:26850 (the Devnet 3 observer node's)
|
||||
// EXPLORER_BACKFILL chain blocks to index behind the tip on an empty table, default 0 = from the first block
|
||||
// EXPLORER_RETAIN_HOURS hours of transactions kept, default 168 (7 days); older rows are deleted once an hour
|
||||
// EXPLORER_BATCH chain blocks per round while catching up, default 25
|
||||
// EXPLORER_TICK_MS the idle tick, default 2000
|
||||
//
|
||||
// Tables (created on start if missing; the prefix is written as T below):
|
||||
// T live_txs one row per EVM transaction: hash, block_hash, block_number, tx_index, from_addr, to_addr, value_wei,
|
||||
// nonce, gas, gas_price, gas_used, effective_gas_price, status, contract_address, input_bytes, selector,
|
||||
// logs_count, logs (jsonb), igneum (jsonb: the receipt's igneum section), tx_type, ts_ms, received_at
|
||||
// T live_accounts one row per address seen: balance_wei, nonce, code_bytes, is_contract, sent, received, first_seen_number,
|
||||
// last_seen_number, miner (mined a block in the observer's table), refreshed_at
|
||||
// T explorer_blocks one row per indexed chain block: number, hash, tx_count, gas_used, ts_ms, indexed_at (the reorg check
|
||||
// re-reads the newest three numbers every tick and re-indexes a number whose hash moved)
|
||||
// T explorer_state one row: indexed_number, tip_number, chain_id, node_info (igneum_getNodeInfo: params, digest, genesis),
|
||||
// txs_total, accounts_total, backlog, blocks_per_s, last_error, started_at, updated_at
|
||||
//
|
||||
// Reorgs: Devnet 3's rule is no reorg over depth 3 (release-0.3.22 gates), so three numbers are re-read each tick; a deeper
|
||||
// one shows as a block hash in live_txs that the observer's live_blocks no longer calls a chain block, which the API reports.
|
||||
|
||||
import { readFileSync } from 'node:fs';
|
||||
import { homedir } from 'node:os';
|
||||
|
||||
const ONCE = process.argv.includes('--once');
|
||||
const T = (process.env.LIVE_TABLE_PREFIX === undefined ? 'dn3_' : process.env.LIVE_TABLE_PREFIX).replace(/[^a-z0-9_]/gi, '');
|
||||
const EVM = process.env.IGNEUM_EVM_RPC || 'http://127.0.0.1:26850';
|
||||
const BACKFILL = Math.max(0, Number(process.env.EXPLORER_BACKFILL) || 0);
|
||||
const RETAIN_H = Math.max(1, Number(process.env.EXPLORER_RETAIN_HOURS) || 168);
|
||||
const BATCH = Math.max(1, Math.min(100, Number(process.env.EXPLORER_BATCH) || 25));
|
||||
const TICK_MS = Math.max(500, Number(process.env.EXPLORER_TICK_MS) || 2000);
|
||||
const ACCOUNTS_PER_TICK = 20;
|
||||
const RPC_CONCURRENCY = 6;
|
||||
|
||||
// ---- pure helpers (tested by explorer-indexer.test.mjs) --------------------------------------------------------------
|
||||
export const hexInt = h => (h === null || h === undefined ? null : Number(BigInt(h)));
|
||||
export const hexBig = h => (h === null || h === undefined ? null : BigInt(h).toString());
|
||||
export const lower = a => (a ? String(a).toLowerCase() : null);
|
||||
|
||||
/** One live_txs row from the block's transaction object and its receipt (both as the node's JSON-RPC returns them). */
|
||||
export function txRow(tx, rc, block) {
|
||||
const input = String(tx.input || '0x');
|
||||
return {
|
||||
hash: lower(tx.hash), block_hash: lower(block.hash), block_number: hexInt(block.number), tx_index: hexInt(tx.transactionIndex) ?? 0,
|
||||
from_addr: lower(tx.from), to_addr: lower(tx.to), value_wei: hexBig(tx.value) || '0', nonce: hexInt(tx.nonce), gas: hexInt(tx.gas),
|
||||
gas_price: hexBig(tx.gasPrice ?? tx.maxFeePerGas) , gas_used: rc ? hexInt(rc.gasUsed) : null, effective_gas_price: rc ? hexBig(rc.effectiveGasPrice) : null,
|
||||
status: rc && rc.status !== undefined && rc.status !== null ? hexInt(rc.status) : null, contract_address: rc ? lower(rc.contractAddress) : null,
|
||||
input_bytes: Math.max(0, (input.length - 2) / 2), selector: input.length >= 10 ? input.slice(0, 10).toLowerCase() : null,
|
||||
logs_count: rc && rc.logs ? rc.logs.length : 0,
|
||||
logs: rc && rc.logs ? rc.logs.map(l => ({ i: hexInt(l.logIndex), address: lower(l.address), topics: l.topics || [], data: l.data || '0x' })) : [],
|
||||
igneum: rc && rc.igneum ? rc.igneum : null, tx_type: hexInt(tx.type) ?? 0, ts_ms: hexInt(block.timestamp) === null ? null : hexInt(block.timestamp) * 1000,
|
||||
};
|
||||
}
|
||||
|
||||
/** The addresses a block touches, each with why (sender, receiver, contract, miner). */
|
||||
export function touched(rows, block) {
|
||||
const m = new Map();
|
||||
const add = (a, kind) => { if (!a) return; const k = lower(a); if (!m.has(k)) m.set(k, new Set()); m.get(k).add(kind); };
|
||||
for (const r of rows) { add(r.from_addr, 'sender'); add(r.to_addr, 'receiver'); add(r.contract_address, 'contract'); }
|
||||
if (block && block.igneum && Array.isArray(block.igneum.rewards)) for (const r of block.igneum.rewards) add(r.miner, 'miner');
|
||||
if (block && block.miner) add(block.miner, 'miner');
|
||||
return m;
|
||||
}
|
||||
|
||||
// ---- Neon over HTTP, the same pattern as the observer ---------------------------------------------------------------
|
||||
function envFile() {
|
||||
try { return Object.fromEntries(readFileSync(`${homedir()}/.config/igneum/env`, 'utf8').split('\n').filter(l => /^[A-Z_]+=/.test(l)).map(l => { const i = l.indexOf('='); return [l.slice(0, i), l.slice(i + 1).replace(/^"|"$/g, '')]; })); } catch { return {}; }
|
||||
}
|
||||
function neon(url) {
|
||||
if (!url) throw new Error('DATABASE_URL is not set');
|
||||
const host = new URL(url).hostname.replace('-pooler', '');
|
||||
return async (query, params = []) => {
|
||||
const r = await fetch(`https://${host}/sql`, { method: 'POST', headers: { 'Neon-Connection-String': url, 'Content-Type': 'application/json' }, body: JSON.stringify({ query, params }), signal: AbortSignal.timeout(30000) });
|
||||
const j = await r.json();
|
||||
if (!r.ok) throw new Error(j.message || JSON.stringify(j));
|
||||
return j.rows || [];
|
||||
};
|
||||
}
|
||||
|
||||
// ---- JSON-RPC with a small concurrency cap ----------------------------------------------------------------------------
|
||||
let inFlight = 0; const waiters = [];
|
||||
async function rpc(method, params = []) {
|
||||
if (inFlight >= RPC_CONCURRENCY) await new Promise(r => waiters.push(r));
|
||||
inFlight++;
|
||||
try {
|
||||
const r = await fetch(EVM, { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ jsonrpc: '2.0', id: 1, method, params }), signal: AbortSignal.timeout(15000) });
|
||||
const j = await r.json();
|
||||
if (j.error) throw new Error(`${method}: ${j.error.message || JSON.stringify(j.error)}`);
|
||||
return j.result;
|
||||
} finally { inFlight--; const w = waiters.shift(); if (w) w(); }
|
||||
}
|
||||
const hex = n => '0x' + Number(n).toString(16);
|
||||
|
||||
// ---- schema ----------------------------------------------------------------------------------------------------------
|
||||
async function schema(sql) {
|
||||
await sql(`CREATE TABLE IF NOT EXISTS ${T}live_txs (
|
||||
hash text PRIMARY KEY, block_hash text NOT NULL, block_number bigint NOT NULL, tx_index integer NOT NULL DEFAULT 0,
|
||||
from_addr text, to_addr text, value_wei text NOT NULL DEFAULT '0', nonce bigint, gas bigint, gas_price text, gas_used bigint,
|
||||
effective_gas_price text, status smallint, contract_address text, input_bytes integer NOT NULL DEFAULT 0, selector text,
|
||||
logs_count integer NOT NULL DEFAULT 0, logs jsonb NOT NULL DEFAULT '[]'::jsonb, igneum jsonb, tx_type smallint NOT NULL DEFAULT 0,
|
||||
ts_ms bigint, received_at timestamptz NOT NULL DEFAULT now())`);
|
||||
await sql(`CREATE INDEX IF NOT EXISTS ${T}live_txs_number ON ${T}live_txs (block_number DESC, tx_index)`);
|
||||
await sql(`CREATE INDEX IF NOT EXISTS ${T}live_txs_from ON ${T}live_txs (from_addr, block_number DESC)`);
|
||||
await sql(`CREATE INDEX IF NOT EXISTS ${T}live_txs_to ON ${T}live_txs (to_addr, block_number DESC)`);
|
||||
await sql(`CREATE INDEX IF NOT EXISTS ${T}live_txs_block ON ${T}live_txs (block_hash)`);
|
||||
await sql(`CREATE INDEX IF NOT EXISTS ${T}live_txs_ts ON ${T}live_txs (ts_ms)`);
|
||||
await sql(`CREATE TABLE IF NOT EXISTS ${T}live_accounts (
|
||||
address text PRIMARY KEY, balance_wei text, nonce bigint, code_bytes integer, is_contract boolean, sent integer NOT NULL DEFAULT 0,
|
||||
received integer NOT NULL DEFAULT 0, first_seen_number bigint, last_seen_number bigint, miner boolean NOT NULL DEFAULT false,
|
||||
refreshed_at timestamptz, updated_at timestamptz NOT NULL DEFAULT now())`);
|
||||
await sql(`CREATE INDEX IF NOT EXISTS ${T}live_accounts_seen ON ${T}live_accounts (last_seen_number DESC)`);
|
||||
await sql(`CREATE TABLE IF NOT EXISTS ${T}explorer_blocks (number bigint PRIMARY KEY, hash text NOT NULL, tx_count integer NOT NULL DEFAULT 0,
|
||||
gas_used bigint, ts_ms bigint, indexed_at timestamptz NOT NULL DEFAULT now())`);
|
||||
await sql(`CREATE TABLE IF NOT EXISTS ${T}explorer_state (id integer PRIMARY KEY, indexed_number bigint, tip_number bigint, chain_id integer,
|
||||
node_info jsonb, txs_total bigint, accounts_total bigint, backlog integer, blocks_per_s double precision, last_error text,
|
||||
rpc text, started_at timestamptz, updated_at timestamptz)`);
|
||||
await sql(`INSERT INTO ${T}explorer_state (id, started_at, updated_at, rpc) VALUES (1, now(), now(), $1) ON CONFLICT (id) DO UPDATE SET started_at = now(), rpc = $1`, [EVM]);
|
||||
}
|
||||
|
||||
// ---- the work ----------------------------------------------------------------------------------------------------------
|
||||
const log = (...a) => console.log(new Date().toISOString(), ...a);
|
||||
const accountQueue = new Map(); // address -> {kinds:Set, number}
|
||||
const contractKnown = new Map(); // address -> is_contract (code fetched once)
|
||||
let txsPerMin = 0, blocksPerMin = 0, accountsPerMin = 0, lastError = null;
|
||||
|
||||
async function indexNumbers(sql, numbers) {
|
||||
const blocks = await Promise.all(numbers.map(n => rpc('eth_getBlockByNumber', [hex(n), true])));
|
||||
const receipts = await Promise.all(numbers.map((n, i) => (blocks[i] && blocks[i].transactions && blocks[i].transactions.length ? rpc('eth_getBlockReceipts', [hex(n)]).catch(() => null) : Promise.resolve([]))));
|
||||
const rows = []; const blockRows = [];
|
||||
for (let i = 0; i < numbers.length; i++) {
|
||||
const b = blocks[i]; if (!b) continue;
|
||||
const byHash = new Map((receipts[i] || []).map(r => [lower(r.transactionHash), r]));
|
||||
const txs = (b.transactions || []).map(t => txRow(t, byHash.get(lower(t.hash)) || null, b));
|
||||
rows.push(...txs);
|
||||
blockRows.push({ number: hexInt(b.number), hash: lower(b.hash), tx_count: txs.length, gas_used: hexInt(b.gasUsed), ts_ms: hexInt(b.timestamp) * 1000 });
|
||||
for (const [a, kinds] of touched(txs, b)) { const q = accountQueue.get(a) || { kinds: new Set(), number: 0 }; for (const k of kinds) q.kinds.add(k); q.number = Math.max(q.number, hexInt(b.number)); accountQueue.set(a, q); }
|
||||
}
|
||||
// a number re-indexed after a reorg drops its old rows first
|
||||
await sql(`DELETE FROM ${T}live_txs WHERE block_number = ANY($1::bigint[])`, [numbers]);
|
||||
for (let off = 0; off < rows.length; off += 200) {
|
||||
const chunk = rows.slice(off, off + 200); const vals = []; const params = [];
|
||||
chunk.forEach((r, k) => {
|
||||
const base = k * 22;
|
||||
vals.push(`(${Array.from({ length: 22 }, (_, j) => '$' + (base + j + 1)).join(',')})`);
|
||||
params.push(r.hash, r.block_hash, r.block_number, r.tx_index, r.from_addr, r.to_addr, r.value_wei, r.nonce, r.gas, r.gas_price, r.gas_used, r.effective_gas_price, r.status, r.contract_address, r.input_bytes, r.selector, r.logs_count, JSON.stringify(r.logs), r.igneum ? JSON.stringify(r.igneum) : null, r.tx_type, r.ts_ms, new Date().toISOString());
|
||||
});
|
||||
await sql(`INSERT INTO ${T}live_txs (hash, block_hash, block_number, tx_index, from_addr, to_addr, value_wei, nonce, gas, gas_price, gas_used, effective_gas_price, status, contract_address, input_bytes, selector, logs_count, logs, igneum, tx_type, ts_ms, received_at)
|
||||
VALUES ${vals.join(',')} ON CONFLICT (hash) DO UPDATE SET block_hash = EXCLUDED.block_hash, block_number = EXCLUDED.block_number, tx_index = EXCLUDED.tx_index, gas_used = EXCLUDED.gas_used,
|
||||
effective_gas_price = EXCLUDED.effective_gas_price, status = EXCLUDED.status, contract_address = EXCLUDED.contract_address, logs_count = EXCLUDED.logs_count, logs = EXCLUDED.logs, igneum = EXCLUDED.igneum, ts_ms = EXCLUDED.ts_ms`, params);
|
||||
}
|
||||
if (blockRows.length) {
|
||||
const vals = []; const params = [];
|
||||
blockRows.forEach((b, k) => { vals.push(`($${k * 5 + 1},$${k * 5 + 2},$${k * 5 + 3},$${k * 5 + 4},$${k * 5 + 5})`); params.push(b.number, b.hash, b.tx_count, b.gas_used, b.ts_ms); });
|
||||
await sql(`INSERT INTO ${T}explorer_blocks (number, hash, tx_count, gas_used, ts_ms) VALUES ${vals.join(',')} ON CONFLICT (number) DO UPDATE SET hash = EXCLUDED.hash, tx_count = EXCLUDED.tx_count, gas_used = EXCLUDED.gas_used, ts_ms = EXCLUDED.ts_ms, indexed_at = now()`, params);
|
||||
}
|
||||
// sent and received counts and first/last seen, in one statement per batch
|
||||
const seen = new Map();
|
||||
for (const r of rows) {
|
||||
if (r.from_addr) { const s = seen.get(r.from_addr) || { sent: 0, received: 0, min: r.block_number, max: r.block_number }; s.sent++; s.min = Math.min(s.min, r.block_number); s.max = Math.max(s.max, r.block_number); seen.set(r.from_addr, s); }
|
||||
if (r.to_addr) { const s = seen.get(r.to_addr) || { sent: 0, received: 0, min: r.block_number, max: r.block_number }; s.received++; s.min = Math.min(s.min, r.block_number); s.max = Math.max(s.max, r.block_number); seen.set(r.to_addr, s); }
|
||||
}
|
||||
const ents = [...seen.entries()];
|
||||
for (let off = 0; off < ents.length; off += 300) {
|
||||
const chunk = ents.slice(off, off + 300); const vals = []; const params = [];
|
||||
chunk.forEach(([a, s], k) => { vals.push(`($${k * 5 + 1},$${k * 5 + 2},$${k * 5 + 3},$${k * 5 + 4},$${k * 5 + 5})`); params.push(a, s.sent, s.received, s.min, s.max); });
|
||||
await sql(`INSERT INTO ${T}live_accounts (address, sent, received, first_seen_number, last_seen_number) VALUES ${vals.join(',')}
|
||||
ON CONFLICT (address) DO UPDATE SET sent = ${T}live_accounts.sent + EXCLUDED.sent, received = ${T}live_accounts.received + EXCLUDED.received,
|
||||
first_seen_number = LEAST(${T}live_accounts.first_seen_number, EXCLUDED.first_seen_number), last_seen_number = GREATEST(${T}live_accounts.last_seen_number, EXCLUDED.last_seen_number), updated_at = now()`, params);
|
||||
}
|
||||
txsPerMin += rows.length; blocksPerMin += blockRows.length;
|
||||
return blockRows;
|
||||
}
|
||||
|
||||
async function refreshAccounts(sql, max) {
|
||||
if (!accountQueue.size) return 0;
|
||||
const picks = [...accountQueue.entries()].sort((a, b) => b[1].number - a[1].number).slice(0, max);
|
||||
for (const [a] of picks) accountQueue.delete(a);
|
||||
const reads = await Promise.all(picks.map(async ([a, q]) => {
|
||||
try {
|
||||
const [bal, nonce] = await Promise.all([rpc('eth_getBalance', [a, 'latest']), rpc('eth_getTransactionCount', [a, 'latest'])]);
|
||||
let code = contractKnown.get(a);
|
||||
if (code === undefined || q.kinds.has('contract')) { const c = await rpc('eth_getCode', [a, 'latest']); code = Math.max(0, (String(c || '0x').length - 2) / 2); contractKnown.set(a, code); }
|
||||
return { a, bal: hexBig(bal), nonce: hexInt(nonce), code, miner: q.kinds.has('miner') };
|
||||
} catch (e) { lastError = String(e.message || e).slice(0, 200); accountQueue.set(a, q); return null; }
|
||||
}));
|
||||
const ok = reads.filter(Boolean);
|
||||
if (ok.length) {
|
||||
const vals = []; const params = [];
|
||||
ok.forEach((r, k) => { vals.push(`($${k * 6 + 1},$${k * 6 + 2},$${k * 6 + 3},$${k * 6 + 4},$${k * 6 + 5},$${k * 6 + 6})`); params.push(r.a, r.bal, r.nonce, r.code, r.code > 0, r.miner); });
|
||||
await sql(`INSERT INTO ${T}live_accounts (address, balance_wei, nonce, code_bytes, is_contract, miner, refreshed_at) VALUES ${vals.map(v => v.replace(/\)$/, ', now())')).join(',')}
|
||||
ON CONFLICT (address) DO UPDATE SET balance_wei = EXCLUDED.balance_wei, nonce = EXCLUDED.nonce, code_bytes = EXCLUDED.code_bytes, is_contract = EXCLUDED.is_contract,
|
||||
miner = ${T}live_accounts.miner OR EXCLUDED.miner, refreshed_at = now(), updated_at = now()`, params);
|
||||
}
|
||||
accountsPerMin += ok.length;
|
||||
return ok.length;
|
||||
}
|
||||
|
||||
async function queueMiners(sql) {
|
||||
// the observer's miners of the last 10 minutes get a balance read every minute, so a miner page is never older than that
|
||||
const rows = await sql(`SELECT DISTINCT evm_miner AS a FROM ${T}live_blocks WHERE received_at > now() - interval '10 minutes' AND evm_miner IS NOT NULL`).catch(() => []);
|
||||
for (const r of rows) { const q = accountQueue.get(lower(r.a)) || { kinds: new Set(), number: 0 }; q.kinds.add('miner'); accountQueue.set(lower(r.a), q); }
|
||||
}
|
||||
|
||||
async function main() {
|
||||
const env = envFile();
|
||||
const sql = neon(process.env.DATABASE_URL || env.DATABASE_URL);
|
||||
await schema(sql);
|
||||
for (const r of await sql(`SELECT address, code_bytes FROM ${T}live_accounts WHERE code_bytes IS NOT NULL`)) contractKnown.set(r.address, Number(r.code_bytes));
|
||||
const st = (await sql(`SELECT indexed_number FROM ${T}explorer_state WHERE id = 1`))[0] || {};
|
||||
let indexed = st.indexed_number === null || st.indexed_number === undefined ? null : Number(st.indexed_number);
|
||||
const chainId = hexInt(await rpc('eth_chainId'));
|
||||
log(`explorer indexer: prefix ${T}, rpc ${EVM}, chain id ${chainId}, indexed ${indexed === null ? 'nothing yet' : indexed}, batch ${BATCH}, retain ${RETAIN_H} h`);
|
||||
let lastInfo = 0, lastMiners = 0, lastRetain = 0, lastLog = Date.now(), lastRate = Date.now(), rateBlocks = 0;
|
||||
for (;;) {
|
||||
const t0 = Date.now();
|
||||
try {
|
||||
const tip = hexInt(await rpc('eth_blockNumber'));
|
||||
if (indexed === null) { indexed = BACKFILL > 0 ? Math.max(-1, tip - BACKFILL) : -1; log(`starting at chain block ${indexed + 1}, tip ${tip}`); }
|
||||
// the reorg check: the newest three indexed numbers must still carry the hash we stored
|
||||
if (indexed >= 0) {
|
||||
const lo = Math.max(0, indexed - 2);
|
||||
const stored = await sql(`SELECT number, hash FROM ${T}explorer_blocks WHERE number BETWEEN $1 AND $2`, [lo, indexed]);
|
||||
const live = await Promise.all(stored.map(s => rpc('eth_getBlockByNumber', [hex(Number(s.number)), false]).catch(() => null)));
|
||||
const moved = stored.filter((s, i) => live[i] && lower(live[i].hash) !== s.hash).map(s => Number(s.number));
|
||||
if (moved.length) { const from = Math.min(...moved); log(`reorg: chain block ${from} moved, re-indexing from it`); indexed = from - 1; }
|
||||
}
|
||||
const behind = tip - indexed;
|
||||
if (behind > 0) {
|
||||
const n = Math.min(BATCH, behind);
|
||||
const numbers = Array.from({ length: n }, (_, i) => indexed + 1 + i);
|
||||
await indexNumbers(sql, numbers);
|
||||
indexed += n; rateBlocks += n;
|
||||
}
|
||||
if (Date.now() - lastMiners > 60000) { await queueMiners(sql); lastMiners = Date.now(); }
|
||||
// while catching up, only the newest addresses are refreshed; the rest wait for the tip
|
||||
await refreshAccounts(sql, behind > BATCH ? 4 : ACCOUNTS_PER_TICK);
|
||||
let info = null;
|
||||
if (Date.now() - lastInfo > 300000) { info = await rpc('igneum_getNodeInfo').catch(() => null); lastInfo = Date.now(); }
|
||||
if (Date.now() - lastRetain > 3600000) {
|
||||
const cut = Date.now() - RETAIN_H * 3600000;
|
||||
const d = await sql(`WITH d AS (DELETE FROM ${T}live_txs WHERE ts_ms < $1 RETURNING 1) SELECT count(*)::int AS n FROM d`, [cut]);
|
||||
await sql(`DELETE FROM ${T}explorer_blocks WHERE ts_ms < $1`, [cut]);
|
||||
if (d[0] && d[0].n) log(`retention: ${d[0].n} transactions older than ${RETAIN_H} h removed`);
|
||||
lastRetain = Date.now();
|
||||
}
|
||||
const dt = (Date.now() - lastRate) / 1000; let bps = null; if (dt >= 10) { bps = rateBlocks / dt; rateBlocks = 0; lastRate = Date.now(); }
|
||||
await sql(`UPDATE ${T}explorer_state SET indexed_number = $1, tip_number = $2, chain_id = $3, backlog = $4, blocks_per_s = COALESCE($5, blocks_per_s), last_error = $6, updated_at = now(),
|
||||
node_info = COALESCE($7::jsonb, node_info), txs_total = CASE WHEN $8 THEN (SELECT count(*) FROM ${T}live_txs) ELSE txs_total END,
|
||||
accounts_total = CASE WHEN $8 THEN (SELECT count(*) FROM ${T}live_accounts) ELSE accounts_total END WHERE id = 1`,
|
||||
[indexed, tip, chainId, Math.max(0, tip - indexed), bps, lastError, info ? JSON.stringify(info) : null, Date.now() - lastLog > 60000 || ONCE]);
|
||||
if (Date.now() - lastLog > 60000) { log(`indexed ${indexed} of tip ${tip} (behind ${tip - indexed}); last minute: ${blocksPerMin} blocks, ${txsPerMin} txs, ${accountsPerMin} account reads, queue ${accountQueue.size}${lastError ? `, last error: ${lastError}` : ''}`); txsPerMin = blocksPerMin = accountsPerMin = 0; lastError = null; lastLog = Date.now(); }
|
||||
if (ONCE) { log(`once: indexed ${indexed}, tip ${tip}`); return; }
|
||||
if (tip - indexed <= 0) await new Promise(r => setTimeout(r, Math.max(0, TICK_MS - (Date.now() - t0))));
|
||||
} catch (e) {
|
||||
lastError = String(e.message || e).slice(0, 200);
|
||||
log(`tick failed: ${lastError}`);
|
||||
if (ONCE) { process.exitCode = 1; return; }
|
||||
await new Promise(r => setTimeout(r, 5000));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (process.argv[1] && /explorer-indexer\.mjs$/.test(process.argv[1])) main().catch(e => { console.error(e); process.exit(1); });
|
||||
38
tools/observer/explorer-indexer.test.mjs
Normal file
38
tools/observer/explorer-indexer.test.mjs
Normal file
|
|
@ -0,0 +1,38 @@
|
|||
// The explorer indexer's pure parts: a transaction row from the node's JSON, the touched-address set, the hex helpers.
|
||||
// node --test tools/observer/explorer-indexer.test.mjs
|
||||
import { test } from 'node:test';
|
||||
import assert from 'node:assert/strict';
|
||||
import { txRow, touched, hexInt, hexBig } from './explorer-indexer.mjs';
|
||||
|
||||
const block = { hash: '0xE61CFFB1D306AD8D3F4F9FEC6D750100517572F54A14A82B40EBFD393531320A', number: '0x654d', timestamp: '0x6acd', miner: '0x8DB0505bbb03dffb58875897eeedc4e790a68d18', igneum: { rewards: [{ miner: '0x78fb14e43d3170dcf9abfdc70f15685c8ea31678', wei: '0x1' }] } };
|
||||
const tx = { hash: '0x4E7C2A9D54208B70A557184CB35E5847CDA4E90FFCA9578C6FC05B9C0DC4EF25', from: '0xE617a7f966fc009934e943fe0453821ff0c8716f', to: '0x7fe2f74b45dd8f0a739c906817f197631fdeff85', value: '0x0', nonce: '0x1a', gas: '0x38936', gasPrice: '0x178411b200', input: '0xa9059cbb000000000000000000000000', transactionIndex: '0x3', type: '0x2' };
|
||||
const rc = { transactionHash: tx.hash, gasUsed: '0x6e5e', effectiveGasPrice: '0x178411b200', status: '0x1', contractAddress: null, logs: [{ logIndex: '0x0', address: '0x7FE2f74b45dd8f0a739c906817f197631fdeff85', topics: ['0xaa'], data: '0x01' }], igneum: { minerTip: '0x148eb7b4f000', pgasUsed: '0x56c' } };
|
||||
|
||||
test('hex helpers: null stays null, big values keep every digit', () => {
|
||||
assert.equal(hexInt(null), null); assert.equal(hexInt('0x654d'), 25933);
|
||||
assert.equal(hexBig('0x18c55612df9b9f9a3600'), '116977006436888600000000'); assert.equal(hexBig(undefined), null);
|
||||
});
|
||||
|
||||
test('txRow: lower-case addresses and hashes, numbers decoded, the receipt folded in, the selector and the log list', () => {
|
||||
const r = txRow(tx, rc, block);
|
||||
assert.equal(r.hash, tx.hash.toLowerCase()); assert.equal(r.block_hash, block.hash.toLowerCase()); assert.equal(r.block_number, 25933); assert.equal(r.tx_index, 3);
|
||||
assert.equal(r.from_addr, '0xe617a7f966fc009934e943fe0453821ff0c8716f'); assert.equal(r.to_addr, '0x7fe2f74b45dd8f0a739c906817f197631fdeff85');
|
||||
assert.equal(r.value_wei, '0'); assert.equal(r.nonce, 26); assert.equal(r.gas, 231734); assert.equal(r.gas_used, 28254); assert.equal(r.status, 1);
|
||||
assert.equal(r.selector, '0xa9059cbb'); assert.equal(r.input_bytes, 16); assert.equal(r.logs_count, 1); assert.equal(r.logs[0].address, '0x7fe2f74b45dd8f0a739c906817f197631fdeff85');
|
||||
assert.equal(r.igneum.minerTip, '0x148eb7b4f000'); assert.equal(r.tx_type, 2); assert.equal(r.ts_ms, 0x6acd * 1000);
|
||||
});
|
||||
|
||||
test('txRow without a receipt: the receipt fields are null, nothing throws', () => {
|
||||
const r = txRow({ ...tx, input: '0x', to: null }, null, block);
|
||||
assert.equal(r.gas_used, null); assert.equal(r.status, null); assert.equal(r.selector, null); assert.equal(r.to_addr, null); assert.deepEqual(r.logs, []); assert.equal(r.igneum, null);
|
||||
});
|
||||
|
||||
test('touched: senders, receivers, the created contract and the rewarded miners, each lower case with its reasons', () => {
|
||||
const rows = [txRow(tx, { ...rc, contractAddress: '0xABCDEF0000000000000000000000000000000001' }, block)];
|
||||
const m = touched(rows, block);
|
||||
assert.deepEqual([...m.get('0xe617a7f966fc009934e943fe0453821ff0c8716f')], ['sender']);
|
||||
assert.deepEqual([...m.get('0x7fe2f74b45dd8f0a739c906817f197631fdeff85')], ['receiver']);
|
||||
assert.deepEqual([...m.get('0xabcdef0000000000000000000000000000000001')], ['contract']);
|
||||
assert.deepEqual([...m.get('0x78fb14e43d3170dcf9abfdc70f15685c8ea31678')], ['miner']);
|
||||
assert.deepEqual([...m.get('0x8db0505bbb03dffb58875897eeedc4e790a68d18')], ['miner']);
|
||||
});
|
||||
Loading…
Reference in a new issue