relay: /wake long-poll for the apps' remote jobs (public GET held 45 s, authenticated POST of the stamp)

GET /wake?since=<stamp> is public (the apps hold no token) and rate limited (30 a minute per IP). It holds up to
45 s, re-reading the stamp every 2 s, and answers {stamp, at, added, changed, held_ms} the moment the stored stamp
differs from since, else the unchanged stamp at the deadline. POST /r/<token>/wake {stamp, added} (the relay's
auth, also x-relay-token or x-igneum-key on /wake) records a stamp; one row per stamp in relay_wake, created by the
first POST. maxDuration 60 s for api/wake.mjs in vercel.json. api/relay.mjs is untouched.

The handler lives in lib/wake.mjs with its dependencies injected; relay/test/wake.test.mjs drives it with a fake
database, a fake clock and a fake sleep (the hold, the change, the deadline, the rate limit, the hold cap, auth, a
database error). CI's site job runs it with the other relay tests.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
igneum-labs 2026-10-05 08:22:02 +00:00
parent b3ba9ac1ee
commit 0550f58857
6 changed files with 298 additions and 2 deletions

View file

@ -61,5 +61,5 @@ jobs:
run: node tools/ci/link-check.mjs
- name: identity grep of the public export list
run: bash tools/ci/identity-check.sh
- name: relay unit tests (parsers, secret compare)
run: node --test relay/test/parse.test.mjs relay/test/auth.test.mjs
- name: relay unit tests (parsers, secret compare, the wake endpoint)
run: node --test relay/test/parse.test.mjs relay/test/auth.test.mjs relay/test/wake.test.mjs

View file

@ -53,6 +53,8 @@ Kinds: `text` (a note), `file`, `task` (for a person or a Claude session on a PC
| `POST register {hostname,info}` | a machine checks in; returns its name, role and whether it is named |
| `POST name {hostname,name}` `POST role {name,role}` | naming and roles, from the Mac |
Wake (`api/wake.mjs`, 0.3.6, 5 October 2026): `GET /wake?since=<stamp>` is public (the apps hold no token) and rate limited, 30 a minute per IP. It holds up to 45 s and answers `{stamp, at, added, changed, held_ms}` the moment the stored stamp differs from `since`, else the unchanged stamp at the deadline; without `since` it answers at once. `POST /r/<token>/wake {stamp, added}` (or `POST /wake` with `x-relay-token` or `x-igneum-key`) records the stamp; `packaging/ota/publish-jobs.sh` sends it after every verified deploy, with the ids it added. One row per stamp in Neon table `relay_wake` (created by the first POST); `tools/jobs.mjs status` reads the rows for the woken latency. The function's `maxDuration` is 60 s (`vercel.json`). Tests: `relay/test/wake.test.mjs` drives the handler with a fake database and clock.
## Mac
`node tools/relay.mjs` (feed), `read <id>`, `drop "<text>"|<file>`, `task PC2 "title" [file]`, `run PC2 "title" script.ps1 [--elevated] [--reboot-continue]`, `watch`, `inbox PC1`, `machines`, `role PC2 prover`, `name DESKTOP-XYZ PC2`, `ack|done|rm <id>`, `url`. Playbooks live in `relay/playbooks/`; `run` fills `__DL_BASE__` in from `~/.config/igneum/dl-token`.

18
relay/api/wake.mjs Normal file
View file

@ -0,0 +1,18 @@
// Igneum wake: the apps' long-poll for a new jobs file (app/igneum-app/src/jobrun.rs, 0.3.6). Reached at /wake
// (public: the apps hold no token) and at /r/<token>/wake (the same function; the token matters only to POST).
// The contract and the whole handler live in ../lib/wake.mjs so the test can drive it without a database.
// GET wake?since=<stamp>[&hold=45] up to 45 s; {stamp, at, added, changed, held_ms}; 30 a minute per IP
// POST wake {stamp, added?} the relay's auth (token in the path or x-relay-token, or x-igneum-key)
// maxDuration 60 s for this file is set in vercel.json.
import { neon, authed, readJson } from '../lib/relay.mjs';
import { makeHandler, RateLimit } from '../lib/wake.mjs';
const json = (res, status, obj) => { res.status(status).setHeader('Content-Type', 'application/json; charset=utf-8'); res.end(JSON.stringify(obj)); };
const limiter = new RateLimit();
const holds = { n: 0 };
export default async function handler(req, res) {
let sql;
try { sql = neon(); } catch (e) { res.setHeader('Cache-Control', 'no-store'); return json(res, 500, { ok: false, error: e.message }); }
return makeHandler({ sql, authed, readJson, json, limiter, holds })(req, res);
}

111
relay/lib/wake.mjs Normal file
View file

@ -0,0 +1,111 @@
// The wake signal for the apps' remote jobs (api/wake.mjs, Igneum Miner 0.3.6). The apps hold no token, so the
// GET is public and rate limited; only the POST that sets the stamp is authenticated. Dependency-free, so
// `node --test relay/test/wake.test.mjs` drives the whole handler with a fake database and a fake clock.
//
// GET wake?since=<stamp>[&hold=45] holds up to 45 s; answers as soon as the stored stamp differs from `since`,
// else the unchanged stamp at the deadline; without `since` the current stamp
// at once. Reply: {ok, stamp, at, added, changed, held_ms}.
// POST wake {stamp, added?} records the stamp (publish-jobs.sh, after a verified deploy). Reply:
// {ok, stamp, at, added}.
// Storage: table relay_wake, one row per stamp, the newest row is the stamp; the first POST creates the table.
export const HOLD_MS_MAX = 45_000; // the function's maxDuration is 60 s (relay/vercel.json)
export const STEP_MS = 2_000; // how often a held request re-reads the stamp
export const STAMP_MAX = 120;
export const RATE_PER_MIN = 30; // per IP; a healthy app makes about 2 requests a minute
export const HOLDS_MAX = 64; // held requests per instance; above it the answer is immediate
const EMPTY = { stamp: '', at: null, added: [] };
export const validStamp = s => typeof s === 'string' && s.length > 0 && s.length <= STAMP_MAX && /^[\w.:-]+$/.test(s);
const isoOf = v => { if (!v) return null; const d = new Date(String(v).replace(' ', 'T').replace(/([+-]\d\d)$/, '$1:00')); return isNaN(d) ? String(v) : d.toISOString(); };
const clip = (v, max) => (v === undefined || v === null ? '' : String(v)).slice(0, max);
/** Holds until read() gives a stamp other than `since`, or the deadline. Pure apart from read, now and sleep. */
export async function waitForChange({ read, since, holdMs = HOLD_MS_MAX, stepMs = STEP_MS, now = Date.now, sleep = ms => new Promise(r => setTimeout(r, ms)) }) {
const hold = Math.max(0, Math.min(HOLD_MS_MAX, Number(holdMs) || 0));
const t0 = now();
const deadline = t0 + hold;
for (;;) {
const cur = (await read()) || EMPTY;
const stamp = cur.stamp ? String(cur.stamp) : '';
const changed = !!since && stamp !== since;
if (!since || changed) return { ...cur, stamp, changed, held_ms: now() - t0 };
const left = deadline - now();
if (left <= 0) return { ...cur, stamp, changed: false, held_ms: now() - t0 };
await sleep(Math.min(stepMs, left));
}
}
/** Fixed one-minute windows per key; take() answers null when allowed, else the seconds to wait. */
export class RateLimit {
constructor({ perMinute = RATE_PER_MIN, now = Date.now } = {}) { this.perMinute = perMinute; this.now = now; this.hits = new Map(); }
take(key) {
const t = this.now();
const k = String(key || '?');
let h = this.hits.get(k);
if (!h || t - h.start >= 60_000) { h = { start: t, n: 0 }; this.hits.set(k, h); }
if (h.n >= this.perMinute) return Math.max(1, Math.ceil((h.start + 60_000 - t) / 1000));
h.n += 1;
if (this.hits.size > 10_000) for (const [key2, v] of this.hits) if (t - v.start >= 60_000) this.hits.delete(key2);
return null;
}
}
export async function ensureTable(sql) {
await sql(`CREATE TABLE IF NOT EXISTS relay_wake (id bigserial PRIMARY KEY, stamp text NOT NULL UNIQUE, at timestamptz NOT NULL DEFAULT now(), meta jsonb NOT NULL DEFAULT '{}'::jsonb)`);
}
/** The newest stamp, or the empty one before the first POST (the table does not exist yet). */
export async function latest(sql) {
try {
const rows = await sql(`SELECT stamp, at, meta FROM relay_wake ORDER BY id DESC LIMIT 1`);
if (!rows.length) return EMPTY;
const m = rows[0].meta && typeof rows[0].meta === 'object' ? rows[0].meta : {};
return { stamp: String(rows[0].stamp), at: isoOf(rows[0].at), added: Array.isArray(m.added) ? m.added : [] };
} catch (e) {
if (/relay_wake/.test(String(e.message)) && /does not exist/.test(String(e.message))) return EMPTY;
throw e;
}
}
export const ipOf = req => String(req.headers['x-forwarded-for'] || (req.socket && req.socket.remoteAddress) || '').split(',')[0].trim();
/**
* The handler with its dependencies injected: sql (Neon), authed (the relay's check), readJson, json (the reply
* writer), limiter, holds (a counter shared by the instance), now and sleep (the test's fake clock).
*/
export function makeHandler({ sql, authed, readJson, json, limiter = new RateLimit(), holds = { n: 0 }, now = Date.now, sleep }) {
return async function handler(req, res) {
res.setHeader('Cache-Control', 'no-store');
const q = req.query || {};
try {
if (req.method === 'GET') {
const wait = limiter.take(ipOf(req));
if (wait !== null) { res.setHeader('Retry-After', String(wait)); return json(res, 429, { ok: false, error: 'rate limited', retry_after: wait }); }
const since = clip(q.since, STAMP_MAX * 2);
const asked = q.hold === undefined ? HOLD_MS_MAX : (Number(q.hold) || 0) * 1000;
const holdMs = holds.n >= HOLDS_MAX ? 0 : Math.min(HOLD_MS_MAX, asked);
holds.n += 1;
try {
const r = await waitForChange({ read: () => latest(sql), since, holdMs, now, sleep });
return json(res, 200, { ok: true, ...r });
} finally { holds.n -= 1; }
}
if (req.method !== 'POST') return json(res, 405, { ok: false, error: 'method' });
if (!authed(req)) return json(res, 401, { ok: false, error: 'no token' });
let body;
try { body = await readJson(req); } catch { return json(res, 400, { ok: false, error: 'bad json' }); }
const stamp = clip(body.stamp !== undefined ? body.stamp : q.stamp, STAMP_MAX * 2).trim();
if (!validStamp(stamp)) return json(res, 400, { ok: false, error: `stamp required: 1 to ${STAMP_MAX} of A-Z a-z 0-9 _ . : -` });
const added = (Array.isArray(body.added) ? body.added : []).map(x => clip(x, 80)).filter(Boolean).slice(0, 50);
await ensureTable(sql);
const rows = await sql(`INSERT INTO relay_wake (stamp, meta) VALUES ($1, $2::jsonb)
ON CONFLICT (stamp) DO UPDATE SET at = now(), meta = EXCLUDED.meta RETURNING stamp, at`, [stamp, JSON.stringify({ added })]);
return json(res, 200, { ok: true, stamp: String(rows[0].stamp), at: isoOf(rows[0].at), added });
} catch (e) {
return json(res, 500, { ok: false, error: String(e.message || e) });
}
};
}

152
relay/test/wake.test.mjs Normal file
View file

@ -0,0 +1,152 @@
// node --test relay/test/wake.test.mjs (no dependencies, no database: a fake sql, a fake clock, a fake sleep)
import { test } from 'node:test';
import assert from 'node:assert/strict';
import { waitForChange, RateLimit, validStamp, makeHandler, HOLD_MS_MAX, HOLDS_MAX } from '../lib/wake.mjs';
// a clock that only moves when something sleeps on it
function clock(start = 1_000_000) {
let t = start;
return { now: () => t, sleep: async ms => { t += ms; }, set: v => { t = v; } };
}
// the smallest Neon stand-in: one table relay_wake, newest row first, no table until the first CREATE
function fakeDb() {
const rows = [];
let created = false;
let seq = 0;
const sql = async (query, params = []) => {
if (/CREATE TABLE/.test(query)) { created = true; return []; }
if (!created) throw new Error('relation "relay_wake" does not exist');
if (/INSERT INTO relay_wake/.test(query)) {
const [stamp, meta] = params;
const hit = rows.find(r => r.stamp === stamp);
if (hit) { hit.at = '2026-10-05 12:00:30+00'; hit.meta = JSON.parse(meta); return [{ stamp, at: hit.at }]; }
const row = { id: ++seq, stamp, at: '2026-10-05 12:00:00+00', meta: JSON.parse(meta) };
rows.push(row);
return [{ stamp, at: row.at }];
}
if (/SELECT stamp, at, meta FROM relay_wake/.test(query)) { const r = rows[rows.length - 1]; return r ? [r] : []; }
throw new Error('unexpected query ' + query);
};
return { sql, rows, set: (stamp, added = []) => { rows.push({ id: ++seq, stamp, at: '2026-10-05 12:00:00+00', meta: { added } }); created = true; } };
}
function res() {
const r = { status: null, headers: {}, body: null };
r.setHeader = (k, v) => { r.headers[k] = v; };
r.status = s => { r.code = s; return r; };
r.end = s => { r.body = s; };
return r;
}
const json = (r, status, obj) => { r.status(status); r.end(JSON.stringify(obj)); };
const reply = r => ({ code: r.code, ...JSON.parse(r.body) });
const req = (method, query = {}, { headers = {}, body } = {}) => ({ method, query, headers: { 'x-forwarded-for': '203.0.113.9', ...headers }, body });
test('validStamp: the publish stamp shape, nothing else', () => {
assert.equal(validStamp('2026-10-05T11:02:17Z.5e7b56f5'), true);
assert.equal(validStamp(''), false);
assert.equal(validStamp('a b'), false);
assert.equal(validStamp('x'.repeat(121)), false);
assert.equal(validStamp(42), false);
assert.equal(validStamp('has/slash'), false);
});
test('waitForChange: no since answers at once; a different stamp answers at once with changed', async () => {
const c = clock();
const r = await waitForChange({ read: async () => ({ stamp: 'A', at: null, added: [] }), since: '', now: c.now, sleep: c.sleep });
assert.deepEqual([r.stamp, r.changed, r.held_ms], ['A', false, 0]);
const r2 = await waitForChange({ read: async () => ({ stamp: 'B', at: null, added: [] }), since: 'A', now: c.now, sleep: c.sleep });
assert.deepEqual([r2.stamp, r2.changed, r2.held_ms], ['B', true, 0]);
});
test('waitForChange: holds in 2 s steps, answers the moment the stamp moves, else the unchanged stamp at the deadline', async () => {
const c = clock();
let cur = 'A';
let reads = 0;
const read = async () => { reads++; if (c.now() >= 1_000_000 + 7_000) cur = 'B'; return { stamp: cur, at: null, added: [] }; };
const r = await waitForChange({ read, since: 'A', now: c.now, sleep: c.sleep });
assert.deepEqual([r.stamp, r.changed, r.held_ms], ['B', true, 8_000]);
assert.equal(reads, 5);
const c2 = clock();
let reads2 = 0;
const r2 = await waitForChange({ read: async () => { reads2++; return { stamp: 'A', at: null, added: [] }; }, since: 'A', now: c2.now, sleep: c2.sleep });
assert.deepEqual([r2.stamp, r2.changed, r2.held_ms], ['A', false, HOLD_MS_MAX]);
assert.equal(reads2, Math.ceil(HOLD_MS_MAX / 2_000) + 1);
// a hold over the cap is clamped; a zero hold is one read
const c3 = clock();
const r3 = await waitForChange({ read: async () => ({ stamp: 'A' }), since: 'A', holdMs: 600_000, now: c3.now, sleep: c3.sleep });
assert.equal(r3.held_ms, HOLD_MS_MAX);
const c4 = clock();
const r4 = await waitForChange({ read: async () => ({ stamp: 'A' }), since: 'A', holdMs: 0, now: c4.now, sleep: c4.sleep });
assert.equal(r4.held_ms, 0);
});
test('waitForChange: an empty store answers the empty stamp (the app then waits its own floor)', async () => {
const c = clock();
const r = await waitForChange({ read: async () => null, since: 'A', now: c.now, sleep: c.sleep });
assert.deepEqual([r.stamp, r.changed, r.held_ms], ['', true, 0]);
});
test('RateLimit: the window, the reset, independent keys', () => {
const c = clock();
const rl = new RateLimit({ perMinute: 3, now: c.now });
assert.equal(rl.take('a'), null); assert.equal(rl.take('a'), null); assert.equal(rl.take('a'), null);
assert.equal(rl.take('a'), 60);
assert.equal(rl.take('b'), null);
c.set(c.now() + 59_000);
assert.equal(rl.take('a'), 1);
c.set(c.now() + 1_000);
assert.equal(rl.take('a'), null);
});
test('handler: GET before any POST answers the empty stamp; POST needs auth and a stamp; GET then sees it', async () => {
const db = fakeDb();
const c = clock();
const h = makeHandler({ sql: db.sql, authed: r => r.headers['x-relay-token'] === 'T' ? 'token' : null, readJson: async r => r.body || {}, json, now: c.now, sleep: c.sleep });
let r = res(); await h(req('GET', {}), r);
assert.deepEqual(reply(r), { code: 200, ok: true, stamp: '', at: null, added: [], changed: false, held_ms: 0 });
r = res(); await h(req('POST', {}, { body: { stamp: 'S1' } }), r);
assert.equal(reply(r).code, 401);
r = res(); await h(req('POST', {}, { headers: { 'x-relay-token': 'T' }, body: { stamp: '' } }), r);
assert.equal(reply(r).code, 400);
r = res(); await h(req('POST', {}, { headers: { 'x-relay-token': 'T' }, body: { stamp: 'S1', added: ['run-1', 7, ''] } }), r);
assert.deepEqual(reply(r), { code: 200, ok: true, stamp: 'S1', at: '2026-10-05T12:00:00.000Z', added: ['run-1', '7'] });
r = res(); await h(req('GET', { since: 'S0' }), r);
assert.deepEqual(reply(r), { code: 200, ok: true, stamp: 'S1', at: '2026-10-05T12:00:00.000Z', added: ['run-1', '7'], changed: true, held_ms: 0 });
r = res(); await h(req('PUT', {}), r);
assert.equal(reply(r).code, 405);
// the stamp may also come as ?stamp= with the token in the path (the /r/<token>/wake rewrite)
r = res(); await h(req('POST', { token: 'T', stamp: 'S2' }, { headers: { 'x-relay-token': 'T' }, body: {} }), r);
assert.equal(reply(r).stamp, 'S2');
});
test('handler: a held GET answers when a POST lands, with the held time', async () => {
const db = fakeDb();
db.set('S1');
const c = clock();
const h = makeHandler({ sql: db.sql, authed: () => 'token', readJson: async r => r.body || {}, json, now: c.now, sleep: async ms => { c.sleep(ms); if (c.now() >= 1_000_000 + 6_000 && !db.rows.some(x => x.stamp === 'S2')) db.set('S2', ['job-9']); } });
const r = res(); await h(req('GET', { since: 'S1' }), r);
const j = reply(r);
assert.deepEqual([j.stamp, j.changed, j.added, j.held_ms], ['S2', true, ['job-9'], 6_000]);
});
test('handler: the rate limit answers 429 with Retry-After and never holds; over the hold cap the answer is immediate', async () => {
const db = fakeDb();
db.set('S1');
const c = clock();
const h = makeHandler({ sql: db.sql, authed: () => null, readJson: async () => ({}), json, limiter: new RateLimit({ perMinute: 1, now: c.now }), now: c.now, sleep: c.sleep });
let r = res(); await h(req('GET', { since: 'S1', hold: 2 }), r);
assert.equal(reply(r).held_ms, 2_000);
r = res(); await h(req('GET', { since: 'S1' }), r);
assert.deepEqual([reply(r).code, reply(r).error, r.headers['Retry-After']], [429, 'rate limited', '58']);
const h2 = makeHandler({ sql: db.sql, authed: () => null, readJson: async () => ({}), json, holds: { n: HOLDS_MAX }, now: c.now, sleep: c.sleep });
r = res(); await h2(req('GET', { since: 'S1' }), r);
assert.deepEqual([reply(r).code, reply(r).held_ms], [200, 0]);
});
test('handler: a database error is a 500 with its message, not a hang', async () => {
const c = clock();
const h = makeHandler({ sql: async () => { throw new Error('connection refused'); }, authed: () => 'token', readJson: async () => ({ stamp: 'S1' }), json, now: c.now, sleep: c.sleep });
const r = res(); await h(req('GET', { since: 'S1' }), r);
assert.deepEqual([reply(r).code, reply(r).error], [500, 'connection refused']);
});

View file

@ -13,8 +13,21 @@
{
"source": "/r/:token/api/:fn",
"destination": "/api/relay?token=:token&fn=:fn"
},
{
"source": "/wake",
"destination": "/api/wake"
},
{
"source": "/r/:token/wake",
"destination": "/api/wake?token=:token"
}
],
"functions": {
"api/wake.mjs": {
"maxDuration": 60
}
},
"headers": [
{
"source": "/(.*)",