#!/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) // // Tables (created on start if missing): live_blocks, live_state, live_events, live_checkpoints, live_certificates. See README.md. // 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. 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`; 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)`, `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`, `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 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); } } function noteArrival(now) { recentArrivals.push(now); const minute = Math.floor(now / 60_000) * 60_000; minuteCounts.set(minute, (minuteCounts.get(minute) || 0) + 1); for (const k of minuteCounts.keys()) if (k < now - 61 * 60_000) minuteCounts.delete(k); } function blocks60s(now) { while (recentArrivals.length && recentArrivals[0] < now - 60_000) recentArrivals.shift(); return recentArrivals.length; } 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; 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); 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); } catch (e) { log('block insert failed', e.message); pendingBlocks.unshift(...rows.slice(0, 50)); } if (pendingBlocks.length) await flushBlocks(); } 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); } // 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) { for (const [id, addr] of ids) if (!knownPeers.has(id)) recordEvent('peer_joined', `Peer joined: ${addr}`); for (const [id, addr] of knownPeers) if (!ids.has(id)) recordEvent('peer_left', `Peer left: ${addr}`); } 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); 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) VALUES (1, $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11::jsonb, $12, now(), $13::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`, [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]); } 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), key_hash: 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 => ({ key_hash: k.keyHash, 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`); } async function prune() { try { await sql(`DELETE FROM ${TB} 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(extract(epoch from received_at) / 60) * 60000)::bigint AS m, count(*)::int AS n FROM ${TB} WHERE received_at > now() - interval '61 minutes' GROUP BY 1`); for (const r of mins) minuteCounts.set(Number(r.m), r.n); 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); 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}, keeping ${RETAIN_HOURS} h of blocks, ${miners.size} miners known`); const rpc = new Rpc(RPC); rpc.onNotification = (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); } for (const h of inner.removedChainBlockHashes || []) { pendingChain.add.delete(h); pendingChain.remove.add(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 () => { await flushBlocks(); await flushChain(); await flushCertificates(rpc); }, FLUSH_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); });