//! Run orchestration — the new core (adlermon 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, } 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(), )), } } Run { cfg, workloads, skipped, current: 0, state: State::Cooldown { remaining: 0, }, state_start: Instant::now(), report: None, } } /// 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. pub fn remaining_secs(&self) -> Option { match &self.state { State::Cooldown { remaining } => { let dur = Duration::from_secs(*remaining); Some((dur - self.state_start.elapsed()).as_secs_f64().max(0.0)) } State::Running { duration_secs } => { let dur = Duration::from_secs(*duration_secs); Some((dur - self.state_start.elapsed()).as_secs_f64().max(0.0)) } 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 } /// 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 { .. }) } /// 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 || elapsed >= Duration::from_secs(*duration_secs) { // Workload finished (on its own or we stop it now). return self.finish_current_workload(); } 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() { 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; 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 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(); Run { cfg, workloads, skipped: Vec::new(), current: 0, state: State::Cooldown { remaining: 0 }, state_start: Instant::now(), report: None, } } #[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 stopped_verdict_when_duration_elapses() { // MockWorkload that never self-exits (ticks_until_exit huge); // the run's duration timer must stop it → Stopped verdict. let _ = std::fs::remove_dir_all("/tmp/opencode/run-tests"); let cfg = mock_cfg(vec![ WorkloadSpec { name: "never-exits".to_string(), duration_secs: 0, cores: "0-15".to_string() }, ]); 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 }, ]); let mut verdict = None; for _ in 0..10 { if let TickOutcome::WorkloadFinished { verdict: v, .. } = run.tick().unwrap() { verdict = Some(v); break; } } // duration_secs=0 → first Running tick sees elapsed >= 0 → stops. assert_eq!(verdict, Some(Verdict::Stopped)); } #[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(); } }