igneum/tools/observer/observer.mjs
igneum-labs 771e00c3e1 Live page: shards per block (proving v0), three strips on one time axis (blocks, finality bar, proving), observer proof feed, hero glint
Observer (tools/observer/observer.mjs): reads the execution layer's JSON-RPC of a node on the proving build
(IGNEUM_EVM_RPC, default the Mac app's node 26800): every chain block's shard plan as it joins the chain
(igneum_getShardPlan by blockHash, one live_proofs row per shard, planned), the proof records of the chain
blocks of the last 10 minutes polled in rotation (igneum_getProofRecords, four in flight, 40 blocks per tick
while active, 10 before activation): proving (in the pool), verified (SP1 proof verified, or carried and checked
by consensus), paid (a carrying segment paid it), with the prover's id8, the carrier, lag in DAA and the payout.
live_state.proving = {supported, active, activation_daa, tip_daa, verifier, pool, blocks_10m,
blocks_fully_proven_10m, shards_proven_10m, shards_paid_10m, median_proof_lag_s, provers_10m}. A node without
the RPCs gives supported false (rechecked every 5 min); an unreachable endpoint is retried every 20 s. Events:
proving (activation, first paid shard), prover_seen. Additive schema (live_proofs, live_state.proving).

API (site/api/live.mjs): proving, and per block shards: [{i, n, state, prover, lag, payout (IGN), pgas}] and
proven; ?window=N (30 to 300 s) for the page's diagnostic long view; LIVE_TABLE_PREFIX reads a test observer's
tables.

Live page (site/live.html), the design change of 4 Oct 2026: three thin strips sharing one time axis, newest at
the right. BLOCKS keeps the per-miner lanes, chain path, blue/red/pending colouring, arrival glow and tooltips;
the lock ring, dashed lock line and final band leave it. FINALITY is an 18 px bar: ember wash = final (up to the
newest locked checkpoint on screen), molten tick = locked checkpoint, faint = proposed, one label at the newest
lock ("locked #522, 12 s ago"); while finality is not active it reads "finality paused: N% of weight silent" and
nothing else (R4.6.3). PROVING shows one cell per shard under each chain block, outline (planned), molten
(proving), prover colour (verified), tick (paid), a dashed "proofs land N s behind the tip" line, or the one
honest line before activation ("Proving layer: not yet activated on this devnet; activation at DAA N" / "node
without proving"). Header stats: on screen, chain, identities, last lock, proven. Legend: one line per strip.
Hover and tap tooltips on blocks and cells (block, shard, prover, lag, payout). Lanes snap on resize (they used
to ease from a zero-height layout). Phone width, no horizontal scroll; draw 0.6 ms avg, 1 ms max with 110 blocks
on screen (playwright, 1280 px).

Hero (site/index.html): a faint second glint behind a real block once every shard of it is verified, only while
the proving layer is active; pace and sampling untouched.

Verified on the private 3-node proving network (tools/proving-v0/run.mjs --network-only, activation 60) with a
CPU prover loop signing as v0/v1/v2: records relayed, verified on node 0, carried and paid (block 155 by 405,
lag 259 DAA, 0.634 IGN); screenshots in docs/design/live-proving (devnet before activation at 1280 and 375 px,
test network active, the ?window=300 view with paid cells). The live devnet shows the "not yet activated;
activation not set" line once the observer runs this build.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-04 17:39:28 +00:00

965 lines
58 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.
// 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';
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'`,
`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`,
// 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 = {}) {
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);
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;
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,
});
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'];
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}${c === 'parent_hashes' ? '::text[]' : ''}`); }
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) {
const now = Date.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); return; }
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)
VALUES (1, $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11::jsonb, $12, now(), $13::jsonb, $14, $15, $16::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`,
[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)]);
} catch (e) { log('state write failed', e.message); }
}
// ---------- Finality v2 (spec 03) ----------
const checkpointStates = new Map(); // index -> last state written
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);
if (prev !== cp.state || prev === undefined) {
await upsertCheckpoint(cp);
if (cp.state === 'locked' && prev !== 'locked') await recordLock(cp);
checkpointStates.set(cp.index, cp.state);
} 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);
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).
const PLAN_BATCH = 24, RECORD_BATCH_ACTIVE = 80, RECORD_BATCH_INACTIVE = 10, PROOF_WINDOW_MS = 10 * 60_000;
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 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 = []) {
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;
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(() => { });
}
} 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 }));
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; }
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,
};
}
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 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();
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);
tick(rpc); prune();
}
process.on('unhandledRejection', e => log('unhandled', e && e.message));
main().catch(e => { console.error(e); process.exit(1); });