159 lines
9.8 KiB
JavaScript
159 lines
9.8 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}.
|
|
// [&machine=<id8>&v=<version>&job=<id>] the ping (MF-11, 0.3.21): the app names itself on every request and the
|
|
// relay keeps one row per machine (relay_wake_seen: last poll, version, last job);
|
|
// the console shows "silent since <time>, last job <name>" after PING_SILENT_S
|
|
// without a poll. No secret: id8 is what the console already shows.
|
|
// 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);
|
|
export const validMachine = s => typeof s === 'string' && /^[0-9a-f]{8}$/.test(s);
|
|
export const PING_SILENT_S = 900; // 15 minutes without a poll from a machine: its job channel is dead
|
|
|
|
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;
|
|
}
|
|
}
|
|
|
|
/** The ping table: one row per machine, the last poll, the app version and the last job it named. */
|
|
export async function ensureSeenTable(sql, state = {}) {
|
|
if (state.seenTable) return;
|
|
await sql(`CREATE TABLE IF NOT EXISTS relay_wake_seen (machine text PRIMARY KEY, last_seen timestamptz NOT NULL DEFAULT now(), version text NOT NULL DEFAULT '', last_job text NOT NULL DEFAULT '')`);
|
|
state.seenTable = true;
|
|
}
|
|
|
|
/** Records one poll from a machine (the GET's machine, v and job); invalid or absent ids record nothing. */
|
|
export async function recordSeen(sql, { machine, version, job }, state = {}) {
|
|
if (!validMachine(machine)) return false;
|
|
await ensureSeenTable(sql, state);
|
|
await sql(`INSERT INTO relay_wake_seen (machine, last_seen, version, last_job) VALUES ($1, now(), $2, $3)
|
|
ON CONFLICT (machine) DO UPDATE SET last_seen = now(), version = EXCLUDED.version, last_job = CASE WHEN EXCLUDED.last_job = '' THEN relay_wake_seen.last_job ELSE EXCLUDED.last_job END`,
|
|
[machine, clip(version, 32), clip(job, 80)]);
|
|
return true;
|
|
}
|
|
|
|
/** Every machine's last poll: [{machine, last_seen, version, last_job}]; [] before the first ping. */
|
|
export async function seenList(sql) {
|
|
try {
|
|
const rows = await sql(`SELECT machine, last_seen, version, last_job FROM relay_wake_seen ORDER BY machine`);
|
|
return rows.map(r => ({ machine: String(r.machine), last_seen: isoOf(r.last_seen), version: String(r.version || ''), last_job: String(r.last_job || '') }));
|
|
} catch (e) {
|
|
if (/relay_wake_seen/.test(String(e.message)) && /does not exist/.test(String(e.message))) return [];
|
|
throw e;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* What the console says about a machine's job channel from its last ping: {silent, silent_s, since, last_job, version,
|
|
* never}. Silent after PING_SILENT_S without a poll; `since` is the last poll's time (what the card prints after
|
|
* "silent since"); never = no ping on record.
|
|
*/
|
|
export function pingState(seen, now = Date.now(), silentAfterS = PING_SILENT_S) {
|
|
if (!seen || !seen.last_seen) return { silent: true, never: true, silent_s: null, since: null, last_job: '', version: '' };
|
|
const t = Date.parse(seen.last_seen);
|
|
const silent_s = isNaN(t) ? null : Math.max(0, Math.round((now - t) / 1000));
|
|
return { silent: silent_s === null || silent_s > silentAfterS, never: false, silent_s, since: isNaN(t) ? seen.last_seen : new Date(t).toISOString(), last_job: seen.last_job || '', version: seen.version || '' };
|
|
}
|
|
|
|
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, state = {} }) {
|
|
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 }); }
|
|
// the ping lands before the hold, so a machine is "seen" the moment it asks, not 45 s later
|
|
if (q.machine !== undefined) { try { await recordSeen(sql, { machine: clip(q.machine, 16), version: q.v, job: q.job }, state); } catch (e) { /* the wake still answers */ } }
|
|
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) });
|
|
}
|
|
};
|
|
}
|