igneum/tools/workers/collect.mjs
igneum-labs 942a58a0ef Worker dashboard: collector on igneum-build-1, Mac pusher, workers.html next to the fleet page (the project lead, 6 October 2026)
tools/workers/collect.mjs runs every 30 s on the box (infra/build-server/workers/igneum-workers.{service,timer},
install.sh) and writes /srv/workers/workers.json from the machine itself: /proc/stat deltas per core, meminfo, df on
/srv, net bytes, hwmon temperatures, the flock state of /srv/builds/_locks, cargo processes with worktree, target and
start time, flock waiters, /srv/builds/_log/builds.jsonl (the format agreed with the build-server agent), sccache
--show-stats, headline.json. tools/workers/push.mjs (launchd every 60 s) adds the Mac's with-lock slots and waiters
and the two PCs' relay job states from the intake, drops them on the box so its file is whole, merges the box's file
and writes the fleet folder's workers.json; it deploys only when the live copy is over 6 min old, otherwise the
fleet orchestrator's 5-minute deploy carries it. The page (tools/workers/page/workers.html, shape.js) tries
https://build.igneum.network/workers.json first and falls back to the folder copy; core strip, arcs, now building,
queue, recently done with closing lines, headline timings, analytics; UK time with UTC tooltips; phone width.
node --test tools/workers/test: 13 tests over /proc/stat, lock dir, JSONL and sccache fixtures and the page shaping.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-06 17:59:07 +00:00

160 lines
11 KiB
JavaScript

#!/usr/bin/env node
// The worker dashboard's collector. Runs on igneum-build-1 every 30 s (infra/build-server/workers/igneum-workers.timer)
// as the build user and writes /srv/workers/workers.json: what the MACHINE says (uptime, load, /proc/stat deltas per core,
// memory, the /srv disk, network bytes, hwmon temperatures), the slot and lock state under /srv/builds/_locks, the cargo
// processes with their worktree and target, the slot waiters, the recent builds from /srv/builds/_log/builds.jsonl and
// sccache's own counters. Nothing here is a reported number: every field is read from /proc, /sys, the lock files, the
// process table or the log the build itself appended.
//
// node tools/workers/collect.mjs write /srv/workers/workers.json (IGNEUM_WORKERS_DIR overrides)
// node tools/workers/collect.mjs --stdout print instead of writing
// node tools/workers/collect.mjs --sample-ms 1000 the /proc/stat and network delta window (default 1000)
//
// The Mac side (tools/workers/push.mjs) drops its own facts at /srv/workers/sources/mac.json (the Mac's with-lock slots) and
// /srv/workers/sources/pcs.json (the two PCs' relay job states); when present and under 10 minutes old they are merged
// into the same workers.json, so one file answers for the whole shop whichever side serves it. Needs Node 22, no deps.
import { readFileSync, readdirSync, writeFileSync, renameSync, mkdirSync, existsSync, statSync, readlinkSync, openSync, closeSync } from 'node:fs';
import { spawnSync } from 'node:child_process';
import { hostname } from 'node:os';
import { join, dirname } from 'node:path';
import { fileURLToPath } from 'node:url';
import { parseProcStat, cpuBusy, parseMeminfo, parseNetDev, parseLoadavg, parseLockDir, parseBuildsJsonl, parseSccacheStats, pickTemps, cargoProcess, deriveKind, BASELINES } from './lib.mjs';
const HERE = dirname(fileURLToPath(import.meta.url));
const argv = process.argv.slice(2);
const flag = (n, d) => { const i = argv.indexOf(`--${n}`); return i >= 0 ? (argv[i + 1] ?? true) : d; };
const OUT_DIR = process.env.IGNEUM_WORKERS_DIR || '/srv/workers';
const LOCKS = process.env.IGNEUM_BUILD_SLOTS_DIR || '/srv/builds/_locks';
const BUILDS = process.env.IGNEUM_BUILD_ROOT || '/srv/builds';
const LOG = process.env.IGNEUM_BUILDS_LOG || join(BUILDS, '_log', 'builds.jsonl');
const SAMPLE_MS = Number(flag('sample-ms', 1000));
const iso = ms => new Date(ms).toISOString().replace(/\.\d{3}Z$/, 'Z');
const read = p => { try { return readFileSync(p, 'utf8'); } catch { return ''; } };
const sh = (cmd, args, opts = {}) => { const r = spawnSync(cmd, args, { encoding: 'utf8', timeout: 8000, ...opts }); return r.status === 0 ? r.stdout : ''; };
// ---- machine facts ----------------------------------------------------------------------------------------------------
function machine() {
const stat0 = parseProcStat(read('/proc/stat')); const net0 = parseNetDev(read('/proc/net/dev')); const t0 = Date.now();
const end = t0 + SAMPLE_MS; while (Date.now() < end) Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, end - Date.now());
const stat1 = parseProcStat(read('/proc/stat')); const net1 = parseNetDev(read('/proc/net/dev')); const dt = (Date.now() - t0) / 1000;
const cpu = cpuBusy(stat0, stat1);
const up = Number(read('/proc/uptime').split(' ')[0]) || 0;
const la = parseLoadavg(read('/proc/loadavg'));
const mem = parseMeminfo(read('/proc/meminfo'));
const df = sh('df', ['-B1', '--output=target,size,used,avail', BUILDS]).trim().split('\n').pop() || '';
const dfp = df.trim().split(/\s+/);
const disk = dfp.length >= 4 ? { mount: dfp[0], total_bytes: Number(dfp[1]), used_bytes: Number(dfp[2]), avail_bytes: Number(dfp[3]), used_pct: Number(dfp[1]) ? Math.round((Number(dfp[2]) / Number(dfp[1])) * 1000) / 10 : 0 } : null;
// the busiest interface carries the build traffic (rsync in, artefacts out)
const nets = net1.map(n => { const b = net0.find(x => x.iface === n.iface) || n; return { ...n, rx_bps: Math.max(0, Math.round((n.rx_bytes - b.rx_bytes) * 8 / dt)), tx_bps: Math.max(0, Math.round((n.tx_bytes - b.tx_bytes) * 8 / dt)) }; }).filter(n => n.rx_bytes + n.tx_bytes > 0).sort((a, b) => (b.rx_bytes + b.tx_bytes) - (a.rx_bytes + a.tx_bytes));
const sensors = [];
for (const h of safeList('/sys/class/hwmon')) {
const dir = join('/sys/class/hwmon', h); const name = read(join(dir, 'name')).trim();
for (const f of safeList(dir)) { const m = /^(temp\d+)_input$/.exec(f); if (!m) continue; sensors.push({ name, label: read(join(dir, `${m[1]}_label`)).trim() || null, millic: Number(read(join(dir, f))) }); }
}
const osr = read('/etc/os-release'); const pretty = (/^PRETTY_NAME="?([^"\n]+)"?/m.exec(osr) || [])[1] || null;
const clk = Number(sh('getconf', ['CLK_TCK']).trim()) || 100;
return { name: hostname(), os: pretty, cores: stat1.cores.length, uptime_s: Math.round(up), boot_at: iso(Date.now() - up * 1000), ...la, cpu, mem, disk, net: nets[0] || null, temps: pickTemps(sensors), clk, kernel: read('/proc/sys/kernel/osrelease').trim() || null };
}
const safeList = d => { try { return readdirSync(d); } catch { return []; } };
const alive = pid => { try { process.kill(pid, 0); return true; } catch (e) { return e.code === 'EPERM'; } };
// flock probe without python: `flock -n <file> true` fails while the slot is held (the slot is held by another process's fd)
const flockHeld = p => { const r = spawnSync('flock', ['-n', p, 'true'], { timeout: 3000 }); return r.status !== 0; };
// ---- slots, waiters, running builds ---------------------------------------------------------------------------------------
function slots(now) {
const files = safeList(LOCKS).map(n => ({ name: n, text: read(join(LOCKS, n)), held: /^(build-\d+|run-\d+|build|measure)$/.test(n) ? flockHeld(join(LOCKS, n)) : undefined }));
return parseLockDir(files, { now, aliveFn: alive });
}
function procs(bootSec, clk, now) {
const running = [], waiters = []; let compilers = 0;
for (const d of safeList('/proc')) {
if (!/^\d+$/.test(d)) continue;
const pid = Number(d); const comm = read(`/proc/${d}/comm`).trim();
if (['rustc', 'cc1plus', 'cc1', 'ld', 'ld.lld', 'x86_64-w64-mingw32-ld'].includes(comm)) compilers++;
if (comm === 'cargo') {
let cwd = null; try { cwd = readlinkSync(`/proc/${d}/cwd`); } catch { cwd = null; }
const p = cargoProcess({ pid, cwd, cmdline: read(`/proc/${d}/cmdline`), stat: read(`/proc/${d}/stat`), bootSec, clk, now, root: BUILDS });
if (p) running.push(p);
}
if (comm === 'flock') {
const args = read(`/proc/${d}/cmdline`).split('\0').filter(Boolean);
if (args[1] === '-w' && /^\d+$/.test(args[2] || '')) {
const stat = read(`/proc/${d}/stat`); const rest = stat.slice(stat.lastIndexOf(')') + 2).split(' ');
const ppid = Number(rest[1]); const startTicks = Number(rest[19]);
const since = Number.isFinite(startTicks) ? iso((bootSec + startTicks / clk) * 1000) : null;
const wait = read(join(LOCKS, `wait-${ppid}`)).trim(); // a waiter's own line, when the appender writes one
waiters.push({ pid: ppid, since, waiting_s: since ? Math.max(0, Math.round((now - Date.parse(since)) / 1000)) : null, label: wait || null, ...(wait ? { kind: deriveKind(parseLabelSafe(wait)) } : {}) });
}
}
}
// the innermost cargo per worktree: a `cargo build` spawns no nested cargo, but a `cargo test` runs test binaries, not cargo
running.sort((a, b) => String(a.started_at).localeCompare(String(b.started_at)));
return { running, waiters, compilers };
}
function parseLabelSafe(line) { const r = /: (.*)$/.exec(line); return { command: r ? r[1] : line }; }
function sccache() {
const bin = existsSync('/home/build/.cargo/bin/sccache') ? '/home/build/.cargo/bin/sccache' : 'sccache';
const out = sh(bin, ['--show-stats'], { env: { ...process.env, SCCACHE_DIR: process.env.SCCACHE_DIR || '/srv/sccache' } });
return out ? parseSccacheStats(out) : null;
}
// ---- sources dropped by the Mac ------------------------------------------------------------------------------------------
function source(name, now) {
const p = join(OUT_DIR, 'sources', `${name}.json`);
try {
const st = statSync(p); const age = Math.round((now - st.mtimeMs) / 1000);
const j = JSON.parse(readFileSync(p, 'utf8'));
return { data: j, meta: { at: iso(st.mtimeMs), age_s: age, stale: age > 600, ok: true } };
} catch (e) { return { data: null, meta: { ok: false, error: e.code === 'ENOENT' ? 'the Mac has not pushed yet' : String(e.message || e) } }; }
}
// ---- main -------------------------------------------------------------------------------------------------------------------
export function collect() {
const now = Date.now();
const m = machine();
const bootSec = Math.round((now - m.uptime_s * 1000) / 1000);
const s = slots(now);
const p = procs(bootSec, m.clk, now);
// attach the slot holder to its build (the label names the worktree and crate the cargo runs in)
for (const r of p.running) {
const h = s.held.find(h => h.worktree === r.worktree && (!h.crate || h.crate === r.crate || !r.crate));
r.slot = h ? h.slot : null; r.agent = h ? h.agent : null; r.queued_at = h ? h.since : null; r.waited_s = h ? h.waited_s : null;
}
const logText = read(LOG);
const recent = parseBuildsJsonl(logText, { limit: 200 });
const mac = source('mac', now), pcs = source('pcs', now);
let headline = null; try { headline = JSON.parse(readFileSync(join(OUT_DIR, 'headline.json'), 'utf8')); } catch { headline = null; }
const doc = {
v: 1,
generated_at: iso(now),
generated_by: 'collect.mjs on ' + m.name,
box: {
...m, clk: undefined,
collected_at: iso(now),
sccache: sccache(),
slots: { count: s.count, held: s.held, free: s.free },
running: p.running,
queue: p.waiters,
compilers: p.compilers,
recent: recent.rows,
log: { path: LOG, present: !!logText, bad_lines: recent.bad, total: recent.total, note: logText ? null : 'builds.jsonl is not written yet: the build-server agent adds the append to bs_remote_run; until then this lane fills from nothing' },
},
mac: mac.data, pcs: pcs.data,
headline, baselines: BASELINES,
sources: { box: { ok: true, at: iso(now) }, mac: mac.meta, pcs: pcs.meta },
};
return doc;
}
if (process.argv[1] && fileURLToPath(import.meta.url) === process.argv[1]) {
const doc = collect();
const text = JSON.stringify(doc);
if (flag('stdout', false)) { process.stdout.write(text + '\n'); }
else {
mkdirSync(OUT_DIR, { recursive: true }); mkdirSync(join(OUT_DIR, 'sources'), { recursive: true });
const tmp = join(OUT_DIR, '.workers.json.tmp'); writeFileSync(tmp, text); renameSync(tmp, join(OUT_DIR, 'workers.json'));
const fd = openSync(join(OUT_DIR, 'last-run'), 'w'); closeSync(fd);
}
}