1use crate::metrics::CpuMetrics;
2use procfs::prelude::*;
3use procfs::process::all_processes;
4use procfs::{CpuTime, KernelStats};
5use std::collections::{HashMap, HashSet};
6use std::time::Instant;
7
8type Result<T> = std::result::Result<T, Box<dyn std::error::Error>>;
9
10#[derive(Debug, Clone, Copy, PartialEq)]
16enum CpuSource {
17 CgroupV2,
19 CgroupV1,
21 ProcStat,
23}
24
25impl CpuSource {
26 fn is_cgroup(self) -> bool {
27 !matches!(self, CpuSource::ProcStat)
28 }
29}
30
31#[derive(Debug, Clone, Copy)]
33struct CfsQuota {
34 max_cores: Option<f64>,
36}
37
38#[allow(clippy::collapsible_if)]
41fn detect_cpu_source() -> CpuSource {
42 if let Ok(contents) = std::fs::read_to_string("/sys/fs/cgroup/cpu.stat") {
44 if contents.contains("usage_usec") {
45 return CpuSource::CgroupV2;
46 }
47 }
48 for path in &[
50 "/sys/fs/cgroup/cpuacct/cpuacct.usage",
51 "/sys/fs/cgroup/cpu,cpuacct/cpuacct.usage",
52 "/sys/fs/cgroup/cpu/cpuacct.usage",
53 ] {
54 if std::fs::read_to_string(path).is_ok() {
55 return CpuSource::CgroupV1;
56 }
57 }
58 CpuSource::ProcStat
59}
60
61#[allow(clippy::collapsible_if)]
63fn detect_cfs_quota() -> CfsQuota {
64 if let Ok(contents) = std::fs::read_to_string("/sys/fs/cgroup/cpu.max") {
66 let parts: Vec<&str> = contents.split_whitespace().collect();
67 if parts.len() == 2 && parts[0] != "max" {
68 if let (Ok(quota), Ok(period)) = (parts[0].parse::<f64>(), parts[1].parse::<f64>()) {
69 if period > 0.0 {
70 return CfsQuota {
71 max_cores: Some(quota / period),
72 };
73 }
74 }
75 }
76 }
77 for prefix in &[
79 "/sys/fs/cgroup/cpu",
80 "/sys/fs/cgroup/cpu,cpuacct",
81 "/sys/fs/cgroup/cpuacct",
82 ] {
83 let quota_path = format!("{}/cpu.cfs_quota_us", prefix);
84 let period_path = format!("{}/cpu.cfs_period_us", prefix);
85 if let (Ok(q_str), Ok(p_str)) = (
86 std::fs::read_to_string("a_path),
87 std::fs::read_to_string(&period_path),
88 ) {
89 if let (Ok(quota), Ok(period)) =
90 (q_str.trim().parse::<i64>(), p_str.trim().parse::<i64>())
91 {
92 if quota > 0 && period > 0 {
94 return CfsQuota {
95 max_cores: Some(quota as f64 / period as f64),
96 };
97 }
98 }
99 }
100 }
101 CfsQuota { max_cores: None }
102}
103
104fn read_cgroupv2_usage_usec() -> Option<u64> {
106 let contents = std::fs::read_to_string("/sys/fs/cgroup/cpu.stat").ok()?;
107 for line in contents.lines() {
108 if let Some(val) = line.strip_prefix("usage_usec ") {
109 return val.trim().parse().ok();
110 }
111 }
112 None
113}
114
115#[allow(clippy::collapsible_if)]
117fn read_cgroupv1_usage_ns() -> Option<u64> {
118 for path in &[
119 "/sys/fs/cgroup/cpuacct/cpuacct.usage",
120 "/sys/fs/cgroup/cpu,cpuacct/cpuacct.usage",
121 "/sys/fs/cgroup/cpu/cpuacct.usage",
122 ] {
123 if let Ok(contents) = std::fs::read_to_string(path)
124 && let Ok(val) = contents.trim().parse()
125 {
126 return Some(val);
127 }
128 }
129 None
130}
131
132fn read_cgroup_usage_secs(source: CpuSource) -> Option<f64> {
135 match source {
136 CpuSource::CgroupV2 => read_cgroupv2_usage_usec().map(|usec| usec as f64 / 1_000_000.0),
137 CpuSource::CgroupV1 => read_cgroupv1_usage_ns().map(|ns| ns as f64 / 1_000_000_000.0),
138 CpuSource::ProcStat => None,
139 }
140}
141
142fn cpu_total(c: &CpuTime) -> u64 {
147 c.user
148 + c.nice
149 + c.system
150 + cpu_idle(c)
151 + c.irq.unwrap_or(0)
152 + c.softirq.unwrap_or(0)
153 + c.steal.unwrap_or(0)
154}
155
156fn cpu_idle(c: &CpuTime) -> u64 {
157 c.idle + c.iowait.unwrap_or(0)
158}
159
160fn core_util_pct(prev: &CpuTime, curr: &CpuTime) -> f64 {
162 util_pct_from_ticks(
163 cpu_total(prev),
164 cpu_idle(prev),
165 cpu_total(curr),
166 cpu_idle(curr),
167 )
168 .clamp(0.0, 100.0)
169}
170
171fn aggregate_util_cores(prev: &CpuTime, curr: &CpuTime, n_cores: usize) -> f64 {
174 util_pct_from_ticks(
175 cpu_total(prev),
176 cpu_idle(prev),
177 cpu_total(curr),
178 cpu_idle(curr),
179 ) / 100.0
180 * n_cores as f64
181}
182
183fn util_pct_from_ticks(prev_total: u64, prev_idle: u64, curr_total: u64, curr_idle: u64) -> f64 {
187 let delta_total = curr_total.saturating_sub(prev_total) as f64;
188 let delta_idle = curr_idle.saturating_sub(prev_idle) as f64;
189 if delta_total == 0.0 {
190 return 0.0;
191 }
192 (delta_total - delta_idle) / delta_total * 100.0
193}
194
195fn process_tree_ticks(root_pid: i32) -> HashMap<i32, (u64, u64)> {
203 let all: Vec<_> = match all_processes() {
205 Ok(iter) => iter.filter_map(|r| r.ok()).collect(),
206 Err(_) => return HashMap::new(),
207 };
208
209 let mut children: HashMap<i32, Vec<i32>> = HashMap::new();
223 let ticks_for: HashMap<i32, (u64, u64)> = all
224 .iter()
225 .filter_map(|proc| {
226 proc.stat().ok().map(|s| {
227 children.entry(s.ppid).or_default().push(proc.pid);
228 let user = s.utime + u64::try_from(s.cutime).unwrap_or(0);
229 let system = s.stime + u64::try_from(s.cstime).unwrap_or(0);
230 (proc.pid, (user, system))
231 })
232 })
233 .collect();
234
235 let mut result = HashMap::new();
237 let mut queue = vec![root_pid];
238 while let Some(pid) = queue.pop() {
239 if let Some(&ticks) = ticks_for.get(&pid) {
240 result.insert(pid, ticks);
241 }
242 if let Some(kids) = children.get(&pid) {
243 queue.extend(kids);
244 }
245 }
246 result
247}
248
249fn process_tree_memory_mib(pids: &[i32]) -> (u64, u64) {
253 let mut pss_kib = 0u64;
254 let mut rss_kib = 0u64;
255 for &pid in pids {
256 let Some(proc_) = procfs::process::Process::new(pid).ok() else {
257 continue;
258 };
259 if let Ok(rollup) = proc_.smaps_rollup()
260 && let Some(bytes) = rollup
261 .memory_map_rollup
262 .iter()
263 .find_map(|m| m.extension.map.get("Pss").copied())
264 {
265 pss_kib += bytes / 1024;
266 }
267 if let Ok(status) = proc_.status()
268 && let Some(vmrss) = status.vmrss
269 {
270 rss_kib += vmrss;
271 }
272 }
273 (pss_kib / 1024, rss_kib / 1024)
274}
275
276fn process_tree_io(pids: &[i32]) -> HashMap<i32, (u64, u64)> {
281 pids.iter()
282 .filter_map(|&pid| {
283 let io = procfs::process::Process::new(pid).ok()?.io().ok()?;
284 Some((pid, (io.read_bytes, io.write_bytes)))
285 })
286 .collect()
287}
288
289struct Snapshot {
294 total: CpuTime,
296 per_core: Vec<CpuTime>,
298 instant: Instant,
301 cgroup_usage_secs: Option<f64>,
303 proc_ticks: HashMap<i32, (u64, u64)>,
306 proc_io: HashMap<i32, (u64, u64)>,
309}
310
311pub struct CpuCollector {
312 pid: Option<i32>,
314 prev: Option<Snapshot>,
315 cpu_source: CpuSource,
317 cfs_quota: CfsQuota,
319 effective_cores: f64,
322 carried_forward: HashSet<i32>,
326 aggregate_cpu_steal: bool,
328}
329
330impl CpuCollector {
331 pub fn new(pid: Option<i32>, aggregate_cpu_steal: bool) -> Self {
332 let cpu_source = detect_cpu_source();
333 let cfs_quota = detect_cfs_quota();
334
335 let physical_cores = KernelStats::current()
337 .map(|s| s.cpu_time.len())
338 .unwrap_or(1) as f64;
339 let effective_cores = match cfs_quota.max_cores {
340 Some(quota) => physical_cores.min(quota),
341 None => physical_cores,
342 };
343
344 Self {
345 pid,
346 prev: None,
347 cpu_source,
348 cfs_quota,
349 effective_cores,
350 carried_forward: HashSet::new(),
351 aggregate_cpu_steal,
352 }
353 }
354
355 pub fn set_tracked_pid(&mut self, pid: Option<i32>) {
357 self.pid = pid;
358 }
359
360 pub fn collect(&mut self) -> Result<CpuMetrics> {
361 let tps = procfs::ticks_per_second() as f64;
362 let process_count = self.read_process_count();
363
364 let stats = KernelStats::current()?;
366 let cgroup_usage_secs = read_cgroup_usage_secs(self.cpu_source);
367 let proc_ticks = self.read_process_ticks();
368 let now = Instant::now();
369
370 let proc_io = self.read_process_io(&proc_ticks);
371 let (process_pss_mib, process_rss_mib) = self.read_process_memory(&proc_ticks);
372
373 let mut curr = Snapshot {
374 total: stats.total,
375 per_core: stats.cpu_time,
376 instant: now,
377 cgroup_usage_secs,
378 proc_ticks,
379 proc_io,
380 };
381
382 let metrics = match &self.prev {
383 None => {
384 self.build_first_metrics(&curr, process_count, process_pss_mib, process_rss_mib)
385 }
386 Some(prev) => self.build_subsequent_metrics_with_deltas(
387 prev,
388 &mut curr,
389 process_count,
390 tps,
391 process_pss_mib,
392 process_rss_mib,
393 ),
394 };
395
396 self.carry_forward_entries(&mut curr);
397 self.prev = Some(curr);
398 Ok(metrics)
399 }
400
401 fn read_process_count(&self) -> u32 {
402 let Ok(proc_dir) = std::fs::read_dir("/proc") else {
403 return 0;
404 };
405
406 let process_count = proc_dir
407 .filter_map(|entry| entry.ok())
408 .filter(|e| Self::is_pid_directory(e))
409 .count();
410
411 u32::try_from(process_count).unwrap_or(0)
412 }
413
414 fn is_pid_directory(entry: &std::fs::DirEntry) -> bool {
415 entry
416 .file_name()
417 .to_string_lossy()
418 .chars()
419 .all(|c| c.is_ascii_digit())
420 }
421
422 fn read_process_ticks(&self) -> HashMap<i32, (u64, u64)> {
423 match self.pid {
424 Some(root) => process_tree_ticks(root),
425 None => HashMap::new(),
426 }
427 }
428
429 fn read_process_io(&self, proc_ticks: &HashMap<i32, (u64, u64)>) -> HashMap<i32, (u64, u64)> {
430 if self.pid.is_some() {
431 let pids: Vec<i32> = proc_ticks.keys().copied().collect();
432 process_tree_io(&pids)
433 } else {
434 HashMap::new()
435 }
436 }
437
438 fn read_process_memory(
439 &self,
440 proc_ticks: &HashMap<i32, (u64, u64)>,
441 ) -> (Option<u64>, Option<u64>) {
442 if self.pid.is_some() {
443 let pids: Vec<i32> = proc_ticks.keys().copied().collect();
444 let (pss, rss) = process_tree_memory_mib(&pids);
445 (Some(pss), Some(rss))
446 } else {
447 (None, None)
448 }
449 }
450
451 fn build_first_metrics(
452 &self,
453 curr: &Snapshot,
454 process_count: u32,
455 process_pss_mib: Option<u64>,
456 process_rss_mib: Option<u64>,
457 ) -> CpuMetrics {
458 CpuMetrics {
459 utilization_pct: 0.0,
460 cgroup_utilization_pct: curr
461 .cgroup_usage_secs
462 .filter(|_| self.cpu_source.is_cgroup())
463 .map(|_| 0.0),
464 cgroup_usage_secs: curr
465 .cgroup_usage_secs
466 .filter(|_| self.cpu_source.is_cgroup())
467 .map(|_| 0.0),
468 per_core_pct: vec![0.0; curr.per_core.len()],
469 utime_secs: 0.0,
470 stime_secs: 0.0,
471 steal_time_secs: 0.0,
472 steal_time_pct: 0.0,
473 per_core_steal_time_pct: if self.aggregate_cpu_steal {
474 vec![]
475 } else {
476 vec![0.0; curr.per_core.len()]
477 },
478 process_count,
479 process_cores_used: self.pid.map(|_| 0.0),
480 process_child_count: self
481 .pid
482 .map(|_| u32::try_from(curr.proc_ticks.len().saturating_sub(1)).unwrap_or(0)),
483 process_utime_secs: self.pid.map(|_| 0.0),
484 process_stime_secs: self.pid.map(|_| 0.0),
485 process_pss_mib,
486 process_rss_mib,
487 process_disk_read_bytes: self.pid.map(|_| 0),
488 process_disk_write_bytes: self.pid.map(|_| 0),
489 process_gpu_usage: None,
490 process_gpu_vram_mib: None,
491 process_gpu_utilized: None,
492 process_tree_pids: curr.proc_ticks.keys().copied().collect(),
493 }
494 }
495
496 fn build_subsequent_metrics_with_deltas(
497 &self,
498 prev: &Snapshot,
499 curr: &mut Snapshot,
500 process_count: u32,
501 tps: f64,
502 process_pss_mib: Option<u64>,
503 process_rss_mib: Option<u64>,
504 ) -> CpuMetrics {
505 let n_cores = curr.per_core.len();
506 let elapsed = (curr.instant - prev.instant).as_secs_f64().max(0.001);
507
508 let (utime_secs, stime_secs) = self.calculate_system_cpu_deltas(prev, curr, tps);
509 let per_core_pct = self.calculate_per_core_utilization(prev, curr);
510 let utilization_pct = aggregate_util_cores(&prev.total, &curr.total, n_cores);
511 let (cgroup_usage_secs, cgroup_utilization_pct) =
512 self.calculate_cgroup_usage(prev, curr, elapsed);
513 let (exited_utime, exited_stime) = self.calculate_exited_child_ticks(prev, curr);
514
515 let process_utime_secs = self.calculate_process_utime_delta(prev, curr, exited_utime, tps);
516 let process_stime_secs = self.calculate_process_stime_delta(prev, curr, exited_stime, tps);
517
518 let steal_time_secs = self.calculate_steal_time_secs(prev, curr, tps);
519 let steal_time_pct = self.calculate_steal_time_pct(prev, curr);
520 let per_core_steal_time_pct = if self.aggregate_cpu_steal {
521 vec![]
522 } else {
523 self.calculate_per_core_steal_pct(prev, curr)
524 };
525
526 let process_cores_used = self.calculate_process_cores_used_with_caps(
527 &process_utime_secs,
528 &process_stime_secs,
529 prev,
530 curr,
531 elapsed,
532 n_cores,
533 );
534
535 let process_disk_read_bytes = self.calculate_disk_read_delta(prev, curr);
536 let process_disk_write_bytes = self.calculate_disk_write_delta(prev, curr);
537
538 CpuMetrics {
539 utilization_pct,
540 cgroup_utilization_pct,
541 cgroup_usage_secs,
542 per_core_pct,
543 utime_secs,
544 stime_secs,
545 steal_time_secs,
546 steal_time_pct,
547 per_core_steal_time_pct,
548 process_count,
549 process_cores_used,
550 process_child_count: self
551 .pid
552 .map(|_| u32::try_from(curr.proc_ticks.len().saturating_sub(1)).unwrap_or(0)),
553 process_utime_secs,
554 process_stime_secs,
555 process_pss_mib,
556 process_rss_mib,
557 process_disk_read_bytes,
558 process_disk_write_bytes,
559 process_gpu_usage: None,
560 process_gpu_vram_mib: None,
561 process_gpu_utilized: None,
562 process_tree_pids: curr.proc_ticks.keys().copied().collect(),
563 }
564 }
565
566 fn calculate_system_cpu_deltas(
567 &self,
568 prev: &Snapshot,
569 curr: &Snapshot,
570 tps: f64,
571 ) -> (f64, f64) {
572 let utime_secs = (curr.total.user + curr.total.nice)
573 .saturating_sub(prev.total.user + prev.total.nice) as f64
574 / tps;
575 let stime_secs = curr.total.system.saturating_sub(prev.total.system) as f64 / tps;
576 (utime_secs, stime_secs)
577 }
578
579 fn calculate_per_core_utilization(&self, prev: &Snapshot, curr: &Snapshot) -> Vec<f64> {
580 prev.per_core
581 .iter()
582 .zip(curr.per_core.iter())
583 .map(|(p, c)| core_util_pct(p, c))
584 .collect()
585 }
586
587 fn calculate_cgroup_usage(
588 &self,
589 prev: &Snapshot,
590 curr: &Snapshot,
591 elapsed: f64,
592 ) -> (Option<f64>, Option<f64>) {
593 match (curr.cgroup_usage_secs, prev.cgroup_usage_secs) {
594 (Some(curr_cg), Some(prev_cg)) => {
595 let delta = (curr_cg - prev_cg).max(0.0);
596 let cores_used = delta / elapsed;
597 (Some(delta), Some(cores_used.min(self.effective_cores)))
598 }
599 _ => (None, None),
600 }
601 }
602
603 fn calculate_exited_child_ticks(&self, prev: &Snapshot, curr: &Snapshot) -> (u64, u64) {
605 if self.pid.is_none() {
606 return (0, 0);
607 }
608
609 prev.proc_ticks
610 .iter()
611 .filter(|(pid, _)| !curr.proc_ticks.contains_key(*pid))
612 .fold((0, 0), |(user_sum, sys_sum), (_, &(user, sys))| {
613 (user_sum + user, sys_sum + sys)
614 })
615 }
616 fn calculate_process_utime_delta(
617 &self,
618 prev: &Snapshot,
619 curr: &Snapshot,
620 exited_utime: u64,
621 tps: f64,
622 ) -> Option<f64> {
623 if self.pid.is_none() {
624 return None;
625 }
626
627 let raw: u64 = curr
628 .proc_ticks
629 .iter()
630 .map(|(pid, &(cu, _))| {
631 let pu = prev.proc_ticks.get(pid).map(|&(u, _)| u).unwrap_or(cu);
632 cu.saturating_sub(pu)
633 })
634 .sum();
635
636 let adjusted = if exited_utime <= raw {
637 raw - exited_utime
638 } else {
639 raw
640 };
641
642 Some(adjusted as f64 / tps)
643 }
644
645 fn calculate_process_stime_delta(
646 &self,
647 prev: &Snapshot,
648 curr: &Snapshot,
649 exited_stime: u64,
650 tps: f64,
651 ) -> Option<f64> {
652 if self.pid.is_none() {
653 return None;
654 }
655
656 let mut raw: u64 = 0;
657 for (pid, &(_, cs)) in &curr.proc_ticks {
658 let ps = prev.proc_ticks.get(pid).map(|&(_, s)| s).unwrap_or(cs);
659 raw += cs.saturating_sub(ps);
660 }
661
662 let adjusted = if exited_stime <= raw {
663 raw - exited_stime
664 } else {
665 raw
666 };
667
668 Some(adjusted as f64 / tps)
669 }
670
671 fn calculate_steal_time_secs(&self, prev: &Snapshot, curr: &Snapshot, tps: f64) -> f64 {
672 let curr_steal = curr.total.steal.unwrap_or(0);
673 let prev_steal = prev.total.steal.unwrap_or(0);
674
675 curr_steal.saturating_sub(prev_steal) as f64 / tps
676 }
677
678 fn calculate_steal_time_pct(&self, prev: &Snapshot, curr: &Snapshot) -> f64 {
679 let prev_total = cpu_total(&prev.total);
680 let curr_total = cpu_total(&curr.total);
681 let prev_steal = prev.total.steal.unwrap_or(0);
682 let curr_steal = curr.total.steal.unwrap_or(0);
683
684 let delta_total = curr_total.saturating_sub(prev_total) as f64;
685 let delta_steal = curr_steal.saturating_sub(prev_steal) as f64;
686
687 if delta_total == 0.0 {
688 0.0
689 } else {
690 (delta_steal / delta_total * 100.0).clamp(0.0, 100.0)
691 }
692 }
693
694 fn calculate_per_core_steal_pct(&self, prev: &Snapshot, curr: &Snapshot) -> Vec<f64> {
695 prev.per_core
696 .iter()
697 .zip(curr.per_core.iter())
698 .map(|(p, c)| {
699 let p_total = cpu_total(p);
700 let c_total = cpu_total(c);
701 let p_steal = p.steal.unwrap_or(0);
702 let c_steal = c.steal.unwrap_or(0);
703
704 let delta_total = c_total.saturating_sub(p_total) as f64;
705 let delta_steal = c_steal.saturating_sub(p_steal) as f64;
706
707 if delta_total == 0.0 {
708 0.0
709 } else {
710 (delta_steal / delta_total * 100.0).clamp(0.0, 100.0)
711 }
712 })
713 .collect()
714 }
715
716 fn calculate_process_cores_used_with_caps(
717 &self,
718 process_utime_secs: &Option<f64>,
719 process_stime_secs: &Option<f64>,
720 prev: &Snapshot,
721 curr: &Snapshot,
722 elapsed: f64,
723 n_cores: usize,
724 ) -> Option<f64> {
725 match (self.pid, process_utime_secs, process_stime_secs) {
726 (Some(_), Some(u), Some(s)) => {
727 let raw_cores = ((u + s) / elapsed).max(0.0);
728
729 let sys_total_delta = cpu_total(&curr.total).saturating_sub(cpu_total(&prev.total));
731 let sys_idle_delta = cpu_idle(&curr.total).saturating_sub(cpu_idle(&prev.total));
732 let sys_busy_secs = if sys_total_delta > 0 {
733 (sys_total_delta - sys_idle_delta.min(sys_total_delta)) as f64
734 / procfs::ticks_per_second() as f64
735 } else {
736 f64::MAX
737 };
738 let tick_ratio_cap = sys_busy_secs / elapsed;
739
740 let quota_cap = self.cfs_quota.max_cores.unwrap_or(n_cores as f64);
742
743 Some(raw_cores.min(tick_ratio_cap).min(quota_cap))
744 }
745 _ => None,
746 }
747 }
748
749 fn calculate_disk_read_delta(&self, prev: &Snapshot, curr: &Snapshot) -> Option<u64> {
750 self.pid.map(|_| {
751 curr.proc_io
752 .iter()
753 .map(|(pid, &(cr, _))| {
754 let pr = prev.proc_io.get(pid).map(|&(r, _)| r).unwrap_or(cr);
755 cr.saturating_sub(pr)
756 })
757 .sum()
758 })
759 }
760
761 fn calculate_disk_write_delta(&self, prev: &Snapshot, curr: &Snapshot) -> Option<u64> {
762 self.pid.map(|_| {
763 curr.proc_io
764 .iter()
765 .map(|(pid, &(_, cw))| {
766 let pw = prev.proc_io.get(pid).map(|&(_, w)| w).unwrap_or(cw);
767 cw.saturating_sub(pw)
768 })
769 .sum()
770 })
771 }
772
773 fn carry_forward_entries(&mut self, curr: &mut Snapshot) {
775 let mut new_carried = HashSet::new();
776
777 if let Some(ref prev_snap) = self.prev {
778 for (&pid, &ticks) in &prev_snap.proc_ticks {
780 if !curr.proc_ticks.contains_key(&pid) && !self.carried_forward.contains(&pid) {
781 curr.proc_ticks.insert(pid, ticks);
782 new_carried.insert(pid);
783 }
784 }
785
786 for (&pid, &io) in &prev_snap.proc_io {
788 if !curr.proc_io.contains_key(&pid) && !self.carried_forward.contains(&pid) {
789 curr.proc_io.insert(pid, io);
790 }
791 }
792 }
793
794 self.carried_forward = new_carried;
795 }
796}
797
798#[cfg(test)]
803mod tests {
804 use super::*;
805
806 #[test]
814 fn test_util_pct_all_idle_is_zero() {
815 assert_eq!(util_pct_from_ticks(0, 0, 1600, 1600), 0.0);
817 }
818
819 #[test]
820 fn test_util_pct_fully_busy_is_100() {
821 let pct = util_pct_from_ticks(0, 0, 1600, 0);
823 assert!((pct - 100.0).abs() < 0.01, "expected 100.0, got {pct}");
824 }
825
826 #[test]
827 fn test_util_pct_half_busy_is_50() {
828 let pct = util_pct_from_ticks(0, 0, 1600, 800);
830 assert!((pct - 50.0).abs() < 0.01, "expected 50.0, got {pct}");
831 }
832
833 #[test]
834 fn test_util_pct_no_delta_is_zero() {
835 assert_eq!(util_pct_from_ticks(100, 50, 100, 50), 0.0);
837 }
838
839 #[test]
842 fn test_aggregate_util_cores_no_clamp() {
843 let pct = util_pct_from_ticks(0, 0, 1000, 1);
845 let cores = pct / 100.0 * 4.0_f64;
846 assert!(cores > 3.9, "expected close to 4.0, got {cores}");
847 assert!(
848 cores < 4.05,
849 "should not greatly exceed n_cores, got {cores}"
850 );
851 }
852
853 #[test]
856 fn test_util_pct_raw_is_not_clamped() {
857 let raw = util_pct_from_ticks(0, 0, 1000, 0);
859 assert!((raw - 100.0).abs() < 0.01);
860 assert_eq!(raw.clamp(0.0, 100.0), 100.0);
862 }
863
864 #[test]
868 fn test_first_collect_returns_zero_for_delta_fields() {
869 let mut collector = CpuCollector::new(None, false);
870 let metrics = collector.collect().expect("first collect failed");
871 assert_eq!(
872 metrics.utilization_pct, 0.0,
873 "utilization_pct must be 0.0 on first collect, got {}",
874 metrics.utilization_pct
875 );
876 assert!(
877 metrics.per_core_pct.iter().all(|&v| v == 0.0),
878 "per_core_pct must be all-zero on first collect: {:?}",
879 metrics.per_core_pct
880 );
881 assert_eq!(
882 metrics.utime_secs, 0.0,
883 "utime_secs must be 0.0 on first collect, got {}",
884 metrics.utime_secs
885 );
886 assert_eq!(
887 metrics.stime_secs, 0.0,
888 "stime_secs must be 0.0 on first collect, got {}",
889 metrics.stime_secs
890 );
891 }
892
893 #[test]
895 fn test_first_collect_with_pid_returns_some_process_fields() {
896 let pid = i32::try_from(std::process::id()).expect("PID too large");
897 let mut collector = CpuCollector::new(Some(pid), false);
898 let m = collector.collect().expect("collect() failed");
899 assert!(
900 m.process_cores_used.is_some(),
901 "process_cores_used must be Some when PID is tracked"
902 );
903 assert!(
904 m.process_child_count.is_some(),
905 "process_child_count must be Some when PID is tracked"
906 );
907 assert!(
908 m.process_pss_mib.is_some(),
909 "process_pss_mib must be Some when PID is tracked"
910 );
911 assert!(
912 m.process_rss_mib.is_some(),
913 "process_rss_mib must be Some when PID is tracked"
914 );
915 assert!(
916 m.process_utime_secs.is_some(),
917 "process_utime_secs must be Some when PID is tracked"
918 );
919 assert!(
920 m.process_stime_secs.is_some(),
921 "process_stime_secs must be Some when PID is tracked"
922 );
923 assert!(
924 m.process_disk_read_bytes.is_some(),
925 "process_disk_read_bytes must be Some when PID is tracked"
926 );
927 assert!(
928 m.process_disk_write_bytes.is_some(),
929 "process_disk_write_bytes must be Some when PID is tracked"
930 );
931 }
932
933 #[test]
935 fn test_process_tree_memory_nonzero_for_self() {
936 let pid = i32::try_from(std::process::id()).expect("PID too large");
937 let (pss, rss) = process_tree_memory_mib(&[pid]);
938 assert!(
939 pss > 0,
940 "PSS for the current process should be > 0, got {pss}"
941 );
942 assert!(
943 rss > 0,
944 "RSS for the current process should be > 0, got {rss}"
945 );
946 assert!(
947 pss <= rss,
948 "PSS ({pss}) should not exceed RSS ({rss}) for a single process"
949 );
950 }
951
952 #[test]
958 fn test_process_tree_ticks_contains_root_pid() {
959 let ticks = process_tree_ticks(1);
960 assert!(
961 ticks.contains_key(&1),
962 "process_tree_ticks(1) must contain PID 1 (init/systemd is always present)"
963 );
964 }
965
966 #[test]
968 fn test_second_collect_with_pid_nonneg_cores() {
969 let pid = i32::try_from(std::process::id()).expect("PID too large");
970 let mut collector = CpuCollector::new(Some(pid), false);
971 let _ = collector.collect().expect("first collect() failed");
972 let m = collector.collect().expect("second collect() failed");
973 let cores = m
974 .process_cores_used
975 .expect("process_cores_used must be Some");
976 assert!(
977 cores >= 0.0,
978 "process_cores_used must be >= 0.0, got {cores}"
979 );
980 }
981
982 #[test]
984 fn test_second_collect_no_pid_all_process_fields_none() {
985 let mut collector = CpuCollector::new(None, false);
986 let _ = collector.collect().expect("first collect() failed");
987 let m = collector.collect().expect("second collect() failed");
988 assert!(
989 m.process_cores_used.is_none(),
990 "process_cores_used must be None when not tracking"
991 );
992 assert!(
993 m.process_child_count.is_none(),
994 "process_child_count must be None when not tracking"
995 );
996 assert!(
997 m.process_pss_mib.is_none(),
998 "process_pss_mib must be None when not tracking"
999 );
1000 assert!(
1001 m.process_rss_mib.is_none(),
1002 "process_rss_mib must be None when not tracking"
1003 );
1004 assert!(
1005 m.process_utime_secs.is_none(),
1006 "process_utime_secs must be None when not tracking"
1007 );
1008 assert!(
1009 m.process_stime_secs.is_none(),
1010 "process_stime_secs must be None when not tracking"
1011 );
1012 assert!(
1013 m.process_disk_read_bytes.is_none(),
1014 "process_disk_read_bytes must be None when not tracking"
1015 );
1016 assert!(
1017 m.process_disk_write_bytes.is_none(),
1018 "process_disk_write_bytes must be None when not tracking"
1019 );
1020 }
1021
1022 #[test]
1024 fn test_process_count_positive() {
1025 let mut collector = CpuCollector::new(None, false);
1026 let m = collector.collect().expect("collect() failed");
1027 assert!(
1028 m.process_count > 0,
1029 "process_count must be > 0, got {}",
1030 m.process_count
1031 );
1032 }
1033
1034 #[test]
1046 fn test_cutime_correction_cancels_exited_child_ticks() {
1047 let prev: HashMap<i32, (u64, u64)> = [
1048 (200, (50, 0)), (100, (500, 0)), ]
1051 .iter()
1052 .cloned()
1053 .collect();
1054
1055 let curr: HashMap<i32, (u64, u64)> =
1059 [(200, (50 + 250 + 2500, 0))].iter().cloned().collect();
1060
1061 let raw: u64 = curr
1062 .iter()
1063 .map(|(pid, &(cu, cs))| {
1064 let (pu, ps) = prev.get(pid).copied().unwrap_or((cu, cs));
1065 cu.saturating_sub(pu) + cs.saturating_sub(ps)
1066 })
1067 .sum();
1068 assert_eq!(
1069 raw, 2750,
1070 "raw delta must include the double-counted pre-snapshot child ticks"
1071 );
1072
1073 let exited: u64 = prev
1074 .iter()
1075 .filter(|(pid, _)| !curr.contains_key(pid))
1076 .map(|(_, &(pu, ps))| pu + ps)
1077 .sum();
1078 assert_eq!(
1079 exited, 500,
1080 "exited ticks must equal the child's pre-snapshot tick count"
1081 );
1082
1083 let corrected = raw.saturating_sub(exited);
1084 assert_eq!(
1086 corrected, 2250,
1087 "corrected delta must exclude the child's pre-snapshot ticks"
1088 );
1089 }
1090
1091 #[test]
1098 fn test_cutime_correction_handles_cascaded_exits() {
1099 let prev: HashMap<i32, (u64, u64)> = [
1100 (7, (0, 0)), (8, (100, 0)), (9, (200, 0)), ]
1104 .iter()
1105 .cloned()
1106 .collect();
1107
1108 let curr: HashMap<i32, (u64, u64)> = [(7, (30 + 400, 0))].iter().cloned().collect();
1115
1116 let raw: u64 = curr
1117 .iter()
1118 .map(|(pid, &(cu, cs))| {
1119 let (pu, ps) = prev.get(pid).copied().unwrap_or((cu, cs));
1120 cu.saturating_sub(pu) + cs.saturating_sub(ps)
1121 })
1122 .sum();
1123 assert_eq!(raw, 430);
1125
1126 let exited: u64 = prev
1127 .iter()
1128 .filter(|(pid, _)| !curr.contains_key(pid))
1129 .map(|(_, &(pu, ps))| pu + ps)
1130 .sum();
1131 assert_eq!(
1132 exited, 300,
1133 "exited = child pre-snap (100) + grandchild pre-snap (200)"
1134 );
1135
1136 let corrected = raw.saturating_sub(exited);
1137 assert_eq!(corrected, 130);
1139 }
1140
1141 #[test]
1159 fn test_process_cores_used_does_not_exceed_system_utilization() {
1160 let pid = i32::try_from(std::process::id()).expect("PID too large");
1161 let mut collector = CpuCollector::new(Some(pid), false);
1162
1163 let mut child = std::process::Command::new("sh")
1166 .args(["-c", "while true; do :; done"])
1167 .spawn()
1168 .expect("failed to spawn sh busy-loop -- required for T-CPU-15");
1169
1170 std::thread::sleep(std::time::Duration::from_millis(200));
1174
1175 let _ = collector.collect().expect("warm-up collect failed");
1177
1178 child.kill().ok();
1183 child.wait().ok();
1184
1185 let m = collector.collect().expect("second collect failed");
1186
1187 let proc_utime = m
1188 .process_utime_secs
1189 .expect("process_utime_secs must be Some");
1190 let proc_stime = m
1191 .process_stime_secs
1192 .expect("process_stime_secs must be Some");
1193 let proc_cpu = proc_utime + proc_stime;
1194 let sys_cpu = m.utime_secs + m.stime_secs;
1195
1196 let tolerance = sys_cpu * 0.15 + 0.05;
1201 assert!(
1202 proc_cpu <= sys_cpu + tolerance,
1203 "process CPU ({proc_cpu:.3}s = {proc_utime:.3}s utime + {proc_stime:.3}s stime) \
1204 must not exceed system CPU ({sys_cpu:.3}s) -- cutime double-counting regression \
1205 for issue #20"
1206 );
1207 }
1208
1209 #[test]
1217 fn test_process_utime_no_double_count_after_child_exits() {
1218 let pid = i32::try_from(std::process::id()).expect("PID too large");
1219 let mut collector = CpuCollector::new(Some(pid), false);
1220
1221 let mut child = std::process::Command::new("sh")
1224 .args(["-c", "for i in $(seq 1 20000); do :; done"])
1225 .spawn()
1226 .expect("failed to spawn sh -- required for T-CPU-16");
1227
1228 std::thread::sleep(std::time::Duration::from_millis(20));
1231
1232 let _ = collector.collect().expect("warm-up collect failed");
1234
1235 let _ = child.wait().expect("failed to wait for child");
1237
1238 let m = collector.collect().expect("second collect failed");
1242
1243 let proc_utime = m
1244 .process_utime_secs
1245 .expect("process_utime_secs must be Some when a PID is tracked");
1246 let sys_utime = m.utime_secs;
1247
1248 let tolerance = sys_utime * 0.05 + 0.05;
1250 assert!(
1251 proc_utime <= sys_utime + tolerance,
1252 "process_utime_secs ({proc_utime:.3}s) exceeds system utime_secs ({sys_utime:.3}s) -- \
1253 cutime double-counting regression (issue #20)"
1254 );
1255 }
1256
1257 #[test]
1274 fn test_cutime_correction_multi_interval_child_exit() {
1275 let pid = i32::try_from(std::process::id()).expect("PID too large");
1276 let mut collector = CpuCollector::new(Some(pid), false);
1277
1278 let mut child = std::process::Command::new("sh")
1280 .args(["-c", "while true; do :; done"])
1281 .spawn()
1282 .expect("failed to spawn sh busy-loop -- required for T-CPU-17");
1283
1284 std::thread::sleep(std::time::Duration::from_millis(100));
1286 let _ = collector.collect().expect("warm-up collect failed");
1287
1288 std::thread::sleep(std::time::Duration::from_millis(100));
1291 let _ = collector.collect().expect("intermediate collect failed");
1292
1293 child.kill().ok();
1297 child.wait().ok();
1298
1299 let m = collector.collect().expect("final collect failed");
1300
1301 let proc_utime = m
1302 .process_utime_secs
1303 .expect("process_utime_secs must be Some");
1304 let proc_stime = m
1305 .process_stime_secs
1306 .expect("process_stime_secs must be Some");
1307 let proc_cpu = proc_utime + proc_stime;
1308 let sys_cpu = m.utime_secs + m.stime_secs;
1309
1310 let tolerance = sys_cpu * 1.0 + 0.50;
1316 assert!(
1317 proc_cpu <= sys_cpu + tolerance,
1318 "process CPU ({proc_cpu:.3}s = {proc_utime:.3}s utime + {proc_stime:.3}s stime) \
1319 must not exceed system CPU ({sys_cpu:.3}s) across multiple intervals -- \
1320 cutime multi-interval regression for issue #20"
1321 );
1322 }
1323
1324 #[test]
1345 fn test_pss_tracks_file_backed_mapping() {
1346 use std::fs;
1347 use std::io::Write as _;
1348 use std::os::unix::io::AsRawFd;
1349
1350 const MAPPING_MIB: usize = 4;
1351 const MAPPING_SIZE: usize = MAPPING_MIB * 1024 * 1024;
1352
1353 let pid = i32::try_from(std::process::id()).expect("PID too large");
1354 let path = format!("/tmp/rt_test_pss_{}", std::process::id());
1355
1356 let (pss_before, rss_before) = process_tree_memory_mib(&[pid]);
1357
1358 {
1360 let mut f = fs::File::create(&path).expect("cannot create temp file for T-CPU-18");
1361 let chunk = vec![0xABu8; 64 * 1024];
1362 for _ in 0..(MAPPING_SIZE / chunk.len()) {
1363 f.write_all(&chunk).expect("write failed");
1364 }
1365 }
1366
1367 let file = fs::File::open(&path).expect("cannot open temp file for T-CPU-18");
1368 let ptr = unsafe {
1369 libc::mmap(
1370 std::ptr::null_mut(),
1371 MAPPING_SIZE,
1372 libc::PROT_READ,
1373 libc::MAP_PRIVATE,
1374 file.as_raw_fd(),
1375 0,
1376 )
1377 };
1378 assert_ne!(ptr, libc::MAP_FAILED, "mmap failed in T-CPU-18");
1379
1380 let slice = unsafe { std::slice::from_raw_parts(ptr as *const u8, MAPPING_SIZE) };
1382 let mut checksum = 0u64;
1383 for offset in (0..MAPPING_SIZE).step_by(4096) {
1384 checksum = checksum.wrapping_add(u64::from(slice[offset]));
1385 }
1386 let _ = checksum;
1387
1388 let (pss_after, rss_after) = process_tree_memory_mib(&[pid]);
1389
1390 unsafe { libc::munmap(ptr, MAPPING_SIZE) };
1392 fs::remove_file(&path).ok();
1393
1394 let pss_delta = pss_after.saturating_sub(pss_before);
1395 let rss_delta = rss_after.saturating_sub(rss_before);
1396
1397 const TRUNC_SLACK_MIB: u64 = 1;
1402
1403 assert!(
1404 rss_delta + TRUNC_SLACK_MIB >= MAPPING_MIB as u64,
1405 "RSS must increase by >= {MAPPING_MIB} MiB after touching the mapping: \
1406 before={rss_before} MiB, after={rss_after} MiB (delta={rss_delta} MiB)"
1407 );
1408 assert!(
1409 pss_delta + TRUNC_SLACK_MIB >= MAPPING_MIB as u64,
1410 "PSS must increase by >= {MAPPING_MIB} MiB as sole mapper of the file: \
1411 before={pss_before} MiB, after={pss_after} MiB (delta={pss_delta} MiB)"
1412 );
1413 assert!(
1414 pss_after <= rss_after + TRUNC_SLACK_MIB,
1415 "PSS ({pss_after} MiB) must not exceed RSS ({rss_after} MiB)"
1416 );
1417 let skew = pss_delta.abs_diff(rss_delta);
1421 assert!(
1422 skew <= 1 + TRUNC_SLACK_MIB,
1423 "PSS delta ({pss_delta} MiB) and RSS delta ({rss_delta} MiB) must agree within \
1424 2 MiB for a sole mapper -- larger skew indicates smaps_rollup is not being read"
1425 );
1426 }
1427
1428 #[test]
1436 fn test_mib_truncation_can_underreport_pss_delta() {
1437 let pss_before_bytes: u64 = (8 * 1024 + 1) * 1024; let pss_after_bytes: u64 = (12 * 1024 - 1) * 1024; let before_mib = (pss_before_bytes / 1024) / 1024; let after_mib = (pss_after_bytes / 1024) / 1024; let delta_mib = after_mib.saturating_sub(before_mib); assert_eq!(before_mib, 8);
1446 assert_eq!(after_mib, 11);
1447 assert_eq!(
1448 delta_mib, 3,
1449 "truncation makes ~4 MiB delta appear as 3 MiB"
1450 );
1451
1452 assert!(delta_mib < 4);
1454
1455 const TRUNC_SLACK_MIB: u64 = 1;
1457 assert!(delta_mib + TRUNC_SLACK_MIB >= 4);
1458 }
1459
1460 fn read_pss_kib(pid: i32) -> u64 {
1463 let proc_ = procfs::process::Process::new(pid).expect("process not found");
1464 proc_
1465 .smaps_rollup()
1466 .expect("smaps_rollup unavailable")
1467 .memory_map_rollup
1468 .iter()
1469 .find_map(|m| m.extension.map.get("Pss").copied())
1470 .unwrap_or(0)
1471 / 1024
1472 }
1473
1474 #[test]
1479 fn test_pss_tracks_file_backed_mapping_kib_resolution() {
1480 use std::fs;
1481 use std::io::Write as _;
1482 use std::os::unix::io::AsRawFd;
1483
1484 const MAPPING_MIB: usize = 4;
1485 const MAPPING_SIZE: usize = MAPPING_MIB * 1024 * 1024;
1486 const EXPECTED_DELTA_KIB: u64 = (MAPPING_MIB * 1024) as u64;
1487 const SLACK_KIB: u64 = 64;
1488
1489 let pid = i32::try_from(std::process::id()).expect("PID too large");
1490 let path = format!("/tmp/rt_test_pss_kib_{}", std::process::id());
1491
1492 let pss_before_kib = read_pss_kib(pid);
1493
1494 {
1495 let mut f = fs::File::create(&path).expect("cannot create temp file for T-CPU-18b");
1496 let chunk = vec![0xCDu8; 64 * 1024];
1497 for _ in 0..(MAPPING_SIZE / chunk.len()) {
1498 f.write_all(&chunk).expect("write failed");
1499 }
1500 }
1501
1502 let file = fs::File::open(&path).expect("cannot open temp file for T-CPU-18b");
1503 let ptr = unsafe {
1504 libc::mmap(
1505 std::ptr::null_mut(),
1506 MAPPING_SIZE,
1507 libc::PROT_READ,
1508 libc::MAP_PRIVATE,
1509 file.as_raw_fd(),
1510 0,
1511 )
1512 };
1513 assert_ne!(ptr, libc::MAP_FAILED, "mmap failed in T-CPU-18b");
1514
1515 let slice = unsafe { std::slice::from_raw_parts(ptr as *const u8, MAPPING_SIZE) };
1516 let mut checksum = 0u64;
1517 for offset in (0..MAPPING_SIZE).step_by(4096) {
1518 checksum = checksum.wrapping_add(u64::from(slice[offset]));
1519 }
1520 let _ = checksum;
1521
1522 let pss_after_kib = read_pss_kib(pid);
1523
1524 unsafe { libc::munmap(ptr, MAPPING_SIZE) };
1525 fs::remove_file(&path).ok();
1526
1527 let pss_delta_kib = pss_after_kib.saturating_sub(pss_before_kib);
1528
1529 assert!(
1530 pss_delta_kib + SLACK_KIB >= EXPECTED_DELTA_KIB,
1531 "PSS must increase by >= {EXPECTED_DELTA_KIB} KiB (±{SLACK_KIB} KiB) as sole \
1532 mapper: before={pss_before_kib} KiB, after={pss_after_kib} KiB \
1533 (delta={pss_delta_kib} KiB)"
1534 );
1535 }
1536
1537 #[test]
1546 fn test_cutime_correction_skipped_when_exited_exceeds_raw() {
1547 let prev: HashMap<i32, (u64, u64)> =
1548 [(1, (500, 0)), (2, (50000, 0))].iter().cloned().collect();
1549
1550 let curr: HashMap<i32, (u64, u64)> = [(1, (600, 0))].iter().cloned().collect();
1551
1552 let raw: u64 = curr
1553 .iter()
1554 .map(|(pid, &(cu, cs))| {
1555 let (pu, ps) = prev.get(pid).copied().unwrap_or((cu, cs));
1556 cu.saturating_sub(pu) + cs.saturating_sub(ps)
1557 })
1558 .sum();
1559 assert_eq!(raw, 100, "raw delta is parent's own 100 ticks");
1560
1561 let exited: u64 = prev
1562 .iter()
1563 .filter(|(pid, _)| !curr.contains_key(pid))
1564 .map(|(_, &(pu, ps))| pu + ps)
1565 .sum();
1566 assert_eq!(exited, 50000);
1567
1568 assert_eq!(raw.saturating_sub(exited), 0);
1570
1571 let corrected = if exited <= raw { raw - exited } else { raw };
1573 assert_eq!(
1574 corrected, 100,
1575 "must preserve raw delta when correction is implausible"
1576 );
1577 }
1578
1579 #[test]
1583 fn test_carry_forward_spans_gap_for_reappearing_pid() {
1584 let prev: HashMap<i32, (u64, u64)> =
1585 [(1, (500, 0)), (2, (10000, 0))].iter().cloned().collect();
1586
1587 let mut stored_prev: HashMap<i32, (u64, u64)> = [(1, (600, 0))].iter().cloned().collect();
1589 for (&pid, &ticks) in &prev {
1590 stored_prev.entry(pid).or_insert(ticks);
1591 }
1592 assert_eq!(
1593 stored_prev.get(&2),
1594 Some(&(10000, 0)),
1595 "child must be carried forward with prev ticks"
1596 );
1597
1598 let curr: HashMap<i32, (u64, u64)> =
1600 [(1, (700, 0)), (2, (11000, 0))].iter().cloned().collect();
1601
1602 let delta_with_cf: u64 = curr
1603 .iter()
1604 .map(|(pid, &(cu, cs))| {
1605 let (pu, ps) = stored_prev.get(pid).copied().unwrap_or((cu, cs));
1606 cu.saturating_sub(pu) + cs.saturating_sub(ps)
1607 })
1608 .sum();
1609 assert_eq!(
1610 delta_with_cf, 1100,
1611 "with carry-forward: parent delta (100) + child delta spanning gap (1000)"
1612 );
1613
1614 let no_cf_prev: HashMap<i32, (u64, u64)> = [(1, (600, 0))].iter().cloned().collect();
1616 let delta_without_cf: u64 = curr
1617 .iter()
1618 .map(|(pid, &(cu, cs))| {
1619 let (pu, ps) = no_cf_prev.get(pid).copied().unwrap_or((cu, cs));
1620 cu.saturating_sub(pu) + cs.saturating_sub(ps)
1621 })
1622 .sum();
1623 assert_eq!(
1624 delta_without_cf, 100,
1625 "without carry-forward: only parent delta (100), child contribution lost"
1626 );
1627 }
1628
1629 #[test]
1633 fn test_carry_forward_limited_to_one_hop() {
1634 let mut carried_forward: HashSet<i32> = HashSet::new();
1635
1636 let prev_ticks: HashMap<i32, (u64, u64)> =
1638 [(1, (500, 0)), (2, (10000, 0))].iter().cloned().collect();
1639 let mut curr_ticks: HashMap<i32, (u64, u64)> = [(1, (600, 0))].iter().cloned().collect();
1640
1641 let mut new_carried = HashSet::new();
1642 for (&pid, &ticks) in &prev_ticks {
1643 if !curr_ticks.contains_key(&pid) && !carried_forward.contains(&pid) {
1644 curr_ticks.insert(pid, ticks);
1645 new_carried.insert(pid);
1646 }
1647 }
1648 carried_forward = new_carried;
1649
1650 assert!(
1651 curr_ticks.contains_key(&2),
1652 "child must be carried forward in interval N"
1653 );
1654 assert!(
1655 carried_forward.contains(&2),
1656 "child must be in the carried-forward set"
1657 );
1658
1659 let prev_ticks_n1 = curr_ticks.clone();
1661 let mut curr_ticks_n1: HashMap<i32, (u64, u64)> = [(1, (700, 0))].iter().cloned().collect();
1662
1663 let mut new_carried_n1 = HashSet::new();
1664 for (&pid, &ticks) in &prev_ticks_n1 {
1665 if !curr_ticks_n1.contains_key(&pid) && !carried_forward.contains(&pid) {
1666 curr_ticks_n1.insert(pid, ticks);
1667 new_carried_n1.insert(pid);
1668 }
1669 }
1670
1671 assert!(
1672 !curr_ticks_n1.contains_key(&2),
1673 "child must NOT be carried forward a second time"
1674 );
1675 assert!(
1676 !new_carried_n1.contains(&2),
1677 "child must NOT be in the new carried-forward set"
1678 );
1679 }
1680}