Relay (X23, X27): the intake key is its own tier (upload and file drops only, RELAY_INTAKE_COMPAT=0 closes it);
a run task needs an Ed25519 signature by the Mac run key over {to, nonce, body sha256, flags} (RELAY_RUN_PUB,
401 without) and an HMAC tag with the target's machine secret that the agent verifies before anything runs;
results and registration are bound to the machine the secret proves (403 on a forged from).
X24: every client and Mac tool sends x-relay-token as a header to /api/relay?fn=; the path token stays for the
phone page only. X25: the agent arms the logon task only for a restart a task asked for and disarms on start
and exit. X26: 30-day retention with blob deletion, feed capped at 100, the dl base as RELAY_DL_BASE held by the
agent, never in a body. X28: GET inbox never acks (POST inbox does), RELAY-REBOOT on its own line and only with a
reboot flag, 120/min and 10 failed auths/min per IP, no username or folder on register, WSL sudo scoped to
apt-get and dpkg with SETENV, no password on a command line. X29: the intake key reaches curl through -K in
upload.sh and both upload-log.bat; tools/ci/curl-header-check.sh fails the class. G14: TZ=UTC in ship-app.mjs
and publish-jobs.sh; tools/ci/commit-tz-check.sh fails the class; history-rewrite.md names the .old-2026-10-05
files as the values in the history. The handler moved to relay/lib/handler.mjs with injected sql and blobs
(relay/lib/blob.mjs holds @vercel/blob) so relay/test/handler.test.mjs drives it without a database:
47 tests across 6 suites, all green.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
295 lines
21 KiB
JavaScript
295 lines
21 KiB
JavaScript
// 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=<fn> with x-relay-token or x-igneum-key, or through the
|
|
// rewrite /r/<token>/api/<fn> from the phone's page (the only caller that keeps the token in the path, X24).
|
|
//
|
|
// GET feed ?since=<id> | ?before=<id> | ?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 });
|
|
}
|
|
};
|
|
}
|