1use crate::Result;
2#[cfg(unix)]
3use crate::settings::settings;
4use miette::IntoDiagnostic;
5use once_cell::sync::Lazy;
6use std::collections::HashMap;
7#[cfg(target_os = "linux")]
8use std::os::fd::{AsRawFd, FromRawFd, OwnedFd};
9#[cfg(windows)]
10use std::os::windows::process::CommandExt;
11use std::sync::Mutex;
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
26pub struct Procs {
27 system: Mutex<sysinfo::System>,
28}
29
30pub static PROCS: Lazy<Procs> = Lazy::new(Procs::new);
31
32impl Default for Procs {
33 fn default() -> Self {
34 Self::new()
35 }
36}
37
38impl Procs {
39 pub fn new() -> Self {
40 Self {
53 system: Mutex::new(sysinfo::System::new()),
54 }
55 }
56
57 fn lock_system(&self) -> std::sync::MutexGuard<'_, sysinfo::System> {
58 self.system.lock().unwrap_or_else(|poisoned| {
59 warn!("System mutex was poisoned, recovering");
60 poisoned.into_inner()
61 })
62 }
63
64 pub fn title(&self, pid: u32) -> Option<String> {
65 self.lock_system()
66 .process(sysinfo::Pid::from_u32(pid))
67 .map(|p| p.name().to_string_lossy().to_string())
68 }
69
70 pub fn boot_time(&self) -> u64 {
77 sysinfo::System::boot_time()
78 }
79
80 pub fn start_time(&self, pid: u32) -> Option<u64> {
86 process_start_token(pid)
87 }
88
89 #[cfg(any(target_os = "linux", windows))]
90 fn start_time_matches(&self, pid: u32, expected: u64) -> bool {
91 self.start_time(pid) == Some(expected)
92 }
93
94 pub fn is_running(&self, pid: u32) -> bool {
95 #[cfg(unix)]
101 {
102 unsafe {
103 if libc::kill(pid as i32, 0) == 0 {
104 return true;
105 }
106 std::io::Error::last_os_error().raw_os_error() != Some(libc::ESRCH)
107 }
108 }
109 #[cfg(not(unix))]
110 {
111 self.refresh_pids(&[pid]);
112 self.lock_system()
113 .process(sysinfo::Pid::from_u32(pid))
114 .is_some()
115 }
116 }
117
118 #[allow(dead_code)]
121 pub fn all_children(&self, pid: u32) -> Vec<u32> {
122 let system = self.lock_system();
123 let all = system.processes();
124 let mut children = vec![];
125 for (child_pid, process) in all {
126 let mut process = process;
127 while let Some(parent) = process.parent() {
128 if parent == sysinfo::Pid::from_u32(pid) {
129 children.push(child_pid.as_u32());
130 break;
131 }
132 match system.process(parent) {
133 Some(p) => process = p,
134 None => break,
135 }
136 }
137 }
138 children
139 }
140 pub fn collect_process_tree_info(&self) -> (ParentToChildren, ProcessNames) {
145 let system = self.lock_system();
146 let all = system.processes();
147 let mut parent_to_children: ParentToChildren = HashMap::new();
148 let mut process_info: ProcessNames = HashMap::new();
149
150 for (pid, proc) in all {
151 let pid_u32 = pid.as_u32();
152 process_info.insert(
153 pid_u32,
154 (
155 proc.name().to_string_lossy().to_string(),
156 proc.exe().map(|e| e.to_string_lossy().to_string()),
157 ),
158 );
159
160 if let Some(ppid) = proc.parent() {
161 parent_to_children
162 .entry(ppid.as_u32())
163 .or_default()
164 .push(pid_u32);
165 }
166 }
167
168 (parent_to_children, process_info)
169 }
170 pub async fn kill_process_group_async(
171 &self,
172 pid: u32,
173 stop_signal: i32,
174 stop_timeout: Option<std::time::Duration>,
175 ) -> Result<bool> {
176 tokio::task::spawn_blocking(move || {
177 PROCS.kill_process_group(pid, stop_signal, stop_timeout, None)
178 })
179 .await
180 .into_diagnostic()?
181 }
182
183 pub async fn kill_process_group_if_start_time_matches_async(
190 &self,
191 pid: u32,
192 expected_start_time: u64,
193 stop_signal: i32,
194 stop_timeout: Option<std::time::Duration>,
195 ) -> Result<bool> {
196 tokio::task::spawn_blocking(move || {
197 PROCS.kill_process_group(pid, stop_signal, stop_timeout, Some(expected_start_time))
198 })
199 .await
200 .into_diagnostic()?
201 }
202
203 #[cfg(unix)]
220 fn kill_process_group(
221 &self,
222 pid: u32,
223 stop_signal: i32,
224 stop_timeout: Option<std::time::Duration>,
225 expected_start_time: Option<u64>,
226 ) -> Result<bool> {
227 let pgid = pid as i32;
228 let signal_name = signal_name(stop_signal);
229
230 #[cfg(target_os = "linux")]
231 if let Some(expected) = expected_start_time {
232 return self.kill_process_group_with_pidfds(pid, expected, stop_signal, stop_timeout);
233 }
234
235 #[cfg(not(target_os = "linux"))]
239 if expected_start_time.is_some() {
240 warn!(
241 "cannot securely identify process group {pgid} on this platform; refusing to signal it"
242 );
243 return Ok(false);
244 }
245
246 debug!("killing process group {pgid} with {signal_name}");
247
248 let ret = unsafe { libc::killpg(pgid, stop_signal) };
253 if ret == -1 {
254 let err = std::io::Error::last_os_error();
255 if err.raw_os_error() == Some(libc::ESRCH) {
256 debug!("process group {pgid} no longer exists");
257 return Ok(false);
258 }
259 if err.raw_os_error() == Some(libc::EPERM) {
260 return Err(miette::miette!(
261 "failed to send {signal_name} to process group {pgid}: permission denied"
262 ));
263 }
264 warn!("failed to send {signal_name} to process group {pgid}: {err}");
265 }
266
267 let stop_timeout = stop_timeout.unwrap_or_else(|| settings().supervisor_stop_timeout());
277 let fast_ms = 10u64;
278 let slow_ms = 50u64;
279 let total_ms = stop_timeout.as_millis().max(1) as u64;
280 let fast_count = ((total_ms / fast_ms) as usize).min(10);
281 let fast_total_ms = fast_ms * fast_count as u64;
282 let remaining_ms = total_ms.saturating_sub(fast_total_ms);
283 let slow_count = (remaining_ms / slow_ms) as usize;
284
285 let fast_checks =
286 std::iter::repeat_n(std::time::Duration::from_millis(fast_ms), fast_count);
287 let slow_checks =
288 std::iter::repeat_n(std::time::Duration::from_millis(slow_ms), slow_count);
289 let mut elapsed_ms = 0u64;
290
291 for sleep_duration in fast_checks.chain(slow_checks) {
292 std::thread::sleep(sleep_duration);
293 elapsed_ms += sleep_duration.as_millis() as u64;
294 if process_group_terminated(pgid) {
295 debug!("process group {pgid} terminated after {signal_name} ({elapsed_ms} ms)",);
296 return Ok(true);
297 }
298 }
299
300 warn!(
302 "process group {pgid} did not respond to {signal_name} after {}ms, sending SIGKILL",
303 stop_timeout.as_millis()
304 );
305 let ret = unsafe { libc::killpg(pgid, libc::SIGKILL) };
306 if ret == -1 {
307 let err = std::io::Error::last_os_error();
308 if err.raw_os_error() != Some(libc::ESRCH) {
309 warn!("failed to send SIGKILL to process group {pgid}: {err}");
310 }
311 }
312
313 for _ in 0..40 {
317 std::thread::sleep(std::time::Duration::from_millis(50));
318 if process_group_terminated(pgid) {
319 return Ok(true);
320 }
321 }
322 Err(miette::miette!(
325 "process group {pgid} still has members after SIGKILL \
326 (possibly stuck in uninterruptible sleep)"
327 ))
328 }
329
330 pub fn process_group_alive(&self, pid: u32) -> bool {
333 #[cfg(unix)]
334 {
335 !process_group_terminated(pid as i32)
336 }
337 #[cfg(not(unix))]
338 {
339 self.is_running(pid)
340 }
341 }
342
343 #[cfg(target_os = "linux")]
344 fn kill_process_group_with_pidfds(
345 &self,
346 pid: u32,
347 expected_start_time: u64,
348 _stop_signal: i32,
349 stop_timeout: Option<std::time::Duration>,
350 ) -> Result<bool> {
351 let leader = match open_pidfd(pid) {
352 Ok(pidfd) => pidfd,
353 Err(err) => {
354 warn!("cannot securely identify process group {pid}: {err}");
355 return Ok(false);
356 }
357 };
358 if !self.start_time_matches(pid, expected_start_time) {
359 debug!("process group {pid} leader identity changed before signaling");
360 return Ok(false);
361 }
362
363 let mut members = vec![(pid, leader)];
364 if let Err(err) = stop_pidfds(&members) {
365 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
366 return Err(err);
367 }
368 if !pidfd_is_running(&members[0].1) {
369 debug!("process group {pid} leader exited before it could be frozen");
370 return Ok(false);
371 }
372
373 loop {
377 let known_members = members.len();
378 let added = match extend_process_group_pidfds(pid as i32, &mut members) {
379 Ok(added) => added,
380 Err(err) => {
381 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
382 return Err(miette::miette!(
383 "failed to scan pinned process group {pid}: {err}"
384 ));
385 }
386 };
387 if added == 0 {
388 break;
389 }
390 if let Err(err) = stop_pidfds(&members[known_members..]) {
391 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
392 return Err(err);
393 }
394 }
395
396 warn!(
397 "force-terminating {} pinned orphan process(es) in group {pid}",
398 members.len()
399 );
400 if let Err(err) = signal_pidfds(&members, libc::SIGKILL, "SIGKILL") {
401 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
402 return Err(err);
403 }
404
405 let exit_timeout = stop_timeout.unwrap_or_else(|| settings().supervisor_stop_timeout());
406 let checks = exit_timeout.as_millis().max(1).div_ceil(50) as usize;
407 for _ in 0..checks {
408 if members.iter().all(|(_, pidfd)| !pidfd_is_running(pidfd)) {
409 return Ok(true);
410 }
411 std::thread::sleep(std::time::Duration::from_millis(50));
412 }
413 if members.iter().all(|(_, pidfd)| !pidfd_is_running(pidfd)) {
414 return Ok(true);
415 }
416
417 warn!("one or more pinned processes in orphan group {pid} remained alive after SIGKILL");
418 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
419 Ok(false)
420 }
421
422 #[cfg(not(unix))]
423 fn kill_process_group(
424 &self,
425 pid: u32,
426 _stop_signal: i32,
427 _stop_timeout: Option<std::time::Duration>,
428 expected_start_time: Option<u64>,
429 ) -> Result<bool> {
430 #[cfg(windows)]
433 let _identity_handle = if let Some(expected) = expected_start_time {
434 let handle = match open_process_handle(pid) {
435 Ok(handle) => handle,
436 Err(err) => {
437 warn!("cannot securely identify process {pid}: {err}");
438 return Ok(false);
439 }
440 };
441 if process_start_token_from_handle(handle.0) != Some(expected) {
442 debug!("process {pid} identity changed before taskkill");
443 return Ok(false);
444 }
445 Some(handle)
446 } else {
447 None
448 };
449
450 #[cfg(not(windows))]
451 if let Some(expected) = expected_start_time
452 && !self.start_time_matches(pid, expected)
453 {
454 debug!("process {pid} identity changed before termination");
455 return Ok(false);
456 }
457
458 self.kill(pid, 0, None)
459 }
460
461 pub async fn kill_async(
462 &self,
463 pid: u32,
464 stop_signal: i32,
465 stop_timeout: Option<std::time::Duration>,
466 ) -> Result<bool> {
467 tokio::task::spawn_blocking(move || PROCS.kill(pid, stop_signal, stop_timeout))
468 .await
469 .into_diagnostic()?
470 }
471
472 fn kill(
482 &self,
483 pid: u32,
484 stop_signal: i32,
485 stop_timeout: Option<std::time::Duration>,
486 ) -> Result<bool> {
487 let sysinfo_pid = sysinfo::Pid::from_u32(pid);
488
489 debug!("killing process {pid}");
490
491 #[cfg(windows)]
492 {
493 let _ = (stop_signal, stop_timeout);
494 let output = std::process::Command::new("taskkill")
499 .args(["/F", "/T", "/PID"])
500 .arg(pid.to_string())
501 .creation_flags(0x08000000) .output();
503 let taskkill_succeeded = match output {
504 Ok(o) if o.status.success() => {
505 debug!("taskkill /F /T /PID {pid} succeeded");
506 true
507 }
508 Ok(o) => {
509 debug!(
510 "taskkill /F /T /PID {pid} exited with status {}: {}",
511 o.status,
512 String::from_utf8_lossy(&o.stderr).trim()
513 );
514 false
515 }
516 Err(e) => {
517 debug!("failed to spawn taskkill for pid {pid}: {e}");
518 false
519 }
520 };
521 std::thread::sleep(std::time::Duration::from_millis(200));
525 if !taskkill_succeeded && self.is_running(pid) {
526 return Err(miette::miette!(
527 "taskkill failed and process {pid} is still running"
528 ));
529 }
530 Ok(true)
531 }
532
533 #[cfg(unix)]
534 {
535 let signal_name = signal_name(stop_signal);
536 debug!("sending {signal_name} to process {pid}");
540 let ret = unsafe { libc::kill(pid as i32, stop_signal) };
541 if ret == -1 {
542 let err = std::io::Error::last_os_error();
543 if err.raw_os_error() == Some(libc::ESRCH) {
544 debug!("process {pid} no longer exists");
545 return Ok(false);
546 }
547 if err.raw_os_error() == Some(libc::EPERM) {
548 return Err(miette::miette!(
549 "failed to send {signal_name} to process {pid}: permission denied"
550 ));
551 }
552 return Err(miette::miette!(
553 "failed to send {signal_name} to process {pid}: {err}"
554 ));
555 }
556
557 let stop_timeout = stop_timeout.unwrap_or_else(|| settings().supervisor_stop_timeout());
560 let fast_ms = 10u64;
561 let slow_ms = 50u64;
562 let total_ms = stop_timeout.as_millis().max(1) as u64;
563 let fast_count = ((total_ms / fast_ms) as usize).min(10);
564 let fast_total_ms = fast_ms * fast_count as u64;
565 let remaining_ms = total_ms.saturating_sub(fast_total_ms);
566 let slow_count = (remaining_ms / slow_ms) as usize;
567
568 for i in 0..fast_count {
569 std::thread::sleep(std::time::Duration::from_millis(fast_ms));
570 self.refresh_pids(&[pid]);
571 if self.is_terminated_or_zombie(sysinfo_pid) {
572 debug!(
573 "process {pid} terminated after {signal_name} ({} ms)",
574 (i + 1) * fast_ms as usize
575 );
576 return Ok(true);
577 }
578 }
579
580 for i in 0..slow_count {
582 std::thread::sleep(std::time::Duration::from_millis(slow_ms));
583 self.refresh_pids(&[pid]);
584 if self.is_terminated_or_zombie(sysinfo_pid) {
585 debug!(
586 "process {pid} terminated after {signal_name} ({} ms)",
587 fast_total_ms + (i + 1) as u64 * slow_ms
588 );
589 return Ok(true);
590 }
591 }
592
593 warn!(
595 "process {pid} did not respond to {signal_name} after {}ms, sending SIGKILL",
596 stop_timeout.as_millis()
597 );
598 let ret = unsafe { libc::kill(pid as i32, libc::SIGKILL) };
599 if ret == -1 {
600 let err = std::io::Error::last_os_error();
601 if err.raw_os_error() != Some(libc::ESRCH) {
602 warn!("failed to send SIGKILL to process {pid}: {err}");
603 }
604 }
605
606 std::thread::sleep(std::time::Duration::from_millis(100));
608 Ok(true)
609 }
610 }
611
612 #[cfg(unix)]
616 fn is_terminated_or_zombie(&self, sysinfo_pid: sysinfo::Pid) -> bool {
617 let system = self.lock_system();
618 match system.process(sysinfo_pid) {
619 None => true,
620 Some(process) => {
621 matches!(process.status(), sysinfo::ProcessStatus::Zombie)
622 }
623 }
624 }
625
626 pub(crate) fn refresh_processes(&self) {
627 let mut system = self.lock_system();
628 system.refresh_processes(ProcessesToUpdate::All, true);
629 #[cfg(windows)]
634 system.refresh_cpu_usage();
635 }
636
637 pub(crate) fn refresh_pids(&self, pids: &[u32]) {
640 let sysinfo_pids: Vec<sysinfo::Pid> =
641 pids.iter().map(|p| sysinfo::Pid::from_u32(*p)).collect();
642 self.lock_system()
643 .refresh_processes(ProcessesToUpdate::Some(&sysinfo_pids), true);
644 }
645
646 pub fn get_batch_group_stats(&self, pids: &[u32]) -> Vec<(u32, Option<ProcessStats>)> {
652 if pids.is_empty() {
653 return Vec::new();
654 }
655
656 let system = self.lock_system();
657 let processes = system.processes();
658
659 let now = std::time::SystemTime::now()
660 .duration_since(std::time::UNIX_EPOCH)
661 .map(|d| d.as_secs())
662 .unwrap_or(0);
663
664 let mut children_map: std::collections::HashMap<sysinfo::Pid, Vec<sysinfo::Pid>> =
666 std::collections::HashMap::new();
667 for (child_pid, child) in processes {
668 if child.thread_kind().is_some() {
671 continue;
672 }
673 if let Some(ppid) = child.parent() {
674 children_map.entry(ppid).or_default().push(*child_pid);
675 }
676 }
677
678 pids.iter()
679 .map(|&pid| {
680 let root_pid = sysinfo::Pid::from_u32(pid);
681 let Some(root) = processes.get(&root_pid) else {
682 return (pid, None);
683 };
684
685 let root_disk = root.disk_usage();
686 let mut stats = ProcessStats {
687 cpu_percent: root.cpu_usage(),
688 memory_bytes: root.memory(),
689 uptime_secs: now.saturating_sub(root.start_time()),
690 disk_read_bytes: root_disk.read_bytes,
691 disk_write_bytes: root_disk.written_bytes,
692 };
693
694 let mut queue = std::collections::VecDeque::new();
696 if let Some(direct_children) = children_map.get(&root_pid) {
697 queue.extend(direct_children);
698 }
699 while let Some(child_pid) = queue.pop_front() {
700 if let Some(child) = processes.get(&child_pid) {
701 let disk = child.disk_usage();
702 stats.cpu_percent += child.cpu_usage();
703 stats.memory_bytes += child.memory();
704 stats.disk_read_bytes += disk.read_bytes;
705 stats.disk_write_bytes += disk.written_bytes;
706 }
707 if let Some(grandchildren) = children_map.get(&child_pid) {
708 queue.extend(grandchildren);
709 }
710 }
711
712 (pid, Some(stats))
713 })
714 .collect()
715 }
716 pub fn refresh_and_get_batch_stats(&self, pids: &[u32]) -> HashMap<u32, ProcessStats> {
723 self.refresh_processes();
724 self.get_batch_group_stats(pids)
725 .into_iter()
726 .filter_map(|(pid, stats)| stats.map(|s| (pid, s)))
727 .collect()
728 }
729
730 pub fn get_batch_tree_stats_map(&self, pids: &[u32]) -> HashMap<u32, ProcessStats> {
732 self.get_batch_group_stats(pids)
733 .into_iter()
734 .filter_map(|(pid, stats)| stats.map(|stats| (pid, stats)))
735 .collect()
736 }
737
738 pub fn get_stats(&self, pid: u32) -> Option<ProcessStats> {
740 self.get_batch_group_stats(&[pid])
741 .into_iter()
742 .next()
743 .and_then(|(_, stats)| stats)
744 }
745
746 pub fn get_extended_stats(&self, pid: u32) -> Option<ExtendedProcessStats> {
748 let system = self.lock_system();
749 let processes = system.processes();
750 let root_pid = sysinfo::Pid::from_u32(pid);
751 let p = processes.get(&root_pid)?;
752
753 let now = std::time::SystemTime::now()
754 .duration_since(std::time::UNIX_EPOCH)
755 .map(|d| d.as_secs())
756 .unwrap_or(0);
757
758 let root_disk = p.disk_usage();
759 let mut aggregate_stats = ProcessStats {
760 cpu_percent: p.cpu_usage(),
761 memory_bytes: p.memory(),
762 uptime_secs: now.saturating_sub(p.start_time()),
763 disk_read_bytes: root_disk.read_bytes,
764 disk_write_bytes: root_disk.written_bytes,
765 };
766
767 let mut children_map: HashMap<sysinfo::Pid, Vec<sysinfo::Pid>> = HashMap::new();
768 for (child_pid, child) in processes {
769 if let Some(ppid) = child.parent() {
770 children_map.entry(ppid).or_default().push(*child_pid);
771 }
772 }
773
774 let mut queue = std::collections::VecDeque::new();
775 if let Some(direct_children) = children_map.get(&root_pid) {
776 queue.extend(direct_children);
777 }
778 while let Some(child_pid) = queue.pop_front() {
779 if let Some(child) = processes.get(&child_pid) {
780 let disk = child.disk_usage();
781 aggregate_stats.cpu_percent += child.cpu_usage();
782 aggregate_stats.memory_bytes += child.memory();
783 aggregate_stats.disk_read_bytes += disk.read_bytes;
784 aggregate_stats.disk_write_bytes += disk.written_bytes;
785 }
786 if let Some(grandchildren) = children_map.get(&child_pid) {
787 queue.extend(grandchildren);
788 }
789 }
790
791 Some(ExtendedProcessStats {
792 name: p.name().to_string_lossy().to_string(),
793 status: format!("{:?}", p.status()),
794 cpu_percent: aggregate_stats.cpu_percent,
795 memory_bytes: aggregate_stats.memory_bytes,
796 virtual_memory_bytes: p.virtual_memory(),
797 uptime_secs: aggregate_stats.uptime_secs,
798 thread_count: p.tasks().map(|t| t.len()).unwrap_or(0),
799 })
800 }
801}
802
803#[cfg(target_os = "linux")]
804fn process_start_token(pid: u32) -> Option<u64> {
805 let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
806 let command_end = stat.rfind(')')?;
807 stat.get(command_end + 1..)?
809 .split_whitespace()
810 .nth(19)?
811 .parse()
812 .ok()
813}
814
815#[cfg(target_os = "macos")]
816fn process_start_token(pid: u32) -> Option<u64> {
817 let mut info = std::mem::MaybeUninit::<libc::proc_bsdinfo>::zeroed();
818 let size = std::mem::size_of::<libc::proc_bsdinfo>() as i32;
819 let read = unsafe {
820 libc::proc_pidinfo(
821 pid as i32,
822 libc::PROC_PIDTBSDINFO,
823 0,
824 info.as_mut_ptr().cast(),
825 size,
826 )
827 };
828 if read != size {
829 return None;
830 }
831 let info = unsafe { info.assume_init() };
832 info.pbi_start_tvsec
833 .checked_mul(1_000_000)?
834 .checked_add(info.pbi_start_tvusec)
835}
836
837#[cfg(windows)]
838fn process_start_token(pid: u32) -> Option<u64> {
839 let handle = open_process_handle(pid).ok()?;
840 process_start_token_from_handle(handle.0)
841}
842
843#[cfg(windows)]
844fn process_start_token_from_handle(handle: HANDLE) -> Option<u64> {
845 let mut creation = FILETIME {
846 dwLowDateTime: 0,
847 dwHighDateTime: 0,
848 };
849 let mut exit = creation;
850 let mut kernel = creation;
851 let mut user = creation;
852 let ok = unsafe { GetProcessTimes(handle, &mut creation, &mut exit, &mut kernel, &mut user) };
853 if ok == 0 {
854 return None;
855 }
856
857 Some((u64::from(creation.dwHighDateTime) << 32) | u64::from(creation.dwLowDateTime))
858}
859
860#[cfg(windows)]
861struct ProcessHandle(HANDLE);
862
863#[cfg(windows)]
864impl Drop for ProcessHandle {
865 fn drop(&mut self) {
866 unsafe {
867 CloseHandle(self.0);
868 }
869 }
870}
871
872#[cfg(windows)]
873fn open_process_handle(pid: u32) -> std::io::Result<ProcessHandle> {
874 let handle = unsafe { OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, 0, pid) };
875 if handle.is_null() {
876 return Err(std::io::Error::last_os_error());
877 }
878 Ok(ProcessHandle(handle))
879}
880
881#[cfg(not(any(target_os = "linux", target_os = "macos", windows)))]
882fn process_start_token(pid: u32) -> Option<u64> {
883 let mut system = sysinfo::System::new();
884 let sysinfo_pid = sysinfo::Pid::from_u32(pid);
885 system.refresh_processes(ProcessesToUpdate::Some(&[sysinfo_pid]), true);
886 system
887 .process(sysinfo_pid)
888 .map(|process| process.start_time())
889}
890
891#[cfg(target_os = "linux")]
892fn open_pidfd(pid: u32) -> std::io::Result<OwnedFd> {
893 let fd = unsafe { libc::syscall(libc::SYS_pidfd_open, pid, 0) };
894 if fd < 0 {
895 return Err(std::io::Error::last_os_error());
896 }
897 Ok(unsafe { OwnedFd::from_raw_fd(fd as i32) })
898}
899
900#[cfg(target_os = "linux")]
901fn pidfd_is_running(pidfd: &OwnedFd) -> bool {
902 match try_pidfd_is_running(pidfd) {
903 Ok(running) => running,
904 Err(err) => {
905 warn!("failed to poll pidfd {}: {err}", pidfd.as_raw_fd());
906 true
907 }
908 }
909}
910
911#[cfg(target_os = "linux")]
912fn try_pidfd_is_running(pidfd: &OwnedFd) -> std::io::Result<bool> {
913 let mut pollfd = libc::pollfd {
914 fd: pidfd.as_raw_fd(),
915 events: libc::POLLIN,
916 revents: 0,
917 };
918 let result = unsafe { libc::poll(&mut pollfd, 1, 0) };
919 if result < 0 {
920 return Err(std::io::Error::last_os_error());
921 }
922 Ok(result == 0)
923}
924
925#[cfg(target_os = "linux")]
926fn signal_pidfds(members: &[(u32, OwnedFd)], signal: i32, signal_name: &str) -> Result<()> {
927 for (pid, pidfd) in members {
928 if !pidfd_is_running(pidfd) {
929 continue;
930 }
931 let result = unsafe {
932 libc::syscall(
933 libc::SYS_pidfd_send_signal,
934 pidfd.as_raw_fd(),
935 signal,
936 std::ptr::null::<libc::siginfo_t>(),
937 0,
938 )
939 };
940 if result == -1 {
941 let err = std::io::Error::last_os_error();
942 if err.raw_os_error() == Some(libc::ESRCH) {
943 continue;
944 }
945 return Err(miette::miette!(
946 "failed to send {signal_name} to pinned process {pid}: {err}"
947 ));
948 }
949 }
950 Ok(())
951}
952
953#[cfg(target_os = "linux")]
954fn stop_pidfds(members: &[(u32, OwnedFd)]) -> Result<()> {
955 signal_pidfds(members, libc::SIGSTOP, "SIGSTOP")?;
956 for _ in 0..200 {
957 if members.iter().all(|(pid, pidfd)| {
958 !pidfd_is_running(pidfd) || matches!(linux_process_state(*pid), Some('T' | 't'))
959 }) {
960 return Ok(());
961 }
962 std::thread::sleep(std::time::Duration::from_millis(5));
963 }
964 Err(miette::miette!(
965 "timed out while freezing orphan process group"
966 ))
967}
968
969#[cfg(target_os = "linux")]
970fn extend_process_group_pidfds(
971 pgid: i32,
972 members: &mut Vec<(u32, OwnedFd)>,
973) -> std::io::Result<usize> {
974 let entries = std::fs::read_dir("/proc")?;
975 let mut added = 0;
976 for entry in entries {
977 let entry = entry?;
978 let Some(pid) = entry
979 .file_name()
980 .to_str()
981 .and_then(|name| name.parse::<u32>().ok())
982 else {
983 continue;
984 };
985 let Some(observed_identity) = linux_process_identity(pid) else {
986 continue;
987 };
988 if observed_identity.0 != pgid {
989 continue;
990 }
991 let mut already_pinned = false;
992 for (known_pid, pidfd) in members.iter() {
993 if *known_pid == pid && try_pidfd_is_running(pidfd)? {
994 already_pinned = true;
995 break;
996 }
997 }
998 if already_pinned {
999 continue;
1000 }
1001
1002 let pidfd = match open_pidfd(pid) {
1003 Ok(pidfd) => pidfd,
1004 Err(err) if err.raw_os_error() == Some(libc::ESRCH) => continue,
1005 Err(err) => return Err(err),
1006 };
1007 if linux_process_identity(pid) != Some(observed_identity) {
1008 return Err(std::io::Error::other(format!(
1009 "process {pid} identity changed while pinning group {pgid}"
1010 )));
1011 }
1012 members.push((pid, pidfd));
1013 added += 1;
1014 }
1015 Ok(added)
1016}
1017
1018#[cfg(target_os = "linux")]
1019fn linux_process_identity(pid: u32) -> Option<(i32, u64)> {
1020 let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
1021 let command_end = stat.rfind(')')?;
1022 let fields: Vec<_> = stat.get(command_end + 1..)?.split_whitespace().collect();
1023 Some((fields.get(2)?.parse().ok()?, fields.get(19)?.parse().ok()?))
1026}
1027
1028#[cfg(target_os = "linux")]
1029fn linux_process_state(pid: u32) -> Option<char> {
1030 let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
1031 let command_end = stat.rfind(')')?;
1032 stat.get(command_end + 1..)?
1033 .split_whitespace()
1034 .next()?
1035 .chars()
1036 .next()
1037}
1038
1039#[derive(Debug, Clone, Copy)]
1040pub struct ProcessStats {
1041 pub cpu_percent: f32,
1042 pub memory_bytes: u64,
1043 pub uptime_secs: u64,
1044 pub disk_read_bytes: u64,
1045 pub disk_write_bytes: u64,
1046}
1047
1048impl ProcessStats {
1049 pub fn memory_display(&self) -> String {
1050 format_bytes(self.memory_bytes)
1051 }
1052
1053 pub fn cpu_display(&self) -> String {
1054 format!("{:.1}%", self.cpu_percent)
1055 }
1056
1057 pub fn uptime_display(&self) -> String {
1058 format_duration(self.uptime_secs)
1059 }
1060
1061 pub fn disk_read_display(&self) -> String {
1062 format_bytes_per_sec(self.disk_read_bytes)
1063 }
1064
1065 pub fn disk_write_display(&self) -> String {
1066 format_bytes_per_sec(self.disk_write_bytes)
1067 }
1068}
1069
1070#[derive(Debug, Clone)]
1071pub struct ExtendedProcessStats {
1072 pub name: String,
1073 pub status: String,
1074 pub cpu_percent: f32,
1075 pub memory_bytes: u64,
1076 pub virtual_memory_bytes: u64,
1077 pub uptime_secs: u64,
1078 pub thread_count: usize,
1079}
1080
1081fn format_bytes(bytes: u64) -> String {
1082 humanbyte::to_string(bytes, humanbyte::Format::IEC)
1083}
1084
1085fn format_duration(secs: u64) -> String {
1086 if secs < 60 {
1087 format!("{secs}s")
1088 } else if secs < 3600 {
1089 format!("{}m {}s", secs / 60, secs % 60)
1090 } else if secs < 86400 {
1091 let hours = secs / 3600;
1092 let mins = (secs % 3600) / 60;
1093 format!("{hours}h {mins}m")
1094 } else {
1095 let days = secs / 86400;
1096 let hours = (secs % 86400) / 3600;
1097 format!("{days}d {hours}h")
1098 }
1099}
1100
1101fn format_bytes_per_sec(bytes: u64) -> String {
1102 format!("{}/s", humanbyte::to_string(bytes, humanbyte::Format::IEC))
1103}
1104
1105#[cfg(unix)]
1121fn process_group_terminated(pgid: i32) -> bool {
1122 unsafe { libc::killpg(pgid, 0) != 0 }
1123}
1124
1125#[cfg(unix)]
1126fn signal_name(sig: i32) -> &'static str {
1127 match sig {
1128 libc::SIGHUP => "SIGHUP",
1129 libc::SIGINT => "SIGINT",
1130 libc::SIGQUIT => "SIGQUIT",
1131 libc::SIGTERM => "SIGTERM",
1132 libc::SIGUSR1 => "SIGUSR1",
1133 libc::SIGUSR2 => "SIGUSR2",
1134 libc::SIGKILL => "SIGKILL",
1135 _ => "UNKNOWN",
1136 }
1137}
1138
1139#[cfg(test)]
1140mod format_tests {
1141 use super::*;
1142
1143 #[test]
1144 fn process_start_time_check_rejects_mismatch() {
1145 let procs = Procs::new();
1146 let pid = std::process::id();
1147 procs.refresh_pids(&[pid]);
1148 let actual = procs
1149 .start_time(pid)
1150 .expect("current process should have a start time");
1151
1152 assert_ne!(procs.start_time(pid), Some(actual.saturating_add(1)));
1153 }
1154
1155 #[test]
1156 fn test_format_bytes() {
1157 assert_eq!(format_bytes(512), "512 B");
1158 assert_eq!(format_bytes(1024), "1.0 KiB");
1159 assert_eq!(format_bytes(1536), "1.5 KiB");
1160 assert_eq!(format_bytes(50 * 1024 * 1024), "50.0 MiB");
1161 assert_eq!(format_bytes(3 * 1024 * 1024 * 1024), "3.0 GiB");
1162 assert_eq!(format_bytes(1100 * 1024 * 1024 * 1024), "1.1 TiB");
1164 }
1165
1166 #[test]
1167 fn test_format_bytes_per_sec() {
1168 assert_eq!(format_bytes_per_sec(512), "512 B/s");
1169 assert_eq!(format_bytes_per_sec(1536), "1.5 KiB/s");
1170 assert_eq!(format_bytes_per_sec(2 * 1024 * 1024), "2.0 MiB/s");
1171 }
1172}
1173
1174#[cfg(all(test, unix))]
1175mod tests {
1176 use super::*;
1177 use std::os::unix::process::CommandExt;
1178 use std::process::{Child, Command, Stdio};
1179 use std::time::{Duration, Instant};
1180
1181 struct ChildGuard(Child);
1182
1183 impl Drop for ChildGuard {
1184 fn drop(&mut self) {
1185 let pid = self.0.id() as i32;
1186 let _ = unsafe { libc::killpg(pid, libc::SIGKILL) };
1188 let _ = self.0.wait();
1189 }
1190 }
1191
1192 #[tokio::test]
1193 async fn orphan_identity_checked_group_kill_rejects_mismatch() {
1194 let mut command = Command::new("sleep");
1195 command
1196 .arg("30")
1197 .stdin(Stdio::null())
1198 .stdout(Stdio::null())
1199 .stderr(Stdio::null());
1200 unsafe {
1201 command.pre_exec(|| {
1202 if libc::setsid() == -1 {
1203 return Err(std::io::Error::last_os_error());
1204 }
1205 Ok(())
1206 });
1207 }
1208
1209 let child = command.spawn().expect("failed to spawn test process");
1210 let pid = child.id();
1211 let _child = ChildGuard(child);
1212
1213 PROCS.refresh_pids(&[pid]);
1214 let actual_start_time = PROCS
1215 .start_time(pid)
1216 .expect("test process should have a start time");
1217
1218 let killed = PROCS
1219 .kill_process_group_if_start_time_matches_async(
1220 pid,
1221 actual_start_time.saturating_add(1),
1222 libc::SIGTERM,
1223 Some(Duration::from_millis(100)),
1224 )
1225 .await
1226 .expect("identity-checked kill should not error");
1227
1228 assert!(!killed);
1229 assert!(PROCS.is_running(pid), "mismatched process must survive");
1230 }
1231
1232 #[cfg(not(target_os = "linux"))]
1233 #[tokio::test]
1234 async fn orphan_identity_checked_group_kill_fails_closed_without_pidfd() {
1235 let mut command = Command::new("sleep");
1236 command
1237 .arg("30")
1238 .stdin(Stdio::null())
1239 .stdout(Stdio::null())
1240 .stderr(Stdio::null());
1241 unsafe {
1242 command.pre_exec(|| {
1243 if libc::setsid() == -1 {
1244 return Err(std::io::Error::last_os_error());
1245 }
1246 Ok(())
1247 });
1248 }
1249
1250 let child = command.spawn().expect("failed to spawn test process");
1251 let pid = child.id();
1252 let _child = ChildGuard(child);
1253
1254 PROCS.refresh_pids(&[pid]);
1255 let actual_start_time = PROCS
1256 .start_time(pid)
1257 .expect("test process should have a start time");
1258
1259 let killed = PROCS
1260 .kill_process_group_if_start_time_matches_async(
1261 pid,
1262 actual_start_time,
1263 libc::SIGTERM,
1264 Some(Duration::from_millis(100)),
1265 )
1266 .await
1267 .expect("identity-checked kill should not error");
1268
1269 assert!(!killed);
1270 assert!(
1271 PROCS.is_running(pid),
1272 "process must survive when identity cannot be pinned"
1273 );
1274 }
1275
1276 #[test]
1277 fn get_stats_includes_descendant_rss() {
1278 let mut command = Command::new("sh");
1279 command
1280 .args(["-c", "sleep 30 & wait"])
1281 .stdin(Stdio::null())
1282 .stdout(Stdio::null())
1283 .stderr(Stdio::null());
1284 unsafe {
1285 command.pre_exec(|| {
1286 if libc::setsid() == -1 {
1287 return Err(std::io::Error::last_os_error());
1288 }
1289 Ok(())
1290 });
1291 }
1292
1293 let parent = command.spawn().expect("failed to spawn process tree");
1294 let parent_pid = parent.id();
1295 let _parent = ChildGuard(parent);
1296
1297 let procs = Procs::new();
1298 let deadline = Instant::now() + Duration::from_secs(5);
1299 let mut child_pids = Vec::new();
1300 while Instant::now() < deadline {
1301 procs.refresh_processes();
1302 child_pids = procs.all_children(parent_pid);
1303 if !child_pids.is_empty() {
1304 break;
1305 }
1306 std::thread::sleep(Duration::from_millis(50));
1307 }
1308 assert!(
1309 !child_pids.is_empty(),
1310 "test process tree did not appear under parent pid {parent_pid}"
1311 );
1312
1313 procs.refresh_processes();
1314 child_pids = procs.all_children(parent_pid);
1315 assert!(
1316 !child_pids.is_empty(),
1317 "test process tree disappeared under parent pid {parent_pid}"
1318 );
1319 let root_pid = sysinfo::Pid::from_u32(parent_pid);
1320 let direct_memory = {
1321 let system = procs.lock_system();
1322 system
1323 .process(root_pid)
1324 .expect("parent process should exist")
1325 .memory()
1326 };
1327 let descendant_memory = {
1328 let system = procs.lock_system();
1329 child_pids
1330 .iter()
1331 .filter_map(|pid| system.process(sysinfo::Pid::from_u32(*pid)))
1332 .map(|process| process.memory())
1333 .sum::<u64>()
1334 };
1335 assert!(
1336 descendant_memory > 0,
1337 "descendants {child_pids:?} should have nonzero RSS"
1338 );
1339
1340 let stats = procs
1341 .get_stats(parent_pid)
1342 .expect("parent process should have aggregate stats");
1343
1344 assert_eq!(
1345 stats.memory_bytes,
1346 direct_memory + descendant_memory,
1347 "get_stats should include descendant RSS for parent pid {parent_pid}; \
1348 descendants: {child_pids:?}, direct RSS: {direct_memory}, \
1349 descendant RSS: {descendant_memory}, reported RSS: {}",
1350 stats.memory_bytes
1351 );
1352 }
1353}