282 lines
23 KiB
JavaScript
282 lines
23 KiB
JavaScript
// 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_NETWORK the network the node must report in igneum_getNodeInfo, default igneum-devnet-3; another name refuses
|
|
// to index (8 October 2026, 10:09 UTC: a unit pointed at the first devnet's node for 26 s wrote 4,797 of
|
|
// that chain's blocks into the dn3_ tables; the tables were wiped and re-filled)
|
|
// 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 (0x), block_hash (the DAG block's hash, no 0x, as live_blocks), 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 NETWORK = process.env.EXPLORER_NETWORK || 'igneum-devnet-3';
|
|
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).replace(/^0x/, ''), 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).replace(/^0x/, ''), 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'));
|
|
const info0 = await rpc('igneum_getNodeInfo').catch(() => null);
|
|
if (!info0 || info0.network !== NETWORK) throw new Error(`the node at ${EVM} reports network ${info0 ? info0.network : 'unknown'}, this indexer writes ${NETWORK} (EXPLORER_NETWORK); refusing to index`);
|
|
log(`explorer indexer: prefix ${T}, rpc ${EVM} (${info0.network}, digest ${String(info0.digest).slice(0, 8)}), 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).replace(/^0x/, '') !== 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 (info && info.network !== NETWORK) throw new Error(`the node now reports network ${info.network}, not ${NETWORK}; stopping`); }
|
|
if (indexed > tip + 100) { log(`the node's tip ${tip} is behind the indexed number ${indexed}: a node that moved back; re-indexing from ${tip - 3}`); indexed = Math.max(-1, tip - 3); }
|
|
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); });
|