1use crate::Result;
2use crate::settings::settings;
3#[cfg(windows)]
4use crate::shell::HideConsoleWindow;
5use miette::IntoDiagnostic;
6use once_cell::sync::Lazy;
7use std::collections::HashMap;
8#[cfg(target_os = "linux")]
9use std::os::fd::{AsRawFd, FromRawFd, OwnedFd};
10use std::sync::Mutex;
11use std::time::{Duration, Instant};
12use sysinfo::ProcessesToUpdate;
13#[cfg(windows)]
14use windows_sys::Win32::Foundation::{CloseHandle, FILETIME, HANDLE};
15#[cfg(windows)]
16use windows_sys::Win32::System::Threading::{
17 GetProcessTimes, OpenProcess, PROCESS_QUERY_LIMITED_INFORMATION,
18};
19
20type ParentToChildren = HashMap<u32, Vec<u32>>;
22
23type ProcessNames = HashMap<u32, (String, Option<String>)>;
25
26const FULL_REFRESH_INTERVAL: Duration = Duration::from_secs(5);
32
33pub struct Procs {
34 system: Mutex<sysinfo::System>,
35 last_full_refresh: Mutex<Option<Instant>>,
38}
39
40pub static PROCS: Lazy<Procs> = Lazy::new(Procs::new);
41
42impl Default for Procs {
43 fn default() -> Self {
44 Self::new()
45 }
46}
47
48impl Procs {
49 pub fn new() -> Self {
50 Self {
63 system: Mutex::new(sysinfo::System::new()),
64 last_full_refresh: Mutex::new(None),
65 }
66 }
67
68 fn lock_system(&self) -> std::sync::MutexGuard<'_, sysinfo::System> {
69 self.system.lock().unwrap_or_else(|poisoned| {
70 warn!("System mutex was poisoned, recovering");
71 poisoned.into_inner()
72 })
73 }
74
75 pub fn title(&self, pid: u32) -> Option<String> {
76 self.lock_system()
77 .process(sysinfo::Pid::from_u32(pid))
78 .map(|p| p.name().to_string_lossy().to_string())
79 }
80
81 pub fn boot_time(&self) -> u64 {
88 sysinfo::System::boot_time()
89 }
90
91 pub fn start_time(&self, pid: u32) -> Option<u64> {
97 process_start_token(pid)
98 }
99
100 #[cfg(not(windows))]
101 fn start_time_matches(&self, pid: u32, expected: u64) -> bool {
102 self.start_time(pid) == Some(expected)
103 }
104
105 pub fn is_running(&self, pid: u32) -> bool {
106 #[cfg(unix)]
112 {
113 unsafe {
114 if libc::kill(pid as i32, 0) == 0 {
115 return true;
116 }
117 std::io::Error::last_os_error().raw_os_error() != Some(libc::ESRCH)
118 }
119 }
120 #[cfg(not(unix))]
121 {
122 self.refresh_pids(&[pid]);
123 self.lock_system()
124 .process(sysinfo::Pid::from_u32(pid))
125 .is_some()
126 }
127 }
128
129 #[allow(dead_code)]
132 pub fn all_children(&self, pid: u32) -> Vec<u32> {
133 let system = self.lock_system();
134 let all = system.processes();
135 let mut children = vec![];
136 for (child_pid, process) in all {
137 let mut process = process;
138 while let Some(parent) = process.parent() {
139 if parent == sysinfo::Pid::from_u32(pid) {
140 children.push(child_pid.as_u32());
141 break;
142 }
143 match system.process(parent) {
144 Some(p) => process = p,
145 None => break,
146 }
147 }
148 }
149 children
150 }
151 pub fn collect_process_tree_info(&self) -> (ParentToChildren, ProcessNames) {
156 let system = self.lock_system();
157 let all = system.processes();
158 let mut parent_to_children: ParentToChildren = HashMap::new();
159 let mut process_info: ProcessNames = HashMap::new();
160
161 for (pid, proc) in all {
162 let pid_u32 = pid.as_u32();
163 process_info.insert(
164 pid_u32,
165 (
166 proc.name().to_string_lossy().to_string(),
167 proc.exe().map(|e| e.to_string_lossy().to_string()),
168 ),
169 );
170
171 if let Some(ppid) = proc.parent() {
172 parent_to_children
173 .entry(ppid.as_u32())
174 .or_default()
175 .push(pid_u32);
176 }
177 }
178
179 (parent_to_children, process_info)
180 }
181 pub async fn kill_process_group_async(
182 &self,
183 pid: u32,
184 stop_signal: i32,
185 stop_timeout: Option<std::time::Duration>,
186 ) -> Result<bool> {
187 tokio::task::spawn_blocking(move || {
188 PROCS.kill_process_group(pid, stop_signal, stop_timeout, None)
189 })
190 .await
191 .into_diagnostic()?
192 }
193
194 pub async fn kill_process_group_if_start_time_matches_async(
204 &self,
205 pid: u32,
206 expected_start_time: Option<u64>,
207 stop_signal: i32,
208 stop_timeout: Option<std::time::Duration>,
209 ) -> Result<bool> {
210 let Some(expected_start_time) = expected_start_time else {
211 warn!(
212 "no recorded start time for pid {pid}; refusing to signal it (identity cannot be bound to a process generation)"
213 );
214 return Ok(false);
215 };
216 tokio::task::spawn_blocking(move || {
217 PROCS.kill_process_group(pid, stop_signal, stop_timeout, Some(expected_start_time))
218 })
219 .await
220 .into_diagnostic()?
221 }
222
223 pub async fn kill_if_start_time_matches_async(
243 &self,
244 pid: u32,
245 expected_start_time: Option<u64>,
246 stop_signal: i32,
247 stop_timeout: Option<std::time::Duration>,
248 ) -> Result<bool> {
249 let Some(expected_start_time) = expected_start_time else {
250 warn!(
251 "no recorded start time for pid {pid}; refusing to signal it (identity cannot be bound to a process generation)"
252 );
253 return Ok(false);
254 };
255 tokio::task::spawn_blocking(move || {
256 PROCS.kill_if_start_time_matches(pid, expected_start_time, stop_signal, stop_timeout)
257 })
258 .await
259 .into_diagnostic()?
260 }
261
262 #[cfg(target_os = "linux")]
263 fn kill_if_start_time_matches(
264 &self,
265 pid: u32,
266 expected_start_time: u64,
267 stop_signal: i32,
268 stop_timeout: Option<std::time::Duration>,
269 ) -> Result<bool> {
270 let pidfd = match open_pidfd(pid) {
271 Ok(pidfd) => pidfd,
272 Err(err) if err.raw_os_error() == Some(libc::ESRCH) => {
273 debug!("process {pid} no longer exists");
274 return Ok(false);
275 }
276 Err(err) => {
277 return Err(miette::miette!(
278 "cannot securely identify process {pid}: {err}"
279 ));
280 }
281 };
282 if !self.verify_start_time_before_signal(pid, expected_start_time)? {
287 return Ok(false);
288 }
289 let target = [(pid, pidfd)];
290 let signal_name = signal_name(stop_signal);
291 debug!("sending {signal_name} to pinned process {pid}");
292 signal_pidfds(&target, stop_signal, signal_name)?;
293
294 let stop_timeout = stop_timeout.unwrap_or_else(|| settings().supervisor_stop_timeout());
297 let fast_ms = 10u64;
298 let slow_ms = 50u64;
299 let total_ms = stop_timeout.as_millis().max(1) as u64;
300 let fast_count = ((total_ms / fast_ms) as usize).min(10);
301 let fast_total_ms = fast_ms * fast_count as u64;
302 let slow_count = (total_ms.saturating_sub(fast_total_ms) / slow_ms) as usize;
303 for i in 0..fast_count {
304 std::thread::sleep(std::time::Duration::from_millis(fast_ms));
305 if !pidfd_is_running(&target[0].1) {
306 debug!(
307 "process {pid} terminated after {signal_name} ({} ms)",
308 (i + 1) as u64 * fast_ms
309 );
310 return Ok(true);
311 }
312 }
313 for i in 0..slow_count {
314 std::thread::sleep(std::time::Duration::from_millis(slow_ms));
315 if !pidfd_is_running(&target[0].1) {
316 debug!(
317 "process {pid} terminated after {signal_name} ({} ms)",
318 fast_total_ms + (i + 1) as u64 * slow_ms
319 );
320 return Ok(true);
321 }
322 }
323
324 warn!(
325 "process {pid} did not respond to {signal_name} after {}ms, sending SIGKILL",
326 stop_timeout.as_millis()
327 );
328 signal_pidfds(&target, libc::SIGKILL, "SIGKILL")?;
329 std::thread::sleep(std::time::Duration::from_millis(100));
331 Ok(true)
332 }
333
334 #[cfg(not(target_os = "linux"))]
335 fn kill_if_start_time_matches(
336 &self,
337 pid: u32,
338 expected_start_time: u64,
339 stop_signal: i32,
340 stop_timeout: Option<std::time::Duration>,
341 ) -> Result<bool> {
342 #[cfg(windows)]
346 let _pin = open_process_handle(pid).ok();
347
348 if !self.verify_start_time_before_signal(pid, expected_start_time)? {
349 return Ok(false);
350 }
351 self.kill(pid, stop_signal, stop_timeout, Some(expected_start_time))
352 }
353
354 fn verify_start_time_before_signal(&self, pid: u32, expected_start_time: u64) -> Result<bool> {
366 match self.start_time(pid) {
367 Some(current) if current == expected_start_time => Ok(true),
368 Some(current) => {
369 warn!(
370 "pid {pid} is not the recorded process (start time {current}, expected {expected_start_time}); not signalling it"
371 );
372 Ok(false)
373 }
374 None if self.is_running(pid) => Err(miette::miette!(
375 "cannot verify the identity of pid {pid}: its start time is unreadable; not signalling it"
376 )),
377 None => {
378 debug!("process {pid} no longer exists");
379 Ok(false)
380 }
381 }
382 }
383
384 #[cfg(unix)]
401 fn kill_process_group(
402 &self,
403 pid: u32,
404 stop_signal: i32,
405 stop_timeout: Option<std::time::Duration>,
406 expected_start_time: Option<u64>,
407 ) -> Result<bool> {
408 let pgid = pid as i32;
409 let signal_name = signal_name(stop_signal);
410
411 #[cfg(target_os = "linux")]
412 if let Some(expected) = expected_start_time {
413 return self.kill_process_group_with_pidfds(pid, expected, stop_signal, stop_timeout);
414 }
415
416 #[cfg(not(target_os = "linux"))]
428 if let Some(expected) = expected_start_time
429 && !self.start_time_matches(pid, expected)
430 {
431 debug!("process {pid} identity changed before killpg; refusing to signal it");
432 return Ok(false);
433 }
434
435 debug!("killing process group {pgid} with {signal_name}");
436
437 let ret = unsafe { libc::killpg(pgid, stop_signal) };
442 if ret == -1 {
443 let err = std::io::Error::last_os_error();
444 if err.raw_os_error() == Some(libc::ESRCH) {
445 debug!("process group {pgid} no longer exists");
446 return Ok(false);
447 }
448 if err.raw_os_error() == Some(libc::EPERM) {
449 return Err(miette::miette!(
450 "failed to send {signal_name} to process group {pgid}: permission denied"
451 ));
452 }
453 warn!("failed to send {signal_name} to process group {pgid}: {err}");
454 }
455
456 let stop_timeout = stop_timeout.unwrap_or_else(|| settings().supervisor_stop_timeout());
466 let fast_ms = 10u64;
467 let slow_ms = 50u64;
468 let total_ms = stop_timeout.as_millis().max(1) as u64;
469 let fast_count = ((total_ms / fast_ms) as usize).min(10);
470 let fast_total_ms = fast_ms * fast_count as u64;
471 let remaining_ms = total_ms.saturating_sub(fast_total_ms);
472 let slow_count = (remaining_ms / slow_ms) as usize;
473
474 let fast_checks =
475 std::iter::repeat_n(std::time::Duration::from_millis(fast_ms), fast_count);
476 let slow_checks =
477 std::iter::repeat_n(std::time::Duration::from_millis(slow_ms), slow_count);
478 let mut elapsed_ms = 0u64;
479
480 for sleep_duration in fast_checks.chain(slow_checks) {
481 std::thread::sleep(sleep_duration);
482 elapsed_ms += sleep_duration.as_millis() as u64;
483 if process_group_terminated(pgid) {
484 debug!("process group {pgid} terminated after {signal_name} ({elapsed_ms} ms)",);
485 return Ok(true);
486 }
487 }
488
489 warn!(
491 "process group {pgid} did not respond to {signal_name} after {}ms, sending SIGKILL",
492 stop_timeout.as_millis()
493 );
494 let ret = unsafe { libc::killpg(pgid, libc::SIGKILL) };
495 if ret == -1 {
496 let err = std::io::Error::last_os_error();
497 if err.raw_os_error() != Some(libc::ESRCH) {
498 warn!("failed to send SIGKILL to process group {pgid}: {err}");
499 }
500 }
501
502 for _ in 0..40 {
506 std::thread::sleep(std::time::Duration::from_millis(50));
507 if process_group_terminated(pgid) {
508 return Ok(true);
509 }
510 }
511 Err(miette::miette!(
514 "process group {pgid} still has members after SIGKILL \
515 (possibly stuck in uninterruptible sleep)"
516 ))
517 }
518
519 pub fn process_group_alive(&self, pid: u32) -> bool {
522 #[cfg(unix)]
523 {
524 !process_group_terminated(pid as i32)
525 }
526 #[cfg(not(unix))]
527 {
528 self.is_running(pid)
529 }
530 }
531
532 #[cfg(target_os = "linux")]
533 fn kill_process_group_with_pidfds(
534 &self,
535 pid: u32,
536 expected_start_time: u64,
537 _stop_signal: i32,
538 stop_timeout: Option<std::time::Duration>,
539 ) -> Result<bool> {
540 let leader = match open_pidfd(pid) {
541 Ok(pidfd) => pidfd,
542 Err(err) => {
543 warn!("cannot securely identify process group {pid}: {err}");
544 return Ok(false);
545 }
546 };
547 if !self.start_time_matches(pid, expected_start_time) {
548 debug!("process group {pid} leader identity changed before signaling");
549 return Ok(false);
550 }
551
552 let mut members = vec![(pid, leader)];
553 if let Err(err) = stop_pidfds(&members) {
554 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
555 return Err(err);
556 }
557 if !pidfd_is_running(&members[0].1) {
558 debug!("process group {pid} leader exited before it could be frozen");
559 return Ok(false);
560 }
561
562 loop {
566 let known_members = members.len();
567 let added = match extend_process_group_pidfds(pid as i32, &mut members) {
568 Ok(added) => added,
569 Err(err) => {
570 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
571 return Err(miette::miette!(
572 "failed to scan pinned process group {pid}: {err}"
573 ));
574 }
575 };
576 if added == 0 {
577 break;
578 }
579 if let Err(err) = stop_pidfds(&members[known_members..]) {
580 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
581 return Err(err);
582 }
583 }
584
585 warn!(
586 "force-terminating {} pinned orphan process(es) in group {pid}",
587 members.len()
588 );
589 if let Err(err) = signal_pidfds(&members, libc::SIGKILL, "SIGKILL") {
590 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
591 return Err(err);
592 }
593
594 let exit_timeout = stop_timeout.unwrap_or_else(|| settings().supervisor_stop_timeout());
595 let checks = exit_timeout.as_millis().max(1).div_ceil(50) as usize;
596 for _ in 0..checks {
597 if members.iter().all(|(_, pidfd)| !pidfd_is_running(pidfd)) {
598 return Ok(true);
599 }
600 std::thread::sleep(std::time::Duration::from_millis(50));
601 }
602 if members.iter().all(|(_, pidfd)| !pidfd_is_running(pidfd)) {
603 return Ok(true);
604 }
605
606 warn!("one or more pinned processes in orphan group {pid} remained alive after SIGKILL");
607 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
608 Ok(false)
609 }
610
611 #[cfg(not(unix))]
612 fn kill_process_group(
613 &self,
614 pid: u32,
615 stop_signal: i32,
616 stop_timeout: Option<std::time::Duration>,
617 expected_start_time: Option<u64>,
618 ) -> Result<bool> {
619 #[cfg(windows)]
622 let _identity_handle = if let Some(expected) = expected_start_time {
623 let handle = match open_process_handle(pid) {
624 Ok(handle) => handle,
625 Err(err) => {
626 warn!("cannot securely identify process {pid}: {err}");
627 return Ok(false);
628 }
629 };
630 if process_start_token_from_handle(handle.0) != Some(expected) {
631 debug!("process {pid} identity changed before taskkill");
632 return Ok(false);
633 }
634 Some(handle)
635 } else {
636 None
637 };
638
639 #[cfg(not(windows))]
640 if let Some(expected) = expected_start_time
641 && !self.start_time_matches(pid, expected)
642 {
643 debug!("process {pid} identity changed before termination");
644 return Ok(false);
645 }
646
647 self.kill(pid, stop_signal, stop_timeout, expected_start_time)
648 }
649
650 #[cfg(not(target_os = "linux"))]
668 fn kill(
669 &self,
670 pid: u32,
671 stop_signal: i32,
672 stop_timeout: Option<std::time::Duration>,
673 expected_start_time: Option<u64>,
674 ) -> Result<bool> {
675 debug!("killing process {pid}");
676
677 #[cfg(windows)]
678 {
679 let _ = expected_start_time;
682 if stop_signal == crate::config_types::StopSignal::SIGINT {
692 let stop_timeout =
693 stop_timeout.unwrap_or_else(|| settings().supervisor_stop_timeout());
694 if interrupt_and_wait(pid, stop_timeout) {
695 debug!("process {pid} exited after Ctrl+C");
696 return Ok(true);
697 }
698 }
699 let output = std::process::Command::new("taskkill")
704 .args(["/F", "/T", "/PID"])
705 .arg(pid.to_string())
706 .hide_console_window()
707 .output();
708 let taskkill_succeeded = match output {
709 Ok(o) if o.status.success() => {
710 debug!("taskkill /F /T /PID {pid} succeeded");
711 true
712 }
713 Ok(o) => {
714 debug!(
715 "taskkill /F /T /PID {pid} exited with status {}: {}",
716 o.status,
717 String::from_utf8_lossy(&o.stderr).trim()
718 );
719 false
720 }
721 Err(e) => {
722 debug!("failed to spawn taskkill for pid {pid}: {e}");
723 false
724 }
725 };
726 std::thread::sleep(std::time::Duration::from_millis(200));
730 if !taskkill_succeeded && self.is_running(pid) {
731 return Err(miette::miette!(
732 "taskkill failed and process {pid} is still running"
733 ));
734 }
735 Ok(true)
736 }
737
738 #[cfg(unix)]
739 {
740 let sysinfo_pid = sysinfo::Pid::from_u32(pid);
741 let signal_name = signal_name(stop_signal);
742 if let Some(expected) = expected_start_time
747 && !self.verify_start_time_before_signal(pid, expected)?
748 {
749 return Ok(false);
750 }
751 debug!("sending {signal_name} to process {pid}");
755 let ret = unsafe { libc::kill(pid as i32, stop_signal) };
756 if ret == -1 {
757 let err = std::io::Error::last_os_error();
758 if err.raw_os_error() == Some(libc::ESRCH) {
759 debug!("process {pid} no longer exists");
760 return Ok(false);
761 }
762 if err.raw_os_error() == Some(libc::EPERM) {
763 return Err(miette::miette!(
764 "failed to send {signal_name} to process {pid}: permission denied"
765 ));
766 }
767 return Err(miette::miette!(
768 "failed to send {signal_name} to process {pid}: {err}"
769 ));
770 }
771
772 let stop_timeout = stop_timeout.unwrap_or_else(|| settings().supervisor_stop_timeout());
775 let fast_ms = 10u64;
776 let slow_ms = 50u64;
777 let total_ms = stop_timeout.as_millis().max(1) as u64;
778 let fast_count = ((total_ms / fast_ms) as usize).min(10);
779 let fast_total_ms = fast_ms * fast_count as u64;
780 let remaining_ms = total_ms.saturating_sub(fast_total_ms);
781 let slow_count = (remaining_ms / slow_ms) as usize;
782
783 for i in 0..fast_count {
784 std::thread::sleep(std::time::Duration::from_millis(fast_ms));
785 self.refresh_pids(&[pid]);
786 if self.is_terminated_or_zombie(sysinfo_pid) {
787 debug!(
788 "process {pid} terminated after {signal_name} ({} ms)",
789 (i + 1) * fast_ms as usize
790 );
791 return Ok(true);
792 }
793 }
794
795 for i in 0..slow_count {
797 std::thread::sleep(std::time::Duration::from_millis(slow_ms));
798 self.refresh_pids(&[pid]);
799 if self.is_terminated_or_zombie(sysinfo_pid) {
800 debug!(
801 "process {pid} terminated after {signal_name} ({} ms)",
802 fast_total_ms + (i + 1) as u64 * slow_ms
803 );
804 return Ok(true);
805 }
806 }
807
808 if let Some(expected) = expected_start_time
813 && !self.start_time_matches(pid, expected)
814 {
815 debug!(
816 "process {pid} exited during the stop timeout and its PID was recycled; not sending SIGKILL"
817 );
818 return Ok(true);
819 }
820 warn!(
821 "process {pid} did not respond to {signal_name} after {}ms, sending SIGKILL",
822 stop_timeout.as_millis()
823 );
824 let ret = unsafe { libc::kill(pid as i32, libc::SIGKILL) };
825 if ret == -1 {
826 let err = std::io::Error::last_os_error();
827 if err.raw_os_error() != Some(libc::ESRCH) {
828 warn!("failed to send SIGKILL to process {pid}: {err}");
829 }
830 }
831
832 std::thread::sleep(std::time::Duration::from_millis(100));
834 Ok(true)
835 }
836 }
837
838 #[cfg(all(unix, not(target_os = "linux")))]
842 fn is_terminated_or_zombie(&self, sysinfo_pid: sysinfo::Pid) -> bool {
843 let system = self.lock_system();
844 match system.process(sysinfo_pid) {
845 None => true,
846 Some(process) => {
847 matches!(process.status(), sysinfo::ProcessStatus::Zombie)
848 }
849 }
850 }
851
852 pub(crate) fn refresh_processes(&self) {
853 let mut system = self.lock_system();
854 system.refresh_processes(ProcessesToUpdate::All, true);
855 #[cfg(windows)]
860 system.refresh_cpu_usage();
861 }
862
863 pub(crate) fn refresh_pids(&self, pids: &[u32]) {
866 let sysinfo_pids: Vec<sysinfo::Pid> =
867 pids.iter().map(|p| sysinfo::Pid::from_u32(*p)).collect();
868 self.lock_system()
869 .refresh_processes(ProcessesToUpdate::Some(&sysinfo_pids), true);
870 }
871
872 pub(crate) fn refresh_if_stale(&self) {
889 let mut last = self
890 .last_full_refresh
891 .lock()
892 .unwrap_or_else(|poisoned| poisoned.into_inner());
893 if last.is_none_or(|t| t.elapsed() >= FULL_REFRESH_INTERVAL) {
894 self.refresh_processes();
895 *last = Some(Instant::now());
896 }
897 }
898
899 pub fn get_batch_group_stats(&self, pids: &[u32]) -> Vec<(u32, Option<ProcessStats>)> {
905 if pids.is_empty() {
906 return Vec::new();
907 }
908
909 let system = self.lock_system();
910 let processes = system.processes();
911
912 let now = std::time::SystemTime::now()
913 .duration_since(std::time::UNIX_EPOCH)
914 .map(|d| d.as_secs())
915 .unwrap_or(0);
916
917 let mut children_map: std::collections::HashMap<sysinfo::Pid, Vec<sysinfo::Pid>> =
919 std::collections::HashMap::new();
920 for (child_pid, child) in processes {
921 if child.thread_kind().is_some() {
924 continue;
925 }
926 if let Some(ppid) = child.parent() {
927 children_map.entry(ppid).or_default().push(*child_pid);
928 }
929 }
930
931 pids.iter()
932 .map(|&pid| {
933 let root_pid = sysinfo::Pid::from_u32(pid);
934 let Some(root) = processes.get(&root_pid) else {
935 return (pid, None);
936 };
937
938 let root_disk = root.disk_usage();
939 let mut stats = ProcessStats {
940 cpu_percent: root.cpu_usage(),
941 memory_bytes: root.memory(),
942 uptime_secs: now.saturating_sub(root.start_time()),
943 disk_read_bytes: root_disk.read_bytes,
944 disk_write_bytes: root_disk.written_bytes,
945 };
946
947 let mut queue = std::collections::VecDeque::new();
949 if let Some(direct_children) = children_map.get(&root_pid) {
950 queue.extend(direct_children);
951 }
952 while let Some(child_pid) = queue.pop_front() {
953 if let Some(child) = processes.get(&child_pid) {
954 let disk = child.disk_usage();
955 stats.cpu_percent += child.cpu_usage();
956 stats.memory_bytes += child.memory();
957 stats.disk_read_bytes += disk.read_bytes;
958 stats.disk_write_bytes += disk.written_bytes;
959 }
960 if let Some(grandchildren) = children_map.get(&child_pid) {
961 queue.extend(grandchildren);
962 }
963 }
964
965 (pid, Some(stats))
966 })
967 .collect()
968 }
969 pub fn refresh_and_get_batch_stats(&self, pids: &[u32]) -> HashMap<u32, ProcessStats> {
979 self.refresh_processes();
980 self.get_batch_group_stats(pids)
981 .into_iter()
982 .filter_map(|(pid, stats)| stats.map(|s| (pid, s)))
983 .collect()
984 }
985
986 pub fn refresh_and_get_batch_stats_if_stale(&self, pids: &[u32]) -> HashMap<u32, ProcessStats> {
993 self.refresh_if_stale();
994 self.get_batch_group_stats(pids)
995 .into_iter()
996 .filter_map(|(pid, stats)| stats.map(|s| (pid, s)))
997 .collect()
998 }
999
1000 pub fn get_batch_tree_stats_map(&self, pids: &[u32]) -> HashMap<u32, ProcessStats> {
1002 self.get_batch_group_stats(pids)
1003 .into_iter()
1004 .filter_map(|(pid, stats)| stats.map(|stats| (pid, stats)))
1005 .collect()
1006 }
1007
1008 pub fn get_stats(&self, pid: u32) -> Option<ProcessStats> {
1010 self.get_batch_group_stats(&[pid])
1011 .into_iter()
1012 .next()
1013 .and_then(|(_, stats)| stats)
1014 }
1015
1016 pub fn get_extended_stats(&self, pid: u32) -> Option<ExtendedProcessStats> {
1018 let system = self.lock_system();
1019 let processes = system.processes();
1020 let root_pid = sysinfo::Pid::from_u32(pid);
1021 let p = processes.get(&root_pid)?;
1022
1023 let now = std::time::SystemTime::now()
1024 .duration_since(std::time::UNIX_EPOCH)
1025 .map(|d| d.as_secs())
1026 .unwrap_or(0);
1027
1028 let root_disk = p.disk_usage();
1029 let mut aggregate_stats = ProcessStats {
1030 cpu_percent: p.cpu_usage(),
1031 memory_bytes: p.memory(),
1032 uptime_secs: now.saturating_sub(p.start_time()),
1033 disk_read_bytes: root_disk.read_bytes,
1034 disk_write_bytes: root_disk.written_bytes,
1035 };
1036
1037 let mut children_map: HashMap<sysinfo::Pid, Vec<sysinfo::Pid>> = HashMap::new();
1038 for (child_pid, child) in processes {
1039 if let Some(ppid) = child.parent() {
1040 children_map.entry(ppid).or_default().push(*child_pid);
1041 }
1042 }
1043
1044 let mut queue = std::collections::VecDeque::new();
1045 if let Some(direct_children) = children_map.get(&root_pid) {
1046 queue.extend(direct_children);
1047 }
1048 while let Some(child_pid) = queue.pop_front() {
1049 if let Some(child) = processes.get(&child_pid) {
1050 let disk = child.disk_usage();
1051 aggregate_stats.cpu_percent += child.cpu_usage();
1052 aggregate_stats.memory_bytes += child.memory();
1053 aggregate_stats.disk_read_bytes += disk.read_bytes;
1054 aggregate_stats.disk_write_bytes += disk.written_bytes;
1055 }
1056 if let Some(grandchildren) = children_map.get(&child_pid) {
1057 queue.extend(grandchildren);
1058 }
1059 }
1060
1061 Some(ExtendedProcessStats {
1062 name: p.name().to_string_lossy().to_string(),
1063 status: format!("{:?}", p.status()),
1064 cpu_percent: aggregate_stats.cpu_percent,
1065 memory_bytes: aggregate_stats.memory_bytes,
1066 virtual_memory_bytes: p.virtual_memory(),
1067 uptime_secs: aggregate_stats.uptime_secs,
1068 thread_count: p.tasks().map(|t| t.len()).unwrap_or(0),
1069 })
1070 }
1071}
1072
1073#[cfg(target_os = "linux")]
1074fn process_start_token(pid: u32) -> Option<u64> {
1075 let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
1076 let command_end = stat.rfind(')')?;
1077 stat.get(command_end + 1..)?
1079 .split_whitespace()
1080 .nth(19)?
1081 .parse()
1082 .ok()
1083}
1084
1085#[cfg(target_os = "macos")]
1086fn process_start_token(pid: u32) -> Option<u64> {
1087 let mut info = std::mem::MaybeUninit::<libc::proc_bsdinfo>::zeroed();
1088 let size = std::mem::size_of::<libc::proc_bsdinfo>() as i32;
1089 let read = unsafe {
1090 libc::proc_pidinfo(
1091 pid as i32,
1092 libc::PROC_PIDTBSDINFO,
1093 0,
1094 info.as_mut_ptr().cast(),
1095 size,
1096 )
1097 };
1098 if read != size {
1099 return None;
1100 }
1101 let info = unsafe { info.assume_init() };
1102 info.pbi_start_tvsec
1103 .checked_mul(1_000_000)?
1104 .checked_add(info.pbi_start_tvusec)
1105}
1106
1107#[cfg(windows)]
1108fn process_start_token(pid: u32) -> Option<u64> {
1109 let handle = open_process_handle(pid).ok()?;
1110 process_start_token_from_handle(handle.0)
1111}
1112
1113#[cfg(windows)]
1114fn process_start_token_from_handle(handle: HANDLE) -> Option<u64> {
1115 let mut creation = FILETIME {
1116 dwLowDateTime: 0,
1117 dwHighDateTime: 0,
1118 };
1119 let mut exit = creation;
1120 let mut kernel = creation;
1121 let mut user = creation;
1122 let ok = unsafe { GetProcessTimes(handle, &mut creation, &mut exit, &mut kernel, &mut user) };
1123 if ok == 0 {
1124 return None;
1125 }
1126
1127 Some((u64::from(creation.dwHighDateTime) << 32) | u64::from(creation.dwLowDateTime))
1128}
1129
1130#[cfg(windows)]
1131struct ProcessHandle(HANDLE);
1132
1133#[cfg(windows)]
1134impl Drop for ProcessHandle {
1135 fn drop(&mut self) {
1136 unsafe {
1137 CloseHandle(self.0);
1138 }
1139 }
1140}
1141
1142#[cfg(windows)]
1143fn open_process_handle(pid: u32) -> std::io::Result<ProcessHandle> {
1144 let handle = unsafe { OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, 0, pid) };
1145 if handle.is_null() {
1146 return Err(std::io::Error::last_os_error());
1147 }
1148 Ok(ProcessHandle(handle))
1149}
1150
1151#[cfg(windows)]
1157fn interrupt_and_wait(pid: u32, timeout: Duration) -> bool {
1158 use crate::console_ctrl::EXIT_SHARED_CONSOLE;
1159 use windows_sys::Win32::Foundation::WAIT_OBJECT_0;
1160 use windows_sys::Win32::System::Threading::{PROCESS_SYNCHRONIZE, WaitForSingleObject};
1161
1162 let handle = unsafe { OpenProcess(PROCESS_SYNCHRONIZE, 0, pid) };
1165 if handle.is_null() {
1166 debug!(
1167 "cannot wait for process {pid}: {}",
1168 std::io::Error::last_os_error()
1169 );
1170 return false;
1171 }
1172 let handle = ProcessHandle(handle);
1173
1174 let deadline = Instant::now().checked_add(timeout);
1179 let millis_left = || {
1180 let Some(deadline) = deadline else {
1181 return u32::MAX - 1;
1182 };
1183 let left = deadline.saturating_duration_since(Instant::now());
1184 u32::try_from(left.as_millis()).unwrap_or(u32::MAX - 1)
1185 };
1186
1187 let mut helper = match std::process::Command::new(&*crate::env::PITCHFORK_BIN)
1188 .args(["interrupt", "--pid"])
1189 .arg(pid.to_string())
1190 .arg("--supervisor-pid")
1191 .arg(std::process::id().to_string())
1192 .stdin(std::process::Stdio::null())
1193 .stdout(std::process::Stdio::null())
1194 .stderr(std::process::Stdio::piped())
1195 .hide_console_window()
1196 .spawn()
1197 {
1198 Ok(helper) => helper,
1199 Err(e) => {
1200 debug!("failed to spawn pitchfork interrupt for pid {pid}: {e}");
1201 return false;
1202 }
1203 };
1204 let helper_handle = std::os::windows::io::AsRawHandle::as_raw_handle(&helper);
1205 if unsafe { WaitForSingleObject(helper_handle as HANDLE, millis_left()) } != WAIT_OBJECT_0 {
1206 debug!("pitchfork interrupt for pid {pid} did not finish within {timeout:?}");
1207 let _ = helper.kill();
1208 let _ = helper.wait();
1209 return false;
1210 }
1211 match helper.wait_with_output() {
1212 Ok(o) if o.status.success() => {}
1213 Ok(o) if o.status.code() == Some(EXIT_SHARED_CONSOLE) => {
1214 debug!("process {pid} shares the supervisor's console; not sending Ctrl+C");
1215 return false;
1216 }
1217 Ok(o) => {
1218 debug!(
1219 "pitchfork interrupt for pid {pid} exited with status {}: {}",
1220 o.status,
1221 String::from_utf8_lossy(&o.stderr).trim()
1222 );
1223 return false;
1224 }
1225 Err(e) => {
1226 debug!("failed to wait for pitchfork interrupt for pid {pid}: {e}");
1227 return false;
1228 }
1229 }
1230
1231 debug!("sent Ctrl+C to process {pid}, waiting up to {timeout:?} in all");
1232 unsafe { WaitForSingleObject(handle.0, millis_left()) == WAIT_OBJECT_0 }
1233}
1234
1235#[cfg(not(any(target_os = "linux", target_os = "macos", windows)))]
1236fn process_start_token(pid: u32) -> Option<u64> {
1237 let mut system = sysinfo::System::new();
1238 let sysinfo_pid = sysinfo::Pid::from_u32(pid);
1239 system.refresh_processes(ProcessesToUpdate::Some(&[sysinfo_pid]), true);
1240 system
1241 .process(sysinfo_pid)
1242 .map(|process| process.start_time())
1243}
1244
1245#[cfg(target_os = "linux")]
1246fn open_pidfd(pid: u32) -> std::io::Result<OwnedFd> {
1247 let fd = unsafe { libc::syscall(libc::SYS_pidfd_open, pid, 0) };
1248 if fd < 0 {
1249 return Err(std::io::Error::last_os_error());
1250 }
1251 Ok(unsafe { OwnedFd::from_raw_fd(fd as i32) })
1252}
1253
1254#[cfg(target_os = "linux")]
1255fn pidfd_is_running(pidfd: &OwnedFd) -> bool {
1256 match try_pidfd_is_running(pidfd) {
1257 Ok(running) => running,
1258 Err(err) => {
1259 warn!("failed to poll pidfd {}: {err}", pidfd.as_raw_fd());
1260 true
1261 }
1262 }
1263}
1264
1265#[cfg(target_os = "linux")]
1266fn try_pidfd_is_running(pidfd: &OwnedFd) -> std::io::Result<bool> {
1267 let mut pollfd = libc::pollfd {
1268 fd: pidfd.as_raw_fd(),
1269 events: libc::POLLIN,
1270 revents: 0,
1271 };
1272 let result = unsafe { libc::poll(&mut pollfd, 1, 0) };
1273 if result < 0 {
1274 return Err(std::io::Error::last_os_error());
1275 }
1276 Ok(result == 0)
1277}
1278
1279#[cfg(target_os = "linux")]
1280fn signal_pidfds(members: &[(u32, OwnedFd)], signal: i32, signal_name: &str) -> Result<()> {
1281 for (pid, pidfd) in members {
1282 if !pidfd_is_running(pidfd) {
1283 continue;
1284 }
1285 let result = unsafe {
1286 libc::syscall(
1287 libc::SYS_pidfd_send_signal,
1288 pidfd.as_raw_fd(),
1289 signal,
1290 std::ptr::null::<libc::siginfo_t>(),
1291 0,
1292 )
1293 };
1294 if result == -1 {
1295 let err = std::io::Error::last_os_error();
1296 if err.raw_os_error() == Some(libc::ESRCH) {
1297 continue;
1298 }
1299 return Err(miette::miette!(
1300 "failed to send {signal_name} to pinned process {pid}: {err}"
1301 ));
1302 }
1303 }
1304 Ok(())
1305}
1306
1307#[cfg(target_os = "linux")]
1308fn stop_pidfds(members: &[(u32, OwnedFd)]) -> Result<()> {
1309 signal_pidfds(members, libc::SIGSTOP, "SIGSTOP")?;
1310 for _ in 0..200 {
1311 if members.iter().all(|(pid, pidfd)| {
1312 !pidfd_is_running(pidfd) || matches!(linux_process_state(*pid), Some('T' | 't'))
1313 }) {
1314 return Ok(());
1315 }
1316 std::thread::sleep(std::time::Duration::from_millis(5));
1317 }
1318 Err(miette::miette!(
1319 "timed out while freezing orphan process group"
1320 ))
1321}
1322
1323#[cfg(target_os = "linux")]
1324fn extend_process_group_pidfds(
1325 pgid: i32,
1326 members: &mut Vec<(u32, OwnedFd)>,
1327) -> std::io::Result<usize> {
1328 let entries = std::fs::read_dir("/proc")?;
1329 let mut added = 0;
1330 for entry in entries {
1331 let entry = entry?;
1332 let Some(pid) = entry
1333 .file_name()
1334 .to_str()
1335 .and_then(|name| name.parse::<u32>().ok())
1336 else {
1337 continue;
1338 };
1339 let Some(observed_identity) = linux_process_identity(pid) else {
1340 continue;
1341 };
1342 if observed_identity.0 != pgid {
1343 continue;
1344 }
1345 let mut already_pinned = false;
1346 for (known_pid, pidfd) in members.iter() {
1347 if *known_pid == pid && try_pidfd_is_running(pidfd)? {
1348 already_pinned = true;
1349 break;
1350 }
1351 }
1352 if already_pinned {
1353 continue;
1354 }
1355
1356 let pidfd = match open_pidfd(pid) {
1357 Ok(pidfd) => pidfd,
1358 Err(err) if err.raw_os_error() == Some(libc::ESRCH) => continue,
1359 Err(err) => return Err(err),
1360 };
1361 if linux_process_identity(pid) != Some(observed_identity) {
1362 return Err(std::io::Error::other(format!(
1363 "process {pid} identity changed while pinning group {pgid}"
1364 )));
1365 }
1366 members.push((pid, pidfd));
1367 added += 1;
1368 }
1369 Ok(added)
1370}
1371
1372#[cfg(target_os = "linux")]
1373fn linux_process_identity(pid: u32) -> Option<(i32, u64)> {
1374 let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
1375 let command_end = stat.rfind(')')?;
1376 let fields: Vec<_> = stat.get(command_end + 1..)?.split_whitespace().collect();
1377 Some((fields.get(2)?.parse().ok()?, fields.get(19)?.parse().ok()?))
1380}
1381
1382#[cfg(target_os = "linux")]
1383fn linux_process_state(pid: u32) -> Option<char> {
1384 let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
1385 let command_end = stat.rfind(')')?;
1386 stat.get(command_end + 1..)?
1387 .split_whitespace()
1388 .next()?
1389 .chars()
1390 .next()
1391}
1392
1393#[derive(Debug, Clone, Copy)]
1394pub struct ProcessStats {
1395 pub cpu_percent: f32,
1396 pub memory_bytes: u64,
1397 pub uptime_secs: u64,
1398 pub disk_read_bytes: u64,
1399 pub disk_write_bytes: u64,
1400}
1401
1402impl ProcessStats {
1403 pub fn memory_display(&self) -> String {
1404 format_bytes(self.memory_bytes)
1405 }
1406
1407 pub fn cpu_display(&self) -> String {
1408 format!("{:.1}%", self.cpu_percent)
1409 }
1410
1411 pub fn uptime_display(&self) -> String {
1412 format_duration(self.uptime_secs)
1413 }
1414
1415 pub fn disk_read_display(&self) -> String {
1416 format_bytes_per_sec(self.disk_read_bytes)
1417 }
1418
1419 pub fn disk_write_display(&self) -> String {
1420 format_bytes_per_sec(self.disk_write_bytes)
1421 }
1422}
1423
1424#[derive(Debug, Clone)]
1425pub struct ExtendedProcessStats {
1426 pub name: String,
1427 pub status: String,
1428 pub cpu_percent: f32,
1429 pub memory_bytes: u64,
1430 pub virtual_memory_bytes: u64,
1431 pub uptime_secs: u64,
1432 pub thread_count: usize,
1433}
1434
1435fn format_bytes(bytes: u64) -> String {
1436 humanbyte::to_string(bytes, humanbyte::Format::IEC)
1437}
1438
1439pub(crate) fn format_duration(secs: u64) -> String {
1440 if secs < 60 {
1441 format!("{secs}s")
1442 } else if secs < 3600 {
1443 format!("{}m {}s", secs / 60, secs % 60)
1444 } else if secs < 86400 {
1445 let hours = secs / 3600;
1446 let mins = (secs % 3600) / 60;
1447 format!("{hours}h {mins}m")
1448 } else {
1449 let days = secs / 86400;
1450 let hours = (secs % 86400) / 3600;
1451 format!("{days}d {hours}h")
1452 }
1453}
1454
1455fn format_bytes_per_sec(bytes: u64) -> String {
1456 format!("{}/s", humanbyte::to_string(bytes, humanbyte::Format::IEC))
1457}
1458
1459#[cfg(unix)]
1475fn process_group_terminated(pgid: i32) -> bool {
1476 unsafe { libc::killpg(pgid, 0) != 0 }
1477}
1478
1479#[cfg(unix)]
1480fn signal_name(sig: i32) -> &'static str {
1481 match sig {
1482 libc::SIGHUP => "SIGHUP",
1483 libc::SIGINT => "SIGINT",
1484 libc::SIGQUIT => "SIGQUIT",
1485 libc::SIGTERM => "SIGTERM",
1486 libc::SIGUSR1 => "SIGUSR1",
1487 libc::SIGUSR2 => "SIGUSR2",
1488 libc::SIGKILL => "SIGKILL",
1489 _ => "UNKNOWN",
1490 }
1491}
1492
1493#[cfg(test)]
1494mod format_tests {
1495 use super::*;
1496
1497 #[test]
1498 fn process_start_time_check_rejects_mismatch() {
1499 let procs = Procs::new();
1500 let pid = std::process::id();
1501 procs.refresh_pids(&[pid]);
1502 let actual = procs
1503 .start_time(pid)
1504 .expect("current process should have a start time");
1505
1506 assert_ne!(procs.start_time(pid), Some(actual.saturating_add(1)));
1507 }
1508
1509 #[test]
1510 fn test_format_bytes() {
1511 assert_eq!(format_bytes(512), "512 B");
1512 assert_eq!(format_bytes(1024), "1.0 KiB");
1513 assert_eq!(format_bytes(1536), "1.5 KiB");
1514 assert_eq!(format_bytes(50 * 1024 * 1024), "50.0 MiB");
1515 assert_eq!(format_bytes(3 * 1024 * 1024 * 1024), "3.0 GiB");
1516 assert_eq!(format_bytes(1100 * 1024 * 1024 * 1024), "1.1 TiB");
1518 }
1519
1520 #[test]
1521 fn test_format_bytes_per_sec() {
1522 assert_eq!(format_bytes_per_sec(512), "512 B/s");
1523 assert_eq!(format_bytes_per_sec(1536), "1.5 KiB/s");
1524 assert_eq!(format_bytes_per_sec(2 * 1024 * 1024), "2.0 MiB/s");
1525 }
1526}
1527
1528#[cfg(all(test, unix))]
1529mod tests {
1530 use super::*;
1531 use std::os::unix::process::CommandExt;
1532 use std::process::{Child, Command, Stdio};
1533 use std::time::{Duration, Instant};
1534
1535 struct ChildGuard(Child);
1536
1537 impl Drop for ChildGuard {
1538 fn drop(&mut self) {
1539 let pid = self.0.id() as i32;
1540 let _ = unsafe { libc::killpg(pid, libc::SIGKILL) };
1542 let _ = self.0.wait();
1543 }
1544 }
1545
1546 #[tokio::test]
1547 async fn orphan_identity_checked_group_kill_rejects_mismatch() {
1548 let mut command = Command::new("sleep");
1549 command
1550 .arg("30")
1551 .stdin(Stdio::null())
1552 .stdout(Stdio::null())
1553 .stderr(Stdio::null());
1554 unsafe {
1555 command.pre_exec(|| {
1556 if libc::setsid() == -1 {
1557 return Err(std::io::Error::last_os_error());
1558 }
1559 Ok(())
1560 });
1561 }
1562
1563 let child = command.spawn().expect("failed to spawn test process");
1564 let pid = child.id();
1565 let _child = ChildGuard(child);
1566
1567 PROCS.refresh_pids(&[pid]);
1568 let actual_start_time = PROCS
1569 .start_time(pid)
1570 .expect("test process should have a start time");
1571
1572 let killed = PROCS
1573 .kill_process_group_if_start_time_matches_async(
1574 pid,
1575 Some(actual_start_time.saturating_add(1)),
1576 libc::SIGTERM,
1577 Some(Duration::from_millis(100)),
1578 )
1579 .await
1580 .expect("identity-checked kill should not error");
1581
1582 assert!(!killed);
1583 assert!(PROCS.is_running(pid), "mismatched process must survive");
1584 }
1585
1586 #[cfg(all(unix, not(target_os = "linux")))]
1591 #[tokio::test]
1592 async fn identity_checked_group_kill_reverifies_inside_blocking_op() {
1593 let mut command = Command::new("sleep");
1594 command
1595 .arg("30")
1596 .stdin(Stdio::null())
1597 .stdout(Stdio::null())
1598 .stderr(Stdio::null());
1599 unsafe {
1600 command.pre_exec(|| {
1601 if libc::setsid() == -1 {
1602 return Err(std::io::Error::last_os_error());
1603 }
1604 Ok(())
1605 });
1606 }
1607
1608 let child = command.spawn().expect("failed to spawn test process");
1609 let pid = child.id();
1610 let mut child = ChildGuard(child);
1611
1612 PROCS.refresh_pids(&[pid]);
1613 let actual_start_time = PROCS
1614 .start_time(pid)
1615 .expect("test process should have a start time");
1616
1617 let killed = PROCS
1618 .kill_process_group_if_start_time_matches_async(
1619 pid,
1620 Some(actual_start_time),
1621 libc::SIGTERM,
1622 Some(Duration::from_millis(100)),
1623 )
1624 .await
1625 .expect("identity-checked kill should not error");
1626
1627 assert!(killed, "matching generation must be signalled");
1628 let deadline = Instant::now() + Duration::from_secs(2);
1631 while child.0.try_wait().unwrap().is_none() {
1632 assert!(Instant::now() < deadline, "signalled child must exit");
1633 tokio::time::sleep(Duration::from_millis(10)).await;
1634 }
1635 assert!(
1636 !PROCS.is_running(pid),
1637 "signalled process group must be gone"
1638 );
1639 }
1640
1641 #[test]
1642 fn get_stats_includes_descendant_rss() {
1643 let mut command = Command::new("sh");
1644 command
1645 .args(["-c", "sleep 30 & wait"])
1646 .stdin(Stdio::null())
1647 .stdout(Stdio::null())
1648 .stderr(Stdio::null());
1649 unsafe {
1650 command.pre_exec(|| {
1651 if libc::setsid() == -1 {
1652 return Err(std::io::Error::last_os_error());
1653 }
1654 Ok(())
1655 });
1656 }
1657
1658 let parent = command.spawn().expect("failed to spawn process tree");
1659 let parent_pid = parent.id();
1660 let _parent = ChildGuard(parent);
1661
1662 let procs = Procs::new();
1663 let deadline = Instant::now() + Duration::from_secs(5);
1664 let mut child_pids = Vec::new();
1665 while Instant::now() < deadline {
1666 procs.refresh_processes();
1667 child_pids = procs.all_children(parent_pid);
1668 if !child_pids.is_empty() {
1669 break;
1670 }
1671 std::thread::sleep(Duration::from_millis(50));
1672 }
1673 assert!(
1674 !child_pids.is_empty(),
1675 "test process tree did not appear under parent pid {parent_pid}"
1676 );
1677
1678 procs.refresh_processes();
1679 child_pids = procs.all_children(parent_pid);
1680 assert!(
1681 !child_pids.is_empty(),
1682 "test process tree disappeared under parent pid {parent_pid}"
1683 );
1684 let root_pid = sysinfo::Pid::from_u32(parent_pid);
1685 let direct_memory = {
1686 let system = procs.lock_system();
1687 system
1688 .process(root_pid)
1689 .expect("parent process should exist")
1690 .memory()
1691 };
1692 let descendant_memory = {
1693 let system = procs.lock_system();
1694 child_pids
1695 .iter()
1696 .filter_map(|pid| system.process(sysinfo::Pid::from_u32(*pid)))
1697 .map(|process| process.memory())
1698 .sum::<u64>()
1699 };
1700 assert!(
1701 descendant_memory > 0,
1702 "descendants {child_pids:?} should have nonzero RSS"
1703 );
1704
1705 let stats = procs
1706 .get_stats(parent_pid)
1707 .expect("parent process should have aggregate stats");
1708
1709 assert_eq!(
1710 stats.memory_bytes,
1711 direct_memory + descendant_memory,
1712 "get_stats should include descendant RSS for parent pid {parent_pid}; \
1713 descendants: {child_pids:?}, direct RSS: {direct_memory}, \
1714 descendant RSS: {descendant_memory}, reported RSS: {}",
1715 stats.memory_bytes
1716 );
1717 }
1718
1719 #[test]
1720 fn full_refresh_is_throttled_by_ttl() {
1721 let procs = Procs::new();
1722
1723 assert!(
1725 procs.last_full_refresh.lock().unwrap().is_none(),
1726 "fresh Procs should have no recorded refresh"
1727 );
1728 procs.refresh_if_stale();
1729 let first = procs.last_full_refresh.lock().unwrap().unwrap();
1730 assert!(first.elapsed() < FULL_REFRESH_INTERVAL);
1731
1732 procs.refresh_if_stale();
1735 let second = procs.last_full_refresh.lock().unwrap().unwrap();
1736 assert_eq!(
1737 second, first,
1738 "refresh_if_stale within TTL must skip the refresh and keep the timestamp"
1739 );
1740
1741 let expired = Instant::now()
1744 .checked_sub(FULL_REFRESH_INTERVAL + Duration::from_secs(1))
1745 .expect("system has been up long enough to backdate by 6s");
1746 *procs.last_full_refresh.lock().unwrap() = Some(expired);
1747 procs.refresh_if_stale();
1748 let third = procs.last_full_refresh.lock().unwrap().unwrap();
1749 assert!(
1750 third > expired,
1751 "refresh_if_stale after expired TTL must refresh and advance the timestamp"
1752 );
1753 assert!(
1754 third.elapsed() < FULL_REFRESH_INTERVAL,
1755 "fresh timestamp after expired-TTL refresh should be recent"
1756 );
1757 }
1758}