From 3eb7af819f372b9c671cdbe6641077b2d9605576 Mon Sep 17 00:00:00 2001 From: igneum-labs <337424239+igneum-labs@users.noreply.github.com> Date: Sun, 4 Oct 2026 13:43:04 +0000 Subject: [PATCH] Igneum Miner: the NVIDIA power cap is judged by reading nvidia-smi back (anything but the requested watts = NOT applied), the card tile and Settings say 'power cap: N W applied' or 'power cap NOT applied (needs the administrator prompt)' with a Retry, the elevated step runs through the window host (ShellExecuteEx runas, a UI context) with the PowerShell path as the 150 s fallback, events either way, cap state in the stability line (PC 2 mined uncapped at 118 MH/s on 0.3.1) Co-Authored-By: Claude Fable 5.1 --- app/igneum-app/src/config.rs | 18 + app/igneum-app/src/detect.rs | 5 + app/igneum-app/src/engine.rs | 156 ++++- app/igneum-app/src/jobrun.rs | 1104 ++++++++++++++++++++++++++++++++++ app/igneum-app/src/jobs.rs | 833 +++++++++++++++++++++++++ app/igneum-app/src/main.rs | 5 +- app/igneum-app/src/server.rs | 4 + app/igneum-app/ui/app.css | 3 + app/igneum-app/ui/app.js | 19 +- app/windows/host.cpp | 27 + 10 files changed, 2146 insertions(+), 28 deletions(-) create mode 100644 app/igneum-app/src/jobrun.rs create mode 100644 app/igneum-app/src/jobs.rs diff --git a/app/igneum-app/src/config.rs b/app/igneum-app/src/config.rs index f968ce40e..3aa285e6a 100644 --- a/app/igneum-app/src/config.rs +++ b/app/igneum-app/src/config.rs @@ -108,6 +108,11 @@ pub struct Packaged { pub live_page: String, #[serde(default)] pub download_page: String, + /// Consensus parameters the packager pins for the bundled node: the engine writes them to + /// /override-params.json and starts igneumd with --override-params-file (the difficulty v2 activation + /// height, 4 October 2026: `{"difficulty_v2_activation_daa": N}`). Absent or empty = no override file. + #[serde(default)] + pub node_override_params: Option, } impl Packaged { @@ -183,3 +188,16 @@ impl Runtime { format!("grpc://127.0.0.1:{}", self.rpc_port) } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn packaged_carries_the_node_override_params() { + let p: Packaged = serde_json::from_str(r#"{"update_manifest":"","node_override_params":{"difficulty_v2_activation_daa":123456}}"#).unwrap(); + assert_eq!(p.node_override_params.as_ref().unwrap()["difficulty_v2_activation_daa"], 123456); + let p: Packaged = serde_json::from_str(r#"{"update_manifest":""}"#).unwrap(); + assert!(p.node_override_params.is_none()); + } +} diff --git a/app/igneum-app/src/detect.rs b/app/igneum-app/src/detect.rs index 566580c6e..318ec820a 100644 --- a/app/igneum-app/src/detect.rs +++ b/app/igneum-app/src/detect.rs @@ -100,6 +100,11 @@ pub fn apply_defaults(c: &mut CardState) { } /// NVIDIA power limits per card index: (default, current, min, max) in watts. +#[cfg(target_os = "macos")] +pub fn nvidia_power_limits() -> std::collections::HashMap { + std::collections::HashMap::new() +} + #[cfg(not(target_os = "macos"))] pub fn nvidia_power_limits() -> std::collections::HashMap { let mut out = std::collections::HashMap::new(); diff --git a/app/igneum-app/src/engine.rs b/app/igneum-app/src/engine.rs index 061fcb805..344f378eb 100644 --- a/app/igneum-app/src/engine.rs +++ b/app/igneum-app/src/engine.rs @@ -48,8 +48,12 @@ pub enum Cmd { ClockCheck, ClockSync, ClockSynced(Result), - /// the elevated nvidia-smi -pl step finished: (what was asked, result) - PowerApplied(String, Result<(), String>), + /// the elevated nvidia-smi -pl step finished: (what was asked, how it ran, the limits read back per device) + PowerApplied(String, Result<(), String>, std::collections::HashMap), + /// the window host ran the elevated command line (Windows): Ok or the reason + ElevatedDone(Result<(), String>), + /// the user asked for the cap again (the Retry button) + ApplyPower, Quit, } @@ -357,6 +361,8 @@ pub struct Engine { telemetry_retry_at: Instant, power_busy: bool, power_restore_pending: bool, + /// an elevated step handed to the window host: (command line, what, requested watts per device, since) + power_via_host: Option<(String, String, std::collections::HashMap, Instant)>, stability: std::collections::HashMap, last_stability: Instant, last_settings_save: Instant, @@ -422,6 +428,7 @@ impl Engine { telemetry_retry_at: now, power_busy: false, power_restore_pending: false, + power_via_host: None, stability: std::collections::HashMap::new(), last_stability: now, last_settings_save: now, @@ -647,26 +654,56 @@ impl Engine { self.resolve_clock(); self.clock_next_https = Instant::now() + Duration::from_secs(6); } - Cmd::PowerApplied(what, r) => { + Cmd::PowerApplied(what, r, readback) => { self.power_busy = false; - match r { - Ok(()) => { - self.shared.event("ok", &format!("GPU power cap set: {what}")); - let mut st = self.st(); - for c in st.mining.cards.iter_mut().filter(|c| c.vendor == "nvidia" && c.enabled) { - c.power_applied = true; - c.power_note = String::new(); + self.power_via_host = None; + // the truth is what nvidia-smi reads back, not whether the prompt said yes + let mut applied = Vec::new(); + let mut missing = Vec::new(); + { + let mut st = self.st(); + for c in st.mining.cards.iter_mut().filter(|c| c.vendor == "nvidia" && c.enabled && c.power_default_w > 0.0) { + let want = requested_watts(c); + let got = readback.get(&c.device).copied().unwrap_or(c.power_limit_w); + if got > 0.0 { + c.power_limit_w = got; } - } - Err(e) => { - self.shared.event("error", &format!("GPU power cap not applied ({e}); mining at the card's current limit")); - let mut st = self.st(); - for c in st.mining.cards.iter_mut().filter(|c| c.vendor == "nvidia" && c.enabled) { + if got > 0.0 && (got - want).abs() < 1.0 { + c.power_applied = true; + c.power_note = format!("power cap: {} W applied", want as u64); + applied.push(format!("{} {} W", c.name, want as u64)); + } else { c.power_applied = false; - c.power_note = format!("power cap not applied: {e}"); + c.power_note = format!("power cap NOT applied (needs the administrator prompt): card reports {} W, wanted {} W", got as u64, want as u64); + missing.push(format!("{} reports {} W, wanted {} W", c.name, got as u64, want as u64)); } } } + if !applied.is_empty() { + self.shared.event("ok", &format!("GPU power cap in force: {}", applied.join(", "))); + } + if !missing.is_empty() { + let why = match &r { + Ok(()) => "the step ran but the card did not take it".to_string(), + Err(e) => e.clone(), + }; + self.shared.event("error", &format!("GPU power cap NOT applied ({why}): {}. Retry from the card tile.", missing.join("; "))); + } + if applied.is_empty() && missing.is_empty() { + self.shared.log(&format!("power cap: nothing to read back for {what}")); + } + } + Cmd::ElevatedDone(r) => { + if let Some((line, what, want, _)) = self.power_via_host.take() { + self.shared.log(&format!("window host ran the elevated step ({}): {line}", match &r { Ok(()) => "ok".to_string(), Err(e) => e.clone() })); + self.finish_power(what, r, want); + } + } + Cmd::ApplyPower => { + for c in self.st().mining.cards.iter_mut().filter(|c| c.vendor == "nvidia") { + c.power_applied = false; // force the step again + } + self.apply_power_limits("retry"); } Cmd::Quit => { self.quitting = true; @@ -727,6 +764,9 @@ impl Engine { for p in &r.peers { a.push(format!("--addpeer={p}")); } + if let Some(path) = self.node_override_file() { + a.push(format!("--override-params-file={}", path.display())); + } a.extend(["--nodnsseed", "--disable-upnp", "--nologfiles", "--yes"].iter().map(|s| s.to_string())); if r.unsynced_mining { a.push("--enable-unsynced-mining".into()); @@ -734,6 +774,29 @@ impl Engine { a } + /// The packager's consensus parameters (igneum-app.json `node_override_params`, for example the difficulty v2 + /// activation height) written to /override-params.json for --override-params-file. None when the + /// package pins nothing, so the node runs on the network's defaults as before. + fn node_override_file(&self) -> Option { + let v = self.shared.packaged.node_override_params.as_ref()?; + if v.as_object().map(|o| o.is_empty()).unwrap_or(true) { + return None; + } + let path = self.shared.runtime.app_dir.join("override-params.json"); + let text = serde_json::to_string_pretty(v).ok()?; + if std::fs::read_to_string(&path).ok().as_deref() != Some(text.as_str()) { + if let Some(d) = path.parent() { + let _ = std::fs::create_dir_all(d); + } + if let Err(e) = std::fs::write(&path, &text) { + self.shared.log(&format!("could not write {}: {e}; the node starts without the override file", path.display())); + return None; + } + self.shared.log(&format!("node override params written to {}: {}", path.display(), text.replace('\n', " "))); + } + Some(path) + } + fn start_node(&mut self) { self.node_starts += 1; let seg = if self.node_starts > 1 { format!("-r{}", self.node_starts) } else { String::new() }; @@ -1046,14 +1109,7 @@ impl Engine { for c in st.mining.cards.iter_mut().filter(|c| c.vendor == "nvidia" && c.enabled && c.power_default_w > 0.0) { let pct = if c.power_pct == 0 { 80 } else { c.power_pct.clamp(60, 100) }; c.power_pct = pct; - let mut watts = c.power_default_w * pct as f64 / 100.0; - if c.power_min_w > 0.0 { - watts = watts.max(c.power_min_w); - } - if c.power_max_w > 0.0 { - watts = watts.min(c.power_max_w); - } - let watts = watts.round(); + let watts = requested_watts(c); if (c.power_limit_w - watts).abs() < 1.0 && c.power_applied { continue; } @@ -1070,10 +1126,30 @@ impl Engine { self.shared.log(&format!("power cap ({why}): {}", cmds.join(" & "))); let line = cmds.join(" & "); let what = what.join(", "); + let want: std::collections::HashMap = self.st().mining.cards.iter().filter(|c| c.vendor == "nvidia" && c.enabled && c.power_default_w > 0.0).map(|c| (c.device.clone(), requested_watts(c))).collect(); + if self.wrapper && cfg!(windows) { + // the window host has a UI context: it shows the administrator prompt and reports back on stdin + self.power_via_host = Some((line.clone(), what.clone(), want, Instant::now())); + println!("ELEVATE {line}"); + let _ = std::io::stdout().flush(); + return; + } let shared = self.shared.clone(); std::thread::spawn(move || { let r = crate::platform::run_elevated(&line); - shared.send(Cmd::PowerApplied(what, r)); + std::thread::sleep(Duration::from_millis(800)); + let back: std::collections::HashMap = crate::detect::nvidia_power_limits().into_iter().map(|(k, v)| (k, v.1)).collect(); + shared.send(Cmd::PowerApplied(what, r, back)); + }); + } + + /// After an elevated step (ours or the host's): read the limits back and judge. + fn finish_power(&mut self, what: String, r: Result<(), String>, _want: std::collections::HashMap) { + let shared = self.shared.clone(); + std::thread::spawn(move || { + std::thread::sleep(Duration::from_millis(800)); + let back: std::collections::HashMap = crate::detect::nvidia_power_limits().into_iter().map(|(k, v)| (k, v.1)).collect(); + shared.send(Cmd::PowerApplied(what, r, back)); }); } @@ -1125,6 +1201,20 @@ impl Engine { self.last_stability = now; self.stability_line(); } + if let Some((line, what, want, since)) = self.power_via_host.clone() { + if now.duration_since(since) > Duration::from_secs(150) { + self.power_via_host = None; + self.shared.log("the window host did not answer the elevated step in 150 s; running it through PowerShell"); + let shared = self.shared.clone(); + std::thread::spawn(move || { + let r = crate::platform::run_elevated(&line); + std::thread::sleep(Duration::from_millis(800)); + let back: std::collections::HashMap = crate::detect::nvidia_power_limits().into_iter().map(|(k, v)| (k, v.1)).collect(); + let _ = want; + shared.send(Cmd::PowerApplied(what, r, back)); + }); + } + } } /// "index, draw, gpu temp, mem temp, limit" every 5 s. @@ -1171,11 +1261,12 @@ impl Engine { v.sort_by(|a, b| a.partial_cmp(b).unwrap()); let p95 = v[((v.len() as f64 * 0.95) as usize).min(v.len() - 1)]; self.shared.log(&format!( - "stability: {}: draw p95 {:.0} W, max {:.0} W (cap {:.0} W), max GPU {:.0} C, max memory {:.0} C, {} samples", + "stability: {}: draw p95 {:.0} W, max {:.0} W, limit {:.0} W ({}), max GPU {:.0} C, max memory {:.0} C, {} samples", c.name, p95, v[v.len() - 1], c.power_limit_w, + if c.power_applied { format!("cap in force, {}% of default", c.power_pct) } else { "cap NOT applied".to_string() }, s.max_tgpu, s.max_tmem, v.len() @@ -1913,6 +2004,19 @@ impl Engine { } } +/// The watts a card's cap asks for: power_pct of the default limit, inside the card's min and max. +fn requested_watts(c: &CardState) -> f64 { + let pct = if c.power_pct == 0 { 80 } else { c.power_pct.clamp(60, 100) }; + let mut w = c.power_default_w * pct as f64 / 100.0; + if c.power_min_w > 0.0 { + w = w.max(c.power_min_w); + } + if c.power_max_w > 0.0 { + w = w.min(c.power_max_w); + } + w.round() +} + /// Per-card session statistics for the stability line. #[derive(Default)] struct Stability { diff --git a/app/igneum-app/src/jobrun.rs b/app/igneum-app/src/jobrun.rs new file mode 100644 index 000000000..255ef3c0a --- /dev/null +++ b/app/igneum-app/src/jobrun.rs @@ -0,0 +1,1104 @@ +//! 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 +//! 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) +//! -> kinds: run (a script, optionally elevated, optionally with the miners stopped first), fetch (a file by +//! https and sha256 into the app data dir), collect (files by glob or a command's output to the log intake), +//! restart (miners, node or the app), update-now (the updater checks and installs at once), shard-benchmark +//! (miners stopped, GPU idle, the prove package fetched, prove-shard.sh inside WSL under a time cap, every +//! RESULT and STAGE line and the results/*.json to the intake, miners back) +//! -> every job reports to the log intake as run_id job--: the first line is a JSON summary +//! (SUMMARY {...}), then the captured output; long jobs report every 5 minutes while they run +//! -> 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. + +use crate::engine::{Cmd, Shared}; +use crate::jobs::{self, Eligibility, Job, Ledger}; +use crate::manifest; +use serde_json::{json, Value}; +use std::io::Write; +use std::path::{Path, PathBuf}; +use std::process::{Command, Stdio}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant}; + +const CHECK_EVERY_S: u64 = 600; +const RETRY_AFTER_ERROR_S: u64 = 300; +const PROGRESS_REPORT_EVERY_S: u64 = 300; +const GPU_IDLE_PCT: f64 = 5.0; +const GPU_IDLE_WAIT_S: u64 = 180; +const DEFAULT_DISTRO: &str = "Ubuntu-24.04"; +const DEFAULT_FIXTURES: &[&str] = &["block-338-shard1", "block-341-shards2", "block-344-shards4"]; +const HISTORY_SHOWN: usize = 20; + +pub enum Event { + Fetched(Result), + Progress { id: String, stage: String }, + Finished { id: String, outcome: Outcome }, +} + +pub struct Fetched { + pub published_at: String, + pub total: usize, + pub runnable: Vec, + /// jobs for this machine whose requirements are not met yet: (id, missing) + pub pending: Vec<(String, Vec)>, +} + +#[derive(Clone, Debug, Default)] +pub struct Outcome { + /// done | failed | timeout | aborted + pub status: String, + pub exit: i64, + pub summary: String, + pub results: Vec, + pub uploaded: bool, +} + +/// What the engine must do for the running job. +pub enum Action { + /// Stop the miners (the node keeps running) and hold them; then call `miners_stopped`. + StopMiners(String), + RestartMiners, + RestartNode, + /// The relaunch helper was started; the engine quits now. + RestartApp, + UpdateNow, +} + +/// Shared with the job's thread: the abort flag and the child's pid, so the engine can end a job on quit. +#[derive(Default)] +pub struct Ctl { + abort: AtomicBool, + pid: Mutex>, +} + +impl Ctl { + fn aborted(&self) -> bool { + self.abort.load(Ordering::Relaxed) + } +} + +struct Active { + job: Job, + run_id: String, + started: Instant, + started_unix: u64, + waiting_for_miners: bool, + holds_miners: bool, + ctl: Arc, +} + +pub struct Jobs { + url: String, + allowed: bool, + app_dir: PathBuf, + data_root: PathBuf, + dir: PathBuf, + ledger: Ledger, + ledger_path: PathBuf, + busy: bool, + next_check: Instant, + queue: Vec, + active: Option, + needs_logged: std::collections::HashSet, + fingerprint: String, +} + +impl Jobs { + pub fn new(shared: &Arc) -> Jobs { + let env = |k: &str| std::env::var(k).ok().filter(|v| !v.is_empty()); + let url = env("IGNEUM_APP_JOBS_URL").unwrap_or_else(|| jobs::jobs_url_from_manifest(&shared.packaged.update_manifest)); + let app_dir = shared.runtime.app_dir.clone(); + let dir = app_dir.join("jobs"); + let _ = std::fs::create_dir_all(&dir); + let ledger_path = app_dir.join("jobs-state.json"); + let mut ledger = Ledger::load(&ledger_path); + let now_unix = crate::platform::unix_now(); + for id in ledger.abandon_running(now_unix) { + shared.log(&format!("job {id}: was running when the app last stopped; marked aborted (it does not run again)")); + } + let _ = ledger.save(&ledger_path); + let jitter = shared.runtime.machine_id.bytes().fold(0u64, |a, b| a.wrapping_mul(31).wrapping_add(b as u64)); + let first = env("IGNEUM_APP_JOBS_FIRST_SECS").and_then(|v| v.parse().ok()).unwrap_or(40 + jitter % 20); + let allowed = shared.settings.lock().unwrap().remote_jobs; + let data_root = crate::platform::data_root(); + let j = Jobs { + url, + allowed, + app_dir, + data_root, + dir, + ledger, + ledger_path, + busy: false, + next_check: Instant::now() + Duration::from_secs(first), + queue: Vec::new(), + active: None, + needs_logged: std::collections::HashSet::new(), + fingerprint: manifest::fingerprint(manifest::OTA_PUBLIC_KEY_HEX), + }; + j.publish(shared); + j + } + + pub fn active(&self) -> bool { + self.active.is_some() + } + + /// True while a job has the miners stopped; the engine restarts them when this turns false. + pub fn holds_miners(&self) -> bool { + self.active.as_ref().map(|a| a.holds_miners).unwrap_or(false) + } + + // ---- state for the dashboard --------------------------------------------------------------------------------- + + fn publish(&self, shared: &Arc) { + let mut st = shared.state.lock().unwrap(); + let j = &mut st.jobs; + j.allowed = self.allowed; + j.key_fingerprint = self.fingerprint.clone(); + j.url_set = !self.url.is_empty(); + j.queued = self.queue.len() as u32; + if let Some(a) = &self.active { + j.active = true; + j.id = a.job.id.clone(); + j.kind = a.job.kind.clone(); + j.title = a.job.label(); + j.run_id = a.run_id.clone(); + j.started_at = a.started_unix as f64; + if a.waiting_for_miners { + j.stage = "stopping the miners".into(); + } + } else { + j.active = false; + j.id = String::new(); + j.kind = String::new(); + j.title = String::new(); + j.run_id = String::new(); + j.stage = String::new(); + j.started_at = 0.0; + j.results.clear(); + } + j.history = self + .ledger + .history(HISTORY_SHOWN) + .into_iter() + .map(|r| crate::state::JobHistory { id: r.id, kind: r.kind, title: r.title, status: r.status, started_at: r.started_at as f64, finished_at: r.finished_at as f64, exit: r.exit, run_id: r.run_id, uploaded: r.uploaded, summary: r.summary }) + .collect(); + } + + pub fn set_allowed(&mut self, shared: &Arc, on: bool) { + self.allowed = on; + { + let mut s = shared.settings.lock().unwrap(); + s.remote_jobs = on; + s.save(&shared.settings_path); + } + shared.state.lock().unwrap().settings.remote_jobs = on; + if on { + shared.event("info", "remote jobs allowed: signed jobs from Igneum run on this machine"); + self.next_check = Instant::now() + Duration::from_secs(2); + } else { + shared.event("info", "remote jobs switched off: nothing from the jobs file runs here"); + self.queue.clear(); + self.abort(shared, "remote jobs switched off"); + } + self.publish(shared); + } + + pub fn check_now(&mut self, shared: &Arc) { + if !self.allowed || self.busy { + return; + } + self.next_check = Instant::now(); + self.tick(shared); + } + + /// Ends the running job (quit, or the switch turned off). The engine releases the miners itself. + pub fn abort(&mut self, shared: &Arc, why: &str) { + let Some(a) = self.active.as_ref() else { return }; + a.ctl.abort.store(true, Ordering::Relaxed); + if let Some(pid) = *a.ctl.pid.lock().unwrap() { + kill_tree(pid); + } + shared.log(&format!("job {}: aborted ({why})", a.job.id)); + let id = a.job.id.clone(); + self.ledger.finish(&id, "aborted", -1, why, false, crate::platform::unix_now()); + let _ = self.ledger.save(&self.ledger_path); + self.active = None; + self.publish(shared); + } + + // ---- the tick ------------------------------------------------------------------------------------------------ + + pub fn tick(&mut self, shared: &Arc) -> Option { + let now = Instant::now(); + if self.allowed && !self.busy && now >= self.next_check { + self.start_fetch(shared); + } + if self.active.is_none() && self.allowed { + if let Some(job) = self.queue.first().cloned() { + self.queue.remove(0); + let a = self.start(shared, job); + self.publish(shared); + return a; + } + } + None + } + + 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); + 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(); + return; + } + self.busy = true; + let url = self.url.clone(); + let dir = self.dir.clone(); + let seen: std::collections::HashSet = self.ledger.records.keys().cloned().collect(); + let machine_id = shared.runtime.machine_id.clone(); + let shared2 = shared.clone(); + std::thread::spawn(move || { + let r = fetch_jobs(&url, &dir).map(|f| { + let now = crate::platform::unix_now(); + let platform = manifest::platform_name(); + let probes: Mutex> = Mutex::new(std::collections::HashMap::new()); + let have = |req: &str| -> bool { + if let Some(v) = probes.lock().unwrap().get(req) { + return *v; + } + let v = probe(req); + probes.lock().unwrap().insert(req.to_string(), v); + v + }; + let mut runnable = Vec::new(); + let mut pending = Vec::new(); + for j in &f.jobs { + if seen.contains(&j.id) { + continue; + } + match jobs::eligibility(j, &machine_id, platform, now, &have) { + Eligibility::Run => runnable.push(j.clone()), + Eligibility::Needs(m) => pending.push((j.id.clone(), m)), + _ => {} + } + } + Fetched { published_at: f.published_at, total: f.jobs.len(), runnable, pending } + }); + shared2.send(Cmd::Job(Event::Fetched(r))); + }); + } + + pub fn event(&mut self, shared: &Arc, ev: Event) -> Option { + match ev { + Event::Fetched(r) => { + self.busy = false; + shared.state.lock().unwrap().jobs.checked_at = crate::platform::unix_now_f(); + match r { + Err(e) => { + let quiet = e.starts_with("no jobs file"); + if !quiet { + self.next_check = Instant::now() + Duration::from_secs(RETRY_AFTER_ERROR_S); + } + shared.log(&format!("jobs: {e}")); + shared.state.lock().unwrap().jobs.error = if quiet { String::new() } else { e }; + } + Ok(f) => { + shared.state.lock().unwrap().jobs.error = String::new(); + for (id, needs) in &f.pending { + if self.needs_logged.insert(id.clone()) { + shared.log(&format!("job {id}: for this machine but needs {} (checked again every {} min)", needs.join(", "), CHECK_EVERY_S / 60)); + } + } + let mut added = 0; + for j in f.runnable { + if !self.ledger.seen(&j.id) && !self.queue.iter().any(|q| q.id == j.id) && self.active.as_ref().map(|a| a.job.id != j.id).unwrap_or(true) { + self.queue.push(j); + 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())); + if added > 0 { + shared.event("info", &format!("{added} remote job{} received from Igneum", if added == 1 { "" } else { "s" })); + } + } + } + self.publish(shared); + 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; + } + None + } + Event::Finished { id, outcome } => { + let Some(a) = self.active.as_ref() else { return None }; + if a.job.id != id { + return None; + } + let a = self.active.take().unwrap(); + let now = crate::platform::unix_now(); + self.ledger.finish(&id, &outcome.status, outcome.exit, &outcome.summary, outcome.uploaded, now); + let _ = self.ledger.save(&self.ledger_path); + let mins = a.started.elapsed().as_secs() / 60; + let kind = if outcome.status == "done" { "ok" } else { "error" }; + shared.event(kind, &format!("job {} ({}) {} after {} min, exit {}{}: {}", a.job.id, a.job.label(), outcome.status, mins, outcome.exit, if outcome.uploaded { ", report uploaded" } else { ", report NOT uploaded" }, short(&outcome.summary, 200))); + for r in outcome.results.iter().take(12) { + shared.log(&format!(" {r}")); + } + { + let mut st = shared.state.lock().unwrap(); + st.jobs.last = crate::state::JobHistory { id: id.clone(), kind: a.job.kind.clone(), title: a.job.label(), status: outcome.status.clone(), started_at: a.started_unix as f64, finished_at: now as f64, exit: outcome.exit, run_id: a.run_id.clone(), uploaded: outcome.uploaded, summary: outcome.summary.clone() }; + st.jobs.last_results = outcome.results.clone(); + } + self.publish(shared); + None + } + } + } + + /// Called by the engine once the miners are stopped for a job that asked for it. + pub fn miners_stopped(&mut self, shared: &Arc) { + let Some(a) = self.active.as_mut() else { return }; + if !a.waiting_for_miners { + return; + } + a.waiting_for_miners = false; + a.holds_miners = true; + let (job, run_id, ctl, started) = (a.job.clone(), a.run_id.clone(), a.ctl.clone(), a.started_unix); + self.spawn_run(shared, job, run_id, ctl, started); + self.publish(shared); + } + + fn start(&mut self, shared: &Arc, job: Job) -> Option { + let now = crate::platform::unix_now(); + let run_id = job.run_id(&shared.runtime.machine_id); + if !self.ledger.start(&job, &run_id, now) { + return None; + } + if let Err(e) = self.ledger.save(&self.ledger_path) { + shared.log(&format!("job {}: cannot write jobs-state.json ({e}); not running it", job.id)); + return None; + } + shared.event("info", &format!("job {} ({}) starts: {}", job.id, job.kind, job.label())); + let ctl = Arc::new(Ctl::default()); + let needs_miners_stopped = job.kind == "shard-benchmark" || (job.kind == "run" && job.bool_param("stop_miners_first")); + self.active = Some(Active { job: job.clone(), run_id: run_id.clone(), started: Instant::now(), started_unix: now, waiting_for_miners: needs_miners_stopped, holds_miners: false, ctl: ctl.clone() }); + match job.kind.as_str() { + "restart" | "update-now" => { + // engine-side; the report says what was asked and the ledger closes at once + let what = job.str_param("what"); + let (summary, action) = match (job.kind.as_str(), what.as_str()) { + ("update-now", _) => ("update check and install asked".to_string(), Action::UpdateNow), + (_, "miners") => ("miners restarted".to_string(), Action::RestartMiners), + (_, "node") => ("node restarted (the miners follow)".to_string(), Action::RestartNode), + _ => ("app restarting".to_string(), Action::RestartApp), + }; + if let Action::RestartApp = action { + if let Err(e) = spawn_relaunch_helper(shared) { + let o = Outcome { status: "failed".into(), exit: 1, summary: format!("could not start the relaunch helper: {e}"), ..Default::default() }; + self.event(shared, Event::Finished { id: job.id.clone(), outcome: o }); + return None; + } + } + let sink = Sink::new(shared, &job, &self.dir); + sink.line(&format!("{}: {summary}", job.kind)); + let uploaded = report(shared, &job, &sink, "done", 0, now, crate::platform::unix_now(), &summary, json!({})); + let o = Outcome { status: "done".into(), exit: 0, summary, results: vec![], uploaded }; + self.event(shared, Event::Finished { id: job.id.clone(), outcome: o }); + Some(action) + } + _ if needs_miners_stopped => Some(Action::StopMiners(format!("job {}: {}", job.id, job.label()))), + _ => { + self.spawn_run(shared, job, run_id, ctl, now); + None + } + } + } + + fn spawn_run(&self, shared: &Arc, job: Job, run_id: String, ctl: Arc, started: u64) { + let shared2 = shared.clone(); + let dir = self.dir.clone(); + let data_root = self.data_root.clone(); + let jobs_url = self.url.clone(); + std::thread::spawn(move || { + let sink = Sink::new(&shared2, &job, &dir); + sink.line(&format!("job {} ({}) on {} machine {} run {run_id}, started {}", job.id, job.kind, shared2.runtime.host, shared2.runtime.machine_id, jobs::format_time(started))); + let r = match job.kind.as_str() { + "run" => run_script(&shared2, &job, &sink, &dir, &data_root, &ctl, started), + "fetch" => run_fetch(&job, &sink, &dir, &data_root, &shared2.runtime.app_dir, &ctl), + "collect" => run_collect(&shared2, &job, &sink, &data_root, &ctl), + "shard-benchmark" => run_shard_benchmark(&shared2, &job, &sink, &data_root, &jobs_url, &ctl, started), + k => Err(format!("kind {k} is not run on this side")), + }; + let finished = crate::platform::unix_now(); + let (status, exit, summary, extra) = match r { + Ok(done) => (done.status.clone(), done.exit, done.summary.clone(), done.extra.clone()), + Err(e) => (if ctl.aborted() { "aborted".to_string() } else { "failed".to_string() }, -1, e, json!({})), + }; + sink.line(&format!("job {}: {status} (exit {exit}) after {} s: {}", job.id, finished.saturating_sub(started), summary)); + let uploaded = report(&shared2, &job, &sink, &status, exit, started, finished, &summary, extra); + let results = sink.results.lock().unwrap().clone(); + shared2.send(Cmd::Job(Event::Finished { id: job.id.clone(), outcome: Outcome { status, exit, summary, results, uploaded } })); + }); + } +} + +/// What a kind returns when it ran to the end (its own status: a non-zero exit is "failed", a cap is "timeout"). +struct Done { + status: String, + exit: i64, + summary: String, + extra: Value, +} + +// ---- the output sink ------------------------------------------------------------------------------------------------ + +/// Every line a job produces: appended to //job.log, the last one on the dashboard, RESULT and STAGE +/// lines kept for the summary and shown as they arrive. +pub struct Sink { + shared: Arc, + id: String, + pub dir: PathBuf, + pub log_path: PathBuf, + file: Mutex>, + pub results: Mutex>, +} + +impl Sink { + fn new(shared: &Arc, job: &Job, jobs_dir: &Path) -> Sink { + let dir = jobs_dir.join(&job.id); + let _ = std::fs::create_dir_all(&dir); + let log_path = dir.join("job.log"); + let file = std::fs::OpenOptions::new().create(true).append(true).open(&log_path).ok(); + Sink { shared: shared.clone(), id: job.id.clone(), dir, log_path, file: Mutex::new(file), results: Mutex::new(Vec::new()) } + } + fn line(&self, text: &str) { + let text = text.replace('\0', ""); + if let Some(f) = self.file.lock().unwrap().as_mut() { + let _ = writeln!(f, "{text}"); + } + let t = text.trim_end(); + let is_result = t.starts_with("RESULT") || t.starts_with("STAGE") || t.starts_with("BUILD FAILED"); + if is_result { + let mut r = self.results.lock().unwrap(); + r.push(t.to_string()); + if r.len() > 200 { + r.remove(0); + } + self.shared.log(&format!("job {}: {t}", self.id)); + } + let mut st = self.shared.state.lock().unwrap(); + if st.jobs.id == self.id { + st.jobs.message = short(t, 160); + if is_result { + st.jobs.results.push(t.to_string()); + if st.jobs.results.len() > 40 { + st.jobs.results.remove(0); + } + } + } + } + fn stage(&self, text: &str) { + self.line(&format!("== {text} ==")); + self.shared.log(&format!("job {}: {text}", self.id)); + self.shared.send(Cmd::Job(Event::Progress { id: self.id.clone(), stage: text.to_string() })); + } +} + +/// The report file (SUMMARY line, then the last 200 KB of the job log) to the intake. True when the intake took it. +fn report(shared: &Arc, job: &Job, sink: &Sink, status: &str, exit: i64, started: u64, finished: u64, summary: &str, extra: Value) -> bool { + let p = &shared.packaged; + let results = sink.results.lock().unwrap().clone(); + let head = jobs::summary_line(job, &shared.runtime.machine_id, &shared.runtime.host, status, exit, started, finished, summary, &results, extra); + let mut body = std::fs::read(&sink.log_path).unwrap_or_default(); + if body.len() > 200_000 { + body = body.split_off(body.len() - 200_000); + } + let path = sink.dir.join("report.log"); + let mut out = head.into_bytes(); + out.push(b'\n'); + out.extend_from_slice(&body); + if std::fs::write(&path, &out).is_err() { + return false; + } + if p.log_intake_url.is_empty() || p.log_intake_key.is_empty() { + return false; + } + let machine = format!("{}-{}", shared.runtime.host, shared.runtime.id8()); + crate::update::upload_log(&p.log_intake_url, &p.log_intake_key, &format!("job-{}", job.kind), &machine, &job.run_id(&shared.runtime.machine_id), &path) +} + +/// One file to the intake under the job's run id, labelled file-. The intake keeps the last 256 KB. +fn upload_file(shared: &Arc, job: &Job, path: &Path, label_prefix: &str) -> bool { + let p = &shared.packaged; + if p.log_intake_url.is_empty() || p.log_intake_key.is_empty() { + return false; + } + let name = path.file_name().map(|n| n.to_string_lossy().into_owned()).unwrap_or_default(); + let label: String = format!("{label_prefix}-{name}").chars().take(80).collect(); + let machine = format!("{}-{}", shared.runtime.host, shared.runtime.id8()); + crate::update::upload_log(&p.log_intake_url, &p.log_intake_key, &label, &machine, &job.run_id(&shared.runtime.machine_id), path) +} + +// ---- fetching the jobs file and probing requirements ----------------------------------------------------------------- + +fn curl(args: &[&str], limit: Duration) -> Result<(), String> { + let mut c = Command::new("curl"); + c.args(args); + let (code, out) = run_capture(&mut c, limit); + let t = out.trim(); + if code == Some(0) { Ok(()) } else { Err(if t.is_empty() { format!("curl exit {code:?}") } else { t.lines().last().unwrap_or("curl failed").to_string() }) } +} + +fn fetch_jobs(url: &str, dir: &Path) -> Result { + let jf = dir.join("jobs.json.new"); + let sf = dir.join("jobs.json.sig.new"); + let _ = std::fs::remove_file(&jf); + let _ = std::fs::remove_file(&sf); + curl(&["-fsSL", "--max-time", "20", "-o", &jf.display().to_string(), url], Duration::from_secs(25)).map_err(|e| if e.contains("404") { "no jobs file published".to_string() } else { format!("jobs file: {e}") })?; + curl(&["-fsSL", "--max-time", "20", "-o", &sf.display().to_string(), &format!("{url}.sig")], Duration::from_secs(25)).map_err(|e| format!("jobs signature: {e}"))?; + let bytes = std::fs::read(&jf).map_err(|e| e.to_string())?; + let sig = std::fs::read_to_string(&sf).map_err(|e| e.to_string())?; + let f = jobs::verify_and_parse(&bytes, sig.trim(), manifest::OTA_PUBLIC_KEY_HEX)?; + let _ = std::fs::rename(&jf, dir.join("jobs.json")); + let _ = std::fs::rename(&sf, dir.join("jobs.json.sig")); + Ok(f) +} + +/// Does this machine meet a requirement? Probed on the fetch thread, once per poll. +fn probe(req: &str) -> bool { + match req { + "nvidia" => run_capture(Command::new("nvidia-smi").arg("-L"), Duration::from_secs(20)).0 == Some(0), + "wsl" => cfg!(windows) && run_capture(Command::new("wsl.exe").arg("--status"), Duration::from_secs(30)).0 == Some(0), + "wsl-prover" => { + cfg!(windows) && run_capture(Command::new("wsl.exe").args(["-d", DEFAULT_DISTRO, "--", "bash", "-lc", "command -v cargo && ls ~/.sp1 && ls ~/igneum-prove"]), Duration::from_secs(90)).0 == Some(0) + } + _ => false, + } +} + +/// Runs a command to the end or the limit; stdout and stderr folded, no window on Windows. (None, ...) = killed. +fn run_capture(cmd: &mut Command, limit: Duration) -> (Option, String) { + cmd.stdin(Stdio::null()).stdout(Stdio::piped()).stderr(Stdio::piped()); + crate::platform::quiet(cmd); + let Ok(mut child) = cmd.spawn() else { return (None, "could not start".into()) }; + let out = child.stdout.take(); + let err = child.stderr.take(); + let reader = std::thread::spawn(move || { + let mut s = String::new(); + if let Some(o) = out { + let _ = std::io::Read::read_to_string(&mut std::io::BufReader::new(o), &mut s); + } + let mut e = String::new(); + if let Some(x) = err { + let _ = std::io::Read::read_to_string(&mut std::io::BufReader::new(x), &mut e); + } + (s, e) + }); + let deadline = Instant::now() + limit; + let mut code = None; + loop { + match child.try_wait() { + Ok(Some(st)) => { + code = st.code(); + break; + } + Ok(None) if Instant::now() < deadline => std::thread::sleep(Duration::from_millis(50)), + _ => { + let _ = child.kill(); + let _ = child.wait(); + break; + } + } + } + let (s, e) = reader.join().unwrap_or_default(); + (code, format!("{s}{e}").replace('\0', "")) +} + +// ---- running a child with the output streamed into the sink, under a cap and the abort flag ---------------------------- + +struct Ran { + /// None when killed (timeout or abort) + code: Option, + timed_out: bool, +} + +fn run_streamed(cmd: &mut Command, sink: &Sink, ctl: &Ctl, limit: Duration, shared: &Arc, job: &Job, started: u64, summary_while_running: &str) -> Result { + cmd.stdin(Stdio::null()).stdout(Stdio::piped()).stderr(Stdio::piped()); + crate::platform::quiet(cmd); + #[cfg(unix)] + { + use std::os::unix::process::CommandExt; + cmd.process_group(0); + } + let mut child = cmd.spawn().map_err(|e| format!("could not start: {e}"))?; + *ctl.pid.lock().unwrap() = Some(child.id()); + let mut readers = Vec::new(); + for (stream, is_err) in [(child.stdout.take().map(|s| Box::new(s) as Box), false), (child.stderr.take().map(|s| Box::new(s) as Box), true)] { + let Some(stream) = stream else { continue }; + let sink_shared = sink.shared.clone(); + let sink_id = sink.id.clone(); + let sink_dir = sink.dir.clone(); + let sink_log = sink.log_path.clone(); + readers.push(std::thread::spawn(move || { + // a light sink clone: the same file (append), the same state slot + let s = Sink { shared: sink_shared, id: sink_id, dir: sink_dir, log_path: sink_log.clone(), file: Mutex::new(std::fs::OpenOptions::new().append(true).open(&sink_log).ok()), results: Mutex::new(Vec::new()) }; + let reader = std::io::BufReader::new(stream); + for line in std::io::BufRead::split(reader, b'\n') { + let Ok(bytes) = line else { break }; + let text = String::from_utf8_lossy(&bytes).trim_end_matches('\r').to_string(); + s.line(&if is_err { format!("! {text}") } else { text }); + } + s.results.into_inner().unwrap_or_default() + })); + } + let deadline = Instant::now() + limit; + let mut next_report = Instant::now() + Duration::from_secs(PROGRESS_REPORT_EVERY_S); + let mut timed_out = false; + let code = loop { + match child.try_wait() { + Ok(Some(st)) => break st.code(), + Ok(None) => { + if ctl.aborted() || Instant::now() >= deadline { + timed_out = !ctl.aborted(); + sink.line(&format!("{}: ending the process", if timed_out { "time cap reached" } else { "aborted" })); + kill_tree(child.id()); + let _ = child.kill(); + let _ = child.wait(); + break None; + } + if Instant::now() >= next_report { + next_report = Instant::now() + Duration::from_secs(PROGRESS_REPORT_EVERY_S); + report(shared, job, sink, "running", 0, started, 0, summary_while_running, json!({})); + } + std::thread::sleep(Duration::from_millis(200)); + } + Err(e) => { + sink.line(&format!("wait failed: {e}")); + break None; + } + } + }; + *ctl.pid.lock().unwrap() = None; + for r in readers { + if let Ok(lines) = r.join() { + let mut all = sink.results.lock().unwrap(); + for l in lines { + if !all.contains(&l) { + all.push(l); + } + } + } + } + Ok(Ran { code, timed_out }) +} + +/// Ends a process and everything it started. Windows: taskkill /T; unix: the process group (set at spawn). +fn kill_tree(pid: u32) { + #[cfg(windows)] + { + let mut c = Command::new("taskkill"); + c.args(["/T", "/F", "/PID", &pid.to_string()]); + crate::platform::quiet(&mut c); + let _ = c.output(); + } + #[cfg(unix)] + unsafe { + libc::kill(-(pid as i32), libc::SIGTERM); + std::thread::sleep(Duration::from_millis(500)); + libc::kill(-(pid as i32), libc::SIGKILL); + } +} + +fn file_name_from_url(url: &str) -> String { + let name = url.rsplit('/').next().unwrap_or("file").split('?').next().unwrap_or("file"); + let clean: String = name.chars().filter(|c| c.is_ascii_alphanumeric() || *c == '.' || *c == '-' || *c == '_').collect(); + if clean.is_empty() || clean.starts_with('.') { "file".into() } else { clean } +} + +/// Downloads with resume into .part, checks the size (when given) and the sha256, renames. +fn fetch_file(url: &str, final_path: &Path, sha256: &str, size: Option, sink: &Sink) -> Result { + let part = PathBuf::from(format!("{}.part", final_path.display())); + if final_path.is_file() && manifest::sha256_file(final_path).map(|s| s == sha256).unwrap_or(false) { + sink.line(&format!("{} is already there with the right sha256", final_path.display())); + return Ok(std::fs::metadata(final_path).map(|m| m.len()).unwrap_or(0)); + } + let _ = std::fs::remove_file(final_path); + if let Some(d) = final_path.parent() { + let _ = std::fs::create_dir_all(d); + } + sink.line(&format!("downloading {url}")); + curl(&["-fsSL", "--retry", "3", "--retry-delay", "5", "-C", "-", "--max-time", "3600", "-o", &part.display().to_string(), url], Duration::from_secs(3660))?; + let got = std::fs::metadata(&part).map(|m| m.len()).unwrap_or(0); + if let Some(want) = size { + if want > 0 && got != want { + let _ = std::fs::remove_file(&part); + return Err(format!("size mismatch: got {got} bytes, the job says {want}")); + } + } + let sum = manifest::sha256_file(&part).map_err(|e| e.to_string())?; + if sum != sha256.to_ascii_lowercase() { + let _ = std::fs::remove_file(&part); + return Err("sha256 mismatch: the file is not what the job signed".into()); + } + std::fs::rename(&part, final_path).map_err(|e| e.to_string())?; + sink.line(&format!("downloaded {} bytes, sha256 ok", got)); + Ok(got) +} + +/// tar -xf (bsdtar on Windows 10 1803+ and macOS reads zip); PowerShell Expand-Archive as the Windows fallback. +fn extract(archive: &Path, into: &Path, sink: &Sink) -> Result<(), String> { + let _ = std::fs::create_dir_all(into); + let (code, out) = run_capture(Command::new("tar").args(["-xf", &archive.display().to_string(), "-C", &into.display().to_string()]), Duration::from_secs(600)); + if code == Some(0) { + sink.line(&format!("extracted {} into {}", archive.display(), into.display())); + return Ok(()); + } + if cfg!(windows) { + let script = format!("Expand-Archive -LiteralPath '{}' -DestinationPath '{}' -Force", archive.display().to_string().replace('\'', "''"), into.display().to_string().replace('\'', "''")); + let (c2, o2) = run_capture(Command::new("powershell").args(["-NoProfile", "-ExecutionPolicy", "Bypass", "-Command", &script]), Duration::from_secs(600)); + if c2 == Some(0) { + sink.line(&format!("extracted {} into {} (Expand-Archive)", archive.display(), into.display())); + return Ok(()); + } + return Err(format!("extract failed: tar: {}; Expand-Archive: {}", out.lines().last().unwrap_or(""), o2.lines().last().unwrap_or(""))); + } + Err(format!("extract failed: {}", out.lines().last().unwrap_or("tar failed"))) +} + +fn shell_for(job: &Job) -> String { + let s = job.str_param("shell"); + if !s.is_empty() { + return s; + } + if cfg!(windows) { "powershell".into() } else { "bash".into() } +} + +fn job_env(cmd: &mut Command, shared: &Arc, job: &Job, job_dir: &Path, data_root: &Path) { + cmd.env("IGNEUM_JOB_ID", &job.id); + cmd.env("IGNEUM_JOB_DIR", job_dir); + cmd.env("IGNEUM_JOB_RUN_ID", job.run_id(&shared.runtime.machine_id)); + cmd.env("IGNEUM_APP_DATA", data_root); + cmd.env("IGNEUM_APP_DIR", &shared.runtime.app_dir); + cmd.env("IGNEUM_LOG_DIR", &shared.runtime.log_dir); + cmd.env("IGNEUM_MACHINE_ID", &shared.runtime.machine_id); + cmd.env("IGNEUM_INTAKE_URL", &shared.packaged.log_intake_url); + cmd.env("IGNEUM_INTAKE_KEY", &shared.packaged.log_intake_key); + cmd.env("IGNEUM_APP_VERSION", crate::engine::VERSION); +} + +// ---- kind: run --------------------------------------------------------------------------------------------------------- + +fn run_script(shared: &Arc, job: &Job, sink: &Sink, jobs_dir: &Path, data_root: &Path, ctl: &Ctl, started: u64) -> Result { + let dir = jobs_dir.join(&job.id); + let shell = shell_for(job); + let body = job.str_param("script").replace("\r\n", "\n"); + let script = dir.join(if shell == "powershell" { "script.ps1" } else { "script.sh" }); + let text = if shell == "powershell" { body.replace('\n', "\r\n") } else { body }; + std::fs::write(&script, if shell == "powershell" { [b"\xEF\xBB\xBF".as_slice(), text.as_bytes()].concat() } else { text.into_bytes() }).map_err(|e| format!("cannot write the script: {e}"))?; + let elevated = cfg!(windows) && job.bool_param("elevated"); + let limit = Duration::from_secs(job.timeout_minutes() * 60); + sink.stage(&format!("running the {shell} script{} (cap {} min)", if elevated { " as administrator" } else { "" }, job.timeout_minutes())); + let mut cmd; + let out_file = dir.join("elevated-output.log"); + if elevated { + // the elevated process has its own environment: a wrapper sets the IGNEUM_* values, runs the script and + // sends everything to a file this side reads afterwards. One UAC prompt appears on the PC; with nobody + // there it times out and the job fails. + let env_lines: String = [("IGNEUM_JOB_ID", job.id.clone()), ("IGNEUM_JOB_DIR", dir.display().to_string()), ("IGNEUM_JOB_RUN_ID", job.run_id(&shared.runtime.machine_id)), ("IGNEUM_APP_DATA", data_root.display().to_string()), ("IGNEUM_APP_DIR", shared.runtime.app_dir.display().to_string()), ("IGNEUM_LOG_DIR", shared.runtime.log_dir.display().to_string()), ("IGNEUM_MACHINE_ID", shared.runtime.machine_id.clone()), ("IGNEUM_INTAKE_URL", shared.packaged.log_intake_url.clone()), ("IGNEUM_INTAKE_KEY", shared.packaged.log_intake_key.clone()), ("IGNEUM_APP_VERSION", crate::engine::VERSION.to_string())] + .iter() + .map(|(k, v)| format!("$env:{k} = '{}'\r\n", v.replace('\'', "''"))) + .collect(); + let wrapper = dir.join("elevated.ps1"); + let w = format!("{env_lines}& '{}' *>&1 | Out-File -FilePath '{}' -Encoding utf8\r\nexit $LASTEXITCODE\r\n", script.display().to_string().replace('\'', "''"), out_file.display().to_string().replace('\'', "''")); + std::fs::write(&wrapper, [b"\xEF\xBB\xBF".as_slice(), w.as_bytes()].concat()).map_err(|e| e.to_string())?; + let _ = std::fs::remove_file(&out_file); + let inner = format!("-NoProfile -ExecutionPolicy Bypass -File \"{}\"", wrapper.display()); + let ps = format!("$p = Start-Process -FilePath powershell.exe -ArgumentList '{}' -Verb RunAs -Wait -WindowStyle Hidden -PassThru; exit $p.ExitCode", inner.replace('\'', "''")); + cmd = Command::new("powershell"); + cmd.args(["-NoProfile", "-ExecutionPolicy", "Bypass", "-Command", &ps]); + } else if shell == "powershell" { + cmd = Command::new("powershell"); + cmd.args(["-NoProfile", "-ExecutionPolicy", "Bypass", "-File", &script.display().to_string()]); + } else { + cmd = Command::new("bash"); + cmd.arg(&script); + } + cmd.current_dir(&dir); + job_env(&mut cmd, shared, job, &dir, data_root); + let ran = run_streamed(&mut cmd, sink, ctl, limit, shared, job, started, "script running")?; + if elevated { + if let Ok(t) = std::fs::read_to_string(&out_file) { + for l in t.lines() { + sink.line(l); + } + } + } + finish_ran(ran, "script") +} + +fn finish_ran(ran: Ran, what: &str) -> Result { + match ran.code { + Some(0) => Ok(Done { status: "done".into(), exit: 0, summary: format!("{what} finished, exit 0"), extra: json!({}) }), + Some(c) => Ok(Done { status: "failed".into(), exit: c as i64, summary: format!("{what} exited with code {c}"), extra: json!({}) }), + None if ran.timed_out => Ok(Done { status: "timeout".into(), exit: -1, summary: format!("{what} hit the time cap and was ended"), extra: json!({}) }), + None => Err(format!("{what} was ended")), + } +} + +// ---- kind: fetch ------------------------------------------------------------------------------------------------------- + +fn fetch_base(dir: &str, job: &Job, jobs_dir: &Path, data_root: &Path, app_dir: &Path) -> PathBuf { + match dir { + "prove" => data_root.join("prove"), + "packs" => app_dir.join("packs"), + "updates" => app_dir.join("updates"), + _ => jobs_dir.join(&job.id), + } +} + +fn run_fetch(job: &Job, sink: &Sink, jobs_dir: &Path, data_root: &Path, app_dir: &Path, ctl: &Ctl) -> Result { + let url = job.str_param("url"); + let base = fetch_base(&job.str_param("dir"), job, jobs_dir, data_root, app_dir); + let to = job.str_param("to"); + let name = if to.is_empty() { PathBuf::from(file_name_from_url(&url)) } else { jobs::safe_rel_path(&to).ok_or("bad 'to'")? }; + let dest = base.join(&name); + sink.stage(&format!("fetching {} to {}", file_name_from_url(&url), dest.display())); + let size = fetch_file(&url, &dest, &job.str_param("sha256"), job.u64_param("size"), sink)?; + if ctl.aborted() { + return Err("aborted".into()); + } + let mut summary = format!("{} ({size} bytes, sha256 ok)", dest.display()); + if job.bool_param("extract") { + let sub = job.str_param("extract_dir"); + let into = if sub.is_empty() { base.clone() } else { base.join(jobs::safe_rel_path(&sub).ok_or("bad extract_dir")?) }; + if job.bool_param("fresh") && !sub.is_empty() && into.starts_with(&base) && into != base { + sink.line(&format!("fresh: removing {}", into.display())); + let _ = std::fs::remove_dir_all(&into); + } + sink.stage("extracting"); + extract(&dest, &into, sink)?; + summary.push_str(&format!(", extracted into {}", into.display())); + } + Ok(Done { status: "done".into(), exit: 0, summary, extra: json!({ "path": dest.display().to_string(), "size": size }) }) +} + +// ---- kind: collect ----------------------------------------------------------------------------------------------------- + +fn run_collect(shared: &Arc, job: &Job, sink: &Sink, data_root: &Path, ctl: &Ctl) -> Result { + let mut uploaded = 0u32; + let mut failed = 0u32; + let mut names = Vec::new(); + let globs = job.list_param("globs"); + if !globs.is_empty() { + sink.stage(&format!("collecting {} pattern{}", globs.len(), if globs.len() == 1 { "" } else { "s" })); + for g in &globs { + let files = jobs::glob_files(data_root, g); + sink.line(&format!("{g}: {} file{}", files.len(), if files.len() == 1 { "" } else { "s" })); + for f in files.iter().take(50) { + if ctl.aborted() { + return Err("aborted".into()); + } + let ok = upload_file(shared, job, f, "file"); + let name = f.strip_prefix(data_root).unwrap_or(f).display().to_string(); + sink.line(&format!(" {} {name} ({} bytes)", if ok { "uploaded" } else { "NOT uploaded" }, std::fs::metadata(f).map(|m| m.len()).unwrap_or(0))); + if ok { + uploaded += 1; + names.push(name); + } else { + failed += 1; + } + } + } + } + let command = job.str_param("command"); + if !command.trim().is_empty() { + sink.stage(&format!("running: {}", short(&command, 120))); + let mut cmd = if cfg!(windows) { + let mut c = Command::new("powershell"); + c.args(["-NoProfile", "-ExecutionPolicy", "Bypass", "-Command", &command]); + c + } else { + let mut c = Command::new("bash"); + c.args(["-lc", &command]); + c + }; + cmd.current_dir(data_root); + job_env(&mut cmd, shared, job, &sink.dir, data_root); + let ran = run_streamed(&mut cmd, sink, ctl, Duration::from_secs(600), shared, job, crate::platform::unix_now(), "collect command running")?; + sink.line(&format!("command exit {:?}", ran.code)); + } + let summary = format!("{uploaded} file{} uploaded{}{}", if uploaded == 1 { "" } else { "s" }, if failed > 0 { format!(", {failed} failed") } else { String::new() }, if command.trim().is_empty() { String::new() } else { ", command output in the report".into() }); + Ok(Done { status: if failed > 0 { "failed".into() } else { "done".into() }, exit: failed as i64, summary, extra: json!({ "uploaded_files": names }) }) +} + +// ---- kind: shard-benchmark ---------------------------------------------------------------------------------------------- + +fn gpu_utilisation() -> Option> { + let (code, out) = run_capture(Command::new("nvidia-smi").args(["--query-gpu=utilization.gpu", "--format=csv,noheader,nounits"]), Duration::from_secs(15)); + if code != Some(0) { + return None; + } + let v: Vec = out.lines().filter_map(|l| l.trim().parse().ok()).collect(); + if v.is_empty() { None } else { Some(v) } +} + +fn wait_gpu_idle(sink: &Sink, ctl: &Ctl) { + let deadline = Instant::now() + Duration::from_secs(GPU_IDLE_WAIT_S); + let mut idle_streak = 0; + while Instant::now() < deadline && !ctl.aborted() { + match gpu_utilisation() { + None => { + sink.line("nvidia-smi gave no utilisation reading; waiting 10 s and going on"); + std::thread::sleep(Duration::from_secs(10)); + return; + } + Some(v) => { + let max = v.iter().cloned().fold(0.0, f64::max); + if max < GPU_IDLE_PCT { + idle_streak += 1; + if idle_streak >= 2 { + sink.line(&format!("GPU idle: utilisation {max:.0}%")); + return; + } + } else { + idle_streak = 0; + sink.line(&format!("GPU still busy: utilisation {max:.0}%")); + } + } + } + std::thread::sleep(Duration::from_secs(3)); + } + sink.line("GPU did not go idle within the wait; going on anyway"); +} + +fn run_shard_benchmark(shared: &Arc, job: &Job, sink: &Sink, data_root: &Path, jobs_url: &str, ctl: &Ctl, started: u64) -> Result { + let distro = { let d = job.str_param("distro"); if d.is_empty() { DEFAULT_DISTRO.to_string() } else { d } }; + let user = job.str_param("wsl_user"); + let mut fixtures = job.list_param("fixtures"); + if fixtures.is_empty() { + fixtures = DEFAULT_FIXTURES.iter().map(|s| s.to_string()).collect(); + } + let zip_url = { let u = job.str_param("zip_url"); if u.is_empty() { jobs::jobs_url_from_manifest(jobs_url).replace(jobs::JOBS_FILE, "igneum-prove-wsl2.zip") } else { u } }; + if !zip_url.starts_with("https://") && !zip_url.starts_with("http://127.0.0.1:") { + return Err("no https zip_url and no jobs URL to derive it from".into()); + } + let base = data_root.join("prove"); + let pkg = base.join("igneum-prove-wsl2"); + let cap_min = job.timeout_minutes(); + let overall = Instant::now() + Duration::from_secs(cap_min * 60); + + sink.stage("waiting for the GPU to go idle"); + wait_gpu_idle(sink, ctl); + if ctl.aborted() { + return Err("aborted".into()); + } + + sink.stage("fetching the prove package"); + let zip = base.join("igneum-prove-wsl2.zip"); + fetch_file(&zip_url, &zip, &job.str_param("sha256"), job.u64_param("size"), sink)?; + if pkg.exists() { + sink.line(&format!("fresh: removing {}", pkg.display())); + std::fs::remove_dir_all(&pkg).map_err(|e| format!("cannot remove the old package: {e}"))?; + } + extract(&zip, &base, sink)?; + let script = pkg.join("prove-shard.sh"); + if !script.is_file() { + return Err(format!("{} is missing after the extract", script.display())); + } + let wsl_script = jobs::to_wsl_path(&script.display().to_string()).ok_or("the package path has no drive letter; WSL cannot see it")?; + let shard = fixtures[0].clone(); + let blocks = fixtures[1..].join(" "); + sink.stage(&format!("proving {shard} then {} inside {distro} (cap {cap_min} min)", if blocks.is_empty() { "nothing else".to_string() } else { blocks.clone() })); + let mut cmd = Command::new("wsl.exe"); + cmd.args(["-d", &distro]); + if !user.is_empty() { + cmd.args(["-u", &user]); + } + cmd.args(["--", "bash", &wsl_script, &shard, &blocks]); + cmd.current_dir(&pkg); + let remaining = overall.saturating_duration_since(Instant::now()).max(Duration::from_secs(60)); + let ran = run_streamed(&mut cmd, sink, ctl, remaining, shared, job, started, "shard benchmark running")?; + if ran.code.is_none() { + // the Linux side outlives wsl.exe: end the prover there too + let _ = run_capture(Command::new("wsl.exe").args(["-d", &distro, "--", "bash", "-c", "pkill -f igneum-prove-host; pkill -f prove-shard.sh; true"]), Duration::from_secs(30)); + } + sink.stage("uploading the results"); + let mut uploaded = Vec::new(); + let results_dir = pkg.join("results"); + let mut files: Vec = std::fs::read_dir(&results_dir).map(|rd| rd.flatten().map(|e| e.path()).filter(|p| p.extension().map(|x| x == "json").unwrap_or(false)).collect()).unwrap_or_default(); + files.sort(); + for f in &files { + let ok = upload_file(shared, job, f, "result"); + sink.line(&format!(" {} {}", if ok { "uploaded" } else { "NOT uploaded" }, f.display())); + if ok { + uploaded.push(f.file_name().map(|n| n.to_string_lossy().into_owned()).unwrap_or_default()); + } + } + let mut logs: Vec = std::fs::read_dir(&pkg).map(|rd| rd.flatten().map(|e| e.path()).filter(|p| p.file_name().map(|n| n.to_string_lossy().starts_with("prove-shards-")).unwrap_or(false)).collect()).unwrap_or_default(); + logs.sort(); + if let Some(l) = logs.last() { + let ok = upload_file(shared, job, l, "prove-log"); + sink.line(&format!(" {} {}", if ok { "uploaded" } else { "NOT uploaded" }, l.display())); + } + let results = sink.results.lock().unwrap().clone(); + let n_results = results.iter().filter(|r| r.starts_with("RESULT")).count(); + let mut done = finish_ran(ran, "prove-shard.sh")?; + done.summary = format!("{}; {n_results} RESULT line{}, {} result file{} uploaded", done.summary, if n_results == 1 { "" } else { "s" }, uploaded.len(), if uploaded.len() == 1 { "" } else { "s" }); + done.extra = json!({ "fixtures": fixtures, "distro": distro, "result_files": uploaded, "package": pkg.display().to_string() }); + Ok(done) +} + +// ---- kind: restart app ----------------------------------------------------------------------------------------------- + +/// A detached helper that starts the app again a few seconds after this engine has gone. +fn spawn_relaunch_helper(shared: &Arc) -> Result<(), String> { + let mut c; + #[cfg(target_os = "macos")] + { + let b = crate::platform::bundle_path().ok_or("not running from Igneum Miner.app")?; + c = Command::new("nohup"); + c.args(["bash", "-c", &format!("sleep 8; open -n '{}'", b.display().to_string().replace('\'', "'\\''"))]); + } + #[cfg(windows)] + { + let dir = std::env::current_exe().ok().and_then(|p| p.parent().map(|d| d.to_path_buf())).ok_or("cannot find the install folder")?; + let exe = dir.join("igneum-app.exe"); + let ps = format!("Start-Sleep 8; Start-Process -FilePath '{}' -ArgumentList '--launch' -WorkingDirectory '{}'", exe.display().to_string().replace('\'', "''"), dir.display().to_string().replace('\'', "''")); + c = Command::new("powershell"); + c.args(["-NoProfile", "-ExecutionPolicy", "Bypass", "-WindowStyle", "Hidden", "-Command", &ps]); + } + #[cfg(not(any(target_os = "macos", windows)))] + { + let exe = std::env::current_exe().map_err(|e| e.to_string())?; + c = Command::new("nohup"); + c.args(["bash", "-c", &format!("sleep 8; '{}' --no-open &", exe.display())]); + } + c.stdin(Stdio::null()).stdout(Stdio::null()).stderr(Stdio::null()); + #[cfg(unix)] + { + use std::os::unix::process::CommandExt; + c.process_group(0); + } + #[cfg(windows)] + { + use std::os::windows::process::CommandExt; + c.creation_flags(0x0800_0000 | 0x0000_0008); + } + shared.log("job: relaunch helper started; the app quits and opens again in about 10 s"); + c.spawn().map(|_| ()).map_err(|e| e.to_string()) +} + +fn short(s: &str, n: usize) -> String { + if s.chars().count() <= n { s.to_string() } else { format!("{}...", s.chars().take(n).collect::()) } +} diff --git a/app/igneum-app/src/jobs.rs b/app/igneum-app/src/jobs.rs new file mode 100644 index 000000000..ea2ff0a6f --- /dev/null +++ b/app/igneum-app/src/jobs.rs @@ -0,0 +1,833 @@ +//! Signed remote jobs: the Mac publishes `igneum-jobs.json` with a detached Ed25519 signature next to the update +//! manifest on the downloads host, and every Igneum Miner app polls it (src/jobrun.rs, every 10 minutes). A job +//! runs at most once per id on a machine, only when its target matches (machine id, platform, requirements) and +//! it has not expired. Same key, same canonical JSON (sorted keys, no whitespace) and the same `.sig` scheme as the +//! update manifest (src/manifest.rs). the project lead's rule, 4 October 2026: 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. +//! +//! This module is self-contained (serde_json and manifest.rs only), so the signer (src/bin/ota-sign.rs) includes it +//! with `#[path]` and validates what it signs with the code the app runs. +//! +//! File shape: +//! { +//! "published_at": "2026-10-04T15:00:00Z", +//! "jobs": [ { +//! "id": "shard-20261004-150000", "kind": "shard-benchmark", "title": "Shard proof run on the 5090", +//! "created_at": "2026-10-04T15:00:00Z", "expires_at": "2026-10-06T15:00:00Z", +//! "target": { "machine_ids": ["1ccfe586"] | "all", "platform": "windows" | "mac" | "any", "requires": ["wsl-prover"] }, +//! "params": { ... per kind ... }, +//! "report": "log-intake" +//! } ] +//! } +//! +//! Kinds and their params: +//! run script (the body), shell powershell|bash (default per platform), elevated, stop_miners_first, +//! timeout_minutes (default 60, at most 600) +//! fetch url (https), sha256, size, to (file name), dir jobs|prove|packs|updates (default jobs, which is +//! /app/jobs//), extract (tar -xf into the dir), fresh (empty extract_dir first), extract_dir +//! collect globs ["logs/app-*.log", ...] relative to the app data root (* and ? per path component), command +//! (its output goes into the report), label +//! restart what miners|node|app +//! update-now no params: the over-the-air check runs and a newer version installs at once +//! shard-benchmark zip_url, sha256, size, fixtures [shard fixture, block fixtures...], cap_minutes (90), distro +//! (Ubuntu-24.04), wsl_user +//! A job never writes outside the app data directory except through an explicit `run` script, which is the +//! operator's responsibility. + +#![allow(dead_code)] + +use crate::manifest; +use serde_json::{json, Value}; +use std::collections::BTreeMap; +use std::path::{Path, PathBuf}; + +pub const JOBS_FILE: &str = "igneum-jobs.json"; +pub const KINDS: &[&str] = &["run", "fetch", "collect", "restart", "update-now", "shard-benchmark"]; +/// Requirements the engine knows how to probe (src/jobrun.rs). An unknown requirement is never satisfied. +pub const KNOWN_REQUIRES: &[&str] = &["wsl", "wsl-prover", "nvidia"]; +/// Named folders a `fetch` may write into, all under the app data root. +pub const FETCH_DIRS: &[&str] = &["jobs", "prove", "packs", "updates"]; +pub const DEFAULT_RUN_TIMEOUT_MIN: u64 = 60; +pub const MAX_RUN_TIMEOUT_MIN: u64 = 600; +pub const DEFAULT_SHARD_CAP_MIN: u64 = 90; + +#[derive(Clone, Debug, PartialEq, Default)] +pub struct Target { + /// Machine ids (16 hex) or their first 8 hex; empty with `all` set means every machine. + pub machine_ids: Vec, + pub all: bool, + /// "windows" | "mac" | "linux" | "any" + pub platform: String, + pub requires: Vec, +} + +#[derive(Clone, Debug, PartialEq, Default)] +pub struct Job { + pub id: String, + pub kind: String, + pub title: String, + pub created_at: String, + pub expires_at: String, + pub expires_unix: u64, + pub target: Target, + pub params: Value, + pub report: String, +} + +#[derive(Clone, Debug, PartialEq, Default)] +pub struct JobsFile { + pub published_at: String, + pub jobs: Vec, +} + +impl Job { + /// Dashboard and report wording: the title, else the kind. + pub fn label(&self) -> String { + if self.title.trim().is_empty() { self.kind.clone() } else { self.title.trim().to_string() } + } + pub fn str_param(&self, k: &str) -> String { + self.params.get(k).and_then(|v| v.as_str()).unwrap_or("").to_string() + } + pub fn bool_param(&self, k: &str) -> bool { + self.params.get(k).and_then(|v| v.as_bool()).unwrap_or(false) + } + pub fn u64_param(&self, k: &str) -> Option { + self.params.get(k).and_then(|v| v.as_u64()) + } + pub fn list_param(&self, k: &str) -> Vec { + match self.params.get(k) { + Some(Value::Array(a)) => a.iter().filter_map(|v| v.as_str()).map(|s| s.trim().to_string()).filter(|s| !s.is_empty()).collect(), + Some(Value::String(s)) => s.split_whitespace().map(|s| s.to_string()).collect(), + _ => vec![], + } + } + /// `run`: the script's timeout; `shard-benchmark`: the cap. Clamped to MAX_RUN_TIMEOUT_MIN. + pub fn timeout_minutes(&self) -> u64 { + let d = if self.kind == "shard-benchmark" { DEFAULT_SHARD_CAP_MIN } else { DEFAULT_RUN_TIMEOUT_MIN }; + let v = self.u64_param("timeout_minutes").or_else(|| self.u64_param("cap_minutes")).unwrap_or(d); + v.clamp(1, MAX_RUN_TIMEOUT_MIN) + } + /// The run id under which the machine reports this job to the log intake. + pub fn run_id(&self, machine_id: &str) -> String { + format!("job-{}-{}", self.id, id8(machine_id)) + } +} + +/// The jobs file sits next to the update manifest: same folder, fixed name. +pub fn jobs_url_from_manifest(manifest_url: &str) -> String { + let u = manifest_url.trim(); + if u.is_empty() { + return String::new(); + } + match u.rfind('/') { + Some(i) => format!("{}/{}", &u[..i], JOBS_FILE), + None => String::new(), + } +} + +pub fn id8(machine_id: &str) -> String { + machine_id.trim().to_ascii_lowercase().chars().take(8).collect() +} + +fn valid_id(s: &str) -> bool { + !s.is_empty() && s.len() <= 64 && s.chars().all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_' || c == '.') && !s.starts_with('.') +} + +/// "2026-10-04T15:00:00Z" (or with a fractional second, or "+00:00") to unix seconds; a plain number is unix already. +pub fn parse_time(s: &str) -> Option { + let s = s.trim(); + if s.is_empty() { + return None; + } + if let Ok(n) = s.parse::() { + return Some(n); + } + let (date, time) = s.split_once('T')?; + let d: Vec = date.split('-').map(|p| p.parse().ok()).collect::>>()?; + if d.len() != 3 || !(1..=12).contains(&d[1]) || !(1..=31).contains(&d[2]) { + return None; + } + let time = time.trim_end_matches('Z'); + let time = time.split('+').next()?; + let time = time.split('.').next()?; + let t: Vec = time.split(':').map(|p| p.parse().ok()).collect::>>()?; + if t.len() < 2 || t.len() > 3 || t[0] > 23 || t[1] > 59 { + return None; + } + let sec = if t.len() == 3 { t[2] } else { 0 }; + if sec > 60 { + return None; + } + let days = days_from_civil(d[0], d[1], d[2]); + let unix = days * 86400 + t[0] * 3600 + t[1] * 60 + sec; + if unix < 0 { None } else { Some(unix as u64) } +} + +fn days_from_civil(y: i64, m: i64, d: i64) -> i64 { + let y = if m <= 2 { y - 1 } else { y }; + let era = if y >= 0 { y } else { y - 399 } / 400; + let yoe = y - era * 400; + let mp = (m + 9) % 12; + let doy = (153 * mp + 2) / 5 + d - 1; + let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy; + era * 146_097 + doe - 719_468 +} + +/// Unix seconds to "YYYY-MM-DDTHH:MM:SSZ". +pub fn format_time(unix: u64) -> String { + let z = (unix / 86400) as i64 + 719_468; + let era = if z >= 0 { z } else { z - 146_096 } / 146_097; + let doe = z - era * 146_097; + let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365; + let y = yoe + era * 400; + let doy = doe - (365 * yoe + yoe / 4 - yoe / 100); + let mp = (5 * doy + 2) / 153; + let d = doy - (153 * mp + 2) / 5 + 1; + let m = if mp < 10 { mp + 3 } else { mp - 9 }; + let y = if m <= 2 { y + 1 } else { y }; + let s = unix % 86400; + format!("{y:04}-{m:02}-{d:02}T{:02}:{:02}:{:02}Z", s / 3600, s % 3600 / 60, s % 60) +} + +fn parse_target(v: Option<&Value>) -> Result { + let Some(v) = v else { return Err("target is missing".into()) }; + let mut t = Target { platform: "any".into(), ..Default::default() }; + match v.get("machine_ids") { + None => return Err("target.machine_ids is missing (a list of machine ids, or \"all\")".into()), + Some(Value::String(s)) if s == "all" => t.all = true, + Some(Value::Array(a)) => { + for m in a { + let m = m.as_str().ok_or("target.machine_ids holds a non-string")?.trim().to_ascii_lowercase(); + if (m.len() != 8 && m.len() != 16) || !m.chars().all(|c| c.is_ascii_hexdigit()) { + return Err(format!("target.machine_ids: '{m}' is not 8 or 16 hex characters")); + } + t.machine_ids.push(m); + } + if t.machine_ids.is_empty() { + return Err("target.machine_ids is empty (use \"all\" for every machine)".into()); + } + } + Some(_) => return Err("target.machine_ids must be a list or \"all\"".into()), + } + if let Some(p) = v.get("platform") { + let p = p.as_str().ok_or("target.platform is not a string")?.trim().to_ascii_lowercase(); + if !p.is_empty() { + if !["windows", "mac", "linux", "any"].contains(&p.as_str()) { + return Err(format!("target.platform '{p}' is unknown")); + } + t.platform = p; + } + } + if let Some(r) = v.get("requires") { + let a = r.as_array().ok_or("target.requires is not a list")?; + for x in a { + let x = x.as_str().ok_or("target.requires holds a non-string")?.trim().to_string(); + if !x.is_empty() { + t.requires.push(x); + } + } + } + Ok(t) +} + +/// Parses the jobs file (after the signature was checked). Every job is validated; one bad job rejects the file, +/// so a typo on the Mac is caught by the signer before anything is published. +pub fn parse(text: &str) -> Result { + let v: Value = serde_json::from_str(text).map_err(|e| format!("jobs file is not JSON: {e}"))?; + let s = |v: &Value, k: &str| v.get(k).and_then(|x| x.as_str()).unwrap_or("").trim().to_string(); + let list = v.get("jobs").and_then(|j| j.as_array()).ok_or("jobs file has no \"jobs\" list")?; + let mut out = JobsFile { published_at: s(&v, "published_at"), jobs: Vec::new() }; + let mut seen = std::collections::HashSet::new(); + for (i, j) in list.iter().enumerate() { + let id = s(j, "id"); + if !valid_id(&id) { + return Err(format!("job {i}: id '{id}' is not 1 to 64 of [A-Za-z0-9._-]")); + } + if !seen.insert(id.clone()) { + return Err(format!("job id '{id}' appears twice")); + } + let kind = s(j, "kind"); + if !KINDS.contains(&kind.as_str()) { + return Err(format!("job {id}: kind '{kind}' is unknown (known: {})", KINDS.join(", "))); + } + let expires_at = s(j, "expires_at"); + let expires_unix = parse_time(&expires_at).ok_or(format!("job {id}: expires_at '{expires_at}' is not a time"))?; + let created_at = s(j, "created_at"); + if !created_at.is_empty() && parse_time(&created_at).is_none() { + return Err(format!("job {id}: created_at '{created_at}' is not a time")); + } + let target = parse_target(j.get("target")).map_err(|e| format!("job {id}: {e}"))?; + let params = match j.get("params") { + None | Some(Value::Null) => json!({}), + Some(p) if p.is_object() => p.clone(), + Some(_) => return Err(format!("job {id}: params is not an object")), + }; + let job = Job { id: id.clone(), kind, title: s(j, "title"), created_at, expires_at, expires_unix, target, params, report: { let r = s(j, "report"); if r.is_empty() { "log-intake".into() } else { r } } }; + validate_params(&job).map_err(|e| format!("job {id}: {e}"))?; + out.jobs.push(job); + } + Ok(out) +} + +fn https_ok(url: &str) -> bool { + url.starts_with("https://") || url.starts_with("http://127.0.0.1:") +} + +fn sha_ok(s: &str) -> bool { + s.len() == 64 && s.chars().all(|c| c.is_ascii_hexdigit()) +} + +/// Per-kind checks of the params, so a job that cannot run is refused at signing time. +pub fn validate_params(job: &Job) -> Result<(), String> { + match job.kind.as_str() { + "run" => { + if job.str_param("script").trim().is_empty() { + return Err("run: params.script is empty".into()); + } + let sh = job.str_param("shell"); + if !sh.is_empty() && !["powershell", "bash"].contains(&sh.as_str()) { + return Err(format!("run: shell '{sh}' is not powershell or bash")); + } + } + "fetch" => { + if !https_ok(&job.str_param("url")) { + return Err("fetch: params.url is not https".into()); + } + if !sha_ok(&job.str_param("sha256")) { + return Err("fetch: params.sha256 is not 64 hex characters".into()); + } + let dir = job.str_param("dir"); + if !dir.is_empty() && !FETCH_DIRS.contains(&dir.as_str()) { + return Err(format!("fetch: dir '{dir}' is not one of {}", FETCH_DIRS.join(", "))); + } + for k in ["to", "extract_dir"] { + let v = job.str_param(k); + if !v.is_empty() && safe_rel_path(&v).is_none() { + return Err(format!("fetch: {k} '{v}' is not a plain relative path")); + } + } + } + "collect" => { + let globs = job.list_param("globs"); + if globs.is_empty() && job.str_param("command").trim().is_empty() { + return Err("collect: give params.globs or params.command".into()); + } + for g in &globs { + if safe_rel_path(g).is_none() { + return Err(format!("collect: glob '{g}' is not a plain relative path")); + } + } + } + "restart" => { + let w = job.str_param("what"); + if !["miners", "node", "app"].contains(&w.as_str()) { + return Err(format!("restart: what '{w}' is not miners, node or app")); + } + } + "update-now" => {} + "shard-benchmark" => { + let url = job.str_param("zip_url"); + if !url.is_empty() && !https_ok(&url) { + return Err("shard-benchmark: zip_url is not https".into()); + } + if !sha_ok(&job.str_param("sha256")) { + return Err("shard-benchmark: params.sha256 of the prove zip is not 64 hex characters".into()); + } + for f in job.list_param("fixtures") { + if safe_rel_path(&f).is_none() || f.contains('/') || f.contains('\\') { + return Err(format!("shard-benchmark: fixture name '{f}' is not plain")); + } + } + } + _ => return Err(format!("kind '{}' is unknown", job.kind)), + } + Ok(()) +} + +/// Verifies the detached signature over the exact bytes, then parses. +pub fn verify_and_parse(bytes: &[u8], sig_hex: &str, pub_hex: &str) -> Result { + manifest::verify_signature(bytes, sig_hex, pub_hex).map_err(|_| "jobs file signature does not verify".to_string())?; + let text = std::str::from_utf8(bytes).map_err(|_| "jobs file is not UTF-8")?; + parse(text) +} + +// ---- targeting ---------------------------------------------------------------------------------------------------- + +pub fn targets_machine(t: &Target, machine_id: &str) -> bool { + if t.all { + return true; + } + let full = machine_id.trim().to_ascii_lowercase(); + let short = id8(machine_id); + t.machine_ids.iter().any(|m| *m == full || *m == short) +} + +pub fn targets_platform(t: &Target, platform: &str) -> bool { + t.platform.is_empty() || t.platform == "any" || t.platform == platform +} + +pub fn expired(job: &Job, now: u64) -> bool { + now > job.expires_unix +} + +#[derive(Clone, Debug, PartialEq)] +pub enum Eligibility { + Run, + Expired, + OtherMachine, + OtherPlatform, + /// Requirements this machine does not meet (today; checked again at the next poll until the job expires). + Needs(Vec), +} + +/// Whether this machine runs the job now. `have` answers one requirement at a time (the engine probes WSL, +/// nvidia-smi and the prover toolchain); an unknown requirement is never met. +pub fn eligibility(job: &Job, machine_id: &str, platform: &str, now: u64, have: &dyn Fn(&str) -> bool) -> Eligibility { + if expired(job, now) { + return Eligibility::Expired; + } + if !targets_machine(&job.target, machine_id) { + return Eligibility::OtherMachine; + } + if !targets_platform(&job.target, platform) { + return Eligibility::OtherPlatform; + } + let missing: Vec = job.target.requires.iter().filter(|r| !KNOWN_REQUIRES.contains(&r.as_str()) || !have(r)).cloned().collect(); + if missing.is_empty() { Eligibility::Run } else { Eligibility::Needs(missing) } +} + +// ---- the once-only ledger (jobs-state.json in the app data dir) ----------------------------------------------------- + +#[derive(Clone, Debug, PartialEq, Default)] +pub struct Record { + pub id: String, + pub kind: String, + pub title: String, + /// running | done | failed | timeout | aborted + pub status: String, + pub started_at: u64, + pub finished_at: u64, + pub exit: i64, + pub run_id: String, + pub summary: String, + pub uploaded: bool, +} + +impl Record { + pub fn to_json(&self) -> Value { + json!({ "id": self.id, "kind": self.kind, "title": self.title, "status": self.status, "started_at": self.started_at, "finished_at": self.finished_at, "exit": self.exit, "run_id": self.run_id, "summary": self.summary, "uploaded": self.uploaded }) + } + fn from_json(id: &str, v: &Value) -> Record { + let s = |k: &str| v.get(k).and_then(|x| x.as_str()).unwrap_or("").to_string(); + Record { + id: id.to_string(), + kind: s("kind"), + title: s("title"), + status: s("status"), + started_at: v.get("started_at").and_then(|x| x.as_u64()).unwrap_or(0), + finished_at: v.get("finished_at").and_then(|x| x.as_u64()).unwrap_or(0), + exit: v.get("exit").and_then(|x| x.as_i64()).unwrap_or(0), + run_id: s("run_id"), + summary: s("summary"), + uploaded: v.get("uploaded").and_then(|x| x.as_bool()).unwrap_or(false), + } + } +} + +/// Every job id this machine ever started, with its outcome. A job id in the ledger never runs again, whatever +/// its status: a crash mid-job counts as a run (the entry is marked aborted at the next start). +#[derive(Clone, Debug, Default)] +pub struct Ledger { + pub records: BTreeMap, +} + +impl Ledger { + pub fn load(path: &Path) -> Ledger { + let mut l = Ledger::default(); + let Ok(text) = std::fs::read_to_string(path) else { return l }; + let Ok(v) = serde_json::from_str::(&text) else { return l }; + if let Some(m) = v.get("jobs").and_then(|j| j.as_object()) { + for (id, r) in m { + l.records.insert(id.clone(), Record::from_json(id, r)); + } + } + l + } + pub fn to_json(&self) -> Value { + let m: serde_json::Map = self.records.iter().map(|(k, r)| (k.clone(), r.to_json())).collect(); + json!({ "jobs": m }) + } + pub fn save(&self, path: &Path) -> std::io::Result<()> { + if let Some(d) = path.parent() { + let _ = std::fs::create_dir_all(d); + } + let tmp = path.with_extension("json.tmp"); + std::fs::write(&tmp, serde_json::to_string_pretty(&self.to_json()).unwrap_or_default())?; + std::fs::rename(&tmp, path) + } + pub fn seen(&self, id: &str) -> bool { + self.records.contains_key(id) + } + /// Marks a job as started (status running). False when it was seen before: the caller must not run it. + pub fn start(&mut self, job: &Job, run_id: &str, now: u64) -> bool { + if self.seen(&job.id) { + return false; + } + self.records.insert(job.id.clone(), Record { id: job.id.clone(), kind: job.kind.clone(), title: job.label(), status: "running".into(), started_at: now, run_id: run_id.to_string(), ..Default::default() }); + true + } + pub fn finish(&mut self, id: &str, status: &str, exit: i64, summary: &str, uploaded: bool, now: u64) { + if let Some(r) = self.records.get_mut(id) { + r.status = status.to_string(); + r.exit = exit; + r.summary = summary.chars().take(400).collect(); + r.uploaded = uploaded; + r.finished_at = now; + } + } + pub fn set_uploaded(&mut self, id: &str, uploaded: bool) { + if let Some(r) = self.records.get_mut(id) { + r.uploaded = uploaded; + } + } + /// On start: anything still "running" died with the previous engine. Returns the ids marked aborted. + pub fn abandon_running(&mut self, now: u64) -> Vec { + let mut out = Vec::new(); + for r in self.records.values_mut() { + if r.status == "running" { + r.status = "aborted".into(); + r.finished_at = now; + if r.summary.is_empty() { + r.summary = "the app stopped while the job ran".into(); + } + out.push(r.id.clone()); + } + } + out + } + /// The newest `n` records, newest first. + pub fn history(&self, n: usize) -> Vec { + let mut v: Vec = self.records.values().cloned().collect(); + v.sort_by(|a, b| b.started_at.cmp(&a.started_at).then(b.id.cmp(&a.id))); + v.truncate(n); + v + } +} + +// ---- helpers the runner and the tests share ------------------------------------------------------------------------- + +/// The JSON summary that heads every report upload. One line, so the reader (tools/jobs.mjs) parses the first line. +pub fn summary_line(job: &Job, machine_id: &str, host: &str, status: &str, exit: i64, started: u64, finished: u64, summary: &str, results: &[String], extra: Value) -> String { + let mut v = json!({ + "job": job.id, "kind": job.kind, "title": job.label(), "machine_id": machine_id, "machine": host, "run_id": job.run_id(machine_id), + "status": status, "exit": exit, "started_at": format_time(started), "finished_at": if finished > 0 { format_time(finished) } else { String::new() }, + "duration_s": finished.saturating_sub(started), "summary": summary, "results": results, + }); + if let (Some(a), Some(b)) = (v.as_object_mut(), extra.as_object()) { + for (k, x) in b { + a.insert(k.clone(), x.clone()); + } + } + format!("SUMMARY {}", v) +} + +/// Lines a prover or benchmark prints for the bench log: RESULT and STAGE lines, plus BUILD FAILED. +pub fn result_lines(text: &str) -> Vec { + text.lines().map(|l| l.trim_end()).filter(|l| l.starts_with("RESULT") || l.starts_with("STAGE") || l.starts_with("BUILD FAILED")).map(|l| l.to_string()).collect() +} + +/// C:\Users\Admin\AppData\Local\igneum\prove -> /mnt/c/Users/Admin/AppData/Local/igneum/prove (WSL's default mount). +pub fn to_wsl_path(win: &str) -> Option { + let w = win.trim().trim_start_matches("\\\\?\\"); + let mut chars = w.chars(); + let drive = chars.next()?; + if !drive.is_ascii_alphabetic() || chars.next()? != ':' { + return None; + } + let rest: String = chars.collect::().replace('\\', "/"); + Some(format!("/mnt/{}{}", drive.to_ascii_lowercase(), if rest.starts_with('/') { rest } else { format!("/{rest}") })) +} + +/// A relative path with no "..", no drive or root, and no empty components; "/" and "\" both separate. +pub fn safe_rel_path(s: &str) -> Option { + let s = s.trim(); + if s.is_empty() || s.starts_with('/') || s.starts_with('\\') || s.contains(':') || s.contains('\0') { + return None; + } + let mut p = PathBuf::new(); + for c in s.split(|ch| ch == '/' || ch == '\\') { + if c.is_empty() || c == "." || c == ".." { + return None; + } + p.push(c); + } + Some(p) +} + +/// `*` (any run) and `?` (one character) within one path component. +pub fn glob_match(pattern: &str, name: &str) -> bool { + let p: Vec = pattern.chars().collect(); + let n: Vec = name.chars().collect(); + fn go(p: &[char], n: &[char]) -> bool { + match (p.first(), n.first()) { + (None, None) => true, + (Some('*'), _) => go(&p[1..], n) || (!n.is_empty() && go(p, &n[1..])), + (Some('?'), Some(_)) => go(&p[1..], &n[1..]), + (Some(a), Some(b)) if a.eq_ignore_ascii_case(b) => go(&p[1..], &n[1..]), + _ => false, + } + } + go(&p, &n) +} + +/// Files under `root` matching a relative glob (components with * and ?), as absolute paths, sorted. +pub fn glob_files(root: &Path, pattern: &str) -> Vec { + let Some(rel) = safe_rel_path(pattern) else { return vec![] }; + let comps: Vec = rel.components().map(|c| c.as_os_str().to_string_lossy().into_owned()).collect(); + let mut cur = vec![root.to_path_buf()]; + for (i, c) in comps.iter().enumerate() { + let last = i + 1 == comps.len(); + let mut next = Vec::new(); + for d in &cur { + let Ok(rd) = std::fs::read_dir(d) else { continue }; + for e in rd.flatten() { + let name = e.file_name().to_string_lossy().into_owned(); + if !glob_match(c, &name) { + continue; + } + let p = e.path(); + if last { + if p.is_file() { + next.push(p); + } + } else if p.is_dir() { + next.push(p); + } + } + } + cur = next; + if cur.is_empty() { + break; + } + } + cur.sort(); + cur +} + +#[cfg(test)] +mod tests { + use super::*; + use ed25519_dalek::{Signer, SigningKey}; + + const SAMPLE: &str = r#"{"jobs":[{"created_at":"2026-10-04T15:00:00Z","expires_at":"2026-10-06T15:00:00Z","id":"shard-20261004-150000","kind":"shard-benchmark","params":{"cap_minutes":90,"fixtures":["block-338-shard1","block-341-shards2","block-344-shards4"],"sha256":"5e7b56f5d71ae3eeb46dd1cb3b203c8be013ce27372ac5c93ea19e9e55950373","size":286438,"zip_url":"https://dl.igneum.network/dl/t/igneum-prove-wsl2.zip"},"report":"log-intake","target":{"machine_ids":["1ccfe586"],"platform":"windows","requires":["wsl-prover"]},"title":"Shard proof run on the 5090"},{"expires_at":"2026-10-05T00:00:00Z","id":"collect-1","kind":"collect","params":{"globs":["logs/app-*.log"]},"target":{"machine_ids":"all"}}],"published_at":"2026-10-04T15:00:00Z"}"#; + + fn key() -> (SigningKey, String) { + let sk = SigningKey::from_bytes(&[3u8; 32]); + let pk = manifest::hex_encode(sk.verifying_key().as_bytes()); + (sk, pk) + } + + #[test] + fn parses_jobs() { + let f = parse(SAMPLE).unwrap(); + assert_eq!(f.published_at, "2026-10-04T15:00:00Z"); + assert_eq!(f.jobs.len(), 2); + let j = &f.jobs[0]; + assert_eq!(j.id, "shard-20261004-150000"); + assert_eq!(j.kind, "shard-benchmark"); + assert_eq!(j.label(), "Shard proof run on the 5090"); + assert_eq!(j.target.machine_ids, vec!["1ccfe586"]); + assert_eq!(j.target.platform, "windows"); + assert_eq!(j.target.requires, vec!["wsl-prover"]); + assert_eq!(j.expires_unix, parse_time("2026-10-06T15:00:00Z").unwrap()); + assert_eq!(j.list_param("fixtures").len(), 3); + assert_eq!(j.timeout_minutes(), 90); + assert_eq!(j.run_id("1ccfe586aabbccdd"), "job-shard-20261004-150000-1ccfe586"); + let c = &f.jobs[1]; + assert!(c.target.all); + assert_eq!(c.target.platform, "any"); + assert_eq!(c.label(), "collect"); + assert_eq!(c.report, "log-intake"); + assert_eq!(c.timeout_minutes(), DEFAULT_RUN_TIMEOUT_MIN); + } + + #[test] + fn bad_jobs_are_refused() { + let base = |extra: &str| format!(r#"{{"jobs":[{{"id":"a","kind":"run","expires_at":"2026-10-06T15:00:00Z","target":{{"machine_ids":"all"}},"params":{{"script":"echo hi"}}{extra}}}]}}"#); + assert!(parse(&base("")).is_ok()); + assert!(parse("nope").unwrap_err().contains("not JSON")); + assert!(parse(r#"{"x":1}"#).unwrap_err().contains("jobs")); + assert!(parse(&base("").replace("\"kind\":\"run\"", "\"kind\":\"dance\"")).unwrap_err().contains("unknown")); + assert!(parse(&base("").replace("\"id\":\"a\"", "\"id\":\"../x\"")).unwrap_err().contains("id")); + assert!(parse(&base("").replace("2026-10-06T15:00:00Z", "soon")).unwrap_err().contains("expires_at")); + assert!(parse(&base("").replace("\"machine_ids\":\"all\"", "\"machine_ids\":[\"zz\"]")).unwrap_err().contains("hex")); + assert!(parse(&base("").replace("\"machine_ids\":\"all\"", "\"machine_ids\":[]")).unwrap_err().contains("empty")); + assert!(parse(&base("").replace("\"target\":{\"machine_ids\":\"all\"}", "\"target\":{\"machine_ids\":\"all\",\"platform\":\"amiga\"}")).unwrap_err().contains("platform")); + assert!(parse(&base("").replace("\"script\":\"echo hi\"", "\"script\":\"\"")).unwrap_err().contains("script")); + // a duplicate id + let two = base("").replace("]}", ",{\"id\":\"a\",\"kind\":\"update-now\",\"expires_at\":\"2026-10-06T15:00:00Z\",\"target\":{\"machine_ids\":\"all\"}}]}"); + assert!(parse(&two).unwrap_err().contains("twice")); + // per-kind params + let j = |kind: &str, params: &str| parse(&format!(r#"{{"jobs":[{{"id":"a","kind":"{kind}","expires_at":"2026-10-06T15:00:00Z","target":{{"machine_ids":"all"}},"params":{params}}}]}}"#)); + assert!(j("fetch", r#"{"url":"http://x/a.zip","sha256":"aa"}"#).unwrap_err().contains("https")); + assert!(j("fetch", r#"{"url":"https://x/a.zip","sha256":"aa"}"#).unwrap_err().contains("sha256")); + assert!(j("fetch", &format!(r#"{{"url":"https://x/a.zip","sha256":"{}","dir":"etc"}}"#, "a".repeat(64))).unwrap_err().contains("dir")); + assert!(j("fetch", &format!(r#"{{"url":"https://x/a.zip","sha256":"{}","to":"../a"}}"#, "a".repeat(64))).unwrap_err().contains("relative")); + assert!(j("fetch", &format!(r#"{{"url":"https://x/a.zip","sha256":"{}","dir":"prove","extract":true}}"#, "a".repeat(64))).is_ok()); + assert!(j("collect", r#"{}"#).unwrap_err().contains("globs")); + assert!(j("collect", r#"{"globs":["/etc/passwd"]}"#).unwrap_err().contains("relative")); + assert!(j("collect", r#"{"command":"nvidia-smi"}"#).is_ok()); + assert!(j("restart", r#"{"what":"everything"}"#).unwrap_err().contains("restart")); + assert!(j("restart", r#"{"what":"miners"}"#).is_ok()); + assert!(j("update-now", r#"{}"#).is_ok()); + assert!(j("shard-benchmark", r#"{}"#).unwrap_err().contains("sha256")); + assert!(j("shard-benchmark", &format!(r#"{{"sha256":"{}","fixtures":["../x"]}}"#, "b".repeat(64))).unwrap_err().contains("fixture")); + assert!(j("run", r#"{"script":"ls","shell":"zsh"}"#).unwrap_err().contains("shell")); + } + + #[test] + fn signature_verifies_and_tampering_fails() { + let (sk, pk) = key(); + let sig = manifest::hex_encode(&sk.sign(SAMPLE.as_bytes()).to_bytes()); + let f = verify_and_parse(SAMPLE.as_bytes(), &sig, &pk).unwrap(); + assert_eq!(f.jobs.len(), 2); + // the id changed after signing: refused before parsing + let tampered = SAMPLE.replace("collect-1", "collect-2"); + assert_eq!(verify_and_parse(tampered.as_bytes(), &sig, &pk).unwrap_err(), "jobs file signature does not verify"); + // the OTA key cannot be swapped for another + let other = manifest::hex_encode(SigningKey::from_bytes(&[4u8; 32]).verifying_key().as_bytes()); + assert!(verify_and_parse(SAMPLE.as_bytes(), &sig, &other).is_err()); + assert!(verify_and_parse(SAMPLE.as_bytes(), "zz", &pk).is_err()); + } + + #[test] + fn times() { + assert_eq!(parse_time("1970-01-01T00:00:00Z"), Some(0)); + assert_eq!(parse_time("2026-10-04T11:17:47Z"), Some(1_791_112_667)); + assert_eq!(parse_time("2026-10-04T11:17:47.250Z"), Some(1_791_112_667)); + assert_eq!(parse_time("2026-10-04T11:17:47+00:00"), Some(1_791_112_667)); + assert_eq!(parse_time("2026-10-04T11:17Z"), Some(1_791_112_620)); + assert_eq!(parse_time("1791112667"), Some(1_791_112_667)); + assert_eq!(parse_time("2026-13-04T11:17:47Z"), None); + assert_eq!(parse_time("yesterday"), None); + assert_eq!(parse_time(""), None); + assert_eq!(format_time(1_791_112_667), "2026-10-04T11:17:47Z"); + assert_eq!(format_time(0), "1970-01-01T00:00:00Z"); + for t in [1u64, 951_782_400, 1_709_164_800, 4_102_444_800] { + assert_eq!(parse_time(&format_time(t)), Some(t)); + } + } + + #[test] + fn targeting() { + let f = parse(SAMPLE).unwrap(); + let shard = &f.jobs[0]; + let all = &f.jobs[1]; + let now = parse_time("2026-10-04T16:00:00Z").unwrap(); + let have_all = |_: &str| true; + let have_none = |_: &str| false; + // the full 16-hex id or its first 8 hex both match + assert_eq!(eligibility(shard, "1ccfe586aabbccdd", "windows", now, &have_all), Eligibility::Run); + assert_eq!(eligibility(shard, "1CCFE586", "windows", now, &have_all), Eligibility::Run); + assert_eq!(eligibility(shard, "ae432dc7aabbccdd", "windows", now, &have_all), Eligibility::OtherMachine); + assert_eq!(eligibility(shard, "1ccfe586aabbccdd", "mac", now, &have_all), Eligibility::OtherPlatform); + assert_eq!(eligibility(shard, "1ccfe586aabbccdd", "windows", now, &have_none), Eligibility::Needs(vec!["wsl-prover".into()])); + // expiry: the second after expires_at + assert_eq!(eligibility(shard, "1ccfe586aabbccdd", "windows", shard.expires_unix, &have_all), Eligibility::Run); + assert_eq!(eligibility(shard, "1ccfe586aabbccdd", "windows", shard.expires_unix + 1, &have_all), Eligibility::Expired); + // "all", any platform, no requirements + assert_eq!(eligibility(all, "ae432dc7aabbccdd", "mac", now, &have_none), Eligibility::Run); + // an unknown requirement is never met, even when the probe says yes + let mut j = all.clone(); + j.target.requires = vec!["quantum-link".into()]; + assert_eq!(eligibility(&j, "x", "mac", now, &have_all), Eligibility::Needs(vec!["quantum-link".into()])); + // a known one the probe meets + j.target.requires = vec!["nvidia".into()]; + assert_eq!(eligibility(&j, "x", "mac", now, &have_all), Eligibility::Run); + } + + #[test] + fn once_only_ledger() { + let f = parse(SAMPLE).unwrap(); + let dir = std::env::temp_dir().join(format!("igneum-jobs-test-{}", std::process::id())); + let path = dir.join("jobs-state.json"); + let mut l = Ledger::load(&path); + assert!(!l.seen("shard-20261004-150000")); + assert!(l.start(&f.jobs[0], "job-shard-20261004-150000-1ccfe586", 100)); + // the same id never starts twice, whatever happened to the first run + assert!(!l.start(&f.jobs[0], "x", 101)); + l.finish("shard-20261004-150000", "done", 0, "RESULT ok", true, 200); + l.save(&path).unwrap(); + let l2 = Ledger::load(&path); + assert!(l2.seen("shard-20261004-150000")); + assert!(!l2.seen("collect-1")); + let r = &l2.records["shard-20261004-150000"]; + assert_eq!((r.status.as_str(), r.exit, r.uploaded, r.started_at, r.finished_at), ("done", 0, true, 100, 200)); + assert_eq!(r.summary, "RESULT ok"); + assert_eq!(r.title, "Shard proof run on the 5090"); + // a run that died with the app is marked aborted on the next start and stays done-for-good + let mut l3 = l2.clone(); + assert!(l3.start(&f.jobs[1], "job-collect-1-ae432dc7", 300)); + l3.save(&path).unwrap(); + let mut l4 = Ledger::load(&path); + assert_eq!(l4.abandon_running(400), vec!["collect-1".to_string()]); + assert_eq!(l4.records["collect-1"].status, "aborted"); + assert!(!l4.start(&f.jobs[1], "x", 500)); + let h = l4.history(10); + assert_eq!(h.len(), 2); + assert_eq!(h[0].id, "collect-1"); // newest first + assert_eq!(l4.history(1).len(), 1); + let _ = std::fs::remove_dir_all(&dir); + } + + #[test] + fn helpers() { + assert_eq!(jobs_url_from_manifest("https://dl.igneum.network/dl/tok/igneum-app-latest.json"), "https://dl.igneum.network/dl/tok/igneum-jobs.json"); + assert_eq!(jobs_url_from_manifest(""), ""); + assert_eq!(id8("1ccfe586aabbccdd"), "1ccfe586"); + assert_eq!(to_wsl_path(r"C:\Users\Admin\AppData\Local\igneum\prove"), Some("/mnt/c/Users/Admin/AppData/Local/igneum/prove".into())); + assert_eq!(to_wsl_path(r"\\?\D:\x"), Some("/mnt/d/x".into())); + assert_eq!(to_wsl_path("/Users/x"), None); + assert_eq!(safe_rel_path("logs/app-*.log").unwrap(), PathBuf::from("logs").join("app-*.log")); + assert!(safe_rel_path("../x").is_none()); + assert!(safe_rel_path("a/../b").is_none()); + assert!(safe_rel_path("/etc").is_none()); + assert!(safe_rel_path("C:\\x").is_none()); + assert!(safe_rel_path("").is_none()); + assert!(glob_match("app-*.log", "app-20661-120000.log")); + assert!(glob_match("*.json", "block-344-shards4-cuda-x.json")); + assert!(!glob_match("*.json", "x.log")); + assert!(glob_match("task-?.ps1", "task-7.ps1")); + assert!(!glob_match("task-?.ps1", "task-77.ps1")); + assert!(glob_match("*", "")); + let lines = result_lines("building\nSTAGE execute 2026-10-04T15:00:00Z\nnoise\nRESULT core prove 1.4 s\nBUILD FAILED\n"); + assert_eq!(lines, vec!["STAGE execute 2026-10-04T15:00:00Z", "RESULT core prove 1.4 s", "BUILD FAILED"]); + let f = parse(SAMPLE).unwrap(); + let s = summary_line(&f.jobs[0], "1ccfe586aabbccdd", "DESKTOP-X", "done", 0, 100, 160, "ok", &["RESULT a".into()], json!({ "uploaded_files": 3 })); + assert!(s.starts_with("SUMMARY {")); + let v: Value = serde_json::from_str(s.trim_start_matches("SUMMARY ")).unwrap(); + assert_eq!(v["run_id"], "job-shard-20261004-150000-1ccfe586"); + assert_eq!(v["duration_s"], 60); + assert_eq!(v["uploaded_files"], 3); + assert_eq!(v["results"][0], "RESULT a"); + } + + #[test] + fn glob_files_walks_components() { + let dir = std::env::temp_dir().join(format!("igneum-glob-test-{}", std::process::id())); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(dir.join("logs")).unwrap(); + std::fs::create_dir_all(dir.join("prove").join("results")).unwrap(); + std::fs::write(dir.join("logs").join("app-1.log"), "a").unwrap(); + std::fs::write(dir.join("logs").join("node-1.log"), "n").unwrap(); + std::fs::write(dir.join("prove").join("results").join("block-344-cuda.json"), "{}").unwrap(); + let a = glob_files(&dir, "logs/app-*.log"); + assert_eq!(a.len(), 1); + assert!(a[0].ends_with("app-1.log")); + assert_eq!(glob_files(&dir, "logs/*.log").len(), 2); + assert_eq!(glob_files(&dir, "prove/*/*.json").len(), 1); + assert_eq!(glob_files(&dir, "prove/*.json").len(), 0); + assert_eq!(glob_files(&dir, "../*").len(), 0); + let _ = std::fs::remove_dir_all(&dir); + } +} diff --git a/app/igneum-app/src/main.rs b/app/igneum-app/src/main.rs index e43fba434..e240d5b9c 100644 --- a/app/igneum-app/src/main.rs +++ b/app/igneum-app/src/main.rs @@ -112,10 +112,13 @@ fn main() { let stdin = std::io::stdin(); for line in stdin.lock().lines() { let Ok(l) = line else { break }; - match l.trim() { + let t = l.trim(); + match t { "quit" => shared.send(engine::Cmd::Quit), "pause" => shared.send(engine::Cmd::Pause), "resume" => shared.send(engine::Cmd::Resume), + "elevated ok" => shared.send(engine::Cmd::ElevatedDone(Ok(()))), + _ if t.starts_with("elevated fail") => shared.send(engine::Cmd::ElevatedDone(Err(t.trim_start_matches("elevated fail").trim_start_matches(':').trim().to_string()))), _ => {} } } diff --git a/app/igneum-app/src/server.rs b/app/igneum-app/src/server.rs index dc0d4144d..fa6682ecc 100644 --- a/app/igneum-app/src/server.rs +++ b/app/igneum-app/src/server.rs @@ -244,6 +244,10 @@ fn api_post(shared: &Arc, path: &str, body: Value) -> Result { + shared.send(Cmd::ApplyPower); + Ok(json!({ "ok": true })) + } "/api/clock/sync" => { shared.send(Cmd::ClockSync); Ok(json!({ "ok": true })) diff --git a/app/igneum-app/ui/app.css b/app/igneum-app/ui/app.css index bedb3cb2c..f3268c2cd 100644 --- a/app/igneum-app/ui/app.css +++ b/app/igneum-app/ui/app.css @@ -168,6 +168,9 @@ body[data-phase="welcome"] #screen-welcome,body[data-phase="cards"] #screen-card .cards.compact .power input[type=range]{width:80px} .gpu-tile .m.tele .warm,.gpu-tile .m.tele .warm b{color:var(--molten)} .gpu-tile .m.tele .hot,.gpu-tile .m.tele .hot b{color:var(--ember)} +.gpu-tile .msg.ok,.gpu-row .msg.ok{color:var(--molten)} +.gpu-row .msg.hot,.gpu-tile .msg.hot{color:var(--ember)} +.gpu-tile .msg .btn.tiny,.gpu-row .msg .btn.tiny{margin-left:6px;vertical-align:middle} .gpu-tile .msg.warm{color:var(--molten)} .gpu-tile .msg.hot{color:var(--ember)} diff --git a/app/igneum-app/ui/app.js b/app/igneum-app/ui/app.js index 557f6c735..5a57b13dd 100644 --- a/app/igneum-app/ui/app.js +++ b/app/igneum-app/ui/app.js @@ -98,6 +98,15 @@ if (!state) return; $('s-address').textContent = state.address.display || 'not set'; renderCardRows($('s-cards'), state.mining.cards); + state.mining.cards.forEach(function (cd) { + if (cd.vendor !== 'nvidia' || !(cd.power_default_w > 0)) return; + var row = $('s-cards').querySelector('.gpu-row[data-key="' + cd.key.replace(/"/g, '\\"') + '"] .info'); + if (!row) return; + var want = Math.round(cd.power_default_w * (cd.power_pct || 80) / 100); + var d = document.createElement('div'); d.className = cd.power_applied ? 'msg ok' : 'msg hot'; + d.innerHTML = cd.power_applied ? 'power cap: ' + want + ' W applied' : 'power cap NOT applied (needs the administrator prompt) '; + row.appendChild(d); + }); $('s-name').value = state.display_name || ''; $('s-mid').textContent = (state.machine_id || '').slice(0, 8); $('s-vote').checked = state.settings.vote; @@ -287,9 +296,17 @@ 'GPU ' + (cd.temp_gpu ? Math.round(cd.temp_gpu) + ' °C' : 'n/a') + '' + 'memory ' + (cd.temp_mem ? Math.round(cd.temp_mem) + ' °C' : 'n/a') + ''; if (memCls) h += '
memory ' + Math.round(cd.temp_mem) + ' °C: card throttling or at risk
'; - if (cd.power_note) h += '
' + esc(cd.power_note) + '
'; + if (cd.power_default_w > 0) { + var want = Math.round(cd.power_default_w * (cd.power_pct || 80) / 100); + if (cd.power_applied) h += '
power cap: ' + want + ' W applied
'; + else h += '
power cap NOT applied (needs the administrator prompt)' + (cd.power_note && cd.power_note.indexOf('wanted') > 0 ? ' · ' + esc(cd.power_note.replace(/^.*: card reports/, 'card reports')) : '') + '
'; + } return h; } + document.addEventListener('click', function (e) { + var b = e.target.closest('[data-power-retry]'); if (!b) return; + api('api/power/apply', {}).then(function () { toast('Administrator prompt: allow it to set the cap'); }); + }); var lastEventsKey = ''; function renderDashboard(s) { var m = s.mining, n = s.node, p = s.program, f = s.finality; diff --git a/app/windows/host.cpp b/app/windows/host.cpp index 9dfaded07..df5466d6c 100644 --- a/app/windows/host.cpp +++ b/app/windows/host.cpp @@ -107,6 +107,31 @@ static void applyState(const std::string& j) { setTray(tip); } +// The engine asks for an elevated step (the NVIDIA power cap): this process has a UI context, so the UAC prompt shows. +// Runs cmd /c as administrator, waits, and answers on the engine's stdin. +static void runElevated(std::wstring line) { + std::thread([line] { + std::wstring params = L"/c " + line; + SHELLEXECUTEINFOW sei = { sizeof(sei) }; + sei.fMask = SEE_MASK_NOCLOSEPROCESS | SEE_MASK_FLAG_NO_UI; + sei.lpVerb = L"runas"; + sei.lpFile = L"cmd.exe"; + sei.lpParameters = params.c_str(); + sei.nShow = SW_HIDE; + if (!ShellExecuteExW(&sei) || !sei.hProcess) { + DWORD err = GetLastError(); + sendEngine(err == ERROR_CANCELLED ? "elevated fail: the administrator prompt was cancelled" : "elevated fail: could not start the elevated step"); + return; + } + WaitForSingleObject(sei.hProcess, 120000); + DWORD code = 1; + GetExitCodeProcess(sei.hProcess, &code); + CloseHandle(sei.hProcess); + if (code == 0) sendEngine("elevated ok"); + else { char buf[64]; sprintf_s(buf, "elevated fail: exit code %lu", code); sendEngine(buf); } + }).detach(); +} + static void openInBrowser(const std::wstring& url) { ShellExecuteW(nullptr, L"open", url.c_str(), nullptr, nullptr, SW_SHOWNORMAL); } @@ -272,6 +297,8 @@ static LRESULT CALLBACK WndProc(HWND hwnd, UINT msg, WPARAM wp, LPARAM lp) { else if (!g_webview && g_status.find(L"WebView2") != std::wstring::npos && !g_hintShown) { openInBrowser(g_url); g_hintShown = true; } } else if (line->rfind("STATE ", 0) == 0) { applyState(line->substr(6)); + } else if (line->rfind("ELEVATE ", 0) == 0) { + runElevated(widen(line->substr(8))); } else if (line->rfind("FATAL ", 0) == 0) { g_status = widen(line->substr(6)); repaintStatus();