From fabd725bc953b7fe366fa2cf58be08125c4eb85b Mon Sep 17 00:00:00 2001 From: igneum-labs <337424239+igneum-labs@users.noreply.github.com> Date: Sat, 3 Oct 2026 19:22:06 +0000 Subject: [PATCH] Live devnet: observer, /api/live and the /live page tools/observer: Node 22 observer on the node's wRPC JSON port (blockAdded and virtualChainChanged subscriptions, 2 s state ticks) writing live_blocks, live_state and live_events to Neon, miner address decoded from the coinbase payload, events for new or quiet miners, peers and difficulty steps. site/api/live.mjs: three indexed queries, max-age=1. site/live.html: status strip, DAG stream with real parent edges and a lane per miner, miners table, events feed, blocks-per-minute sparkline, OFFLINE freeze. The homepage live path and the /api/live cache header went in with 8a222c2. Co-Authored-By: Claude Fable 5.1 --- site/api/live.mjs | 108 ++++++++++ site/live.html | 332 +++++++++++++++++++++++++++++++ tools/observer/README.md | 47 +++++ tools/observer/observer.mjs | 382 ++++++++++++++++++++++++++++++++++++ 4 files changed, 869 insertions(+) create mode 100644 site/api/live.mjs create mode 100644 site/live.html create mode 100644 tools/observer/README.md create mode 100644 tools/observer/observer.mjs diff --git a/site/api/live.mjs b/site/api/live.mjs new file mode 100644 index 000000000..68be617ee --- /dev/null +++ b/site/api/live.mjs @@ -0,0 +1,108 @@ +// Igneum live devnet feed. GET /api/live returns what tools/observer wrote to Neon: +// {now, state, blocks (last 90 s), miners (last 10 min), events (last 30)}. +// Three queries, all on indexed timestamps. Cached one second at the edge. +// Zero dependencies: Neon's HTTP SQL endpoint over Node's built-in fetch. + +const STALE_AFTER_S = 30; + +function neon() { + const url = process.env.DATABASE_URL; + if (!url) throw new Error('DATABASE_URL is not set'); + const host = new URL(url).hostname.replace('-pooler', ''); + return async (query, params = []) => { + const r = await fetch(`https://${host}/sql`, { + method: 'POST', + headers: { 'Neon-Connection-String': 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 || []; + }; +} + +const num = v => (v === null || v === undefined ? null : Number(v)); +const short = h => (h ? String(h).slice(0, 16) : null); +const minerId = vk => (vk ? String(vk).slice(0, 8) : null); + +export default async function handler(req, res) { + res.setHeader('Cache-Control', 'public, max-age=1'); + res.setHeader('Access-Control-Allow-Origin', '*'); + if (req.method !== 'GET') { + res.setHeader('Allow', 'GET'); + return res.status(405).json({ ok: false, error: 'method not allowed' }); + } + try { + const sql = neon(); + const [head, blockRows, minerRows] = await Promise.all([ + sql(`SELECT now() AS now, + (SELECT row_to_json(s) FROM live_state s WHERE s.id = 1) AS state, + (SELECT json_agg(e) FROM (SELECT ts, kind, text FROM live_events ORDER BY ts DESC LIMIT 30) e) AS events`), + sql(`SELECT hash, blue_score, daa_score, timestamp_ms, parent_hashes, is_chain_block, vote_key_hash, received_at + FROM live_blocks + WHERE received_at > now() - interval '90 seconds' + ORDER BY received_at DESC + LIMIT 600`), + sql(`SELECT vote_key_hash, count(*)::int AS blocks, max(received_at) AS last_seen, max(engine) AS engine, max(miner_address) AS address + FROM live_blocks + WHERE received_at > now() - interval '10 minutes' AND vote_key_hash IS NOT NULL + GROUP BY vote_key_hash + ORDER BY blocks DESC, last_seen DESC`), + ]); + + const now = new Date(head[0].now).getTime(); + const s = head[0].state || null; + const updated = s && s.updated_at ? new Date(s.updated_at).getTime() : null; + const age_s = updated === null ? null : Math.max(0, (now - updated) / 1000); + const stale = age_s === null || age_s > STALE_AFTER_S; + const total10m = minerRows.reduce((a, m) => a + m.blocks, 0); + + const state = { + stale, + age_s: age_s === null ? null : Math.round(age_s * 10) / 10, + network: s ? s.network : null, + node_version: s ? s.node_version : null, + block_count: s ? num(s.block_count) : null, + header_count: s ? num(s.header_count) : null, + blue_score: s ? num(s.blue_score) : null, + difficulty: s ? num(s.difficulty) : null, + hashes_per_second_estimate: s ? num(s.hashes_per_second_estimate) : null, + peers: s ? num(s.peers) : null, + mempool: s ? num(s.mempool) : null, + blocks_60s: s ? num(s.blocks_60s) : null, + blocks_per_second_60s: s && s.blocks_60s !== null ? Math.round(num(s.blocks_60s) / 60 * 1000) / 1000 : null, + blocks_per_minute: s ? s.blocks_per_minute || [] : [], + miners_10m: minerRows.length, + observer_started_at: s ? s.observer_started_at : null, + updated_at: s ? s.updated_at : null, + }; + + const blocks = blockRows.reverse().map(b => ({ + hash: short(b.hash), + blue_score: num(b.blue_score), + daa: num(b.daa_score), + ts: num(b.timestamp_ms), + rx: new Date(b.received_at).getTime(), + parents: (b.parent_hashes || []).map(short), + chain: !!b.is_chain_block, + miner: minerId(b.vote_key_hash), + })); + + const miners = minerRows.map(m => ({ + id: minerId(m.vote_key_hash), + vote_key_hash: m.vote_key_hash, + blocks: m.blocks, + share: total10m ? Math.round(m.blocks / total10m * 1000) / 10 : 0, + last_seen: m.last_seen, + engine: m.engine || null, + address: m.address || null, + })); + + const events = (head[0].events || []).map(e => ({ ts: e.ts, kind: e.kind, text: e.text })); + + return res.status(200).json({ ok: true, now: new Date(now).toISOString(), state, blocks, miners, events }); + } catch (e) { + res.setHeader('Cache-Control', 'no-store'); + return res.status(500).json({ ok: false, error: String(e.message || e) }); + } +} diff --git a/site/live.html b/site/live.html new file mode 100644 index 000000000..be1efe32e --- /dev/null +++ b/site/live.html @@ -0,0 +1,332 @@ + + + + + +Live devnet + + + + + + + + + + + + + + +
+
+
devnet v0
+

Live devnet

+

A node is read every two seconds. Every block below is real: its parents, its miner and whether it sits on the selected chain. Proving and checkpoints are not in the devnet yet.

+
+ +
+
Status
OFFLINE
no observer update yet
+
Network
igneum-devnet
node
+
Blocks / s
0.00
over 60 s
+
Blue score
0
0 blocks
+
Difficulty
0
target per block
+
Hash rate
0
network estimate
+
Miners
0
active in 10 min
+
Peers
0
mempool 0
+
+ +
+
+
block dag
+
on screen 0chain 0side 0newest blue 0
+
+ +
+ chain block + side block (blue or red, off the selected chain) + just arrived + parent edge + 2ea5116eminer short id, one lane each +
+
+ +
+
+

Miners

last 10 min
+
+ + + +
MinerBlocksShareLast seenEngine
No blocks in the last 10 minutes.
+
+

The short id is the first 8 hex characters of the vote key hash in each header. The engine is what the miner writes into the coinbase; the devnet miner writes nothing yet.

+
+
+

Events

newest first
+
No events yet.
+
+
+ +
+

Blocks per minute

last hour
+ +

Counted by the observer as blocks arrive.

+
+
+ + + + + + diff --git a/tools/observer/README.md b/tools/observer/README.md new file mode 100644 index 000000000..904c61c20 --- /dev/null +++ b/tools/observer/README.md @@ -0,0 +1,47 @@ +# Igneum devnet observer + +Watches one Igneum node over wRPC JSON and writes what it sees to Neon, so `/api/live` and `/live` on the site can show the devnet as it runs. Node 22 or newer, no dependencies. + +## Run + +``` +node tools/observer/observer.mjs +``` + +Environment, every value optional: + +| Variable | Default | Meaning | +|---|---|---| +| `IGNEUM_RPC` | `ws://127.0.0.1:28610` | The node's wRPC JSON url. 28610 is the igneum-devnet wRPC JSON port (`consensus/core/src/network.rs`). Start the node with `--rpclisten-json=127.0.0.1:28610` or the port of your choice. | +| `DATABASE_URL` | read from `~/.config/igneum/env` | Neon connection string. Never commit it. | +| `LIVE_RETAIN_HOURS` | `24` | Hours of blocks kept in `live_blocks`. Older rows are deleted once a minute. | + +A node started by another tool may listen on gRPC only. Then run your own non-mining peer with a JSON listener, on ports that do not clash with the devnet's (gRPC 26610, P2P 26611): + +``` +vendor/igneum-node/target/release/kaspad --devnet --nodnsseed --disable-upnp \ + --appdir=/tmp/igneum-obsnode --rpclisten=127.0.0.1:26640 --rpclisten-json=127.0.0.1:28640 \ + --listen=127.0.0.1:26641 --connect=127.0.0.1:26611 --nologfiles +IGNEUM_RPC=ws://127.0.0.1:28640 node tools/observer/observer.mjs +``` + +## What it does + +- Subscribes to `blockAdded` and `virtualChainChanged`; polls `getBlockDagInfo`, `getInfo`, `getConnectedPeerInfo`, `getSinkBlueScore` and `estimateNetworkHashesPerSecond` every 2 s. +- Decodes the miner address from the coinbase payload script (same bech32 variant as `crypto/addresses`). The fork's `vote_key_hash` header field is stored per block; the first 8 hex characters are the miner's short id on the site. +- `engine` is the miner's tag in the coinbase extra data after the node's version prefix. The node exposes no engine name over RPC, so this is null on devnet v0. +- When the node refuses the hash-rate estimate (it needs a 1,000-block window) the observer reports blue work added per second over the last 10 minutes instead. + +## Tables + +Created on start if missing. + +| Table | Rows | Columns | +|---|---|---| +| `live_blocks` | one per block, kept `LIVE_RETAIN_HOURS` | `hash`, `blue_score`, `daa_score`, `timestamp_ms`, `parents` (count), `parent_hashes`, `is_chain_block`, `vote_key_hash`, `miner_address`, `engine`, `received_at`. Indexes on `received_at`, `timestamp_ms`, `(vote_key_hash, received_at)`. | +| `live_state` | one row, updated every 2 s | `block_count`, `header_count`, `blue_score`, `difficulty`, `hashes_per_second_estimate`, `peers`, `mempool`, `node_version`, `network`, `blocks_60s`, `blocks_per_minute` (60 pairs of minute epoch ms and count), `observer_started_at`, `updated_at` | +| `live_events` | one per event, kept 7 days | `ts`, `kind`, `text`. Kinds: `observer`, `miner_seen`, `miner_quiet`, `miner_back`, `peer_joined`, `peer_left`, `difficulty` (step over 5%). Checkpoints will be added when the finality layer lands. | + +## Reading it + +`site/api/live.mjs` serves `/api/live` from these tables in three indexed queries. `site/live.html` polls it every 2 s. The site shows OFFLINE when `live_state.updated_at` is older than 30 s. diff --git a/tools/observer/observer.mjs b/tools/observer/observer.mjs new file mode 100644 index 000000000..c8f7f2ccf --- /dev/null +++ b/tools/observer/observer.mjs @@ -0,0 +1,382 @@ +#!/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 +// +// Tables (created on start if missing): live_blocks, live_state, live_events. See README.md. +// 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 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 live_blocks ( + 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 live_blocks_received_at ON live_blocks (received_at)`, + `CREATE INDEX IF NOT EXISTS live_blocks_timestamp_ms ON live_blocks (timestamp_ms)`, + `CREATE INDEX IF NOT EXISTS live_blocks_vote_key_received ON live_blocks (vote_key_hash, received_at)`, + `CREATE TABLE IF NOT EXISTS live_state ( + 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)`, + `CREATE TABLE IF NOT EXISTS live_events ( + id bigserial PRIMARY KEY, + ts timestamptz NOT NULL DEFAULT now(), + kind text NOT NULL, + text text NOT NULL)`, + `CREATE INDEX IF NOT EXISTS live_events_ts ON live_events (ts)`, + ]; + 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. + const text = Buffer.from(extra).toString('utf8').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 live_events (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 live_blocks (${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 live_blocks 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 live_blocks 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)`); } + + try { + await sql(`INSERT INTO live_state (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) + VALUES (1, $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11::jsonb, $12, now()) + 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()`, + [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()]); + } catch (e) { log('state write failed', e.message); } +} + +async function prune() { + try { + await sql(`DELETE FROM live_blocks WHERE received_at < now() - ($1 || ' hours')::interval`, [String(RETAIN_HOURS)]); + await sql(`DELETE FROM live_events WHERE ts < 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 live_blocks 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 live_blocks 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 live_state WHERE id = 1`); + 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 === '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 } }); + 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); });