From 0550f58857029ad69bacb7c3bd256bd6deb1ff38 Mon Sep 17 00:00:00 2001 From: igneum-labs <337424239+igneum-labs@users.noreply.github.com> Date: Mon, 5 Oct 2026 08:22:02 +0000 Subject: [PATCH 1/3] 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 38a82d683..d9c6673c2 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": "/(.*)", From 7e0fc7da657cc23739f8f555b4ef9e8da116153f Mon Sep 17 00:00:00 2001 From: igneum-labs <337424239+igneum-labs@users.noreply.github.com> Date: Mon, 5 Oct 2026 08:22:02 +0000 Subject: [PATCH 2/3] publish-jobs: wake the apps after a verified deploy; jobs.mjs status shows the woken latency; 0.3.6 plan publish-jobs.sh --deploy POSTs the new stamp (published_at plus 8 hex of the file's sha256) and the added id to the relay's /wake once the live file verifies. The relay token goes in a 600-mode header file, never on the command line or the screen. Prints "woke the apps (stamp ...)" or a one-line warning; the apps' 2-minute poll still catches it. tools/jobs.mjs status reads relay_wake (one row per publish with the ids it added) and prints "woken +N s after the publish" for a machine's latest job that a publish added; nothing when the table does not exist yet. docs/plans/release-0.3.6.md: "Instant jobs" section with the design, the expected latency and a TODO row per machine for the measured number once 0.3.6 is live. packaging/ota/README.md: the 10-minute poll is history. Co-Authored-By: Claude Fable 5.1 --- docs/plans/release-0.3.6.md | 30 +++++++++++++++++++++++++++++ packaging/ota/README.md | 4 +++- packaging/ota/publish-jobs.sh | 36 ++++++++++++++++++++++++++++++----- tools/jobs.mjs | 17 +++++++++++++++-- 4 files changed, 79 insertions(+), 8 deletions(-) diff --git a/docs/plans/release-0.3.6.md b/docs/plans/release-0.3.6.md index 4c2a0a448..d8ab5a320 100644 --- a/docs/plans/release-0.3.6.md +++ b/docs/plans/release-0.3.6.md @@ -17,3 +17,33 @@ Written 5 October 2026, 08:45 BST, while proving v0 went live on the devnet at D - Consensus override changes must land on every node at once: a hand node restarted early with a different `proving_v0_activation_daa` was refused by every peer (digest handshake) and sat isolated at a lower height for 20 minutes. Order that works: publish the manifest override, `update-now` to every app, wait for every app node to log the new parameters, then restart the hand nodes and the seed with the same file. - `scratchpad/restart-hand-nodes.sh` died silently after `igneumd --version` (the 0.3.5 binary exits 1 after printing) under `set -e`; the restart it reported never happened. Every restart script ends by printing the new pids and their start times. - Switching proving on needs no app restart: `POST /api/prove {"on":true}` (the job `prove-on-pc2-84100` does this after installing the CUDA host into `/opt/igneum` for the app's WSL user). + +## Instant jobs (5 October 2026) + +the project lead: "why is it taking so long for pc2 and pc1s tasks to spin up? can we speed it up?". Before 0.3.6 every app polled +`igneum-jobs.json` every 10 minutes (`CHECK_EVERY_S = 600`), so a job published from the Mac waited up to 10 minutes +on each PC. The PCs have no inbound ports and one miner is off the LAN, so the fix is a wake signal the app pulls +over an outbound connection. Branch `job-wake`. + +| Piece | What it does | Where | +|---|---|---| +| Wake endpoint | `GET /wake?since=`: public (the apps hold no token), 30 a minute per IP, held up to 45 s, 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) records a stamp. Rows in Neon `relay_wake`, created by the first POST. `maxDuration` 60 s. | `relay/api/wake.mjs`, `relay/lib/wake.mjs`, `relay/vercel.json`, test `relay/test/wake.test.mjs` (fake database and clock; in CI's site job) | +| App waker | One thread per app long-polls the endpoint with the last stamp; a changed stamp sends `Event::Wake` and the engine fetches the jobs file at once (signature check unchanged). First reply seeds the stamp. Backoff 5, 15, 60 s while the relay is unreachable, 10 s floor between requests, a full hold when the relay has no stamp yet. One log line per wake, one when it falls back, one when it recovers. `IGNEUM_APP_JOBS_WAKE_URL` overrides the URL (set and empty: no waker). | `app/igneum-app/src/jobrun.rs` (`WakeState`, `wake_loop`), unit tests for the stamp and backoff logic | +| Fallback poll | 2 minutes (was 10), the first poll still 40 to 60 s after start. An unchanged file is no longer logged on every poll. | `jobrun.rs` `CHECK_EVERY_S` | +| Publisher | After the live file verifies, `--deploy` POSTs the stamp (`published_at` plus 8 hex of the file's sha256) and the added id; the token goes in a header file, never on the command line; prints "woke the apps" or a one-line warning. | `packaging/ota/publish-jobs.sh` `wake_apps` | +| Status | `node tools/jobs.mjs status` shows "woken +N s after the publish" per machine for a job a publish added. | `tools/jobs.mjs` | + +Expected latency, publish to job start: the relay re-reads the stamp every 2 s inside the hold, the app's curl returns +at once, the engine fetches the file (two small downloads) and starts the job on its next tick. About 3 to 6 s when +the app is mid-hold, plus up to 10 s if the app was inside its floor after an earlier reply; the 2-minute poll is the +ceiling when the relay is down. Numbers below are measured, not estimated. + +| Machine | Publish to started_at | Date | Source | +|---|---|---|---| +| PC 1 (ae432dc7) | TODO (measure with `node tools/jobs.mjs status` after the first 0.3.6 job) | | | +| PC 2 (1ccfe586) | TODO | | | +| Mac | TODO | | | + +Not yet done on this branch: the relay deploy (`relay/`, by the owner; the `/wake` route and the 60 s `maxDuration` +go live with it, the `relay_wake` table appears on the first POST), the first live publish, and the Windows curl path +of the long-poll (curl.exe 8.x in System32; the 58 s `--max-time` was reviewed, not run). diff --git a/packaging/ota/README.md b/packaging/ota/README.md index 74b64be3a..428ec3eca 100644 --- a/packaging/ota/README.md +++ b/packaging/ota/README.md @@ -133,7 +133,9 @@ screenshots are `docs/design/app-screens/update-*.png`. Windows: reviewed only, the project lead's rule: one app on both PCs that the Mac can send commands and files to over the line, so everything is tested and built without a person at the PC. The channel is `igneum-jobs.json` plus `igneum-jobs.json.sig`, next to the -update manifest, signed with the same OTA key and verified by the same code; the apps poll it every 10 minutes. +update manifest, signed with the same OTA key and verified by the same code. Since 0.3.6 (5 October 2026) every app +holds a long-poll on the relay's public `/wake` and fetches the file within seconds of `--deploy` (the script records +the new stamp there); a poll every 2 minutes is the fallback (it was 10 minutes before 0.3.6). The relay (`relay/`) stays for the Mac and for humans; its PC agent is replaced by this. | Piece | Where | diff --git a/packaging/ota/publish-jobs.sh b/packaging/ota/publish-jobs.sh index b102c72f3..d5defb547 100755 --- a/packaging/ota/publish-jobs.sh +++ b/packaging/ota/publish-jobs.sh @@ -1,10 +1,11 @@ #!/usr/bin/env bash # Publishes signed remote jobs for the Igneum Miner apps: igneum-jobs.json (canonical JSON, sorted keys, no # whitespace) and its detached Ed25519 signature igneum-jobs.json.sig, next to the update manifest in the downloads -# folder (dl//), signed on this Mac with the OTA key ~/.config/igneum/ota-signing-key. Every app polls the -# file every 10 minutes (app/igneum-app/src/jobrun.rs), verifies it with the public key compiled into -# src/manifest.rs, runs each job that targets it ONCE per id, and reports to the log intake as -# run_id job-- (read back with tools/jobs.mjs). +# folder (dl//), signed on this Mac with the OTA key ~/.config/igneum/ota-signing-key. Every app holds a +# long-poll on the relay's public /wake (relay/api/wake.mjs) and fetches the file the moment --deploy records the new +# stamp there (0.3.6, app/igneum-app/src/jobrun.rs; a poll every 2 minutes is the fallback), verifies it with the +# public key compiled into src/manifest.rs, runs each job that targets it ONCE per id, and reports to the log intake +# as run_id job-- (read back with tools/jobs.mjs). # # packaging/ota/publish-jobs.sh add --kind run --target 1ccfe586 --script path.ps1 [--elevated] [--stop-miners] \ # [--timeout-minutes 60] [--shell powershell|bash] --title "..." [--expires-hours 48] [--deploy] @@ -153,6 +154,30 @@ verify_live() { # (5 s apart); 0 = verified, 1 = not, with the reason on return 1 } +# After a verified deploy: records the new jobs-file stamp on the relay (relay/api/wake.mjs), where every app holds a +# long-poll and fetches the file the moment the stamp moves (0.3.6). The stamp is published_at plus 8 hex of the +# file's sha256, so a re-signed file wakes the apps too. The relay token (~/.config/igneum/relay-token) travels in a +# header file, never on the command line or the screen. A failure here is a warning: the apps poll every 2 minutes. +wake_apps() { # + local tf="$HOME/.config/igneum/relay-token" rurl="https://relay.igneum.network" sum size pub stamp hdr out rc=0 added='[]' + [ -f "$tf" ] || { echo "warning: no $tf; the apps were not woken (they poll every 2 minutes)" >&2; return 0; } + [ -f "$HOME/.config/igneum/relay-url" ] && rurl="$(tr -d '[:space:]' < "$HOME/.config/igneum/relay-url")" + rurl="${rurl%/}" + read -r sum size < <("$SIGNER" sha256 "$JOBS") + pub="$(python3 -c 'import json,sys; print(json.load(open(sys.argv[1])).get("published_at", ""))' "$JOBS")" + stamp="$pub.${sum:0:8}" + if [ -n "${1:-}" ]; then added="[\"$1\"]"; fi + hdr="$(mktemp)"; chmod 600 "$hdr" + printf 'x-relay-token: %s\n' "$(tr -d '[:space:]' < "$tf")" > "$hdr" + out="$(curl -sS --max-time 20 -X POST "$rurl/wake" -H 'Content-Type: application/json' -H @"$hdr" --data-binary "{\"stamp\":\"$stamp\",\"added\":$added}" 2>&1)" || rc=$? + rm -f "$hdr" + if [ "$rc" = 0 ] && printf '%s' "$out" | grep -q '"ok":true'; then + echo "woke the apps (stamp $stamp)" + else + echo "warning: the wake call to $rurl/wake failed (${out:0:160}); the apps poll every 2 minutes and catch it" >&2 + fi +} + if [ "$CMD" = verify ]; then [ -f "$JOBS" ] || { echo "no local jobs file at $JOBS" >&2; exit 1; } verify_live "$TRIES"; exit $? @@ -315,7 +340,8 @@ if [ "$DEPLOY" = 1 ]; then (cd "$DLSITE" && npx --yes vercel@latest --global-config "$HOME/.config/igneum/vercel" deploy --prod --yes 2>&1 | sed "s#$TOKEN##g"; exit "${PIPESTATUS[0]}") \ || { echo "the deploy failed (the Vercel CLI's exit status above); nothing verified" >&2; exit 1; } verify_live "$TRIES" || exit 1 - echo "the apps pick it up within 10 minutes (Settings > remote jobs > Check now at once); results: node tools/jobs.mjs ${ID:-}" + wake_apps "$ID" + echo "the apps fetch it within seconds when woken, else within 2 minutes (Settings > remote jobs > Check now at once); results: node tools/jobs.mjs ${ID:-}" else if [ -n "$DLSITE" ]; then echo "not deployed: cd $DLSITE && npx --yes vercel@latest --global-config ~/.config/igneum/vercel deploy --prod --yes (or re-run with --deploy)" diff --git a/tools/jobs.mjs b/tools/jobs.mjs index d124169a9..5b1fadc55 100755 --- a/tools/jobs.mjs +++ b/tools/jobs.mjs @@ -1,7 +1,9 @@ #!/usr/bin/env node // Mac side of the remote jobs (app/igneum-app/src/jobs.rs, published by packaging/ota/publish-jobs.sh). // node tools/jobs.mjs the published jobs file: fetched from the downloads host, signature checked -// node tools/jobs.mjs status per machine, from the log intake: the latest job run and its SUMMARY line +// node tools/jobs.mjs status per machine, from the log intake: the latest job run and its SUMMARY line, +// plus the woken latency (job started_at minus the publish that added it, +// from relay_wake, written by publish-jobs.sh --deploy since 0.3.6) // node tools/jobs.mjs the result: the newest upload per label under run_id job-- // node tools/jobs.mjs --all every upload, oldest first (the 5-minute progress reports of a long job) // node tools/jobs.mjs watch poll the intake every 30 s until every reporting machine is final @@ -79,9 +81,20 @@ if (a === 'status') { SELECT DISTINCT ON (machine) machine, run_id, label, received_at, lines FROM miner_logs WHERE run_id LIKE 'job-%' AND label LIKE 'job-%' ORDER BY machine, received_at DESC`); if (!rows.length) { console.log('no job reports in the intake yet'); process.exit(0); } + // the woken latency: each publish is one relay_wake row (stamp, at, meta.added = the ids it added); a machine's + // latest job that a publish added shows its started_at minus that publish's at. No table yet: nothing shown. + let publishes = []; + try { publishes = await sql('SELECT stamp, at, meta FROM relay_wake ORDER BY id DESC LIMIT 200'); } catch { publishes = []; } + const tsOf = v => Date.parse(String(v || '').replace(' ', 'T').replace(/([+-]\d\d)$/, '$1:00')); + const publishedAt = id => { for (const p of publishes) { const a = p.meta && Array.isArray(p.meta.added) ? p.meta.added : []; if (a.includes(id)) return tsOf(p.at); } return NaN; }; + const woken = s => { + if (!s || !s.job || !s.started_at) return ''; + const d = Math.round((tsOf(s.started_at) - publishedAt(s.job)) / 1000); + return Number.isFinite(d) && d >= 0 && d < 86400 ? ` woken +${d} s after the publish` : ''; + }; for (const r of rows) { const s = summaryOf(r.lines); - console.log(`${r.machine.padEnd(28)} ${r.run_id.padEnd(44)} ${when(r.received_at)} ${s ? `${s.status} exit ${s.exit} after ${s.duration_s} s: ${s.summary}` : '(no SUMMARY line)'}`); + console.log(`${r.machine.padEnd(28)} ${r.run_id.padEnd(44)} ${when(r.received_at)} ${s ? `${s.status} exit ${s.exit} after ${s.duration_s} s: ${s.summary}` : '(no SUMMARY line)'}${woken(s)}`); if (s && s.errors && s.errors.length) for (const e of s.errors) console.log(`${''.padEnd(28)} ERROR ${e.slice(0, 160)}`); } process.exit(0); From c40aa303a6fa4fb22022cc13dd73262c2c288a8c Mon Sep 17 00:00:00 2001 From: igneum-labs <337424239+igneum-labs@users.noreply.github.com> Date: Mon, 5 Oct 2026 08:25:29 +0000 Subject: [PATCH 3/3] app jobs: wake long-poll on the relay, fetch within seconds of a publish; safety-net poll 2 minutes (was 10) the project lead, 5 October 2026: "why is it taking so long for pc2 and pc1s tasks to spin up?". The PCs have no inbound ports and one miner is off the LAN, so one thread per app now long-polls the relay's public /wake with the last stamp (GET ?since=, held up to 45 s there). A changed stamp sends Event::Wake and the engine fetches the jobs file at once (signature check unchanged); a wake during a fetch in flight fetches again right after it. The first reply only seeds the stamp. Backoff 5, 15, then 60 s while the relay is unreachable, a 10 s floor between requests, a full hold when the relay has no stamp yet: never a busy loop. One log line per wake, one when it falls back, one on recovery. IGNEUM_APP_JOBS_WAKE_URL overrides the URL (set and empty: no waker). CHECK_EVERY_S 600 -> 120: the poll is the fallback now; the first poll stays 40 to 60 s after start. An unchanged file is no longer logged on every poll. Remote jobs off stops the waker too. Unit tests: the stamp logic (seed, same, changed, empty), the backoff sequence and the recovery flag, the floor, the URL and query. cargo test -p igneum-app: 60 + 22 pass; cargo build --release -p igneum-app: finished. Co-Authored-By: Claude Fable 5.1 --- app/igneum-app/src/jobrun.rs | 236 ++++++++++++++++++++++++++++++++++- 1 file changed, 232 insertions(+), 4 deletions(-) diff --git a/app/igneum-app/src/jobrun.rs b/app/igneum-app/src/jobrun.rs index 6d2873559..8b2fa9e84 100644 --- a/app/igneum-app/src/jobrun.rs +++ b/app/igneum-app/src/jobrun.rs @@ -1,5 +1,8 @@ //! The remote-job runner (the model is src/jobs.rs). Driven from the engine's tick like the updater: -//! poll (40 s after start, then every 10 minutes): fetch igneum-jobs.json and its .sig from the folder of the +//! poll (40 s after start, then every 2 minutes) and, since 0.3.6, a wake: one thread long-polls the relay's public +//! /wake and the engine fetches the moment publish-jobs.sh records a new jobs-file stamp (seconds, not minutes; +//! backoff 5, 15, 60 s while the relay is unreachable, the 2-minute poll carries on). The fetch: igneum-jobs.json +//! and its .sig from the folder of the //! update manifest, verify with the OTA public key, parse, keep the jobs this machine has not run that target it //! (machine id, platform, requirements probed here: wsl, wsl-prover, nvidia) //! -> run them one at a time, in file order; each id at most once (jobs-state.json, written before the run) @@ -15,7 +18,7 @@ //! -> the dashboard shows the running job (strip) and the history (Settings), and the switch "Allow remote jobs //! from Igneum (signed)" with the key fingerprint; off aborts the running job and stops polling //! Environment (tests): IGNEUM_APP_JOBS_URL overrides the jobs file URL, IGNEUM_APP_JOBS_CHECK_SECS the interval, -//! IGNEUM_APP_JOBS_FIRST_SECS the first delay. +//! IGNEUM_APP_JOBS_FIRST_SECS the first delay, IGNEUM_APP_JOBS_WAKE_URL the wake endpoint (set and empty: no waker). use crate::engine::{Cmd, Shared}; use crate::jobbuild as jb; @@ -29,8 +32,19 @@ use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; -const CHECK_EVERY_S: u64 = 600; +/// The safety-net poll. Before 0.3.6 this was 600 s and a published job waited up to 10 minutes on every PC (the project lead, +/// 5 October 2026: "why is it taking so long for pc2 and pc1s tasks to spin up?"); the wake below makes it seconds. +const CHECK_EVERY_S: u64 = 120; const RETRY_AFTER_ERROR_S: u64 = 300; +/// The wake endpoint (relay/api/wake.mjs): a public, rate-limited long-poll that answers the moment publish-jobs.sh +/// records a new jobs-file stamp. No token: the apps hold none. IGNEUM_APP_JOBS_WAKE_URL overrides it. +const WAKE_URL: &str = "https://relay.igneum.network/wake"; +/// The relay holds a request this long (its function is capped at 60 s); curl's max-time sits 13 s above it. +const WAKE_HOLD_S: u64 = 45; +/// Two wake requests are never closer than this (the relay allows 30 a minute per address, and both PCs share one). +const WAKE_FLOOR_S: u64 = 10; +/// Waits after the 1st, 2nd and every later failed wake request. +const WAKE_BACKOFF_S: [u64; 3] = [5, 15, 60]; const PROGRESS_REPORT_EVERY_S: u64 = 300; const GPU_IDLE_PCT: f64 = 5.0; const GPU_IDLE_WAIT_S: u64 = 180; @@ -43,6 +57,8 @@ const HISTORY_SHOWN: usize = 20; pub enum Event { Fetched(Result), + /// the relay's stamp moved (the waker thread): fetch the jobs file now + Wake(String), Progress { id: String, stage: String }, Finished { id: String, outcome: Outcome }, } @@ -112,6 +128,11 @@ pub struct Jobs { active: Option, needs_logged: std::collections::HashSet, fingerprint: String, + wake: Arc, + /// a wake arrived while a fetch was in flight: fetch again as soon as it ends + wake_pending: bool, + /// the published_at of the last fetched file: an unchanged file is not logged every 2 minutes + last_published: String, } impl Jobs { @@ -145,8 +166,18 @@ impl Jobs { active: None, needs_logged: std::collections::HashSet::new(), fingerprint: manifest::fingerprint(manifest::OTA_PUBLIC_KEY_HEX), + wake: Arc::new(WakeCtl { on: AtomicBool::new(allowed) }), + wake_pending: false, + last_published: String::new(), }; j.publish(shared); + if !j.url.is_empty() { + let wake_url = wake_url_of(std::env::var("IGNEUM_APP_JOBS_WAKE_URL").ok()); + if !wake_url.is_empty() { + let (sh, ctl) = (shared.clone(), j.wake.clone()); + std::thread::spawn(move || wake_loop(sh, wake_url, ctl)); + } + } let sh = shared.clone(); std::thread::spawn(move || { let ctx = account_context(); @@ -208,6 +239,7 @@ impl Jobs { pub fn set_allowed(&mut self, shared: &Arc, on: bool) { self.allowed = on; + self.wake.on.store(on, Ordering::Relaxed); { let mut s = shared.settings.lock().unwrap(); s.remote_jobs = on; @@ -269,6 +301,7 @@ impl Jobs { fn start_fetch(&mut self, shared: &Arc) { let every = std::env::var("IGNEUM_APP_JOBS_CHECK_SECS").ok().and_then(|v| v.parse().ok()).unwrap_or(CHECK_EVERY_S); self.next_check = Instant::now() + Duration::from_secs(every); + self.wake_pending = false; if self.url.is_empty() { let mut st = shared.state.lock().unwrap(); st.jobs.error = "no jobs URL in this build (no update manifest configured)".into(); @@ -356,15 +389,32 @@ impl Jobs { added += 1; } } - shared.log(&format!("jobs: file of {} (published {}): {} new for this machine, {} queued, {} waiting on requirements", f.total, f.published_at, added, self.queue.len(), f.pending.len())); + // polled every 2 minutes since 0.3.6: the line is for a changed file or new work, not every poll + if added > 0 || f.published_at != self.last_published { + shared.log(&format!("jobs: file of {} (published {}): {} new for this machine, {} queued, {} waiting on requirements", f.total, f.published_at, added, self.queue.len(), f.pending.len())); + } + self.last_published = f.published_at.clone(); if added > 0 { shared.event("info", &format!("{added} remote job{} received from Igneum", if added == 1 { "" } else { "s" })); } } } + if self.wake_pending { + // the file moved while this fetch ran: the next tick fetches again, whatever the retry delay + self.wake_pending = false; + self.next_check = Instant::now(); + } self.publish(shared); None } + Event::Wake(_stamp) => { + if self.busy { + self.wake_pending = true; + } else { + self.next_check = Instant::now(); + } + None + } Event::Progress { id, stage } => { if self.active.as_ref().map(|a| a.job.id == id).unwrap_or(false) { shared.state.lock().unwrap().jobs.stage = stage; @@ -588,6 +638,137 @@ fn upload_file(shared: &Arc, job: &Job, path: &Path, label_prefix: &str) crate::update::upload_log(&p.log_intake_url, &p.log_intake_key, &label, &machine, &job.run_id(&shared.runtime.machine_id), path, &shared.upload_header()) } +// ---- the wake signal (0.3.6) ---------------------------------------------------------------------------------------- +// The PCs have no inbound ports and may sit off the LAN (one miner is in the US), so a new jobs file is announced over +// an outbound long-poll: GET ?since= is held up to 45 s by the relay and answered the moment the +// stamp recorded by publish-jobs.sh changes. One thread per app; the engine fetches on Event::Wake. The thread never +// busy-loops: a reply that came back early is followed by the rest of a 10 s floor, a failing endpoint backs off +// 5, 15, then 60 s, and an empty stamp (nothing published yet) waits a full hold. + +/// Shared with the waker thread: it polls only while remote jobs are allowed. +pub struct WakeCtl { + on: AtomicBool, +} + +/// The stamp logic, free of I/O for the tests. +#[derive(Default, Debug)] +struct WakeState { + stamp: String, + seeded: bool, + failures: u32, +} + +#[derive(Debug, PartialEq)] +enum WakeReply { + /// the first stamp seen: remembered, nothing to fetch (the first poll covers the start) + Seeded, + Same, + /// the stamp moved: fetch now + Changed { from: String }, + /// the relay holds no stamp yet + Empty, +} + +impl WakeState { + /// A successful reply. The bool says the endpoint had been failing (one recovery line in the log). + fn reply(&mut self, stamp: &str) -> (WakeReply, bool) { + let recovered = self.failures > 0; + self.failures = 0; + if stamp.is_empty() { + return (WakeReply::Empty, recovered); + } + if !self.seeded { + self.seeded = true; + self.stamp = stamp.to_string(); + return (WakeReply::Seeded, recovered); + } + if stamp == self.stamp { + return (WakeReply::Same, recovered); + } + let from = std::mem::replace(&mut self.stamp, stamp.to_string()); + (WakeReply::Changed { from }, recovered) + } + + /// A failed request: how long to wait (5, 15, 60, 60, ... s) and whether this is the first failure of a run. + fn failed(&mut self) -> (Duration, bool) { + self.failures += 1; + let i = (self.failures as usize - 1).min(WAKE_BACKOFF_S.len() - 1); + (Duration::from_secs(WAKE_BACKOFF_S[i]), self.failures == 1) + } + + /// What to wait after a reply so two requests are never closer than the floor. + fn pause_after(elapsed: Duration) -> Duration { + Duration::from_secs(WAKE_FLOOR_S).saturating_sub(elapsed) + } +} + +/// The wake URL: the environment overrides the built-in one; set and empty switches the waker off. +fn wake_url_of(env: Option) -> String { + match env { + None => WAKE_URL.to_string(), + Some(u) => u.trim().to_string(), + } +} + +fn wake_query(url: &str, since: &str) -> String { + if since.is_empty() { + url.to_string() + } else { + format!("{url}{}since={since}", if url.contains('?') { '&' } else { '?' }) + } +} + +/// One long-poll. Ok carries the relay's stamp (empty when it holds none). +fn wake_request(url: &str, since: &str) -> Result { + let full = wake_query(url, since); + let max_time = (WAKE_HOLD_S + 13).to_string(); + let (code, out) = run_capture(Command::new(crate::platform::tool("curl")).args(["-fsS", "--max-time", &max_time, &full]), Duration::from_secs(WAKE_HOLD_S + 20)); + if code != Some(0) { + return Err(format!("curl exit {code:?}: {}", short_out(&out))); + } + let v: Value = serde_json::from_str(out.trim()).map_err(|e| format!("bad reply ({e}): {}", short_out(&out)))?; + if v.get("ok").and_then(|b| b.as_bool()) != Some(true) { + return Err(format!("relay: {}", short_out(&out))); + } + Ok(v.get("stamp").and_then(|s| s.as_str()).unwrap_or("").to_string()) +} + +fn wake_loop(shared: Arc, url: String, ctl: Arc) { + let mut st = WakeState::default(); + loop { + if !ctl.on.load(Ordering::Relaxed) { + std::thread::sleep(Duration::from_secs(5)); + continue; + } + let t0 = Instant::now(); + match wake_request(&url, &st.stamp) { + Ok(stamp) => { + let (r, recovered) = st.reply(&stamp); + if recovered { + shared.log("wake: the relay answers again; a new jobs file is fetched within seconds from here"); + } + match r { + WakeReply::Seeded => shared.log(&format!("wake: listening at {url} (jobs stamp {stamp}); a new jobs file is fetched within seconds")), + WakeReply::Changed { from } => { + shared.log(&format!("wake: jobs stamp {from} -> {stamp}; fetching the jobs file now")); + shared.send(Cmd::Job(Event::Wake(stamp))); + } + WakeReply::Same => {} + WakeReply::Empty => std::thread::sleep(Duration::from_secs(WAKE_HOLD_S)), + } + std::thread::sleep(WakeState::pause_after(t0.elapsed())); + } + Err(e) => { + let (wait, first) = st.failed(); + if first { + shared.log(&format!("wake: {url} unreachable ({e}); the {}-minute poll carries on; retrying in {} s, then {} s, then every {} s", CHECK_EVERY_S / 60, WAKE_BACKOFF_S[0], WAKE_BACKOFF_S[1], WAKE_BACKOFF_S[2])); + } + std::thread::sleep(wait); + } + } + } +} + // ---- fetching the jobs file and probing requirements ----------------------------------------------------------------- fn curl(args: &[&str], limit: Duration) -> Result<(), String> { @@ -1080,6 +1261,53 @@ mod tests { let d = collect_done(0, 0, vec![], Some(Ran { code: None, timed_out: false })); assert_eq!((d.status.as_str(), d.exit), ("failed", -1)); } + + #[test] + fn wake_state_seeds_then_reports_each_change_once() { + let mut w = WakeState::default(); + assert_eq!(w.reply(""), (WakeReply::Empty, false)); + assert_eq!(w.reply("S1"), (WakeReply::Seeded, false)); + assert_eq!(w.reply("S1"), (WakeReply::Same, false)); + assert_eq!(w.reply("S2"), (WakeReply::Changed { from: "S1".into() }, false)); + assert_eq!(w.reply("S2"), (WakeReply::Same, false)); + assert_eq!(w.stamp, "S2"); + // an empty reply after seeding keeps the last stamp, so the next request still carries it + assert_eq!(w.reply(""), (WakeReply::Empty, false)); + assert_eq!(w.stamp, "S2"); + } + + #[test] + fn wake_backoff_is_5_15_60_and_stays_there_until_a_reply() { + let mut w = WakeState::default(); + w.reply("S1"); + assert_eq!(w.failed(), (Duration::from_secs(5), true)); + assert_eq!(w.failed(), (Duration::from_secs(15), false)); + assert_eq!(w.failed(), (Duration::from_secs(60), false)); + assert_eq!(w.failed(), (Duration::from_secs(60), false)); + // the reply after an outage reports the recovery once and is still compared with the stamp from before it + assert_eq!(w.reply("S2"), (WakeReply::Changed { from: "S1".into() }, true)); + assert_eq!(w.reply("S2"), (WakeReply::Same, false)); + assert_eq!(w.failed(), (Duration::from_secs(5), true)); + } + + #[test] + fn wake_never_busy_loops() { + assert_eq!(WakeState::pause_after(Duration::from_millis(300)), Duration::from_millis(9_700)); + assert_eq!(WakeState::pause_after(Duration::from_secs(45)), Duration::ZERO); + assert!(WAKE_BACKOFF_S.iter().all(|s| *s >= 5)); + assert!(WAKE_HOLD_S + 13 < 60, "curl's max-time must stay under the relay function's 60 s cap"); + assert!(CHECK_EVERY_S <= 120, "the safety-net poll is the fallback when the relay is down"); + } + + #[test] + fn wake_url_and_query() { + assert_eq!(wake_url_of(None), WAKE_URL); + assert_eq!(wake_url_of(Some(String::new())), ""); + assert_eq!(wake_url_of(Some(" http://127.0.0.1:4180/wake ".into())), "http://127.0.0.1:4180/wake"); + assert_eq!(wake_query("https://r/wake", ""), "https://r/wake"); + assert_eq!(wake_query("https://r/wake", "2026-10-05T11:02:17Z.5e7b56f5"), "https://r/wake?since=2026-10-05T11:02:17Z.5e7b56f5"); + assert_eq!(wake_query("https://r/api/wake?x=1", "S"), "https://r/api/wake?x=1&since=S"); + } } // ---- kind: shard-benchmark ----------------------------------------------------------------------------------------------