// Igneum relay: one function, dispatched on ?fn=. Reached through the rewrite /r//api/. // Auth: the token in the path (or x-relay-token), or the intake key in x-igneum-key. Nothing else. // GET feed ?since= | ?before= | ?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) }); } }