igneum/tools/workers/collect.mjs
igneum-labs f9b5b26f94 Worker dashboard: scheduling classes on the cards (kind, nice, cores), gate-queued marker, wait and gate-pending files
The slot line now ends "; kind=<k> nice=<n> cores=<c>; agent=<a>" (build box scheduling, 51304fa5): lib.mjs lifts the
three into their own fields and keeps the bare label; the old line shape still derives its kind. The collector reads
/srv/builds/_locks/wait-<pid> and gate-pending-<pid> directly (dead pids skipped), marks a gate with priority gate
and sorts it first; the page shows the class line on Now building and Queue cards, a "gate queued" pill, and the class
in the closing lines. JSONL nice, cores, jobs and priority carried through. 14 tests pass.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
(cherry picked from commit 9fce5742a6)
2026-10-07 10:35:07 +00:00

180 lines
12 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, parseSlotLine, 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;
const flockParents = new Map(); // ppid -> since, from the flock -w processes the waiters run
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]);
flockParents.set(ppid, Number.isFinite(startTicks) ? iso((bootSec + startTicks / clk) * 1000) : null);
}
}
}
// the waiters' own lines: wait-<pid> (a queued build) and gate-pending-<pid> (a gate that takes the next free slot
// ahead of queued suites and benches), written by remote-run.sh in the slot-file format and removed on take or give-up;
// a file whose pid is gone is a leftover, not a waiter
const seen = new Set();
for (const f of safeList(LOCKS)) {
const m = /^(wait|gate-pending)-(\d+)$/.exec(f); if (!m) continue;
const pid = Number(m[2]); if (!alive(pid) || seen.has(pid)) continue; seen.add(pid);
const line = parseSlotLine(read(join(LOCKS, f)).split('\n')[0], { now });
const since = (line && line.since) || flockParents.get(pid) || null;
waiters.push({ pid, since, waiting_s: since ? Math.max(0, Math.round((now - Date.parse(since)) / 1000)) : null, label: line ? line.label : null, kind: line ? deriveKind(line) : null, agent: line ? line.agent : null, worktree: line ? line.worktree : null, crate: line ? line.crate : null, command: line ? line.command : null, nice: line ? line.nice : null, cores: line ? line.cores : null, priority: m[1] === 'gate-pending' ? 'gate' : 'normal' });
}
for (const [ppid, since] of flockParents) if (!seen.has(ppid)) waiters.push({ pid: ppid, since, waiting_s: since ? Math.max(0, Math.round((now - Date.parse(since)) / 1000)) : null, label: null, priority: 'normal' });
waiters.sort((a, b) => (a.priority === 'gate' ? 0 : 1) - (b.priority === 'gate' ? 0 : 1) || String(a.since || '').localeCompare(String(b.since || '')));
running.sort((a, b) => String(a.started_at).localeCompare(String(b.started_at)));
return { running, waiters, compilers };
}
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;
if (h && h.class) r.kind = h.class; r.nice = h ? h.nice : null; r.cores = h ? h.cores : 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; }
// the background capacity layer (infra/build-server/capacity) writes capacity.json beside this file; show it when
// present and under 10 min old (the layer's controller stamps it every slice and on every pause/resume)
let background = null;
try {
const p = join(OUT_DIR, 'capacity.json'); const st = statSync(p); const age = Math.round((now - st.mtimeMs) / 1000);
const j = JSON.parse(readFileSync(p, 'utf8'));
background = { ...j, meta: { at: iso(st.mtimeMs), age_s: age, stale: age > 600, ok: true } };
} catch (e) { background = { meta: { ok: false, error: e.code === 'ENOENT' ? 'the capacity layer has not written yet' : String(e.message || e) } }; }
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,
background,
sources: { box: { ok: true, at: iso(now) }, mac: mac.meta, pcs: pcs.meta, background: background && background.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);
}
}