tools/observer/detector.mjs: per-program implied rate per miner id from blue work over wall seconds (the chain's own estimate rule restricted to one id), excess spread above Poisson, epoch-start share, nonce chi-square and increasing-fraction tests, card bands from the log intake, two-way residual correlations and cliques; a design_candidate clique held 6 net windows is the alert, written to live_state.detector and live_events kind detector. One hook in observer.mjs (every 60 s) and one jsonb column. node:test file with a fabricated fixed design (fires) and a fabricated honest population (quiet); the live devnet in --dry mode is quiet with its baseline recorded in the README. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
346 lines
26 KiB
JavaScript
346 lines
26 KiB
JavaScript
// Share-pattern detector (Counter ASIC 3.0 item 4a, 6 October 2026). Reads what the observer already stores per block
|
|
// (live_blocks: vote_key_hash as the miner id, daa_score, timestamp_ms, color, detail.bits, detail.nonce) and what the log
|
|
// intake knows about card models (miner_logs: the app's "GPUs:" line and the workers' STATUS lines), and answers one
|
|
// question once a minute: does any group of miner ids behave like one fixed design? Written to live_state.detector
|
|
// (jsonb) and, on a change, to live_events as kind `detector`. The chain does nothing with it: the detector is the
|
|
// trigger for people (docs/plans/epoch-length.md section 11), not for consensus.
|
|
//
|
|
// Everything that decides is a pure function over rows (aggregate, analyse, cardBands), tested in detector.test.mjs
|
|
// with a fabricated fixed-design population and a fabricated honest one. The live path (run) only fetches rows and
|
|
// writes the result. `node tools/observer/detector.mjs --dry` runs the live tables read-only and prints the state.
|
|
//
|
|
// Baseline numbers from the devnet (6 October 2026) and the thresholds' reasons: tools/observer/README.md, section
|
|
// "Detector". No em dashes anywhere in this file by the copy law.
|
|
|
|
import { readFileSync } from 'node:fs';
|
|
import { homedir } from 'node:os';
|
|
|
|
export const DEFAULTS = {
|
|
epochLen: 3600, // DAA s per program (spec 01 section 1.12; epoch-length.md changes this by signal only)
|
|
windowEpochs: 6, // the rolling window, closed epochs only (DETECTOR_WINDOW_EPOCHS)
|
|
settleDaa: 1200, // an epoch is re-read until the tip is this far past its end (colours settle; merge depth scale)
|
|
minBlue: 30, // an id is "present" in an epoch with at least this many blue blocks (Poisson sd 18%)
|
|
minEpochsSpread: 4, // epochs an id must be present in for the spread statistic
|
|
minEpochsCorr: 5, // epochs two ids must share for a correlation (n = 3 or 4 is noise: see README)
|
|
stepMax: 0.5, // |log rate change| between consecutive present epochs above this = an operational step, not a program
|
|
excessSpreadMax: 0.10, // excess (above Poisson) per-program spread above 10% is a design flag (honest devnet max 5.1%)
|
|
earlyShare: 0.10, // the first tenth of each epoch (DAA) carries a tenth of an honest miner's blocks
|
|
earlyShareMin: 0.02, // under 2% with earlyMinBlue blocks = the miner cannot mine the start of an epoch (compile per program)
|
|
earlyMinBlue: 200, // P(X <= 4 | n = 200, p = 0.1) is about 1e-6
|
|
nonceMin: 32, // nonces needed for the nonce test
|
|
chi2Crit: 37.70, // chi-square, 15 degrees of freedom, p = 0.001
|
|
incMin: 0.35, incMax: 0.65, incMinN: 200, // fraction of increasing consecutive nonces (honest 0.50; a counter 1.0)
|
|
bandTol: 0.30, // a chain-implied rate within 30% of a band / divisor matches it
|
|
identityDivisors: [1, 2, 8], // the app mines a card under 1, 2 or 8 vote keys (app/igneum-app/src/detect.rs)
|
|
winsor: 0.30, // residual log rates are clipped at +-30% before correlating
|
|
r: 0.8, // pairwise residual correlation above this = an edge (at n = 6 a true 0.95 reads 0.85 to 0.99 with noise)
|
|
k: 3, // a clique of this many ids with design flags = a candidate alert
|
|
holdWindows: 6, // windows (one per closed epoch) the candidate must hold before the alert is active; a window without it counts one down
|
|
maxIds: 400, // the correlation graph is bounded (ids with the most blue blocks)
|
|
};
|
|
|
|
// ---------- the chain's work rule (vendor/igneum-node/consensus/src/processes/difficulty.rs calc_work) ----------
|
|
export function targetFromBits(bits) {
|
|
bits = Number(bits);
|
|
const size = bits >>> 24;
|
|
const word = BigInt(bits & 0x007fffff);
|
|
return size <= 3 ? word >> BigInt(8 * (3 - size)) : word << BigInt(8 * (size - 3));
|
|
}
|
|
const U256_MAX = (1n << 256n) - 1n;
|
|
export function calcWork(bits) { const t = targetFromBits(bits); return ((U256_MAX - t) / (t + 1n)) + 1n; }
|
|
// The same quantity as a double for sums (2^256 / (target + 1), exact to 53 bits; the SQL aggregate uses this form)
|
|
export function workDouble(bits) { bits = Number(bits); const size = bits >>> 24; const word = bits & 0x007fffff; return 2 ** (256 - 8 * (size - 3)) / word; }
|
|
export const epochOf = (daa, L = DEFAULTS.epochLen) => Math.floor(Number(daa) / L);
|
|
|
|
// ---------- aggregate: block rows -> per (id, epoch) cells ----------
|
|
// rows: {vote_key_hash, daa_score, timestamp_ms, color, bits, nonce (string, u64)}; any order. Output cells carry what the
|
|
// tests need: blue count, summed work (double), blocks in the first tenth of the epoch, and the nonce histograms.
|
|
export function aggregate(rows, opts = {}) {
|
|
const o = { ...DEFAULTS, ...opts };
|
|
const sorted = [...rows].filter(r => r.vote_key_hash && r.bits != null).sort((a, b) => Number(a.daa_score) - Number(b.daa_score) || Number(a.timestamp_ms) - Number(b.timestamp_ms));
|
|
const epochs = new Map(); const cells = new Map(); const lastNonce = new Map();
|
|
for (const r of sorted) {
|
|
const daa = Number(r.daa_score), e = epochOf(daa, o.epochLen), t = Number(r.timestamp_ms);
|
|
const ep = epochs.get(e) || { e, t0: Infinity, t1: -Infinity, blocks: 0, blue: 0, work: 0, maxDaa: 0 };
|
|
ep.t0 = Math.min(ep.t0, t); ep.t1 = Math.max(ep.t1, t); ep.blocks++; ep.maxDaa = Math.max(ep.maxDaa, daa); epochs.set(e, ep);
|
|
const id = String(r.vote_key_hash).slice(0, 8), key = `${id}:${e}`;
|
|
const c = cells.get(key) || { id, e, blue: 0, work: 0, early: 0, n: 0, lo: new Array(16).fill(0), hi: new Array(16).fill(0), inc: 0, pairs: 0 };
|
|
if (r.nonce !== undefined && r.nonce !== null) {
|
|
let n = null; try { n = BigInt(String(r.nonce)); } catch { n = null; }
|
|
if (n !== null) {
|
|
c.n++; c.lo[Number(n & 15n)]++; c.hi[Number((n >> 60n) & 15n)]++;
|
|
const prev = lastNonce.get(id); if (prev !== undefined) { c.pairs++; if (n > prev) c.inc++; } lastNonce.set(id, n);
|
|
}
|
|
}
|
|
if (r.color === 'blue') { const w = workDouble(r.bits); c.blue++; c.work += w; ep.blue++; ep.work += w; if (daa % o.epochLen < o.epochLen / 10) c.early++; }
|
|
cells.set(key, c);
|
|
}
|
|
return { epochs, cells };
|
|
}
|
|
// Merge cells of the same shape (the live path keeps settled epochs cached and re-reads the open ones)
|
|
export function mergeAggregates(list) {
|
|
const epochs = new Map(), cells = new Map();
|
|
for (const a of list) { for (const [e, ep] of a.epochs) epochs.set(e, ep); for (const [k, c] of a.cells) cells.set(k, c); }
|
|
return { epochs, cells };
|
|
}
|
|
|
|
// ---------- small statistics ----------
|
|
const mean = a => a.reduce((x, y) => x + y, 0) / a.length;
|
|
const sd = a => { if (a.length < 2) return 0; const m = mean(a); return Math.sqrt(a.reduce((x, y) => x + (y - m) ** 2, 0) / (a.length - 1)); };
|
|
const median = a => { const b = [...a].sort((x, y) => x - y); const h = b.length >> 1; return b.length % 2 ? b[h] : (b[h - 1] + b[h]) / 2; };
|
|
export function chi2Uniform(counts) { const n = counts.reduce((a, b) => a + b, 0); if (!n) return 0; const exp = n / counts.length; return counts.reduce((a, c) => a + (c - exp) ** 2 / exp, 0); }
|
|
export function pearson(a, b) { const ma = mean(a), mb = mean(b); let sab = 0, saa = 0, sbb = 0; for (let j = 0; j < a.length; j++) { sab += (a[j] - ma) * (b[j] - mb); saa += (a[j] - ma) ** 2; sbb += (b[j] - mb) ** 2; } return saa > 0 && sbb > 0 ? sab / Math.sqrt(saa * sbb) : 0; }
|
|
const round = (x, d = 1) => x === null || x === undefined || !Number.isFinite(x) ? null : Math.round(x * 10 ** d) / 10 ** d;
|
|
|
|
// Maximal cliques of size >= k in a small graph (Bron-Kerbosch without pivoting; the graph is bounded by maxIds and
|
|
// by the edge rule, and ids with degree under k - 1 are dropped first)
|
|
export function cliques(nodes, edges, k) {
|
|
const adj = new Map(nodes.map(n => [n, new Set()]));
|
|
for (const [a, b] of edges) { adj.get(a).add(b); adj.get(b).add(a); }
|
|
const keep = nodes.filter(n => adj.get(n).size >= k - 1);
|
|
const out = [];
|
|
const bk = (R, P, X) => {
|
|
if (!P.size && !X.size) { if (R.length >= k) out.push([...R]); return; }
|
|
for (const v of [...P]) {
|
|
const nv = adj.get(v);
|
|
bk([...R, v], new Set([...P].filter(x => nv.has(x))), new Set([...X].filter(x => nv.has(x))));
|
|
P.delete(v); X.add(v);
|
|
if (out.length > 50) return;
|
|
}
|
|
};
|
|
bk([], new Set(keep), new Set());
|
|
return out.sort((a, b) => b.length - a.length);
|
|
}
|
|
|
|
// ---------- analyse: cells -> the detector state ----------
|
|
// bands: [{model, p5, p50, p95}] in MH/s (cardBands); prev: the previous state (for the hold counter); tipDaa: the tip
|
|
export function analyse(agg, opts = {}, bands = [], prev = null, tipDaa = null) {
|
|
const o = { ...DEFAULTS, ...opts };
|
|
const L = o.epochLen;
|
|
const all = [...agg.epochs.values()].sort((a, b) => a.e - b.e);
|
|
const tip = tipDaa ?? (all.length ? all[all.length - 1].maxDaa : 0);
|
|
// closed = the whole epoch is in the past of the tip; settled = colours are final
|
|
const closed = all.filter(ep => (ep.e + 1) * L <= tip && ep.blocks >= 0.5 * L);
|
|
const win = closed.slice(-o.windowEpochs);
|
|
const W = win.map(ep => ep.e);
|
|
const secs = new Map(win.map(ep => [ep.e, Math.max(1, (ep.t1 - ep.t0) / 1000)]));
|
|
const network = win.map(ep => ({ epoch: ep.e, secs: round(secs.get(ep.e), 0), blue: ep.blue, mhs: round(ep.work / secs.get(ep.e) / 1e6, 1), settled: tip >= (ep.e + 1) * L + o.settleDaa }));
|
|
// per id
|
|
const byId = new Map();
|
|
for (const c of agg.cells.values()) { if (!W.includes(c.e)) continue; if (!byId.has(c.id)) byId.set(c.id, []); byId.get(c.id).push(c); }
|
|
const ids = [...byId.keys()].sort((a, b) => byId.get(b).reduce((x, c) => x + c.blue, 0) - byId.get(a).reduce((x, c) => x + c.blue, 0)).slice(0, o.maxIds);
|
|
const logRate = new Map(); // id -> Map(epoch -> log rate)
|
|
const miners = new Map();
|
|
for (const id of ids) {
|
|
const cs = byId.get(id);
|
|
const present = cs.filter(c => c.blue >= o.minBlue).sort((a, b) => a.e - b.e);
|
|
const lr = new Map(present.map(c => [c.e, Math.log(c.work / secs.get(c.e))]));
|
|
logRate.set(id, lr);
|
|
const rates = [...lr.values()];
|
|
const blueTotal = cs.reduce((x, c) => x + c.blue, 0), early = cs.reduce((x, c) => x + c.early, 0);
|
|
// steady: no operational step between consecutive present epochs
|
|
let steady = present.length >= 2; const steps = [];
|
|
for (let i = 1; i < present.length; i++) { const d = lr.get(present[i].e) - lr.get(present[i - 1].e); steps.push(d); if (Math.abs(d) > o.stepMax) steady = false; }
|
|
const spreadSd = rates.length >= 2 ? sd(rates) : null;
|
|
const poissonSd = present.length ? Math.sqrt(mean(present.map(c => 1 / c.blue))) : null;
|
|
const excess = spreadSd === null ? null : Math.sqrt(Math.max(0, spreadSd ** 2 - poissonSd ** 2));
|
|
// nonces, pooled over the window
|
|
const lo = new Array(16).fill(0), hi = new Array(16).fill(0); let n = 0, inc = 0, pairs = 0;
|
|
for (const c of cs) { n += c.n; inc += c.inc; pairs += c.pairs; for (let i = 0; i < 16; i++) { lo[i] += c.lo[i]; hi[i] += c.hi[i]; } }
|
|
const chiLo = chi2Uniform(lo), chiHi = chi2Uniform(hi), incFrac = pairs ? inc / pairs : null;
|
|
const flags = [], notes = [];
|
|
if (present.length >= o.minEpochsSpread && steady && excess > o.excessSpreadMax) flags.push('spread');
|
|
if (!steady && present.length >= 2) notes.push('unsteady');
|
|
const earlyFrac = blueTotal ? early / blueTotal : null;
|
|
if (blueTotal >= o.earlyMinBlue && earlyFrac < o.earlyShareMin) flags.push('late_start');
|
|
if (n >= o.nonceMin && (chiLo > o.chi2Crit || chiHi > o.chi2Crit || (pairs >= o.incMinN && (incFrac < o.incMin || incFrac > o.incMax)))) flags.push('nonce');
|
|
// band: the window-median chain rate against every band / identity divisor
|
|
const mhs = rates.length ? Math.exp(median(rates)) / 1e6 : null;
|
|
let band = null;
|
|
if (mhs !== null && bands.length) {
|
|
const matches = [];
|
|
for (const b of bands) for (const d of o.identityDivisors) if (mhs >= b.p5 / d * (1 - o.bandTol) && mhs <= b.p95 / d * (1 + o.bandTol)) matches.push(`${b.model}/${d}`);
|
|
const top = Math.max(...bands.map(b => b.p95)), bottom = Math.min(...bands.map(b => b.p5)) / Math.max(...o.identityDivisors);
|
|
band = { matches, high: mhs > top * (1 + o.bandTol), small: mhs < bottom * (1 - o.bandTol) };
|
|
if (band.high) flags.push('band_high');
|
|
else if (!matches.length && !band.small && steady && present.length >= o.minEpochsSpread) flags.push('band');
|
|
if (band.small) notes.push('below_every_band');
|
|
}
|
|
miners.set(id, {
|
|
id, epochs_present: present.length, blue: blueTotal, mhs: round(mhs, 2),
|
|
steady, max_step_pct: steps.length ? round(Math.max(...steps.map(Math.abs)) * 100, 0) : null,
|
|
spread_sd_pct: round(spreadSd === null ? null : spreadSd * 100), poisson_sd_pct: round(poissonSd === null ? null : poissonSd * 100), excess_spread_pct: round(excess === null ? null : excess * 100),
|
|
early_share_pct: round(earlyFrac === null ? null : earlyFrac * 100), nonce: { n, chi2_low4: round(chiLo), chi2_high4: round(chiHi), inc_frac: round(incFrac, 3) },
|
|
band, flags, notes,
|
|
});
|
|
}
|
|
// two-way residuals: id mean and the epoch common factor over the ids present in it, clipped
|
|
const idMean = new Map([...logRate].map(([id, lr]) => [id, lr.size ? mean([...lr.values()]) : 0]));
|
|
const epochFactor = new Map(W.map(e => { const devs = ids.filter(id => logRate.get(id).has(e)).map(id => logRate.get(id).get(e) - idMean.get(id)); return [e, devs.length ? mean(devs) : 0]; }));
|
|
const resid = new Map(ids.map(id => [id, new Map([...logRate.get(id)].map(([e, v]) => [e, Math.max(-o.winsor, Math.min(o.winsor, v - idMean.get(id) - epochFactor.get(e)))]))]));
|
|
// correlation graph
|
|
const corrIds = ids.filter(id => resid.get(id).size >= o.minEpochsCorr);
|
|
const edges = []; let pairsTested = 0, maxR = null;
|
|
for (let a = 0; a < corrIds.length; a++) for (let b = a + 1; b < corrIds.length; b++) {
|
|
const ra = resid.get(corrIds[a]), rb = resid.get(corrIds[b]);
|
|
const common = [...ra.keys()].filter(e => rb.has(e));
|
|
if (common.length < o.minEpochsCorr) continue;
|
|
const r = pearson(common.map(e => ra.get(e)), common.map(e => rb.get(e)));
|
|
pairsTested++; if (maxR === null || r > maxR) maxR = r;
|
|
if (r > o.r) edges.push([corrIds[a], corrIds[b], r]);
|
|
}
|
|
const groups = cliques(corrIds, edges.map(e => [e[0], e[1]]), o.k).map(members => {
|
|
const rs = edges.filter(e => members.includes(e[0]) && members.includes(e[1])).map(e => e[2]);
|
|
const designFlags = [...new Set(members.flatMap(id => miners.get(id).flags))];
|
|
return { ids: members, size: members.length, min_r: round(Math.min(...rs), 2), design_flags: designFlags, kind: designFlags.length ? 'design_candidate' : 'machine_group' };
|
|
});
|
|
// the alert: a clique with design flags, held over consecutive windows (one window per newly closed epoch)
|
|
const windowEnd = W.length ? W[W.length - 1] : null;
|
|
const candidate = groups.find(g => g.kind === 'design_candidate') || null;
|
|
const prevAlert = (prev && prev.alert) || { held: 0, window_end: null, active: false, since: null };
|
|
// one count per newly closed epoch: up with a candidate, one down without (a clique at the noise edge may drop a
|
|
// pair for one window and come back; the alert needs holdWindows net)
|
|
let held = prevAlert.held || 0;
|
|
if (W.length >= o.minEpochsCorr && windowEnd !== prevAlert.window_end) held = candidate ? held + 1 : Math.max(0, held - 1);
|
|
else if (W.length < o.minEpochsCorr) held = 0;
|
|
const active = held >= o.holdWindows;
|
|
const alert = { active, held, hold_windows: o.holdWindows, window_end: windowEnd, since: active ? (prevAlert.active ? prevAlert.since : new Date().toISOString()) : null, candidate };
|
|
// events: transitions and new per-id flags
|
|
const events = [];
|
|
const prevFlags = new Map(((prev && prev.miners) || []).map(m => [m.id, m.flags || []]));
|
|
for (const m of miners.values()) for (const f of m.flags) if (!(prevFlags.get(m.id) || []).includes(f)) events.push(`Detector: miner ${m.id} flagged ${f} (${flagEvidence(m, f)})`);
|
|
if (candidate && !(prevAlert.candidate && sameIds(prevAlert.candidate.ids, candidate.ids))) events.push(`Detector: ${candidate.size} ids move as one machine with design flags ${candidate.design_flags.join(', ')} (min r ${candidate.min_r}); held ${held} of ${o.holdWindows} windows`);
|
|
if (active && !prevAlert.active) events.push(`Detector ALERT: a population behaves like one fixed design: ${candidate.ids.join(', ')} (${candidate.design_flags.join(', ')}, min r ${candidate.min_r}) over ${o.holdWindows} windows of ${o.windowEpochs} epochs; see epoch-length.md section 11`);
|
|
if (!active && prevAlert.active) events.push('Detector: the alert cleared');
|
|
const state = {
|
|
computed_at: new Date().toISOString(), epoch_len: L, tip_daa: tip,
|
|
window: { epochs: W, closed_epochs: closed.length, ids: ids.length, ids_correlated: corrIds.length },
|
|
network, miners: [...miners.values()],
|
|
correlation: { pairs_tested: pairsTested, max_r: round(maxR, 2), edges: edges.length, groups },
|
|
alert, thresholds: { r: o.r, k: o.k, hold_windows: o.holdWindows, excess_spread_max_pct: o.excessSpreadMax * 100, early_share_min_pct: o.earlyShareMin * 100, chi2_crit: o.chi2Crit, inc_range: [o.incMin, o.incMax], band_tol_pct: o.bandTol * 100, min_blue: o.minBlue, min_epochs_corr: o.minEpochsCorr },
|
|
bands: bands.map(b => ({ model: b.model, p5: round(b.p5), p50: round(b.p50), p95: round(b.p95), n: b.n })),
|
|
};
|
|
return { state, events };
|
|
}
|
|
const sameIds = (a, b) => a.length === b.length && a.every(x => b.includes(x));
|
|
function flagEvidence(m, f) {
|
|
if (f === 'spread') return `excess per-program spread ${m.excess_spread_pct}% over ${m.epochs_present} epochs, Poisson ${m.poisson_sd_pct}%`;
|
|
if (f === 'late_start') return `${m.early_share_pct}% of ${m.blue} blocks in the first tenth of each epoch`;
|
|
if (f === 'nonce') return `chi2 low ${m.nonce.chi2_low4}, high ${m.nonce.chi2_high4}, increasing ${m.nonce.inc_frac} over ${m.nonce.n} nonces`;
|
|
if (f === 'band_high') return `${m.mhs} MH/s, above every known card`;
|
|
if (f === 'band') return `${m.mhs} MH/s matches no known card band at a steady rate`;
|
|
return '';
|
|
}
|
|
|
|
// ---------- card bands from the log intake ----------
|
|
// rows: {label, line} where line is an app "GPUs:" line (label win-<id8> or mac-<id8>) or a worker STATUS line
|
|
// (label miner-<vendor>-<id8>-<n>). The worker label's trailing index is the 1-based position in the GPUs list of the
|
|
// same machine id (app/igneum-app/src/engine.rs names cards <vendor>-<machine>-<index>).
|
|
export function modelOf(text) {
|
|
const t = String(text);
|
|
if (/RTX 5090/.test(t)) return '5090';
|
|
if (/RX 9070 XT|gfx1201/.test(t)) return '9070 XT';
|
|
if (/Apple M5 Max/.test(t)) return 'M5 Max';
|
|
if (/Apple M4 Max/.test(t)) return 'M4 Max';
|
|
if (/Apple M\d/.test(t)) return t.match(/Apple M\d[^,;()]*/)[0].trim();
|
|
if (/gfx1036|Radeon\(TM\) Graphics/.test(t)) return 'gfx1036';
|
|
if (/UHD/.test(t)) return 'Intel UHD';
|
|
const m = /GeForce (RTX \d+[^,;()]*)|Radeon (RX [^,;()]*)/.exec(t); if (m) return (m[1] || m[2]).trim();
|
|
return null;
|
|
}
|
|
export function cardBands(rows, opts = {}) {
|
|
const o = { minUptime: 120, ...opts };
|
|
const gpus = new Map(); // machine id8 -> [model...]
|
|
for (const r of rows) { const m = /^(?:win|mac)-([0-9a-f]{8})$/.exec(r.label); const g = / GPUs: (.*)$/.exec(r.line || ''); if (m && g) gpus.set(m[1], g[1].split(';').map(s => modelOf(s))); }
|
|
const samples = new Map();
|
|
for (const r of rows) {
|
|
const m = /^miner-[a-z]+-([0-9a-f]{8})-(\d+)$/.exec(r.label); if (!m) continue;
|
|
const s = /STATUS '[^']+' \[worker\]: (\d+)s .*? now=([\d.]+) MH\/s wall/.exec(r.line || ''); if (!s) continue;
|
|
const uptime = Number(s[1]), now = Number(s[2]); if (!(now > 0) || uptime < o.minUptime) continue;
|
|
const list = gpus.get(m[1]); const model = (list && list[Number(m[2]) - 1]) || modelOf(r.line) || null; if (!model) continue;
|
|
if (!samples.has(model)) samples.set(model, []); samples.get(model).push(now);
|
|
}
|
|
const q = (a, p) => a[Math.min(a.length - 1, Math.floor(p * a.length))];
|
|
return [...samples].map(([model, v]) => { v.sort((a, b) => a - b); return { model, n: v.length, p5: q(v, 0.05), p50: q(v, 0.5), p95: q(v, 0.95) }; }).sort((a, b) => b.p50 - a.p50);
|
|
}
|
|
|
|
// ---------- the live path ----------
|
|
// ctx: {sql, TB (blocks table), TS (state table), recordEvent(kind, text), log}; opts from env. Settled epochs are read
|
|
// once and cached as cells; the open and unsettled epochs are re-read each run (two epochs at most, under 7,500 rows).
|
|
const PAGE = 5000;
|
|
export function makeRunner(ctx, opts = {}) {
|
|
const o = { ...DEFAULTS, ...opts };
|
|
const cache = new Map(); // epoch -> aggregate of that epoch (settled only)
|
|
let bands = [], bandsAt = 0, prev = null;
|
|
async function readEpoch(e) {
|
|
const rows = [];
|
|
for (let off = 0; ; off += PAGE) {
|
|
const page = await ctx.sql(`SELECT vote_key_hash, daa_score::bigint AS daa_score, timestamp_ms::bigint AS timestamp_ms, color, (detail->>'bits')::bigint AS bits, detail->>'nonce' AS nonce
|
|
FROM ${ctx.TB} WHERE daa_score >= $1 AND daa_score < $2 AND detail IS NOT NULL AND vote_key_hash IS NOT NULL
|
|
ORDER BY daa_score, timestamp_ms LIMIT ${PAGE} OFFSET ${off}`, [String(e * o.epochLen), String((e + 1) * o.epochLen)]);
|
|
rows.push(...page); if (page.length < PAGE) break;
|
|
}
|
|
return aggregate(rows, o);
|
|
}
|
|
async function readBands() {
|
|
// The app logs its "GPUs:" line once at start, so it is read from any upload of the last 7 days (rare line, small
|
|
// result); the workers' STATUS lines (every 10 s) come from the uploads of the last 6 hours, distinct, because the
|
|
// 262 KB rolling uploads overlap. Both queries are server-side regexp extractions: no whole upload crosses the wire.
|
|
const gpus = await ctx.sql(`SELECT DISTINCT label, m[1] AS line FROM miner_logs, LATERAL regexp_matches(lines, '([^\n]* GPUs: [^\n]*)', 'g') AS m
|
|
WHERE received_at > now() - interval '7 days' AND (label LIKE 'win-%' OR label LIKE 'mac-%')`);
|
|
const status = await ctx.sql(`SELECT DISTINCT label, m[1] AS line FROM miner_logs, LATERAL regexp_matches(lines, '([^\n]* STATUS ''[^'']+'' \\[worker\\]: [^\n]*)', 'g') AS m
|
|
WHERE received_at > now() - interval '6 hours' AND label LIKE 'miner-%'`);
|
|
return cardBands([...gpus, ...status], o);
|
|
}
|
|
return async function run({ dry = false, tipDaa = null } = {}) {
|
|
const tipRow = tipDaa === null ? await ctx.sql(`SELECT max(daa_score)::bigint AS tip FROM ${ctx.TB}`) : null;
|
|
const tip = tipDaa ?? Number(tipRow[0].tip || 0);
|
|
const tipEpoch = epochOf(tip, o.epochLen);
|
|
const wanted = []; for (let e = tipEpoch - o.windowEpochs - 1; e <= tipEpoch; e++) if (e >= 0) wanted.push(e);
|
|
const parts = [];
|
|
for (const e of wanted) {
|
|
const settled = tip >= (e + 1) * o.epochLen + o.settleDaa;
|
|
if (settled && cache.has(e)) { parts.push(cache.get(e)); continue; }
|
|
const a = await readEpoch(e); if (settled) cache.set(e, a); parts.push(a);
|
|
}
|
|
for (const e of [...cache.keys()]) if (!wanted.includes(e)) cache.delete(e);
|
|
if (Date.now() - bandsAt > 30 * 60_000) { try { bands = await readBands(); bandsAt = Date.now(); } catch (e) { ctx.log && ctx.log('detector bands read failed', e.message); } }
|
|
const { state, events } = analyse(mergeAggregates(parts), o, bands, prev, tip);
|
|
prev = state;
|
|
if (!dry) {
|
|
await ctx.sql(`UPDATE ${ctx.TS} SET detector = $1::jsonb WHERE id = 1`, [JSON.stringify(state)]);
|
|
for (const t of events) await ctx.recordEvent('detector', t);
|
|
}
|
|
return { state, events };
|
|
};
|
|
}
|
|
export function optsFromEnv(env = process.env) {
|
|
const o = {};
|
|
if (env.DETECTOR_WINDOW_EPOCHS) o.windowEpochs = Number(env.DETECTOR_WINDOW_EPOCHS);
|
|
if (env.DETECTOR_R) o.r = Number(env.DETECTOR_R);
|
|
if (env.DETECTOR_K) o.k = Number(env.DETECTOR_K);
|
|
if (env.DETECTOR_HOLD) o.holdWindows = Number(env.DETECTOR_HOLD);
|
|
if (env.DETECTOR_EPOCH_LEN) o.epochLen = Number(env.DETECTOR_EPOCH_LEN);
|
|
return o;
|
|
}
|
|
|
|
// ---------- dry run (read-only against the live tables) ----------
|
|
if (process.argv[1] && process.argv[1].endsWith('detector.mjs') && process.argv.includes('--dry')) {
|
|
const env = readFileSync(`${homedir()}/.config/igneum/env`, 'utf8');
|
|
const m = /^DATABASE_URL=(.*)$/m.exec(env); if (!m) { console.error('DATABASE_URL not found in ~/.config/igneum/env'); process.exit(1); }
|
|
const url = m[1].trim().replace(/^['"]|['"]$/g, ''); const host = new URL(url).hostname.replace('-pooler', '');
|
|
const sql = async (query, params = []) => { if (!/^\s*(select|with)/i.test(query)) throw new Error('dry run: read-only'); 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 T = (process.env.LIVE_TABLE_PREFIX || '').replace(/[^a-z0-9_]/gi, '');
|
|
const run = makeRunner({ sql, TB: `${T}live_blocks`, TS: `${T}live_state`, recordEvent: async () => { }, log: console.log }, optsFromEnv());
|
|
const { state, events } = await run({ dry: true });
|
|
if (process.argv.includes('--json')) { console.log(JSON.stringify(state, null, 1)); }
|
|
else {
|
|
console.log(`window epochs ${state.window.epochs.join(' ')} (tip DAA ${state.tip_daa}, ${state.window.ids} ids, ${state.window.ids_correlated} correlated)`);
|
|
console.table(state.network);
|
|
console.table(state.miners.map(m => ({ id: m.id, epochs: m.epochs_present, blue: m.blue, mhs: m.mhs, steady: m.steady, step_pct: m.max_step_pct, spread_pct: m.spread_sd_pct, poisson_pct: m.poisson_sd_pct, excess_pct: m.excess_spread_pct, early_pct: m.early_share_pct, nonce_n: m.nonce.n, chi2_lo: m.nonce.chi2_low4, chi2_hi: m.nonce.chi2_high4, inc: m.nonce.inc_frac, band: m.band ? (m.band.matches.join('|') || (m.band.high ? 'HIGH' : m.band.small ? 'small' : 'none')) : '-', flags: m.flags.join(',') + (m.notes.length ? ' (' + m.notes.join(',') + ')' : '') })));
|
|
console.log('bands:', state.bands.map(b => `${b.model} ${b.p5}/${b.p50}/${b.p95} MH/s (n ${b.n})`).join('; ') || 'none');
|
|
console.log('correlation:', JSON.stringify(state.correlation));
|
|
console.log('alert:', JSON.stringify(state.alert));
|
|
console.log('events:', events.length ? events : 'none');
|
|
}
|
|
}
|