From 5ee7909d5cfb07ddf92846cce686165ab809e99b Mon Sep 17 00:00:00 2001 From: igneum-josh <337424239+igneum-josh@users.noreply.github.com> Date: Mon, 5 Oct 2026 09:22:02 +0100 Subject: [PATCH] relay: /wake long-poll for the apps' remote jobs (public GET held 45 s, authenticated POST of the stamp) GET /wake?since= 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//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 --- .github/workflows/ci.yml | 4 +- relay/README.md | 2 + relay/api/wake.mjs | 18 +++++ relay/lib/wake.mjs | 111 ++++++++++++++++++++++++++++ relay/test/wake.test.mjs | 152 +++++++++++++++++++++++++++++++++++++++ relay/vercel.json | 13 ++++ 6 files changed, 298 insertions(+), 2 deletions(-) create mode 100644 relay/api/wake.mjs create mode 100644 relay/lib/wake.mjs create mode 100644 relay/test/wake.test.mjs diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9d1cb6931..c0562bf13 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -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 diff --git a/relay/README.md b/relay/README.md index 9758cfa33..e22ab9082 100644 --- a/relay/README.md +++ b/relay/README.md @@ -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=` 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//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 `, `drop ""|`, `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 `, `url`. Playbooks live in `relay/playbooks/`; `run` fills `__DL_BASE__` in from `~/.config/igneum/dl-token`. diff --git a/relay/api/wake.mjs b/relay/api/wake.mjs new file mode 100644 index 000000000..21736cb6b --- /dev/null +++ b/relay/api/wake.mjs @@ -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//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=[&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); +} diff --git a/relay/lib/wake.mjs b/relay/lib/wake.mjs new file mode 100644 index 000000000..25cf7c26d --- /dev/null +++ b/relay/lib/wake.mjs @@ -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=[&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) }); + } + }; +} diff --git a/relay/test/wake.test.mjs b/relay/test/wake.test.mjs new file mode 100644 index 000000000..a17189fbf --- /dev/null +++ b/relay/test/wake.test.mjs @@ -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//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']); +}); diff --git a/relay/vercel.json b/relay/vercel.json index 783c74e94..818273be4 100644 --- a/relay/vercel.json +++ b/relay/vercel.json @@ -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": "/(.*)",