1use std::collections::{HashMap, HashSet, VecDeque};
37use std::time::{Duration, Instant};
38
39#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
44pub struct Id {
45 pub pid: u32,
46 pub start: u64,
47}
48
49#[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#[derive(Debug, PartialEq, Eq)]
63pub enum Snapshot {
64 Complete(Vec<Proc>),
65 Partial,
67 Unavailable,
69}
70
71#[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(500),
90 }
91 }
92}
93
94pub enum Read {
96 Proc(Proc),
97 Gone,
99 Failed,
101}
102
103pub 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 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
183pub 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#[derive(Clone, Copy, Debug, PartialEq, Eq)]
220pub struct Window {
221 pub start: Instant,
222 pub end: Instant,
223 pub gain_ns: u64,
224 pub gapped: bool,
228}
229
230impl Window {
231 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#[derive(Clone, Copy, Debug, PartialEq, Eq)]
242pub enum Observation {
243 Baseline,
246 Window(Window),
248 Skipped,
252 Unmeasured,
256}
257
258pub const PARTIALS_BEFORE_RESET: u32 = 3;
261
262#[derive(Default)]
265pub struct Tracker {
266 prev: Option<(Instant, HashMap<Id, Proc>)>,
267 partials: u32,
269}
270
271impl Tracker {
272 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
315pub 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; }
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
362pub fn snapshot(root: u32, seen: &[Id], limits: &Limits) -> Snapshot {
364 platform::snapshot(root, seen, limits)
365}
366
367pub 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 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; fn tick_ns() -> u64 {
385 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 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; };
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 #[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 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 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 let rc = unsafe { proc_pid_rusage_raw(pid, RUSAGE_INFO_V2, &mut ri) };
498 (rc == 0).then_some(ri)
499 }
500
501 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 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 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 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 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 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 }; 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 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 #[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)]); 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 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 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 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)]); 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 #[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 #[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 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 #[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 #[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 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(); 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}