igneum/tools/observer/observer.mjs
igneum-labs 2c58912211 Live proving feed: the page says when its shard states are stale, and a chain block with an empty plan reads "no shards in this block" (8 October 2026, main's order: the project lead wants to watch proving and shards in real time; a stale snapshot must read as stale). Reading of the day: /api/live's proven and paid states come from the Devnet 3 observer's own node on igneum-build-1 (igneum-observer-node-dn3, EVM 26850), not from hub-1; that node stopped crediting payouts at 23:57Z on 7 October (paidShards frozen at 7,304 while blocks carried 8,457 proof records by 03:00Z and hub-1 counted 11,861 by 08:56Z), and the page kept showing the paid total as fresh. site/lib/proving-health.mjs (tests known-failed first: "Cannot find module", then 10 green on build-2): stale when the observer's last status read is over 60 s old (status_at, new in the observer's live_state.proving), when the node's tip trails the newest block by over 300 DAA, or when blocks carried proof records in the last 10 minutes and no shard moved to verified or paid for 10 minutes (the blind-to-records class); a quiet chain is not stale for lack of credits. /api/live carries stale, stale_reason, status_age_s, tip_lag_daa, records_10m and last_change_at on proving (two more indexed readings; scene/feed-contract.json lists the keys; the recorded fixture carries them). /live: the Proven tile's line reads "proving feed stale: ..." in ember, else "... last credit N s ago"; the inspector's count adds "feed stale". Scene 2.0.7: a chain block whose plan is known (block.number set, carried as plan on the normalised block) and empty reads "no shards in this block"; loading and late stay for an unread plan (shard-words test, known-failed on 2.0.6). Gate: the health test joins the site unit-test line.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-08 09:12:56 +00:00

1171 lines
75 KiB
JavaScript

#!/usr/bin/env node
// Igneum devnet observer. Connects to a node's wRPC JSON endpoint, subscribes to new blocks and
// writes what it sees to Neon so /api/live and /live can show it.
//
// node tools/observer/observer.mjs
//
// Environment (every value has a default):
// IGNEUM_RPC wRPC JSON url of the node default ws://127.0.0.1:28610 (igneum-devnet wRPC JSON port)
// DATABASE_URL Neon connection string default: read from ~/.config/igneum/env
// LIVE_RETAIN_HOURS hours of blocks to keep default 24
// LIVE_TABLE_PREFIX prefix for every table name default '' (a test observer can write fintest_live_* instead)
// IGNEUM_EVM_RPC execution-layer JSON-RPC (http) of a node on the proving build, for the shard plans and proof
// records default http://127.0.0.1:26800 (the Mac app's node; the observer node has no EVM listener yet)
//
// Tables (created on start if missing): live_blocks, live_state, live_events, live_checkpoints, live_certificates,
// live_proofs. See README.md.
// Proving v0 (4 Oct 2026, spec 7.7): every chain block's shard plan is read over the EVM RPC (igneum_getShardPlan) as
// the block joins the selected chain, one live_proofs row per shard (planned); the proof records of the chain blocks
// of the last 10 minutes are polled in rotation (igneum_getProofRecords, 80 blocks per 2 s tick, four calls in flight) and move a shard to
// proving (a record in the pool), verified (the SP1 proof verified, or the record carried and checked) or paid (a
// carrying block paid it). live_state.proving carries the activation state and the 10-minute counts. A node without
// the proving RPCs (method not found) gives proving = {supported: false}; the observer rechecks every 5 minutes.
// Finality v2: subscribes to FinalityLock notifications and polls getFinalityCheckpoints and getFinalityWeights
// every 2 s; writes the checkpoints table and emits "checkpoint N locked (xx% of weight)" events.
// Light client (site/api/checkpoint.mjs, site/verify/): every certificate a block carries is stored with the voter
// table the node reports and the header chain from the previous locked checkpoint, so a browser can verify it.
// Explorer (5 October 2026, docs/plans/explorer.md): every block row also carries what /explorer, /block/<hash> and
// /address/<addr> show: the EVM transaction count and hashes, the miner's EVM payout address (the coinbase's IGNA
// tag, else the vote key hash's low 20 bytes as consensus/core/src/evm.rs falls back to), the number of proof records
// carried (the IGNP section), the subsidy the payload declares and the outputs paid, the header fields, the mergeset
// and the certificate indices, all read from the blockAdded notification itself: no extra RPC per block. The chain
// block number comes with the shard plan the proving feed already fetches. live_state.rpc_load counts this
// process's RPC calls per minute (wRPC and EVM) so the load of a change can be measured; live_state.supply_check is
// the hourly comparison of the coinbase sums with the emission rule (site/lib/emission.mjs, spec 2.5).
// Zero dependencies: Node 22 WebSocket and fetch, Neon's HTTP SQL endpoint.
//
// Ingest path (4 Oct 2026, after the observer fell 80 min behind under machine load 200 to 300): notifications
// only enqueue; a drain loop turns them into rows; the block flush, the colour marking (mergesets, from the
// notification's verbose data with getBlock only on a cache miss, four at a time) and the certificate work each run
// on their own timer and never wait on one another. live_state carries observer_lag_s (now minus the newest stored
// block's header time) and queue_depth. Per-minute counts are bucketed by the block's own timestamp, so a catch-up
// fills past minutes instead of painting a spike. If no blockAdded arrives for 60 s while the node's block_count
// advances, the observer resubscribes; after two failed attempts it exits 2 so tools/observer/run.sh restarts it.
import { readFileSync } from 'node:fs';
import { homedir } from 'node:os';
import { blockSubsidy } from '../../site/lib/emission.mjs';
import { keccak256 } from '../../site/lib/eth.mjs';
import { makeRunner as makeDetector, optsFromEnv as detectorOpts } from './detector.mjs'; // Counter ASIC 3.0 item 4a: the share-pattern detector
import { vendorShareTick } from './vendor-share.mjs'; // Counter ASIC 3.0 item 7: hash rate by vendor
const RPC = process.env.IGNEUM_RPC || 'ws://127.0.0.1:28610';
const RETAIN_HOURS = Number(process.env.LIVE_RETAIN_HOURS || 24);
const T = (process.env.LIVE_TABLE_PREFIX || '').replace(/[^a-z0-9_]/gi, '');
const TB = `${T}live_blocks`, TS = `${T}live_state`, TE = `${T}live_events`, TC = `${T}live_checkpoints`, TX = `${T}live_certificates`, TP = `${T}live_proofs`;
const EVM = process.env.IGNEUM_EVM_RPC || 'http://127.0.0.1:26800';
const STATE_EVERY_MS = 2000;
const FLUSH_EVERY_MS = 500;
const PRUNE_EVERY_MS = 60_000;
const QUIET_AFTER_MS = 5 * 60_000;
const DIFFICULTY_STEP = 0.05;
const STARTED_AT = new Date();
// ---------- Neon ----------
function databaseUrl() {
if (process.env.DATABASE_URL) return process.env.DATABASE_URL;
let env = '';
try { env = readFileSync(`${homedir()}/.config/igneum/env`, 'utf8'); } catch { }
const m = /^DATABASE_URL=(.*)$/m.exec(env);
if (!m) { console.error('DATABASE_URL is not set and ~/.config/igneum/env has none'); process.exit(1); }
return m[1].trim().replace(/^['"]|['"]$/g, '');
}
const DB_URL = databaseUrl();
const DB_HOST = new URL(DB_URL).hostname.replace('-pooler', '');
async function sql(query, params = []) {
const r = await fetch(`https://${DB_HOST}/sql`, {
method: 'POST',
headers: { 'Neon-Connection-String': DB_URL, 'Content-Type': 'application/json' },
body: JSON.stringify({ query, params }),
});
const j = await r.json();
if (!r.ok) throw new Error(j.message || JSON.stringify(j));
return j.rows || [];
}
async function setupSchema() {
const stmts = [
`CREATE TABLE IF NOT EXISTS ${TB} (
hash text PRIMARY KEY,
blue_score bigint NOT NULL,
daa_score bigint NOT NULL,
timestamp_ms bigint NOT NULL,
parents int NOT NULL,
parent_hashes text[] NOT NULL DEFAULT '{}',
is_chain_block boolean NOT NULL DEFAULT false,
vote_key_hash text,
miner_address text,
engine text,
received_at timestamptz NOT NULL DEFAULT now())`,
`CREATE INDEX IF NOT EXISTS ${TB}_received_at ON ${TB} (received_at)`,
`CREATE INDEX IF NOT EXISTS ${TB}_timestamp_ms ON ${TB} (timestamp_ms)`,
`CREATE INDEX IF NOT EXISTS ${TB}_vote_key_received ON ${TB} (vote_key_hash, received_at)`,
// GHOSTDAG colour (4 Oct 2026, additive): pending until a chain block merges it, then blue (paid) or red (excluded)
`ALTER TABLE ${TB} ADD COLUMN IF NOT EXISTS color text NOT NULL DEFAULT 'pending'`,
// Explorer (5 Oct 2026, additive): what the block pages show, from the notification; number from the shard plan
`ALTER TABLE ${TB} ADD COLUMN IF NOT EXISTS tx_count int`,
`ALTER TABLE ${TB} ADD COLUMN IF NOT EXISTS evm_miner text`,
`ALTER TABLE ${TB} ADD COLUMN IF NOT EXISTS proof_records int`,
`ALTER TABLE ${TB} ADD COLUMN IF NOT EXISTS subsidy_sompi bigint`,
`ALTER TABLE ${TB} ADD COLUMN IF NOT EXISTS paid_sompi bigint`,
`ALTER TABLE ${TB} ADD COLUMN IF NOT EXISTS selected_parent text`,
`ALTER TABLE ${TB} ADD COLUMN IF NOT EXISTS number bigint`,
`ALTER TABLE ${TB} ADD COLUMN IF NOT EXISTS detail jsonb`,
`CREATE INDEX IF NOT EXISTS ${TB}_evm_miner_received ON ${TB} (evm_miner, received_at)`,
`CREATE INDEX IF NOT EXISTS ${TB}_miner_address_received ON ${TB} (miner_address, received_at)`,
`CREATE INDEX IF NOT EXISTS ${TB}_number ON ${TB} (number)`,
`CREATE TABLE IF NOT EXISTS ${TS} (
id int PRIMARY KEY DEFAULT 1 CHECK (id = 1),
block_count bigint, header_count bigint, blue_score bigint, difficulty double precision,
hashes_per_second_estimate bigint, peers int, mempool int,
node_version text, network text,
blocks_60s int, blocks_per_minute jsonb,
observer_started_at timestamptz, updated_at timestamptz)`,
`ALTER TABLE ${TS} ADD COLUMN IF NOT EXISTS finality jsonb`,
`ALTER TABLE ${TS} ADD COLUMN IF NOT EXISTS observer_lag_s double precision`,
`ALTER TABLE ${TS} ADD COLUMN IF NOT EXISTS queue_depth int`,
// Proving v0 (4 Oct 2026, additive): {supported, active, activation_daa, tip_daa, blocks_fully_proven_10m, ...}
`ALTER TABLE ${TS} ADD COLUMN IF NOT EXISTS proving jsonb`,
// Explorer: this process's RPC calls per minute, and the hourly coinbase-versus-rule comparison
`ALTER TABLE ${TS} ADD COLUMN IF NOT EXISTS rpc_load jsonb`,
`ALTER TABLE ${TS} ADD COLUMN IF NOT EXISTS supply_check jsonb`,
`ALTER TABLE ${TS} ADD COLUMN IF NOT EXISTS detector jsonb`,
// Vendor share (Counter ASIC 3.0 item 7, tools/observer/vendor-share.mjs): hash rate by vendor, fleet-reported and chain-attributed
`ALTER TABLE ${TS} ADD COLUMN IF NOT EXISTS vendor_share jsonb`,
// One row per planned shard of a chain block (spec 7.7 item 8): the plan as the block joins the chain, then the
// record's progress. prover is the first 8 hex characters of the record's vote key hash (never the full key, R4.6.2).
`CREATE TABLE IF NOT EXISTS ${TP} (
block_hash text NOT NULL,
shard int NOT NULL,
shards int NOT NULL,
block_number bigint,
block_daa bigint,
block_ts bigint,
pgas bigint,
state text NOT NULL DEFAULT 'planned',
prover text,
verified boolean,
carried_by text,
carrier_number bigint,
carrier_daa bigint,
lag_daa int,
payout_wei text,
received_at timestamptz NOT NULL DEFAULT now(),
updated_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (block_hash, shard))`,
`CREATE INDEX IF NOT EXISTS ${TP}_received_at ON ${TP} (received_at)`,
`CREATE TABLE IF NOT EXISTS ${TE} (
id bigserial PRIMARY KEY,
ts timestamptz NOT NULL DEFAULT now(),
kind text NOT NULL,
text text NOT NULL)`,
`CREATE INDEX IF NOT EXISTS ${TE}_ts ON ${TE} (ts)`,
`CREATE TABLE IF NOT EXISTS ${TC} (
index bigint PRIMARY KEY,
hash text NOT NULL,
blue_score bigint NOT NULL,
daa_score bigint NOT NULL,
state text NOT NULL,
signed_weight bigint NOT NULL DEFAULT 0,
active_weight bigint NOT NULL DEFAULT 0,
total_weight bigint NOT NULL DEFAULT 0,
fraction_active double precision NOT NULL DEFAULT 0,
fraction_total double precision NOT NULL DEFAULT 0,
votes_seen int NOT NULL DEFAULT 0,
voters int NOT NULL DEFAULT 0,
aggregators text[] NOT NULL DEFAULT '{}',
locked_at timestamptz,
first_seen_at timestamptz NOT NULL DEFAULT now(),
updated_at timestamptz NOT NULL DEFAULT now())`,
`CREATE INDEX IF NOT EXISTS ${TC}_state ON ${TC} (state, index)`,
// One row per certified checkpoint: the certificate bytes as carried in a block, the voter table (canonical
// order: above dust, not stripped, sorted by key hash) with public keys, and the selected-chain headers from
// the previous locked checkpoint to this one, raw as the node's RPC gives them.
`CREATE TABLE IF NOT EXISTS ${TX} (
index bigint PRIMARY KEY,
hash text NOT NULL,
blue_score bigint,
daa_score bigint,
chain_id text NOT NULL,
voter_count int NOT NULL,
bitmap_hex text NOT NULL,
signature_hex text NOT NULL,
aggregator text NOT NULL,
aggregator_proof_hex text NOT NULL,
certificate_hex text NOT NULL,
carrier text NOT NULL,
voters jsonb NOT NULL,
voters_index bigint NOT NULL,
total_weight bigint NOT NULL,
active_weight double precision NOT NULL,
prev_index bigint,
prev_hash text,
headers jsonb NOT NULL,
headers_complete boolean NOT NULL DEFAULT false,
created_at timestamptz NOT NULL DEFAULT now())`,
];
for (const s of stmts) await sql(s);
}
// ---------- Address decoding (crypto/addresses/src/bech32.rs, ported) ----------
const CHARSET = 'qpzry9x8gf2tvdw0s3jn54khce6mua7l';
function polymod(values) {
let c = 1n;
for (const d of values) {
const c0 = c >> 35n;
c = ((c & 0x07ffffffffn) << 5n) ^ BigInt(d);
if (c0 & 0x01n) c ^= 0x98f2bc8e61n;
if (c0 & 0x02n) c ^= 0x79b76d99e2n;
if (c0 & 0x04n) c ^= 0xf33e5fb3c4n;
if (c0 & 0x08n) c ^= 0xae2eabe2a8n;
if (c0 & 0x10n) c ^= 0x1e4f43e470n;
}
return c ^ 1n;
}
function conv8to5(bytes) {
const out = []; let buff = 0, bits = 0;
for (const b of bytes) {
buff = ((buff << 8) | b) & 0xffff; bits += 8;
while (bits >= 5) { bits -= 5; out.push((buff >> bits) & 31); buff &= (1 << bits) - 1; }
}
if (bits > 0) out.push((buff << (5 - bits)) & 31);
return out;
}
export function encodeAddress(prefix, version, payload) {
const five = conv8to5([version, ...payload]);
const pre = [...prefix].map(ch => ch.charCodeAt(0) & 0x1f);
const sum = polymod([...pre, 0, ...five, 0, 0, 0, 0, 0, 0, 0, 0]);
const sumBytes = []; for (let i = 4; i >= 0; i--) sumBytes.push(Number((sum >> BigInt(i * 8)) & 0xffn));
return `${prefix}:${[...five, ...conv8to5(sumBytes)].map(v => CHARSET[v]).join('')}`;
}
function scriptToAddress(prefix, script) {
if (script.length === 34 && script[0] === 0x20 && script[33] === 0xac) return encodeAddress(prefix, 0, script.slice(1, 33));
if (script.length === 35 && script[0] === 0x21 && script[34] === 0xab) return encodeAddress(prefix, 1, script.slice(1, 34));
if (script.length === 35 && script[0] === 0xaa && script[1] === 0x20 && script[34] === 0x87) return encodeAddress(prefix, 8, script.slice(2, 34));
return null;
}
function payloadBytes(p) {
if (Array.isArray(p)) return Uint8Array.from(p);
if (typeof p === 'string') return Uint8Array.from(Buffer.from(p, 'hex'));
return new Uint8Array();
}
// Coinbase payload: blue score u64 LE, subsidy u64 LE, script version u16 LE, script length u8, script, extra data.
function minerFromCoinbase(block, prefix) {
const tx = block.transactions && block.transactions[0];
if (!tx) return { address: null, extra: null };
const b = payloadBytes(tx.payload);
if (b.length < 19) return { address: null, extra: null };
const len = b[18];
const script = b.slice(19, 19 + len);
const extra = b.slice(19 + len);
// The node prefixes the extra data with its own version tag ("2.1.0/"). Anything after that is the
// miner's tag (engine or label). The node exposes no engine name over RPC, so this is null on devnet v0.
// Finality v2: the miner's key reveal (IGNK ...) and the node's finality section (... IGNF) are binary; only the
// text before the reveal is the miner's tag.
let text = Buffer.from(extra).toString('latin1');
const reveal = text.indexOf('IGNK');
if (reveal >= 0) text = text.slice(0, reveal);
text = text.replace(/[^\x20-\x7e]/g, '').replace(/^\d+\.\d+\.\d+(-[\w.]+)?\//, '').trim();
return { address: scriptToAddress(prefix, script), extra: text ? text.slice(0, 64) : null };
}
// ---------- RPC client ----------
class Rpc {
constructor(url) { this.url = url; this.id = 0; this.pending = new Map(); this.onNotification = () => { }; this.ws = null; this.open = false; }
connect() {
return new Promise((resolve) => {
const ws = new WebSocket(this.url); this.ws = ws;
ws.onopen = () => { this.open = true; log('rpc connected', this.url); resolve(true); };
ws.onmessage = (e) => {
// A header nonce is a full u64; JSON.parse would round it above 2^53, so it is kept as a string
let m; try { m = JSON.parse(String(e.data).replace(/"nonce":(\d+)/g, '"nonce":"$1"')); } catch { return; }
if (m.id !== undefined && m.id !== null && this.pending.has(m.id)) {
const p = this.pending.get(m.id); this.pending.delete(m.id);
m.error ? p.reject(new Error(m.error.message || JSON.stringify(m.error))) : p.resolve(m.params);
} else if (m.method) this.onNotification(m.method, m.params);
};
ws.onerror = () => { };
ws.onclose = () => {
const wasOpen = this.open; this.open = false;
for (const p of this.pending.values()) p.reject(new Error('rpc closed'));
this.pending.clear();
if (wasOpen) log('rpc closed'); else resolve(false);
this.onClose && this.onClose();
};
});
}
call(method, params = {}) {
rpcLoad.wrpc++;
return new Promise((resolve, reject) => {
if (!this.open) return reject(new Error('rpc not connected'));
const id = ++this.id; this.pending.set(id, { resolve, reject });
this.ws.send(JSON.stringify({ id, method, params }));
setTimeout(() => { if (this.pending.has(id)) { this.pending.delete(id); reject(new Error(`${method} timed out`)); } }, 10_000);
});
}
}
// ---------- Observer ----------
const log = (...a) => console.log(new Date().toISOString(), ...a);
const short = h => (h || '').slice(0, 8);
// RPC load (5 Oct 2026): calls this process makes, counted per wall-clock minute, the last five minutes kept. Logged
// once a minute and written to live_state.rpc_load, so a change to the observer can be measured against the rule
// in docs/plans/explorer.md (the explorer must stay under 2x of the load before it).
const rpcLoad = { wrpc: 0, evm: 0, minutes: [], minuteStart: Date.now() };
function rollRpcLoad(now) {
if (now - rpcLoad.minuteStart < 60_000) return;
const secs = (now - rpcLoad.minuteStart) / 1000;
const m = { at: new Date(rpcLoad.minuteStart).toISOString(), wrpc: Math.round(rpcLoad.wrpc * 60 / secs), evm: Math.round(rpcLoad.evm * 60 / secs) };
rpcLoad.minutes.push(m); if (rpcLoad.minutes.length > 5) rpcLoad.minutes.shift();
log(`rpc load: wrpc ${m.wrpc}/min, evm ${m.evm}/min`);
rpcLoad.wrpc = 0; rpcLoad.evm = 0; rpcLoad.minuteStart = now;
}
function rpcLoadState() {
const ms = rpcLoad.minutes; if (!ms.length) return null;
const avg = k => Math.round(ms.reduce((a, m) => a + m[k], 0) / ms.length);
return { wrpc_per_min: ms[ms.length - 1].wrpc, evm_per_min: ms[ms.length - 1].evm, wrpc_per_min_5m: avg('wrpc'), evm_per_min_5m: avg('evm'), minutes: ms.length, evm_rpc: EVM };
}
// Explorer detail of a block, all from the notification (header, coinbase, verbose data, EVM transactions).
// Coinbase payload: blue score u64 LE, subsidy u64 LE (the full E(daa) of this block, coinbase.rs), script version
// u16 LE, script length u8, script, extra data. Extra data: the node's version tag, the miner's tag, the key reveal
// (IGNK, finality.rs), the payout address (IGNA || 40 hex, evm.rs), the proof-record section (records || len_le32 ||
// IGNP, proving.rs, 274 bytes each) and the finality section (items || len_le32 || IGNF).
const PROOF_RECORD_LEN = 274;
function coinbaseDetail(block, prefix) {
const tx = block.transactions && block.transactions[0];
if (!tx) return null;
const b = payloadBytes(tx.payload);
if (b.length < 19) return null;
const dv = new DataView(b.buffer, b.byteOffset, b.byteLength);
const subsidy = dv.getBigUint64(8, true);
const extra = b.subarray(19 + b[18]);
const text = Buffer.from(extra).toString('latin1');
const a = text.indexOf('IGNA');
const evmMiner = a >= 0 && /^[0-9a-f]{40}$/.test(text.slice(a + 4, a + 44)) ? '0x' + text.slice(a + 4, a + 44) : null;
// proof records: strip the finality section, then read the IGNP trailer
let body = extra; let n = body.length;
if (n >= 8 && String.fromCharCode(...body.subarray(n - 4)) === 'IGNF') { const len = new DataView(body.buffer, body.byteOffset, n).getUint32(n - 8, true); if (len + 8 <= n) body = body.subarray(0, n - 8 - len); }
n = body.length; let records = 0;
if (n >= 8 && String.fromCharCode(...body.subarray(n - 4)) === 'IGNP') { const len = new DataView(body.buffer, body.byteOffset, n).getUint32(n - 8, true); if (len + 8 <= n) records = Math.floor(len / PROOF_RECORD_LEN); }
const outputs = (tx.outputs || []).map(o => {
const spk = typeof o.scriptPublicKey === 'string' ? o.scriptPublicKey : (o.scriptPublicKey && o.scriptPublicKey.script) || '';
const script = Uint8Array.from(Buffer.from(String(spk).slice(4), 'hex')); // 2-byte version prefix, then the script
const addr = (o.verboseData && o.verboseData.scriptPublicKeyAddress) || scriptToAddress(prefix, script);
const pool = !addr && script.length > 2 && script[0] === 0x6a; // OP_RETURN "igneum-proving-pool-v0": the pool output burns on the UTXO side
return { address: addr, value: String(o.value), pool };
});
const paid = outputs.reduce((t, o) => t + BigInt(o.value), 0n);
return { id: tx.verboseData && tx.verboseData.transactionId, subsidy: subsidy.toString(), paid: paid.toString(), outputs, evm_miner: evmMiner, records, reveal: text.includes('IGNK') };
}
// EVM transactions ride in the block as raw bytes (RpcBlock.evm_transactions, Vec<Vec<u8>>): the hash is keccak-256 of
// the envelope. Decoding fields from the RLP is left to the execution layer (/api/explorer reads them from the EVM RPC
// when one is configured); here the count and the hashes are enough for the table and the block page.
function evmTxs(block) {
const list = block.evmTransactions || [];
return list.map(t => {
const raw = typeof t === 'string' ? Buffer.from(t.replace(/^0x/, ''), 'hex') : Array.isArray(t) ? Buffer.from(t) : null;
if (!raw) return { hash: null, size: null };
return { hash: '0x' + Buffer.from(keccak256(raw)).toString('hex'), size: raw.length };
});
}
function blockDetail(block, prefix) {
const h = block.header, vd = block.verboseData || {};
const cb = coinbaseDetail(block, prefix);
const txs = evmTxs(block);
const certs = certificatesIn(block).map(c => c.index);
const vk = h.voteKeyHash || null;
const evmMiner = (cb && cb.evm_miner) || (vk ? '0x' + vk.slice(24, 64) : null); // evm.rs miner_evm_address: bytes 12..32 of the vote key hash
return {
row: {
tx_count: txs.length, evm_miner: evmMiner, proof_records: cb ? cb.records : 0,
subsidy_sompi: cb ? cb.subsidy : null, paid_sompi: cb ? cb.paid : null, selected_parent: vd.selectedParentHash || null,
},
detail: {
version: h.version, bits: h.bits, nonce: String(h.nonce), blue_work: h.blueWork, hash_merkle_root: h.hashMerkleRoot,
accepted_id_merkle_root: h.acceptedIdMerkleRoot, utxo_commitment: h.utxoCommitment, pruning_point: h.pruningPoint,
difficulty: vd.difficulty ?? null, parents_by_level: Array.isArray(h.parentsByLevel) ? h.parentsByLevel.length : null,
mergeset: { blues: vd.mergeSetBluesHashes || [], reds: vd.mergeSetRedsHashes || [] },
coinbase: cb, txs, certificates: certs,
},
};
}
const pendingBlocks = []; // rows waiting for the next flush
const pendingChain = { add: new Set(), remove: new Set() };
const pendingMerge = new Map(); // chain block hash -> { blues, reds } from its verbose data, or null to fetch with getBlock
const pendingUnmerge = new Set(); // chain blocks removed by a reorg: their mergeset goes back to pending
const mergesets = new Map(); // block hash -> { blues, reds } from the notification's verbose data (bounded cache)
const MERGESET_CACHE = 4000;
const inbox = []; // raw notifications, drained off the socket path
let draining = false;
let newestStoredTs = 0; // header timestamp of the newest block written to the table
let lastBlockAt = 0; // wall clock of the last blockAdded notification
let blockCountAtLastBlock = 0; // the node's block_count when a block last arrived
let lastDagBlockCount = 0;
let lastResubscribeAt = 0;
let resubscribeAttempts = 0;
let flushBusy = false;
const miners = new Map(); // vote_key_hash -> { lastSeen ms, quiet bool }
const minuteCounts = new Map(); // minute epoch -> count (last 60 minutes)
const recentArrivals = []; // ms timestamps of arrivals in the last 60 s
const workSamples = []; // { t: ms, w: BigInt blue work } of the heaviest header seen, last 10 min
let lastDifficultyEvent = null;
// Fallback hash-rate estimate when the node refuses estimateNetworkHashesPerSecond (it needs a 1,000-block
// window): blue work added per second over the last 10 minutes, the same quantity the node measures.
function noteWork(now, blueWorkHex) {
if (!blueWorkHex) return;
let w; try { w = BigInt('0x' + blueWorkHex); } catch { return; }
const last = workSamples[workSamples.length - 1];
if (last && w <= last.w) return;
workSamples.push({ t: now, w });
while (workSamples.length && workSamples[0].t < now - 10 * 60_000) workSamples.shift();
}
function localHashesPerSecond() {
if (workSamples.length < 2) return null;
const a = workSamples[0], b = workSamples[workSamples.length - 1];
const secs = (b.t - a.t) / 1000;
if (secs < 20) return null;
return Number((b.w - a.w) / BigInt(Math.round(secs)));
}
let knownPeers = null; // Set of peer ids from the previous tick
let network = 'igneum-devnet';
let addressPrefix = 'igneumdev';
let nodeVersion = null;
let sinkHash = null;
async function recordEvent(kind, text) {
log('event', kind, text);
try { await sql(`INSERT INTO ${TE} (kind, text) VALUES ($1, $2)`, [kind, text]); } catch (e) { log('event write failed', e.message); }
}
// Counts are bucketed by the block's own header timestamp, never by when the observer processed it: a catch-up
// after a stall fills in the past minutes instead of painting a spike into the minute it happened to run.
function noteArrival(now, ts) {
recentArrivals.push(ts);
const minute = Math.floor(ts / 60_000) * 60_000;
minuteCounts.set(minute, (minuteCounts.get(minute) || 0) + 1);
if (minuteCounts.size > 70) for (const k of minuteCounts.keys()) if (k < now - 61 * 60_000) minuteCounts.delete(k);
}
function blocks60s(now) {
if (recentArrivals.length > 2000) { const keep = recentArrivals.filter(t => t > now - 120_000); recentArrivals.length = 0; recentArrivals.push(...keep); }
let n = 0; for (const t of recentArrivals) if (t > now - 60_000) n++;
return n;
}
function blocksPerMinute(now) {
const out = []; const start = Math.floor(now / 60_000) * 60_000 - 59 * 60_000;
for (let m = start; m <= now; m += 60_000) out.push([m, minuteCounts.get(m) || 0]);
return out;
}
function onBlock(block) {
const h = block.header; if (!h || !h.hash) return;
const now = Date.now();
const parents = (h.parentsByLevel && h.parentsByLevel[0]) || [];
const miner = minerFromCoinbase(block, addressPrefix);
for (const cert of certificatesIn(block)) pendingCerts.push({ cert, carrier: h.hash });
const vk = h.voteKeyHash || null;
const vd = block.verboseData;
if (vd && Array.isArray(vd.mergeSetBluesHashes)) {
const ms = { blues: vd.mergeSetBluesHashes, reds: vd.mergeSetRedsHashes || [] };
mergesets.set(h.hash, ms);
if (mergesets.size > MERGESET_CACHE) mergesets.delete(mergesets.keys().next().value);
if (vd.isChainBlock) pendingMerge.set(h.hash, ms);
}
if (vd && vd.isChainBlock) queuePlan(h.hash);
lastBlockAt = now; blockCountAtLastBlock = lastDagBlockCount; resubscribeAttempts = 0;
let ex = null; try { ex = blockDetail(block, addressPrefix); } catch (e) { log('block detail failed', short(h.hash), e.message); }
pendingBlocks.push({
hash: h.hash, blue_score: h.blueScore, daa_score: h.daaScore, timestamp_ms: h.timestamp,
parents: parents.length, parent_hashes: parents,
is_chain_block: !!(block.verboseData && block.verboseData.isChainBlock),
vote_key_hash: vk, miner_address: miner.address, engine: miner.extra,
tx_count: ex ? ex.row.tx_count : null, evm_miner: ex ? ex.row.evm_miner : null, proof_records: ex ? ex.row.proof_records : null,
subsidy_sompi: ex ? ex.row.subsidy_sompi : null, paid_sompi: ex ? ex.row.paid_sompi : null, selected_parent: ex ? ex.row.selected_parent : null,
detail: ex ? JSON.stringify(ex.detail) : null,
});
noteArrival(now, Number(h.timestamp) || now);
noteWork(now, h.blueWork);
if (vk) {
const m = miners.get(vk);
if (!m) { miners.set(vk, { lastSeen: now, quiet: false }); recordEvent('miner_seen', `New miner ${short(vk)} seen${miner.address ? ` (${miner.address.slice(0, 24)}...)` : ''}`); }
else { if (m.quiet) { m.quiet = false; recordEvent('miner_back', `Miner ${short(vk)} is back after a quiet spell`); } m.lastSeen = now; }
}
}
function pgArray(list) { return `{${list.map(s => `"${String(s).replace(/["\\]/g, '')}"`).join(',')}}`; }
async function flushBlocks() {
if (!pendingBlocks.length) return;
const rows = pendingBlocks.splice(0, 200);
const cols = ['hash', 'blue_score', 'daa_score', 'timestamp_ms', 'parents', 'parent_hashes', 'is_chain_block', 'vote_key_hash', 'miner_address', 'engine',
'tx_count', 'evm_miner', 'proof_records', 'subsidy_sompi', 'paid_sompi', 'selected_parent', 'detail'];
const cast = { parent_hashes: '::text[]', detail: '::jsonb', subsidy_sompi: '::bigint', paid_sompi: '::bigint' };
const params = []; const values = [];
for (const r of rows) {
const ph = [];
for (const c of cols) { params.push(c === 'parent_hashes' ? pgArray(r[c]) : r[c]); ph.push(`$${params.length}${cast[c] || ''}`); }
values.push(`(${ph.join(',')})`);
}
try {
await sql(`INSERT INTO ${TB} (${cols.join(',')}) VALUES ${values.join(',')} ON CONFLICT (hash) DO NOTHING`, params);
for (const r of rows) if (Number(r.timestamp_ms) > newestStoredTs) newestStoredTs = Number(r.timestamp_ms);
} catch (e) { log('block insert failed', e.message); if (pendingBlocks.length < 5000) pendingBlocks.unshift(...rows); return; }
if (pendingBlocks.length) await flushBlocks();
}
// Colour: every chain block's mergeset is blue (in the selected chain's past, paid) or red (excluded). Blocks no chain
// block has merged yet stay pending. A reorg puts the removed chain blocks' mergesets back to pending first.
async function mergesetOf(rpc, hash) {
const c = mergesets.get(hash); if (c) return c;
const b = (await rpc.call('getBlock', { hash, includeTransactions: false })).block;
const vd = (b && b.verboseData) || {};
return { blues: vd.mergeSetBluesHashes || [], reds: vd.mergeSetRedsHashes || [] };
}
// Up to four RPC lookups in flight; the cache answers nearly all of them, so this is normally no RPC at all
async function mergesetsOf(rpc, hashes, onFail) {
const out = new Map();
for (let i = 0; i < hashes.length; i += 4) {
await Promise.all(hashes.slice(i, i + 4).map(async h => { try { out.set(h, await mergesetOf(rpc, h)); } catch (e) { onFail(h, e); } }));
}
return out;
}
let colorsBusy = false;
async function flushColors(rpc) {
if (colorsBusy || (!pendingMerge.size && !pendingUnmerge.size)) return;
colorsBusy = true;
try {
if (pendingUnmerge.size) {
const hs = [...pendingUnmerge].slice(0, 50); const back = [];
const got = await mergesetsOf(rpc, hs, (h, e) => log('unmerge fetch failed', h.slice(0, 8), e.message));
for (const h of hs) { pendingUnmerge.delete(h); const m = got.get(h); if (m) back.push(...m.blues, ...m.reds); }
if (back.length) await sql(`UPDATE ${TB} SET color = 'pending' WHERE hash = ANY($1::text[]) AND color <> 'pending'`, [pgArray(back)]);
}
const blues = [], reds = [], need = [];
const batch = [...pendingMerge].slice(0, 200);
for (const [h, m] of batch) { if (m) { blues.push(...m.blues); reds.push(...m.reds); } else need.push(h); pendingMerge.delete(h); }
if (need.length) { const got = await mergesetsOf(rpc, need, (h, e) => log('mergeset fetch failed', h.slice(0, 8), e.message)); for (const m of got.values()) { blues.push(...m.blues); reds.push(...m.reds); } }
if (blues.length) await sql(`UPDATE ${TB} SET color = 'blue' WHERE hash = ANY($1::text[]) AND color <> 'blue'`, [pgArray(blues)]);
if (reds.length) await sql(`UPDATE ${TB} SET color = 'red' WHERE hash = ANY($1::text[]) AND color <> 'red'`, [pgArray(reds)]);
} catch (e) { log('color update failed', e.message); }
colorsBusy = false;
}
async function flushChain() {
if (pendingChain.add.size) { const a = [...pendingChain.add]; pendingChain.add.clear(); try { await sql(`UPDATE ${TB} SET is_chain_block = true WHERE hash = ANY($1::text[]) AND NOT is_chain_block`, [pgArray(a)]); } catch (e) { log('chain update failed', e.message); } }
if (pendingChain.remove.size) { const r = [...pendingChain.remove]; pendingChain.remove.clear(); try { await sql(`UPDATE ${TB} SET is_chain_block = false WHERE hash = ANY($1::text[]) AND is_chain_block`, [pgArray(r)]); } catch (e) { log('chain update failed', e.message); } }
}
async function tick(rpc) {
try { await tickInner(rpc); rpc.noteTick && rpc.noteTick(true); } catch (e) { rpc.noteTick && rpc.noteTick(false); }
}
async function tickInner(rpc) {
const now = Date.now();
rollRpcLoad(now);
let dag, info, peersRes, hps;
try {
[dag, info, peersRes, hps] = await Promise.all([
rpc.call('getBlockDagInfo', {}),
rpc.call('getInfo', {}),
rpc.call('getConnectedPeerInfo', {}),
rpc.call('estimateNetworkHashesPerSecond', { windowSize: 1000, startHash: null }).catch(() => null),
]);
} catch (e) { log('tick failed', e.message); throw e; }
if (dag.network) network = String(dag.network).startsWith('igneum') ? dag.network : `igneum-${dag.network}`;
nodeVersion = info.serverVersion || nodeVersion;
if (dag.sink && dag.sink !== sinkHash) { sinkHash = dag.sink; pendingChain.add.add(dag.sink); }
lastDagBlockCount = Number(dag.blockCount) || 0;
// self-check: the node keeps adding blocks but none reach us, so the subscription is dead; resubscribe, then give up
if (lastBlockAt && now - lastBlockAt > 60_000 && lastDagBlockCount > blockCountAtLastBlock && now - lastResubscribeAt > 60_000) {
lastResubscribeAt = now; resubscribeAttempts++;
const msg = `no blockAdded for ${Math.round((now - lastBlockAt) / 1000)} s while the node's block_count rose ${blockCountAtLastBlock} -> ${lastDagBlockCount}`;
if (resubscribeAttempts > 2) { log(msg + '; resubscribed twice without effect, exiting for a restart'); await recordEvent('observer', 'Observer lost the block feed; restarting'); process.exit(2); }
log(msg + `; resubscribing (attempt ${resubscribeAttempts})`);
try { await rpc.call('subscribe', { BlockAdded: {} }); await rpc.call('subscribe', { VirtualChainChanged: { include_accepted_transaction_ids: false } }); }
catch (e) { log('resubscribe failed', e.message); resubscribeAttempts++; }
}
const lagS = newestStoredTs ? Math.max(0, (now - newestStoredTs) / 1000) : null;
const queueDepth = inbox.length + pendingBlocks.length + pendingMerge.size;
// peers joined or left
const peers = peersRes.peerInfo || [];
const ids = new Map(peers.map(p => [p.id, `${p.address && p.address.ip}:${p.address && p.address.port}`]));
if (knownPeers) {
// never an address in a public event (review round 4, R4.6.2): the count is the information
for (const [id] of ids) if (!knownPeers.has(id)) recordEvent('peer_joined', `Peer joined (${ids.size} peers)`);
for (const [id] of knownPeers) if (!ids.has(id)) recordEvent('peer_left', `Peer left (${ids.size} peers)`);
}
knownPeers = ids;
// difficulty step over 5%
const diff = Number(dag.difficulty) || 0;
if (lastDifficultyEvent === null) lastDifficultyEvent = diff;
else if (lastDifficultyEvent > 0 && Math.abs(diff - lastDifficultyEvent) / lastDifficultyEvent > DIFFICULTY_STEP) {
const pct = ((diff - lastDifficultyEvent) / lastDifficultyEvent * 100).toFixed(1);
recordEvent('difficulty', `Difficulty ${diff > lastDifficultyEvent ? 'up' : 'down'} ${Math.abs(pct)}% to ${diff.toLocaleString('en-GB', { maximumFractionDigits: 0 })}`);
lastDifficultyEvent = diff;
}
// miners gone quiet
for (const [vk, m] of miners) if (!m.quiet && now - m.lastSeen > QUIET_AFTER_MS) { m.quiet = true; recordEvent('miner_quiet', `Miner ${short(vk)} has gone quiet (no block for 5 min)`); }
// finality v2: checkpoints and weights (a node from before the finality layer answers with an error; then null)
const finality = await finalityTick(rpc);
// proving v0: activation state and the 10-minute counts (a node without the RPCs gives {supported: false})
await provingTick();
const proving = await provingState();
try {
await sql(`INSERT INTO ${TS} (id, block_count, header_count, blue_score, difficulty, hashes_per_second_estimate, peers, mempool,
node_version, network, blocks_60s, blocks_per_minute, observer_started_at, updated_at, finality, observer_lag_s, queue_depth, proving, rpc_load, supply_check)
VALUES (1, $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11::jsonb, $12, now(), $13::jsonb, $14, $15, $16::jsonb, $17::jsonb, $18::jsonb)
ON CONFLICT (id) DO UPDATE SET block_count = EXCLUDED.block_count, header_count = EXCLUDED.header_count, blue_score = EXCLUDED.blue_score,
difficulty = EXCLUDED.difficulty, hashes_per_second_estimate = EXCLUDED.hashes_per_second_estimate, peers = EXCLUDED.peers,
mempool = EXCLUDED.mempool, node_version = EXCLUDED.node_version, network = EXCLUDED.network, blocks_60s = EXCLUDED.blocks_60s,
blocks_per_minute = EXCLUDED.blocks_per_minute, observer_started_at = EXCLUDED.observer_started_at, updated_at = now(),
finality = EXCLUDED.finality, observer_lag_s = EXCLUDED.observer_lag_s, queue_depth = EXCLUDED.queue_depth, proving = EXCLUDED.proving,
rpc_load = EXCLUDED.rpc_load, supply_check = EXCLUDED.supply_check`,
[dag.blockCount, dag.headerCount, (await rpc.call('getSinkBlueScore', {}).catch(() => ({}))).blueScore ?? null, diff,
hps ? hps.networkHashesPerSecond : localHashesPerSecond(), peers.length, info.mempoolSize ?? 0,
nodeVersion, network, blocks60s(now), JSON.stringify(blocksPerMinute(now)), STARTED_AT.toISOString(),
finality ? JSON.stringify(finality) : null, lagS === null ? null : Math.round(lagS * 10) / 10, queueDepth, JSON.stringify(proving),
JSON.stringify(rpcLoadState()), lastSupplyCheck ? JSON.stringify(lastSupplyCheck) : null]);
} catch (e) { log('state write failed', e.message); }
}
// ---------- Finality v2 (spec 03) ----------
const checkpointStates = new Map(); // index -> last state written; set BEFORE any await, so the two lock paths never both record
const checkpointDetailed = new Set(); // indices written with votesSeen (the FinalityLock notification carries none; the poll fills it in once)
let finalitySupported = null; // null = unknown, false = the node has no finality RPC
let lastWeightsAt = 0;
let lastWeights = null;
function pct(x) { return (Number(x) * 100).toFixed(1); }
async function recordLock(cp) {
// The lockEvent text the site shows: "checkpoint N locked (xx% of weight)"
const votes = cp.votesSeen !== undefined ? `${cp.votesSeen} votes of ${cp.voters} voters` : `${cp.voters} voters`;
recordEvent('checkpoint_locked', `checkpoint ${cp.index} locked (${pct(cp.fractionTotal)}% of weight, ${pct(cp.fractionActive)}% of active, ${votes}) at block ${short(cp.hash)}`);
}
async function upsertCheckpoint(cp) {
const locked = cp.state === 'locked';
try {
await sql(`INSERT INTO ${TC} (index, hash, blue_score, daa_score, state, signed_weight, active_weight, total_weight, fraction_active, fraction_total,
votes_seen, voters, aggregators, locked_at, updated_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13::text[], $14, now())
ON CONFLICT (index) DO UPDATE SET hash = EXCLUDED.hash, blue_score = EXCLUDED.blue_score, daa_score = EXCLUDED.daa_score, state = EXCLUDED.state,
signed_weight = EXCLUDED.signed_weight, active_weight = EXCLUDED.active_weight, total_weight = EXCLUDED.total_weight,
fraction_active = EXCLUDED.fraction_active, fraction_total = EXCLUDED.fraction_total, votes_seen = EXCLUDED.votes_seen, voters = EXCLUDED.voters,
aggregators = EXCLUDED.aggregators, locked_at = COALESCE(${TC}.locked_at, EXCLUDED.locked_at), updated_at = now()`,
// The FinalityLock notification carries no votesSeen; the 2-s poll fills it in
[cp.index, cp.hash, cp.blueScore, cp.daaScore, cp.state, cp.signedWeight, cp.activeWeight, cp.totalWeight, cp.fractionActive, cp.fractionTotal,
cp.votesSeen ?? 0, cp.voters ?? 0, pgArray(cp.aggregators || []), locked ? new Date().toISOString() : null]);
} catch (e) { log('checkpoint write failed', e.message); }
}
async function finalityTick(rpc) {
if (finalitySupported === false) return null;
let cps;
try { cps = await rpc.call('getFinalityCheckpoints', { last: 60 }); }
catch (e) {
if (finalitySupported === null) { log('node has no finality RPC (pre-v2 node):', e.message); finalitySupported = false; }
return null;
}
finalitySupported = true;
lastCheckpointsReport = cps;
if (!certBackfillDone) backfillCertificate(rpc, cps).catch(e => log('certificate backfill failed', e.message));
for (const cp of cps.checkpoints || []) {
const prev = checkpointStates.get(cp.index);
const changed = prev !== cp.state || prev === undefined;
const detail = cp.votesSeen !== undefined && !checkpointDetailed.has(cp.index);
if (changed || detail) {
// claim the transition before the first await: the FinalityLock notification handler checks this map synchronously,
// and until 4 October 2026 both paths recorded the same lock (two "locked" events 30 ms apart on the live feed)
checkpointStates.set(cp.index, cp.state);
if (cp.votesSeen !== undefined) checkpointDetailed.add(cp.index);
await upsertCheckpoint(cp);
if (changed && cp.state === 'locked' && prev !== 'locked') {
// a lock recorded over a known earlier state is a revisit (4 October 2026, 20:15:59: index 1319, locked 14 min
// before, was recorded again right after a restart; this line says what the seed held for it)
if (prev !== undefined) log(`checkpoint ${cp.index}: locked after a seeded or earlier state '${prev}'`);
await recordLock(cp);
}
} else if (cp.state !== 'locked' && (cp.index % 1 === 0)) {
// vote counts move while a checkpoint is open: refresh it
await upsertCheckpoint(cp);
}
}
for (const k of checkpointStates.keys()) if (k + 500 < Number(cps.nextIndex)) checkpointStates.delete(k);
for (const k of checkpointDetailed) if (k + 500 < Number(cps.nextIndex)) checkpointDetailed.delete(k);
const now = Date.now();
if (now - lastWeightsAt > 10_000) {
lastWeightsAt = now;
try {
const w = await rpc.call('getFinalityWeights', {});
lastWeights = {
checkpoint_index: w.checkpointIndex, checkpoint_hash: w.checkpointHash, daa_score: w.daaScore,
total_weight: w.totalWeight, active_weight: w.activeWeight, voters: w.voters,
keys: (w.keys || []).slice(0, 64).map(k => ({ id: short(k.keyHash), blocks: k.blocks, voter: k.voter, participation: k.participation, stripped_until_daa: k.strippedUntilDaa, revealed: !!k.pubkey })),
};
} catch (e) { if (!String(e.message).includes('no checkpoint')) log('getFinalityWeights failed', e.message); }
}
const locked = (cps.checkpoints || []).filter(c => c.state === 'locked');
const newest = locked.length ? locked[locked.length - 1] : null;
return {
params: cps.params, chain_id: cps.chainId, next_index: cps.nextIndex, finality_active: cps.finalityActive,
latest_locked_index: cps.latestLockedIndex, latest_locked_hash: cps.latestLockedHash,
latest_locked_blue_score: newest ? newest.blueScore : null,
weights: lastWeights,
};
}
// ---------- Finality v2 certificates (spec 03 C3): what the light client verifies ----------
// A block's coinbase extra data ends with the node's finality section `items || len_le32 || "IGNF"`; items are
// tagged 1 vote (280 bytes), 2 certificate, 3 evidence (560 bytes). A certificate is `index_le64 || checkpoint(32)
// || voter_count_le32 || bitmap_len_le32 || bitmap || signature(96) || aggregator(32) || proof(96)`
// (vendor/igneum-node/consensus/core/src/finality.rs, Certificate::write). Bit p of the bitmap (byte p/8, bit
// p%8) is position p of the canonical voter list at the checkpoint.
const pendingCerts = []; // { cert, carrier } seen in blocks, waiting for the next flush
const storedCerts = new Map(); // index -> checkpoint hash already in the table
let lastCheckpointsReport = null; // the last getFinalityCheckpoints answer: chain id and the locked list
let certBackfillDone = false;
let certsBusy = false;
const MAX_CHAIN_HEADERS = 200;
const ZERO_HASH = '0'.repeat(64);
const hexOf = u8 => Buffer.from(u8).toString('hex');
function certificatesIn(block) {
const tx = block.transactions && block.transactions[0];
if (!tx) return [];
const p = payloadBytes(tx.payload);
if (p.length < 19) return [];
const extra = p.subarray(19 + p[18]);
const n = extra.length;
if (n < 8 || String.fromCharCode(...extra.subarray(n - 4)) !== 'IGNF') return [];
const len = new DataView(extra.buffer, extra.byteOffset, n).getUint32(n - 8, true);
if (len + 8 > n) return [];
const body = extra.subarray(n - 8 - len, n - 8);
const out = []; let o = 0;
while (o < body.length) {
const tag = body[o];
if (tag === 1) o += 1 + 280;
else if (tag === 3) o += 1 + 560;
else if (tag === 2) { const c = parseCertificate(body, o + 1); if (!c) break; out.push(c); o += 1 + c.length; }
else break;
}
return out;
}
function parseCertificate(b, s) {
if (b.length < s + 48) return null;
const dv = new DataView(b.buffer, b.byteOffset, b.byteLength);
const voterCount = dv.getUint32(s + 40, true);
const bl = dv.getUint32(s + 44, true);
const end = s + 48 + bl + 96 + 32 + 96;
if (bl > 1 << 20 || b.length < end) return null;
return {
index: Number(dv.getBigUint64(s, true)), checkpoint: hexOf(b.subarray(s + 8, s + 40)), voterCount,
bitmap: hexOf(b.subarray(s + 48, s + 48 + bl)), signature: hexOf(b.subarray(s + 48 + bl, s + 144 + bl)),
aggregator: hexOf(b.subarray(s + 144 + bl, s + 176 + bl)), aggregatorProof: hexOf(b.subarray(s + 176 + bl, end)),
bytes: hexOf(b.subarray(s, end)), length: end - s,
};
}
// Stores a certificate with everything a verifier needs: the voter table and the header chain back to the previous
// locked checkpoint. The node reports the voter table at its latest determined checkpoint; captured as the lock
// lands (the next determination is 30 blue score away) it is the table at this checkpoint, and `voters_index`
// records which checkpoint it was read at either way.
async function completeCertificate(rpc, cert, carrier) {
if (storedCerts.get(cert.index) === cert.checkpoint) return;
const report = lastCheckpointsReport || await rpc.call('getFinalityCheckpoints', { last: 60 });
const cp = (report.checkpoints || []).find(c => Number(c.index) === cert.index);
if (cp && cp.hash !== cert.checkpoint) { log(`certificate at ${cert.index} is for ${short(cert.checkpoint)}, the node's checkpoint is ${short(cp.hash)}; not stored`); return; }
const w = await rpc.call('getFinalityWeights', {});
const voters = (w.keys || []).filter(k => k.voter)
.sort((a, b) => (a.keyHash < b.keyHash ? -1 : a.keyHash > b.keyHash ? 1 : 0))
.map(k => ({ pubkey: k.pubkey, weight: Number(k.blocks), participation: Number(k.participation) }));
if (voters.length !== cert.voterCount) log(`certificate at ${cert.index} names ${cert.voterCount} voters, the node's table at checkpoint ${w.checkpointIndex} has ${voters.length}`);
const prev = (report.checkpoints || []).filter(c => c.state === 'locked' && Number(c.index) < cert.index).sort((a, b) => Number(b.index) - Number(a.index))[0] || null;
const headers = []; let hash = cert.checkpoint; let complete = false;
while (headers.length < MAX_CHAIN_HEADERS) {
const b = (await rpc.call('getBlock', { hash, includeTransactions: false })).block;
headers.push(b.header);
if (prev && hash === prev.hash) { complete = true; break; }
hash = (b.verboseData && b.verboseData.selectedParentHash) || ZERO_HASH;
if (hash === ZERO_HASH) { complete = !prev; break; }
}
headers.reverse();
const top = headers[headers.length - 1];
try {
await sql(`INSERT INTO ${TX} (index, hash, blue_score, daa_score, chain_id, voter_count, bitmap_hex, signature_hex, aggregator, aggregator_proof_hex,
certificate_hex, carrier, voters, voters_index, total_weight, active_weight, prev_index, prev_hash, headers, headers_complete)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13::jsonb, $14, $15, $16, $17, $18, $19::jsonb, $20)
ON CONFLICT (index) DO NOTHING`,
[cert.index, cert.checkpoint, Number(top.blueScore), Number(top.daaScore), report.chainId || network, cert.voterCount, cert.bitmap, cert.signature,
cert.aggregator, cert.aggregatorProof, cert.bytes, carrier, JSON.stringify(voters), Number(w.checkpointIndex), Number(w.totalWeight), Number(w.activeWeight),
prev ? Number(prev.index) : null, prev ? prev.hash : null, JSON.stringify(headers), complete]);
storedCerts.set(cert.index, cert.checkpoint);
log(`certificate ${cert.index} stored: ${voters.length} voters (table at checkpoint ${w.checkpointIndex}), ${headers.length} headers${complete ? '' : ' (chain incomplete)'}, carried by ${short(carrier)}`);
} catch (e) { log('certificate write failed', e.message); }
}
async function flushCertificates(rpc) {
if (certsBusy || !pendingCerts.length) return;
certsBusy = true;
try {
while (pendingCerts.length) {
const { cert, carrier } = pendingCerts.shift();
try { await completeCertificate(rpc, cert, carrier); } catch (e) { log('certificate store failed', e.message); }
}
} finally { certsBusy = false; }
}
// On start: the newest locked checkpoint without a stored certificate is looked up in the blocks after it (a
// certificate is carried by the first templates after the lock). One pass, up to three pages of getBlocks.
async function backfillCertificate(rpc, report) {
certBackfillDone = true;
const newest = (report.checkpoints || []).filter(c => c.state === 'locked').sort((a, b) => Number(b.index) - Number(a.index))[0];
if (!newest || storedCerts.get(Number(newest.index)) === newest.hash) return;
let low = newest.hash; let seen = 0;
for (let page = 0; page < 3; page++) {
const blocks = (await rpc.call('getBlocks', { lowHash: low, includeBlocks: true, includeTransactions: true })).blocks || [];
for (const b of blocks) {
seen++;
for (const c of certificatesIn(b)) if (c.index === Number(newest.index)) { await completeCertificate(rpc, c, b.header.hash); return; }
}
if (blocks.length < 2) break;
low = blocks[blocks.length - 1].header.hash;
}
log(`no block within ${seen} of checkpoint ${newest.index} carries its certificate`);
}
// ---------- Proving v0 (spec 7.7): shard plans and proof records over the execution layer's JSON-RPC ----------
// The chain block hash is the same hash on both layers (the executor keys its segments by it), so every lookup is
// {blockHash}. The EVM node may trail the observer's node by a segment or two: a plan that is "not found" is retried
// with a growing delay, up to 30 attempts. A node without the RPCs answers -32601 once and the feed is marked
// unsupported (rechecked every 5 minutes, so an upgraded node is picked up without a restart).
// PROOF_WINDOW_MS: how long a planned block keeps being read for its records. 30 minutes (7 October 2026, was 10): the fleet's
// proofs land a median 258 s after the block and the tail runs past 10 minutes, and a block dropped from the open set never
// moves past 'proving' on the page. The open set at 0.45 blocks a second is about 800 blocks; the batch of 80 reads each one
// every 10 ticks.
const PLAN_BATCH = 24, RECORD_BATCH_ACTIVE = 80, RECORD_BATCH_INACTIVE = 10, PROOF_WINDOW_MS = 30 * 60_000;
// OBSERVER_CLAIMS: '1' reads open claims through igneum_getProofClaims; a method name reads them through that method; unset
// or '0' reads none (the records reply's own claims field is used whenever the node sends one, flag or not)
let CLAIMS_RPC = process.env.OBSERVER_CLAIMS === '1' ? 'igneum_getProofClaims' : /^[a-z_]+$/i.test(process.env.OBSERVER_CLAIMS || '') && process.env.OBSERVER_CLAIMS !== '0' ? process.env.OBSERVER_CLAIMS : null;
const pendingPlans = new Map(); // chain block hash -> { tries, nextAt }
const openBlocks = new Map(); // chain block hash -> { number, daa, ts, shards: [state...], provers, polledAt, done }
const daaOfNumber = new Map(); // chain block number -> daa score, from the plans fetched (carriers are later chain blocks)
const provers = new Set(); // prover id8s seen this run (one event each)
let provingSupported = null; // null unknown, false = no RPC, true
let provingRecheckAt = 0;
let provingStatus = null; // the last igneum_getProvingStatus answer
let provingStatusAt = 0; // when it was read (ms); live_state.proving.status_at, the API's status_age_s
let provingReason = null; // why supported is false
let provingWasActive = null;
let plansBusy = false, recordsBusy = false, firstPaidSeen = false;
let lastProvingStats = null, lastProvingStatsAt = 0;
const hx = v => (v === null || v === undefined ? null : (typeof v === 'string' && v.startsWith('0x') ? parseInt(v, 16) : Number(v)));
const weiStr = v => { if (v === null || v === undefined) return null; try { return BigInt(v).toString(); } catch { return null; } };
async function evm(method, params = []) {
rpcLoad.evm++;
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(8000) });
const j = await r.json();
if (j.error) { const e = new Error(`${method}: ${j.error.message}`); e.code = j.error.code; throw e; }
return j.result;
}
function queuePlan(hash) { if (!openBlocks.has(hash) && !pendingPlans.has(hash)) pendingPlans.set(hash, { tries: 0, nextAt: 0 }); }
function dropPlan(hash) { pendingPlans.delete(hash); openBlocks.delete(hash); }
async function provingTick() {
const now = Date.now();
if (provingSupported === false && now < provingRecheckAt) return;
try {
const s = await evm('igneum_getProvingStatus');
if (provingSupported !== true) log(`proving RPCs available at ${EVM} (activation ${s.activationDaa === null ? 'not set' : hx(s.activationDaa)}, verifier ${s.verifier})`);
provingSupported = true; provingReason = null; provingStatus = s; provingStatusAt = Date.now();
if (provingWasActive === false && s.active) recordEvent('proving', `Proving layer active: payouts start at DAA ${hx(s.activationDaa)}`);
provingWasActive = !!s.active;
} catch (e) {
const absent = e.code === -32601;
if (provingSupported !== false || !absent) log(absent ? `node at ${EVM} has no proving RPCs` : `proving status failed: ${e.message}`);
provingSupported = false; provingStatus = null; provingReason = absent ? 'node without proving' : `evm rpc unreachable (${e.message.slice(0, 60)})`;
provingRecheckAt = now + (absent ? 5 * 60_000 : 20_000);
}
}
async function flushPlans() {
if (plansBusy || provingSupported !== true || !pendingPlans.size) return;
plansBusy = true;
try {
const now = Date.now();
const due = [...pendingPlans].filter(([, p]) => p.nextAt <= now).slice(0, PLAN_BATCH).map(([h]) => h);
const rows = [];
for (let i = 0; i < due.length; i += 4) {
await Promise.all(due.slice(i, i + 4).map(async h => {
const p = pendingPlans.get(h); if (!p) return;
try {
const plan = await evm('igneum_getShardPlan', [{ blockHash: '0x' + h }]);
pendingPlans.delete(h);
const number = hx(plan.number), daa = hx(plan.daaScore);
daaOfNumber.set(number, daa);
const shards = (plan.shards || []).map(() => 'planned');
openBlocks.set(h, { number, daa, ts: now, shards, provers: shards.map(() => null), polledAt: 0, done: false });
(plan.shards || []).forEach((s, i) => rows.push({ hash: h, shard: i, shards: shards.length, number, daa, pgas: hx(s.pgas) }));
} catch (e) {
p.tries++;
if (p.tries > 30 || !/not found/.test(e.message)) { pendingPlans.delete(h); if (!/not found/.test(e.message)) log('shard plan failed', h.slice(0, 8), e.message); }
else p.nextAt = now + Math.min(30_000, 1000 * p.tries);
}
}));
}
if (rows.length) {
const params = [], values = [];
for (const r of rows) {
params.push(r.hash, r.shard, r.shards, r.number, r.daa, r.pgas);
const n = params.length;
values.push(`($${n - 5},$${n - 4},$${n - 3},$${n - 2},$${n - 1},$${n})`);
}
// block_ts comes from the block table (the header time the page's axis uses), one update for the batch
await sql(`INSERT INTO ${TP} (block_hash, shard, shards, block_number, block_daa, pgas) VALUES ${values.join(',')} ON CONFLICT (block_hash, shard) DO NOTHING`, params);
await sql(`UPDATE ${TP} p SET block_ts = b.timestamp_ms FROM ${TB} b WHERE b.hash = p.block_hash AND p.block_ts IS NULL AND p.block_hash = ANY($1::text[])`, [pgArray([...new Set(rows.map(r => r.hash))])]).catch(() => { });
// explorer: the chain block's number (the EVM block number), one update for the batch, no extra RPC
const nums = [...new Map(rows.map(r => [r.hash, r.number]))];
await sql(`UPDATE ${TB} b SET number = v.n::bigint FROM (SELECT unnest($1::text[]) AS h, unnest($2::text[]) AS n) v WHERE b.hash = v.h AND b.number IS NULL`, [pgArray(nums.map(x => x[0])), pgArray(nums.map(x => String(x[1])))]).catch(e => log('number update failed', e.message));
}
} catch (e) { log('plan flush failed', e.message); }
finally { plansBusy = false; }
}
// Per shard: paid (a carrying segment paid it) > verified (the proof verified in the pool, or the record carried and
// checked by consensus, which is what happens before activation) > proving (a record in the pool) > planned.
function shardStates(block, r) {
const n = block.shards.length, states = block.shards.map(() => 'planned'), out = states.map(() => ({ prover: null, verified: null, carrier: null, carrierNumber: null, payout: null }));
// an open claim (a prover has taken the shard, no proof yet) reads 'proving' from the claim until the proof lands: the
// records reply's own claims field when the node carries one, or the claims RPC behind OBSERVER_CLAIMS (7 October 2026)
for (const c of r.claims || []) { const i = c.shard; if (i >= n || states[i] !== 'planned') continue; states[i] = 'proving'; out[i] = { ...out[i], prover: short(String(c.keyHash || c.prover || '').replace(/^0x/, '')) }; }
for (const e of r.pool || []) { const i = e.shard; if (i >= n) continue; const rank = e.verified === true ? 2 : 1; if (rank > ['planned', 'proving', 'verified', 'paid'].indexOf(states[i])) { states[i] = rank === 2 ? 'verified' : 'proving'; out[i] = { ...out[i], prover: short(String(e.keyHash || '').replace(/^0x/, '')), verified: e.verified, carrier: e.includedIn ? String(e.includedIn).replace(/^0x/, '') : null }; } }
for (const c of r.carried || []) { const i = c.shard; if (i >= n || !c.valid) continue; if (states[i] !== 'paid') { states[i] = 'verified'; out[i] = { ...out[i], prover: short(String(c.keyHash || '').replace(/^0x/, '')), carrier: String(c.carrier || '').replace(/^0x/, ''), carrierNumber: hx(c.carrierNumber) }; } }
(r.paid || []).forEach((p, i) => { if (!p || i >= n) return; states[i] = 'paid'; out[i] = { ...out[i], prover: short(String(p.keyHash || '').replace(/^0x/, '')), carrierNumber: hx(p.carrierNumber), payout: weiStr(p.wei) }; });
return { states, out };
}
async function pollRecords() {
if (recordsBusy || provingSupported !== true) return;
recordsBusy = true;
try {
const now = Date.now();
for (const [h, b] of openBlocks) if (b.done || now - b.ts > PROOF_WINDOW_MS) openBlocks.delete(h);
const active = !!(provingStatus && provingStatus.active);
const batch = [...openBlocks].sort((a, b) => a[1].polledAt - b[1].polledAt).slice(0, active ? RECORD_BATCH_ACTIVE : RECORD_BATCH_INACTIVE);
const updates = [];
for (let i = 0; i < batch.length; i += 4) {
await Promise.all(batch.slice(i, i + 4).map(async ([h, b]) => {
b.polledAt = now;
let r; try { r = await evm('igneum_getProofRecords', [{ blockHash: '0x' + h }]); } catch (e) { if (!/not found/.test(e.message)) log('proof records failed', h.slice(0, 8), e.message); return; }
// the claims read, behind the flag until the node ships the RPC (main's ask of 7 October 2026 on the 0.3.20 line); a
// method the node lacks turns the flag off for the run instead of logging on every block
if (CLAIMS_RPC && !b.done && !(r && r.claims)) { try { const c = await evm(CLAIMS_RPC, [{ blockHash: '0x' + h }]); r = { ...(r || {}), claims: Array.isArray(c) ? c : (c && c.claims) || [] }; } catch (e) { if (/method not found|unknown method|not supported/i.test(e.message)) { log('claims RPC', CLAIMS_RPC, 'not on this node; claims off for the run'); CLAIMS_RPC = null; } else if (!/not found/.test(e.message)) log('proof claims failed', h.slice(0, 8), e.message); } }
const { states, out } = shardStates(b, r);
for (let s = 0; s < states.length; s++) {
const key = `${states[s]}|${out[s].prover}|${out[s].carrierNumber}|${out[s].carrier}`;
if (b.provers[s] === key) continue;
b.provers[s] = key; b.shards[s] = states[s];
let carrierDaa = out[s].carrierNumber !== null ? daaOfNumber.get(out[s].carrierNumber) ?? null : null;
if (carrierDaa === null && out[s].carrierNumber !== null) { try { const p = await evm('igneum_getShardPlan', ['0x' + out[s].carrierNumber.toString(16)]); carrierDaa = hx(p.daaScore); daaOfNumber.set(out[s].carrierNumber, carrierDaa); } catch { } }
updates.push({ hash: h, shard: s, state: states[s], prover: out[s].prover, verified: out[s].verified, carrier: out[s].carrier, carrierNumber: out[s].carrierNumber, carrierDaa, lag: carrierDaa !== null && b.daa !== null ? carrierDaa - b.daa : null, payout: out[s].payout });
if (out[s].prover && !provers.has(out[s].prover)) { provers.add(out[s].prover); recordEvent('prover_seen', `Prover ${out[s].prover} seen (shard ${s} of block ${short(h)})`); }
if (states[s] === 'paid' && !firstPaidSeen) { firstPaidSeen = true; recordEvent('proving', `First paid shard seen: block ${short(h)} shard ${s}, carried ${updates[updates.length - 1].lag ?? '?'} DAA later`); }
}
if (states.every(s => s === 'paid')) b.done = true;
}));
}
for (const u of updates) {
await sql(`UPDATE ${TP} SET state = $3, prover = $4, verified = $5, carried_by = $6, carrier_number = $7, carrier_daa = $8, lag_daa = $9, payout_wei = COALESCE($10, payout_wei), updated_at = now()
WHERE block_hash = $1 AND shard = $2`, [u.hash, u.shard, u.state, u.prover, u.verified, u.carrier, u.carrierNumber, u.carrierDaa, u.lag, u.payout]).catch(e => log('proof row update failed', e.message));
}
} catch (e) { log('records poll failed', e.message); }
finally { recordsBusy = false; }
}
// live_state.proving: activation and the last 10 minutes (counts refreshed every 10 s; the status every tick)
async function provingState() {
if (provingSupported !== true) return { supported: false, reason: provingReason || 'not checked yet', evm_rpc: EVM };
const s = provingStatus, now = Date.now();
if (now - lastProvingStatsAt > 10_000) {
lastProvingStatsAt = now;
try {
const rows = await sql(`WITH w AS (SELECT * FROM ${TP} WHERE received_at > now() - interval '10 minutes')
SELECT (SELECT count(DISTINCT block_hash) FROM w)::int AS blocks,
(SELECT count(*) FROM (SELECT block_hash FROM w GROUP BY block_hash HAVING bool_and(state IN ('verified','paid'))) f)::int AS blocks_full,
(SELECT count(*) FROM w WHERE state IN ('verified','paid'))::int AS shards_proven,
(SELECT count(*) FROM w WHERE state = 'paid')::int AS shards_paid,
(SELECT count(DISTINCT prover) FROM w WHERE prover IS NOT NULL AND state IN ('verified','paid'))::int AS provers,
(SELECT percentile_cont(0.5) WITHIN GROUP (ORDER BY lag_daa) FROM w WHERE lag_daa IS NOT NULL)::float AS median_lag`);
lastProvingStats = rows[0] || null;
} catch (e) { log('proving stats failed', e.message); }
}
const st = lastProvingStats || {};
return {
supported: true, active: !!s.active, activation_daa: s.activationDaa === null ? null : hx(s.activationDaa), tip_daa: hx(s.tipDaa),
verifier: s.verifier, pool: s.pool || null, paid_shards_total: s.paidShards ?? null, shard_budget_pgas: hx(s.shardBudget),
blocks_10m: st.blocks ?? 0, blocks_fully_proven_10m: st.blocks_full ?? 0, shards_proven_10m: st.shards_proven ?? 0, shards_paid_10m: st.shards_paid ?? 0,
// DAA score of the carrying chain block minus the proven block's; the devnet targets one DAA step per second
median_proof_lag_s: st.median_lag === null || st.median_lag === undefined ? null : Math.round(st.median_lag * 10) / 10,
provers_10m: st.provers ?? 0, open_blocks: openBlocks.size, pending_plans: pendingPlans.size, evm_rpc: EVM,
// 8 October 2026: when the status above was read, so the API can call an old snapshot stale (site/lib/proving-health.mjs)
status_at: provingStatusAt ? new Date(provingStatusAt).toISOString() : null,
};
}
// ---------- Supply check (spec 2.5 against the chain) ----------
// Two comparisons over the newest stored blocks, once an hour and 20 s after start: (a) the subsidy each payload
// declares equals blockSubsidy(daa, 1) of site/lib/emission.mjs; (b) each block's outputs sum to the declared subsidy
// of every block it merges (utxo_validation.rs:176 pays a merged block what its own payload declares; blues and reds
// inside the DAA window alike). Written to live_state.supply_check; /api/supply shows it beside the rule's total.
const SUPPLY_CHECK_EVERY_MS = 60 * 60_000, SUPPLY_SAMPLE = 500;
const VENDOR_SHARE_EVERY_MS = 60_000; // live_state.vendor_share, written by tools/observer/vendor-share.mjs
let lastSupplyCheck = null;
async function supplyCheck() {
try {
const rows = await sql(`SELECT hash, daa_score, subsidy_sompi, paid_sompi, detail->'mergeset' AS mergeset FROM ${TB}
WHERE subsidy_sompi IS NOT NULL ORDER BY received_at DESC LIMIT ${SUPPLY_SAMPLE}`);
const byHash = new Map(rows.map(r => [r.hash, r]));
let ruleOk = 0, ruleBad = 0, sumOk = 0, sumBad = 0, sumSkipped = 0; const examples = [];
for (const r of rows) {
const want = blockSubsidy(Number(r.daa_score), 1);
if (BigInt(r.subsidy_sompi) === want) ruleOk++; else { ruleBad++; if (examples.length < 3) examples.push({ kind: 'rule', block: short(r.hash), daa: Number(r.daa_score), declared: String(r.subsidy_sompi), rule: want.toString() }); }
const ms = r.mergeset || { blues: [], reds: [] };
const merged = [...(ms.blues || []), ...(ms.reds || [])];
if (!merged.length || merged.some(h => !byHash.has(h))) { sumSkipped++; continue; } // a merged block outside the sample: no verdict
const expect = merged.reduce((t, h) => t + BigInt(byHash.get(h).subsidy_sompi), 0n);
if (BigInt(r.paid_sompi) === expect) sumOk++; else { sumBad++; if (examples.length < 3) examples.push({ kind: 'sum', block: short(r.hash), paid: String(r.paid_sompi), expected: expect.toString(), merged: merged.length }); }
}
lastSupplyCheck = { checked_at: new Date().toISOString(), sampled: rows.length, rule_match: ruleOk, rule_mismatch: ruleBad, sum_match: sumOk, sum_mismatch: sumBad, sum_skipped: sumSkipped, examples, bps: 1 };
log(`supply check: ${rows.length} blocks, payload subsidy = rule ${ruleOk}/${ruleOk + ruleBad}, outputs = merged subsidies ${sumOk}/${sumOk + sumBad} (${sumSkipped} without a verdict)`);
} catch (e) { log('supply check failed', e.message); }
}
async function prune() {
try {
await sql(`DELETE FROM ${TB} WHERE received_at < now() - ($1 || ' hours')::interval`, [String(RETAIN_HOURS)]);
await sql(`DELETE FROM ${TP} WHERE received_at < now() - ($1 || ' hours')::interval`, [String(RETAIN_HOURS)]);
await sql(`DELETE FROM ${TE} WHERE ts < now() - interval '7 days'`);
await sql(`DELETE FROM ${TC} WHERE updated_at < now() - interval '7 days'`);
} catch (e) { log('prune failed', e.message); }
}
async function seedFromDb() {
try {
const rows = await sql(`SELECT vote_key_hash, max(received_at) AS last FROM ${TB} WHERE vote_key_hash IS NOT NULL AND received_at > now() - interval '1 day' GROUP BY 1`);
for (const r of rows) miners.set(r.vote_key_hash, { lastSeen: new Date(r.last).getTime(), quiet: Date.now() - new Date(r.last).getTime() > QUIET_AFTER_MS });
const mins = await sql(`SELECT (floor(timestamp_ms / 60000.0) * 60000)::bigint AS m, count(*)::int AS n FROM ${TB} WHERE timestamp_ms > (extract(epoch from now()) * 1000 - 61 * 60000) GROUP BY 1`);
for (const r of mins) minuteCounts.set(Number(r.m), r.n);
const recent = await sql(`SELECT timestamp_ms FROM ${TB} WHERE timestamp_ms > (extract(epoch from now()) * 1000 - 120000)`);
for (const r of recent) recentArrivals.push(Number(r.timestamp_ms));
const newest = await sql(`SELECT max(timestamp_ms)::bigint AS t FROM ${TB}`);
if (newest.length && newest[0].t) newestStoredTs = Number(newest[0].t);
const st = await sql(`SELECT difficulty FROM ${TS} WHERE id = 1`);
const cps = await sql(`SELECT index, state FROM ${TC} WHERE updated_at > now() - interval '1 day'`);
for (const r of cps) checkpointStates.set(Number(r.index), r.state);
{ const by = {}; for (const r of cps) by[r.state] = (by[r.state] || 0) + 1; log(`seeded ${cps.length} checkpoint states: ${Object.entries(by).map(([k, v]) => `${k} ${v}`).join(', ') || 'none'}`); }
const certs = await sql(`SELECT index, hash FROM ${TX}`);
for (const r of certs) storedCerts.set(Number(r.index), r.hash);
// proving: keep polling the blocks of the last 10 minutes that are not fully paid, so a restart loses nothing
const open = await sql(`SELECT block_hash, block_number, block_daa, extract(epoch from received_at) * 1000 AS ts, array_agg(state ORDER BY shard) AS states, array_agg(prover ORDER BY shard) AS provers
FROM ${TP} WHERE received_at > now() - interval '10 minutes' GROUP BY 1, 2, 3, 4`);
for (const r of open) {
daaOfNumber.set(Number(r.block_number), Number(r.block_daa));
openBlocks.set(r.block_hash, { number: Number(r.block_number), daa: Number(r.block_daa), ts: Number(r.ts), shards: r.states, provers: r.states.map(() => null), polledAt: 0, done: r.states.every(s => s === 'paid') });
for (const p of r.provers) if (p) provers.add(p);
}
if (open.some(r => r.states.some(s => s === 'paid'))) firstPaidSeen = true;
if (st.length && st[0].difficulty) lastDifficultyEvent = Number(st[0].difficulty);
} catch (e) { log('seed failed', e.message); }
}
async function main() {
await setupSchema();
await seedFromDb();
log(`observer started, rpc ${RPC}, evm rpc ${EVM}, keeping ${RETAIN_HOURS} h of blocks, ${miners.size} miners known, ${openBlocks.size} blocks open for proofs`);
const rpc = new Rpc(RPC);
// The socket handler only enqueues; the drain loop does the work in bounded batches so the socket is always read
rpc.onNotification = (method, params) => { inbox.push([method, params]); if (!draining) { draining = true; setImmediate(drain); } };
function drain() {
const batch = inbox.splice(0, 500);
for (const [m, p] of batch) { try { handle(m, p); } catch (e) { log('notification failed', e.message); } }
if (inbox.length) setImmediate(drain); else draining = false;
}
function handle(method, params) {
const inner = params && (params.BlockAdded || params.VirtualChainChanged || params);
if (method === 'blockAddedNotification' && inner && inner.block) onBlock(inner.block);
else if (method === 'finalityLockNotification' && inner) {
// The node's lockEvent: write the checkpoint and the event at once; the 2 s poll fills in the rest
const cp = inner.FinalityLock || inner;
if (cp.index !== undefined && checkpointStates.get(Number(cp.index)) !== 'locked') {
checkpointStates.set(Number(cp.index), 'locked');
upsertCheckpoint({ ...cp, state: 'locked', aggregators: [] });
recordLock(cp);
}
}
else if (method === 'virtualChainChangedNotification' && inner) {
for (const h of inner.addedChainBlockHashes || []) { pendingChain.remove.delete(h); pendingChain.add.add(h); pendingUnmerge.delete(h); if (!pendingMerge.has(h)) pendingMerge.set(h, mergesets.get(h) || null); queuePlan(h); }
for (const h of inner.removedChainBlockHashes || []) { pendingChain.add.delete(h); pendingChain.remove.add(h); pendingMerge.delete(h); pendingUnmerge.add(h); dropPlan(h); }
}
}
let connected = false;
const connect = async () => {
while (!(await rpc.connect())) { log('rpc connect failed, retrying in 3 s'); await new Promise(r => setTimeout(r, 3000)); }
try {
const srv = await rpc.call('getServerInfo', {});
nodeVersion = srv.serverVersion || null;
const net = String(srv.networkId || 'devnet');
network = net.startsWith('igneum') ? net : `igneum-${net}`;
addressPrefix = network.includes('devnet') ? 'igneumdev' : network.includes('testnet') ? 'igneumtest' : network.includes('simnet') ? 'igneumsim' : 'igneum';
await rpc.call('subscribe', { BlockAdded: {} });
await rpc.call('subscribe', { VirtualChainChanged: { include_accepted_transaction_ids: false } });
await rpc.call('subscribe', { FinalityLock: {} }).catch(e => log('FinalityLock subscription refused (pre-v2 node?):', e.message));
log(`subscribed on ${network}, node ${nodeVersion}`);
if (connected) recordEvent('observer', 'Observer reconnected to the node'); else recordEvent('observer', `Observer started against ${network}`);
connected = true;
} catch (e) { log('subscribe failed', e.message); }
};
rpc.onClose = () => { setTimeout(connect, 2000); };
await connect();
// Watchdog (5 October 2026): after the observer node was restarted for 0.3.9 the reporter logged "rpc not connected" for
// an hour without reconnecting (the socket's close never resolved into a new connect). Ten failed ticks in a row force
// the socket shut and a fresh connect; the live page must never go stale while the node answers.
let failedTicks = 0;
rpc.noteTick = (ok) => {
failedTicks = ok ? 0 : failedTicks + 1;
if (failedTicks === 10) {
log('watchdog: 10 failed ticks, forcing a reconnect');
failedTicks = 0;
try { rpc.ws && rpc.ws.close(); } catch { }
rpc.open = false;
setTimeout(connect, 1000);
}
};
setInterval(async () => { if (flushBusy) return; flushBusy = true; try { await flushBlocks(); await flushChain(); } finally { flushBusy = false; } }, FLUSH_EVERY_MS);
setInterval(() => flushColors(rpc), FLUSH_EVERY_MS);
setInterval(() => flushCertificates(rpc), FLUSH_EVERY_MS);
setInterval(flushPlans, FLUSH_EVERY_MS);
setInterval(pollRecords, STATE_EVERY_MS);
setInterval(() => tick(rpc), STATE_EVERY_MS);
setInterval(prune, PRUNE_EVERY_MS);
setInterval(supplyCheck, SUPPLY_CHECK_EVERY_MS); setTimeout(supplyCheck, 20_000);
// Detector (6 Oct 2026): once a minute, live_state.detector and live_events kind `detector`; tools/observer/detector.mjs
const detector = makeDetector({ sql, TB, TS, recordEvent, log }, detectorOpts());
setInterval(() => detector().catch(e => log('detector failed', e.message)), 60_000); setTimeout(() => detector().catch(e => log('detector failed', e.message)), 40_000);
setInterval(() => vendorShareTick(sql, T, log), VENDOR_SHARE_EVERY_MS); setTimeout(() => vendorShareTick(sql, T, log), 30_000);
tick(rpc); prune();
}
process.on('unhandledRejection', e => log('unhandled', e && e.message));
main().catch(e => { console.error(e); process.exit(1); });