//! Run orchestration — the new core (adlermon has no equivalent). Owns
//! the sweep lifecycle: for each workload in the config, cooldown →
//! workload.start() → tick loop (sample sensors, add to report) →
//! workload.stop() / wait for exit → flush report. Front-end-agnostic:
//! the TUI drives it by calling `tick()` every `poll_ms`; a future
//! `--headless` mode 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::sensors;
use crate::workload::{self, 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 },
/// A sensor sample was taken and added to the report. The latest
/// snapshot is included so the TUI can render it without borrowing
/// into the run's history.
SampleTaken { latest: sensors::Snapshot },
/// The workload exited (either it self-finished or the duration timer
/// elapsed and we stopped it). The verdict + report path are included
/// for the TUI's result table + stdout summary.
WorkloadFinished {
idx: usize,
name: String,
verdict: Verdict,
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 sample was due). 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<Option<Box<dyn Workload>>>,
/// 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,
/// When the last sensor sample was taken (for the dt_secs in power
/// delta math + pacing samples within a tick).
last_sample: Instant,
/// Previous RAPL energy reading, for the power-delta computation in
/// `sensors::snapshot`. None on the first sample of a workload.
prev_energy_uj: Option<u64>,
/// The report being written for the current workload.
report: Option<Report>,
}
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<Option<Box<dyn Workload>>> = 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(),
last_sample: Instant::now(),
prev_energy_uj: None,
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()
}
/// One tick of the state machine. Call every `poll_ms`. Reads sensors
/// (when running), advances the state, and returns what happened.
pub fn tick(&mut self) -> io::Result<TickOutcome> {
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();
}
// Still running — take a sensor sample.
let now = Instant::now();
let dt = now.duration_since(self.last_sample).as_secs_f64();
self.last_sample = now;
let snap = sensors::snapshot(self.prev_energy_uj, dt);
self.prev_energy_uj = snap.energy_uj();
// Stamp the sample with the workload's elapsed time.
let t = elapsed.as_secs_f64();
let mut snap = snap;
snap.t = t;
if let Some(r) = &mut self.report {
r.add_sample(snap.clone())?;
}
Ok(TickOutcome::SampleTaken { latest: snap })
}
State::Done => Ok(TickOutcome::SweepDone),
}
}
fn start_current_workload(&mut self) -> io::Result<TickOutcome> {
// 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();
self.last_sample = Instant::now();
self.prev_energy_uj = None;
Ok(TickOutcome::WorkloadStarted {
idx: self.current,
name: spec.name.clone(),
duration_secs,
})
}
fn finish_current_workload(&mut self) -> io::Result<TickOutcome> {
let verdict = {
let w = self.workloads[self.current].as_mut();
if let Some(w) = w {
if w.is_running() {
let _ = w.stop();
}
w.wait()
} else {
Verdict::Error("no workload".to_string())
}
};
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())?;
}
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,
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 mock_cfg(workloads: Vec<WorkloadSpec>) -> Config {
Config {
workloads,
poll_ms: 10,
cool_secs: 0, // no cooldown in tests
vcore_limit: 1.4,
temp_crit: 95.0,
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<MockWorkload>) -> Run {
let workloads: Vec<Option<Box<dyn Workload>>> =
mocks.into_iter().map(|m| Some(Box::new(m) as Box<dyn Workload>)).collect();
Run {
cfg,
workloads,
skipped: Vec::new(),
current: 0,
state: State::Cooldown { remaining: 0 },
state_start: Instant::now(),
last_sample: Instant::now(),
prev_energy_uj: None,
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 },
]);
// Tick 1: cooldown (0s) elapses → start mock-a.
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::SampleTaken { .. } => {}
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 {
match run.tick().unwrap() {
TickOutcome::WorkloadFinished { verdict: v, .. } => {
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());
// Can't easily inject skipped entries via the public ctor without
// an unknown workload name in the config; just verify the accessor
// exists + returns a slice.
let _ = run.skipped();
}
}