diff --git a/site/api/live.mjs b/site/api/live.mjs index 60ab8c432..b6e5ceef4 100644 --- a/site/api/live.mjs +++ b/site/api/live.mjs @@ -79,6 +79,9 @@ export default async function handler(req, res) { blocks_per_minute: s ? s.blocks_per_minute || [] : [], miners_10m: minerRows.length, observer_started_at: s ? s.observer_started_at : null, + // 4 Oct 2026: how far the observer's newest stored block is behind the clock, and what it still has queued + observer_lag_s: s ? num(s.observer_lag_s) : null, + queue_depth: s ? num(s.queue_depth) : null, updated_at: s ? s.updated_at : null, }; diff --git a/site/live.html b/site/live.html index a9087a176..584029433 100644 --- a/site/live.html +++ b/site/live.html @@ -271,7 +271,7 @@ footer{border-top:1px solid var(--line);padding-block:48px 32px} var blocks=[],byHash={},miners={},lanes=[],laneN=1,othersN=0; var LV=[{S:22,s:14,pps:40,npps:14},{S:16,s:12,pps:26,npps:10},{S:12,s:10,pps:18,npps:8}]; var pxPerSec=40,targetPx=40,bps=0,live=false,serverNow=0,serverAt=0,frozenAt=null,note=''; - var running=false,needSize=true,inView=true,ptrX=-1,ptrY=-1,ptrTouch=false,pin=null,laneLast=[],laneRun=[]; + var running=false,needSize=true,inView=true,ptrX=-1,ptrY=-1,ptrTouch=false,pin=null,laneLast=[],laneRun=[],lagNote=''; var TOP=20,BOT=12; // 12 hues that sit beside ember on obsidian; the index comes from the short id, so a miner keeps its colour across reloads var PAL=['#5CB8FF','#3FD39A','#F5C542','#C49BFF','#FF7AA8','#8FD15F','#FFB08A','#6F8CFF','#D9E25A','#FF8FD6','#4FD8D8','#E0A96D']; @@ -402,7 +402,8 @@ footer{border-top:1px solid var(--line);padding-block:48px 32px} for(i=0;i30?'observer '+(lag<120?Math.round(lag)+' s':Math.round(lag/60)+' min')+' behind':''; + setStatus(on,s.age_s===null?'no observer update yet':(lagNote?lagNote+(s.queue_depth?', '+s.queue_depth+' queued':''):'observer updated '+rel(s.age_s*1000))); $('st-network').textContent=s.network||'igneum-devnet';$('st-node').textContent=s.node_version?'node '+s.node_version:'node'; $('st-bps').textContent=(s.blocks_per_second_60s||0).toFixed(2); $('st-blue').textContent=compact(s.blue_score);$('st-count').textContent=compact(s.block_count)+' blocks'; diff --git a/tools/observer/observer.mjs b/tools/observer/observer.mjs index f717473a9..29d1ec353 100644 --- a/tools/observer/observer.mjs +++ b/tools/observer/observer.mjs @@ -16,6 +16,14 @@ // 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'; @@ -81,6 +89,8 @@ async function setupSchema() { 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(), @@ -239,6 +249,17 @@ 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 @@ -273,15 +294,18 @@ async function recordEvent(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; +// 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); - for (const k of minuteCounts.keys()) if (k < now - 61 * 60_000) minuteCounts.delete(k); + if (minuteCounts.size > 70) 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; + 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; @@ -297,14 +321,20 @@ function onBlock(block) { for (const cert of certificatesIn(block)) pendingCerts.push({ cert, carrier: h.hash }); const vk = h.voteKeyHash || null; const vd = block.verboseData; - if (vd && vd.isChainBlock && Array.isArray(vd.mergeSetBluesHashes)) pendingMerge.set(h.hash, { blues: vd.mergeSetBluesHashes, reds: vd.mergeSetRedsHashes || [] }); + 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); + noteArrival(now, Number(h.timestamp) || now); noteWork(now, h.blueWork); if (vk) { const m = miners.get(vk); @@ -327,33 +357,42 @@ async function flushBlocks() { } 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)); } + 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 back = []; - for (const h of [...pendingUnmerge].slice(0, 50)) { try { const m = await mergesetOf(rpc, h); back.push(...m.blues, ...m.reds); pendingUnmerge.delete(h); } catch (e) { log('unmerge fetch failed', h.slice(0, 8), e.message); pendingUnmerge.delete(h); } } + 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 = []; - for (const [h, m] of [...pendingMerge].slice(0, 100)) { - let ms = m; - if (!ms) { try { ms = await mergesetOf(rpc, h); } catch (e) { log('mergeset fetch failed', h.slice(0, 8), e.message); pendingMerge.delete(h); continue; } } - blues.push(...ms.blues); reds.push(...ms.reds); pendingMerge.delete(h); - } + 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); } @@ -379,6 +418,19 @@ async function tick(rpc) { 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 || []; @@ -406,17 +458,17 @@ async function tick(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) + 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`, + 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]); + finality ? JSON.stringify(finality) : null, lagS === null ? null : Math.round(lagS * 10) / 10, queueDepth]); } catch (e) { log('state write failed', e.message); } } @@ -626,8 +678,12 @@ 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`); + 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); @@ -642,7 +698,14 @@ async function main() { 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) => { + // 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) { @@ -655,10 +718,10 @@ async function main() { } } 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, null); } + 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)); } @@ -679,7 +742,9 @@ async function main() { rpc.onClose = () => { setTimeout(connect, 2000); }; await connect(); - setInterval(async () => { await flushBlocks(); await flushChain(); await flushColors(rpc); await flushCertificates(rpc); }, FLUSH_EVERY_MS); + 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(); diff --git a/tools/observer/run.sh b/tools/observer/run.sh new file mode 100755 index 000000000..d5ad6b66f --- /dev/null +++ b/tools/observer/run.sh @@ -0,0 +1,20 @@ +#!/bin/sh +# Restart loop for the devnet observer. An exit (the observer exits 2 when it loses the block feed and two +# resubscribes do not bring it back) is a restart three seconds later. Start it detached from the repo root: +# +# nohup tools/observer/run.sh >/dev/null 2>&1 & +# +# Environment: IGNEUM_RPC (default ws://127.0.0.1:28640, the Mac's non-mining peer with a JSON listener), +# OBSERVER_LOG (default /tmp/igneum-devnet/observer-mjs-v4.out, appended). Everything else as observer.mjs. +cd "$(dirname "$0")/../.." || exit 1 +export IGNEUM_RPC="${IGNEUM_RPC:-ws://127.0.0.1:28640}" +LOG="${OBSERVER_LOG:-/tmp/igneum-devnet/observer-mjs-v4.out}" +mkdir -p "$(dirname "$LOG")" +exec >>"$LOG" 2>&1 +echo "$(date -u +%FT%TZ) run.sh: observer loop started, rpc $IGNEUM_RPC, pid $$" +while :; do + node tools/observer/observer.mjs + code=$? + echo "$(date -u +%FT%TZ) run.sh: observer exited with $code, restarting in 3 s" + sleep 3 +done