josie / alder-tools

//! 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");
        }
    }
}