igneum/relay/api/relay.mjs
igneum-labs ec0ede62b9 Relay: run and task posts need the console token (round 4, X23); prove host saves proofs buffered (ledger P20, second gap)
The intake key sits in every miner package, so the relay now lets it report only (drop text and files, ack, done,
register, upload). Posting a run or task, or renaming and re-roling a machine, needs the console token.
The prove host wrote proofs through SP1's unbuffered save: on WSL2 under /mnt/c the 18 MB core proof of a shard
took longer to save than to prove. Proofs now go through a 4 MB buffer with a timed 'saved' line, and
prove-shard.sh keeps results on the Linux side and copies them per stage. Ledger P20 updated.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-04 18:14:10 +00:00

201 lines
14 KiB
JavaScript

// Igneum relay: one function, dispatched on ?fn=. Reached through the rewrite /r/<token>/api/<fn>.
// Auth: the token in the path (or x-relay-token), or the intake key in x-igneum-key. Nothing else.
// GET feed ?since=<id> | ?before=<id> | ?machine=X | ?limit=N items newest first + machines + unread counts
// GET item ?id= one item with its full body
// GET file ?id= 302 to the Blob URL (or the inline bytes)
// GET inbox ?machine=PC1&kind=run|task|all&ack=1 unread tasks for a machine; ack marks them read
// GET machines every machine with role, hostname, last_seen
// POST drop JSON {from,to,kind,title,body,file_name,file_url,size,task_id,flags} or raw octet-stream (x-file-name, x-from)
// POST task JSON {to,title,body,file_name,file_url,size,kind:'task'|'run',flags:{elevated,reboot_continue},from}
// POST upload JSON {name,size} -> {token, put_url, api_version} client token for a direct PUT to Vercel Blob (50 MB)
// POST ack JSON {ids:[...]} mark read
// POST done JSON {id, exit_code} mark done (runner finished)
// POST register JSON {hostname, info, role?} machine checks in; returns its name and role
// POST name JSON {hostname, name} name an unknown machine
// POST role JSON {name, role} set a machine's role
// POST delete JSON {id}
import { neon, authed, readJson, readRaw, str, safeName, storeBuffer, clientUploadToken, ITEM_COLS, rowOut, iso, touch, KINDS, ROLES, MAX_INLINE, MAX_BODY } from '../lib/relay.mjs';
const machineOut = m => ({ ...m, last_seen: iso(m.last_seen) });
const json = (res, status, obj) => { res.status(status).setHeader('Content-Type', 'application/json; charset=utf-8'); res.end(JSON.stringify(obj)); };
async function insertItem(sql, o) {
let kind = KINDS.has(o.kind) ? o.kind : (o.file_url || o.file_b64 ? 'file' : 'text');
if (kind === 'text' && o.file_url) kind = 'file';
const flags = o.flags && typeof o.flags === 'object' ? o.flags : {};
const rows = await sql(
`INSERT INTO relay_items (from_machine, to_machine, kind, title, body, file_name, file_url, file_b64, size, flags, task_id)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10::jsonb,$11) RETURNING id, ts`,
[str(o.from, 80) || 'unknown', str(o.to, 80) || 'all', kind, str(o.title, 300), str(o.body, MAX_BODY),
o.file_name ? safeName(o.file_name) : null, o.file_url ? str(o.file_url, 1000) : null, o.file_b64 || null,
Number(o.size) || 0, JSON.stringify(flags), o.task_id ? Number(o.task_id) : null]);
await touch(sql, o.from);
return { id: Number(rows[0].id), ts: rows[0].ts, kind };
}
export default async function handler(req, res) {
res.setHeader('Cache-Control', 'no-store');
const via = authed(req);
if (!via) return json(res, 401, { ok: false, error: 'no token' });
const fn = String(req.query.fn || '');
const q = req.query;
let sql;
try { sql = neon(); } catch (e) { return json(res, 500, { ok: false, error: e.message }); }
try {
if (req.method === 'GET') {
if (fn === 'feed') {
const limit = Math.min(500, Math.max(1, Number(q.limit) || 200));
const where = []; const params = [];
if (q.since) { params.push(Number(q.since)); where.push(`id > $${params.length}`); }
if (q.before) { params.push(Number(q.before)); where.push(`id < $${params.length}`); }
if (q.machine) { params.push(str(q.machine, 80)); where.push(`(from_machine = $${params.length} OR to_machine = $${params.length})`); }
const items = await sql(`SELECT ${ITEM_COLS} FROM relay_items ${where.length ? 'WHERE ' + where.join(' AND ') : ''} ORDER BY id DESC LIMIT ${limit}`, params);
const machines = await sql(`SELECT m.name, m.hostname, m.role, m.named, m.info, m.last_seen,
(SELECT count(*) FROM relay_items i WHERE NOT i.read AND i.kind IN ('task','run') AND (i.to_machine = m.name OR (i.to_machine = 'all' AND i.kind = 'task')))::int AS unread
FROM relay_machines m ORDER BY m.last_seen DESC NULLS LAST, m.name`);
return json(res, 200, { ok: true, items: items.map(rowOut), machines: machines.map(machineOut), now: new Date().toISOString() });
}
if (fn === 'machines') {
const machines = await sql(`SELECT name, hostname, role, named, info, last_seen FROM relay_machines ORDER BY name`);
return json(res, 200, { ok: true, machines: machines.map(machineOut) });
}
if (fn === 'item') {
const rows = await sql(`SELECT ${ITEM_COLS} FROM relay_items WHERE id = $1`, [Number(q.id)]);
if (!rows.length) return json(res, 404, { ok: false, error: 'no such item' });
return json(res, 200, { ok: true, item: rowOut(rows[0]) });
}
if (fn === 'file') {
const rows = await sql(`SELECT file_name, file_url, file_b64 FROM relay_items WHERE id = $1`, [Number(q.id)]);
if (!rows.length || (!rows[0].file_url && !rows[0].file_b64)) return json(res, 404, { ok: false, error: 'no file' });
if (rows[0].file_url) { res.statusCode = 302; res.setHeader('Location', rows[0].file_url + (q.download ? '?download=1' : '')); return res.end(); }
const buf = Buffer.from(rows[0].file_b64, 'base64');
res.setHeader('Content-Type', 'application/octet-stream');
res.setHeader('Content-Disposition', `attachment; filename="${safeName(rows[0].file_name)}"`);
return res.status(200).end(buf);
}
if (fn === 'inbox') {
const machine = str(q.machine, 80);
if (!machine) return json(res, 400, { ok: false, error: 'machine required' });
const kind = q.kind === 'run' ? ['run'] : q.kind === 'task' ? ['task'] : ['task', 'run'];
// run items only ever go to one named machine; task items may be addressed to all
const rows = await sql(`SELECT ${ITEM_COLS} FROM relay_items
WHERE NOT read AND NOT done AND kind = ANY($2) AND (to_machine = $1 OR (to_machine = 'all' AND kind = 'task'))
ORDER BY id ASC LIMIT 50`, [machine, kind]);
if (q.ack && rows.length) await sql(`UPDATE relay_items SET read = true, read_at = now() WHERE id = ANY($1)`, [rows.map(r => Number(r.id))]);
await touch(sql, machine);
return json(res, 200, { ok: true, machine, items: rows.map(rowOut) });
}
return json(res, 404, { ok: false, error: `unknown fn ${fn}` });
}
if (req.method !== 'POST') return json(res, 405, { ok: false, error: 'method' });
const ct = String(req.headers['content-type'] || '');
if (fn === 'drop' && !ct.includes('json')) {
// raw bytes through the function: curl --data-binary @file -H 'Content-Type: application/octet-stream' -H 'x-file-name: a.zip'
const buf = await readRaw(req);
if (!buf.length) return json(res, 400, { ok: false, error: 'empty body' });
if (buf.length > MAX_INLINE) return json(res, 413, { ok: false, error: `raw upload over ${MAX_INLINE} bytes; use fn=upload for a Blob client token` });
const name = safeName(req.headers['x-file-name'] || 'file.bin');
const stored = await storeBuffer(name, buf, ct || 'application/octet-stream');
const r = await insertItem(sql, { from: req.headers['x-from'] || q.from, to: req.headers['x-to'] || q.to, kind: 'file',
title: str(req.headers['x-title'] || q.title || name, 300), body: '', file_name: name, file_url: stored.url, size: buf.length,
task_id: req.headers['x-task-id'] || q.task_id });
return json(res, 200, { ok: true, ...r, file_url: stored.url });
}
let body;
try { body = await readJson(req); } catch { return json(res, 400, { ok: false, error: 'bad json' }); }
// Review round 4, X23: the intake key is in every miner package, so it may only report (drop text and files, ack,
// done, register, upload). Anything a machine would EXECUTE, and anything that renames or re-roles a machine,
// needs the console token.
if ((fn === 'task' || (fn === 'drop' && (body.kind === 'run' || body.kind === 'task')) || fn === 'name' || fn === 'role' || fn === 'delete') && via !== 'token')
return json(res, 403, { ok: false, error: 'this call needs the console token' });
if (fn === 'drop' || fn === 'task') {
const o = { ...body };
if (fn === 'task') { o.kind = o.kind === 'run' ? 'run' : 'task'; if (!o.to) return json(res, 400, { ok: false, error: 'to required' }); }
if (o.kind === 'run' && (!o.to || o.to === 'all')) return json(res, 400, { ok: false, error: 'a run task needs one named machine' });
if (o.file_b64 && !o.file_url) {
const buf = Buffer.from(String(o.file_b64), 'base64');
if (buf.length > MAX_INLINE) return json(res, 413, { ok: false, error: 'inline file over 4 MB; use fn=upload' });
const stored = await storeBuffer(o.file_name || 'file.bin', buf, o.content_type);
o.file_url = stored.url; o.size = buf.length; delete o.file_b64;
}
if (!o.body && !o.file_url && !o.title) return json(res, 400, { ok: false, error: 'nothing to send' });
if (o.file_url && !/^https:\/\/[a-z0-9.-]+\.public\.blob\.vercel-storage\.com\//i.test(o.file_url)) return json(res, 400, { ok: false, error: 'file_url must be a Vercel Blob URL from fn=upload' });
const r = await insertItem(sql, o);
return json(res, 200, { ok: true, ...r });
}
if (fn === 'upload') {
const name = safeName(body.name || 'file.bin');
const t = await clientUploadToken(name, Number(body.size) || 0);
return json(res, 200, { ok: true, name, ...t });
}
if (fn === 'ack') {
const ids = (Array.isArray(body.ids) ? body.ids : [body.id]).map(Number).filter(Boolean);
if (!ids.length) return json(res, 400, { ok: false, error: 'ids required' });
await sql(`UPDATE relay_items SET read = true, read_at = now() WHERE id = ANY($1)`, [ids]);
return json(res, 200, { ok: true, ids });
}
if (fn === 'done') {
const id = Number(body.id); if (!id) return json(res, 400, { ok: false, error: 'id required' });
const extra = body.exit_code === undefined ? {} : { exit_code: Number(body.exit_code) };
await sql(`UPDATE relay_items SET done = true, done_at = now(), read = true, read_at = COALESCE(read_at, now()), flags = flags || $2::jsonb WHERE id = $1`, [id, JSON.stringify(extra)]);
return json(res, 200, { ok: true, id });
}
if (fn === 'delete') {
const id = Number(body.id); if (!id) return json(res, 400, { ok: false, error: 'id required' });
await sql(`DELETE FROM relay_items WHERE id = $1`, [id]);
return json(res, 200, { ok: true, id });
}
if (fn === 'register') {
const hostname = str(body.hostname, 120).trim();
if (!hostname) return json(res, 400, { ok: false, error: 'hostname required' });
const info = body.info && typeof body.info === 'object' ? body.info : {};
let rows = await sql(`SELECT name, role, named FROM relay_machines WHERE hostname = $1`, [hostname]);
if (!rows.length) {
// first contact from this hostname: it shows up under its own hostname until the Mac names it
rows = await sql(`INSERT INTO relay_machines (name, hostname, role, named, info, last_seen) VALUES ($1, $1, $2, false, $3::jsonb, now())
ON CONFLICT (name) DO UPDATE SET hostname = EXCLUDED.hostname, last_seen = now(), info = EXCLUDED.info RETURNING name, role, named`,
[hostname, ROLES.has(body.role) ? body.role : '', JSON.stringify(info)]);
} else {
await sql(`UPDATE relay_machines SET last_seen = now(), info = $2::jsonb WHERE hostname = $1`, [hostname, JSON.stringify(info)]);
}
return json(res, 200, { ok: true, name: rows[0].name, role: rows[0].role || '', named: !!rows[0].named, hostname });
}
if (fn === 'name') {
const hostname = str(body.hostname, 120).trim(); const name = str(body.name, 80).trim();
if (!hostname || !name) return json(res, 400, { ok: false, error: 'hostname and name required' });
const target = await sql(`SELECT name, hostname FROM relay_machines WHERE name = $1`, [name]);
const old = await sql(`SELECT name FROM relay_machines WHERE hostname = $1`, [hostname]);
if (target.length && target[0].hostname && target[0].hostname !== hostname) return json(res, 409, { ok: false, error: `${name} is already ${target[0].hostname}` });
if (target.length) {
// a pre-seeded name (PC2 with no hostname yet): attach the hostname, drop the placeholder row, move its items
if (old.length && old[0].name !== name) {
await sql(`DELETE FROM relay_machines WHERE hostname = $1 AND name <> $2`, [hostname, name]);
await sql(`UPDATE relay_items SET from_machine = $2 WHERE from_machine = $1`, [old[0].name, name]);
await sql(`UPDATE relay_items SET to_machine = $2 WHERE to_machine = $1`, [old[0].name, name]);
}
await sql(`UPDATE relay_machines SET hostname = $1, named = true, last_seen = COALESCE(last_seen, now()) WHERE name = $2`, [hostname, name]);
} else if (old.length) {
await sql(`UPDATE relay_machines SET name = $2, named = true WHERE hostname = $1`, [hostname, name]);
await sql(`UPDATE relay_items SET from_machine = $2 WHERE from_machine = $1`, [old[0].name, name]);
await sql(`UPDATE relay_items SET to_machine = $2 WHERE to_machine = $1`, [old[0].name, name]);
} else {
await sql(`INSERT INTO relay_machines (name, hostname, role, named) VALUES ($2, $1, '', true)`, [hostname, name]);
}
return json(res, 200, { ok: true, hostname, name });
}
if (fn === 'role') {
const name = str(body.name, 80).trim(); const role = str(body.role, 20).trim();
if (!name || !ROLES.has(role)) return json(res, 400, { ok: false, error: `role must be one of ${[...ROLES].filter(Boolean).join(', ')}` });
await sql(`INSERT INTO relay_machines (name, role, named) VALUES ($1, $2, true) ON CONFLICT (name) DO UPDATE SET role = EXCLUDED.role`, [name, role]);
return json(res, 200, { ok: true, name, role });
}
return json(res, 404, { ok: false, error: `unknown fn ${fn}` });
} catch (e) {
return json(res, 500, { ok: false, error: String(e.message || e) });
}
}