igneum/tools/observer/observer.mjs
igneum-labs adc52b9c83 Observer: decoupled ingest, lag metric, per-minute by header time, feed self-check, restart loop
After the 4 Oct 2026 stall (no block stored from 11:57 to 13:15 UTC, then 7,022
blocks in two minutes). Notifications only enqueue; a drain loop handles them
in bounded batches. Block flush, colour marking and certificate work each run
on their own timer and never wait on one another. Mergesets come from the
notification's verbose data (bounded cache); getBlock only on a miss, four at
a time. live_state gains observer_lag_s and queue_depth; the API serves them;
the page shows "observer N s behind" past 30 s instead of waiting for the
first block. blocks_per_minute and blocks_60s are bucketed by the block's own
timestamp and reseeded from the table on start, 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 and tools/observer/run.sh restarts it.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-04 13:26:46 +00:00

754 lines
43 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)
//
// 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.
//
// 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`;
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`,
`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);
}
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) {
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, observer_lag_s, queue_depth)
VALUES (1, $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11::jsonb, $12, now(), $13::jsonb, $14, $15)
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`,
[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]);
} 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(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);
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);
// 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); }
for (const h of inner.removedChainBlockHashes || []) { pendingChain.add.delete(h); pendingChain.remove.add(h); pendingMerge.delete(h); pendingUnmerge.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 () => { 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(() => 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); });