f24503a3caf6f35f5e69a06ef8ba4312c2acc98f / alderbench/src/run.rs · 24472 bytes · raw
//! 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<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,
/// The report being written for the current workload.
report: Option<Report>,
/// 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<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(),
)),
}
}
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<f64> {
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<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 {
// 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<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();
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();
}
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<WorkloadSpec>) -> 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<MockWorkload>) -> Run {
let workloads: Vec<Option<Box<dyn Workload>>> =
mocks.into_iter().map(|m| Some(Box::new(m) as Box<dyn Workload>)).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}");
}
}