//! Workload abstraction for the ISA sweep. A Workload is something the run
//! loop can start, let run for a duration while sensors sample around it,
//! stop, and read a stability verdict from.
//!
//! Two shapes live behind this trait:
//! - **Spawned** workloads (stress-ng, y-cruncher, 7zip) own a child
//! process. `start()` spawns, `stop()` signals + waits, the verdict
//! comes from the exit code + parsed stdout.
//! - **In-process** workloads (c2c-latency, dram-latency, deferred) run
//! in a thread. `start()` spawns the thread, `stop()` sets a flag,
//! the verdict comes from the measurement itself.
//!
//! The trait abstracts over both so the run loop treats every workload
//! the same. v1 ships the spawned impls; in-process impls come next.
#![allow(dead_code)] // consumed by run.rs + ui.rs — not yet wired
use std::io::{self, Read};
use std::path::PathBuf;
use std::process::{Child, Command, ExitStatus};
/// What a workload reports when it's done. The stability verdict is the
/// product — `Clean` means the workload's own self-check passed (y-cruncher
/// "Passed", 7z exit 0, stress-ng exit 0); anything else is a failure mode
/// the OC stability question cares about.
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum Verdict {
/// Workload self-check passed. The chip was stable for this run.
Clean,
/// Workload self-check failed (y-cruncher "failed", 7z nonzero exit,
/// stress-ng nonzero exit). Instability at this ISA rail.
Failed(String),
/// Workload was stopped early (user quit, timeout, sweep cancelled).
/// Not a stability signal — the run didn't complete.
Stopped,
/// Workload couldn't run (binary missing, spawn error, etc.). Not a
/// stability signal — infrastructure failure.
Error(String),
}
/// The common surface the run loop drives. See the module doc for the two
/// shapes behind this trait.
pub trait Workload {
/// Human-readable name (e.g. "y-cruncher-avx2"). Used in the TUI and
/// the report.
fn name(&self) -> &str;
/// One-line params string (e.g. "60s cores=0-15 FFTv4"). For the report.
fn params(&self) -> String;
/// Start the workload. Called once per run.
fn start(&mut self) -> io::Result<()>;
/// True between `start()` and `stop()`. The run loop polls this each
/// tick to decide whether to keep sampling or move to cooldown.
/// `&mut self` because `try_wait` mutates the Child's internal state.
fn is_running(&mut self) -> bool;
/// Stop the workload. May be called before the duration elapses (user
/// quit, sweep cancelled). Must be safe to call on an already-stopped
/// workload.
fn stop(&mut self) -> io::Result<()>;
/// Block until the workload finishes (or was already stopped). Returns
/// the verdict. Called once, after the run loop decides the workload is
/// done (either `is_running()` went false or the run loop's own
/// duration timer expired and it called `stop()`).
fn wait(&mut self) -> Verdict;
}
// ---------------------------------------------------------------------------
// Spawned-workload helper
// ---------------------------------------------------------------------------
/// Common machinery for workloads that spawn a child process. Owns the
/// Child, captures stdout (for verdict parsing), and provides a generic
/// `stop()` that tries SIGTERM then SIGKILL. Subtypes build the Command in
/// their `start()` and parse stdout in their `wait()`.
struct SpawnedWorkload {
child: Option<Child>,
stdout: Vec<u8>,
/// Set when `stop()` was called by the run loop (vs. the process
/// exiting on its own). Distinguishes `Stopped` from `Clean`/`Failed`.
stopped: bool,
/// Captured stderr (kept for diagnostics; not parsed for verdict).
stderr: Vec<u8>,
}
impl SpawnedWorkload {
fn new() -> Self {
SpawnedWorkload {
child: None,
stdout: Vec::new(),
stderr: Vec::new(),
stopped: false,
}
}
fn spawn(cmd: &mut Command) -> io::Result<Self> {
let child = cmd
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.spawn()?;
Ok(SpawnedWorkload {
child: Some(child),
stdout: Vec::new(),
stderr: Vec::new(),
stopped: false,
})
}
fn is_running(&mut self) -> bool {
match &mut self.child {
Some(child) => child.try_wait().ok().flatten().is_none(),
None => false,
}
}
fn stop(&mut self) -> io::Result<()> {
if let Some(child) = &mut self.child {
self.stopped = true;
// Try a polite SIGTERM first; the run loop will call wait()
// next which reaps. If TERM doesn't take in a couple seconds,
// escalate to SIGKILL — but we don't block here (the run loop
// owns timing). For now just SIGTERM and let wait() handle the
// rest. On Unix, `start_kill` sends SIGKILL; for SIGTERM we
// use the raw pid + nix-less libc. Keep it simple: SIGKILL
// directly. Stress workloads don't have cleanup that warrants
// a graceful TERM, and a stuck child would hang the sweep.
let _ = child.kill();
}
Ok(())
}
/// Drain stdout/stderr from the child into our buffers. Called by
/// subtypes' `wait()` before parsing the verdict. Reads as much as is
/// available without blocking.
fn drain(&mut self) -> io::Result<()> {
let Some(child) = &mut self.child else {
return Ok(());
};
if let Some(stdout) = &mut child.stdout {
let mut buf = [0u8; 4096];
loop {
let n = stdout.read(&mut buf)?;
if n == 0 {
break;
}
self.stdout.extend_from_slice(&buf[..n]);
}
}
if let Some(stderr) = &mut child.stderr {
let mut buf = [0u8; 4096];
loop {
let n = stderr.read(&mut buf)?;
if n == 0 {
break;
}
self.stderr.extend_from_slice(&buf[..n]);
}
}
Ok(())
}
/// Block until the child exits, drain stdout/stderr, and return the
/// raw ExitStatus. Called by subtypes' `wait()`.
fn wait_raw(&mut self) -> io::Result<Option<ExitStatus>> {
let Some(child) = &mut self.child else {
return Ok(None);
};
let status = child.wait()?;
// After wait() returns, the pipes are at EOF — drain anything we
// missed while we weren't polling.
self.drain()?;
Ok(Some(status))
}
fn _take_stdout(&mut self) -> String {
String::from_utf8_lossy(&self.stdout).to_string()
}
}
// ---------------------------------------------------------------------------
// stress-ng --cpu (integer rail)
// ---------------------------------------------------------------------------
/// Integer rail. `stress-ng --cpu N --taskset <range> --timeout T --cpu-method all`.
/// Stress-ng exits 0 on success, nonzero on failure (including bogo-op
/// self-check mismatches).
pub struct StressNg {
duration_secs: u64,
cores: String,
inner: SpawnedWorkload,
}
impl StressNg {
pub fn new(duration_secs: u64, cores: String) -> Self {
StressNg {
duration_secs,
cores,
inner: SpawnedWorkload::new(),
}
}
}
impl Workload for StressNg {
fn name(&self) -> &str {
"stress-ng-cpu"
}
fn params(&self) -> String {
format!("{}s cores={} --cpu --cpu-method all", self.duration_secs, self.cores)
}
fn start(&mut self) -> io::Result<()> {
let mut cmd = Command::new("stress-ng");
cmd.arg("--cpu")
.arg("15") // all logical cpus; --taskset below pins to the range
.arg("--taskset")
.arg(&self.cores)
.arg("--timeout")
.arg(format!("{}s", self.duration_secs))
.arg("--cpu-method")
.arg("all")
.arg("--metrics-brief");
self.inner = SpawnedWorkload::spawn(&mut cmd)?;
Ok(())
}
fn is_running(&mut self) -> bool {
self.inner.is_running()
}
fn stop(&mut self) -> io::Result<()> {
self.inner.stop()
}
fn wait(&mut self) -> Verdict {
let stopped = self.inner.stopped;
match self.inner.wait_raw() {
Ok(Some(status)) => {
if stopped {
return Verdict::Stopped;
}
if status.success() {
Verdict::Clean
} else {
Verdict::Failed(format!("stress-ng exit {}", status))
}
}
Ok(None) => Verdict::Error("no child".to_string()),
Err(e) => Verdict::Error(format!("wait: {e}")),
}
}
}
// ---------------------------------------------------------------------------
// y-cruncher (SSE or AVX2 rail, via per-microarch variant binary)
// ---------------------------------------------------------------------------
/// SSE or AVX2 rail. y-cruncher has NO CLI flag for ISA selection — we
/// force the ISA by invoking the per-microarch VARIANT BINARY directly
/// (NOT the wrapper, which auto-selects). Verified 2026-09-01:
/// SSE = Binaries/11-SNB ~ Hina (Sandy Bridge, SSE4.2 only)
/// AVX2 = Binaries/13-HSW ~ Airi (Haswell, AVX2)
/// Both run standalone, report "Passed"/"failed" per FFT round, "Stop on
/// Error: Enabled" by default. CLI: `<variant> skip-warnings stress -D:s
/// -TL:s [algorithm]`.
pub struct YCruncher {
duration_secs: u64,
cores: String,
variant: &'static str, // "11-SNB ~ Hina" or "13-HSW ~ Airi"
label: &'static str, // "sse" or "avx2"
bin_dir: PathBuf,
inner: SpawnedWorkload,
}
impl YCruncher {
pub fn sse(duration_secs: u64, cores: String, bin_dir: PathBuf) -> Self {
YCruncher {
duration_secs,
cores,
variant: "11-SNB ~ Hina",
label: "sse",
bin_dir,
inner: SpawnedWorkload::new(),
}
}
pub fn avx2(duration_secs: u64, cores: String, bin_dir: PathBuf) -> Self {
YCruncher {
duration_secs,
cores,
variant: "13-HSW ~ Airi",
label: "avx2",
bin_dir,
inner: SpawnedWorkload::new(),
}
}
/// Resolve the variant binary path. y-cruncher ships as a versioned
/// subdir under bin_dir (e.g. `bin/y-cruncher v0.8.7.9547-static/`),
/// so we glob for the first match rather than hardcode the version.
fn variant_path(&self) -> Option<PathBuf> {
let entries = std::fs::read_dir(&self.bin_dir).ok()?;
for entry in entries.flatten() {
let name = entry.file_name();
let name = name.to_str()?;
if name.starts_with("y-cruncher") {
let variant = entry.path().join("Binaries").join(self.variant);
if variant.exists() {
return Some(variant);
}
}
}
None
}
}
impl Workload for YCruncher {
fn name(&self) -> &str {
match self.label {
"sse" => "y-cruncher-sse",
"avx2" => "y-cruncher-avx2",
_ => "y-cruncher",
}
}
fn params(&self) -> String {
format!(
"{}s cores={} variant={} stress FFTv4",
self.duration_secs, self.cores, self.variant
)
}
fn start(&mut self) -> io::Result<()> {
let variant_path = self.variant_path().ok_or_else(|| {
io::Error::new(
io::ErrorKind::NotFound,
format!(
"y-cruncher variant '{}' not found under bin_dir {:?} \
(expected bin/y-cruncher v*/Binaries/{})",
self.variant, self.bin_dir, self.variant,
),
)
})?;
// Canonicalize before setting current_dir — when current_dir is
// set, the kernel resolves the executable path RELATIVE TO THE NEW
// cwd, so a relative variant path would be looked up inside the
// y-cruncher bundle dir and not found. Absolute path survives the
// cwd change.
let variant_path = variant_path.canonicalize().map_err(|e| {
io::Error::new(
io::ErrorKind::NotFound,
format!("y-cruncher variant path canonicalize failed: {e}"),
)
})?;
let mut cmd = Command::new(&variant_path);
cmd.arg("skip-warnings")
.arg("stress")
// y-cruncher expects `-D:3` as ONE arg (not `-D:` + `3`).
.arg(format!("-D:{}", self.duration_secs))
.arg(format!("-TL:{}", self.duration_secs))
.arg("FFTv4");
// y-cruncher reads its Libraries.txt etc. from CWD — run from the
// variant's parent directory (the static bundle dir).
if let Some(dir) = variant_path.parent().and_then(|p| p.parent()) {
cmd.current_dir(dir);
}
self.inner = SpawnedWorkload::spawn(&mut cmd)?;
Ok(())
}
fn is_running(&mut self) -> bool {
self.inner.is_running()
}
fn stop(&mut self) -> io::Result<()> {
self.inner.stop()
}
fn wait(&mut self) -> Verdict {
let stopped = self.inner.stopped;
match self.inner.wait_raw() {
Ok(Some(status)) => {
if stopped {
return Verdict::Stopped;
}
let stdout = String::from_utf8_lossy(&self.inner.stdout).to_string();
// y-cruncher prints "Running FFTv4: Passed" per round and
// "failed" on mismatch. The definitive signal is the exit
// code (nonzero on self-check failure since "Stop on Error"
// is enabled by default), but we also scan stdout for the
// explicit "failed" string in case of partial output.
let failed_text = stdout
.lines()
.any(|l| l.contains("failed") && !l.contains("Running"));
if status.success() && !failed_text {
Verdict::Clean
} else if failed_text {
Verdict::Failed(format!("y-cruncher self-check failed (variant={})", self.variant))
} else {
Verdict::Failed(format!("y-cruncher exit {} (variant={})", status, self.variant))
}
}
Ok(None) => Verdict::Error("no child".to_string()),
Err(e) => Verdict::Error(format!("wait: {e}")),
}
}
}
// ---------------------------------------------------------------------------
// 7-Zip benchmark (the z-7ip bench — NOT y-cruncher)
// ---------------------------------------------------------------------------
/// z-7ip bench. `7z b` runs 7-Zip's built-in compression/decompression
/// benchmark; CPU + memory bandwidth throughput. Reports MIPS + total
/// score (the comparative-perf metric across configs). Exits 0 on success.
/// `7z b -mmt<N>` sets thread count; we pass cores count translated from
/// the range string (0-15 → 16 threads). Defaults to all logical cpus if
/// the range is "0-15".
pub struct SevenZip {
duration_secs: u64,
cores: String,
inner: SpawnedWorkload,
}
impl SevenZip {
pub fn new(duration_secs: u64, cores: String) -> Self {
SevenZip {
duration_secs,
cores,
inner: SpawnedWorkload::new(),
}
}
/// Translate a cores range ("0-15", "0-11", "12-15") into a thread
/// count for `7z b -mmt<N>`. Returns None if the range is malformed
/// — 7z's default (all cpus) is fine in that case.
fn thread_count(&self) -> Option<u32> {
let s = self.cores.trim();
if let Some((a, b)) = s.split_once('-') {
let a: u32 = a.trim().parse().ok()?;
let b: u32 = b.trim().parse().ok()?;
if b >= a {
return Some(b - a + 1);
}
}
None
}
}
impl Workload for SevenZip {
fn name(&self) -> &str {
"7zip-bench"
}
fn params(&self) -> String {
format!(
"{}s cores={} (7z b{})",
self.duration_secs,
self.cores,
match self.thread_count() {
Some(n) => format!(" -mmt{n}"),
None => String::new(),
}
)
}
fn start(&mut self) -> io::Result<()> {
let mut cmd = Command::new("7z");
cmd.arg("b");
if let Some(n) = self.thread_count() {
cmd.arg(format!("-mmt{n}"));
}
// 7z b runs a fixed number of passes by default, not time-bound.
// For a sweep we want a comparable run length across configs —
// pass `-mmt<N>` for thread control and let the duration_secs be
// a soft target (7z finishes when it finishes; the run loop will
// stop() if it overruns). Document this in the report.
self.inner = SpawnedWorkload::spawn(&mut cmd)?;
Ok(())
}
fn is_running(&mut self) -> bool {
self.inner.is_running()
}
fn stop(&mut self) -> io::Result<()> {
self.inner.stop()
}
fn wait(&mut self) -> Verdict {
let stopped = self.inner.stopped;
match self.inner.wait_raw() {
Ok(Some(status)) => {
if stopped {
return Verdict::Stopped;
}
if status.success() {
Verdict::Clean
} else {
Verdict::Failed(format!("7z exit {}", status))
}
}
Ok(None) => Verdict::Error("no child".to_string()),
Err(e) => Verdict::Error(format!("wait: {e}")),
}
}
}
// ---------------------------------------------------------------------------
// Public constructor: map a WorkloadSpec name to an impl
// ---------------------------------------------------------------------------
/// Build a Workload from a config spec. Returns None for unknown names
/// (the run loop skips them with a warning) or names that aren't wired
/// yet (c2c-latency, dram-latency, minecraft-server — deferred).
pub fn from_spec(
spec: &crate::config::WorkloadSpec,
bin_dir: &std::path::Path,
) -> Option<Box<dyn Workload>> {
match spec.name.as_str() {
"stress-ng-cpu" => Some(Box::new(StressNg::new(spec.duration_secs, spec.cores.clone()))),
"y-cruncher-sse" => {
Some(Box::new(YCruncher::sse(spec.duration_secs, spec.cores.clone(), bin_dir.to_path_buf())))
}
"y-cruncher-avx2" => Some(Box::new(YCruncher::avx2(
spec.duration_secs,
spec.cores.clone(),
bin_dir.to_path_buf(),
))),
"7zip-bench" => Some(Box::new(SevenZip::new(spec.duration_secs, spec.cores.clone()))),
// Deferred — hand-rolled in-process workloads, next chunk.
"c2c-latency" | "dram-latency" | "minecraft-server" => None,
_ => None,
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn stressng_params_include_range() {
let w = StressNg::new(60, "0-11".to_string());
assert!(w.params().contains("0-11"));
assert!(w.params().contains("60s"));
}
#[test]
fn ycruncher_variant_paths_resolve() {
let bin_dir = PathBuf::from("bin");
let sse = YCruncher::sse(60, "0-15".to_string(), bin_dir.clone());
let avx2 = YCruncher::avx2(60, "0-15".to_string(), bin_dir);
assert_eq!(sse.variant, "11-SNB ~ Hina");
assert_eq!(avx2.variant, "13-HSW ~ Airi");
// variant_path() should find the actual bundle if present in ./bin
// (smoke-tested below; skip assertion here so tests pass without
// the binary).
let _ = sse.variant_path();
let _ = avx2.variant_path();
}
#[test]
fn ycruncher_names_distinguish_rails() {
let sse = YCruncher::sse(60, "0-15".to_string(), PathBuf::from("bin"));
let avx2 = YCruncher::avx2(60, "0-15".to_string(), PathBuf::from("bin"));
assert_eq!(sse.name(), "y-cruncher-sse");
assert_eq!(avx2.name(), "y-cruncher-avx2");
}
#[test]
fn sevenzip_thread_count_parses_ranges() {
let w = SevenZip::new(60, "0-15".to_string());
assert_eq!(w.thread_count(), Some(16));
let w = SevenZip::new(60, "0-11".to_string());
assert_eq!(w.thread_count(), Some(12));
let w = SevenZip::new(60, "12-15".to_string());
assert_eq!(w.thread_count(), Some(4));
let w = SevenZip::new(60, "junk".to_string());
assert_eq!(w.thread_count(), None);
}
#[test]
fn from_spec_wires_known_names() {
let spec = crate::config::WorkloadSpec {
name: "stress-ng-cpu".to_string(),
duration_secs: 30,
cores: "0-15".to_string(),
};
let w = from_spec(&spec, std::path::Path::new("bin"));
assert!(w.is_some());
assert_eq!(w.unwrap().name(), "stress-ng-cpu");
}
#[test]
fn from_spec_rejects_unknown_names() {
let spec = crate::config::WorkloadSpec {
name: "no-such-workload".to_string(),
duration_secs: 30,
cores: "0-15".to_string(),
};
assert!(from_spec(&spec, std::path::Path::new("bin")).is_none());
}
#[test]
fn from_spec_rejects_deferred_workloads() {
// c2c-latency / dram-latency / minecraft-server are in scope but
// not yet implemented — from_spec returns None so the run loop
// can skip them with a warning rather than panic.
for name in ["c2c-latency", "dram-latency", "minecraft-server"] {
let spec = crate::config::WorkloadSpec {
name: name.to_string(),
duration_secs: 30,
cores: "0-15".to_string(),
};
assert!(from_spec(&spec, std::path::Path::new("bin")).is_none(), "{name} should be None");
}
}
}