igneum/relay/lib/wake.mjs
igneum-josh 5ee7909d5c 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>
2026-10-05 09:22:02 +01:00

111 lines
6.4 KiB
JavaScript

// 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) });
}
};
}