Skip to main content

amont_runtime/
proctree.rs

1//! How much CPU a check's process tree is spending — the second half of
2//! "is this silent check stuck, or just quiet?" (ADR-0008, `hooks.liveness`).
3//!
4//! The silence budget used to be the whole test: a tool that printed nothing
5//! for `amont.idleTimeout` was killed. vitest without a terminal prints its
6//! summary and nothing before it, so a suite with every test passing looked
7//! exactly like a hang. What a hang does NOT do is burn CPU, so this module
8//! measures that, and the caller keeps a silent check alive while its tree is
9//! measurably working.
10//!
11//! Three facts shape everything below, each measured, not assumed:
12//!
13//! - **A sum over LIVE processes goes down while a suite works.** A test
14//!   runner that forks a worker per file loses each worker's CPU the moment
15//!   it exits. So every process is counted WITH the children it has already
16//!   reaped (Linux `cutime`/`cstime`, macOS `ri_child_*`), which only grows.
17//! - **Reaping moves CPU from child to parent.** A child credited while alive
18//!   would be credited again when its parent's reaped-children counter jumps.
19//!   [`Tracker`] subtracts a vanished process's last known total from its
20//!   nearest surviving ancestor, so old work is never counted twice — however
21//!   late the reaping, at any depth.
22//! - **macOS reports mach ticks, not nanoseconds.** On Apple Silicon one tick
23//!   is 41.67 ns; read as nanoseconds a busy suite looks ~42x idle. The
24//!   timebase is applied on every read.
25//!
26//! No process is spawned and no crate is added (the crate is dependency-free,
27//! decisions:ADR-0020): Linux reads `/proc/<pid>/stat`, macOS calls libproc,
28//! which `std` already links. Everything else gets [`Snapshot::Unavailable`]
29//! and the caller falls back to the silence-only rule.
30//!
31//! A snapshot is BOUNDED — a process cap, a retry cap and a wall deadline
32//! ([`Limits`]) — and anything that hits a bound is [`Snapshot::Partial`],
33//! which earns no credit. The figures are a heuristic answering "is anything
34//! working?", not CPU accounting.
35
36use std::collections::{HashMap, HashSet, VecDeque};
37use std::time::{Duration, Instant};
38
39/// A process, as distinct from its pid: pids are reused, and a reused pid
40/// must not inherit the old process's CPU. `start` is the kernel's own start
41/// stamp (Linux `starttime`, macOS `ri_proc_start_abstime`) — compared for
42/// equality only, never converted to wall time.
43#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
44pub struct Id {
45    pub pid: u32,
46    pub start: u64,
47}
48
49/// One process in a snapshot: who it is, whose child it is, and the CPU it
50/// and every child it has reaped have used, in nanoseconds.
51#[derive(Clone, Copy, Debug, PartialEq, Eq)]
52pub struct Proc {
53    pub id: Id,
54    pub ppid: u32,
55    pub cpu_ns: u64,
56}
57
58/// What one look at the tree produced. Only `Complete` is ever measured
59/// from: a snapshot that stopped short would read a missing process as zero
60/// work, which is the one error this module must never make in the
61/// "stuck" direction.
62#[derive(Debug, PartialEq, Eq)]
63pub enum Snapshot {
64    Complete(Vec<Proc>),
65    /// A bound was hit, or part of the tree could not be listed.
66    Partial,
67    /// The root itself could not be read (gone, or not measurable here).
68    Unavailable,
69}
70
71/// The bounds on one snapshot. The sampler runs on its own thread, so these
72/// protect nothing about the ceiling clock; they keep one look at a huge or
73/// churning tree from turning into an unbounded one.
74#[derive(Clone, Copy, Debug)]
75pub struct Limits {
76    pub max_procs: usize,
77    pub max_retries: u32,
78    pub deadline: Duration,
79}
80
81impl Default for Limits {
82    fn default() -> Self {
83        Limits {
84            max_procs: 4096,
85            max_retries: 2,
86            // Half the sampler's stop bound. 100 ms was too tight on a
87            // loaded machine: the walk itself was descheduled and every
88            // sample came back partial (ADR-0009).
89            deadline: Duration::from_millis(500),
90        }
91    }
92}
93
94/// What `read` found at a pid.
95pub enum Read {
96    Proc(Proc),
97    /// Exited between being listed and being read — skipped, not an error.
98    Gone,
99    /// Could not be read for another reason.
100    Failed,
101}
102
103/// Walk the tree under `root`, plus every `seen` process still alive, within
104/// `limits`. Platform code supplies how to list a pid's children and how to
105/// read one pid; the walk, the bounds and the reuse guard live here so they
106/// are tested on every OS.
107///
108/// `read` gets the pid and the parent it was reached from (macOS reads no
109/// ppid of its own). A `seen` process whose pid now carries a different
110/// start stamp was reused and is skipped. A `seen` process reparented away
111/// keeps counting, with its children — which is the only way an orphaned
112/// worker's CPU stays in view.
113pub fn walk(
114    root: u32,
115    seen: &[Id],
116    limits: &Limits,
117    clock: &mut dyn FnMut() -> Instant,
118    children_of: &mut dyn FnMut(u32) -> Option<Vec<u32>>,
119    read: &mut dyn FnMut(u32, u32) -> Read,
120) -> Snapshot {
121    let started = clock();
122    let root_proc = match read(root, 0) {
123        Read::Proc(p) => p,
124        Read::Gone | Read::Failed => return Snapshot::Unavailable,
125    };
126    let mut out = vec![root_proc];
127    let mut have: HashSet<u32> = HashSet::from([root]);
128    let mut queue: VecDeque<u32> = VecDeque::from([root]);
129    let mut orphans = seen.iter().filter(|id| id.pid != root);
130    loop {
131        let parent = match queue.pop_front() {
132            Some(p) => p,
133            None => {
134                // The tree under the root is done; pick up the next seen
135                // process that this walk has not reached.
136                let Some(id) = orphans.by_ref().find(|id| !have.contains(&id.pid)) else {
137                    break;
138                };
139                if clock().duration_since(started) > limits.deadline {
140                    return Snapshot::Partial;
141                }
142                match read(id.pid, 1) {
143                    Read::Proc(p) if p.id == *id => {
144                        have.insert(p.id.pid);
145                        out.push(p);
146                        if out.len() > limits.max_procs {
147                            return Snapshot::Partial;
148                        }
149                        queue.push_back(p.id.pid);
150                    }
151                    Read::Proc(_) | Read::Gone => {}
152                    Read::Failed => return Snapshot::Partial,
153                }
154                continue;
155            }
156        };
157        let Some(kids) = children_of(parent) else {
158            return Snapshot::Partial;
159        };
160        for kid in kids {
161            if !have.insert(kid) {
162                continue;
163            }
164            if clock().duration_since(started) > limits.deadline {
165                return Snapshot::Partial;
166            }
167            match read(kid, parent) {
168                Read::Proc(p) => {
169                    out.push(p);
170                    if out.len() > limits.max_procs {
171                        return Snapshot::Partial;
172                    }
173                    queue.push_back(kid);
174                }
175                Read::Gone => {}
176                Read::Failed => return Snapshot::Partial,
177            }
178        }
179    }
180    Snapshot::Complete(out)
181}
182
183/// Parse one `/proc/<pid>/stat` into a [`Proc`]. Bytes, not a string: the
184/// command name is up to 15 arbitrary bytes and may hold `)`, spaces, a
185/// newline or non-UTF-8, so the fields are found after the LAST `)`.
186/// After it: state [0], ppid [1], utime/stime/cutime/cstime [11..=14],
187/// starttime [19]. cutime/cstime are printed signed; a negative one is 0.
188pub fn parse_proc_stat(stat: &[u8], tick_ns: u64) -> Option<Proc> {
189    let open = stat.iter().position(|&b| b == b'(')?;
190    let close = stat.iter().rposition(|&b| b == b')')?;
191    let pid: u32 = std::str::from_utf8(&stat[..open])
192        .ok()?
193        .trim()
194        .parse()
195        .ok()?;
196    let rest = std::str::from_utf8(stat.get(close + 1..)?).ok()?;
197    let f: Vec<&str> = rest.split_ascii_whitespace().collect();
198    let ppid: u32 = f.get(1)?.parse().ok()?;
199    let unsigned = |i: usize| -> Option<u64> { f.get(i)?.parse().ok() };
200    let signed = |i: usize| -> Option<u64> {
201        let v: i64 = f.get(i)?.parse().ok()?;
202        Some(u64::try_from(v).unwrap_or(0))
203    };
204    let ticks = unsigned(11)?
205        .checked_add(unsigned(12)?)?
206        .checked_add(signed(13)?)?
207        .checked_add(signed(14)?)?;
208    Some(Proc {
209        id: Id {
210            pid,
211            start: unsigned(19)?,
212        },
213        ppid,
214        cpu_ns: ticks.checked_mul(tick_ns)?,
215    })
216}
217
218/// How much work one window between two complete snapshots showed.
219#[derive(Clone, Copy, Debug, PartialEq, Eq)]
220pub struct Window {
221    pub start: Instant,
222    pub end: Instant,
223    pub gain_ns: u64,
224    /// One or more partial snapshots fell inside this window. Work done by
225    /// a child born and reaped during the gap is missed, so a gapped window
226    /// may count as busy but never as measured idle (ADR-0009).
227    pub gapped: bool,
228}
229
230impl Window {
231    /// The window's work in thousandths of one core (1000 = one core busy
232    /// for the whole window).
233    pub fn milli_cores(&self) -> u32 {
234        let wall = self.end.duration_since(self.start).as_nanos().max(1);
235        let milli = u128::from(self.gain_ns) * 1000 / wall;
236        u32::try_from(milli).unwrap_or(u32::MAX)
237    }
238}
239
240/// What one sample told the caller.
241#[derive(Clone, Copy, Debug, PartialEq, Eq)]
242pub enum Observation {
243    /// A complete snapshot with nothing to compare against: the first one,
244    /// or the first after an incomplete one. Credits nothing.
245    Baseline,
246    /// Two consecutive complete snapshots, and the work between them.
247    Window(Window),
248    /// The snapshot was incomplete, but the last complete one is kept: the
249    /// next complete snapshot spans the gap as a window, so one slow walk
250    /// on a loaded machine costs an interval, not the baseline.
251    Skipped,
252    /// Nothing can be measured from here: the root could not be read, or
253    /// too many snapshots in a row came back partial. The next complete one
254    /// starts over as a baseline.
255    Unmeasured,
256}
257
258/// How many partial snapshots in a row the tracker spans before it gives
259/// the baseline up: past this, "the next complete one" is not coming.
260pub const PARTIALS_BEFORE_RESET: u32 = 3;
261
262/// Turns successive snapshots into windows of work, remembering which
263/// processes it has seen so a reparented worker keeps counting.
264#[derive(Default)]
265pub struct Tracker {
266    prev: Option<(Instant, HashMap<Id, Proc>)>,
267    /// Partial snapshots since the last complete one.
268    partials: u32,
269}
270
271impl Tracker {
272    /// The processes of the last complete snapshot: the next walk re-reads
273    /// them even if they are no longer under the root.
274    pub fn seen(&self) -> Vec<Id> {
275        self.prev
276            .as_ref()
277            .map(|(_, m)| m.keys().copied().collect())
278            .unwrap_or_default()
279    }
280
281    pub fn observe(&mut self, now: Instant, snap: Snapshot) -> Observation {
282        let procs = match snap {
283            Snapshot::Complete(procs) => procs,
284            Snapshot::Partial => {
285                self.partials += 1;
286                if self.partials >= PARTIALS_BEFORE_RESET {
287                    self.prev = None;
288                    self.partials = 0;
289                    return Observation::Unmeasured;
290                }
291                return Observation::Skipped;
292            }
293            Snapshot::Unavailable => {
294                self.prev = None;
295                self.partials = 0;
296                return Observation::Unmeasured;
297            }
298        };
299        let gapped = self.partials > 0;
300        self.partials = 0;
301        let cur: HashMap<Id, Proc> = procs.into_iter().map(|p| (p.id, p)).collect();
302        let Some((then, prev)) = self.prev.replace((now, cur)) else {
303            return Observation::Baseline;
304        };
305        let cur = &self.prev.as_ref().expect("just replaced").1;
306        Observation::Window(Window {
307            start: then,
308            end: now,
309            gain_ns: gain_ns(&prev, cur),
310            gapped,
311        })
312    }
313}
314
315/// The CPU spent between two complete snapshots.
316///
317/// - A process in both: what it gained, never negative (macOS zeroes the
318///   reaped-children times on `exec`; a drop is not negative work).
319/// - A process only in `cur`: new since `prev` (nothing joins a tree
320///   later except by being born into it), so all of it counts.
321/// - A process only in `prev` vanished. Its last known total was already
322///   credited; its nearest ancestor still present absorbs that total when it
323///   reaps it, so the total is subtracted from that ancestor's gain. Walking
324///   up through ancestors that vanished too is what keeps a deep tree reaped
325///   bottom-up from being counted once per level.
326///
327/// What stays heuristic: a process born and reaped between two samples was
328/// never seen, so its whole CPU lands on its parent whenever the reaping
329/// happens — at most one interval of its work, possibly in a later window.
330pub fn gain_ns(prev: &HashMap<Id, Proc>, cur: &HashMap<Id, Proc>) -> u64 {
331    let by_pid: HashMap<u32, &Proc> = prev.values().map(|p| (p.id.pid, p)).collect();
332    let mut transferred: HashMap<Id, u64> = HashMap::new();
333    for gone in prev.values().filter(|p| !cur.contains_key(&p.id)) {
334        let mut up = gone.ppid;
335        let mut hops = 0;
336        while let Some(parent) = by_pid.get(&up) {
337            if cur.contains_key(&parent.id) {
338                *transferred.entry(parent.id).or_default() += gone.cpu_ns;
339                break;
340            }
341            up = parent.ppid;
342            hops += 1;
343            if hops > prev.len() {
344                break; // a ppid cycle cannot happen; do not trust that
345            }
346        }
347    }
348    let mut gain: u64 = 0;
349    for p in cur.values() {
350        let add = match prev.get(&p.id) {
351            Some(old) => p
352                .cpu_ns
353                .saturating_sub(old.cpu_ns)
354                .saturating_sub(transferred.get(&p.id).copied().unwrap_or(0)),
355            None => p.cpu_ns,
356        };
357        gain = gain.saturating_add(add);
358    }
359    gain
360}
361
362/// One snapshot of the tree under `root` on this platform.
363pub fn snapshot(root: u32, seen: &[Id], limits: &Limits) -> Snapshot {
364    platform::snapshot(root, seen, limits)
365}
366
367/// Whether [`snapshot`] can measure anything on this platform at all.
368pub const SUPPORTED: bool = cfg!(any(target_os = "linux", target_os = "macos"));
369
370#[cfg(target_os = "linux")]
371mod platform {
372    use super::{parse_proc_stat, walk, Id, Limits, Proc, Read, Snapshot};
373    use std::collections::HashMap;
374    use std::time::Instant;
375
376    // Externs rather than a dependency: `scripts/check-no-deps.sh` keeps the
377    // crate crate-free, and `sysconf` takes and returns plain integers.
378    extern "C" {
379        #[link_name = "sysconf"]
380        fn libc_sysconf_raw(name: i32) -> std::os::raw::c_long;
381    }
382    const SC_CLK_TCK: i32 = 2; // the same value in glibc and musl
383
384    fn tick_ns() -> u64 {
385        // SAFETY: sysconf takes an integer and returns one; no pointers.
386        let hz = unsafe { libc_sysconf_raw(SC_CLK_TCK) };
387        let hz = if hz > 0 { hz as u64 } else { 100 };
388        1_000_000_000 / hz
389    }
390
391    pub(super) fn snapshot(root: u32, seen: &[Id], limits: &Limits) -> Snapshot {
392        let started = Instant::now();
393        let tick = tick_ns();
394        // One pass over /proc: every process, keyed by pid, with children
395        // listed per parent. Only `stat` is read — it never blocks on a
396        // process in uninterruptible sleep, unlike `cmdline`.
397        let Ok(dir) = std::fs::read_dir("/proc") else {
398            return Snapshot::Unavailable;
399        };
400        let mut procs: HashMap<u32, Proc> = HashMap::new();
401        let mut children: HashMap<u32, Vec<u32>> = HashMap::new();
402        for entry in dir.flatten() {
403            if started.elapsed() > limits.deadline {
404                return Snapshot::Partial;
405            }
406            let name = entry.file_name();
407            let Some(pid) = name.to_str().and_then(|n| n.parse::<u32>().ok()) else {
408                continue;
409            };
410            let Ok(bytes) = std::fs::read(format!("/proc/{pid}/stat")) else {
411                continue; // exited since the directory was listed
412            };
413            if let Some(p) = parse_proc_stat(&bytes, tick) {
414                children.entry(p.ppid).or_default().push(pid);
415                procs.insert(pid, p);
416            }
417        }
418        walk(
419            root,
420            seen,
421            limits,
422            &mut || started.max(Instant::now()),
423            &mut |pid| Some(children.get(&pid).cloned().unwrap_or_default()),
424            &mut |pid, _| match procs.get(&pid) {
425                Some(p) => Read::Proc(*p),
426                None => Read::Gone,
427            },
428        )
429    }
430}
431
432#[cfg(target_os = "macos")]
433mod platform {
434    use super::{walk, Id, Limits, Proc, Read, Snapshot};
435    use std::time::Instant;
436
437    /// `struct rusage_info_v2` from `<sys/resource.h>`: a 16-byte uuid then
438    /// eighteen `uint64_t`. Times are in MACH TICKS, not nanoseconds.
439    #[repr(C)]
440    #[derive(Default)]
441    struct RusageInfoV2 {
442        uuid: [u8; 16],
443        user_time: u64,
444        system_time: u64,
445        pkg_idle_wkups: u64,
446        interrupt_wkups: u64,
447        pageins: u64,
448        wired_size: u64,
449        resident_size: u64,
450        phys_footprint: u64,
451        proc_start_abstime: u64,
452        proc_exit_abstime: u64,
453        child_user_time: u64,
454        child_system_time: u64,
455        child_pkg_idle_wkups: u64,
456        child_interrupt_wkups: u64,
457        child_pageins: u64,
458        child_elapsed_abstime: u64,
459        diskio_bytesread: u64,
460        diskio_byteswritten: u64,
461    }
462    const _: () = assert!(std::mem::size_of::<RusageInfoV2>() == 160);
463    const RUSAGE_INFO_V2: i32 = 2;
464
465    #[repr(C)]
466    #[derive(Default)]
467    struct MachTimebaseInfo {
468        numer: u32,
469        denom: u32,
470    }
471
472    // libproc and mach live in libSystem, which std already links.
473    extern "C" {
474        #[link_name = "proc_pid_rusage"]
475        fn proc_pid_rusage_raw(pid: i32, flavor: i32, buffer: *mut RusageInfoV2) -> i32;
476        #[link_name = "proc_listchildpids"]
477        fn proc_listchildpids_raw(ppid: i32, buffer: *mut i32, buffersize: i32) -> i32;
478        #[link_name = "mach_timebase_info"]
479        fn mach_timebase_info_raw(info: *mut MachTimebaseInfo) -> i32;
480    }
481
482    fn timebase() -> (u128, u128) {
483        let mut tb = MachTimebaseInfo::default();
484        // SAFETY: writes one two-u32 struct we own.
485        let rc = unsafe { mach_timebase_info_raw(&mut tb) };
486        if rc != 0 || tb.denom == 0 {
487            return (1, 1);
488        }
489        (u128::from(tb.numer), u128::from(tb.denom))
490    }
491
492    fn rusage(pid: u32) -> Option<RusageInfoV2> {
493        let mut ri = RusageInfoV2::default();
494        let pid = i32::try_from(pid).ok()?;
495        // SAFETY: the buffer is a correctly sized, owned `rusage_info_v2`
496        // (size asserted above) for the flavor we name.
497        let rc = unsafe { proc_pid_rusage_raw(pid, RUSAGE_INFO_V2, &mut ri) };
498        (rc == 0).then_some(ri)
499    }
500
501    /// Children of `pid`, growing the buffer when it came back full, at most
502    /// `retries` times. `proc_listchildpids` returns a COUNT of pids (unlike
503    /// `proc_listpids`, which returns bytes).
504    fn children(pid: u32, retries: u32) -> Option<Vec<u32>> {
505        let ppid = i32::try_from(pid).ok()?;
506        let mut cap: usize = 256;
507        for _ in 0..=retries {
508            let mut buf = vec![0i32; cap];
509            let bytes = i32::try_from(cap * std::mem::size_of::<i32>()).ok()?;
510            // SAFETY: `buf` holds `cap` i32 and we pass its size in bytes.
511            let n = unsafe { proc_listchildpids_raw(ppid, buf.as_mut_ptr(), bytes) };
512            if n < 0 {
513                return None;
514            }
515            let n = n as usize;
516            if n < cap {
517                buf.truncate(n);
518                return Some(
519                    buf.into_iter()
520                        .filter_map(|p| u32::try_from(p).ok())
521                        .collect(),
522                );
523            }
524            cap *= 4;
525        }
526        None
527    }
528
529    pub(super) fn snapshot(root: u32, seen: &[Id], limits: &Limits) -> Snapshot {
530        let (numer, denom) = timebase();
531        let to_ns = |ticks: u64| -> u64 {
532            u64::try_from(u128::from(ticks) * numer / denom).unwrap_or(u64::MAX)
533        };
534        let retries = limits.max_retries;
535        walk(
536            root,
537            seen,
538            limits,
539            &mut Instant::now,
540            &mut |pid| children(pid, retries),
541            &mut |pid, parent| match rusage(pid) {
542                Some(ri) => Read::Proc(Proc {
543                    id: Id {
544                        pid,
545                        start: ri.proc_start_abstime,
546                    },
547                    ppid: parent,
548                    cpu_ns: to_ns(
549                        ri.user_time
550                            .saturating_add(ri.system_time)
551                            .saturating_add(ri.child_user_time)
552                            .saturating_add(ri.child_system_time),
553                    ),
554                }),
555                // ESRCH (reaped) or EPERM (another user's): skip that pid.
556                None => Read::Gone,
557            },
558        )
559    }
560}
561
562#[cfg(not(any(target_os = "linux", target_os = "macos")))]
563mod platform {
564    use super::{Id, Limits, Snapshot};
565
566    pub(super) fn snapshot(_root: u32, _seen: &[Id], _limits: &Limits) -> Snapshot {
567        Snapshot::Unavailable
568    }
569}
570
571#[cfg(test)]
572mod tests {
573    use super::*;
574
575    fn id(pid: u32) -> Id {
576        Id {
577            pid,
578            start: 1000 + u64::from(pid),
579        }
580    }
581    fn p(pid: u32, ppid: u32, cpu_ms: u64) -> Proc {
582        Proc {
583            id: id(pid),
584            ppid,
585            cpu_ns: cpu_ms * 1_000_000,
586        }
587    }
588    fn map(ps: &[Proc]) -> HashMap<Id, Proc> {
589        ps.iter().map(|p| (p.id, *p)).collect()
590    }
591    const MS: u64 = 1_000_000;
592
593    // --- /proc/<pid>/stat parsing ---------------------------------------
594
595    fn stat(pid: u32, comm: &[u8], ppid: u32, times: [i64; 4], start: u64) -> Vec<u8> {
596        let mut v = format!("{pid} (").into_bytes();
597        v.extend_from_slice(comm);
598        v.extend_from_slice(
599            format!(
600                ") S {ppid} 1 1 0 -1 4194304 100 0 0 0 {} {} {} {} 20 0 1 0 {start} 1000 100",
601                times[0], times[1], times[2], times[3]
602            )
603            .as_bytes(),
604        );
605        v
606    }
607
608    #[test]
609    fn a_plain_stat_line_parses_to_self_plus_reaped_cpu() {
610        let s = stat(42, b"node", 7, [100, 50, 30, 20], 555);
611        let got = parse_proc_stat(&s, 10 * MS).unwrap();
612        assert_eq!(
613            got.id,
614            Id {
615                pid: 42,
616                start: 555
617            }
618        );
619        assert_eq!(got.ppid, 7);
620        assert_eq!(got.cpu_ns, 200 * 10 * MS);
621    }
622
623    #[test]
624    fn the_command_name_may_hold_parens_spaces_newlines_and_non_utf8() {
625        for comm in [
626            &b"node (vitest 1)"[..],
627            b"a b\nc",
628            b"x)y) z",
629            &[0xff, 0xfe, b')'],
630        ] {
631            let s = stat(9, comm, 3, [1, 1, 1, 1], 77);
632            let got = parse_proc_stat(&s, MS).unwrap_or_else(|| panic!("{comm:?}"));
633            assert_eq!(
634                (got.id.pid, got.ppid, got.id.start, got.cpu_ns),
635                (9, 3, 77, 4 * MS)
636            );
637        }
638    }
639
640    #[test]
641    fn a_negative_reaped_time_counts_as_zero() {
642        let s = stat(5, b"sh", 1, [10, 0, -3, -1], 1);
643        assert_eq!(parse_proc_stat(&s, MS).unwrap().cpu_ns, 10 * MS);
644    }
645
646    #[test]
647    fn a_truncated_or_garbled_line_is_none() {
648        assert!(parse_proc_stat(b"", MS).is_none());
649        assert!(parse_proc_stat(b"12 (x) S 1", MS).is_none());
650        assert!(
651            parse_proc_stat(b"nope (x) S 1 1 1 0 -1 0 0 0 0 0 1 1 1 1 20 0 1 0 5", MS).is_none()
652        );
653    }
654
655    // --- the bounded walk ------------------------------------------------
656
657    struct World {
658        procs: HashMap<u32, Proc>,
659    }
660    impl World {
661        fn new(ps: &[Proc]) -> Self {
662            World {
663                procs: ps.iter().map(|p| (p.id.pid, *p)).collect(),
664            }
665        }
666        fn kids(&self, pid: u32) -> Vec<u32> {
667            let mut k: Vec<u32> = self
668                .procs
669                .values()
670                .filter(|p| p.ppid == pid)
671                .map(|p| p.id.pid)
672                .collect();
673            k.sort_unstable();
674            k
675        }
676        fn walk(&self, root: u32, seen: &[Id], limits: &Limits) -> Snapshot {
677            let t = Instant::now();
678            walk(
679                root,
680                seen,
681                limits,
682                &mut || t,
683                &mut |pid| Some(self.kids(pid)),
684                &mut |pid, _| match self.procs.get(&pid) {
685                    Some(p) => Read::Proc(*p),
686                    None => Read::Gone,
687                },
688            )
689        }
690    }
691
692    fn pids(s: &Snapshot) -> Vec<u32> {
693        let Snapshot::Complete(v) = s else {
694            panic!("{s:?}")
695        };
696        let mut out: Vec<u32> = v.iter().map(|p| p.id.pid).collect();
697        out.sort_unstable();
698        out
699    }
700
701    #[test]
702    fn the_walk_takes_the_root_and_every_descendant_and_nothing_else() {
703        let w = World::new(&[p(10, 1, 0), p(11, 10, 0), p(12, 11, 0), p(99, 1, 0)]);
704        assert_eq!(pids(&w.walk(10, &[], &Limits::default())), vec![10, 11, 12]);
705    }
706
707    #[test]
708    fn an_unreadable_root_is_unavailable() {
709        let w = World::new(&[p(11, 10, 0)]);
710        assert_eq!(w.walk(10, &[], &Limits::default()), Snapshot::Unavailable);
711    }
712
713    #[test]
714    fn a_seen_orphan_keeps_counting_with_its_children_but_a_reused_pid_does_not() {
715        // 12 was a grandchild; its parent exited and it was reparented to 1.
716        let w = World::new(&[p(10, 1, 0), p(12, 1, 0), p(13, 12, 0), p(20, 1, 0)]);
717        let reused = Id { pid: 20, start: 1 }; // pid 20 now belongs to someone else
718        assert_eq!(
719            pids(&w.walk(10, &[id(12), reused], &Limits::default())),
720            vec![10, 12, 13]
721        );
722    }
723
724    #[test]
725    fn an_orphan_never_seen_before_reparenting_is_invisible() {
726        let w = World::new(&[p(10, 1, 0), p(12, 1, 0)]);
727        assert_eq!(pids(&w.walk(10, &[], &Limits::default())), vec![10]);
728    }
729
730    #[test]
731    fn exceeding_the_process_cap_is_partial() {
732        let mut ps = vec![p(10, 1, 0)];
733        ps.extend((0..5).map(|i| p(100 + i, 10, 0)));
734        let w = World::new(&ps);
735        let limits = Limits {
736            max_procs: 5,
737            ..Limits::default()
738        };
739        assert_eq!(w.walk(10, &[], &limits), Snapshot::Partial);
740        let limits = Limits {
741            max_procs: 6,
742            ..Limits::default()
743        };
744        assert!(matches!(w.walk(10, &[], &limits), Snapshot::Complete(_)));
745    }
746
747    #[test]
748    fn passing_the_deadline_is_partial() {
749        let ps = [p(10, 1, 0), p(11, 10, 0), p(12, 10, 0)];
750        let w = World::new(&ps);
751        let base = Instant::now();
752        let mut ticks = 0u64;
753        // The test is about the deadline, not its default: pin one the
754        // 60 ms ticks below cross on the third process.
755        let limits = Limits {
756            deadline: Duration::from_millis(100),
757            ..Limits::default()
758        };
759        let got = walk(
760            10,
761            &[],
762            &limits,
763            &mut || {
764                ticks += 1;
765                base + Duration::from_millis(60 * ticks)
766            },
767            &mut |pid| Some(w.kids(pid)),
768            &mut |pid, _| Read::Proc(w.procs[&pid]),
769        );
770        assert_eq!(got, Snapshot::Partial);
771    }
772
773    #[test]
774    fn a_child_list_that_cannot_be_read_is_partial() {
775        let w = World::new(&[p(10, 1, 0), p(11, 10, 0)]);
776        let t = Instant::now();
777        let got = walk(
778            10,
779            &[],
780            &Limits::default(),
781            &mut || t,
782            &mut |_| None,
783            &mut |pid, _| Read::Proc(w.procs[&pid]),
784        );
785        assert_eq!(got, Snapshot::Partial);
786    }
787
788    #[test]
789    fn a_child_gone_between_listing_and_reading_is_skipped() {
790        let w = World::new(&[p(10, 1, 0), p(11, 10, 0)]);
791        let t = Instant::now();
792        let got = walk(
793            10,
794            &[],
795            &Limits::default(),
796            &mut || t,
797            &mut |pid| Some(if pid == 10 { vec![11, 12] } else { vec![] }),
798            &mut |pid, _| w.procs.get(&pid).map_or(Read::Gone, |p| Read::Proc(*p)),
799        );
800        assert_eq!(pids(&got), vec![10, 11]);
801    }
802
803    // --- gain: what counts as work -----------------------------------------
804
805    #[test]
806    fn a_process_in_both_snapshots_counts_what_it_gained_never_less_than_zero() {
807        let prev = map(&[p(10, 1, 100), p(11, 10, 500)]);
808        let cur = map(&[p(10, 1, 350), p(11, 10, 0)]); // 11 exec'd: macOS zeroes
809        assert_eq!(gain_ns(&prev, &cur), 250 * MS);
810    }
811
812    #[test]
813    fn a_process_new_since_the_last_snapshot_counts_whole() {
814        let prev = map(&[p(10, 1, 100)]);
815        let cur = map(&[p(10, 1, 100), p(11, 10, 400)]);
816        assert_eq!(gain_ns(&prev, &cur), 400 * MS);
817    }
818
819    #[test]
820    fn fork_per_file_work_is_counted_through_the_parents_reaped_time() {
821        // Workers born and reaped between samples: never seen, but the
822        // parent's reaped-children counter carries their CPU.
823        let prev = map(&[p(10, 1, 1000)]);
824        let cur = map(&[p(10, 1, 9000)]);
825        assert_eq!(gain_ns(&prev, &cur), 8000 * MS);
826    }
827
828    #[test]
829    fn a_child_reaped_late_is_not_credited_again() {
830        // 11 burnt 2 s long ago (credited then), slept, and is reaped now:
831        // the parent's counter jumps by 2 s, and none of it is new work.
832        let prev = map(&[p(10, 1, 100), p(11, 10, 2000)]);
833        let cur = map(&[p(10, 1, 2100)]);
834        assert_eq!(gain_ns(&prev, &cur), 0);
835    }
836
837    #[test]
838    fn a_deep_tree_reaped_bottom_up_is_not_counted_once_per_level() {
839        // Six levels, all credited while alive; all reaped in one window,
840        // each parent carrying its descendants' totals up to the root.
841        let chain: Vec<Proc> = (0..6)
842            .map(|i| p(10 + i, if i == 0 { 1 } else { 9 + i }, 1000))
843            .collect();
844        let prev = map(&chain);
845        let cur = map(&[p(10, 1, 6000)]);
846        assert_eq!(gain_ns(&prev, &cur), 0);
847    }
848
849    #[test]
850    fn new_work_alongside_a_late_reap_still_counts() {
851        let prev = map(&[p(10, 1, 100), p(11, 10, 2000)]);
852        let cur = map(&[p(10, 1, 2600)]); // 500 ms of the parent's own work
853        assert_eq!(gain_ns(&prev, &cur), 500 * MS);
854    }
855
856    #[test]
857    fn an_orphan_reaped_by_init_takes_nothing_from_the_tree() {
858        let prev = map(&[p(10, 1, 100), p(12, 1, 3000)]);
859        let cur = map(&[p(10, 1, 400)]);
860        assert_eq!(gain_ns(&prev, &cur), 300 * MS);
861    }
862
863    // --- the tracker: baselines and resets -----------------------------
864
865    #[test]
866    fn the_first_complete_snapshot_is_only_a_baseline() {
867        let mut t = Tracker::default();
868        let now = Instant::now();
869        assert_eq!(
870            t.observe(now, Snapshot::Complete(vec![p(10, 1, 99_000)])),
871            Observation::Baseline
872        );
873    }
874
875    /// One partial snapshot is skipped, not fatal: the baseline is kept, the
876    /// seen set still names its processes, and the next complete snapshot is
877    /// a window spanning the gap, marked as gapped.
878    #[test]
879    fn a_partial_snapshot_is_skipped_and_the_next_complete_one_spans_the_gap() {
880        let mut t = Tracker::default();
881        let t0 = Instant::now();
882        let s = Duration::from_secs(1);
883        t.observe(t0, Snapshot::Complete(vec![p(10, 1, 0)]));
884        assert_eq!(t.observe(t0 + s, Snapshot::Partial), Observation::Skipped);
885        assert_eq!(t.seen(), vec![id(10)]);
886        match t.observe(t0 + 2 * s, Snapshot::Complete(vec![p(10, 1, 5000)])) {
887            Observation::Window(w) => {
888                assert_eq!(w.gain_ns, 5000 * MS);
889                assert_eq!(w.start, t0);
890                assert_eq!(w.end, t0 + 2 * s);
891                assert_eq!(w.milli_cores(), 2500);
892                assert!(w.gapped);
893            }
894            o => panic!("{o:?}"),
895        }
896        // A window with no gap in it says so.
897        match t.observe(t0 + 3 * s, Snapshot::Complete(vec![p(10, 1, 5500)])) {
898            Observation::Window(w) => {
899                assert_eq!(w.gain_ns, 500 * MS);
900                assert!(!w.gapped);
901            }
902            o => panic!("{o:?}"),
903        }
904    }
905
906    /// Three partials in a row, or an unreadable root, give the baseline up:
907    /// unmeasured, nothing seen, and the next complete snapshot starts over.
908    #[test]
909    fn too_many_partials_or_an_unavailable_root_rebaseline() {
910        let mut t = Tracker::default();
911        let t0 = Instant::now();
912        let s = Duration::from_secs(1);
913        t.observe(t0, Snapshot::Complete(vec![p(10, 1, 0)]));
914        assert_eq!(t.observe(t0 + s, Snapshot::Partial), Observation::Skipped);
915        assert_eq!(
916            t.observe(t0 + 2 * s, Snapshot::Partial),
917            Observation::Skipped
918        );
919        assert_eq!(
920            t.observe(t0 + 3 * s, Snapshot::Partial),
921            Observation::Unmeasured
922        );
923        assert!(t.seen().is_empty());
924        assert_eq!(
925            t.observe(t0 + 4 * s, Snapshot::Complete(vec![p(10, 1, 5000)])),
926            Observation::Baseline
927        );
928
929        let mut t = Tracker::default();
930        t.observe(t0, Snapshot::Complete(vec![p(10, 1, 0)]));
931        assert_eq!(
932            t.observe(t0 + s, Snapshot::Unavailable),
933            Observation::Unmeasured
934        );
935        assert!(t.seen().is_empty());
936    }
937
938    #[test]
939    fn the_tracker_remembers_what_it_saw_for_the_next_walk() {
940        let mut t = Tracker::default();
941        t.observe(
942            Instant::now(),
943            Snapshot::Complete(vec![p(10, 1, 0), p(12, 10, 0)]),
944        );
945        let mut seen = t.seen();
946        seen.sort_by_key(|i| i.pid);
947        assert_eq!(seen, vec![id(10), id(12)]);
948    }
949
950    #[test]
951    fn four_busy_cores_read_as_four_thousand_milli_cores() {
952        let t0 = Instant::now();
953        let w = Window {
954            start: t0,
955            end: t0 + Duration::from_secs(10),
956            gain_ns: 40 * 1_000 * MS,
957            gapped: false,
958        };
959        assert_eq!(w.milli_cores(), 4000);
960    }
961
962    // --- the real platform source ----------------------------------------
963
964    /// A child that burns CPU shows it while alive, and its parent carries
965    /// it after reaping — on the platforms that measure at all. On Apple
966    /// Silicon this is the test that catches reading mach ticks as ns.
967    #[cfg(any(target_os = "linux", target_os = "macos"))]
968    #[test]
969    fn the_platform_sees_a_burning_child_and_its_reaped_time() {
970        use std::process::Command;
971        let me = std::process::id();
972        let before = match snapshot(me, &[], &Limits::default()) {
973            Snapshot::Complete(v) => v.iter().find(|p| p.id.pid == me).unwrap().cpu_ns,
974            s => panic!("{s:?}"),
975        };
976        // `yes` into /dev/null is pure CPU: one second of it is ~one core-
977        // second. (A shell loop that forks `date` to watch the clock spends
978        // most of its wall time creating processes, and on a fast Apple
979        // Silicon runner read as barely 0.2 s — too weak to test with.)
980        let mut hog = Command::new("yes")
981            .stdout(std::process::Stdio::null())
982            .spawn()
983            .unwrap();
984        std::thread::sleep(std::time::Duration::from_secs(1));
985        let _ = hog.kill();
986        let _ = hog.wait(); // reaped: its CPU is now in our children's time
987        let after = match snapshot(me, &[], &Limits::default()) {
988            Snapshot::Complete(v) => v.iter().find(|p| p.id.pid == me).unwrap().cpu_ns,
989            s => panic!("{s:?}"),
990        };
991        let reaped = after.saturating_sub(before);
992        assert!(
993            reaped >= 500 * MS,
994            "reaped child CPU only {} ms",
995            reaped / MS
996        );
997        assert!(
998            reaped < 60_000 * MS,
999            "implausible {} ms — wrong time unit?",
1000            reaped / MS
1001        );
1002    }
1003}