Skip to main content

pitchfork_cli/
procs.rs

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
20/// Map from parent PID to its child PIDs.
21type ParentToChildren = HashMap<u32, Vec<u32>>;
22
23/// Map from PID to process name and optional executable path.
24type ProcessNames = HashMap<u32, (String, Option<String>)>;
25
26/// How often a full process-table refresh may run. High-frequency callers
27/// (web API polling, TUI frames) share one full scan instead of each paying
28/// the ~40ms /proc walk on every call. Targeted `refresh_pids` calls are not
29/// throttled. sysinfo's `cpu_usage()` is a delta between refreshes, so a
30/// longer window also smooths the reported CPU%.
31const FULL_REFRESH_INTERVAL: Duration = Duration::from_secs(5);
32
33pub struct Procs {
34    system: Mutex<sysinfo::System>,
35    /// When the last full process-table refresh ran. `None` until the first
36    /// refresh forces the initial scan (the table starts empty).
37    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        // IMPORTANT: Do NOT call refresh_processes() or System::new_all() here.
51        //
52        // Both refresh the state of every process in the system, which takes
53        // ~500ms on a typical machine. Since PROCS is a Lazy static, the first
54        // access triggers this constructor — and `pitchfork cd` (which only
55        // needs to check if the supervisor PID is alive) would block for that
56        // duration on every directory change.
57        //
58        // See https://github.com/jdx/pitchfork/discussions/439
59        //
60        // Callers that need process info must call refresh_pids() (for specific
61        // PIDs) or refresh_processes() (for full-system stats) explicitly.
62        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    /// Time the system booted, in seconds since the epoch.
82    ///
83    /// Constant for the lifetime of a boot, so recording it alongside a
84    /// daemon's PID makes it possible to tell later whether that record
85    /// belongs to the current boot. Note this is *not* comparable with
86    /// `start_time`, whose units are platform-specific.
87    pub fn boot_time(&self) -> u64 {
88        sysinfo::System::boot_time()
89    }
90
91    /// High-resolution kernel start token for the process.
92    ///
93    /// Combined with the PID this forms a stable identity for the lifetime of a
94    /// process. Unlike sysinfo's seconds-since-epoch value, this preserves the
95    /// native platform resolution so same-second PID reuse cannot compare equal.
96    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        // Use kill(pid, 0) on Unix for an O(1) liveness check that does not
107        // depend on the process cache being populated. This avoids the need
108        // for a full process refresh just to check a single PID.
109        // ESRCH = process does not exist; EPERM = process exists but owned
110        // by another user (still "running" from our perspective).
111        #[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    /// Walk the /proc tree to find all descendant PIDs.
130    /// Kept for diagnostics/status display; no longer used in the kill path.
131    #[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    /// Collect minimal process tree information in a single lock.
152    ///
153    /// Returns a map of parent PID → child PIDs and a map of PID → (name, exe).
154    /// This avoids repeated mutex locking when traversing deep trees.
155    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    /// Kill a process group only while its leader still has `expected_start_time`.
195    ///
196    /// Identity is refreshed inside the blocking kill operation, immediately
197    /// before signaling. Linux holds a pidfd and Windows holds an open process
198    /// handle across termination so the validated PID cannot be recycled.
199    /// Unix platforms without a durable process handle re-verify the start
200    /// time inside the blocking operation, right before the signal — the
201    /// same check-then-signal atomicity tokio affords in one closure. Pass
202    /// `None` to fail closed (refuse to signal anything).
203    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    /// Kill a single process only if it is still the generation identified by
224    /// `expected_start_time`.
225    ///
226    /// This is the single-process counterpart of
227    /// [`Self::kill_process_group_if_start_time_matches_async`], for processes
228    /// that are not session leaders (the supervisor itself is started in the
229    /// caller's session, so signalling its group would hit the caller's shell).
230    /// A recorded PID whose process is gone may have been handed to an
231    /// unrelated process, e.g. after a reboot, so the kill refuses when the
232    /// live start time differs from the recorded one. Returns `Ok(false)` when
233    /// nothing was signalled.
234    ///
235    /// On Linux the identity check binds a pidfd to the expected generation
236    /// and every signal is sent through that pidfd, so a PID recycled after
237    /// the check can never be reached: the pidfd keeps referring to the dead
238    /// generation and the kernel answers `ESRCH`. Windows holds an open
239    /// process handle for the same effect. Other Unix platforms re-verify
240    /// immediately before signalling, leaving a window no wider than the
241    /// daemon kill path has there.
242    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        // The pidfd refers to whichever process owned the PID when it was
283        // opened. Checking the start token now proves that is the expected
284        // generation; from here on only the pidfd is signalled, never the
285        // number, so a later recycle of the PID cannot be reached.
286        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        // Same polling schedule as `kill`: fast checks first, then 50ms steps
295        // for the rest of the stop timeout, then SIGKILL.
296        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        // Brief wait for SIGKILL to take effect
330        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        // An open handle keeps the Windows process object, and with it the
343        // PID, from being recycled until the handle is closed, so the
344        // identity verified below still holds for the taskkill in `kill`.
345        #[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    /// Confirm that `pid` is still the generation `expected_start_time`
355    /// describes, immediately before it is signalled.
356    ///
357    /// Three answers are possible and callers must not conflate them:
358    /// `Ok(true)` when the token matches; `Ok(false)` when the process is
359    /// gone or a different generation now owns the PID, meaning there is
360    /// nothing of ours left to signal; and `Err` when the process is alive
361    /// but its token cannot be read. The last case is *not* proof of
362    /// termination — treating it as such would let `supervisor stop` drop
363    /// the record of a supervisor that keeps running — so it is reported as
364    /// a failure and the record is left alone.
365    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    /// Kill an entire process group with graceful shutdown strategy:
385    /// 1. Send the configured stop signal to the process group (-pgid) and
386    ///    wait up to the stop timeout for the WHOLE group to exit
387    /// 2. If any processes remain, send SIGKILL to the group and verify
388    ///
389    /// Since daemons are spawned with setsid(), the daemon PID == PGID,
390    /// so this atomically signals all descendant processes.
391    ///
392    /// The stop timeout is the graceful-shutdown budget for the entire group
393    /// (like systemd's TimeoutStopSec for a cgroup): members that are still
394    /// alive when it expires are SIGKILLed. Daemons whose teardown legitimately
395    /// takes longer (e.g. draining connections, stopping containers) should
396    /// raise `stop_signal.timeout` rather than rely on outliving the budget.
397    ///
398    /// Returns `Err` if the signal could not be sent (e.g. permission denied)
399    /// or if group members survived even SIGKILL (uninterruptible sleep).
400    #[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        // A start-time check alone cannot prevent the numeric PID/PGID from
417        // being recycled before killpg. Linux closes that race with a pidfd.
418        // Other Unix platforms have no durable process handle, so they do the
419        // next best thing: re-read the start time fresh from the kernel here,
420        // inside the blocking operation and immediately before the signal, so
421        // the check-to-signal gap is just the adjacency of two syscalls in one
422        // thread — no await, no scheduler boundary. For that gap to matter, the
423        // kernel would have to wrap the entire sequential PID space (XNU
424        // allocates monotonically from lastpid and refuses IDs still in use as
425        // a proc, pgrp, or session) and hand the exact PGID to a new group
426        // leader between the two syscalls, which is not reachable in practice.
427        #[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        // Send the stop signal to the entire process group.
438        // killpg sends to all processes in the group atomically.
439        // We intentionally skip the zombie check here because the leader may be
440        // a zombie while children in the group are still running.
441        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        // Wait for graceful shutdown: fast initial check then slower polling.
457        // Per-daemon timeout overrides the global setting.
458        //
459        // The wait must cover the ENTIRE group, not just the leader. The leader
460        // is often a thin shell (`sh -c ...`) that dies within milliseconds of
461        // the signal while its children are still shutting down gracefully
462        // (e.g. `docker compose up` waiting for its container to stop).
463        // Returning as soon as the leader dies lets a force-restart spawn a
464        // replacement that collides with the still-terminating old instance.
465        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        // SIGKILL the entire process group as last resort
490        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        // Wait for SIGKILL to take effect on the whole group (bounded: SIGKILL
503        // cannot be caught, so members disappear as soon as the kernel reaps
504        // them — anything left after this is stuck in uninterruptible sleep).
505        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        // Report failure so callers do not mark the daemon Stopped (and start
512        // a replacement) while a member is still alive.
513        Err(miette::miette!(
514            "process group {pgid} still has members after SIGKILL \
515             (possibly stuck in uninterruptible sleep)"
516        ))
517    }
518
519    /// Whether any signalable member of the daemon's process group is still
520    /// alive. The daemon PID == PGID because daemons are spawned with setsid().
521    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        // Stop newly discovered members before rescanning. Once a scan adds
563        // nothing, every process capable of forking into this group is frozen,
564        // so the pinned set is complete and the PGID cannot be recycled.
565        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        // Keep the Windows process object alive through taskkill so its
620        // numeric PID cannot be recycled after identity validation.
621        #[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    /// Kill a process with graceful shutdown strategy:
651    /// 1. Send the configured stop signal and wait up to ~3s (10ms intervals for first 100ms, then 50ms intervals)
652    /// 2. If still running, send SIGKILL to force termination
653    ///
654    /// This ensures fast-exiting processes don't wait unnecessarily,
655    /// while stubborn processes eventually get forcefully terminated.
656    ///
657    /// Returns `Err` if the signal could not be sent (e.g. permission denied
658    /// when targeting a process owned by another user/root).
659    ///
660    /// Signals by PID number, so callers must have pinned the process
661    /// identity first; on Linux the pidfd path in
662    /// `kill_if_start_time_matches` is used instead. When
663    /// `expected_start_time` is given, the start token is re-checked
664    /// immediately before the first signal and again before the SIGKILL
665    /// escalation: if the process exited and its PID was recycled in either
666    /// window, the newcomer is left alone.
667    #[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            // The caller holds an open process handle while this runs, which
680            // keeps the PID from being recycled, so no re-check is needed.
681            let _ = expected_start_time;
682            // Windows has no signals. SIGINT is the one with a counterpart,
683            // Ctrl+C, so a daemon asking for it gets that first and until its
684            // stop timeout to exit; every other signal, and a daemon that is
685            // still running after the timeout, is terminated outright with
686            // its process tree. The timeout bounds only the graceful part.
687            //
688            // Best effort: once the daemon has exited, its tree can no longer
689            // be walked, so a process it started that outlives it is left
690            // running.
691            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            // Use taskkill /F /T to kill the entire process tree.
700            // sysinfo's process.kill() only kills the main process, leaving
701            // child processes (e.g. python3 spawned by sh -c) orphaned and
702            // still holding ports. The /T flag kills all descendant processes.
703            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            // Brief sleep to let the OS signal the process handle, giving
727            // tokio's child.wait() in the monitor task a chance to detect
728            // the exit and fire on_stop/on_exit hooks.
729            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            // Without pidfds there is no way to bind a signal to a process
743            // generation, so re-verify the start token immediately before
744            // the first signal. This shrinks the check-to-signal window to
745            // the two adjacent syscalls, the tightest this platform allows.
746            if let Some(expected) = expected_start_time
747                && !self.verify_start_time_before_signal(pid, expected)?
748            {
749                return Ok(false);
750            }
751            // Send stop signal for graceful shutdown using libc::kill directly
752            // so we can distinguish EPERM (permission denied) from ESRCH
753            // (process already gone — possible in a narrow race window).
754            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            // Fast check: 10ms intervals, then slower 50ms polling for stop_timeout.
773            // Per-daemon timeout overrides the global setting.
774            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            // Slower check: 50ms intervals for the remainder of stop_timeout
796            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            // SIGKILL as last resort after stop_timeout. The PID may have been
809            // recycled while we waited: the original exited and an unrelated
810            // process took its number, which the liveness polls above cannot
811            // tell apart. Re-check the generation before escalating.
812            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            // Brief wait for SIGKILL to take effect
833            std::thread::sleep(std::time::Duration::from_millis(100));
834            Ok(true)
835        }
836    }
837
838    /// Check if a process is terminated or is a zombie.
839    /// On Linux, zombie processes still have /proc/[pid] entries but are effectively dead.
840    /// This prevents unnecessary signal escalation for processes that have already exited.
841    #[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        // On Windows, refresh_processes() does not update CPU usage.
856        // sysinfo requires a separate refresh_cpu_usage() call to compute
857        // the CPU delta between two samples. The first call stores the
858        // baseline; subsequent calls return the actual percentage.
859        #[cfg(windows)]
860        system.refresh_cpu_usage();
861    }
862
863    /// Refresh only specific PIDs instead of all processes.
864    /// More efficient when you only need to check a small set of known PIDs.
865    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    /// Full process-table refresh, throttled to at most once per
873    /// [`FULL_REFRESH_INTERVAL`] **among the callers that use this gated
874    /// path** (web API polling, TUI frames, process-tree endpoints). Callers
875    /// that bypass the gate — e.g. `refresh_processes()` directly in the
876    /// active-port probe — do not update `last_full_refresh`, so a gated
877    /// caller right after may still pay a scan. The first call always
878    /// refreshes.
879    ///
880    /// Callers that need to observe *new* descendants (e.g. the active-port
881    /// probe right after a daemon spawns) must still call `refresh_processes`
882    /// directly, since a throttled refresh may miss children forked since the
883    /// last scan.
884    ///
885    /// Lock order: this method holds `last_full_refresh` while
886    /// `refresh_processes` takes `system`. Never acquire `system` first and
887    /// then call this method, or the reverse order would deadlock.
888    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    /// Get aggregated stats for multiple process trees in a single pass.
900    ///
901    /// Builds the parent→children map once (O(N)) and then BFS-es from each
902    /// root PID (O(D_i) per daemon). Total cost is O(N + ΣD_i) instead of
903    /// O(D × N) when collecting stats for each daemon separately.
904    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        // Build parent → children map once for all daemons
918        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            // Skip Linux userland threads: they report the same memory as their parent process,
922            // so including them would cause massive double-counting.
923            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                // BFS from root_pid to find all descendants
948                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    /// Refresh the process tree, then call [`Self::get_batch_group_stats`].
970    ///
971    /// Always performs a fresh full refresh. This is the correct entry point
972    /// for callers that need a new sample rather than a cached one — notably
973    /// the resource-limit checks, where counting the same cached CPU sample
974    /// twice as two violations would silently weaken enforcement.
975    ///
976    /// Returns a PID → [`ProcessStats`] map so callers do not have to repeat
977    /// the same `filter_map`/`collect` boilerplate.
978    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    /// Like [`Self::refresh_and_get_batch_stats`], but the full refresh is
987    /// throttled to at most once per [`FULL_REFRESH_INTERVAL`]. Concurrent or
988    /// frequent callers (web API polling, TUI frames) share a single scan and
989    /// read the cached process table in between. Do not use this for resource
990    /// enforcement: a cached CPU sample must not be counted as a new
991    /// violation.
992    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    /// Get process-tree stats for multiple root PIDs, omitting roots that no longer exist.
1001    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    /// Get process-tree stats (cpu%, memory bytes, uptime secs, disk I/O) for a given root PID.
1009    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    /// Get extended process information for a given PID
1017    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    // Fields after the command start at field 3 (state); starttime is field 22.
1078    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/// Send Ctrl+C to the console of `pid` and wait up to `timeout` for it to exit.
1152///
1153/// Returns whether it exited. If not, the caller still has to terminate it,
1154/// including when no Ctrl+C was sent: the daemon shares the supervisor's
1155/// console, or `pitchfork interrupt` could not send it.
1156#[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    // Opened before the interrupt, so the wait below is for this process
1163    // even if it exits before the wait starts.
1164    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    // The stop timeout covers sending Ctrl+C as well as the exit it asks for,
1175    // so a `pitchfork interrupt` that hangs cannot hold up the forced stop.
1176    // A timeout too long for an `Instant` (the setting accepts up to
1177    // u64::MAX seconds) waits as long as Win32 allows short of forever.
1178    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    // Fields after the command start at field 3. pgrp is field 5 and the
1378    // scheduler-tick start token is field 22.
1379    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/// Check whether a process group has no remaining members that we could
1460/// still be waiting on.
1461///
1462/// `killpg(pgid, 0)` probes for members without delivering a signal:
1463/// - success: at least one member remains that we may signal — keep waiting.
1464/// - ESRCH: the group is empty.
1465/// - EPERM: members remain, but none that we may signal (e.g. a root-owned
1466///   helper spawned inside the group). Our stop signal and SIGKILL never
1467///   reached them, so waiting longer cannot make progress — treat the group
1468///   as terminated, matching the previous leader-only behavior.
1469///
1470/// Unreaped zombies still count as members, but they are reaped promptly
1471/// (the leader by the supervisor's monitoring task via `child.wait()`,
1472/// orphans by init or the supervisor's PID-1 reaper), so the group empties
1473/// as soon as every member has actually exited.
1474#[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        // rolls over past GiB instead of showing e.g. "1100.0GB"
1517        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            // The test process is started in its own session, so PID == PGID.
1541            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    /// On Unix platforms without a durable process handle the generation
1587    /// check runs inside the blocking operation, immediately before the
1588    /// signal: a matching generation is signalled, a mismatched one is
1589    /// refused at the last moment.
1590    #[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        // A terminated child remains visible to kill(pid, 0) until its parent
1629        // reaps it. Poll with a deadline so a missed signal still fails.
1630        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        // First call always refreshes: the table starts empty.
1724        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        // A second call within the TTL window must NOT refresh again: the
1733        // recorded timestamp must stay identical (a refresh would advance it).
1734        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        // Simulate an expired TTL: force the recorded timestamp into the past,
1742        // then confirm the next call refreshes and advances the timestamp.
1743        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}