igneum/tools/workers/push.mjs
igneum-josh d1036db7fb Worker dashboard: collector on igneum-build-1, Mac pusher, workers.html next to the fleet page (Josh, 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 18:59:07 +01:00

205 lines
14 KiB
JavaScript

#!/usr/bin/env node
// The worker dashboard's Mac side, run by launchd every 60 s (tools/workers/launchd/com.igneum.workers-push.plist).
// Gathers what only this Mac can see, merges it with the box's file and publishes:
// 1. the Mac's with-lock slots (/tmp/igneum-locks: build-0..2, run-0..2, measure, build) and their waiters
// 2. the two PCs' relay job states from the log intake (Neon, the same reader as tools/jobs.mjs): the running job,
// the queue (signed jobs for that machine with no report yet), the last 20 closing reports with durations
// 3. scp of those two facts to the box (/srv/workers/sources/{mac,pcs}.json) so the box's own file is whole too
// 4. the box's /srv/workers/workers.json pulled back, the Mac and PC facts overlaid, headline.json and the Mac baselines
// added, written to $DLSITE/<fleet-path>/workers.json
// 5. a Vercel deploy ONLY when the live copy is older than 6 minutes: the fleet orchestrator deploys the same folder every
// 5 minutes (tools/fleet/autorun.py), so the file normally rides along with it; one deploy per publish is the
// publish-fleet.sh mechanism and a deploy a minute is not wanted.
//
// node tools/workers/push.mjs one pass (what launchd runs)
// node tools/workers/push.mjs --stdout print the merged document, publish nothing
// node tools/workers/push.mjs --deploy force the deploy this pass
// node tools/workers/push.mjs --no-deploy never deploy this pass (write the folder copy only)
// Reads ~/.config/igneum/{build-server,env,dl-token,dlsite-dir,fleet-path} and ~/.ssh/igneum_ed25519. Prints nothing
// secret. Log: ~/Library/Logs/igneum-workers-push.log (launchd) or stderr.
import { readFileSync, readdirSync, writeFileSync, renameSync, mkdirSync, existsSync, statSync } from 'node:fs';
import { spawnSync, execFileSync } from 'node:child_process';
import { homedir, hostname, tmpdir } from 'node:os';
import { join, dirname } from 'node:path';
import { fileURLToPath } from 'node:url';
import { parseLockDir, summaryOf, deriveKind, BASELINES } from './lib.mjs';
const HERE = dirname(fileURLToPath(import.meta.url));
const argv = process.argv.slice(2);
const has = f => argv.includes(`--${f}`);
const cfg = n => { try { return readFileSync(join(homedir(), '.config', 'igneum', n), 'utf8').trim(); } catch { return ''; } };
const iso = ms => new Date(ms).toISOString().replace(/\.\d{3}Z$/, 'Z');
const log = (...a) => console.error(`${iso(Date.now())} push: ${a.join(' ')}`);
const KEY = process.env.IGNEUM_BUILD_KEY || join(homedir(), '.ssh', 'igneum_ed25519');
const HOST = (cfg('build-server') || '').split('\n')[0].trim(); // build@<ip>
const DLSITE = cfg('dlsite-dir'), FLEET = cfg('fleet-path');
const STALE_LIVE_S = 360;
const SSH_OPTS = ['-i', KEY, '-o', 'BatchMode=yes', '-o', 'StrictHostKeyChecking=accept-new', '-o', 'ConnectTimeout=10', '-o', 'ControlMaster=auto', '-o', `ControlPath=${join(homedir(), '.ssh', 'cm', 'igneum-workers-%r@%h:%p')}`, '-o', 'ControlPersist=600'];
mkdirSync(join(homedir(), '.ssh', 'cm'), { recursive: true });
// ---- 1. the Mac's slots ----------------------------------------------------------------------------------------------------
const alive = pid => { try { process.kill(pid, 0); return true; } catch (e) { return e.code === 'EPERM'; } };
export function macSlots(now = Date.now(), dir = '/tmp/igneum-locks') {
let files = [];
try { files = readdirSync(dir).map(name => ({ name, text: readFileSync(join(dir, name), 'utf8') })); } catch { files = []; }
const r = parseLockDir(files, { now, aliveFn: alive });
let cap = 3; try { cap = Math.max(1, Math.min(3, parseInt(cfg('build-slots'), 10) || 3)); } catch { cap = 3; }
// waiters: the python with-lock.sh exec'd into, still in its wait loop (`python3 - build /tmp/igneum-locks <cmd>`); the
// holder exec'd on into its command, so a process still carrying that argv is waiting. Our own argv never matches.
const waiters = [];
const ps = spawnSync('ps', ['-eo', 'pid=,etimes=,args='], { encoding: 'utf8' }).stdout || '';
for (const line of ps.split('\n')) {
const m = /^\s*(\d+)\s+(\d+)\s+(\S*python3?) - (build|measure|run) \/tmp\/igneum-locks (.*)$/.exec(line);
if (!m || Number(m[1]) === process.pid) continue;
const held = r.held.find(h => h.pid === Number(m[1]));
if (held) continue; // measure mode keeps python in front? no: it execs too; a listed one is a waiter
waiters.push({ pid: Number(m[1]), mode: m[4], waiting_s: Number(m[2]), command: m[5].slice(0, 160), kind: deriveKind({ command: m[5] }) });
}
for (const h of r.held) h.worktree = h.worktree || worktreeOfPid(h.pid);
return { name: hostname().replace(/\.local$/, ''), collected_at: iso(now), build_slots: cap, slots: { count: cap, held: r.held, free: r.free }, queue: waiters };
}
// the holder's cwd names the worktree it builds in (lsof is slow; /proc does not exist here, so `lsof -p` only when needed)
function worktreeOfPid(pid) {
const out = spawnSync('lsof', ['-a', '-p', String(pid), '-d', 'cwd', '-Fn'], { encoding: 'utf8', timeout: 4000 }).stdout || '';
const m = /\n n(\/Users\/[^\n]*)/.exec('\n' + out.replace(/^n/m, ' n'));
const p = m ? m[1] : (out.split('\n').find(l => l.startsWith('n/')) || '').slice(1);
const w = /\/Projects\/(igneum(?:-wt-[\w.-]+)?)/.exec(p);
return w ? w[1] : null;
}
// ---- 2. the PCs through the intake ---------------------------------------------------------------------------------------
const MACHINES = [
{ id: 'ae432dc7', name: 'PC 1', role: 'builds', note: 'Linux and Windows binaries, Josh\'s desk' },
{ id: '1ccfe586', name: 'PC 2', role: 'test suites', note: 'cargo test suites and shard runs' },
];
function db() {
const m = /^DATABASE_URL=(.*)$/m.exec(cfg('env'));
if (!m) return null;
const url = m[1].trim().replace(/^['"]|['"]$/g, '');
const host = new URL(url).hostname.replace('-pooler', '');
return async (query, params = []) => {
const r = await fetch(`https://${host}/sql`, { method: 'POST', headers: { 'Neon-Connection-String': url, 'Content-Type': 'application/json' }, body: JSON.stringify({ query, params }), signal: AbortSignal.timeout(15000) });
const j = await r.json();
if (!r.ok) throw new Error(j.message || `http ${r.status}`);
return j.rows;
};
}
async function signedJobs() {
const tok = cfg('dl-token'); if (!tok) return [];
try {
const r = await fetch(`https://dl.igneum.network/dl/${tok}/igneum-jobs.signed.json?t=${Date.now()}`, { headers: { 'Cache-Control': 'no-cache' }, signal: AbortSignal.timeout(15000) });
if (!r.ok) return [];
const env = await r.json(); const f = JSON.parse(env.file);
return Array.isArray(f.jobs) ? f.jobs : [];
} catch (e) { log('jobs file:', e.message); return []; }
}
export function shapePcs(rows, jobs, now = Date.now()) {
// rows: the newest upload per run_id with its machine; a run_id is job-<id>-<machine id>
const byMachine = new Map(MACHINES.map(m => [m.id, { ...m, machine: null, running: null, queue: [], recent: [], last_report_at: null }]));
const seenJob = new Set();
for (const r of rows) {
const mid = (/-([0-9a-f]{8})$/.exec(r.machine || '') || [])[1] || (/-([0-9a-f]{8})$/.exec(r.run_id || '') || [])[1];
const pc = byMachine.get(mid); if (!pc) continue;
pc.machine = pc.machine || r.machine;
const s = summaryOf(r.lines); if (!s) continue;
const jobId = s.job || (r.run_id || '').replace(/^job-/, '').replace(new RegExp(`-${mid}$`), '');
seenJob.add(jobId);
const kind = s.kind || (jobId.split('-')[0]) || 'job';
const at = Date.parse(String(r.received_at).replace(' ', 'T').replace(/([+-]\d\d)$/, '$1:00'));
if (!pc.last_report_at || at > Date.parse(pc.last_report_at)) pc.last_report_at = iso(at);
if (s.status === 'running') {
// a progress report under 15 minutes old is a live job; older is a job whose closing report never came
const fresh = now - at < 15 * 60_000;
const entry = { job: jobId, kind, title: titleOf(jobs, jobId), started_at: s.started_at || null, stage: s.stage, elapsed_s: s.started_at ? Math.max(0, Math.round((now - Date.parse(s.started_at)) / 1000)) : null, last_report_at: iso(at), stale: !fresh };
if (!pc.running || Date.parse(entry.last_report_at) > Date.parse(pc.running.last_report_at)) pc.running = entry;
} else {
pc.recent.push({ job: jobId, kind, title: titleOf(jobs, jobId), status: s.status, ok: s.status === 'done' && Number(s.exit) === 0, exit: s.exit, secs: s.duration_s, started_at: s.started_at || null, finished_at: s.finished_at && s.finished_at !== 'closed' ? s.finished_at : iso(at), summary: String(s.summary || '').slice(0, 240), errors: (s.errors || []).map(e => e.slice(0, 200)), stage: s.stage });
}
}
for (const pc of byMachine.values()) {
pc.recent.sort((a, b) => String(b.finished_at).localeCompare(String(a.finished_at))); pc.recent = pc.recent.slice(0, 20);
if (pc.running && pc.running.stale) { pc.running = null; }
// the queue: signed, unexpired jobs aimed at this machine that no report has mentioned yet
for (const j of jobs) {
const t = j.target || {}; const ids = Array.isArray(t.machine_ids) ? t.machine_ids : t.machine_ids ? [t.machine_ids] : [];
if (!ids.includes(pc.id)) continue;
if (Date.parse(j.expires_at) < now) continue;
if (seenJob.has(j.id)) continue;
pc.queue.push({ job: j.id, kind: j.kind, title: j.title || '', expires_at: j.expires_at });
}
pc.queue.sort((a, b) => String(a.job).localeCompare(String(b.job)));
}
return { collected_at: iso(now), machines: [...byMachine.values()] };
}
const titleOf = (jobs, id) => { const j = jobs.find(j => j.id === id); return j ? (j.title || '') : ''; };
async function pcs() {
const sql = db(); if (!sql) return { error: 'no DATABASE_URL in ~/.config/igneum/env', collected_at: iso(Date.now()), machines: [] };
const rows = await sql(`SELECT DISTINCT ON (run_id) run_id, machine, received_at, lines FROM miner_logs WHERE run_id LIKE 'job-%' AND label LIKE 'job-%' AND received_at > now() - interval '7 days' ORDER BY run_id, received_at DESC`);
const jobs = await signedJobs();
return shapePcs(rows, jobs);
}
// ---- 3 and 4. the box -------------------------------------------------------------------------------------------------------
function ssh(args, input) { return spawnSync('ssh', [...SSH_OPTS, HOST, ...args], { encoding: 'utf8', input, timeout: 30000, maxBuffer: 16 * 1024 * 1024 }); }
function dropToBox(name, doc) {
const r = ssh([`cat > /srv/workers/sources/.${name}.tmp && mv /srv/workers/sources/.${name}.tmp /srv/workers/sources/${name}.json`], JSON.stringify(doc));
if (r.status !== 0) throw new Error(`scp ${name}: ${(r.stderr || '').trim().slice(0, 200) || 'ssh failed'}`);
}
function pullFromBox() {
const r = ssh(['cat /srv/workers/workers.json; echo; echo __SEP__; cat /srv/workers/headline.json 2>/dev/null || echo null']);
if (r.status !== 0) throw new Error((r.stderr || '').trim().slice(0, 200) || 'ssh failed');
const [w, h] = r.stdout.split('__SEP__');
return { workers: JSON.parse(w), headline: JSON.parse((h || 'null').trim() || 'null') };
}
// ---- 5. publish ----------------------------------------------------------------------------------------------------------------
async function liveAge() {
try {
const r = await fetch(`https://dl.igneum.network/${FLEET}/workers.json?t=${Date.now()}`, { cache: 'no-store', signal: AbortSignal.timeout(10000) });
if (!r.ok) return Infinity;
const j = await r.json(); return (Date.now() - Date.parse(j.generated_at)) / 1000;
} catch { return Infinity; }
}
function deploy() {
const r = spawnSync('npx', ['--yes', 'vercel@latest', '--global-config', join(homedir(), '.config', 'igneum', 'vercel'), 'deploy', '--prod', '--yes'], { cwd: DLSITE, encoding: 'utf8', timeout: 240000, env: { ...process.env, PATH: `/opt/homebrew/bin:/usr/local/bin:${process.env.PATH || ''}` } });
return r.status === 0;
}
async function main() {
const now = Date.now();
const doc = { v: 1, generated_at: iso(now), generated_by: `push.mjs on ${hostname().replace(/\.local$/, '')}`, box: null, mac: null, pcs: null, sources: {}, baselines: BASELINES, headline: null };
// the Mac
try { doc.mac = macSlots(now); doc.sources.mac = { ok: true, at: doc.mac.collected_at }; } catch (e) { doc.sources.mac = { ok: false, error: String(e.message || e) }; }
// the PCs
try { doc.pcs = await pcs(); doc.sources.pcs = doc.pcs.error ? { ok: false, error: doc.pcs.error } : { ok: true, at: doc.pcs.collected_at }; } catch (e) { doc.sources.pcs = { ok: false, error: `intake: ${String(e.message || e).slice(0, 160)}` }; }
// the box: drop ours, pull theirs
if (!HOST) doc.sources.box = { ok: false, error: 'no ~/.config/igneum/build-server' };
else {
try {
if (doc.mac) dropToBox('mac', doc.mac);
if (doc.pcs && !doc.pcs.error) dropToBox('pcs', doc.pcs);
const { workers, headline } = pullFromBox();
doc.box = workers.box; doc.headline = headline;
const age = (now - Date.parse(workers.generated_at)) / 1000;
doc.sources.box = { ok: age < 120, at: workers.generated_at, age_s: Math.round(age), ...(age >= 120 ? { error: `the box's collector last wrote ${Math.round(age)} s ago (timer igneum-workers.timer)` } : {}) };
} catch (e) { doc.sources.box = { ok: false, error: `ssh to the box failed: ${String(e.message || e).slice(0, 160)}` }; }
}
if (!doc.headline) doc.headline = null;
const text = JSON.stringify(doc);
if (has('stdout')) { process.stdout.write(text + '\n'); return; }
if (!DLSITE || !FLEET) { log('no dlsite-dir or fleet-path in ~/.config/igneum; nothing published'); return; }
const dir = join(DLSITE, FLEET); mkdirSync(dir, { recursive: true });
const tmp = join(dir, '.workers.json.tmp'); writeFileSync(tmp, text); renameSync(tmp, join(dir, 'workers.json'));
let why = 'folder copy written; the fleet orchestrator deploys the folder within 5 min';
if (!has('no-deploy')) {
const age = has('deploy') ? Infinity : await liveAge();
if (age > STALE_LIVE_S) { const ok = deploy(); why = ok ? `deployed (live copy was ${age === Infinity ? 'missing' : Math.round(age) + ' s old'})` : 'DEPLOY FAILED'; }
else why = `live copy ${Math.round(age)} s old, no deploy`;
}
log(`box ${doc.sources.box.ok ? 'ok' : 'NO: ' + doc.sources.box.error}; mac ${doc.sources.mac.ok ? `${doc.mac.slots.held.length} held, ${doc.mac.queue.length} waiting` : 'NO'}; pcs ${doc.sources.pcs.ok ? doc.pcs.machines.map(m => `${m.name} ${m.running ? 'running ' + m.running.job : 'idle'} q${m.queue.length}`).join(', ') : 'NO: ' + doc.sources.pcs.error}; ${why}`);
}
if (process.argv[1] && fileURLToPath(import.meta.url) === process.argv[1]) {
main().catch(e => { log('failed:', e.message || e); process.exit(1); });
}