diff --git a/app/igneum-app/src/jobrun.rs b/app/igneum-app/src/jobrun.rs index 10cd8b340..5c2d24327 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 (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), + /// 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 ----------------------------------------------------------------------------------------------