app jobs: wake long-poll on the relay, fetch within seconds of a publish; safety-net poll 2 minutes (was 10)

Josh, 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 <url>?since=<stamp>, 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 <noreply@anthropic.com>
This commit is contained in:
igneum-josh 2026-10-05 09:25:29 +01:00
parent 3da3c88184
commit 7481bcbc96

View file

@ -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 (Josh,
/// 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<Fetched, String>),
/// 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<Active>,
needs_logged: std::collections::HashSet<String>,
fingerprint: String,
wake: Arc<WakeCtl>,
/// 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<Shared>, 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<Shared>) {
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<Shared>, 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 <WAKE_URL>?since=<stamp> 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>) -> 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<String, String> {
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<Shared>, url: String, ctl: Arc<WakeCtl>) {
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 ----------------------------------------------------------------------------------------------