//! Run orchestration — the new core (aldermon has no equivalent). Owns //! the sweep lifecycle: for each workload in the config, cooldown → //! workload.start() → wait for exit (polling is_running) → workload.stop() //! if needed → workload.wait() for verdict → workload.score() for the //! perf number → flush report. Front-end-agnostic: the TUI drives it by //! calling `tick()` every `poll_ms`; `--headless` does the same without //! rendering. //! //! State machine (poll-driven, no async): //! Cooldown(remaining) ── elapsed → Running(idx) //! Running(idx) ── workload exits OR duration elapsed → ReportFinished(idx+1) //! ReportFinished(idx) ── next idx exists → Cooldown(cool_secs) //! └─ no more workloads → Done #![allow(dead_code)] // consumed by main.rs + ui.rs — ui.rs not yet wired use std::io; use std::time::{Duration, Instant}; use crate::config::Config; use crate::report::Report; use crate::workload::{self, Score, Verdict, Workload}; /// What a tick produced. The caller (TUI or headless loop) uses this to /// decide what to render + when to stop driving the run. #[derive(Clone, Debug)] pub enum TickOutcome { /// A workload was just started this tick. `idx` is the workload index, /// `name` is its name, `duration_secs` is the planned run length. WorkloadStarted { idx: usize, name: String, duration_secs: u64, }, /// The workload exited (either it self-finished or the duration timer /// elapsed and we stopped it). The verdict + score + report path are /// included for the TUI's result table + stdout summary. WorkloadFinished { idx: usize, name: String, verdict: Verdict, score: Score, report_path: std::path::PathBuf, }, /// The sweep is complete (all workloads done, final cooldown elapsed). SweepDone, /// Nothing happened this tick (still in cooldown or still running and /// no exit was detected). The caller can just render the existing state. Idle, } /// Run state. The run loop (TUI or headless) owns one of these and calls /// `tick()` every `poll_ms`. pub struct Run { cfg: Config, /// All workloads that could be built from the config, in order. Names /// that `from_spec` couldn't wire (deferred impls, typos) are skipped /// during construction with a warning stored here so the TUI can show /// them. workloads: Vec>>, /// Skipped workload names + why (for the TUI / stdout warning). skipped: Vec<(String, String)>, /// Index into `workloads` for the current/next workload to run. current: usize, state: State, /// When the current state (Cooldown or Running) started. state_start: Instant, /// The report being written for the current workload. report: Option, /// How many loops are left, counting the current one. 1 = this is /// the last (or only) loop. 0 = infinite (loop forever). Decremented /// when a sweep completes + there are more loops to go; the run then /// resets `current` to 0 + re-enters Cooldown. When this hits 0 after /// the decrement (and it wasn't infinite), State::Done fires. loops_left: u32, /// 1-indexed ordinal of the current loop (for the TUI header). 1 on /// the first pass, 2 after one reset, etc. loop_index: u32, } enum State { /// Cooling down before starting workload `current`. `remaining` is the /// cooldown left in seconds (drives the TUI's cooldown display). Cooldown { remaining: u64 }, /// Running workload `current`. `duration_secs` is the planned length /// (the run loop stops the workload when elapsed >= this). Running { duration_secs: u64 }, /// All workloads done; the sweep is over. Further ticks return /// `SweepDone` indefinitely. Done, } impl Run { /// Build a run from a config. Constructs all workloads up front so /// start-time errors (missing binary, etc.) surface before the sweep /// begins rather than mid-sweep. Workloads that can't be built are /// skipped with a recorded warning. pub fn new(cfg: Config) -> Self { let mut workloads: Vec>> = Vec::new(); let mut skipped: Vec<(String, String)> = Vec::new(); for spec in &cfg.workloads { match workload::from_spec(spec, &cfg.bin_dir) { Some(w) => workloads.push(Some(w)), None => skipped.push(( spec.name.clone(), "unknown or not-yet-wired workload".to_string(), )), } } let loops_left = cfg.loop_count; Run { cfg, workloads, skipped, current: 0, state: State::Cooldown { remaining: 0 }, state_start: Instant::now(), report: None, loops_left, loop_index: 1, } } /// Skipped workload names + reasons (for the TUI / stdout warning). pub fn skipped(&self) -> &[(String, String)] { &self.skipped } /// The total number of workloads that will run (skipped ones excluded). pub fn total_workloads(&self) -> usize { self.workloads.iter().filter(|w| w.is_some()).count() } /// 1-indexed ordinal of the workload currently running or about to /// start (skipped entries excluded). 0 before the sweep begins. For /// the TUI header ("workload N of M"). pub fn current_ordinal(&self) -> usize { // Count Some entries up to + including `current`. self.workloads .iter() .take(self.current + 1) .filter(|w| w.is_some()) .count() } /// The current workload's spec (name, duration, cores), or None if /// the sweep hasn't started / is done. For the TUI's current-panel. pub fn current_spec(&self) -> Option<&crate::config::WorkloadSpec> { let w = self.workloads.get(self.current)?; if w.is_some() { self.cfg.workloads.get(self.current) } else { None } } /// Seconds elapsed in the current state (Cooldown or Running), as /// f64 for the TUI's progress display. pub fn elapsed_secs(&self) -> f64 { self.state_start.elapsed().as_secs_f64() } /// Remaining seconds in the current state. During Running this is /// `duration - elapsed` (clamped at 0). During Cooldown it's /// `remaining - elapsed` (clamped at 0). None when Done. Uses /// `saturating_sub` because Duration subtraction panics on overflow /// (happens when a workload overruns its duration — e.g. 7z which /// isn't time-bound + runs until the run loop's stop() kicks in). pub fn remaining_secs(&self) -> Option { let elapsed = self.state_start.elapsed(); match &self.state { State::Cooldown { remaining } => { let dur = Duration::from_secs(*remaining); Some(dur.saturating_sub(elapsed).as_secs_f64()) } State::Running { duration_secs } => { let dur = Duration::from_secs(*duration_secs); Some(dur.saturating_sub(elapsed).as_secs_f64()) } State::Done => None, } } /// True when the sweep is over (all workloads done). For the TUI to /// know when to stop ticking. pub fn is_done(&self) -> bool { matches!(self.state, State::Done) } /// The configured poll period, ms. The TUI uses this for its event /// poll timeout so the run + render stay in lockstep. pub fn poll_ms(&self) -> u64 { self.cfg.poll_ms } /// Poll period while a workload is RUNNING (not cooldown). Slower /// than poll_ms to minimize CPU contention with the workload on /// single-core runs. The TUI uses this instead of poll_ms when the /// state is Running. 500ms is fast enough to detect completion /// (y-cruncher + 7z exit within one tick of finishing) while /// stealing <0.1% of the workload's CPU. pub fn running_poll_ms(&self) -> u64 { 500 } /// True when currently cooling down between workloads (for the TUI /// to render a "cooldown" state vs a "running" state). pub fn is_cooldown(&self) -> bool { matches!(self.state, State::Cooldown { .. }) } /// 1-indexed ordinal of the current loop (for the TUI header, /// "loop 2/5"). 1 on the first pass, 2 after one reset, etc. When /// loop_count is 1 (single pass) this is always 1. pub fn loop_index(&self) -> u32 { self.loop_index } /// The configured loop count (1 = single pass, 0 = infinite). For /// the TUI header ("loop N/M"). 0 renders as "loop N/∞". pub fn loop_count(&self) -> u32 { self.cfg.loop_count } /// One tick of the state machine. Call every `poll_ms`. Polls the /// workload's `is_running()` to detect completion + advances state. pub fn tick(&mut self) -> io::Result { let elapsed = self.state_start.elapsed(); match &mut self.state { State::Cooldown { remaining } => { let cooldown_dur = Duration::from_secs(*remaining); if elapsed < cooldown_dur { // Still cooling. Don't reset state_start — `elapsed` // is measured from the cooldown's start, and comparing // to the full cooldown duration is the clean test. return Ok(TickOutcome::Idle); } // Cooldown done — start the next workload. self.start_current_workload() } State::Running { duration_secs: _ } => { let running = self.workloads[self.current] .as_mut() .map(|w| w.is_running()) .unwrap_or(false); if !running { // Workload self-exited — finish immediately. return self.finish_current_workload(); } // Workload still running. We do NOT hard-stop on the // duration timer — the duration is a display-only estimate // ("this will probably take X seconds"). The workload // runs until it self-exits (y-cruncher finishes its FFT // iteration, 7z finishes its passes) or the user quits. // This avoids false Stopped verdicts on single-core runs // where the TUI process steals CPU from the workload, // extending its wall time past the duration + any grace. Ok(TickOutcome::Idle) } State::Done => Ok(TickOutcome::SweepDone), } } fn start_current_workload(&mut self) -> io::Result { // Skip any None entries (shouldn't happen since `current` only // advances to Some indices, but be defensive). while self.current < self.workloads.len() && self.workloads[self.current].is_none() { self.current += 1; } if self.current >= self.workloads.len() { self.state = State::Done; return Ok(TickOutcome::SweepDone); } let spec = &self.cfg.workloads[self.current]; let report = Report::new( &self.cfg.report_dir, &spec.name, &self.workload_params(spec), )?; self.report = Some(report); let w = self.workloads[self.current].as_mut().unwrap(); w.start()?; let duration_secs = spec.duration_secs; self.state = State::Running { duration_secs }; self.state_start = Instant::now(); Ok(TickOutcome::WorkloadStarted { idx: self.current, name: spec.name.clone(), duration_secs, }) } fn finish_current_workload(&mut self) -> io::Result { let verdict = { let w = self.workloads[self.current].as_mut(); if let Some(w) = w { if w.is_running() { let _ = w.stop(); } let v = w.wait(); (v, w.score()) } else { (Verdict::Error("no workload".to_string()), Score::None) } }; let (verdict, score) = verdict; let report_path = self .report .as_ref() .map(|r| r.path().to_path_buf()) .unwrap_or_default(); if let Some(r) = &mut self.report { r.finish(verdict.clone(), score.clone())?; } let idx = self.current; let name = self.cfg.workloads[idx].name.clone(); // Advance to the next workload + enter cooldown. self.current += 1; // Skip None entries on advance too. while self.current < self.workloads.len() && self.workloads[self.current].is_none() { self.current += 1; } if self.current >= self.workloads.len() { // End of this loop's sweep. Decide whether to loop again. // loops_left == 0 means INFINITE (loop forever); > 1 means // there are more loops to go; == 1 means this was the last. let loop_again = self.loops_left == 0 || self.loops_left > 1; if loop_again { // Decrement (unless infinite — keep 0) + reset to loop start. if self.loops_left > 0 { self.loops_left -= 1; } self.loop_index += 1; self.current = 0; // Skip any leading None entries (defensive — shouldn't // happen since construction skips them, but the same // loop the constructor relies on). while self.current < self.workloads.len() && self.workloads[self.current].is_none() { self.current += 1; } // If the whole workload list was None (shouldn't happen), // fall through to Done. if self.current >= self.workloads.len() { self.state = State::Done; } else { self.state = State::Cooldown { remaining: self.cfg.cool_secs, }; self.state_start = Instant::now(); } } else { self.state = State::Done; } } else { self.state = State::Cooldown { remaining: self.cfg.cool_secs, }; self.state_start = Instant::now(); } Ok(TickOutcome::WorkloadFinished { idx, name, verdict, score, report_path, }) } /// Build the params string for a workload spec (mirrors the Workload's /// own `params()` but available before the workload is started, so the /// report header can be written at start time). fn workload_params(&self, spec: &crate::config::WorkloadSpec) -> String { format!("{}s cores={}", spec.duration_secs, spec.cores) } } #[cfg(test)] mod tests { use super::*; use crate::config::WorkloadSpec; /// A workload that exits after N `is_running()` calls, so tests can /// drive the state machine without real time or real processes. struct MockWorkload { name: String, ticks_until_exit: usize, ticks_seen: usize, started: bool, stopped: bool, } impl Workload for MockWorkload { fn name(&self) -> &str { &self.name } fn params(&self) -> String { format!("mock ticks={}", self.ticks_until_exit) } fn start(&mut self) -> io::Result<()> { self.started = true; self.stopped = false; self.ticks_seen = 0; Ok(()) } fn is_running(&mut self) -> bool { if !self.started || self.stopped { return false; } self.ticks_seen += 1; self.ticks_seen < self.ticks_until_exit } fn stop(&mut self) -> io::Result<()> { self.stopped = true; Ok(()) } fn wait(&mut self) -> Verdict { if self.stopped { Verdict::Stopped } else { Verdict::Clean } } fn score(&mut self) -> Score { Score::None } } fn mock_cfg(workloads: Vec) -> Config { Config { workloads, poll_ms: 10, cool_secs: 0, // no cooldown in tests loop_count: 1, report_dir: std::path::PathBuf::from("/tmp/opencode/run-tests"), bin_dir: std::path::PathBuf::from("bin"), loaded_path: None, } } /// Build a Run with mock workloads in place of `from_spec`. Bypasses /// the real constructor so we don't need real binaries. fn run_with_mocks(cfg: Config, mocks: Vec) -> Run { let workloads: Vec>> = mocks .into_iter() .map(|m| Some(Box::new(m) as Box)) .collect(); let loops_left = cfg.loop_count; Run { cfg, workloads, skipped: Vec::new(), current: 0, state: State::Cooldown { remaining: 0 }, state_start: Instant::now(), report: None, loops_left, loop_index: 1, } } #[test] fn sweep_runs_all_workloads_then_done() { let _ = std::fs::remove_dir_all("/tmp/opencode/run-tests"); let cfg = mock_cfg(vec![ WorkloadSpec { name: "mock-a".to_string(), duration_secs: 60, cores: "0-15".to_string(), }, WorkloadSpec { name: "mock-b".to_string(), duration_secs: 60, cores: "0-15".to_string(), }, ]); let mut run = run_with_mocks( cfg, vec![ MockWorkload { name: "mock-a".to_string(), ticks_until_exit: 3, ticks_seen: 0, started: false, stopped: false, }, MockWorkload { name: "mock-b".to_string(), ticks_until_exit: 2, ticks_seen: 0, started: false, stopped: false, }, ], ); let mut started = 0; let mut finished = 0; let mut done = false; for _ in 0..20 { match run.tick().unwrap() { TickOutcome::WorkloadStarted { name, .. } => { started += 1; assert!(name == "mock-a" || name == "mock-b"); } TickOutcome::WorkloadFinished { name, verdict, .. } => { finished += 1; assert!(name == "mock-a" || name == "mock-b"); assert_eq!(verdict, Verdict::Clean); } TickOutcome::SweepDone => { done = true; break; } TickOutcome::Idle => {} } } assert!(done); assert_eq!(started, 2); assert_eq!(finished, 2); } #[test] fn running_never_hard_stops_on_duration() { // The run loop does NOT hard-stop on the duration timer — the // duration is a display-only estimate. A workload that never // self-exits should keep running indefinitely (the only stop // path is the user quitting via the TUI, which calls stop()). // Verify: a never-exiting mock + a long duration → the run loop // never returns WorkloadFinished, no matter how many ticks. let _ = std::fs::remove_dir_all("/tmp/opencode/run-tests"); let mut cfg = mock_cfg(vec![WorkloadSpec { name: "never-exits".to_string(), duration_secs: 0, cores: "0-15".to_string(), }]); cfg.loop_count = 1; let mut run = run_with_mocks( cfg, vec![MockWorkload { name: "never-exits".to_string(), ticks_until_exit: 1000, ticks_seen: 0, started: false, stopped: false, }], ); // Tick once to start (Cooldown → Running). let _ = run.tick().unwrap(); // Manually rewind state_start so elapsed is huge (well past any // duration + grace). With no hard stop, the run loop should keep // returning Idle, NOT WorkloadFinished. run.state_start = Instant::now() .checked_sub(Duration::from_secs(3600)) .unwrap_or(Instant::now()); for _ in 0..20 { let outcome = run.tick().unwrap(); match outcome { TickOutcome::Idle => {} // expected — still running TickOutcome::WorkloadFinished { .. } => { panic!("run loop should not hard-stop on duration timer"); } _ => {} } } } #[test] fn total_workloads_counts_only_some() { let cfg = mock_cfg(vec![]); let mut run = run_with_mocks(cfg, vec![]); run.workloads.push(None); // a skipped entry run.workloads.push(Some(Box::new(MockWorkload { name: "x".to_string(), ticks_until_exit: 1, ticks_seen: 0, started: false, stopped: false, }))); assert_eq!(run.total_workloads(), 1); } #[test] fn skipped_workloads_reported() { let run = Run::new(Config::default()); let _ = run.skipped(); } #[test] fn loop_resets_current_and_decrements() { // Two workloads, loop_count = 2 → should run each twice (4 starts // total) before SweepDone. let _ = std::fs::remove_dir_all("/tmp/opencode/run-tests"); let mut cfg = mock_cfg(vec![ WorkloadSpec { name: "a".to_string(), duration_secs: 60, cores: "0-15".to_string(), }, WorkloadSpec { name: "b".to_string(), duration_secs: 60, cores: "0-15".to_string(), }, ]); cfg.loop_count = 2; let mut run = run_with_mocks( cfg, vec![ MockWorkload { name: "a".to_string(), ticks_until_exit: 2, ticks_seen: 0, started: false, stopped: false, }, MockWorkload { name: "b".to_string(), ticks_until_exit: 2, ticks_seen: 0, started: false, stopped: false, }, ], ); let mut starts = 0; let mut done = false; for _ in 0..40 { match run.tick().unwrap() { TickOutcome::WorkloadStarted { name, .. } => { starts += 1; assert!(name == "a" || name == "b"); } TickOutcome::WorkloadFinished { .. } => {} TickOutcome::SweepDone => { done = true; break; } TickOutcome::Idle => {} } } assert!(done, "should reach SweepDone"); assert_eq!(starts, 4, "two workloads × two loops = 4 starts"); assert_eq!( run.loop_index(), 2, "loop_index should be 2 after one reset" ); } #[test] fn loop_infinite_never_done() { // loop_count = 0 means infinite. With short mocks the sweep would // complete many times; verify it never reaches SweepDone (we'd // be here forever in real life, but we cap the loop + check that // it's still going after several sweep completions). let _ = std::fs::remove_dir_all("/tmp/opencode/run-tests"); let mut cfg = mock_cfg(vec![WorkloadSpec { name: "a".to_string(), duration_secs: 60, cores: "0-15".to_string(), }]); cfg.loop_count = 0; // infinite let mut run = run_with_mocks( cfg, vec![MockWorkload { name: "a".to_string(), ticks_until_exit: 2, ticks_seen: 0, started: false, stopped: false, }], ); let mut starts = 0; let mut sweep_done = false; for _ in 0..100 { match run.tick().unwrap() { TickOutcome::WorkloadStarted { .. } => starts += 1, TickOutcome::SweepDone => { sweep_done = true; break; } _ => {} } } assert!(!sweep_done, "infinite loop should never reach SweepDone"); // Should have completed several loops (start count >> 1). assert!( starts > 5, "infinite loop should have started many times, got {starts}" ); } }