igneum/tools/observer/observer.mjs
igneum-labs e24b643e01 Finality v2: fork reading guide, spec implementation notes, bench entry, observer checkpoints, live page locks
- docs/fork-divergence.md: "Finality v2" table (every file, risk, merge note), decisions
- docs/spec/03-finality.md: section 3.10 implementation notes, clause by clause
- docs/bench-log.md: test-network results (72 of 72 steady locks, median 0.80 s; equivocation
  strip; partition: 0 locks at 39.6% of total with the floor binding, heal in 30 s), follower
- tools/observer: live_checkpoints table, FinalityLock subscription, "checkpoint N locked" events
- site: /api/live adds checkpoints and locked/final flags; /live draws the lock ring and final line

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-03 21:54:24 +00:00

501 lines
26 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. 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.
// 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`;
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)`,
];
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) => {
let m; try { m = JSON.parse(e.data); } 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);
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()`,
[cp.index, cp.hash, cp.blueScore, cp.daaScore, cp.state, cp.signedWeight, cp.activeWeight, cp.totalWeight, cp.fractionActive, cp.fractionTotal,
cp.votesSeen, cp.voters, 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;
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,
};
}
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);
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(); }, 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); });