Skip to main content

nmbrs_runtime/
sysmon.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Session-level system-performance sampler.
5//!
6//! Enabled per session with `sysmon=<categories>` — `sysmon=all`,
7//! `sysmon=any`, or a comma list of `cpu,io,ram,rambw,storage`. `all` and
8//! `any` differ in exactly one way: `all` ABORTS when a subsystem is
9//! unavailable, `any` enables what the host supports and NAMES what it
10//! skipped — an explicit opt-in to best-effort, so the skip is announced
11//! rather than silent. Reads `/proc` (and statvfs)
12//! on a fixed interval (default 5 s, `sysmon-interval=<seconds>`) and
13//! publishes host-side utilization as session-scoped gauges plus a live
14//! observer callback for display surfaces: measurement of the machine
15//! under the workload, not of the workload.
16//!
17//! The categories:
18//!
19//! - **cpu** — `/proc/stat`. Three separate measures: the MEAN utilization
20//!   from the aggregate `cpu ` line, the MAXIMUM single-core saturation
21//!   from the `cpuN` lines (with the core named), and the QUARTILES of the
22//!   per-core distribution (`p25`/`p50`/`p75`). A pinned compaction thread
23//!   can saturate one core while the mean reads 3% — both facts matter and
24//!   neither substitutes for the other.
25//!
26//!   The quartiles exist because neither the mean nor the max answers "is
27//!   there headroom". The max says 100% for a single hot thread on an idle
28//!   box; the mean says 3% for a box where every core is at 3% and for a
29//!   box with one pegged core, without distinguishing them. `p50` low with
30//!   `max_core` pegged is one hot thread with headroom to spare; `p50` high
31//!   is genuine saturation. Any adaptive guard that throttles on CPU wants
32//!   the quartile, not the max. The per-core vector is already built to
33//!   find the max, so the shape costs one sort per window.
34//! - **io** — `/proc/diskstats`, every line. Utilization per device is
35//!   `Δio_ticks / Δwall` (io_ticks is the 10th stat field: milliseconds the
36//!   device had I/O in flight — the same quantity `iostat -x %util`
37//!   reports). Only the HIGHEST device utilization is recorded, with the
38//!   device named. On NVMe this saturates as an idle-detector rather than a
39//!   capacity meter (dozens of requests run in parallel), which is exactly
40//!   what a "is the disk the bottleneck" glance wants.
41//! - **ram** — `/proc/meminfo`. Two separate measures:
42//!   `(MemTotal − MemAvailable) / MemTotal` (committed: what is actually
43//!   claimed and not readily reclaimable — the kernel's own estimate) and
44//!   `(MemTotal − MemFree) / MemTotal` (everything, page cache included).
45//!   On a database host the second sits near 100% by design; the pair reads
46//!   as "how much is spoken for" vs "how much is touched".
47//! - **rambw** — memory bandwidth via resctrl MBM
48//!   (`/sys/fs/resctrl/mon_data`). REQUIRED-EXPLICIT when enabled: if the
49//!   resctrl interface is not mounted (or `sysmon-membw-gbps=` is not set to
50//!   provide the peak reference), the session ABORTS with instructions —
51//!   never a silent skip. `sysmon=all` includes it, so `all` on a host
52//!   without resctrl aborts; a host without it runs
53//!   `sysmon=cpu,io,ram,storage`.
54//! - **storage** — filesystem SPACE utilization: statvfs over every
55//!   `/dev/`-backed mount in `/proc/mounts` (deduplicated by source device),
56//!   `1 − available/total` per mount, only the highest recorded, mount
57//!   point named. Space is the disk measure `io` cannot see: a device can
58//!   be I/O-idle and one write from full.
59//!
60//! Counters are cumulative, so every rate-like utilization here is a
61//! pairwise delta over the sample window; the published gauge is the latest
62//! window's value and any windowing beyond that belongs to MetricsQL at
63//! query time.
64
65use std::sync::Arc;
66use std::sync::RwLock;
67use std::time::Duration;
68
69use nmbrs_metrics::component::Component;
70use nmbrs_metrics::instruments::gauge::ValueGauge;
71
72/// Which categories a `sysmon=` setting enabled.
73#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
74pub struct Categories {
75    pub cpu: bool,
76    pub io: bool,
77    pub ram: bool,
78    pub rambw: bool,
79    pub storage: bool,
80}
81
82impl Categories {
83    pub const ALL: Categories = Categories {
84        cpu: true,
85        io: true,
86        ram: true,
87        rambw: true,
88        storage: true,
89    };
90
91    pub fn any(&self) -> bool {
92        self.cpu || self.io || self.ram || self.rambw || self.storage
93    }
94}
95
96/// What a `sysmon=` value asked for: a fixed category set, or "whatever
97/// this host supports".
98#[derive(Debug, Clone, Copy, PartialEq, Eq)]
99pub enum Selection {
100    /// Named categories (including `all`): every one is REQUIRED, and an
101    /// unavailable subsystem aborts the session with instructions.
102    Cats(Categories),
103    /// `sysmon=any`: enable every available subsystem, skip-and-announce
104    /// the rest. The stated opt-in is what makes the skip acceptable.
105    Any,
106}
107
108/// Parse a `sysmon=` value: `all`, `any`, or a comma list of category
109/// names. Unknown names are errors that NAME the valid set — a typo
110/// silently monitoring nothing would be the failure mode this surface
111/// exists to avoid.
112pub fn parse_selection(value: &str) -> Result<Selection, String> {
113    let trimmed = value.trim();
114    if trimmed.eq_ignore_ascii_case("any") {
115        return Ok(Selection::Any);
116    }
117    if trimmed
118        .to_ascii_lowercase()
119        .split(',')
120        .any(|t| t.trim() == "any")
121    {
122        return Err("sysmon: `any` stands alone (it means \"every available \
123             subsystem\") — combining it with named categories is \
124             ambiguous. Use `sysmon=any`, or name the categories."
125            .to_string());
126    }
127    parse_categories(trimmed).map(Selection::Cats)
128}
129
130/// Parse a fixed category list (`all` or comma names). See
131/// [`parse_selection`] for the `any` form.
132pub fn parse_categories(value: &str) -> Result<Categories, String> {
133    if value.trim().eq_ignore_ascii_case("all") {
134        return Ok(Categories::ALL);
135    }
136    let mut cats = Categories::default();
137    for token in value.split(',') {
138        match token.trim().to_ascii_lowercase().as_str() {
139            "cpu" => cats.cpu = true,
140            "io" => cats.io = true,
141            "ram" => cats.ram = true,
142            "rambw" => cats.rambw = true,
143            "storage" => cats.storage = true,
144            "" => {}
145            other => {
146                return Err(format!(
147                    "sysmon: unknown category '{other}'. Valid: all, or a \
148                     comma list of cpu, io, ram, rambw, storage"
149                ));
150            }
151        }
152    }
153    if !cats.any() {
154        return Err(
155            "sysmon: no categories enabled. Use `sysmon=all` or a comma \
156             list of cpu, io, ram, rambw, storage"
157                .to_string(),
158        );
159    }
160    Ok(cats)
161}
162
163/// Per-category CPU readings: the mean, the hottest core, and the
164/// SHAPE of the per-core distribution.
165///
166/// The hottest core alone cannot answer "is there headroom": one
167/// pegged core on an otherwise idle box reads identically to a box
168/// that is genuinely out of CPU. Measured 2026-08-21 on a live
169/// `stcs_adaptive` run — `max_core` p90 = 1.000 while `mean` sat at
170/// 0.183 — so any guard keyed on the max would throttle at 18% host
171/// utilization. The quartiles separate the two cases: `p50` low with
172/// `max_core` pegged is a single hot thread (headroom remains); `p50`
173/// high is real saturation.
174///
175/// `max_core` is kept, not replaced — it remains the right signal for
176/// spotting a single-threaded bottleneck, which is a different
177/// question from headroom.
178#[derive(Debug, Clone, Copy, PartialEq, Default)]
179pub struct CpuReading {
180    pub mean: f64,
181    pub max_core: f64,
182    pub top_core: usize,
183    /// Per-core utilization quartiles across all online cores.
184    pub p25: f64,
185    pub p50: f64,
186    pub p75: f64,
187    /// How many cores the quartiles were computed over — the sample
188    /// size behind the shape, so a reader can tell a 2-core box's
189    /// "quartiles" from a 96-core box's.
190    pub cores: usize,
191}
192
193/// One completed sample window. Each field is `Some` exactly when its
194/// category was enabled — a disabled category is absent, not zero.
195#[derive(Debug, Clone, PartialEq, Default)]
196pub struct SysmonSample {
197    pub cpu: Option<CpuReading>,
198    /// (device, utilization) — the busiest device this window.
199    pub io: Option<(String, f64)>,
200    /// (committed, everything-including-page-cache).
201    pub ram: Option<(f64, f64)>,
202    /// Memory-bandwidth utilization against the configured peak.
203    pub rambw: Option<f64>,
204    /// (mount point, space utilization) — the fullest filesystem.
205    pub storage: Option<(String, f64)>,
206}
207
208/// Sampler configuration, resolved by the runner from session params.
209#[derive(Debug, Clone)]
210pub struct SysmonConfig {
211    pub cats: Categories,
212    /// Sample window. Default 5 s (`sysmon-interval=<seconds>`).
213    pub interval: Duration,
214    /// Peak memory bandwidth in bytes/s — the reference that turns rambw
215    /// bytes/s into a utilization. Required when `rambw` is enabled.
216    pub membw_peak_bytes_per_s: Option<f64>,
217}
218
219/// The rambw prerequisites, checked BEFORE the session starts so an
220/// unsupported host aborts with instructions instead of silently monitoring
221/// less than was asked for.
222pub fn check_rambw_requirements(config: &SysmonConfig) -> Result<(), String> {
223    if !config.cats.rambw {
224        return Ok(());
225    }
226    if read_membw_bytes().is_none() {
227        return Err("\
228sysmon: rambw was requested, but the kernel resctrl interface is not \
229available at /sys/fs/resctrl/mon_data.
230
231To enable memory-bandwidth monitoring:
232  1. The CPU must support bandwidth monitoring (Intel RDT / AMD QoS) —
233     check for the `cqm_mbm_total` flag:  grep -m1 cqm_mbm_total /proc/cpuinfo
234  2. The kernel must be built with CONFIG_X86_CPU_RESCTRL (standard on
235     mainstream distro kernels).
236  3. Mount the interface:  sudo mount -t resctrl resctrl /sys/fs/resctrl
237
238If this host cannot support it (most VMs cannot), run without the rambw
239category:  sysmon=cpu,io,ram,storage  — or use  sysmon=any  to enable
240every subsystem this host supports."
241            .to_string());
242    }
243    if config.membw_peak_bytes_per_s.is_none() {
244        return Err("\
245sysmon: rambw needs a peak-bandwidth reference to turn bytes/s into a \
246utilization. Set it with  sysmon-membw-gbps=<peak>  (the host's rated \
247memory bandwidth in GB/s), or run without the rambw category:  \
248sysmon=cpu,io,ram,storage  — or use  sysmon=any"
249            .to_string());
250    }
251    Ok(())
252}
253
254/// Resolve `sysmon=any` against THIS host: every category the host
255/// supports, plus one human-readable line per category that had to be
256/// skipped and why. cpu/io/ram/storage are /proc-backed and always
257/// available on Linux; rambw carries real prerequisites.
258///
259/// The skip lines exist because `any` is best-effort, not silent-effort:
260/// the runner logs each one, so a session that monitored four of five
261/// categories says so.
262pub fn resolve_any(config: &SysmonConfig) -> (Categories, Vec<String>) {
263    let mut cats = Categories::ALL;
264    let mut skipped = Vec::new();
265    let rambw_probe = SysmonConfig {
266        cats,
267        ..config.clone()
268    };
269    if let Err(reason) = check_rambw_requirements(&rambw_probe) {
270        cats.rambw = false;
271        // First line of the full instructions — the one that names the
272        // missing prerequisite. `sysmon=rambw` gets the complete text.
273        let first = reason.lines().next().unwrap_or("unavailable").to_string();
274        skipped.push(format!("{first} (run `sysmon=rambw` for the enable steps)"));
275    }
276    (cats, skipped)
277}
278
279// ---------------------------------------------------------------------------
280// Pure parsers + delta math. Everything below reads STRINGS so the arithmetic
281// is testable on fixtures; only the spawn loop touches the filesystem.
282// ---------------------------------------------------------------------------
283
284/// Per-device cumulative `io_ticks` (ms with I/O in flight) from
285/// `/proc/diskstats`. Field layout per the kernel's Documentation/iostats:
286/// `major minor name <17 stat fields>`; io_ticks is stat field 10, i.e.
287/// whitespace token 12.
288pub fn parse_diskstats(text: &str) -> Vec<(String, u64)> {
289    text.lines()
290        .filter_map(|line| {
291            let t: Vec<&str> = line.split_whitespace().collect();
292            let name = t.get(2)?;
293            let io_ticks: u64 = t.get(12)?.parse().ok()?;
294            Some((name.to_string(), io_ticks))
295        })
296        .collect()
297}
298
299/// The highest per-device utilization between two diskstats snapshots taken
300/// `dt_ms` apart. Devices present in only one snapshot are skipped (hotplug
301/// between samples); an empty intersection yields `None`.
302pub fn max_disk_util(
303    prev: &[(String, u64)],
304    cur: &[(String, u64)],
305    dt_ms: f64,
306) -> Option<(String, f64)> {
307    if dt_ms <= 0.0 {
308        return None;
309    }
310    let mut best: Option<(String, f64)> = None;
311    for (name, cur_ticks) in cur {
312        let Some((_, prev_ticks)) = prev.iter().find(|(n, _)| n == name) else {
313            continue;
314        };
315        let util = (cur_ticks.saturating_sub(*prev_ticks) as f64 / dt_ms).clamp(0.0, 1.0);
316        if best.as_ref().is_none_or(|(_, b)| util > *b) {
317            best = Some((name.clone(), util));
318        }
319    }
320    best
321}
322
323/// Cumulative jiffy counts for one `cpu` line: (busy, total).
324///
325/// total = user+nice+system+idle+iowait+irq+softirq+steal (the first eight
326/// fields; guest time is already accounted inside user/nice). busy = total −
327/// idle − iowait: iowait is idle-with-an-excuse, and counting it as busy
328/// would read an I/O-bound stall as CPU saturation.
329#[derive(Debug, Clone, Copy, PartialEq, Eq)]
330pub struct CpuTicks {
331    pub busy: u64,
332    pub total: u64,
333}
334
335/// Aggregate + per-core cumulative ticks from `/proc/stat`.
336pub fn parse_proc_stat(text: &str) -> Option<(CpuTicks, Vec<CpuTicks>)> {
337    let mut aggregate: Option<CpuTicks> = None;
338    let mut cores: Vec<(usize, CpuTicks)> = Vec::new();
339    for line in text.lines() {
340        let mut t = line.split_whitespace();
341        let Some(name) = t.next() else { continue };
342        if !name.starts_with("cpu") {
343            continue;
344        }
345        let fields: Vec<u64> = t.filter_map(|f| f.parse().ok()).collect();
346        if fields.len() < 8 {
347            continue;
348        }
349        let total: u64 = fields[..8].iter().sum();
350        let idle = fields[3] + fields[4];
351        let ticks = CpuTicks {
352            busy: total - idle,
353            total,
354        };
355        if name == "cpu" {
356            aggregate = Some(ticks);
357        } else if let Ok(n) = name[3..].parse::<usize>() {
358            cores.push((n, ticks));
359        }
360    }
361    cores.sort_by_key(|(n, _)| *n);
362    aggregate.map(|a| (a, cores.into_iter().map(|(_, t)| t).collect()))
363}
364
365/// Utilization between two tick snapshots: Δbusy / Δtotal.
366pub fn cpu_util(prev: CpuTicks, cur: CpuTicks) -> f64 {
367    let dt = cur.total.saturating_sub(prev.total);
368    if dt == 0 {
369        return 0.0;
370    }
371    (cur.busy.saturating_sub(prev.busy) as f64 / dt as f64).clamp(0.0, 1.0)
372}
373
374/// The most saturated core between two snapshots. Core lists of different
375/// lengths (offline/online between samples) compare over the shared prefix.
376pub fn max_core_util(prev: &[CpuTicks], cur: &[CpuTicks]) -> Option<(usize, f64)> {
377    prev.iter()
378        .zip(cur.iter())
379        .enumerate()
380        .map(|(i, (p, c))| (i, cpu_util(*p, *c)))
381        .max_by(|a, b| a.1.total_cmp(&b.1))
382}
383
384/// Per-core utilization quartiles between two snapshots: `(p25, p50,
385/// p75, cores)`. Same shared-prefix rule as [`max_core_util`].
386///
387/// The per-core vector is already built to find the max; the quartiles
388/// are a sort of that same vector, so the shape costs one sort per
389/// sample window and no extra reads of `/proc/stat`.
390///
391/// Nearest-rank on the sorted ascending vector (no interpolation): the
392/// quantile of a core count is a real core's utilization, which is the
393/// honest reading when the population is small — an interpolated "p75"
394/// across 4 cores would name a utilization no core actually had.
395pub fn core_util_quartiles(prev: &[CpuTicks], cur: &[CpuTicks]) -> Option<(f64, f64, f64, usize)> {
396    let mut utils: Vec<f64> = prev
397        .iter()
398        .zip(cur.iter())
399        .map(|(p, c)| cpu_util(*p, *c))
400        .collect();
401    if utils.is_empty() {
402        return None;
403    }
404    utils.sort_by(f64::total_cmp);
405    let n = utils.len();
406    let at = |q: f64| -> f64 {
407        // Nearest-rank, clamped: ceil(q*n) - 1 over 1-based ranks.
408        let rank = (q * n as f64).ceil().max(1.0) as usize;
409        utils[rank.min(n) - 1]
410    };
411    Some((at(0.25), at(0.50), at(0.75), n))
412}
413
414/// The three `/proc/meminfo` fields the two utilization measures need,
415/// in kB as the kernel reports them.
416#[derive(Debug, Clone, Copy, PartialEq, Eq)]
417pub struct MemInfo {
418    pub total_kb: u64,
419    pub free_kb: u64,
420    pub available_kb: u64,
421}
422
423pub fn parse_meminfo(text: &str) -> Option<MemInfo> {
424    let mut total = None;
425    let mut free = None;
426    let mut available = None;
427    for line in text.lines() {
428        let mut t = line.split_whitespace();
429        match t.next() {
430            Some("MemTotal:") => total = t.next()?.parse().ok(),
431            Some("MemFree:") => free = t.next()?.parse().ok(),
432            Some("MemAvailable:") => available = t.next()?.parse().ok(),
433            _ => {}
434        }
435    }
436    Some(MemInfo {
437        total_kb: total?,
438        free_kb: free?,
439        available_kb: available?,
440    })
441}
442
443/// The two memory measures: (committed, everything-including-page-cache).
444pub fn mem_utils(m: MemInfo) -> (f64, f64) {
445    if m.total_kb == 0 {
446        return (0.0, 0.0);
447    }
448    let committed = (m.total_kb.saturating_sub(m.available_kb)) as f64 / m.total_kb as f64;
449    let cached = (m.total_kb.saturating_sub(m.free_kb)) as f64 / m.total_kb as f64;
450    (committed.clamp(0.0, 1.0), cached.clamp(0.0, 1.0))
451}
452
453/// WRITABLE `/dev/`-backed mount points from `/proc/mounts`, deduplicated by
454/// source device (bind mounts and btrfs subvolumes re-list one device many
455/// times; space is a per-DEVICE fact).
456///
457/// Read-only mounts are excluded, and it matters: snap images are
458/// `/dev/loop*` squashfs mounts that are 100% full BY CONSTRUCTION, so one
459/// installed snap would pin the storage item at bright-orange forever.
460/// Verified on this host — `/dev/loop0 /snap/... ro,...` at 100% while the
461/// fullest writable filesystem sat at 40%. A read-only filesystem cannot
462/// fill up, so its fullness is not a utilization.
463pub fn parse_dev_mounts(text: &str) -> Vec<String> {
464    let mut seen_sources: Vec<&str> = Vec::new();
465    let mut mounts = Vec::new();
466    for line in text.lines() {
467        let mut t = line.split_whitespace();
468        let (Some(source), Some(mount), _fstype, Some(options)) =
469            (t.next(), t.next(), t.next(), t.next())
470        else {
471            continue;
472        };
473        if !source.starts_with("/dev/") || seen_sources.contains(&source) {
474            continue;
475        }
476        let read_only = options.split(',').any(|o| o == "ro");
477        if read_only {
478            continue;
479        }
480        seen_sources.push(source);
481        // /proc/mounts octal-escapes spaces in mount points (\040).
482        mounts.push(mount.replace("\\040", " "));
483    }
484    mounts
485}
486
487/// Space utilization of one filesystem: `1 − available/total`, matching what
488/// `df` calls Use%. `None` on statvfs failure or a zero-block pseudo-fs.
489/// Unix-only, like its caller [`max_storage_util`] (which walks
490/// `/proc/mounts` and therefore never yields a mount elsewhere).
491#[cfg(unix)]
492fn statvfs_util(mount: &str) -> Option<f64> {
493    let c_mount = std::ffi::CString::new(mount).ok()?;
494    let mut vfs: libc::statvfs = unsafe { std::mem::zeroed() };
495    if unsafe { libc::statvfs(c_mount.as_ptr(), &mut vfs) } != 0 {
496        return None;
497    }
498    if vfs.f_blocks == 0 {
499        return None;
500    }
501    Some((1.0 - vfs.f_bavail as f64 / vfs.f_blocks as f64).clamp(0.0, 1.0))
502}
503
504/// The fullest `/dev/`-backed filesystem right now: (mount point, util).
505fn max_storage_util() -> Option<(String, f64)> {
506    #[cfg(unix)]
507    {
508        let mounts = std::fs::read_to_string("/proc/mounts")
509            .map(|t| parse_dev_mounts(&t))
510            .unwrap_or_default();
511        mounts
512            .into_iter()
513            .filter_map(|m| statvfs_util(&m).map(|u| (m, u)))
514            .max_by(|a, b| a.1.total_cmp(&b.1))
515    }
516    // No `/proc/mounts` off Unix — storage-utilization sampling
517    // is simply unavailable there.
518    #[cfg(not(unix))]
519    {
520        None
521    }
522}
523
524/// Sum of resctrl MBM total-bytes counters across mon_data groups, when the
525/// resctrl filesystem is mounted with monitoring. `None` when unavailable.
526fn read_membw_bytes() -> Option<u64> {
527    let root = std::path::Path::new("/sys/fs/resctrl/mon_data");
528    let entries = std::fs::read_dir(root).ok()?;
529    let mut sum: u64 = 0;
530    let mut seen = false;
531    for e in entries.flatten() {
532        let f = e.path().join("mbm_total_bytes");
533        if let Ok(text) = std::fs::read_to_string(&f)
534            && let Ok(v) = text.trim().parse::<u64>()
535        {
536            sum += v;
537            seen = true;
538        }
539    }
540    seen.then_some(sum)
541}
542
543// ---------------------------------------------------------------------------
544// The sampler task.
545// ---------------------------------------------------------------------------
546
547/// A gauge family whose series are one-per-SUBJECT — the device, core, or
548/// mount the measurement is about. Each distinct subject value materialises
549/// a dimensional cell under the session component
550/// ([`nmbrs_metrics::cells::resolve_under`]) and registers the family there
551/// once; after that a sample is a hash lookup and a `set`.
552///
553/// This is what "submitted to metrics with the appropriate dimensional
554/// labels, keyed by the session component" means concretely:
555/// `sysmon_io_util{session="…",device="nvme1n1"}` — the subject is a label,
556/// the cell refines the session's identity, and a sweep that shifts between
557/// devices yields one honestly-labeled series per device rather than one
558/// anonymous series whose subject silently changes.
559struct SubjectGauge {
560    family: &'static str,
561    /// The dimension this family's subject occupies (`device`, `core`,
562    /// `mount`).
563    label_key: &'static str,
564    parent: Arc<RwLock<Component>>,
565    /// Instruments already materialised, by subject value.
566    instances: std::collections::HashMap<String, Arc<ValueGauge>>,
567}
568
569impl SubjectGauge {
570    fn new(family: &'static str, label_key: &'static str, parent: Arc<RwLock<Component>>) -> Self {
571        Self {
572            family,
573            label_key,
574            parent,
575            instances: Default::default(),
576        }
577    }
578
579    fn set(&mut self, subject: &str, value: f64) {
580        if let Some(g) = self.instances.get(subject) {
581            g.set(value);
582            return;
583        }
584        let coord = nmbrs_metrics::labels::Labels::of(self.label_key, subject);
585        let cell = nmbrs_metrics::cells::resolve_under(&self.parent, &coord);
586        let labels = {
587            let guard = cell.read().unwrap_or_else(|e| e.into_inner());
588            guard.effective_labels().clone()
589        };
590        let g = Arc::new(ValueGauge::new(labels.with("family", self.family)));
591        let registered = cell
592            .write()
593            .unwrap_or_else(|e| e.into_inner())
594            .register_instrument(
595                self.family,
596                nmbrs_metrics::component::InstrumentRef::Gauge(g.clone()),
597            );
598        if let Err(e) = registered {
599            // One cell, one family: a second registration here is a
600            // programming error worth hearing about once, not per sample.
601            crate::diag!(
602                crate::observer::LogLevel::Warn,
603                "sysmon: {family} cell {subject}: {e}",
604                family = self.family
605            );
606        }
607        g.set(value);
608        self.instances.insert(subject.to_string(), g);
609    }
610}
611
612struct Gauges {
613    /// Host-scalar measures — no subject, so they live on the session root
614    /// itself and carry exactly its labels.
615    cpu_mean: Option<Arc<ValueGauge>>,
616    /// Per-core distribution shape. Plain (not subject-dimensioned)
617    /// gauges: they describe the HOST, the same dimensional cell as
618    /// `cpu_mean`. Fanning them out per core — as `cpu_core_max` does
619    /// via its `core` label — would scatter one number across as many
620    /// instances as the box has cores and make aggregation a join.
621    cpu_core_p25: Option<Arc<ValueGauge>>,
622    cpu_core_p50: Option<Arc<ValueGauge>>,
623    cpu_core_p75: Option<Arc<ValueGauge>>,
624    ram_committed: Option<Arc<ValueGauge>>,
625    ram_cached: Option<Arc<ValueGauge>>,
626    rambw: Option<Arc<ValueGauge>>,
627    /// Subject-dimensioned measures — one cell per device / core / mount.
628    io: Option<SubjectGauge>,
629    cpu_core_max: Option<SubjectGauge>,
630    storage: Option<SubjectGauge>,
631}
632
633/// Register gauges for the ENABLED categories on the session component.
634/// Direct registration on the session root — no child component, no new
635/// labels: the samples describe the whole host, which is exactly the
636/// session's dimensional cell.
637fn register_gauges(component: &Arc<RwLock<Component>>, cats: Categories) -> Result<Gauges, String> {
638    let mut guard = component.write().unwrap_or_else(|e| e.into_inner());
639    let labels = guard.effective_labels().clone();
640    let mut mk = |family: &str| -> Result<Option<Arc<ValueGauge>>, String> {
641        let g = Arc::new(ValueGauge::new(labels.with("family", family)));
642        guard.register_instrument(
643            family,
644            nmbrs_metrics::component::InstrumentRef::Gauge(g.clone()),
645        )?;
646        Ok(Some(g))
647    };
648    let mut gauges = Gauges {
649        cpu_mean: None,
650        cpu_core_p25: None,
651        cpu_core_p50: None,
652        cpu_core_p75: None,
653        ram_committed: None,
654        ram_cached: None,
655        rambw: None,
656        io: None,
657        cpu_core_max: None,
658        storage: None,
659    };
660    if cats.cpu {
661        gauges.cpu_mean = mk("sysmon_cpu_util")?;
662        gauges.cpu_core_p25 = mk("sysmon_cpu_core_p25")?;
663        gauges.cpu_core_p50 = mk("sysmon_cpu_core_p50")?;
664        gauges.cpu_core_p75 = mk("sysmon_cpu_core_p75")?;
665    }
666    if cats.ram {
667        gauges.ram_committed = mk("sysmon_ram_util")?;
668        gauges.ram_cached = mk("sysmon_ram_util_cached")?;
669    }
670    if cats.rambw {
671        gauges.rambw = mk("sysmon_rambw_util")?;
672    }
673    drop(guard);
674    // Subject-dimensioned families register per cell at first sight of each
675    // subject, NOT here — registering on the root as well would claim the
676    // family for the un-refined identity and collide with the first cell.
677    if cats.io {
678        gauges.io = Some(SubjectGauge::new(
679            "sysmon_io_util",
680            "device",
681            component.clone(),
682        ));
683    }
684    if cats.cpu {
685        gauges.cpu_core_max = Some(SubjectGauge::new(
686            "sysmon_cpu_core_max",
687            "core",
688            component.clone(),
689        ));
690    }
691    if cats.storage {
692        gauges.storage = Some(SubjectGauge::new(
693            "sysmon_storage_util",
694            "mount",
695            component.clone(),
696        ));
697    }
698    Ok(gauges)
699}
700
701/// Spawn the sampler. `check_rambw_requirements` must have passed first —
702/// the runner aborts the session on its Err rather than calling this.
703/// Runs until session shutdown; publishes each window to the session gauges
704/// and to `observer.sysmon_update`.
705pub fn spawn(
706    config: SysmonConfig,
707    component: Arc<RwLock<Component>>,
708    observer: Arc<dyn crate::observer::RunObserver>,
709) -> Result<tokio::task::JoinHandle<()>, String> {
710    check_rambw_requirements(&config)?;
711    let mut gauges = register_gauges(&component, config.cats)?;
712    let cats = config.cats;
713
714    let mut shutdown = crate::session_signals::subscribe_shutdown();
715    Ok(tokio::spawn(async move {
716        let mut prev_disks = cats
717            .io
718            .then(|| {
719                std::fs::read_to_string("/proc/diskstats")
720                    .map(|t| parse_diskstats(&t))
721                    .unwrap_or_default()
722            })
723            .unwrap_or_default();
724        let mut prev_cpu = if cats.cpu {
725            std::fs::read_to_string("/proc/stat")
726                .ok()
727                .and_then(|t| parse_proc_stat(&t))
728        } else {
729            None
730        };
731        let mut prev_membw = if cats.rambw { read_membw_bytes() } else { None };
732        let mut prev_at = std::time::Instant::now();
733
734        loop {
735            tokio::select! {
736                _ = tokio::time::sleep(config.interval) => {}
737                _ = shutdown.changed() => break,
738            }
739            let now = std::time::Instant::now();
740            let dt = now.duration_since(prev_at);
741            let dt_ms = dt.as_secs_f64() * 1000.0;
742            prev_at = now;
743            let mut sample = SysmonSample::default();
744
745            if cats.io {
746                let cur = std::fs::read_to_string("/proc/diskstats")
747                    .map(|t| parse_diskstats(&t))
748                    .unwrap_or_default();
749                sample.io = max_disk_util(&prev_disks, &cur, dt_ms);
750                prev_disks = cur;
751            }
752            if cats.cpu {
753                let cur = std::fs::read_to_string("/proc/stat")
754                    .ok()
755                    .and_then(|t| parse_proc_stat(&t));
756                if let (Some((pa, pc)), Some((ca, cc))) = (&prev_cpu, &cur) {
757                    let (top_core, max_core) = max_core_util(pc, cc).unwrap_or((0, 0.0));
758                    let (p25, p50, p75, cores) =
759                        core_util_quartiles(pc, cc).unwrap_or((0.0, 0.0, 0.0, 0));
760                    sample.cpu = Some(CpuReading {
761                        mean: cpu_util(*pa, *ca),
762                        max_core,
763                        top_core,
764                        p25,
765                        p50,
766                        p75,
767                        cores,
768                    });
769                }
770                prev_cpu = cur;
771            }
772            if cats.ram {
773                sample.ram = std::fs::read_to_string("/proc/meminfo")
774                    .ok()
775                    .and_then(|t| parse_meminfo(&t))
776                    .map(mem_utils);
777            }
778            if cats.rambw
779                && let Some(peak) = config.membw_peak_bytes_per_s
780            {
781                let cur = read_membw_bytes();
782                if let (Some(p), Some(c)) = (prev_membw, cur) {
783                    sample.rambw = Some(
784                        ((c.saturating_sub(p)) as f64 / dt.as_secs_f64().max(1e-9) / peak)
785                            .clamp(0.0, 1.0),
786                    );
787                }
788                prev_membw = cur;
789            }
790            if cats.storage {
791                sample.storage = max_storage_util();
792            }
793
794            if let (Some(g), Some((dev, u))) = (&mut gauges.io, &sample.io) {
795                g.set(dev, *u);
796            }
797            if let (Some(g), Some(c)) = (&gauges.cpu_mean, &sample.cpu) {
798                g.set(c.mean);
799            }
800            if let (Some(g), Some(c)) = (&gauges.cpu_core_p25, &sample.cpu) {
801                g.set(c.p25);
802            }
803            if let (Some(g), Some(c)) = (&gauges.cpu_core_p50, &sample.cpu) {
804                g.set(c.p50);
805            }
806            if let (Some(g), Some(c)) = (&gauges.cpu_core_p75, &sample.cpu) {
807                g.set(c.p75);
808            }
809            if let (Some(g), Some(c)) = (&mut gauges.cpu_core_max, &sample.cpu) {
810                g.set(&c.top_core.to_string(), c.max_core);
811            }
812            if let (Some(g), Some((committed, _))) = (&gauges.ram_committed, &sample.ram) {
813                g.set(*committed);
814            }
815            if let (Some(g), Some((_, cached))) = (&gauges.ram_cached, &sample.ram) {
816                g.set(*cached);
817            }
818            if let (Some(g), Some(u)) = (&gauges.rambw, sample.rambw) {
819                g.set(u);
820            }
821            if let (Some(g), Some((mount, u))) = (&mut gauges.storage, &sample.storage) {
822                g.set(mount, *u);
823            }
824
825            observer.sysmon_update(&sample);
826        }
827    }))
828}
829
830#[cfg(test)]
831mod tests {
832    use super::*;
833
834    /// `all` and the full comma list mean the same thing — the user's words.
835    #[test]
836    fn all_equals_the_full_category_list() {
837        assert_eq!(parse_categories("all").unwrap(), Categories::ALL);
838        assert_eq!(
839            parse_categories("cpu,io,ram,rambw,storage").unwrap(),
840            Categories::ALL
841        );
842    }
843
844    #[test]
845    fn category_subsets_parse_and_typos_are_named_errors() {
846        let c = parse_categories("cpu, io").unwrap();
847        assert!(c.cpu && c.io && !c.ram && !c.rambw && !c.storage);
848        let err = parse_categories("cpu,ramb").unwrap_err();
849        assert!(
850            err.contains("ramb") && err.contains("rambw"),
851            "the error names the typo and the valid set: {err}"
852        );
853        assert!(parse_categories("").is_err(), "empty enables nothing");
854    }
855
856    /// rambw enabled on a host without resctrl is an ABORT with instructions,
857    /// not a skip. (This box has no resctrl, so this exercises the real
858    /// probe.)
859    #[test]
860    fn rambw_without_resctrl_aborts_with_instructions() {
861        let config = SysmonConfig {
862            cats: parse_categories("rambw").unwrap(),
863            interval: Duration::from_secs(5),
864            membw_peak_bytes_per_s: Some(100e9),
865        };
866        if std::path::Path::new("/sys/fs/resctrl/mon_data").exists() {
867            // Host actually has resctrl — the gate passes instead; nothing
868            // to assert about instructions here.
869            return;
870        }
871        let err = check_rambw_requirements(&config).unwrap_err();
872        assert!(
873            err.contains("mount -t resctrl"),
874            "the abort must tell the user HOW to enable it: {err}"
875        );
876        assert!(
877            err.contains("sysmon=cpu,io,ram,storage"),
878            "…and how to run without it: {err}"
879        );
880    }
881
882    /// rambw with resctrl but no peak reference is equally an abort — a
883    /// bytes/s figure with no denominator is not a utilization.
884    #[test]
885    fn rambw_without_a_peak_reference_aborts() {
886        let config = SysmonConfig {
887            cats: parse_categories("rambw").unwrap(),
888            interval: Duration::from_secs(5),
889            membw_peak_bytes_per_s: None,
890        };
891        let err = check_rambw_requirements(&config).unwrap_err();
892        assert!(
893            err.contains("sysmon-membw-gbps") || err.contains("resctrl"),
894            "must name the missing prerequisite: {err}"
895        );
896    }
897
898    /// Disabled rambw asks nothing of the host.
899    #[test]
900    fn no_rambw_no_requirements() {
901        let config = SysmonConfig {
902            cats: parse_categories("cpu,io,ram,storage").unwrap(),
903            interval: Duration::from_secs(5),
904            membw_peak_bytes_per_s: None,
905        };
906        assert!(check_rambw_requirements(&config).is_ok());
907    }
908
909    /// `any` stands alone; combined with names it is ambiguous and refused.
910    #[test]
911    fn any_parses_alone_and_refuses_combination() {
912        assert_eq!(parse_selection("any").unwrap(), Selection::Any);
913        assert_eq!(parse_selection("ANY").unwrap(), Selection::Any);
914        assert!(parse_selection("any,cpu").is_err());
915        assert!(matches!(parse_selection("all").unwrap(),
916            Selection::Cats(c) if c == Categories::ALL));
917    }
918
919    /// On a host without resctrl, `any` yields everything but rambw and
920    /// SAYS SO; the skip line points at the full instructions. (This box
921    /// has no resctrl; a host that has it passes the gate instead.)
922    #[test]
923    fn any_downgrades_rambw_with_an_announced_reason() {
924        let config = SysmonConfig {
925            cats: Categories::ALL,
926            interval: Duration::from_secs(5),
927            membw_peak_bytes_per_s: None,
928        };
929        let (cats, skipped) = resolve_any(&config);
930        assert!(cats.cpu && cats.io && cats.ram && cats.storage);
931        if std::path::Path::new("/sys/fs/resctrl/mon_data").exists() {
932            return; // host genuinely supports it; nothing to skip here
933        }
934        assert!(!cats.rambw, "unavailable rambw is disabled under `any`");
935        assert_eq!(skipped.len(), 1);
936        assert!(
937            skipped[0].contains("sysmon=rambw"),
938            "the skip points at the full enable steps: {}",
939            skipped[0]
940        );
941    }
942
943    /// A subject-dimensioned gauge materialises ONE cell per subject under
944    /// the session component, each series carrying the subject as a label —
945    /// and re-setting an existing subject attaches no twin.
946    #[test]
947    fn subject_gauges_dimension_by_cell_under_the_session_component() {
948        let session = Arc::new(RwLock::new(nmbrs_metrics::component::Component::new(
949            nmbrs_metrics::labels::Labels::of("session", "s1"),
950            std::collections::HashMap::new(),
951        )));
952        let mut g = SubjectGauge::new("sysmon_io_util", "device", session.clone());
953        g.set("nvme1n1", 0.97);
954        g.set("nvme2n1", 0.40);
955        g.set("nvme1n1", 0.98);
956
957        let guard = session.read().unwrap();
958        assert_eq!(
959            guard.child_count(),
960            2,
961            "one cell per device, re-sets attach no twin"
962        );
963        drop(guard);
964        assert_eq!(g.instances.len(), 2);
965        // The series carries session + device + family — the session's
966        // identity refined by the subject, never replaced.
967        let labels = g.instances["nvme1n1"].labels().to_prometheus();
968        for owned in ["session=", "device=\"nvme1n1\"", "family="] {
969            assert!(labels.contains(owned), "{owned} missing from {labels}");
970        }
971    }
972
973    /// Real lines from this host's /proc/diskstats — io_ticks is token 12.
974    #[test]
975    fn diskstats_reads_io_ticks_from_the_tenth_stat_field() {
976        let text = "\
977 259       0 nvme0n1 1083809 77036 79071293 1310486 4905083 3299234 633830668 26973654 0 3344862 28284141 0 0 0 0 0 0
978 259       5 nvme1n1 16478852962 684755 139660303464 4041025697 191545776 3994164 30038357911 3695409309 0 588447592 3442404703 29888 5913 38808742648 936992 0 0";
979        let parsed = parse_diskstats(text);
980        assert_eq!(
981            parsed,
982            vec![
983                ("nvme0n1".to_string(), 3_344_862),
984                ("nvme1n1".to_string(), 588_447_592),
985            ]
986        );
987    }
988
989    /// Highest utilization wins and is named; a device busy 1000ms of a
990    /// 2000ms window is 50%.
991    #[test]
992    fn disk_util_is_the_max_across_devices() {
993        let prev = vec![("a".to_string(), 1000_u64), ("b".to_string(), 5000)];
994        let cur = vec![("a".to_string(), 1400), ("b".to_string(), 6000)];
995        let (name, util) = max_disk_util(&prev, &cur, 2000.0).unwrap();
996        assert_eq!(name, "b");
997        assert!((util - 0.5).abs() < 1e-9);
998    }
999
1000    /// Mean and max-core are separate measures: one saturated core among
1001    /// idle ones must not disappear into the mean.
1002    #[test]
1003    fn one_hot_core_shows_in_max_not_mean() {
1004        let stat_t0 = "\
1005cpu  1000 0 0 8000 0 0 0 0 0 0
1006cpu0 1000 0 0 0 0 0 0 0 0 0
1007cpu1 0 0 0 4000 0 0 0 0 0 0
1008cpu2 0 0 0 4000 0 0 0 0 0 0";
1009        let stat_t1 = "\
1010cpu  2000 0 0 16000 0 0 0 0 0 0
1011cpu0 2000 0 0 0 0 0 0 0 0 0
1012cpu1 0 0 0 8000 0 0 0 0 0 0
1013cpu2 0 0 0 8000 0 0 0 0 0 0";
1014        let (a0, c0) = parse_proc_stat(stat_t0).unwrap();
1015        let (a1, c1) = parse_proc_stat(stat_t1).unwrap();
1016        let mean = cpu_util(a0, a1);
1017        let (core, max) = max_core_util(&c0, &c1).unwrap();
1018        assert!((mean - 1000.0 / 9000.0).abs() < 1e-9, "mean {mean}");
1019        assert_eq!(core, 0);
1020        assert!((max - 1.0).abs() < 1e-9, "core0 fully busy, got {max}");
1021
1022        // The quartiles are what separate THIS box — one pegged core, two
1023        // idle — from a box where every core is genuinely busy. Both read
1024        // max_core = 1.0; only the median tells them apart.
1025        let (p25, p50, p75, cores) = core_util_quartiles(&c0, &c1).unwrap();
1026        assert_eq!(cores, 3);
1027        assert!((p25 - 0.0).abs() < 1e-9, "p25 {p25}");
1028        assert!((p50 - 0.0).abs() < 1e-9, "median core is IDLE, got {p50}");
1029        assert!((p75 - 1.0).abs() < 1e-9, "p75 {p75}");
1030    }
1031
1032    /// The saturation case the quartiles must NOT confuse with one hot
1033    /// core: every core busy reads max_core = 1.0 just like a single
1034    /// pegged core, but the median moves with it.
1035    #[test]
1036    fn quartiles_separate_broad_saturation_from_one_hot_core() {
1037        let stat_t0 = "\
1038cpu  0 0 0 12000 0 0 0 0 0 0
1039cpu0 0 0 0 4000 0 0 0 0 0 0
1040cpu1 0 0 0 4000 0 0 0 0 0 0
1041cpu2 0 0 0 4000 0 0 0 0 0 0";
1042        let stat_t1 = "\
1043cpu  12000 0 0 12000 0 0 0 0 0 0
1044cpu0 4000 0 0 4000 0 0 0 0 0 0
1045cpu1 4000 0 0 4000 0 0 0 0 0 0
1046cpu2 4000 0 0 4000 0 0 0 0 0 0";
1047        let (_, c0) = parse_proc_stat(stat_t0).unwrap();
1048        let (_, c1) = parse_proc_stat(stat_t1).unwrap();
1049        let (_, max) = max_core_util(&c0, &c1).unwrap();
1050        let (p25, p50, p75, cores) = core_util_quartiles(&c0, &c1).unwrap();
1051        assert_eq!(cores, 3);
1052        // EXACTLY the max the one-hot-core case reports — 1.0 there, 1.0
1053        // here. On the max alone the two boxes are indistinguishable.
1054        assert!((max - 1.0).abs() < 1e-9, "max {max}");
1055        // The median is what tells them apart: 0.0 when one core is hot
1056        // and the rest idle, 1.0 when the whole box is saturated.
1057        assert!((p25 - 1.0).abs() < 1e-9, "p25 {p25}");
1058        assert!((p50 - 1.0).abs() < 1e-9, "p50 {p50}");
1059        assert!((p75 - 1.0).abs() < 1e-9, "p75 {p75}");
1060    }
1061
1062    /// Nearest-rank, not interpolation: every quartile of a small core
1063    /// count must be a utilization some real core actually had.
1064    #[test]
1065    fn quartiles_are_nearest_rank_over_real_cores() {
1066        let single = vec![CpuTicks {
1067            busy: 0,
1068            total: 1000,
1069        }];
1070        let single_end = vec![CpuTicks {
1071            busy: 250,
1072            total: 2000,
1073        }];
1074        let (p25, p50, p75, cores) = core_util_quartiles(&single, &single_end).unwrap();
1075        assert_eq!(cores, 1);
1076        // One core: every quantile is that core.
1077        assert!((p25 - 0.25).abs() < 1e-9);
1078        assert!((p50 - 0.25).abs() < 1e-9);
1079        assert!((p75 - 0.25).abs() < 1e-9);
1080        // No cores at all is None, not a fabricated zero.
1081        assert!(core_util_quartiles(&[], &[]).is_none());
1082    }
1083
1084    /// iowait is idle-with-an-excuse: an I/O-bound stall must not read as
1085    /// CPU saturation.
1086    #[test]
1087    fn iowait_does_not_count_as_busy() {
1088        let t0 = parse_proc_stat("cpu 0 0 0 0 0 0 0 0 0 0").unwrap().0;
1089        let t1 = parse_proc_stat("cpu 0 0 0 500 500 0 0 0 0 0").unwrap().0;
1090        assert_eq!(cpu_util(t0, t1), 0.0);
1091    }
1092
1093    /// The two memory measures diverge exactly by reclaimable cache.
1094    #[test]
1095    fn committed_and_cached_measures_are_distinct() {
1096        let m =
1097            parse_meminfo("MemTotal: 1000 kB\nMemFree: 100 kB\nMemAvailable: 600 kB\n").unwrap();
1098        let (committed, cached) = mem_utils(m);
1099        assert!((committed - 0.4).abs() < 1e-9, "claimed = 1 - avail/total");
1100        assert!((cached - 0.9).abs() < 1e-9, "touched = 1 - free/total");
1101    }
1102
1103    /// A window with no elapsed time or no shared devices yields nothing
1104    /// rather than a fabricated zero.
1105    #[test]
1106    fn degenerate_windows_yield_none() {
1107        let d = vec![("a".to_string(), 5_u64)];
1108        assert!(max_disk_util(&d, &d, 0.0).is_none());
1109        assert!(max_disk_util(&[], &d, 1000.0).is_none());
1110    }
1111
1112    /// Only WRITABLE /dev/-backed mounts count for storage, deduplicated by
1113    /// source. Pseudo-filesystems are not disks; a read-only squashfs snap
1114    /// image is 100% full by construction and would pin the item forever.
1115    #[test]
1116    fn dev_mounts_are_filtered_and_deduplicated() {
1117        let mounts = "\
1118proc /proc proc rw 0 0
1119/dev/nvme0n1p1 / ext4 rw 0 0
1120tmpfs /tmp tmpfs rw 0 0
1121/dev/loop0 /snap/core22/2045 squashfs ro,nodev 0 0
1122/dev/nvme1n1 /mnt/nvme xfs rw 0 0
1123/dev/nvme1n1 /mnt/alias xfs rw 0 0
1124/dev/mapper/vg-data /data\\040dir ext4 rw 0 0";
1125        assert_eq!(
1126            parse_dev_mounts(mounts),
1127            vec![
1128                "/".to_string(),
1129                "/mnt/nvme".to_string(),
1130                "/data dir".to_string(),
1131            ]
1132        );
1133    }
1134}