#!/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 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- (a queued build) and gate-pending- (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); } }