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