#!/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//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@ 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 `); 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-- 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); }); }