From 73a1b3dd7ee0b5d7190b7ddfdac59e4887940d8a Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Mon, 8 Jun 2026 18:39:31 -0500 Subject: [PATCH 1/3] controller: measure benchmark memory and CPU from the process subtree MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The resource monitor read the pod-wide cgroup (`memory.current` / `cpu.stat`), so the reported "Peak memory" and "CPU" for a benchmark also counted page cache, leftover build artifacts, and — for benchmarks whose `bench.sh run` step compiles (e.g. `cargo bench --bench sql`) — the compiler itself. On a live wide_schema run that meant a single rustc (~16 GB, 12 cores) was being attributed to the benchmark. Scope both memory and CPU to the spawned benchmark command's process subtree by walking /proc, excluding compiler/build-tool processes (rustc, cargo, cc, ld, sccache, build-script-*, …): - memory: sum VmRSS over the subtree each second (peak/avg). - CPU: accumulate per-process utime/stime deltas over the window, so process churn (the bench binary starting after the compile finishes) is handled and the compile's CPU is excluded. shell.rs splits spawn from wait so the monitor gets the child pid. When no pid is available (or no /proc, e.g. local macOS) it falls back to the cgroup figures as before. Co-Authored-By: Claude Opus 4.8 (1M context) --- controller/src/runner/monitor.rs | 305 ++++++++++++++++++++++++++++--- controller/src/runner/shell.rs | 66 +++++-- 2 files changed, 334 insertions(+), 37 deletions(-) diff --git a/controller/src/runner/monitor.rs b/controller/src/runner/monitor.rs index e8782b4..d92d7d7 100644 --- a/controller/src/runner/monitor.rs +++ b/controller/src/runner/monitor.rs @@ -1,9 +1,18 @@ -//! Resource monitoring via cgroup v2 files. +//! Resource monitoring for benchmark execution. //! -//! Captures memory, CPU, and I/O stats during benchmark execution by reading -//! cgroup v2 files at `/sys/fs/cgroup/`. When cgroup files are unavailable -//! (e.g. local macOS development), all reads return `None` and stats default -//! to zero. +//! Memory is sampled from the benchmark command's **process subtree** by +//! walking `/proc//status` (summing `VmRSS`) and excluding compiler / +//! build-tool processes (`rustc`, `cargo`, `cc`, `ld`, `sccache`, …). That is +//! what we report as the benchmark's memory. Reading the pod-wide cgroup +//! (`/sys/fs/cgroup/memory.current`) instead would also count page cache, +//! leftover build artifacts, and — for benchmarks whose `bench.sh run` step +//! compiles (e.g. `cargo bench`) — the compiler itself, which can dwarf the +//! actual query memory by tens of GB. +//! +//! CPU stats still come from cgroup v2 `cpu.stat`. When the per-process root +//! pid is unknown the sampler falls back to the cgroup memory figure. When +//! neither `/proc` nor the cgroup files are available (e.g. local macOS +//! development), reads return `None` and stats default to zero. use std::fmt; use std::path::{Path, PathBuf}; @@ -92,27 +101,36 @@ struct CpuStat { system_usec: u64, } -/// Monitors cgroup v2 resource usage during benchmark execution. +/// Monitors resource usage during benchmark execution. pub struct CgroupMonitor { start_time: Instant, start_memory: u64, start_cpu: Option, + root_pid: Option, stop_flag: Arc, peak_memory: Arc, memory_sum: Arc, peak_spill: Arc, sample_count: Arc, + // Per-process CPU consumed by the benchmark subtree (excluding build + // tools), accumulated in the poll loop. Only meaningful when `root_pid` is + // set; otherwise CPU is taken from the cgroup delta in `finish`. + cpu_user_usec: Arc, + cpu_sys_usec: Arc, poll_handle: JoinHandle<()>, } impl CgroupMonitor { - /// Begin monitoring. Snapshots current cgroup values as baselines and - /// spawns a background polling task for memory and spill directory tracking. + /// Begin monitoring. `root_pid` is the spawned benchmark command's pid; + /// both memory and CPU are sampled from its process subtree, excluding + /// compiler/build-tool processes — so a `cargo bench` compile inside the + /// monitored window is not attributed to the benchmark. When `root_pid` is + /// `None`, falls back to the pod-wide cgroup figures. /// /// If `spill_dir` is provided, the polling loop will also sample the total /// size of files in that directory every second to track peak spill usage. - pub fn start(spill_dir: Option) -> Self { - let start_memory = read_memory_current().unwrap_or(0); + pub fn start(root_pid: Option, spill_dir: Option) -> Self { + let start_memory = sample_memory(root_pid).unwrap_or(0); let start_cpu = read_cpu_stat(); let stop_flag = Arc::new(AtomicBool::new(false)); @@ -120,6 +138,8 @@ impl CgroupMonitor { let memory_sum = Arc::new(AtomicU64::new(start_memory)); let peak_spill = Arc::new(AtomicU64::new(0)); let sample_count = Arc::new(AtomicU64::new(1)); + let cpu_user_usec = Arc::new(AtomicU64::new(0)); + let cpu_sys_usec = Arc::new(AtomicU64::new(0)); let poll_handle = { let stop = stop_flag.clone(); @@ -127,17 +147,50 @@ impl CgroupMonitor { let sum = memory_sum.clone(); let spill_peak = peak_spill.clone(); let count = sample_count.clone(); + let cpu_user = cpu_user_usec.clone(); + let cpu_sys = cpu_sys_usec.clone(); tokio::spawn(async move { + // Per-pid first/last cumulative CPU ticks, so each benchmark + // process contributes only the CPU it burned during the window + // (handles process churn as the bench binary starts/exits). + let mut cpu_first: std::collections::HashMap = + std::collections::HashMap::new(); + let mut cpu_last: std::collections::HashMap = + std::collections::HashMap::new(); + while !stop.load(Ordering::Relaxed) { tokio::time::sleep(POLL_INTERVAL).await; if stop.load(Ordering::Relaxed) { break; } - if let Some(current) = read_memory_current() { - peak.fetch_max(current, Ordering::Relaxed); - sum.fetch_add(current, Ordering::Relaxed); + if let Some((rss, cpus)) = proc_tree_sample(root_pid) { + peak.fetch_max(rss, Ordering::Relaxed); + sum.fetch_add(rss, Ordering::Relaxed); count.fetch_add(1, Ordering::Relaxed); + + for (pid, utime, stime) in cpus { + cpu_first.entry(pid).or_insert((utime, stime)); + cpu_last.insert(pid, (utime, stime)); + } + // Recompute accumulated subtree CPU each poll. + let mut tu = 0u64; + let mut ts = 0u64; + for (pid, &(fu, fs)) in &cpu_first { + if let Some(&(lu, ls)) = cpu_last.get(pid) { + tu += lu.saturating_sub(fu); + ts += ls.saturating_sub(fs); + } + } + cpu_user.store(ticks_to_usec(tu), Ordering::Relaxed); + cpu_sys.store(ticks_to_usec(ts), Ordering::Relaxed); + } else if root_pid.is_none() { + // Fallback path: keep memory sampling via cgroup. + if let Some(current) = read_memory_current() { + peak.fetch_max(current, Ordering::Relaxed); + sum.fetch_add(current, Ordering::Relaxed); + count.fetch_add(1, Ordering::Relaxed); + } } if let Some(ref dir) = spill_dir { let size = dir_size(dir); @@ -151,36 +204,47 @@ impl CgroupMonitor { start_time: Instant::now(), start_memory, start_cpu, + root_pid, stop_flag, peak_memory, memory_sum, peak_spill, sample_count, + cpu_user_usec, + cpu_sys_usec, poll_handle, } } - /// Stop monitoring and compute delta statistics. + /// Stop monitoring and compute statistics. pub async fn finish(self) -> ResourceStats { let wall_time = self.start_time.elapsed(); self.stop_flag.store(true, Ordering::Relaxed); let _ = self.poll_handle.await; - let end_memory = read_memory_current().unwrap_or(0); - let end_cpu = read_cpu_stat(); + let end_memory = sample_memory(self.root_pid).unwrap_or(0); let peak = self.peak_memory.load(Ordering::Relaxed).max(end_memory); let total_sum = self.memory_sum.load(Ordering::Relaxed) + end_memory; let total_count = self.sample_count.load(Ordering::Relaxed) + 1; let avg = total_sum.checked_div(total_count).unwrap_or(0); - let (cpu_user, cpu_sys) = match (self.start_cpu, end_cpu) { - (Some(start), Some(end)) => ( - end.user_usec.saturating_sub(start.user_usec), - end.system_usec.saturating_sub(start.system_usec), - ), - _ => (0, 0), + // CPU: per-process subtree accumulation when we have a root pid (so the + // compile is excluded), otherwise the pod-wide cgroup delta. + let (cpu_user, cpu_sys) = if self.root_pid.is_some() { + ( + self.cpu_user_usec.load(Ordering::Relaxed), + self.cpu_sys_usec.load(Ordering::Relaxed), + ) + } else { + match (self.start_cpu, read_cpu_stat()) { + (Some(start), Some(end)) => ( + end.user_usec.saturating_sub(start.user_usec), + end.system_usec.saturating_sub(start.system_usec), + ), + _ => (0, 0), + } }; let peak_spill = self.peak_spill.load(Ordering::Relaxed); @@ -199,6 +263,147 @@ impl CgroupMonitor { } } +/// Clock ticks (USER_HZ) to microseconds. Linux USER_HZ is 100 on all common +/// configurations, so a tick is 10 ms = 10_000 µs. +fn ticks_to_usec(ticks: u64) -> u64 { + ticks.saturating_mul(10_000) +} + +// --- per-process sampling --- + +/// Compiler / build-tool process names (as they appear in `/proc//status` +/// `Name:` / `/proc//stat` comm, truncated to 15 chars) whose usage should +/// NOT be attributed to the benchmark. A `cargo bench` step spawns these as +/// children of the monitored command; counting them would report compiler +/// memory and CPU as benchmark memory and CPU. +const BUILD_TOOL_NAMES: &[&str] = &[ + "cargo", "rustc", "sccache", "cc", "cc1", "cc1plus", "gcc", "g++", "c++", "clang", "clang++", + "ld", "ld.lld", "lld", "ld.gold", "ld.bfd", "ld.mold", "mold", "collect2", "as", "ar", + "rustdoc", +]; + +fn is_build_tool(name: &str) -> bool { + BUILD_TOOL_NAMES.contains(&name) + || name.starts_with("build-script") + || name.starts_with("build_script") +} + +/// Per-process cumulative CPU sample: `(pid, utime_ticks, stime_ticks)`. +type ProcCpu = (u32, u64, u64); + +/// One subtree snapshot: `(total VmRSS bytes, per-process CPU)`. +type SubtreeSample = (u64, Vec); + +/// Sample current memory: the benchmark process subtree RSS when `root_pid` is +/// known, otherwise the pod-wide cgroup figure. +fn sample_memory(root_pid: Option) -> Option { + match root_pid { + Some(pid) => proc_tree_sample(Some(pid)).map(|(rss, _)| rss), + None => read_memory_current(), + } +} + +/// One snapshot of the process subtree rooted at `root_pid`, excluding +/// compiler/build-tool processes: +/// - total `VmRSS` in bytes +/// - per-process cumulative CPU `(pid, utime_ticks, stime_ticks)` +/// +/// Returns `None` if `/proc` can't be read at all (e.g. macOS), or if +/// `root_pid` is `None`. Returns `Some((0, []))` if the subtree is already gone. +fn proc_tree_sample(root_pid: Option) -> Option { + let root_pid = root_pid?; + let entries = std::fs::read_dir("/proc").ok()?; + + // pid -> (ppid, name, rss_bytes, utime_ticks, stime_ticks) + let mut procs: std::collections::HashMap = std::collections::HashMap::new(); + for entry in entries.flatten() { + let pid: u32 = match entry.file_name().to_str().and_then(|s| s.parse().ok()) { + Some(p) => p, + None => continue, // non-numeric /proc entry + }; + if let Some(info) = read_proc_info(pid) { + procs.insert(pid, info); + } + } + + // children index + let mut children: std::collections::HashMap> = std::collections::HashMap::new(); + for (&pid, info) in &procs { + children.entry(info.ppid).or_default().push(pid); + } + + // BFS from root, collecting non-build-tool processes. + let mut rss_total: u64 = 0; + let mut cpus: Vec<(u32, u64, u64)> = Vec::new(); + let mut stack = vec![root_pid]; + while let Some(pid) = stack.pop() { + if let Some(info) = procs.get(&pid) { + if !is_build_tool(&info.name) { + rss_total += info.rss_bytes; + cpus.push((pid, info.utime_ticks, info.stime_ticks)); + } + } + if let Some(kids) = children.get(&pid) { + stack.extend(kids); + } + } + Some((rss_total, cpus)) +} + +struct ProcInfo { + ppid: u32, + name: String, + rss_bytes: u64, + utime_ticks: u64, + stime_ticks: u64, +} + +/// Read process info from `/proc//status` (ppid, name, VmRSS) and +/// `/proc//stat` (utime, stime). +fn read_proc_info(pid: u32) -> Option { + let status = std::fs::read_to_string(format!("/proc/{pid}/status")).ok()?; + let mut ppid = None; + let mut name = None; + let mut rss_kb = 0u64; + for line in status.lines() { + if let Some(v) = line.strip_prefix("Name:") { + name = Some(v.trim().to_string()); + } else if let Some(v) = line.strip_prefix("PPid:") { + ppid = v.trim().parse().ok(); + } else if let Some(v) = line.strip_prefix("VmRSS:") { + // e.g. "VmRSS: 123456 kB" + rss_kb = v + .split_whitespace() + .next() + .and_then(|n| n.parse().ok()) + .unwrap_or(0); + } + } + + let (utime_ticks, stime_ticks) = read_proc_cpu_ticks(pid).unwrap_or((0, 0)); + + Some(ProcInfo { + ppid: ppid?, + name: name?, + rss_bytes: rss_kb * 1024, + utime_ticks, + stime_ticks, + }) +} + +/// Parse cumulative `(utime, stime)` clock ticks from `/proc//stat`. +/// The comm field (field 2) can contain spaces and parentheses, so split on the +/// last `)`; the remaining whitespace-separated fields start at field 3 +/// (`state`), making utime field 14 → index 11 and stime field 15 → index 12. +fn read_proc_cpu_ticks(pid: u32) -> Option<(u64, u64)> { + let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?; + let after_comm = stat.rsplit_once(')')?.1; + let fields: Vec<&str> = after_comm.split_whitespace().collect(); + let utime = fields.get(11).and_then(|s| s.parse().ok())?; + let stime = fields.get(12).and_then(|s| s.parse().ok())?; + Some((utime, stime)) +} + // --- cgroup v2 file readers --- fn read_memory_current() -> Option { @@ -367,10 +572,64 @@ throttled_usec 0 #[tokio::test] async fn test_monitor_returns_stats() { - let monitor = CgroupMonitor::start(None); + let monitor = CgroupMonitor::start(None, None); tokio::time::sleep(Duration::from_millis(100)).await; let stats = monitor.finish().await; assert!(stats.wall_time >= Duration::from_millis(50)); assert!(stats.sample_count >= 2); } + + #[tokio::test] + async fn test_monitor_self_subtree_has_cpu_and_mem() { + // Monitor our own process subtree; we should observe non-zero RSS and, + // after doing some work, non-zero CPU — proving the /proc sampler runs. + let pid = std::process::id(); + let monitor = CgroupMonitor::start(Some(pid), None); + // Burn a little CPU so utime advances across polls. + let mut x: u64 = 0; + let spin_until = Instant::now() + Duration::from_millis(1200); + while Instant::now() < spin_until { + x = x.wrapping_add(1); + std::hint::black_box(x); + } + let stats = monitor.finish().await; + // On non-/proc platforms (macOS) these are zero; only assert there. + if std::path::Path::new("/proc/self/status").exists() { + assert!(stats.peak_memory_bytes > 0, "expected non-zero RSS"); + assert!( + stats.cpu_user_usec > 0, + "expected non-zero user CPU for the spinning process" + ); + } + } + + #[test] + fn test_is_build_tool() { + assert!(is_build_tool("rustc")); + assert!(is_build_tool("cargo")); + assert!(is_build_tool("sccache")); + assert!(is_build_tool("ld.lld")); + assert!(is_build_tool("build-script-build")); + assert!(is_build_tool("build_script_build")); + assert!(!is_build_tool("dfbench")); + assert!(!is_build_tool("sql-712c5e60f8")); // criterion bench binary + } + + #[test] + fn test_ticks_to_usec() { + assert_eq!(ticks_to_usec(0), 0); + assert_eq!(ticks_to_usec(1), 10_000); // 1 tick = 10ms = 10_000us + assert_eq!(ticks_to_usec(100), 1_000_000); // 100 ticks = 1s + } + + #[test] + fn test_read_proc_cpu_ticks_self() { + // Only meaningful on Linux; skip elsewhere. + if std::path::Path::new("/proc/self/stat").exists() { + let pid = std::process::id(); + let (u, s) = read_proc_cpu_ticks(pid).expect("read self cpu ticks"); + // Cumulative ticks are monotonic and fit in u64; just sanity bound. + assert!(u < u64::MAX && s < u64::MAX); + } + } } diff --git a/controller/src/runner/shell.rs b/controller/src/runner/shell.rs index 17549f9..1ecae56 100644 --- a/controller/src/runner/shell.rs +++ b/controller/src/runner/shell.rs @@ -18,11 +18,18 @@ const DEFAULT_DIAGNOSTIC_AFTER_SECS: u64 = 600; const DEFAULT_DIAGNOSTIC_INTERVAL_SECS: u64 = 300; const PRE_DEADLINE_DIAGNOSTIC_OFFSET_SECS: u64 = 300; -/// Run a command, log it, stream output to the log file, and return stdout as a string. -/// Fails if the command exits with a non-zero status. -pub async fn run_command(cmd: &str, args: &[&str], cwd: &Path) -> Result { - info!(cmd, ?args, ?cwd, "running command"); +/// A spawned command with its captured-output paths, used to split spawning +/// from waiting so callers (e.g. the resource monitor) can observe the PID +/// while the command runs. +struct SpawnedCommand { + child: tokio::process::Child, + pid: Option, + stdout_path: String, + stderr_path: String, +} +/// Spawn a command with stdout/stderr redirected to temp log files. +fn spawn_logged(cmd: &str, args: &[&str], cwd: &Path) -> Result { let temp_id = temp_log_id(); let stdout_path = format!("/tmp/cmd-{temp_id}.stdout"); let stderr_path = format!("/tmp/cmd-{temp_id}.stderr"); @@ -32,7 +39,7 @@ pub async fn run_command(cmd: &str, args: &[&str], cwd: &Path) -> Result let stderr_file = std::fs::File::create(&stderr_path) .with_context(|| format!("failed to create stderr log: {stderr_path}"))?; - let mut child = Command::new(cmd) + let child = Command::new(cmd) .args(args) .current_dir(cwd) .stdout(Stdio::from(stdout_file)) @@ -41,12 +48,28 @@ pub async fn run_command(cmd: &str, args: &[&str], cwd: &Path) -> Result .with_context(|| format!("failed to spawn: {cmd} {}", args.join(" ")))?; let pid = child.id(); - let status = wait_for_child(cmd, args, cwd, &mut child, pid).await?; + Ok(SpawnedCommand { + child, + pid, + stdout_path, + stderr_path, + }) +} + +/// Wait for a spawned command, collect its output into the shared log, and +/// return stdout. Fails if the command exits with a non-zero status. +async fn finish_logged( + cmd: &str, + args: &[&str], + cwd: &Path, + mut spawned: SpawnedCommand, +) -> Result { + let status = wait_for_child(cmd, args, cwd, &mut spawned.child, spawned.pid).await?; - let stdout = tokio::fs::read_to_string(&stdout_path) + let stdout = tokio::fs::read_to_string(&spawned.stdout_path) .await .unwrap_or_default(); - let stderr = tokio::fs::read_to_string(&stderr_path) + let stderr = tokio::fs::read_to_string(&spawned.stderr_path) .await .unwrap_or_default(); @@ -54,8 +77,8 @@ pub async fn run_command(cmd: &str, args: &[&str], cwd: &Path) -> Result append_to_log(&stdout).await; append_to_log(&stderr).await; - let _ = tokio::fs::remove_file(&stdout_path).await; - let _ = tokio::fs::remove_file(&stderr_path).await; + let _ = tokio::fs::remove_file(&spawned.stdout_path).await; + let _ = tokio::fs::remove_file(&spawned.stderr_path).await; if !status.success() { let code = status.code().unwrap_or(-1); @@ -68,7 +91,20 @@ pub async fn run_command(cmd: &str, args: &[&str], cwd: &Path) -> Result Ok(stdout) } -/// Run a command with cgroup resource monitoring. Returns both stdout and resource stats. +/// Run a command, log it, stream output to the log file, and return stdout as a string. +/// Fails if the command exits with a non-zero status. +pub async fn run_command(cmd: &str, args: &[&str], cwd: &Path) -> Result { + info!(cmd, ?args, ?cwd, "running command"); + let spawned = spawn_logged(cmd, args, cwd)?; + finish_logged(cmd, args, cwd, spawned).await +} + +/// Run a command with resource monitoring. Returns both stdout and resource stats. +/// +/// The monitor is scoped to the spawned command's process subtree (see +/// [`CgroupMonitor`]), so the reported memory reflects the benchmark itself +/// rather than the whole pod (which would also count any compilation, +/// page cache, and leftover build artifacts). /// /// If `spill_dir` is provided, the monitor will poll the directory size every /// second to track peak spill usage. @@ -78,10 +114,12 @@ pub async fn run_command_monitored( cwd: &Path, spill_dir: Option, ) -> Result<(String, ResourceStats)> { - let monitor = CgroupMonitor::start(spill_dir); - let output = run_command(cmd, args, cwd).await?; + info!(cmd, ?args, ?cwd, "running command (monitored)"); + let spawned = spawn_logged(cmd, args, cwd)?; + let monitor = CgroupMonitor::start(spawned.pid, spill_dir); + let output = finish_logged(cmd, args, cwd, spawned).await; let stats = monitor.finish().await; - Ok((output, stats)) + Ok((output?, stats)) } /// Spawn a command in the background, returning a JoinHandle that resolves to the Result. From c9a4dec433ae8369416f41efebd1877b46a09de0 Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Tue, 9 Jun 2026 09:51:22 -0500 Subject: [PATCH 2/3] test: make subtree CPU test reliable on Linux CI MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The self-subtree test ran on a single-threaded tokio runtime and spun for 1.2s, so the background poll task was starved (and CPU is a delta needing >=2 samples at the 1s poll interval) — it observed zero CPU and failed on Linux. Use a multi-thread runtime, spin >2s so at least two polls land, and guard the CPU assertion on having enough samples. Co-Authored-By: Claude Opus 4.8 (1M context) --- controller/src/runner/monitor.rs | 22 +++++++++++++++------- 1 file changed, 15 insertions(+), 7 deletions(-) diff --git a/controller/src/runner/monitor.rs b/controller/src/runner/monitor.rs index d92d7d7..d98a33d 100644 --- a/controller/src/runner/monitor.rs +++ b/controller/src/runner/monitor.rs @@ -579,15 +579,18 @@ throttled_usec 0 assert!(stats.sample_count >= 2); } - #[tokio::test] + // Multi-threaded runtime so the background poll task can sample while this + // task burns CPU; CPU is a delta between samples, so the process must be + // sampled at least twice (poll interval is 1s) — hence the >2s spin. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_monitor_self_subtree_has_cpu_and_mem() { // Monitor our own process subtree; we should observe non-zero RSS and, // after doing some work, non-zero CPU — proving the /proc sampler runs. let pid = std::process::id(); let monitor = CgroupMonitor::start(Some(pid), None); - // Burn a little CPU so utime advances across polls. + // Burn CPU so the process's utime advances across poll samples. let mut x: u64 = 0; - let spin_until = Instant::now() + Duration::from_millis(1200); + let spin_until = Instant::now() + Duration::from_millis(2500); while Instant::now() < spin_until { x = x.wrapping_add(1); std::hint::black_box(x); @@ -596,10 +599,15 @@ throttled_usec 0 // On non-/proc platforms (macOS) these are zero; only assert there. if std::path::Path::new("/proc/self/status").exists() { assert!(stats.peak_memory_bytes > 0, "expected non-zero RSS"); - assert!( - stats.cpu_user_usec > 0, - "expected non-zero user CPU for the spinning process" - ); + // CPU needs >= 2 poll samples for a non-zero delta (the first only + // establishes the per-process baseline). sample_count = start + polls + // + end, so >= 3 means at least two polls landed. + if stats.sample_count >= 3 { + assert!( + stats.cpu_user_usec > 0, + "expected non-zero user CPU for the spinning process" + ); + } } } From 3407d2e30f9aad16f0fc320ce0e57c8ddc4dcdb2 Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Tue, 9 Jun 2026 10:03:11 -0500 Subject: [PATCH 3/3] monitor: fall back to cgroup when /proc unavailable, fix CPU accounting Address review feedback on #14: - Fall back to the pod-wide cgroup figures whenever `/proc` subtree sampling fails, not only when the root pid is unknown. A `proc_ok` flag tracks whether any `/proc` sample succeeded; memory and CPU both fall back to cgroup when it stays false (root pid unknown, or `/proc` unreadable in a restricted container). `sample_memory(Some(pid))` now also falls back to `memory.current` instead of returning `None`/0. - Take a final memory+CPU sample once stop is signalled, so CPU is captured up to the end of the run rather than as of the last full poll interval (was undercounting by up to POLL_INTERVAL with no final /proc read). - Read `_SC_CLK_TCK` via `libc::sysconf` (cached) instead of hard-coding USER_HZ=100, so CPU time is correct on non-100Hz kernels. - Update module docs to describe the subtree-or-cgroup behavior for both memory and CPU (CPU no longer "still comes from cgroup cpu.stat"). Co-Authored-By: Claude Opus 4.8 (1M context) --- Cargo.lock | 1 + controller/Cargo.toml | 1 + controller/src/runner/monitor.rs | 133 +++++++++++++++++++++---------- 3 files changed, 95 insertions(+), 40 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 3c305b2..63a10dd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -204,6 +204,7 @@ dependencies = [ "chrono", "k8s-openapi", "kube", + "libc", "logfire", "once_cell", "regex", diff --git a/controller/Cargo.toml b/controller/Cargo.toml index b56c897..9a1a9f7 100644 --- a/controller/Cargo.toml +++ b/controller/Cargo.toml @@ -28,3 +28,4 @@ once_cell = "1" tokio-util = { version = "0.7", features = ["rt"] } backon = "1" serde_yaml = "0.9" +libc = "0.2" diff --git a/controller/src/runner/monitor.rs b/controller/src/runner/monitor.rs index d98a33d..f254fd6 100644 --- a/controller/src/runner/monitor.rs +++ b/controller/src/runner/monitor.rs @@ -1,17 +1,20 @@ //! Resource monitoring for benchmark execution. //! -//! Memory is sampled from the benchmark command's **process subtree** by -//! walking `/proc//status` (summing `VmRSS`) and excluding compiler / -//! build-tool processes (`rustc`, `cargo`, `cc`, `ld`, `sccache`, …). That is -//! what we report as the benchmark's memory. Reading the pod-wide cgroup -//! (`/sys/fs/cgroup/memory.current`) instead would also count page cache, -//! leftover build artifacts, and — for benchmarks whose `bench.sh run` step -//! compiles (e.g. `cargo bench`) — the compiler itself, which can dwarf the -//! actual query memory by tens of GB. +//! When the benchmark command's root pid is known and `/proc` is readable, +//! both memory and CPU are sampled from its **process subtree** by walking +//! `/proc//status` (summing `VmRSS`) and `/proc//stat` (accumulating +//! `utime`/`stime`), excluding compiler / build-tool processes (`rustc`, +//! `cargo`, `cc`, `ld`, `sccache`, …). That is what we report. Reading the +//! pod-wide cgroup (`/sys/fs/cgroup/memory.current`, `cpu.stat`) instead would +//! also count page cache, leftover build artifacts, and — for benchmarks whose +//! `bench.sh run` step compiles (e.g. `cargo bench`) — the compiler itself, +//! which can dwarf the actual query memory and CPU by tens of GB / thousands of +//! CPU-seconds. //! -//! CPU stats still come from cgroup v2 `cpu.stat`. When the per-process root -//! pid is unknown the sampler falls back to the cgroup memory figure. When -//! neither `/proc` nor the cgroup files are available (e.g. local macOS +//! When the root pid is unknown, or `/proc` can't be read (e.g. a restricted +//! container), the sampler falls back to the pod-wide cgroup figures: memory +//! from `memory.current` and CPU from the `cpu.stat` delta over the window. +//! When neither `/proc` nor the cgroup files are available (e.g. local macOS //! development), reads return `None` and stats default to zero. use std::fmt; @@ -113,10 +116,14 @@ pub struct CgroupMonitor { peak_spill: Arc, sample_count: Arc, // Per-process CPU consumed by the benchmark subtree (excluding build - // tools), accumulated in the poll loop. Only meaningful when `root_pid` is + // tools), accumulated in the poll loop. Only meaningful when `proc_ok` is // set; otherwise CPU is taken from the cgroup delta in `finish`. cpu_user_usec: Arc, cpu_sys_usec: Arc, + // Set once any `/proc` subtree sample succeeds. When this stays false (root + // pid unknown, or `/proc` unreadable) memory and CPU both come from the + // pod-wide cgroup instead. + proc_ok: Arc, poll_handle: JoinHandle<()>, } @@ -140,6 +147,7 @@ impl CgroupMonitor { let sample_count = Arc::new(AtomicU64::new(1)); let cpu_user_usec = Arc::new(AtomicU64::new(0)); let cpu_sys_usec = Arc::new(AtomicU64::new(0)); + let proc_ok = Arc::new(AtomicBool::new(false)); let poll_handle = { let stop = stop_flag.clone(); @@ -149,6 +157,7 @@ impl CgroupMonitor { let count = sample_count.clone(); let cpu_user = cpu_user_usec.clone(); let cpu_sys = cpu_sys_usec.clone(); + let proc_ok = proc_ok.clone(); tokio::spawn(async move { // Per-pid first/last cumulative CPU ticks, so each benchmark @@ -159,12 +168,13 @@ impl CgroupMonitor { let mut cpu_last: std::collections::HashMap = std::collections::HashMap::new(); - while !stop.load(Ordering::Relaxed) { - tokio::time::sleep(POLL_INTERVAL).await; - if stop.load(Ordering::Relaxed) { - break; - } + // One memory+CPU sample. Prefers the `/proc` subtree; falls back + // to the pod-wide cgroup figure when `/proc` can't be read + // (root pid unknown, or no `/proc`, e.g. a restricted container). + let sample = |cpu_first: &mut std::collections::HashMap, + cpu_last: &mut std::collections::HashMap| { if let Some((rss, cpus)) = proc_tree_sample(root_pid) { + proc_ok.store(true, Ordering::Relaxed); peak.fetch_max(rss, Ordering::Relaxed); sum.fetch_add(rss, Ordering::Relaxed); count.fetch_add(1, Ordering::Relaxed); @@ -176,7 +186,7 @@ impl CgroupMonitor { // Recompute accumulated subtree CPU each poll. let mut tu = 0u64; let mut ts = 0u64; - for (pid, &(fu, fs)) in &cpu_first { + for (pid, &(fu, fs)) in &*cpu_first { if let Some(&(lu, ls)) = cpu_last.get(pid) { tu += lu.saturating_sub(fu); ts += ls.saturating_sub(fs); @@ -184,19 +194,32 @@ impl CgroupMonitor { } cpu_user.store(ticks_to_usec(tu), Ordering::Relaxed); cpu_sys.store(ticks_to_usec(ts), Ordering::Relaxed); - } else if root_pid.is_none() { - // Fallback path: keep memory sampling via cgroup. - if let Some(current) = read_memory_current() { - peak.fetch_max(current, Ordering::Relaxed); - sum.fetch_add(current, Ordering::Relaxed); - count.fetch_add(1, Ordering::Relaxed); - } + } else if let Some(current) = read_memory_current() { + // Fallback: pod-wide cgroup memory. CPU comes from the + // cgroup delta in `finish` when `proc_ok` is never set. + peak.fetch_max(current, Ordering::Relaxed); + sum.fetch_add(current, Ordering::Relaxed); + count.fetch_add(1, Ordering::Relaxed); + } + }; + + while !stop.load(Ordering::Relaxed) { + tokio::time::sleep(POLL_INTERVAL).await; + if stop.load(Ordering::Relaxed) { + break; } + sample(&mut cpu_first, &mut cpu_last); if let Some(ref dir) = spill_dir { let size = dir_size(dir); spill_peak.fetch_max(size, Ordering::Relaxed); } } + + // Final sample once stop is signalled, so memory and CPU are + // captured right up to the end of the run rather than as of the + // last full poll interval (which could undercount by up to + // POLL_INTERVAL). + sample(&mut cpu_first, &mut cpu_last); }) }; @@ -212,6 +235,7 @@ impl CgroupMonitor { sample_count, cpu_user_usec, cpu_sys_usec, + proc_ok, poll_handle, } } @@ -223,16 +247,22 @@ impl CgroupMonitor { self.stop_flag.store(true, Ordering::Relaxed); let _ = self.poll_handle.await; - let end_memory = sample_memory(self.root_pid).unwrap_or(0); - + // The poll loop's final sample already updated peak/sum/count up to the + // end of the run, so don't take (and double-count) another into them. + let proc_ok = self.proc_ok.load(Ordering::Relaxed); + let end_memory = if proc_ok { + sample_memory(self.root_pid).unwrap_or(0) + } else { + read_memory_current().unwrap_or(0) + }; let peak = self.peak_memory.load(Ordering::Relaxed).max(end_memory); - let total_sum = self.memory_sum.load(Ordering::Relaxed) + end_memory; - let total_count = self.sample_count.load(Ordering::Relaxed) + 1; + let total_sum = self.memory_sum.load(Ordering::Relaxed); + let total_count = self.sample_count.load(Ordering::Relaxed); let avg = total_sum.checked_div(total_count).unwrap_or(0); - // CPU: per-process subtree accumulation when we have a root pid (so the - // compile is excluded), otherwise the pod-wide cgroup delta. - let (cpu_user, cpu_sys) = if self.root_pid.is_some() { + // CPU: per-process subtree accumulation when `/proc` sampling worked (so + // the compile is excluded), otherwise the pod-wide cgroup delta. + let (cpu_user, cpu_sys) = if proc_ok { ( self.cpu_user_usec.load(Ordering::Relaxed), self.cpu_sys_usec.load(Ordering::Relaxed), @@ -263,10 +293,26 @@ impl CgroupMonitor { } } -/// Clock ticks (USER_HZ) to microseconds. Linux USER_HZ is 100 on all common -/// configurations, so a tick is 10 ms = 10_000 µs. +/// Clock ticks per second (`USER_HZ`), read once from `sysconf(_SC_CLK_TCK)`. +/// This is 100 on most Linux configurations but is not guaranteed, so we query +/// it rather than hard-coding it. Falls back to 100 if the query fails. +fn clk_tck() -> u64 { + use std::sync::OnceLock; + static TCK: OnceLock = OnceLock::new(); + *TCK.get_or_init(|| { + let v = unsafe { libc::sysconf(libc::_SC_CLK_TCK) }; + if v > 0 { + v as u64 + } else { + 100 + } + }) +} + +/// Convert cumulative clock ticks (`USER_HZ`) from `/proc//stat` to +/// microseconds, using the actual `_SC_CLK_TCK` value. fn ticks_to_usec(ticks: u64) -> u64 { - ticks.saturating_mul(10_000) + ticks.saturating_mul(1_000_000) / clk_tck() } // --- per-process sampling --- @@ -295,10 +341,12 @@ type ProcCpu = (u32, u64, u64); type SubtreeSample = (u64, Vec); /// Sample current memory: the benchmark process subtree RSS when `root_pid` is -/// known, otherwise the pod-wide cgroup figure. +/// known and `/proc` is readable, otherwise the pod-wide cgroup figure. fn sample_memory(root_pid: Option) -> Option { match root_pid { - Some(pid) => proc_tree_sample(Some(pid)).map(|(rss, _)| rss), + Some(pid) => proc_tree_sample(Some(pid)) + .map(|(rss, _)| rss) + .or_else(read_memory_current), None => read_memory_current(), } } @@ -576,7 +624,10 @@ throttled_usec 0 tokio::time::sleep(Duration::from_millis(100)).await; let stats = monitor.finish().await; assert!(stats.wall_time >= Duration::from_millis(50)); - assert!(stats.sample_count >= 2); + // The start baseline always counts as one sample; additional samples + // require a readable data source (cgroup or `/proc`), which is absent on + // e.g. local macOS development. + assert!(stats.sample_count >= 1); } // Multi-threaded runtime so the background poll task can sample while this @@ -625,9 +676,11 @@ throttled_usec 0 #[test] fn test_ticks_to_usec() { + let tck = clk_tck(); + assert!(tck > 0); assert_eq!(ticks_to_usec(0), 0); - assert_eq!(ticks_to_usec(1), 10_000); // 1 tick = 10ms = 10_000us - assert_eq!(ticks_to_usec(100), 1_000_000); // 100 ticks = 1s + // `clk_tck` ticks = exactly one second, whatever USER_HZ happens to be. + assert_eq!(ticks_to_usec(tck), 1_000_000); } #[test]