// The relay handler with its dependencies injected (api/relay.mjs wires Neon and Vercel Blob; relay/test/handler.test.mjs // wires fakes). Dispatched on ?fn=; reached as /api/relay?fn= with x-relay-token or x-igneum-key, or through the // rewrite /r//api/ from the phone's page (the only caller that keeps the token in the path, X24). // // GET feed ?since= | ?before= | ?machine=X | ?limit=N (cap 100) 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 unread tasks for a machine; a GET never marks anything (X28) // POST inbox JSON {machine, kind, ack:true} the same list, marked read // GET machines every machine with role, hostname, last_seen, bound // 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,kind:'task'|'run',flags:{elevated,reboot_continue,reboot,nonce,sig,mac},from} // run: token AND flags.sig (Ed25519 by the relay run key, verified here with RELAY_RUN_PUB) AND flags.mac // (the machine's HMAC tag, verified by the agent) AND a fresh flags.nonce; anything less is 401 (X23) // 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?} + x-machine-secret machine checks in; the secret names it (X27) // POST name JSON {hostname, name} name an unknown machine (token) // POST role JSON {name, role} set a machine's role (token) // POST secret JSON {name, secret_hash} bind a machine to sha256(its secret) (token) // POST delete JSON {id} the row and its blob (token) // Retention: rows older than 30 days go, blobs with them, checked on a feed read at most every 10 minutes per instance. // Rate limit: 120 calls a minute per IP, 10 failed authentications a minute per IP (X28). import { readJson, readRaw, str, safeName, ITEM_COLS, BLOB_URL_RE, rowOut, iso, touch, KINDS, ROLES, MAX_INLINE, MAX_BODY } from './relay.mjs'; import { POST_ALLOWED, DROP_KINDS, mayRead, checkRun, feedLimit, retentionCutoff, machineForSecret, secretHash, isSecret, RATE_PER_MIN, AUTH_FAIL_PER_MIN } from './guard.mjs'; import { RateLimit, ipOf } from './wake.mjs'; export const EXPIRE_EVERY_MS = 10 * 60_000; const machineOut = m => ({ name: m.name, hostname: m.hostname, role: m.role, named: m.named, info: m.info, last_seen: iso(m.last_seen), bound: !!m.secret_hash, ...(m.unread !== undefined ? { unread: m.unread } : {}) }); 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 }; } /** The secret_hash column, added once per instance (ADD COLUMN IF NOT EXISTS is idempotent and takes no long lock). */ async function ensureSecretColumn(sql, state) { if (state.secretColumn) return; await sql(`ALTER TABLE relay_machines ADD COLUMN IF NOT EXISTS secret_hash text`); state.secretColumn = true; } /** Rows older than the retention and their blobs (X26). */ export async function expire(sql, blob, now = Date.now()) { const rows = await sql(`DELETE FROM relay_items WHERE ts < $1 RETURNING id, file_url`, [retentionCutoff(now)]); const blobs = await blob.deleteBlobs(rows.map(r => r.file_url)); return { rows: rows.length, blobs }; } /** * makeHandler({ sql, blob: {storeBuffer, clientUploadToken, deleteBlobs}, authed, json, env, now, limiter, authLimiter, state }) * json(res, status, obj) writes a reply; state is per instance (the expire clock and the column flag). */ export function makeHandler({ sql, blob, authed, json, env = process.env, now = Date.now, limiter, authLimiter, state = {} }) { const limit = limiter || new RateLimit({ perMinute: RATE_PER_MIN, now }); const authLimit = authLimiter || new RateLimit({ perMinute: AUTH_FAIL_PER_MIN, now }); const bindMachine = async (secret) => { // the machine a presented x-machine-secret names, or an error text await ensureSecretColumn(sql, state); const rows = await sql(`SELECT name, secret_hash FROM relay_machines WHERE secret_hash IS NOT NULL`); return machineForSecret(secret, rows); }; const hasSecret = async (name) => { await ensureSecretColumn(sql, state); const rows = await sql(`SELECT secret_hash FROM relay_machines WHERE name = $1`, [name]); return !!(rows.length && rows[0].secret_hash); }; return async function handler(req, res) { res.setHeader('Cache-Control', 'no-store'); const ip = ipOf(req); const via = authed(req); if (!via) { const wait = authLimit.take(ip); if (wait !== null) { res.setHeader('Retry-After', String(wait)); return json(res, 429, { ok: false, error: `too many failed authentications; retry in ${wait} s` }); } return json(res, 401, { ok: false, error: 'no token' }); } const wait = limit.take(ip); if (wait !== null) { res.setHeader('Retry-After', String(wait)); return json(res, 429, { ok: false, error: `rate limit; retry in ${wait} s` }); } const fn = String(req.query.fn || ''); const q = req.query; try { if (req.method === 'GET') { if (!mayRead(via)) return json(res, 403, { ok: false, error: 'the intake key may only upload and drop files' }); if (fn === 'feed') { if (now() - (state.lastExpire || 0) > EXPIRE_EVERY_MS) { state.lastExpire = now(); try { state.expired = await expire(sql, blob, now()); } catch (e) { state.expireError = String(e.message || e); } } const lim = feedLimit(q); 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 ${lim}`, params); await ensureSecretColumn(sql, state); const machines = await sql(`SELECT m.name, m.hostname, m.role, m.named, m.info, m.last_seen, m.secret_hash, (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), limit: lim, now: new Date(now()).toISOString() }); } if (fn === 'machines') { await ensureSecretColumn(sql, state); const machines = await sql(`SELECT name, hostname, role, named, info, last_seen, secret_hash 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') return inbox(q, false); return json(res, 404, { ok: false, error: `unknown fn ${fn}` }); } if (req.method !== 'POST') return json(res, 405, { ok: false, error: 'method' }); if (!POST_ALLOWED[via].has(fn)) return json(res, via === 'intake' ? 403 : (POST_ALLOWED.token.has(fn) ? 403 : 404), { ok: false, error: POST_ALLOWED.token.has(fn) ? `${fn} needs the console token` : `unknown fn ${fn}` }); 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 blob.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' }); } const presented = req.headers['x-machine-secret']; if (fn === 'inbox') { if (!POST_ALLOWED[via].has('inbox')) return json(res, 403, { ok: false, error: 'inbox needs the token or the relay key' }); return inbox(body, !!body.ack); } 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' }); } const kind = KINDS.has(o.kind) ? o.kind : (o.file_url || o.file_b64 ? 'file' : 'text'); const allowed = DROP_KINDS[via]; if (allowed && !allowed.has(kind)) return json(res, 403, { ok: false, error: `a ${kind} item needs the console token` }); if (kind === 'run') { const why = checkRun(o, env.RELAY_RUN_PUB); if (why) return json(res, why.startsWith('a run task needs one') ? 400 : 401, { ok: false, error: why }); const dup = await sql(`SELECT id FROM relay_items WHERE kind = 'run' AND flags->>'nonce' = $1`, [o.flags.nonce]); if (dup.length) return json(res, 409, { ok: false, error: `nonce already used by #${dup[0].id}` }); } if (kind === 'result') { // X27: a result names the machine its secret proves; a forged `from` is refused if (presented !== undefined) { const m = await bindMachine(presented); if (m.error) return json(res, 403, { ok: false, error: m.error }); if (o.from && o.from !== m.name) return json(res, 403, { ok: false, error: `from ${o.from} does not match the machine secret (${m.name})` }); o.from = m.name; } else if (o.from && await hasSecret(o.from)) { return json(res, 403, { ok: false, error: `results from ${o.from} need its machine secret (x-machine-secret)` }); } else { o.flags = { ...(o.flags && typeof o.flags === 'object' ? o.flags : {}), unbound: true }; } } 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 blob.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 && !BLOB_URL_RE.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 blob.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) }; if (presented !== undefined) { const m = await bindMachine(presented); if (m.error) return json(res, 403, { ok: false, error: m.error }); extra.done_by = m.name; } 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' }); const rows = await sql(`DELETE FROM relay_items WHERE id = $1 RETURNING file_url`, [id]); const blobs = await blob.deleteBlobs(rows.map(r => r.file_url)); return json(res, 200, { ok: true, id, blobs }); } 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 } : {}; delete info.user; delete info.dir; // X28: no username and no secret folder in the registration if (presented !== undefined) { // the secret names the machine; the hostname is recorded against that row (both PCs report the same hostname) const m = await bindMachine(presented); if (m.error) return json(res, 403, { ok: false, error: m.error }); const rows = await sql(`UPDATE relay_machines SET hostname = $2, last_seen = now(), info = $3::jsonb WHERE name = $1 RETURNING name, role, named`, [m.name, hostname, JSON.stringify(info)]); return json(res, 200, { ok: true, name: rows[0].name, role: rows[0].role || '', named: !!rows[0].named, hostname, bound: true }); } await ensureSecretColumn(sql, state); let rows = await sql(`SELECT name, role, named, secret_hash FROM relay_machines WHERE hostname = $1`, [hostname]); if (rows.some(r => r.secret_hash)) return json(res, 403, { ok: false, error: `${hostname} is bound to a machine secret; present x-machine-secret` }); 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, bound: false }); } if (fn === 'secret') { const name = str(body.name, 80).trim(); const hash = str(body.secret_hash, 64).trim(); if (!name || !isSecret(hash)) return json(res, 400, { ok: false, error: 'name and secret_hash (64 hex, sha256 of the secret) required' }); await ensureSecretColumn(sql, state); await sql(`INSERT INTO relay_machines (name, role, named, secret_hash) VALUES ($1, '', true, $2) ON CONFLICT (name) DO UPDATE SET secret_hash = EXCLUDED.secret_hash, named = true`, [name, hash]); return json(res, 200, { ok: true, name, bound: true }); } 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) }); } async function inbox(src, ack) { const machine = str(src.machine, 80); if (!machine) return json(res, 400, { ok: false, error: 'machine required' }); const kind = src.kind === 'run' ? ['run'] : src.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 (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), acked: ack ? rows.length : 0 }); } }; }