igneum/tools/workers/push.mjs
igneum-labs 9a6f2606b8 Worker dashboard: one server section per box, per-box live feeds, installer and pusher take build-2 and build-3
The build-server agent has igneum-build-2 and -3 on order (suites and benches; proving and the fast-time nodes). The
installer takes a host argument (N reads ~/.config/igneum/build-server-N, or build@<ip>); the pusher drops the Mac and
PC facts on every box it has a host file for, pulls each box's file and publishes boxes[] (box stays the first for the
old shape); the page draws one server section and one crew card per box, reads every box's Caddy feed in parallel
(build, build-2, build-3.igneum.network) with the edge copy filling any box that does not answer, and merges every
box's builds into the lanes, the timeline and the analytics. Nothing changes for build-1 until the host files exist.

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

214 lines
15 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');
// every build box: ~/.config/igneum/build-server (build-1), build-server-2, build-server-3 ... one line build@<ip> each
const HOSTS = [['build-server', 'igneum-build-1'], ['build-server-2', 'igneum-build-2'], ['build-server-3', 'igneum-build-3']]
.map(([f, n]) => ({ file: f, name: n, host: (cfg(f) || '').split('\n')[0].trim() })).filter(h => h.host);
const HOST = HOSTS.length ? HOSTS[0].host : '';
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, the project lead\'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, host = HOST) { return spawnSync('ssh', [...SSH_OPTS, host, ...args], { encoding: 'utf8', input, timeout: 30000, maxBuffer: 16 * 1024 * 1024 }); }
function dropToBox(name, doc, host = HOST) {
const r = ssh([`cat > /srv/workers/sources/.${name}.tmp && mv /srv/workers/sources/.${name}.tmp /srv/workers/sources/${name}.json`], JSON.stringify(doc), host);
if (r.status !== 0) throw new Error(`scp ${name}: ${(r.stderr || '').trim().slice(0, 200) || 'ssh failed'}`);
}
function pullFromBox(host = HOST) {
const r = ssh(['cat /srv/workers/workers.json; echo; echo __SEP__; cat /srv/workers/headline.json 2>/dev/null || echo null'], undefined, host);
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 boxes: drop ours on each, pull each one's file; boxes[] carries every box, box stays the first for older readers
doc.boxes = [];
if (!HOSTS.length) doc.sources.box = { ok: false, error: 'no ~/.config/igneum/build-server' };
for (const h of HOSTS) {
try {
if (doc.mac) dropToBox('mac', doc.mac, h.host);
if (doc.pcs && !doc.pcs.error) dropToBox('pcs', doc.pcs, h.host);
const { workers, headline } = pullFromBox(h.host);
const age = (now - Date.parse(workers.generated_at)) / 1000;
const b = { ...workers.box, headline, feed: `https://${h.name === 'igneum-build-1' ? 'build' : h.name.replace('igneum-', '')}.igneum.network/workers.json`, source: { 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)` } : {}) } };
doc.boxes.push(b);
if (!doc.box) { doc.box = workers.box; doc.headline = headline; doc.sources.box = b.source; }
} catch (e) {
const src = { ok: false, error: `ssh to ${h.name} failed: ${String(e.message || e).slice(0, 160)}` };
doc.boxes.push({ name: h.name, source: src });
if (!doc.sources.box) doc.sources.box = src;
}
}
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(`boxes ${doc.boxes.map(b => `${b.name} ${b.source && b.source.ok ? 'ok' : 'NO'}`).join(', ')}; 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); });
}