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 std::time::{Duration, Instant};
13use sysinfo::ProcessesToUpdate;
14#[cfg(windows)]
15use windows_sys::Win32::Foundation::{CloseHandle, FILETIME, HANDLE};
16#[cfg(windows)]
17use windows_sys::Win32::System::Threading::{
18 GetProcessTimes, OpenProcess, PROCESS_QUERY_LIMITED_INFORMATION,
19};
20
21type ParentToChildren = HashMap<u32, Vec<u32>>;
23
24type ProcessNames = HashMap<u32, (String, Option<String>)>;
26
27const FULL_REFRESH_INTERVAL: Duration = Duration::from_secs(5);
33
34pub struct Procs {
35 system: Mutex<sysinfo::System>,
36 last_full_refresh: Mutex<Option<Instant>>,
39}
40
41pub static PROCS: Lazy<Procs> = Lazy::new(Procs::new);
42
43impl Default for Procs {
44 fn default() -> Self {
45 Self::new()
46 }
47}
48
49impl Procs {
50 pub fn new() -> Self {
51 Self {
64 system: Mutex::new(sysinfo::System::new()),
65 last_full_refresh: Mutex::new(None),
66 }
67 }
68
69 fn lock_system(&self) -> std::sync::MutexGuard<'_, sysinfo::System> {
70 self.system.lock().unwrap_or_else(|poisoned| {
71 warn!("System mutex was poisoned, recovering");
72 poisoned.into_inner()
73 })
74 }
75
76 pub fn title(&self, pid: u32) -> Option<String> {
77 self.lock_system()
78 .process(sysinfo::Pid::from_u32(pid))
79 .map(|p| p.name().to_string_lossy().to_string())
80 }
81
82 pub fn boot_time(&self) -> u64 {
89 sysinfo::System::boot_time()
90 }
91
92 pub fn start_time(&self, pid: u32) -> Option<u64> {
98 process_start_token(pid)
99 }
100
101 #[cfg(any(unix, windows))]
102 fn start_time_matches(&self, pid: u32, expected: u64) -> bool {
103 self.start_time(pid) == Some(expected)
104 }
105
106 pub fn is_running(&self, pid: u32) -> bool {
107 #[cfg(unix)]
113 {
114 unsafe {
115 if libc::kill(pid as i32, 0) == 0 {
116 return true;
117 }
118 std::io::Error::last_os_error().raw_os_error() != Some(libc::ESRCH)
119 }
120 }
121 #[cfg(not(unix))]
122 {
123 self.refresh_pids(&[pid]);
124 self.lock_system()
125 .process(sysinfo::Pid::from_u32(pid))
126 .is_some()
127 }
128 }
129
130 #[allow(dead_code)]
133 pub fn all_children(&self, pid: u32) -> Vec<u32> {
134 let system = self.lock_system();
135 let all = system.processes();
136 let mut children = vec![];
137 for (child_pid, process) in all {
138 let mut process = process;
139 while let Some(parent) = process.parent() {
140 if parent == sysinfo::Pid::from_u32(pid) {
141 children.push(child_pid.as_u32());
142 break;
143 }
144 match system.process(parent) {
145 Some(p) => process = p,
146 None => break,
147 }
148 }
149 }
150 children
151 }
152 pub fn collect_process_tree_info(&self) -> (ParentToChildren, ProcessNames) {
157 let system = self.lock_system();
158 let all = system.processes();
159 let mut parent_to_children: ParentToChildren = HashMap::new();
160 let mut process_info: ProcessNames = HashMap::new();
161
162 for (pid, proc) in all {
163 let pid_u32 = pid.as_u32();
164 process_info.insert(
165 pid_u32,
166 (
167 proc.name().to_string_lossy().to_string(),
168 proc.exe().map(|e| e.to_string_lossy().to_string()),
169 ),
170 );
171
172 if let Some(ppid) = proc.parent() {
173 parent_to_children
174 .entry(ppid.as_u32())
175 .or_default()
176 .push(pid_u32);
177 }
178 }
179
180 (parent_to_children, process_info)
181 }
182 pub async fn kill_process_group_async(
183 &self,
184 pid: u32,
185 stop_signal: i32,
186 stop_timeout: Option<std::time::Duration>,
187 ) -> Result<bool> {
188 tokio::task::spawn_blocking(move || {
189 PROCS.kill_process_group(pid, stop_signal, stop_timeout, None)
190 })
191 .await
192 .into_diagnostic()?
193 }
194
195 pub async fn kill_process_group_if_start_time_matches_async(
205 &self,
206 pid: u32,
207 expected_start_time: Option<u64>,
208 stop_signal: i32,
209 stop_timeout: Option<std::time::Duration>,
210 ) -> Result<bool> {
211 let Some(expected_start_time) = expected_start_time else {
212 warn!(
213 "no recorded start time for pid {pid}; refusing to signal it (identity cannot be bound to a process generation)"
214 );
215 return Ok(false);
216 };
217 tokio::task::spawn_blocking(move || {
218 PROCS.kill_process_group(pid, stop_signal, stop_timeout, Some(expected_start_time))
219 })
220 .await
221 .into_diagnostic()?
222 }
223
224 #[cfg(unix)]
241 fn kill_process_group(
242 &self,
243 pid: u32,
244 stop_signal: i32,
245 stop_timeout: Option<std::time::Duration>,
246 expected_start_time: Option<u64>,
247 ) -> Result<bool> {
248 let pgid = pid as i32;
249 let signal_name = signal_name(stop_signal);
250
251 #[cfg(target_os = "linux")]
252 if let Some(expected) = expected_start_time {
253 return self.kill_process_group_with_pidfds(pid, expected, stop_signal, stop_timeout);
254 }
255
256 #[cfg(not(target_os = "linux"))]
268 if let Some(expected) = expected_start_time {
269 if !self.start_time_matches(pid, expected) {
270 debug!("process {pid} identity changed before killpg; refusing to signal it");
271 return Ok(false);
272 }
273 }
274
275 debug!("killing process group {pgid} with {signal_name}");
276
277 let ret = unsafe { libc::killpg(pgid, stop_signal) };
282 if ret == -1 {
283 let err = std::io::Error::last_os_error();
284 if err.raw_os_error() == Some(libc::ESRCH) {
285 debug!("process group {pgid} no longer exists");
286 return Ok(false);
287 }
288 if err.raw_os_error() == Some(libc::EPERM) {
289 return Err(miette::miette!(
290 "failed to send {signal_name} to process group {pgid}: permission denied"
291 ));
292 }
293 warn!("failed to send {signal_name} to process group {pgid}: {err}");
294 }
295
296 let stop_timeout = stop_timeout.unwrap_or_else(|| settings().supervisor_stop_timeout());
306 let fast_ms = 10u64;
307 let slow_ms = 50u64;
308 let total_ms = stop_timeout.as_millis().max(1) as u64;
309 let fast_count = ((total_ms / fast_ms) as usize).min(10);
310 let fast_total_ms = fast_ms * fast_count as u64;
311 let remaining_ms = total_ms.saturating_sub(fast_total_ms);
312 let slow_count = (remaining_ms / slow_ms) as usize;
313
314 let fast_checks =
315 std::iter::repeat_n(std::time::Duration::from_millis(fast_ms), fast_count);
316 let slow_checks =
317 std::iter::repeat_n(std::time::Duration::from_millis(slow_ms), slow_count);
318 let mut elapsed_ms = 0u64;
319
320 for sleep_duration in fast_checks.chain(slow_checks) {
321 std::thread::sleep(sleep_duration);
322 elapsed_ms += sleep_duration.as_millis() as u64;
323 if process_group_terminated(pgid) {
324 debug!("process group {pgid} terminated after {signal_name} ({elapsed_ms} ms)",);
325 return Ok(true);
326 }
327 }
328
329 warn!(
331 "process group {pgid} did not respond to {signal_name} after {}ms, sending SIGKILL",
332 stop_timeout.as_millis()
333 );
334 let ret = unsafe { libc::killpg(pgid, libc::SIGKILL) };
335 if ret == -1 {
336 let err = std::io::Error::last_os_error();
337 if err.raw_os_error() != Some(libc::ESRCH) {
338 warn!("failed to send SIGKILL to process group {pgid}: {err}");
339 }
340 }
341
342 for _ in 0..40 {
346 std::thread::sleep(std::time::Duration::from_millis(50));
347 if process_group_terminated(pgid) {
348 return Ok(true);
349 }
350 }
351 Err(miette::miette!(
354 "process group {pgid} still has members after SIGKILL \
355 (possibly stuck in uninterruptible sleep)"
356 ))
357 }
358
359 pub fn process_group_alive(&self, pid: u32) -> bool {
362 #[cfg(unix)]
363 {
364 !process_group_terminated(pid as i32)
365 }
366 #[cfg(not(unix))]
367 {
368 self.is_running(pid)
369 }
370 }
371
372 #[cfg(target_os = "linux")]
373 fn kill_process_group_with_pidfds(
374 &self,
375 pid: u32,
376 expected_start_time: u64,
377 _stop_signal: i32,
378 stop_timeout: Option<std::time::Duration>,
379 ) -> Result<bool> {
380 let leader = match open_pidfd(pid) {
381 Ok(pidfd) => pidfd,
382 Err(err) => {
383 warn!("cannot securely identify process group {pid}: {err}");
384 return Ok(false);
385 }
386 };
387 if !self.start_time_matches(pid, expected_start_time) {
388 debug!("process group {pid} leader identity changed before signaling");
389 return Ok(false);
390 }
391
392 let mut members = vec![(pid, leader)];
393 if let Err(err) = stop_pidfds(&members) {
394 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
395 return Err(err);
396 }
397 if !pidfd_is_running(&members[0].1) {
398 debug!("process group {pid} leader exited before it could be frozen");
399 return Ok(false);
400 }
401
402 loop {
406 let known_members = members.len();
407 let added = match extend_process_group_pidfds(pid as i32, &mut members) {
408 Ok(added) => added,
409 Err(err) => {
410 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
411 return Err(miette::miette!(
412 "failed to scan pinned process group {pid}: {err}"
413 ));
414 }
415 };
416 if added == 0 {
417 break;
418 }
419 if let Err(err) = stop_pidfds(&members[known_members..]) {
420 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
421 return Err(err);
422 }
423 }
424
425 warn!(
426 "force-terminating {} pinned orphan process(es) in group {pid}",
427 members.len()
428 );
429 if let Err(err) = signal_pidfds(&members, libc::SIGKILL, "SIGKILL") {
430 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
431 return Err(err);
432 }
433
434 let exit_timeout = stop_timeout.unwrap_or_else(|| settings().supervisor_stop_timeout());
435 let checks = exit_timeout.as_millis().max(1).div_ceil(50) as usize;
436 for _ in 0..checks {
437 if members.iter().all(|(_, pidfd)| !pidfd_is_running(pidfd)) {
438 return Ok(true);
439 }
440 std::thread::sleep(std::time::Duration::from_millis(50));
441 }
442 if members.iter().all(|(_, pidfd)| !pidfd_is_running(pidfd)) {
443 return Ok(true);
444 }
445
446 warn!("one or more pinned processes in orphan group {pid} remained alive after SIGKILL");
447 let _ = signal_pidfds(&members, libc::SIGCONT, "SIGCONT");
448 Ok(false)
449 }
450
451 #[cfg(not(unix))]
452 fn kill_process_group(
453 &self,
454 pid: u32,
455 _stop_signal: i32,
456 _stop_timeout: Option<std::time::Duration>,
457 expected_start_time: Option<u64>,
458 ) -> Result<bool> {
459 #[cfg(windows)]
462 let _identity_handle = if let Some(expected) = expected_start_time {
463 let handle = match open_process_handle(pid) {
464 Ok(handle) => handle,
465 Err(err) => {
466 warn!("cannot securely identify process {pid}: {err}");
467 return Ok(false);
468 }
469 };
470 if process_start_token_from_handle(handle.0) != Some(expected) {
471 debug!("process {pid} identity changed before taskkill");
472 return Ok(false);
473 }
474 Some(handle)
475 } else {
476 None
477 };
478
479 #[cfg(not(windows))]
480 if let Some(expected) = expected_start_time
481 && !self.start_time_matches(pid, expected)
482 {
483 debug!("process {pid} identity changed before termination");
484 return Ok(false);
485 }
486
487 self.kill(pid, 0, None)
488 }
489
490 pub async fn kill_async(
491 &self,
492 pid: u32,
493 stop_signal: i32,
494 stop_timeout: Option<std::time::Duration>,
495 ) -> Result<bool> {
496 tokio::task::spawn_blocking(move || PROCS.kill(pid, stop_signal, stop_timeout))
497 .await
498 .into_diagnostic()?
499 }
500
501 fn kill(
511 &self,
512 pid: u32,
513 stop_signal: i32,
514 stop_timeout: Option<std::time::Duration>,
515 ) -> Result<bool> {
516 let sysinfo_pid = sysinfo::Pid::from_u32(pid);
517
518 debug!("killing process {pid}");
519
520 #[cfg(windows)]
521 {
522 let _ = (stop_signal, stop_timeout);
523 let output = std::process::Command::new("taskkill")
528 .args(["/F", "/T", "/PID"])
529 .arg(pid.to_string())
530 .creation_flags(0x08000000) .output();
532 let taskkill_succeeded = match output {
533 Ok(o) if o.status.success() => {
534 debug!("taskkill /F /T /PID {pid} succeeded");
535 true
536 }
537 Ok(o) => {
538 debug!(
539 "taskkill /F /T /PID {pid} exited with status {}: {}",
540 o.status,
541 String::from_utf8_lossy(&o.stderr).trim()
542 );
543 false
544 }
545 Err(e) => {
546 debug!("failed to spawn taskkill for pid {pid}: {e}");
547 false
548 }
549 };
550 std::thread::sleep(std::time::Duration::from_millis(200));
554 if !taskkill_succeeded && self.is_running(pid) {
555 return Err(miette::miette!(
556 "taskkill failed and process {pid} is still running"
557 ));
558 }
559 Ok(true)
560 }
561
562 #[cfg(unix)]
563 {
564 let signal_name = signal_name(stop_signal);
565 debug!("sending {signal_name} to process {pid}");
569 let ret = unsafe { libc::kill(pid as i32, stop_signal) };
570 if ret == -1 {
571 let err = std::io::Error::last_os_error();
572 if err.raw_os_error() == Some(libc::ESRCH) {
573 debug!("process {pid} no longer exists");
574 return Ok(false);
575 }
576 if err.raw_os_error() == Some(libc::EPERM) {
577 return Err(miette::miette!(
578 "failed to send {signal_name} to process {pid}: permission denied"
579 ));
580 }
581 return Err(miette::miette!(
582 "failed to send {signal_name} to process {pid}: {err}"
583 ));
584 }
585
586 let stop_timeout = stop_timeout.unwrap_or_else(|| settings().supervisor_stop_timeout());
589 let fast_ms = 10u64;
590 let slow_ms = 50u64;
591 let total_ms = stop_timeout.as_millis().max(1) as u64;
592 let fast_count = ((total_ms / fast_ms) as usize).min(10);
593 let fast_total_ms = fast_ms * fast_count as u64;
594 let remaining_ms = total_ms.saturating_sub(fast_total_ms);
595 let slow_count = (remaining_ms / slow_ms) as usize;
596
597 for i in 0..fast_count {
598 std::thread::sleep(std::time::Duration::from_millis(fast_ms));
599 self.refresh_pids(&[pid]);
600 if self.is_terminated_or_zombie(sysinfo_pid) {
601 debug!(
602 "process {pid} terminated after {signal_name} ({} ms)",
603 (i + 1) * fast_ms as usize
604 );
605 return Ok(true);
606 }
607 }
608
609 for i in 0..slow_count {
611 std::thread::sleep(std::time::Duration::from_millis(slow_ms));
612 self.refresh_pids(&[pid]);
613 if self.is_terminated_or_zombie(sysinfo_pid) {
614 debug!(
615 "process {pid} terminated after {signal_name} ({} ms)",
616 fast_total_ms + (i + 1) as u64 * slow_ms
617 );
618 return Ok(true);
619 }
620 }
621
622 warn!(
624 "process {pid} did not respond to {signal_name} after {}ms, sending SIGKILL",
625 stop_timeout.as_millis()
626 );
627 let ret = unsafe { libc::kill(pid as i32, libc::SIGKILL) };
628 if ret == -1 {
629 let err = std::io::Error::last_os_error();
630 if err.raw_os_error() != Some(libc::ESRCH) {
631 warn!("failed to send SIGKILL to process {pid}: {err}");
632 }
633 }
634
635 std::thread::sleep(std::time::Duration::from_millis(100));
637 Ok(true)
638 }
639 }
640
641 #[cfg(unix)]
645 fn is_terminated_or_zombie(&self, sysinfo_pid: sysinfo::Pid) -> bool {
646 let system = self.lock_system();
647 match system.process(sysinfo_pid) {
648 None => true,
649 Some(process) => {
650 matches!(process.status(), sysinfo::ProcessStatus::Zombie)
651 }
652 }
653 }
654
655 pub(crate) fn refresh_processes(&self) {
656 let mut system = self.lock_system();
657 system.refresh_processes(ProcessesToUpdate::All, true);
658 #[cfg(windows)]
663 system.refresh_cpu_usage();
664 }
665
666 pub(crate) fn refresh_pids(&self, pids: &[u32]) {
669 let sysinfo_pids: Vec<sysinfo::Pid> =
670 pids.iter().map(|p| sysinfo::Pid::from_u32(*p)).collect();
671 self.lock_system()
672 .refresh_processes(ProcessesToUpdate::Some(&sysinfo_pids), true);
673 }
674
675 pub(crate) fn refresh_if_stale(&self) {
692 let mut last = self
693 .last_full_refresh
694 .lock()
695 .unwrap_or_else(|poisoned| poisoned.into_inner());
696 if last.is_none_or(|t| t.elapsed() >= FULL_REFRESH_INTERVAL) {
697 self.refresh_processes();
698 *last = Some(Instant::now());
699 }
700 }
701
702 pub fn get_batch_group_stats(&self, pids: &[u32]) -> Vec<(u32, Option<ProcessStats>)> {
708 if pids.is_empty() {
709 return Vec::new();
710 }
711
712 let system = self.lock_system();
713 let processes = system.processes();
714
715 let now = std::time::SystemTime::now()
716 .duration_since(std::time::UNIX_EPOCH)
717 .map(|d| d.as_secs())
718 .unwrap_or(0);
719
720 let mut children_map: std::collections::HashMap<sysinfo::Pid, Vec<sysinfo::Pid>> =
722 std::collections::HashMap::new();
723 for (child_pid, child) in processes {
724 if child.thread_kind().is_some() {
727 continue;
728 }
729 if let Some(ppid) = child.parent() {
730 children_map.entry(ppid).or_default().push(*child_pid);
731 }
732 }
733
734 pids.iter()
735 .map(|&pid| {
736 let root_pid = sysinfo::Pid::from_u32(pid);
737 let Some(root) = processes.get(&root_pid) else {
738 return (pid, None);
739 };
740
741 let root_disk = root.disk_usage();
742 let mut stats = ProcessStats {
743 cpu_percent: root.cpu_usage(),
744 memory_bytes: root.memory(),
745 uptime_secs: now.saturating_sub(root.start_time()),
746 disk_read_bytes: root_disk.read_bytes,
747 disk_write_bytes: root_disk.written_bytes,
748 };
749
750 let mut queue = std::collections::VecDeque::new();
752 if let Some(direct_children) = children_map.get(&root_pid) {
753 queue.extend(direct_children);
754 }
755 while let Some(child_pid) = queue.pop_front() {
756 if let Some(child) = processes.get(&child_pid) {
757 let disk = child.disk_usage();
758 stats.cpu_percent += child.cpu_usage();
759 stats.memory_bytes += child.memory();
760 stats.disk_read_bytes += disk.read_bytes;
761 stats.disk_write_bytes += disk.written_bytes;
762 }
763 if let Some(grandchildren) = children_map.get(&child_pid) {
764 queue.extend(grandchildren);
765 }
766 }
767
768 (pid, Some(stats))
769 })
770 .collect()
771 }
772 pub fn refresh_and_get_batch_stats(&self, pids: &[u32]) -> HashMap<u32, ProcessStats> {
782 self.refresh_processes();
783 self.get_batch_group_stats(pids)
784 .into_iter()
785 .filter_map(|(pid, stats)| stats.map(|s| (pid, s)))
786 .collect()
787 }
788
789 pub fn refresh_and_get_batch_stats_if_stale(&self, pids: &[u32]) -> HashMap<u32, ProcessStats> {
796 self.refresh_if_stale();
797 self.get_batch_group_stats(pids)
798 .into_iter()
799 .filter_map(|(pid, stats)| stats.map(|s| (pid, s)))
800 .collect()
801 }
802
803 pub fn get_batch_tree_stats_map(&self, pids: &[u32]) -> HashMap<u32, ProcessStats> {
805 self.get_batch_group_stats(pids)
806 .into_iter()
807 .filter_map(|(pid, stats)| stats.map(|stats| (pid, stats)))
808 .collect()
809 }
810
811 pub fn get_stats(&self, pid: u32) -> Option<ProcessStats> {
813 self.get_batch_group_stats(&[pid])
814 .into_iter()
815 .next()
816 .and_then(|(_, stats)| stats)
817 }
818
819 pub fn get_extended_stats(&self, pid: u32) -> Option<ExtendedProcessStats> {
821 let system = self.lock_system();
822 let processes = system.processes();
823 let root_pid = sysinfo::Pid::from_u32(pid);
824 let p = processes.get(&root_pid)?;
825
826 let now = std::time::SystemTime::now()
827 .duration_since(std::time::UNIX_EPOCH)
828 .map(|d| d.as_secs())
829 .unwrap_or(0);
830
831 let root_disk = p.disk_usage();
832 let mut aggregate_stats = ProcessStats {
833 cpu_percent: p.cpu_usage(),
834 memory_bytes: p.memory(),
835 uptime_secs: now.saturating_sub(p.start_time()),
836 disk_read_bytes: root_disk.read_bytes,
837 disk_write_bytes: root_disk.written_bytes,
838 };
839
840 let mut children_map: HashMap<sysinfo::Pid, Vec<sysinfo::Pid>> = HashMap::new();
841 for (child_pid, child) in processes {
842 if let Some(ppid) = child.parent() {
843 children_map.entry(ppid).or_default().push(*child_pid);
844 }
845 }
846
847 let mut queue = std::collections::VecDeque::new();
848 if let Some(direct_children) = children_map.get(&root_pid) {
849 queue.extend(direct_children);
850 }
851 while let Some(child_pid) = queue.pop_front() {
852 if let Some(child) = processes.get(&child_pid) {
853 let disk = child.disk_usage();
854 aggregate_stats.cpu_percent += child.cpu_usage();
855 aggregate_stats.memory_bytes += child.memory();
856 aggregate_stats.disk_read_bytes += disk.read_bytes;
857 aggregate_stats.disk_write_bytes += disk.written_bytes;
858 }
859 if let Some(grandchildren) = children_map.get(&child_pid) {
860 queue.extend(grandchildren);
861 }
862 }
863
864 Some(ExtendedProcessStats {
865 name: p.name().to_string_lossy().to_string(),
866 status: format!("{:?}", p.status()),
867 cpu_percent: aggregate_stats.cpu_percent,
868 memory_bytes: aggregate_stats.memory_bytes,
869 virtual_memory_bytes: p.virtual_memory(),
870 uptime_secs: aggregate_stats.uptime_secs,
871 thread_count: p.tasks().map(|t| t.len()).unwrap_or(0),
872 })
873 }
874}
875
876#[cfg(target_os = "linux")]
877fn process_start_token(pid: u32) -> Option<u64> {
878 let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
879 let command_end = stat.rfind(')')?;
880 stat.get(command_end + 1..)?
882 .split_whitespace()
883 .nth(19)?
884 .parse()
885 .ok()
886}
887
888#[cfg(target_os = "macos")]
889fn process_start_token(pid: u32) -> Option<u64> {
890 let mut info = std::mem::MaybeUninit::<libc::proc_bsdinfo>::zeroed();
891 let size = std::mem::size_of::<libc::proc_bsdinfo>() as i32;
892 let read = unsafe {
893 libc::proc_pidinfo(
894 pid as i32,
895 libc::PROC_PIDTBSDINFO,
896 0,
897 info.as_mut_ptr().cast(),
898 size,
899 )
900 };
901 if read != size {
902 return None;
903 }
904 let info = unsafe { info.assume_init() };
905 info.pbi_start_tvsec
906 .checked_mul(1_000_000)?
907 .checked_add(info.pbi_start_tvusec)
908}
909
910#[cfg(windows)]
911fn process_start_token(pid: u32) -> Option<u64> {
912 let handle = open_process_handle(pid).ok()?;
913 process_start_token_from_handle(handle.0)
914}
915
916#[cfg(windows)]
917fn process_start_token_from_handle(handle: HANDLE) -> Option<u64> {
918 let mut creation = FILETIME {
919 dwLowDateTime: 0,
920 dwHighDateTime: 0,
921 };
922 let mut exit = creation;
923 let mut kernel = creation;
924 let mut user = creation;
925 let ok = unsafe { GetProcessTimes(handle, &mut creation, &mut exit, &mut kernel, &mut user) };
926 if ok == 0 {
927 return None;
928 }
929
930 Some((u64::from(creation.dwHighDateTime) << 32) | u64::from(creation.dwLowDateTime))
931}
932
933#[cfg(windows)]
934struct ProcessHandle(HANDLE);
935
936#[cfg(windows)]
937impl Drop for ProcessHandle {
938 fn drop(&mut self) {
939 unsafe {
940 CloseHandle(self.0);
941 }
942 }
943}
944
945#[cfg(windows)]
946fn open_process_handle(pid: u32) -> std::io::Result<ProcessHandle> {
947 let handle = unsafe { OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, 0, pid) };
948 if handle.is_null() {
949 return Err(std::io::Error::last_os_error());
950 }
951 Ok(ProcessHandle(handle))
952}
953
954#[cfg(not(any(target_os = "linux", target_os = "macos", windows)))]
955fn process_start_token(pid: u32) -> Option<u64> {
956 let mut system = sysinfo::System::new();
957 let sysinfo_pid = sysinfo::Pid::from_u32(pid);
958 system.refresh_processes(ProcessesToUpdate::Some(&[sysinfo_pid]), true);
959 system
960 .process(sysinfo_pid)
961 .map(|process| process.start_time())
962}
963
964#[cfg(target_os = "linux")]
965fn open_pidfd(pid: u32) -> std::io::Result<OwnedFd> {
966 let fd = unsafe { libc::syscall(libc::SYS_pidfd_open, pid, 0) };
967 if fd < 0 {
968 return Err(std::io::Error::last_os_error());
969 }
970 Ok(unsafe { OwnedFd::from_raw_fd(fd as i32) })
971}
972
973#[cfg(target_os = "linux")]
974fn pidfd_is_running(pidfd: &OwnedFd) -> bool {
975 match try_pidfd_is_running(pidfd) {
976 Ok(running) => running,
977 Err(err) => {
978 warn!("failed to poll pidfd {}: {err}", pidfd.as_raw_fd());
979 true
980 }
981 }
982}
983
984#[cfg(target_os = "linux")]
985fn try_pidfd_is_running(pidfd: &OwnedFd) -> std::io::Result<bool> {
986 let mut pollfd = libc::pollfd {
987 fd: pidfd.as_raw_fd(),
988 events: libc::POLLIN,
989 revents: 0,
990 };
991 let result = unsafe { libc::poll(&mut pollfd, 1, 0) };
992 if result < 0 {
993 return Err(std::io::Error::last_os_error());
994 }
995 Ok(result == 0)
996}
997
998#[cfg(target_os = "linux")]
999fn signal_pidfds(members: &[(u32, OwnedFd)], signal: i32, signal_name: &str) -> Result<()> {
1000 for (pid, pidfd) in members {
1001 if !pidfd_is_running(pidfd) {
1002 continue;
1003 }
1004 let result = unsafe {
1005 libc::syscall(
1006 libc::SYS_pidfd_send_signal,
1007 pidfd.as_raw_fd(),
1008 signal,
1009 std::ptr::null::<libc::siginfo_t>(),
1010 0,
1011 )
1012 };
1013 if result == -1 {
1014 let err = std::io::Error::last_os_error();
1015 if err.raw_os_error() == Some(libc::ESRCH) {
1016 continue;
1017 }
1018 return Err(miette::miette!(
1019 "failed to send {signal_name} to pinned process {pid}: {err}"
1020 ));
1021 }
1022 }
1023 Ok(())
1024}
1025
1026#[cfg(target_os = "linux")]
1027fn stop_pidfds(members: &[(u32, OwnedFd)]) -> Result<()> {
1028 signal_pidfds(members, libc::SIGSTOP, "SIGSTOP")?;
1029 for _ in 0..200 {
1030 if members.iter().all(|(pid, pidfd)| {
1031 !pidfd_is_running(pidfd) || matches!(linux_process_state(*pid), Some('T' | 't'))
1032 }) {
1033 return Ok(());
1034 }
1035 std::thread::sleep(std::time::Duration::from_millis(5));
1036 }
1037 Err(miette::miette!(
1038 "timed out while freezing orphan process group"
1039 ))
1040}
1041
1042#[cfg(target_os = "linux")]
1043fn extend_process_group_pidfds(
1044 pgid: i32,
1045 members: &mut Vec<(u32, OwnedFd)>,
1046) -> std::io::Result<usize> {
1047 let entries = std::fs::read_dir("/proc")?;
1048 let mut added = 0;
1049 for entry in entries {
1050 let entry = entry?;
1051 let Some(pid) = entry
1052 .file_name()
1053 .to_str()
1054 .and_then(|name| name.parse::<u32>().ok())
1055 else {
1056 continue;
1057 };
1058 let Some(observed_identity) = linux_process_identity(pid) else {
1059 continue;
1060 };
1061 if observed_identity.0 != pgid {
1062 continue;
1063 }
1064 let mut already_pinned = false;
1065 for (known_pid, pidfd) in members.iter() {
1066 if *known_pid == pid && try_pidfd_is_running(pidfd)? {
1067 already_pinned = true;
1068 break;
1069 }
1070 }
1071 if already_pinned {
1072 continue;
1073 }
1074
1075 let pidfd = match open_pidfd(pid) {
1076 Ok(pidfd) => pidfd,
1077 Err(err) if err.raw_os_error() == Some(libc::ESRCH) => continue,
1078 Err(err) => return Err(err),
1079 };
1080 if linux_process_identity(pid) != Some(observed_identity) {
1081 return Err(std::io::Error::other(format!(
1082 "process {pid} identity changed while pinning group {pgid}"
1083 )));
1084 }
1085 members.push((pid, pidfd));
1086 added += 1;
1087 }
1088 Ok(added)
1089}
1090
1091#[cfg(target_os = "linux")]
1092fn linux_process_identity(pid: u32) -> Option<(i32, u64)> {
1093 let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
1094 let command_end = stat.rfind(')')?;
1095 let fields: Vec<_> = stat.get(command_end + 1..)?.split_whitespace().collect();
1096 Some((fields.get(2)?.parse().ok()?, fields.get(19)?.parse().ok()?))
1099}
1100
1101#[cfg(target_os = "linux")]
1102fn linux_process_state(pid: u32) -> Option<char> {
1103 let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?;
1104 let command_end = stat.rfind(')')?;
1105 stat.get(command_end + 1..)?
1106 .split_whitespace()
1107 .next()?
1108 .chars()
1109 .next()
1110}
1111
1112#[derive(Debug, Clone, Copy)]
1113pub struct ProcessStats {
1114 pub cpu_percent: f32,
1115 pub memory_bytes: u64,
1116 pub uptime_secs: u64,
1117 pub disk_read_bytes: u64,
1118 pub disk_write_bytes: u64,
1119}
1120
1121impl ProcessStats {
1122 pub fn memory_display(&self) -> String {
1123 format_bytes(self.memory_bytes)
1124 }
1125
1126 pub fn cpu_display(&self) -> String {
1127 format!("{:.1}%", self.cpu_percent)
1128 }
1129
1130 pub fn uptime_display(&self) -> String {
1131 format_duration(self.uptime_secs)
1132 }
1133
1134 pub fn disk_read_display(&self) -> String {
1135 format_bytes_per_sec(self.disk_read_bytes)
1136 }
1137
1138 pub fn disk_write_display(&self) -> String {
1139 format_bytes_per_sec(self.disk_write_bytes)
1140 }
1141}
1142
1143#[derive(Debug, Clone)]
1144pub struct ExtendedProcessStats {
1145 pub name: String,
1146 pub status: String,
1147 pub cpu_percent: f32,
1148 pub memory_bytes: u64,
1149 pub virtual_memory_bytes: u64,
1150 pub uptime_secs: u64,
1151 pub thread_count: usize,
1152}
1153
1154fn format_bytes(bytes: u64) -> String {
1155 humanbyte::to_string(bytes, humanbyte::Format::IEC)
1156}
1157
1158fn format_duration(secs: u64) -> String {
1159 if secs < 60 {
1160 format!("{secs}s")
1161 } else if secs < 3600 {
1162 format!("{}m {}s", secs / 60, secs % 60)
1163 } else if secs < 86400 {
1164 let hours = secs / 3600;
1165 let mins = (secs % 3600) / 60;
1166 format!("{hours}h {mins}m")
1167 } else {
1168 let days = secs / 86400;
1169 let hours = (secs % 86400) / 3600;
1170 format!("{days}d {hours}h")
1171 }
1172}
1173
1174fn format_bytes_per_sec(bytes: u64) -> String {
1175 format!("{}/s", humanbyte::to_string(bytes, humanbyte::Format::IEC))
1176}
1177
1178#[cfg(unix)]
1194fn process_group_terminated(pgid: i32) -> bool {
1195 unsafe { libc::killpg(pgid, 0) != 0 }
1196}
1197
1198#[cfg(unix)]
1199fn signal_name(sig: i32) -> &'static str {
1200 match sig {
1201 libc::SIGHUP => "SIGHUP",
1202 libc::SIGINT => "SIGINT",
1203 libc::SIGQUIT => "SIGQUIT",
1204 libc::SIGTERM => "SIGTERM",
1205 libc::SIGUSR1 => "SIGUSR1",
1206 libc::SIGUSR2 => "SIGUSR2",
1207 libc::SIGKILL => "SIGKILL",
1208 _ => "UNKNOWN",
1209 }
1210}
1211
1212#[cfg(test)]
1213mod format_tests {
1214 use super::*;
1215
1216 #[test]
1217 fn process_start_time_check_rejects_mismatch() {
1218 let procs = Procs::new();
1219 let pid = std::process::id();
1220 procs.refresh_pids(&[pid]);
1221 let actual = procs
1222 .start_time(pid)
1223 .expect("current process should have a start time");
1224
1225 assert_ne!(procs.start_time(pid), Some(actual.saturating_add(1)));
1226 }
1227
1228 #[test]
1229 fn test_format_bytes() {
1230 assert_eq!(format_bytes(512), "512 B");
1231 assert_eq!(format_bytes(1024), "1.0 KiB");
1232 assert_eq!(format_bytes(1536), "1.5 KiB");
1233 assert_eq!(format_bytes(50 * 1024 * 1024), "50.0 MiB");
1234 assert_eq!(format_bytes(3 * 1024 * 1024 * 1024), "3.0 GiB");
1235 assert_eq!(format_bytes(1100 * 1024 * 1024 * 1024), "1.1 TiB");
1237 }
1238
1239 #[test]
1240 fn test_format_bytes_per_sec() {
1241 assert_eq!(format_bytes_per_sec(512), "512 B/s");
1242 assert_eq!(format_bytes_per_sec(1536), "1.5 KiB/s");
1243 assert_eq!(format_bytes_per_sec(2 * 1024 * 1024), "2.0 MiB/s");
1244 }
1245}
1246
1247#[cfg(all(test, unix))]
1248mod tests {
1249 use super::*;
1250 use std::os::unix::process::CommandExt;
1251 use std::process::{Child, Command, Stdio};
1252 use std::time::{Duration, Instant};
1253
1254 struct ChildGuard(Child);
1255
1256 impl Drop for ChildGuard {
1257 fn drop(&mut self) {
1258 let pid = self.0.id() as i32;
1259 let _ = unsafe { libc::killpg(pid, libc::SIGKILL) };
1261 let _ = self.0.wait();
1262 }
1263 }
1264
1265 #[tokio::test]
1266 async fn orphan_identity_checked_group_kill_rejects_mismatch() {
1267 let mut command = Command::new("sleep");
1268 command
1269 .arg("30")
1270 .stdin(Stdio::null())
1271 .stdout(Stdio::null())
1272 .stderr(Stdio::null());
1273 unsafe {
1274 command.pre_exec(|| {
1275 if libc::setsid() == -1 {
1276 return Err(std::io::Error::last_os_error());
1277 }
1278 Ok(())
1279 });
1280 }
1281
1282 let child = command.spawn().expect("failed to spawn test process");
1283 let pid = child.id();
1284 let _child = ChildGuard(child);
1285
1286 PROCS.refresh_pids(&[pid]);
1287 let actual_start_time = PROCS
1288 .start_time(pid)
1289 .expect("test process should have a start time");
1290
1291 let killed = PROCS
1292 .kill_process_group_if_start_time_matches_async(
1293 pid,
1294 Some(actual_start_time.saturating_add(1)),
1295 libc::SIGTERM,
1296 Some(Duration::from_millis(100)),
1297 )
1298 .await
1299 .expect("identity-checked kill should not error");
1300
1301 assert!(!killed);
1302 assert!(PROCS.is_running(pid), "mismatched process must survive");
1303 }
1304
1305 #[cfg(all(unix, not(target_os = "linux")))]
1310 #[tokio::test]
1311 async fn identity_checked_group_kill_reverifies_inside_blocking_op() {
1312 let mut command = Command::new("sleep");
1313 command
1314 .arg("30")
1315 .stdin(Stdio::null())
1316 .stdout(Stdio::null())
1317 .stderr(Stdio::null());
1318 unsafe {
1319 command.pre_exec(|| {
1320 if libc::setsid() == -1 {
1321 return Err(std::io::Error::last_os_error());
1322 }
1323 Ok(())
1324 });
1325 }
1326
1327 let child = command.spawn().expect("failed to spawn test process");
1328 let pid = child.id();
1329 let _child = ChildGuard(child);
1330
1331 PROCS.refresh_pids(&[pid]);
1332 let actual_start_time = PROCS
1333 .start_time(pid)
1334 .expect("test process should have a start time");
1335
1336 let killed = PROCS
1337 .kill_process_group_if_start_time_matches_async(
1338 pid,
1339 Some(actual_start_time),
1340 libc::SIGTERM,
1341 Some(Duration::from_millis(100)),
1342 )
1343 .await
1344 .expect("identity-checked kill should not error");
1345
1346 assert!(killed, "matching generation must be signalled");
1347 assert!(
1348 !PROCS.is_running(pid),
1349 "signalled process group must be gone"
1350 );
1351 }
1352
1353 #[test]
1354 fn get_stats_includes_descendant_rss() {
1355 let mut command = Command::new("sh");
1356 command
1357 .args(["-c", "sleep 30 & wait"])
1358 .stdin(Stdio::null())
1359 .stdout(Stdio::null())
1360 .stderr(Stdio::null());
1361 unsafe {
1362 command.pre_exec(|| {
1363 if libc::setsid() == -1 {
1364 return Err(std::io::Error::last_os_error());
1365 }
1366 Ok(())
1367 });
1368 }
1369
1370 let parent = command.spawn().expect("failed to spawn process tree");
1371 let parent_pid = parent.id();
1372 let _parent = ChildGuard(parent);
1373
1374 let procs = Procs::new();
1375 let deadline = Instant::now() + Duration::from_secs(5);
1376 let mut child_pids = Vec::new();
1377 while Instant::now() < deadline {
1378 procs.refresh_processes();
1379 child_pids = procs.all_children(parent_pid);
1380 if !child_pids.is_empty() {
1381 break;
1382 }
1383 std::thread::sleep(Duration::from_millis(50));
1384 }
1385 assert!(
1386 !child_pids.is_empty(),
1387 "test process tree did not appear under parent pid {parent_pid}"
1388 );
1389
1390 procs.refresh_processes();
1391 child_pids = procs.all_children(parent_pid);
1392 assert!(
1393 !child_pids.is_empty(),
1394 "test process tree disappeared under parent pid {parent_pid}"
1395 );
1396 let root_pid = sysinfo::Pid::from_u32(parent_pid);
1397 let direct_memory = {
1398 let system = procs.lock_system();
1399 system
1400 .process(root_pid)
1401 .expect("parent process should exist")
1402 .memory()
1403 };
1404 let descendant_memory = {
1405 let system = procs.lock_system();
1406 child_pids
1407 .iter()
1408 .filter_map(|pid| system.process(sysinfo::Pid::from_u32(*pid)))
1409 .map(|process| process.memory())
1410 .sum::<u64>()
1411 };
1412 assert!(
1413 descendant_memory > 0,
1414 "descendants {child_pids:?} should have nonzero RSS"
1415 );
1416
1417 let stats = procs
1418 .get_stats(parent_pid)
1419 .expect("parent process should have aggregate stats");
1420
1421 assert_eq!(
1422 stats.memory_bytes,
1423 direct_memory + descendant_memory,
1424 "get_stats should include descendant RSS for parent pid {parent_pid}; \
1425 descendants: {child_pids:?}, direct RSS: {direct_memory}, \
1426 descendant RSS: {descendant_memory}, reported RSS: {}",
1427 stats.memory_bytes
1428 );
1429 }
1430
1431 #[test]
1432 fn full_refresh_is_throttled_by_ttl() {
1433 let procs = Procs::new();
1434
1435 assert!(
1437 procs.last_full_refresh.lock().unwrap().is_none(),
1438 "fresh Procs should have no recorded refresh"
1439 );
1440 procs.refresh_if_stale();
1441 let first = procs.last_full_refresh.lock().unwrap().unwrap();
1442 assert!(first.elapsed() < FULL_REFRESH_INTERVAL);
1443
1444 procs.refresh_if_stale();
1447 let second = procs.last_full_refresh.lock().unwrap().unwrap();
1448 assert_eq!(
1449 second, first,
1450 "refresh_if_stale within TTL must skip the refresh and keep the timestamp"
1451 );
1452
1453 let expired = Instant::now()
1456 .checked_sub(FULL_REFRESH_INTERVAL + Duration::from_secs(1))
1457 .expect("system has been up long enough to backdate by 6s");
1458 *procs.last_full_refresh.lock().unwrap() = Some(expired);
1459 procs.refresh_if_stale();
1460 let third = procs.last_full_refresh.lock().unwrap().unwrap();
1461 assert!(
1462 third > expired,
1463 "refresh_if_stale after expired TTL must refresh and advance the timestamp"
1464 );
1465 assert!(
1466 third.elapsed() < FULL_REFRESH_INTERVAL,
1467 "fresh timestamp after expired-TTL refresh should be recent"
1468 );
1469 }
1470}