igneum/tools/observer/detector.mjs
igneum-josh 9c7838e5b3 Counter ASIC 3.0 items 4 and 5: the epoch common factor is the median over steady core ids
The mean over every present id let one paused-and-resumed card (the Mac, a 178% step) push
every other residual the same way in the epochs it was off, which read as an r = 0.94 edge
between two honest 5090 keys on the merged tree (window 41 to 46). The factor is now the
median over ids that are steady and present in every window epoch (median over all present
when the core is under 3): the same window reads max r 0.10. README: the three calibration
readings (the edge, the factor-of-two from identities=2, the unsteady Mac) answered.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-06 09:28:55 +01:00

354 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, clipped. The factor is the MEDIAN over the core ids (steady
// and present in every epoch of the window); the mean over everyone present let one paused-and-resumed card (the Mac,
// 6 October, a 178% step) push every other id's residual the same way in the epochs it was off, which read as an
// r = 0.94 edge between two honest keys. Fewer than 3 core ids: the median over all present ids.
const idMean = new Map([...logRate].map(([id, lr]) => [id, lr.size ? mean([...lr.values()]) : 0]));
const core = ids.filter(id => miners.get(id).steady && W.every(e => logRate.get(id).has(e)));
const epochFactor = new Map(W.map(e => {
const pool = core.length >= 3 ? core : ids.filter(id => logRate.get(id).has(e));
const devs = pool.filter(id => logRate.get(id).has(e)).map(id => logRate.get(id).get(e) - idMean.get(id));
return [e, devs.length ? median(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');
}
}