Skip to main content

pitchfork_cli/supervisor/
mod.rs

1//! Supervisor module - daemon process supervisor
2//!
3//! This module is split into focused submodules:
4//! - `state`: State access layer (get/set operations)
5//! - `lifecycle`: Daemon start/stop operations
6//! - `adopt`: Re-adoption of orphaned daemons after a supervisor crash
7//! - `log_sink`: Out-of-process capture of daemon output
8//! - `autostop`: Autostop logic and boot daemon startup
9//! - `retry`: Retry logic with backoff
10//! - `watchers`: Background tasks (interval, cron, file watching)
11//! - `ipc_handlers`: IPC request dispatch
12
13mod adopt;
14mod autostop;
15mod health;
16mod hooks;
17mod idle;
18mod ipc_handlers;
19mod lifecycle;
20mod log_sink;
21#[cfg(unix)]
22mod pty;
23mod retry;
24mod state;
25mod watchers;
26
27use crate::daemon_id::DaemonId;
28use crate::daemon_status::DaemonStatus;
29use crate::deps::compute_reverse_stop_order;
30use crate::ipc::server::{IpcServer, IpcServerHandle, StartupLock};
31
32use crate::procs::PROCS;
33use crate::settings::settings;
34use crate::state_file::StateFile;
35use crate::{Result, env};
36#[cfg(unix)]
37use duct::cmd;
38#[cfg(unix)]
39use miette::IntoDiagnostic;
40use once_cell::sync::Lazy;
41use std::collections::HashMap;
42#[cfg(unix)]
43use std::collections::HashSet;
44use std::fs;
45#[cfg(unix)]
46use std::os::unix::fs::PermissionsExt;
47use std::path::PathBuf;
48use std::process::exit;
49use std::sync::atomic;
50use std::sync::atomic::{AtomicBool, AtomicU32};
51use std::time::Duration;
52#[cfg(unix)]
53use tokio::signal::unix::SignalKind;
54use tokio::sync::{Mutex, Notify};
55use tokio::task::JoinHandle;
56use tokio::{signal, time};
57
58/// Exit statuses reaped by the container-mode zombie reaper for managed daemon
59/// PIDs. On non-Linux Unix platforms where `waitid(WNOWAIT)` is unavailable,
60/// `waitpid(None, WNOHANG)` may race with Tokio's `child.wait()`. When the
61/// zombie reaper wins, the exit status is stashed here so the monitoring task
62/// in lifecycle.rs can recover it instead of treating the ECHILD as a failure.
63///
64/// On Linux this map is unused because the reaper uses `waitid` with `WNOWAIT`
65/// to peek before reaping, which avoids the race entirely.
66#[cfg(all(unix, not(target_os = "linux")))]
67pub(crate) static REAPED_STATUSES: Lazy<Mutex<HashMap<u32, i32>>> =
68    Lazy::new(|| Mutex::new(HashMap::new()));
69
70// Re-export types needed by other modules
71pub(crate) use state::UpsertDaemonOpts;
72
73pub struct Supervisor {
74    pub(crate) state_file: Mutex<StateFile>,
75    pub(crate) pending_notifications: Mutex<Vec<(log::LevelFilter, String)>>,
76    pub(crate) last_refreshed_at: Mutex<time::Instant>,
77    /// Daemons whose retry sequence a foreground `run` is already working
78    /// through, each with the flag that asks it to stop. The backoff between
79    /// its attempts leaves the record errored with no PID, which is exactly
80    /// what `check_retry` looks for, so without this the background checker
81    /// would start the next attempt itself and the foreground call would be
82    /// left reporting on a run it does not own. `stop` raises the flag, so the
83    /// sequence ends rather than starting another attempt behind the user's
84    /// back.
85    /// One flag per claim: two starts can be working through the same
86    /// daemon's retries at once, and a stop has to reach all of them.
87    pub(crate) retrying:
88        std::sync::Mutex<HashMap<DaemonId, Vec<std::sync::Arc<std::sync::atomic::AtomicBool>>>>,
89    /// How many times each daemon has been stopped. The retry checker reads
90    /// this when it decides to run an attempt and again when it is about to
91    /// start one, holding the daemon's lock: a stop in between means the
92    /// attempt it approved is one the user has since called off.
93    pub(crate) stop_epochs: std::sync::Mutex<HashMap<DaemonId, u64>>,
94    /// Map of daemon ID to scheduled autostop time
95    pub(crate) pending_autostops: Mutex<HashMap<DaemonId, time::Instant>>,
96    /// Autostop stops that have been spawned as detached tasks but have not
97    /// yet begun stopping. `cancel_pending_autostops_for_dir` flips the flag
98    /// to call off the stop when a shell re-enters the directory while the
99    /// stop task is still in flight.
100    pub(crate) in_flight_autostops:
101        Mutex<HashMap<DaemonId, std::sync::Arc<std::sync::atomic::AtomicBool>>>,
102    /// Handle for graceful IPC server shutdown
103    pub(crate) ipc_shutdown: Mutex<Option<IpcServerHandle>>,
104    /// Tracks in-flight hook tasks so shutdown can wait for them to complete
105    pub(crate) hook_tasks: Mutex<Vec<JoinHandle<()>>>,
106    /// Number of monitoring tasks that are still running (between process exit
107    /// and hook registration completion). Used by `close()` to know when it is
108    /// safe to drain `hook_tasks`.
109    pub(crate) active_monitors: AtomicU32,
110    /// Signalled by each monitoring task after it finishes registering hooks
111    /// (or decides it has nothing to register). `close()` waits on this.
112    pub(crate) monitor_done: Notify,
113    /// Cancellation token for the proxy server — cancelled on shutdown to
114    /// stop accepting new connections and drain in-flight ones.
115    pub(crate) proxy_cancel: Mutex<Option<tokio_util::sync::CancellationToken>>,
116    /// Join handle for the proxy task so shutdown can wait for cleanup.
117    pub(crate) proxy_task: Mutex<Option<JoinHandle<()>>>,
118    /// Join handle for the loopback DNS resolver, cancelled by `proxy_cancel`.
119    pub(crate) dns_task: Mutex<Option<JoinHandle<()>>>,
120    /// mDNS publisher for LAN mode (None if LAN mode is disabled).
121    /// Shared with the LAN IP monitor task so it can re-publish on IP change.
122    pub(crate) mdns_publisher:
123        Mutex<Option<std::sync::Arc<tokio::sync::Mutex<crate::proxy::mdns::MdnsPublisher>>>>,
124    /// Join handle for the LAN IP monitor task.
125    pub(crate) lan_monitor_task: Mutex<Option<JoinHandle<()>>>,
126    /// Cancellation token for the background state flush task.
127    pub(crate) flush_cancel: std::sync::Mutex<Option<tokio_util::sync::CancellationToken>>,
128    /// Daemons that currently have a live monitoring task (child `wait()`
129    /// monitor or adopted-orphan poll monitor), keyed to the PID being
130    /// monitored plus a unique registration token. Lets orphan
131    /// reconciliation tell a supervised daemon from one whose monitor died
132    /// with a previous supervisor process.
133    pub(crate) monitored: std::sync::Mutex<HashMap<DaemonId, adopt::MonitorEntry>>,
134    /// Where to deliver output a log sink reports over IPC.
135    ///
136    /// A daemon whose output is captured by a sink writes nothing this process
137    /// reads, so the sink evaluates the daemon's readiness pattern itself and
138    /// sends back the line that matched. The monitoring task registers here
139    /// before the sink starts and unregisters when it ends, so a line arriving
140    /// from a sink that outlived its daemon has nowhere to go and is dropped.
141    pub(crate) sink_output: std::sync::Mutex<HashMap<DaemonId, log_sink::Relay>>,
142    /// Per-daemon stop locks. A stop holds the daemon's lock for its whole
143    /// duration (which now includes waiting for the entire process group to
144    /// exit), and starts/orphan-cleanup acquire it first — serializing them
145    /// against in-flight stops instead of racing the Stopping window.
146    pub(crate) stop_locks: Mutex<HashMap<DaemonId, std::sync::Arc<tokio::sync::Mutex<()>>>>,
147    /// Set when `close()` begins. From then on this supervisor's own
148    /// state-file record is on its way out and must not be restored.
149    pub(crate) shutting_down: AtomicBool,
150}
151
152/// A line of daemon output on its way to the monitoring task.
153#[derive(Debug, Clone)]
154pub(crate) struct OutputLine {
155    pub(crate) text: String,
156    pub(crate) source: OutputSource,
157}
158
159/// Where a line of daemon output came from, which decides what is left to do
160/// with it.
161#[derive(Debug, Clone, Copy, PartialEq, Eq)]
162pub(crate) enum OutputSource {
163    /// Read by this process from the daemon's stdout, stderr or PTY master.
164    /// Nothing has been done with it yet.
165    Local,
166    /// Reported by the daemon's log sink, which has already written the line to
167    /// the store — storing it again here would duplicate it.
168    ///
169    /// `fires_hook` is the sink's answer to whether this line passed the
170    /// `on_output` hook's filter and debounce. It is carried rather than
171    /// re-derived because a line can be reported for readiness alone, and
172    /// firing a hook that filters for something else would be wrong.
173    Sink { fires_hook: bool },
174}
175
176pub(crate) fn interval_duration() -> Duration {
177    settings().general_interval()
178}
179
180pub static SUPERVISOR: Lazy<Supervisor> =
181    Lazy::new(|| Supervisor::new().expect("Error creating supervisor"));
182
183pub async fn start_if_not_running() -> Result<()> {
184    let sf = StateFile::get();
185    if let Some(d) = sf.daemons.get(&DaemonId::pitchfork())
186        && supervisor_record_is_live(d)
187    {
188        return Ok(());
189    }
190    // The record can be missing or stale while a supervisor keeps serving
191    // IPC, e.g. when state.toml was replaced or rewritten by another tool.
192    // Starting another one would take over the socket and leave the running
193    // supervisor, and every daemon it manages, unreachable.
194    if crate::ipc::supervisor_listening().await {
195        debug!("supervisor is listening on the IPC socket but not recorded in the state file");
196        return Ok(());
197    }
198    start_in_background()
199}
200
201/// How long `supervisor start/run --force` waits for the supervisor it
202/// replaces to stop serving IPC. A supervisor being stopped keeps its socket until its daemons
203/// have stopped, which can take a while.
204pub(crate) const IPC_SOCKET_RELEASE_TIMEOUT: Duration = Duration::from_secs(30);
205
206/// Wait until no supervisor is listening on the IPC socket, failing if one
207/// still is after [`IPC_SOCKET_RELEASE_TIMEOUT`].
208///
209/// Used when replacing a supervisor: the old one keeps serving IPC while it
210/// shuts down, and binding the socket before it lets go would leave clients
211/// talking to whichever supervisor they happen to reach.
212pub(crate) async fn wait_for_ipc_socket_release() -> Result<()> {
213    let deadline = time::Instant::now() + IPC_SOCKET_RELEASE_TIMEOUT;
214    let mut waiting = false;
215    while crate::ipc::supervisor_listening().await {
216        if time::Instant::now() >= deadline {
217            return Err(miette::miette!(
218                "another pitchfork supervisor is still listening on {} after {}s; \
219                 stop it with `pitchfork supervisor stop` before starting a new one",
220                crate::ipc::socket_display(),
221                IPC_SOCKET_RELEASE_TIMEOUT.as_secs()
222            ));
223        }
224        if !waiting {
225            info!(
226                "waiting for the supervisor listening on {} to shut down",
227                crate::ipc::socket_display()
228            );
229            waiting = true;
230        }
231        time::sleep(Duration::from_millis(100)).await;
232    }
233    Ok(())
234}
235
236/// Whether the supervisor's own state-file record still describes a live
237/// pitchfork supervisor, rather than a stale entry whose PID the OS has since
238/// handed to an unrelated process.
239///
240/// The record survives crashes and reboots, so a bare liveness probe on its
241/// PID is not enough: after a reboot low PIDs go to early system daemons, and
242/// a `kill(pid, 0)` on one of those says "alive". The check therefore also
243/// requires the identity the supervisor recorded about itself at startup to
244/// match the live process — see [`supervisor_identity_matches`]. When the
245/// PID is alive but the identity does not match, the record is stale; callers
246/// are free to overwrite it and must never signal that PID.
247pub(crate) fn supervisor_record_is_live(record: &crate::daemon::Daemon) -> bool {
248    let Some(pid) = record.pid else {
249        return false;
250    };
251    if !PROCS.is_running(pid) {
252        return false;
253    }
254    if record.start_time.is_none() && record.boot_time.is_none() {
255        // A record from a supervisor older than v2.18.0 carries no identity
256        // at all, so nothing can contradict it. Rather than trust any live
257        // PID, require the process to at least be a pitchfork binary: that
258        // rejects the reboot case (an unrelated system daemon on the PID)
259        // while a still-running old supervisor stays recognisable.
260        PROCS.refresh_pids(&[pid]);
261        let title = PROCS.title(pid);
262        if legacy_supervisor_title_matches(title.as_deref()) {
263            return true;
264        }
265        warn!(
266            "pid {pid} recorded for the supervisor by an older pitchfork is now {title:?}, not a pitchfork process; treating the record as stale"
267        );
268        return false;
269    }
270    if supervisor_identity_matches(
271        record.start_time,
272        PROCS.start_time(pid),
273        record.boot_time,
274        PROCS.boot_time(),
275    ) {
276        return true;
277    }
278    warn!(
279        "pid {pid} recorded for the supervisor belongs to another process now (recorded start_time {:?} boot_time {:?}, live start_time {:?} boot_time {}); treating the record as stale",
280        record.start_time,
281        record.boot_time,
282        PROCS.start_time(pid),
283        PROCS.boot_time()
284    );
285    false
286}
287
288/// Whether a live process name can be a pitchfork supervisor. Only used for
289/// legacy records that carry no start or boot time (see
290/// [`supervisor_record_is_live`]); an unreadable name is not accepted, since
291/// the record has nothing else vouching for it.
292pub(crate) fn legacy_supervisor_title_matches(title: Option<&str>) -> bool {
293    title.is_some_and(|t| t.to_ascii_lowercase().starts_with("pitchfork"))
294}
295
296/// Whether a live process can be the supervisor a state-file record describes.
297///
298/// The kernel start token is the identity, as in [`process_identity_matches`]:
299/// when both the recorded and the live token can be read, they alone decide.
300/// Equal tokens mean the same process generation; different tokens mean the
301/// PID was recycled, within this boot or across a reboot.
302///
303/// The recorded boot time is only consulted when a token is missing on either
304/// side. It must not veto matching tokens: on Linux and macOS the reported
305/// boot time is derived from the realtime clock, so an NTP step or a
306/// sleep/resume moves it while the supervisor keeps running, and treating that
307/// as a reboot would spawn a second supervisor and orphan the first. Without
308/// tokens, though, a boot time from a previous boot is the one thing that can
309/// still prove the record stale, and a record predating both fields is not
310/// contradicted by anything (callers apply a weaker check to those).
311pub(crate) fn supervisor_identity_matches(
312    recorded_start_time: Option<u64>,
313    current_start_time: Option<u64>,
314    recorded_boot_time: Option<u64>,
315    current_boot_time: u64,
316) -> bool {
317    match (recorded_start_time, current_start_time) {
318        (Some(recorded), Some(current)) => recorded == current,
319        _ => recorded_boot_time.is_none_or(|recorded| {
320            recorded.abs_diff(current_boot_time) <= BOOT_TIME_TOLERANCE_SECS
321        }),
322    }
323}
324
325pub fn start_in_background() -> Result<()> {
326    debug!("starting supervisor in background");
327    // Ensure the log directory exists so we can redirect stderr there.
328    // Panics and other fatal errors from the background supervisor process
329    // would otherwise be silently swallowed.
330    let log_file = &*env::PITCHFORK_LOG_FILE;
331    if let Some(parent) = log_file.parent() {
332        let _ = fs::create_dir_all(parent);
333    }
334    #[cfg(unix)]
335    fix_state_dir_permissions();
336
337    // On Unix, use duct with stderr redirected to the log file.
338    #[cfg(unix)]
339    {
340        let stderr_file = fs::OpenOptions::new()
341            .create(true)
342            .append(true)
343            .open(log_file)
344            .into_diagnostic()?;
345        // Computed here rather than in the child: sysconf is not on the
346        // async-signal-safe list.
347        let max_fd = max_inherited_fd();
348        cmd!(&*env::PITCHFORK_BIN, "supervisor", "run")
349            .env_remove("PITCHFORK_CONFIG")
350            .stdin_null()
351            .stdout_null()
352            .stderr_file(stderr_file)
353            .before_spawn(move |cmd| {
354                use std::os::unix::process::CommandExt;
355                // SAFETY: the hook only issues close_range/fcntl syscalls,
356                // which are async-signal-safe, and does not allocate.
357                unsafe {
358                    cmd.pre_exec(move || {
359                        cloexec_inherited_fds(max_fd);
360                        Ok(())
361                    });
362                }
363                Ok(())
364            })
365            .start()
366            .into_diagnostic()?;
367    }
368
369    // On Windows, use CreateProcessW directly with bInheritHandles=FALSE.
370    // std::process::Command always sets bInheritHandles=TRUE when any stdio
371    // handle is configured (even Stdio::null()), which causes the background
372    // supervisor to inherit ALL inheritable handles from the parent —
373    // including bats' stdout capture pipe. The supervisor keeps the pipe
374    // open after the CLI exits, and bats hangs forever waiting for EOF.
375    //
376    // CreateProcessW with bInheritHandles=FALSE prevents any handle
377    // inheritance. We pass NUL device handles for stdin/stdout/stderr
378    // via STARTUPINFO without inheriting any parent handles.
379    #[cfg(windows)]
380    {
381        use windows_sys::Win32::Foundation::{CloseHandle, FALSE};
382        use windows_sys::Win32::System::Threading::{
383            CREATE_NO_WINDOW, CREATE_UNICODE_ENVIRONMENT, CreateProcessW, DETACHED_PROCESS,
384            PROCESS_INFORMATION, STARTUPINFOW,
385        };
386
387        // With bInheritHandles=FALSE, the child inherits NO parent handles.
388        // The supervisor uses its own internal file-based logger
389        // (PITCHFORK_LOG_FILE), so it doesn't need stdio from the parent.
390        // We don't set STARTF_USESTDHANDLES because that flag requires
391        // bInheritHandles=TRUE to function correctly per Microsoft docs.
392        // Without stdio handles, the detached process gets null stdio by
393        // default, which is exactly what we want.
394        let mut si: STARTUPINFOW = unsafe { std::mem::zeroed() };
395        si.cb = std::mem::size_of::<STARTUPINFOW>() as u32;
396
397        let bin_path = &*env::PITCHFORK_BIN;
398        let mut cmd_line: Vec<u16> = format!("\"{}\" supervisor run\0", bin_path.to_string_lossy())
399            .encode_utf16()
400            .collect();
401
402        use std::os::windows::ffi::OsStrExt;
403        let mut vars: Vec<_> = std::env::vars_os()
404            .filter(|(k, _)| !k.to_string_lossy().eq_ignore_ascii_case("PITCHFORK_CONFIG"))
405            .collect();
406        vars.sort_by_key(|(k, _)| k.to_string_lossy().to_uppercase());
407        let mut environment = Vec::<u16>::new();
408        for (key, value) in vars {
409            environment.extend(key.encode_wide());
410            environment.push(b'=' as u16);
411            environment.extend(value.encode_wide());
412            environment.push(0);
413        }
414        environment.extend([0, 0]);
415        let mut pi: PROCESS_INFORMATION = unsafe { std::mem::zeroed() };
416        let ok = unsafe {
417            CreateProcessW(
418                std::ptr::null(),
419                cmd_line.as_mut_ptr(),
420                std::ptr::null(),
421                std::ptr::null(),
422                FALSE, // bInheritHandles = FALSE — the whole point
423                DETACHED_PROCESS | CREATE_NO_WINDOW | CREATE_UNICODE_ENVIRONMENT,
424                environment.as_ptr().cast(),
425                std::ptr::null(),
426                &si,
427                &mut pi,
428            )
429        };
430
431        if ok == 0 {
432            return Err(miette::miette!(
433                "CreateProcessW failed for supervisor: {}",
434                std::io::Error::last_os_error()
435            ));
436        }
437
438        // Close process/thread handles — we don't need them (detached process).
439        unsafe {
440            CloseHandle(pi.hProcess);
441            CloseHandle(pi.hThread);
442        }
443    }
444
445    Ok(())
446}
447
448/// Upper bound (exclusive) for the fd scan in [`cloexec_fd_scan`], used when
449/// `close_range` is unavailable (macOS, Linux < 5.11).
450#[cfg(unix)]
451fn max_inherited_fd() -> libc::c_int {
452    // SAFETY: sysconf has no preconditions.
453    let open_max = unsafe { libc::sysconf(libc::_SC_OPEN_MAX) };
454    // RLIMIT_NOFILE alone is not a safe bound: a caller can open a
455    // high-numbered descriptor and then lower the limit before running us.
456    // Also list the fds that are actually open (both macOS and Linux provide
457    // /dev/fd) so the scan reaches every one of them.
458    let open_fds = std::fs::read_dir("/dev/fd")
459        .into_iter()
460        .flatten()
461        .filter_map(|entry| entry.ok()?.file_name().to_str()?.parse().ok());
462    fd_scan_end(open_max, open_fds)
463}
464
465#[cfg(unix)]
466fn fd_scan_end(
467    open_max: libc::c_long,
468    open_fds: impl IntoIterator<Item = libc::c_int>,
469) -> libc::c_int {
470    // -1 means "no limit"; cap the scan so an unlimited or huge RLIMIT_NOFILE
471    // doesn't make spawning the supervisor noticeably slow.
472    let limit = if open_max <= 0 {
473        1 << 16
474    } else {
475        open_max.min(1 << 20) as libc::c_int
476    };
477    let past_highest_open = open_fds
478        .into_iter()
479        .max()
480        .map_or(0, |fd| fd.saturating_add(1));
481    limit.max(past_highest_open)
482}
483
484/// Keep descriptors the CLI inherited without O_CLOEXEC from leaking into the
485/// background supervisor (and from there into every daemon it spawns).
486///
487/// Callers routinely hand us such descriptors: bats' fd 3, or pipes from a
488/// wrapping script. The long-lived supervisor would otherwise hold them open
489/// forever, so whoever waits for EOF on the pipe (bats, `$(...)`, CI log
490/// capture) hangs even after the CLI has exited. This is the Unix counterpart
491/// of `bInheritHandles=FALSE` in the Windows branch of `start_in_background`.
492///
493/// Runs in the forked child after stdio has been set up on fds 0-2. It marks
494/// fds >= 3 close-on-exec instead of closing them outright so that std's
495/// internal exec-status pipe (already CLOEXEC) keeps working and exec failures
496/// are still reported to the parent. Must stay async-signal-safe: no
497/// allocation, no locks, syscalls only.
498#[cfg(unix)]
499fn cloexec_inherited_fds(max_fd: libc::c_int) {
500    #[cfg(any(target_os = "linux", target_os = "android"))]
501    {
502        // SAFETY: close_range takes plain integer arguments. CLOSE_RANGE_CLOEXEC
503        // needs Linux 5.11; older kernels return ENOSYS/EINVAL and we fall
504        // through to the scan below.
505        let ret = unsafe {
506            libc::syscall(
507                libc::SYS_close_range,
508                3 as libc::c_uint,
509                libc::c_uint::MAX,
510                libc::CLOSE_RANGE_CLOEXEC as libc::c_uint,
511            )
512        };
513        if ret == 0 {
514            return;
515        }
516    }
517    cloexec_fd_scan(max_fd);
518}
519
520/// Fallback for [`cloexec_inherited_fds`]: mark fds `3..max_fd` close-on-exec
521/// one at a time. Async-signal-safe.
522#[cfg(unix)]
523fn cloexec_fd_scan(max_fd: libc::c_int) {
524    for fd in 3..max_fd {
525        // SAFETY: fcntl on an arbitrary fd number is safe; unused numbers
526        // just fail with EBADF.
527        unsafe {
528            let flags = libc::fcntl(fd, libc::F_GETFD);
529            if flags >= 0 && flags & libc::FD_CLOEXEC == 0 {
530                libc::fcntl(fd, libc::F_SETFD, flags | libc::FD_CLOEXEC);
531            }
532        }
533    }
534}
535
536/// Decide whether a project session should be removed during refresh.
537///
538/// `recorded_title` is the title snapshot taken at the start of refresh.
539/// `session` is the current state entry re-read under the lock. If the state
540/// has been updated since the snapshot (e.g., re-entered with the same
541/// PID/dir but a new title), we must skip removal to avoid deleting the new
542/// session. The host PID lives in the session key now, so there is no
543/// `liveness_pid` field to compare against.
544#[cfg(any(unix, test))]
545fn should_remove_liveness_session(
546    session: &crate::state_file::ProjectSession,
547    recorded_title: &Option<String>,
548    current_title: Option<&str>,
549    is_running: bool,
550) -> bool {
551    // If the state was updated since the snapshot (e.g., re-entered with the
552    // same PID/dir but a new title), skip removal to avoid deleting the new
553    // session.
554    if session.liveness_title.as_ref() != recorded_title.as_ref() {
555        return false;
556    }
557    // Dead host process — evict.
558    if !is_running {
559        return true;
560    }
561    // Host is alive. Evict only on a real title mismatch (PID reuse). If no
562    // title was recorded, or the current title is unavailable, we cannot
563    // reliably detect PID reuse — keep the session rather than risk evicting
564    // a live process.
565    match (recorded_title.as_deref(), current_title) {
566        (Some(recorded), Some(current)) => recorded != current,
567        _ => false,
568    }
569}
570
571impl Supervisor {
572    pub fn new() -> Result<Self> {
573        Ok(Self {
574            state_file: Mutex::new(StateFile::read(&*env::PITCHFORK_STATE_FILE).unwrap_or_else(
575                |e| {
576                    warn!("failed to read state file, starting with empty state: {e}");
577                    StateFile::new(env::PITCHFORK_STATE_FILE.clone())
578                },
579            )),
580            last_refreshed_at: Mutex::new(time::Instant::now()),
581            pending_notifications: Mutex::new(vec![]),
582            retrying: std::sync::Mutex::new(HashMap::new()),
583            stop_epochs: std::sync::Mutex::new(HashMap::new()),
584            pending_autostops: Mutex::new(HashMap::new()),
585            in_flight_autostops: Mutex::new(HashMap::new()),
586            ipc_shutdown: Mutex::new(None),
587            hook_tasks: Mutex::new(Vec::new()),
588            active_monitors: AtomicU32::new(0),
589            monitor_done: Notify::new(),
590            proxy_cancel: Mutex::new(None),
591            proxy_task: Mutex::new(None),
592            dns_task: Mutex::new(None),
593            mdns_publisher: Mutex::new(None),
594            lan_monitor_task: Mutex::new(None),
595            flush_cancel: std::sync::Mutex::new(None),
596            monitored: std::sync::Mutex::new(HashMap::new()),
597            sink_output: std::sync::Mutex::new(HashMap::new()),
598            stop_locks: Mutex::new(HashMap::new()),
599            shutting_down: AtomicBool::new(false),
600        })
601    }
602
603    /// Get (or create) the per-daemon stop lock for `id`.
604    pub(crate) async fn stop_lock(&self, id: &DaemonId) -> std::sync::Arc<tokio::sync::Mutex<()>> {
605        self.stop_locks
606            .lock()
607            .await
608            .entry(id.clone())
609            .or_default()
610            .clone()
611    }
612
613    pub async fn start(
614        &self,
615        is_boot: bool,
616        container: bool,
617        web_port: Option<u16>,
618        web_path: Option<String>,
619    ) -> Result<()> {
620        // Ensure the state directory and its contents are accessible by non-root
621        // users. This is needed when the supervisor is started with `sudo` — all
622        // files it creates are owned by root, which prevents normal CLI clients
623        // from reading/writing state or connecting to the IPC socket.
624        #[cfg(unix)]
625        fix_state_dir_permissions();
626
627        // Refuse to run beside a supervisor that is already listening, and
628        // do so before recording ourselves in the state file or starting any
629        // daemons: taking over its socket would leave it running but
630        // unreachable. (`--force` has already waited for the one it replaced
631        // to let go of the socket.) The lock is held until our own listener
632        // is bound, so a supervisor starting at the same time waits here and
633        // then finds this one listening.
634        let startup_lock = StartupLock::acquire().await?;
635        if crate::ipc::supervisor_listening().await {
636            return Err(miette::miette!(
637                "another pitchfork supervisor is already listening on {}",
638                crate::ipc::socket_display()
639            ));
640        }
641
642        let pid = std::process::id();
643        // Ensure PROCS has data for the supervisor PID before upsert_daemon reads title()
644        PROCS.refresh_pids(&[pid]);
645        // Determine container mode: CLI flag takes priority, then settings.
646        // Running as PID 1 always enables it: orphaned descendants of daemons
647        // re-parent to us, and without the zombie reaper they would accumulate
648        // as unreaped zombies — which also keep their process group alive,
649        // stalling whole-group stop waits indefinitely.
650        let container_mode =
651            container || settings().supervisor.container || std::process::id() == 1;
652        if container_mode {
653            info!("Starting supervisor in container/PID1 mode with pid {pid}");
654        } else {
655            info!("Starting supervisor with pid {pid}");
656        }
657
658        // Whether the previous supervisor exited uncleanly must be read before
659        // we record ourselves in the state file just below (see
660        // `supervisor_exited_uncleanly`); the background cleanup task runs
661        // after that record exists, so it receives the answer instead of
662        // reading it too late.
663        let unclean = supervisor_exited_uncleanly(self).await;
664
665        self.upsert_daemon(
666            UpsertDaemonOpts::builder(DaemonId::pitchfork())
667                .set(|o| {
668                    o.pid = Some(pid);
669                    o.status = DaemonStatus::Running;
670                })
671                .build(),
672        )
673        .await?;
674        #[cfg(unix)]
675        fix_state_dir_permissions();
676
677        // Self-heal: if the boot registration points to a stale binary path
678        // (e.g. after a brew/mise upgrade), re-register with the current path.
679        // Runs in the background — must not block or fail supervisor startup.
680        tokio::task::spawn_blocking(|| {
681            if let Ok(boot_manager) = crate::boot_manager::BootManager::new() {
682                boot_manager.check_and_reregister_if_stale();
683            }
684        });
685
686        // If the previous supervisor died uncleanly, its daemon child processes
687        // may still be alive (orphaned, re-parented to init).  Terminate them
688        // before starting replacements so we don't end up with duplicate
689        // processes holding the same ports.
690        //
691        // This runs in the background: each orphan kill now waits for its
692        // whole process group to exit (seconds per orphan), and doing that
693        // inline would delay IPC socket creation past the CLI's short connect
694        // budget on autostart. Per-daemon stop locks serialize the cleanup
695        // against any Run/Stop requests that arrive for the same daemon in the
696        // meantime, and boot daemons start after cleanup completes so they
697        // cannot observe an orphan as "already running".
698        let boot_after_cleanup = is_boot;
699        tokio::spawn(async move {
700            cleanup_orphaned_daemons(&SUPERVISOR, unclean).await;
701            if boot_after_cleanup {
702                info!("Boot start mode enabled, starting boot_start daemons");
703                if let Err(e) = SUPERVISOR.start_boot_daemons().await {
704                    error!("failed to start boot daemons: {e}");
705                }
706            }
707        });
708
709        self.interval_watch()?;
710
711        // Run the first cron check synchronously before starting the cron
712        // watcher and IPC server. This registers config-only cron daemons and
713        // fires any `immediate=true` triggers in the foreground, so they cannot
714        // race with a concurrent `pitchfork start` IPC. By the time the cron
715        // watcher's first tick runs, `last_cron_triggered` is already anchored
716        // and the immediate daemons are already running.
717        if let Err(e) = self.check_cron_schedules().await {
718            error!("failed to check cron schedules on startup: {e}");
719        }
720
721        self.cron_watch()?;
722        self.signals()?;
723        self.daemon_file_watch()?;
724
725        // In container mode, install SIGCHLD handler to reap orphaned/zombie processes
726        #[cfg(unix)]
727        if container_mode {
728            self.reap_zombies()?;
729        }
730
731        // Start web server: CLI --web-port takes priority, then settings.web.auto_start + bind_port
732        let s = settings();
733        let effective_port = web_port.or_else(|| {
734            if s.web.auto_start {
735                match u16::try_from(s.web.bind_port).ok().filter(|&p| p > 0) {
736                    Some(p) => Some(p),
737                    None => {
738                        error!(
739                            "web.bind_port {} is out of valid port range (1-65535), web UI disabled",
740                            s.web.bind_port
741                        );
742                        None
743                    }
744                }
745            } else {
746                None
747            }
748        });
749        // CLI --web-path takes priority, then settings.web.base_path
750        let effective_path = web_path.or_else(|| {
751            let bp = s.web.base_path.clone();
752            if bp.is_empty() { None } else { Some(bp) }
753        });
754        if let Some(port) = effective_port {
755            tokio::spawn(async move {
756                if let Err(e) = crate::web::serve(port, effective_path).await {
757                    error!("Web server error: {e}");
758                }
759            });
760        }
761
762        // Start standalone API server if configured
763        let api_port = if s.api.auto_start {
764            match u16::try_from(s.api.bind_port).ok().filter(|&p| p > 0) {
765                Some(p) => Some(p),
766                None => {
767                    error!(
768                        "api.bind_port {} is out of valid port range (1-65535), API server disabled",
769                        s.api.bind_port
770                    );
771                    None
772                }
773            }
774        } else {
775            None
776        };
777        if let Some(port) = api_port {
778            tokio::spawn(async move {
779                if let Err(e) = crate::web::serve_api(port, None).await {
780                    error!("API server error: {e}");
781                }
782            });
783        }
784
785        // Start reverse proxy server if enabled
786        if s.proxy.enable {
787            // Pre-generate the TLS certificate synchronously before spawning the proxy
788            // task. This ensures the cert exists immediately after `sup start` returns,
789            // so `proxy trust` can be run right away without waiting for the async task.
790            #[cfg(feature = "proxy-tls")]
791            if s.proxy.https {
792                let proxy_dir = crate::env::PITCHFORK_STATE_DIR.join("proxy");
793                let ca_cert_path = proxy_dir.join("ca.pem");
794                let ca_key_path = proxy_dir.join("ca-key.pem");
795                // Checked and written under the CA lock: `proxy setup` may be
796                // generating the same pair right now.
797                match crate::proxy::server::ensure_ca(&ca_cert_path, &ca_key_path, || {
798                    ca_cert_path.exists() && ca_key_path.exists()
799                }) {
800                    Ok(true) => {
801                        info!(
802                            "Generated local CA certificate at {}",
803                            ca_cert_path.display()
804                        );
805                    }
806                    Ok(false) => {}
807                    Err(e) => {
808                        error!("Failed to generate CA certificate: {e}");
809                    }
810                }
811
812                // Auto-trust: attempt to install the CA certificate into the
813                // system trust store. May fail silently due to permissions;
814                // user can run `pitchfork proxy trust` manually.
815                if s.proxy.auto_trust && ca_cert_path.exists() {
816                    use crate::proxy::trust::{AutoTrustResult, auto_trust};
817                    match auto_trust(&ca_cert_path) {
818                        AutoTrustResult::AlreadyTrusted => {}
819                        AutoTrustResult::Trusted => {
820                            info!("CA certificate auto-trusted in system store");
821                        }
822                        AutoTrustResult::NotTrusted { reason } => {
823                            warn!("Auto-trust skipped: {reason}");
824                            warn!("Run `pitchfork proxy trust` to install manually");
825                        }
826                    }
827                }
828            }
829            // Spawn the proxy server and wait for its bind result via a oneshot
830            // channel.  This avoids the TOCTOU race of a pre-flight bind check
831            // while still surfacing binding failures immediately.
832            let (bind_tx, bind_rx) = tokio::sync::oneshot::channel();
833            let proxy_cancel = tokio_util::sync::CancellationToken::new();
834            let proxy_cancel_clone = proxy_cancel.clone();
835            *self.proxy_cancel.lock().await = Some(proxy_cancel);
836            let proxy_task = tokio::spawn(async move {
837                if let Err(e) = crate::proxy::server::serve(bind_tx, proxy_cancel_clone).await {
838                    error!("Proxy server error: {e}");
839                }
840            });
841            *self.proxy_task.lock().await = Some(proxy_task);
842            match bind_rx.await {
843                Ok(Ok(())) => {
844                    info!("Proxy server bound successfully");
845                    // Resolved once and shared: mDNS and the DNS resolver
846                    // must advertise the same address.
847                    let lan_ip = self.resolve_lan_ip().await;
848                    self.start_mdns(lan_ip).await;
849                    self.start_dns_resolver(lan_ip).await;
850                }
851                Ok(Err(msg)) => {
852                    error!("{msg}");
853                    self.add_notification(log::LevelFilter::Error, msg).await;
854                }
855                Err(_) => {
856                    // Sender dropped without sending — serve() panicked or
857                    // returned before signalling.  Already logged by the
858                    // spawn error handler above.
859                }
860            }
861        }
862
863        // Pre-warm slug cache so the first /api/proxies request is fast.
864        // Spawned as a background task so it does not block startup.
865        tokio::spawn(async {
866            crate::proxy::server::get_cached_slugs().await;
867        });
868
869        let (ipc, ipc_handle) = IpcServer::new(startup_lock).await?;
870        *self.ipc_shutdown.lock().await = Some(ipc_handle);
871        self.start_state_flush_task();
872        self.conn_watch(ipc).await
873    }
874
875    /// Start the loopback DNS resolver for the proxy TLD.
876    ///
877    /// The resolver shares the proxy's cancellation token, so it stops with the
878    /// proxy. A bind failure is a notification rather than a fatal error: the
879    /// proxy still works for anyone who reaches it some other way.
880    async fn start_dns_resolver(&self, lan_ip: Option<std::net::Ipv4Addr>) {
881        let s = crate::settings::settings();
882        if !s.proxy.dns {
883            return;
884        }
885        let cfg = crate::proxy::dns::config_from_settings(&s, lan_ip);
886        let port = crate::proxy::dns::dns_port(&s);
887        // Loopback only: the resolver is for this machine's stub resolver, and
888        // LAN peers are served by mDNS instead.
889        let addr = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, port));
890
891        // The token's guard is held until the handle is stored. `close` takes
892        // the token under this lock before it collects `dns_task`, so it either
893        // runs first — and there is no token to spawn with — or waits until
894        // the handle is in place to be drained. Releasing it earlier left a
895        // window where `close` found no handle and the task outlived shutdown
896        // holding the resolver's sockets.
897        let cancel_guard = self.proxy_cancel.lock().await;
898        let Some(cancel) = cancel_guard.clone() else {
899            return;
900        };
901        let (bind_tx, bind_rx) = tokio::sync::oneshot::channel();
902        let task = tokio::spawn(async move {
903            if let Err(e) = crate::proxy::dns::serve(cfg, addr, bind_tx, cancel).await {
904                error!("DNS resolver error: {e}");
905            }
906        });
907        *self.dns_task.lock().await = Some(task);
908        drop(cancel_guard);
909        match bind_rx.await {
910            Ok(Ok(())) => info!("DNS resolver bound successfully"),
911            Ok(Err(msg)) => {
912                let msg = format!(
913                    "{msg}\nProxy host names will not resolve through pitchfork. \
914                     Choose another port with proxy.dns_port, or set proxy.dns = false."
915                );
916                error!("{msg}");
917                self.add_notification(log::LevelFilter::Error, msg).await;
918            }
919            Err(_) => {}
920        }
921    }
922
923    /// Watch the LAN address and keep the DNS responder, and mDNS when it is
924    /// running, pointed at the current one.
925    ///
926    /// Started even when the mDNS publisher could not be created: the DNS
927    /// responder serves this machine regardless, and an address it keeps
928    /// answering with after the interface has moved is worse than useless.
929    ///
930    /// Does nothing when `proxy.lan_ip` pins an address. That is a choice to
931    /// respect, not a starting point to drift from — the check lives here so
932    /// neither caller can forget it.
933    async fn start_lan_ip_monitor(
934        &self,
935        initial_ip: std::net::Ipv4Addr,
936        port: u16,
937        publisher: Option<std::sync::Arc<tokio::sync::Mutex<crate::proxy::mdns::MdnsPublisher>>>,
938    ) {
939        if !crate::settings::settings().proxy.lan_ip.is_empty() {
940            return;
941        }
942        // No cancellation token means `close` has already taken it, so shutdown
943        // is under way. Spawning here would leave a task polling with no way to
944        // stop it, and past the point where `close` collects the handle.
945        //
946        // Held until the handle is stored, for the reason `start_dns_resolver`
947        // gives: otherwise `close` can collect `lan_monitor_task` in between.
948        let cancel_guard = self.proxy_cancel.lock().await;
949        let Some(cancel) = cancel_guard.clone() else {
950            debug!("Not starting the LAN IP monitor: the supervisor is shutting down");
951            return;
952        };
953        let monitor_cancel = Some(cancel);
954        let task = tokio::spawn(async move {
955            let mut last_ip = initial_ip;
956            let mut ticker = tokio::time::interval(std::time::Duration::from_secs(5));
957            ticker.tick().await; // first tick is immediate
958            loop {
959                // Cancellation is raced against the tick, not checked after
960                // it. Checking afterwards means the task only notices once the
961                // full interval has elapsed, so a shutdown that waits a second
962                // for it always gives up and aborts instead — the graceful
963                // path would never once be taken.
964                match monitor_cancel.as_ref() {
965                    Some(cancel) => {
966                        tokio::select! {
967                            _ = ticker.tick() => {}
968                            _ = cancel.cancelled() => break,
969                        }
970                    }
971                    None => {
972                        ticker.tick().await;
973                    }
974                }
975                if let Some(new_ip) = crate::proxy::lan_ip::detect_lan_ip_if_changed(last_ip).await
976                {
977                    log::info!("LAN IP changed: {last_ip} → {new_ip}");
978                    last_ip = new_ip;
979                    crate::proxy::dns::update_lan_ip(new_ip);
980                    if let Some(publisher) = publisher.as_ref() {
981                        publisher.lock().await.republish_all(new_ip, port);
982                    }
983                }
984            }
985        });
986        *self.lan_monitor_task.lock().await = Some(task);
987        drop(cancel_guard);
988    }
989
990    /// Start mDNS publishing for LAN mode (called after the proxy binds successfully).
991    /// The LAN address mDNS publishes and the DNS resolver answers with.
992    ///
993    /// Resolved once and handed to both. Detecting separately in each let them
994    /// disagree when the interface address changed in between, which would
995    /// advertise one address over mDNS and serve another over DNS, and probed
996    /// the network twice at startup for one answer.
997    ///
998    /// `None` means LAN mode is off, or is on and the address could not be
999    /// determined; either way the caller has nothing to publish. The reason is
1000    /// reported here so it is said once rather than by each caller.
1001    async fn resolve_lan_ip(&self) -> Option<std::net::Ipv4Addr> {
1002        let s = crate::settings::settings();
1003        let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
1004        if !s.proxy.enable || !lan_enabled {
1005            return None;
1006        }
1007        if s.proxy.lan_ip.is_empty() {
1008            let detected = crate::proxy::lan_ip::detect_lan_ip().await;
1009            if detected.is_none() {
1010                error!(
1011                    "LAN mode is enabled but no LAN IP address could be detected. \
1012                     Set proxy.lan_ip to a specific address, or ensure you are connected to a network."
1013                );
1014            }
1015            return detected;
1016        }
1017        match s.proxy.lan_ip.parse::<std::net::Ipv4Addr>() {
1018            Ok(ip) => Some(ip),
1019            Err(e) => {
1020                let msg = format!(
1021                    concat!(
1022                        "proxy.lan_ip {:?} is not a valid IPv4 address: {}. ",
1023                        "LAN mode will not start; fix the setting or clear it ",
1024                        "to auto-detect."
1025                    ),
1026                    s.proxy.lan_ip, e
1027                );
1028                error!("{msg}");
1029                self.add_notification(log::LevelFilter::Error, msg).await;
1030                None
1031            }
1032        }
1033    }
1034
1035    async fn start_mdns(&self, lan_ip: Option<std::net::Ipv4Addr>) {
1036        let s = crate::settings::settings();
1037        let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
1038        if !s.proxy.enable || !lan_enabled {
1039            return;
1040        }
1041
1042        let Some(lan_ip) = lan_ip else { return };
1043        let port = u16::try_from(s.proxy.port).unwrap_or(443);
1044
1045        let Some(mut publisher) = crate::proxy::mdns::MdnsPublisher::new(lan_ip) else {
1046            error!("Failed to start mDNS publisher. Is Avahi (Linux) or Bonjour (macOS) running?");
1047            // The DNS responder hands out this address too, and it is useful on
1048            // this machine whether or not mDNS came up. Keep watching the
1049            // interface so its answers do not go stale.
1050            self.start_lan_ip_monitor(lan_ip, port, None).await;
1051            return;
1052        };
1053
1054        // Publish all registered slugs.
1055        let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
1056        for slug in slugs.keys() {
1057            let hostname = format!("{slug}.local");
1058            publisher.publish(&hostname, port);
1059        }
1060
1061        log::info!(
1062            "LAN mode: mDNS publishing on {lan_ip}, {} slug(s) registered",
1063            slugs.len()
1064        );
1065
1066        let publisher = std::sync::Arc::new(tokio::sync::Mutex::new(publisher));
1067
1068        // Start the IP monitor. It declines on its own when the address is
1069        // pinned rather than auto-detected.
1070        self.start_lan_ip_monitor(lan_ip, port, Some(publisher.clone()))
1071            .await;
1072
1073        *self.mdns_publisher.lock().await = Some(publisher);
1074    }
1075
1076    /// Re-read slugs from config and update mDNS records.
1077    ///
1078    /// Publishes new slugs and unpublishes removed ones. Called via IPC when
1079    /// `proxy add` or `proxy remove` modifies the slug registry.
1080    async fn sync_mdns(&self) {
1081        // Clone the Arc and release the outer lock immediately so we don't
1082        // block close() from taking the publisher during shutdown.
1083        let publisher = {
1084            let guard = self.mdns_publisher.lock().await;
1085            match guard.as_ref() {
1086                Some(p) => p.clone(),
1087                None => {
1088                    debug!("sync_mdns: mDNS publisher not active, skipping");
1089                    return;
1090                }
1091            }
1092        };
1093
1094        let s = crate::settings::settings();
1095        let port = u16::try_from(s.proxy.port).unwrap_or(443);
1096
1097        let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
1098        let mut pub_guard = publisher.lock().await;
1099
1100        // Unpublish slugs that no longer exist in config.
1101        let current_keys: Vec<&String> = slugs.keys().collect();
1102        let registered: Vec<String> = pub_guard.registered_hostnames();
1103        for hostname in &registered {
1104            // hostname is "slug.local" — extract slug part.
1105            let slug = hostname.strip_suffix(".local").unwrap_or(hostname);
1106            if !current_keys.iter().any(|k| k.as_str() == slug) {
1107                log::info!("mDNS: unpublishing removed slug {slug}");
1108                pub_guard.unpublish(hostname);
1109            }
1110        }
1111
1112        // Publish new slugs that aren't yet registered.
1113        for slug in slugs.keys() {
1114            let hostname = format!("{slug}.local");
1115            if !pub_guard.is_published(&hostname) {
1116                log::info!("mDNS: publishing new slug {slug}");
1117                pub_guard.publish(&hostname, port);
1118            }
1119        }
1120    }
1121
1122    /// Spawn a background task that periodically flushes the state file to
1123    /// disk if it has been marked dirty.  Uses debouncing (1s interval) to
1124    /// batch rapid state changes.
1125    fn start_state_flush_task(&self) {
1126        let cancel = tokio_util::sync::CancellationToken::new();
1127        *self.flush_cancel.lock().unwrap() = Some(cancel.clone());
1128        tokio::spawn(async move {
1129            let mut interval = time::interval(Duration::from_secs(1));
1130            interval.set_missed_tick_behavior(time::MissedTickBehavior::Skip);
1131            loop {
1132                tokio::select! {
1133                    _ = interval.tick() => {}
1134                    _ = cancel.cancelled() => {
1135                        debug!("state flush task received shutdown signal");
1136                        break;
1137                    }
1138                }
1139                let state = SUPERVISOR.state_file.lock().await;
1140                if state.is_dirty()
1141                    && let Err(e) = state.write()
1142                {
1143                    warn!("failed to flush state file: {e}");
1144                }
1145            }
1146            debug!("state flush task exiting");
1147        });
1148    }
1149
1150    pub(crate) async fn flush_state(&self) {
1151        let state = self.state_file.lock().await;
1152        if state.is_dirty()
1153            && let Err(e) = state.write()
1154        {
1155            warn!("failed to flush state file: {e}");
1156        }
1157    }
1158
1159    pub(crate) async fn refresh(&self) -> Result<()> {
1160        trace!("refreshing");
1161
1162        // Collect PIDs we need to check (shell PIDs and liveness PIDs)
1163        // This is more efficient than refreshing all processes on the system
1164        let dirs_with_pids = self.get_dirs_with_shell_pids().await;
1165        let liveness_sessions = self.get_liveness_sessions().await;
1166        let pids_to_check: Vec<u32> = dirs_with_pids
1167            .values()
1168            .flatten()
1169            .copied()
1170            .chain(liveness_sessions.iter().map(|(pid, _, _)| *pid))
1171            .collect::<std::collections::HashSet<_>>()
1172            .into_iter()
1173            .collect();
1174
1175        if pids_to_check.is_empty() {
1176            // No PIDs to check, skip the expensive refresh
1177            trace!("no tracked PIDs to check, skipping process refresh");
1178        } else {
1179            debug!("refreshing PIDs: {pids_to_check:?}");
1180            PROCS.refresh_pids(&pids_to_check);
1181        }
1182
1183        let mut last_refreshed_at = self.last_refreshed_at.lock().await;
1184        *last_refreshed_at = time::Instant::now();
1185
1186        self.restore_own_record().await;
1187
1188        #[cfg_attr(not(unix), allow(unused_mut))]
1189        let mut dirs_to_leave: Vec<PathBuf> = Vec::new();
1190
1191        // Prune shell PIDs that are no longer running. This is essential on
1192        // Unix so that exited shells don't keep daemons alive forever.
1193        //
1194        // On Windows, skip this check: Git Bash (MSYS2) PIDs from `$$` are
1195        // Cygwin-internal PIDs that are invisible to sysinfo (which sees
1196        // Windows PIDs). The is_running check would always return false,
1197        // immediately removing every registered shell and breaking autostop.
1198        // Shell registration/deregistration relies on UpdateShellDir IPC
1199        // messages instead.
1200        #[cfg(unix)]
1201        for (dir, pids) in dirs_with_pids {
1202            let to_remove = pids
1203                .iter()
1204                .filter(|pid| !PROCS.is_running(**pid))
1205                .collect::<Vec<_>>();
1206            for pid in &to_remove {
1207                self.remove_shell_pid(**pid).await?
1208            }
1209            if to_remove.len() == pids.len() {
1210                dirs_to_leave.push(dir);
1211            }
1212        }
1213
1214        // Atomically remove project sessions whose host PID has died or whose
1215        // recorded title no longer matches the current process title. Every
1216        // project session carries a host PID in its key, so we iterate all of
1217        // them. Re-reading the sessions under the lock prevents enter/leave
1218        // interleaving from deleting a session that was just replaced with a
1219        // new title snapshot.
1220        //
1221        // Gated to Unix to mirror the shell-PID pruning above: on Windows,
1222        // Git Bash (MSYS2) `$$` PIDs are Cygwin-internal and invisible to
1223        // sysinfo, so the liveness check would immediately revoke every
1224        // freshly-entered session. Windows relies on explicit `project leave`
1225        // (or shell UpdateShellDir) for deregistration instead.
1226        #[cfg(unix)]
1227        {
1228            let mut state = self.state_file.lock().await;
1229            for (pid, dir, recorded_title) in liveness_sessions {
1230                let Some(session) = state.get_project_session(pid, &dir) else {
1231                    continue;
1232                };
1233                let current_title = PROCS.title(pid);
1234                let is_running = PROCS.is_running(pid);
1235                debug!(
1236                    "refresh liveness session pid {pid} dir {} recorded_title={recorded_title:?} current_title={current_title:?} is_running={is_running}",
1237                    dir.display()
1238                );
1239                if should_remove_liveness_session(
1240                    session,
1241                    &recorded_title,
1242                    current_title.as_deref(),
1243                    is_running,
1244                ) {
1245                    warn!(
1246                        "removing project session pid {pid} dir {} (liveness pid title mismatch or dead)",
1247                        dir.display()
1248                    );
1249                    if state.remove_project_session(pid, &dir).is_some() {
1250                        dirs_to_leave.push(dir);
1251                    }
1252                }
1253            }
1254        }
1255
1256        for dir in dirs_to_leave {
1257            self.leave_dir(&dir).await?;
1258        }
1259
1260        // Catch state-`running` daemons that lost their monitor (e.g. the
1261        // monitor died with a previous supervisor): mark dead ones errored
1262        // and re-adopt live ones. Runs before check_retry so a daemon marked
1263        // errored here is retried on this same tick.
1264        self.reconcile_unmonitored_daemons().await;
1265
1266        self.check_retry().await?;
1267        self.process_pending_autostops().await?;
1268
1269        Ok(())
1270    }
1271
1272    /// Install a SIGCHLD handler that reaps orphaned zombie child processes.
1273    ///
1274    /// When running as PID 1 inside a container, orphaned processes are
1275    /// re-parented to PID 1. Without explicit reaping, they accumulate
1276    /// as zombies in the process table indefinitely.
1277    ///
1278    /// Only reaps processes that are NOT managed by the supervisor (i.e.
1279    /// not tracked in the state file). Managed daemon processes are reaped
1280    /// by their monitoring tasks via `child.wait()`.
1281    ///
1282    /// ## Strategy
1283    ///
1284    /// **Linux**: Uses `waitid(Id::All, WNOHANG | WNOWAIT | WEXITED)` to
1285    /// *peek* at the next zombie without consuming its status. If the PID
1286    /// belongs to a managed daemon, the reaper skips it so Tokio's
1287    /// `child.wait()` can collect the status normally. Only unmanaged
1288    /// orphans are actually reaped (via `waitpid(Pid, WNOHANG)`). This
1289    /// eliminates the race entirely.
1290    ///
1291    /// **Non-Linux Unix** (e.g. macOS — mainly for local development;
1292    /// container mode targets Linux): `waitid` is unavailable, so we fall
1293    /// back to `waitpid(None, WNOHANG)`. If the reaper accidentally
1294    /// consumes a managed PID's status, it stashes the exit code in
1295    /// [`REAPED_STATUSES`] for the monitoring task to recover.
1296    #[cfg(unix)]
1297    fn reap_zombies(&self) -> Result<()> {
1298        let mut stream = signal::unix::signal(SignalKind::child())
1299            .map_err(|e| miette::miette!("Failed to register SIGCHLD handler: {e}"))?;
1300        tokio::spawn(async move {
1301            loop {
1302                stream.recv().await;
1303                // Collect PIDs of managed daemons so we don't steal their exit status
1304                let managed_pids: HashSet<u32> = SUPERVISOR
1305                    .state_file
1306                    .lock()
1307                    .await
1308                    .daemons
1309                    .values()
1310                    .filter_map(|d| d.pid)
1311                    .collect();
1312                // Reap all available zombie children that are NOT managed
1313                Self::reap_unmanaged_zombies(&managed_pids).await;
1314            }
1315        });
1316        info!("container mode: SIGCHLD zombie reaper installed");
1317        Ok(())
1318    }
1319
1320    /// Linux implementation: peek with `waitid(WNOWAIT)` then selectively reap.
1321    ///
1322    /// `WNOWAIT` leaves the zombie in the table so we can inspect its PID
1323    /// without consuming the exit status. Only if the PID is *not* managed
1324    /// do we call `waitpid(Pid, WNOHANG)` to actually reap it.
1325    #[cfg(target_os = "linux")]
1326    async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
1327        use nix::sys::wait::{Id, WaitPidFlag, WaitStatus, waitid, waitpid};
1328        use nix::unistd::Pid;
1329
1330        loop {
1331            // Peek at the next zombie without consuming it
1332            let peek_flags = WaitPidFlag::WNOHANG | WaitPidFlag::WNOWAIT | WaitPidFlag::WEXITED;
1333            match waitid(Id::All, peek_flags) {
1334                Ok(WaitStatus::StillAlive) => break,
1335                Ok(status) => {
1336                    let Some(pid_raw) = status.pid().map(|p| p.as_raw() as u32) else {
1337                        break;
1338                    };
1339                    if managed_pids.contains(&pid_raw) {
1340                        // This is a managed daemon — leave it for Tokio's child.wait().
1341                        // We must break out of the loop because waitid(Id::All) would
1342                        // keep returning the same zombie if we don't consume it.
1343                        trace!(
1344                            "zombie reaper: skipping managed daemon pid {pid_raw}, \
1345                             leaving for Tokio to reap"
1346                        );
1347                        break;
1348                    }
1349                    // Not managed — actually reap it
1350                    match waitpid(Pid::from_raw(pid_raw as i32), Some(WaitPidFlag::WNOHANG)) {
1351                        Ok(s) => trace!("reaped orphaned zombie child: {s:?}"),
1352                        Err(nix::errno::Errno::ECHILD) => break,
1353                        Err(e) => {
1354                            trace!("waitpid error reaping pid {pid_raw}: {e}");
1355                            break;
1356                        }
1357                    }
1358                }
1359                Err(nix::errno::Errno::ECHILD) => break, // no children at all
1360                Err(e) => {
1361                    trace!("waitid error in zombie reaper: {e}");
1362                    break;
1363                }
1364            }
1365        }
1366    }
1367
1368    /// Non-Linux fallback: blind `waitpid(None, WNOHANG)` with stash recovery.
1369    ///
1370    /// Since `waitid(WNOWAIT)` is not available, we cannot peek. If we
1371    /// accidentally reap a managed PID, we stash the exit code in
1372    /// [`REAPED_STATUSES`] so the monitoring task can recover it.
1373    #[cfg(all(unix, not(target_os = "linux")))]
1374    async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
1375        use nix::sys::wait::{WaitPidFlag, WaitStatus, waitpid};
1376
1377        loop {
1378            match waitpid(None, Some(WaitPidFlag::WNOHANG)) {
1379                Ok(WaitStatus::StillAlive) => break,
1380                Ok(status) => {
1381                    let Some(pid) = status.pid().map(|p| p.as_raw() as u32) else {
1382                        continue;
1383                    };
1384                    if managed_pids.contains(&pid) {
1385                        // Race lost — stash the exit code for lifecycle recovery
1386                        let exit_code = match status {
1387                            WaitStatus::Exited(_, code) => code,
1388                            WaitStatus::Signaled(_, sig, _) => -(sig as i32),
1389                            _ => -1,
1390                        };
1391                        warn!(
1392                            "zombie reaper reaped managed daemon pid {pid} \
1393                             (exit_code={exit_code}); stashing status for recovery"
1394                        );
1395                        REAPED_STATUSES.lock().await.insert(pid, exit_code);
1396                    } else {
1397                        trace!("reaped orphaned zombie child: {status:?}");
1398                    }
1399                }
1400                Err(nix::errno::Errno::ECHILD) => break, // no more children
1401                Err(e) => {
1402                    trace!("waitpid error in zombie reaper: {e}");
1403                    break;
1404                }
1405            }
1406        }
1407    }
1408
1409    #[cfg(unix)]
1410    fn signals(&self) -> Result<()> {
1411        let signals = [
1412            SignalKind::terminate(),
1413            SignalKind::alarm(),
1414            SignalKind::interrupt(),
1415            SignalKind::quit(),
1416            SignalKind::hangup(),
1417            SignalKind::user_defined1(),
1418            SignalKind::user_defined2(),
1419        ];
1420        static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
1421        for signal in signals {
1422            let stream = match signal::unix::signal(signal) {
1423                Ok(s) => s,
1424                Err(e) => {
1425                    warn!("Failed to register signal handler for {signal:?}: {e}");
1426                    continue;
1427                }
1428            };
1429            tokio::spawn(async move {
1430                let mut stream = stream;
1431                loop {
1432                    stream.recv().await;
1433                    if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
1434                        exit(1);
1435                    } else {
1436                        SUPERVISOR.handle_signal().await;
1437                    }
1438                }
1439            });
1440        }
1441        Ok(())
1442    }
1443
1444    #[cfg(windows)]
1445    fn signals(&self) -> Result<()> {
1446        tokio::spawn(async move {
1447            static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
1448            loop {
1449                if let Err(e) = signal::ctrl_c().await {
1450                    error!("Failed to wait for ctrl-c: {}", e);
1451                    return;
1452                }
1453                if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
1454                    exit(1);
1455                } else {
1456                    SUPERVISOR.handle_signal().await;
1457                }
1458            }
1459        });
1460        Ok(())
1461    }
1462
1463    async fn handle_signal(&self) {
1464        info!("received signal, stopping");
1465        self.close().await;
1466        exit(0)
1467    }
1468
1469    pub(crate) async fn close(&self) {
1470        self.shutting_down.store(true, atomic::Ordering::Release);
1471        // Signal the proxy server to stop accepting new connections
1472        // and drain in-flight ones, *before* stopping daemons so the
1473        // proxy has time to finish forwarding active requests.
1474        if let Some(cancel) = self.proxy_cancel.lock().await.take() {
1475            cancel.cancel();
1476        }
1477
1478        // Stop the LAN IP monitor task. It watches the token that was just
1479        // cancelled, so give it a moment to come back on its own rather than
1480        // cutting it off mid-iteration the instant after asking it to stop.
1481        if let Some(mut monitor_task) = self.lan_monitor_task.lock().await.take()
1482            && tokio::time::timeout(Duration::from_secs(1), &mut monitor_task)
1483                .await
1484                .is_err()
1485        {
1486            monitor_task.abort();
1487        }
1488
1489        // Shutdown the mDNS publisher (sends goodbye packets).
1490        if let Some(publisher) = self.mdns_publisher.lock().await.take() {
1491            publisher.lock().await.shutdown();
1492        }
1493
1494        if let Some(mut dns_task) = self.dns_task.lock().await.take()
1495            && tokio::time::timeout(Duration::from_secs(5), &mut dns_task)
1496                .await
1497                .is_err()
1498        {
1499            // Cancelled but still running: aborted, as the LAN monitor is,
1500            // so its sockets are not left bound to the resolver port.
1501            dns_task.abort();
1502        }
1503
1504        if let Some(proxy_task) = self.proxy_task.lock().await.take() {
1505            // Longer than the proxy's own drain budget, so the task finishes on
1506            // its own terms rather than being cut off mid-drain.
1507            let _ = tokio::time::timeout(
1508                crate::proxy::server::SHUTDOWN_DRAIN_BUDGET + Duration::from_secs(2),
1509                proxy_task,
1510            )
1511            .await;
1512        }
1513
1514        // Clean up /etc/hosts entries managed by pitchfork
1515        let s = settings();
1516        if s.proxy.enable && s.proxy.sync_hosts {
1517            crate::proxy::hosts::clean_hosts_file();
1518        }
1519
1520        let pitchfork_id = DaemonId::pitchfork();
1521        let active = self.active_daemons().await;
1522        let active_ids: Vec<DaemonId> = active
1523            .iter()
1524            .filter(|d| d.id != pitchfork_id)
1525            .map(|d| d.id.clone())
1526            .collect();
1527
1528        // Stop daemons in reverse dependency order.
1529        // If dependency resolution fails (e.g. config changed), fall back to
1530        // stopping in arbitrary order so we still shut down cleanly.
1531        // Daemons within the same level are stopped concurrently.
1532        //
1533        // Each stop waits for the daemon's whole process group (bounded by its
1534        // stop budget) and levels are sequential, so total shutdown time is the
1535        // sum of the slowest stop per level. If an external manager (docker,
1536        // systemd) kills us before this completes, cleanup_orphaned_daemons()
1537        // recovers the leftover processes and stale state on the next start.
1538        let stop_levels = compute_reverse_stop_order(&active_ids);
1539        for level in &stop_levels {
1540            let mut tasks = Vec::new();
1541            for id in level {
1542                let id = id.clone();
1543                tasks.push(tokio::spawn(async move {
1544                    if let Err(err) = SUPERVISOR.stop(&id).await {
1545                        error!("failed to stop daemon {id}: {err}");
1546                    }
1547                }));
1548            }
1549            for task in tasks {
1550                let _ = task.await;
1551            }
1552        }
1553        let _ = self.remove_daemon(&pitchfork_id).await;
1554
1555        // Signal the background state flush task to exit so it doesn't
1556        // keep waking up and acquiring the state mutex after shutdown.
1557        if let Some(cancel) = self.flush_cancel.lock().unwrap().take() {
1558            cancel.cancel();
1559        }
1560
1561        // Force-flush state to disk before shutting down IPC so no
1562        // in-memory-only changes are lost.
1563        {
1564            let state = self.state_file.lock().await;
1565            if state.is_dirty()
1566                && let Err(e) = state.write()
1567            {
1568                warn!("failed to flush state file during shutdown: {e}");
1569            }
1570        }
1571
1572        // Signal IPC server to shut down gracefully
1573        if let Some(mut handle) = self.ipc_shutdown.lock().await.take() {
1574            handle.shutdown().await;
1575        }
1576
1577        // Wait for all in-flight monitoring tasks to finish registering their
1578        // hook handles. Each monitoring task increments `active_monitors` when
1579        // its process exits, and decrements it (+ notifies `monitor_done`)
1580        // after all fire_hook() calls complete. This replaces the old
1581        // yield_now() approach which had a race window.
1582        let drain_timeout = time::sleep(Duration::from_secs(5));
1583        tokio::pin!(drain_timeout);
1584        loop {
1585            if self.active_monitors.load(atomic::Ordering::Acquire) == 0 {
1586                break;
1587            }
1588            tokio::select! {
1589                _ = self.monitor_done.notified() => {}
1590                _ = &mut drain_timeout => {
1591                    warn!("timed out waiting for monitoring tasks to register hooks, proceeding with shutdown");
1592                    break;
1593                }
1594            }
1595        }
1596        let handles: Vec<JoinHandle<()>> = std::mem::take(&mut *self.hook_tasks.lock().await);
1597        let hook_timeout = Duration::from_secs(30);
1598        for handle in handles {
1599            match time::timeout(hook_timeout, handle).await {
1600                Ok(_) => {} // Hook completed (success or error, doesn't matter)
1601                Err(_) => {
1602                    warn!(
1603                        "hook task did not complete within {hook_timeout:?} during shutdown, skipping"
1604                    );
1605                }
1606            }
1607        }
1608
1609        // Unix: remove the socket directory if it is empty. The IPC server
1610        // already removed our socket; anything left belongs to a supervisor
1611        // that replaced this one (e.g. `supervisor run --force`) while we
1612        // were stopping daemons, and must not be deleted.
1613        // Windows: named pipes have no filesystem component.
1614        #[cfg(unix)]
1615        let _ = fs::remove_dir(&*env::IPC_SOCK_DIR);
1616    }
1617
1618    pub(crate) async fn add_notification(&self, level: log::LevelFilter, message: String) {
1619        self.pending_notifications
1620            .lock()
1621            .await
1622            .push((level, message));
1623    }
1624}
1625
1626/// Fix ownership on the state directory so non-root users can access files
1627/// created by a `sudo`-started supervisor.
1628///
1629/// When `[settings.supervisor] user`, a recorded invoking user (see
1630/// [`env::InvokingUser`]), or `SUDO_UID`/`SUDO_GID` are set, we
1631/// `chown` the state directory and safe subdirectories back to that non-root
1632/// runtime user. This is strictly better than `chmod 0o666` because it does not
1633/// widen the permission bits — the files stay owner-only (0o600/0o700) but the
1634/// *owner* is the user that daemon processes and CLI clients need to share.
1635///
1636/// **Security**: The `proxy/` subtree is intentionally skipped. It contains
1637/// `ca-key.pem` which must remain `0o600` and owned by the process that
1638/// generated it. Changing its ownership or permissions would expose the CA
1639/// private key to other local users.
1640///
1641/// If none of these are available (e.g. a direct root login, or a boot service
1642/// installed from a root shell), we fall back to relaxing permissions on only
1643/// the `sock/` and `logs/` subdirectories (plus `state.toml`) so CLI clients
1644/// can still function.
1645#[cfg(unix)]
1646fn fix_state_dir_permissions() {
1647    let state_dir = &*env::PITCHFORK_STATE_DIR;
1648    if let Some((uid, gid)) = state_owner_ids() {
1649        if !state_dir.exists()
1650            && let Err(err) = fs::create_dir_all(state_dir)
1651        {
1652            warn!(
1653                "failed to create state directory for ownership fix at {}: {err}",
1654                state_dir.display()
1655            );
1656            return;
1657        }
1658
1659        // Best path: chown back to the runtime user. Permissions stay tight.
1660        chown_recursive(state_dir, uid, gid, true);
1661        debug!(
1662            "chowned state directory to uid={uid} gid={gid} at {}",
1663            state_dir.display()
1664        );
1665    } else {
1666        if !state_dir.exists() {
1667            return;
1668        }
1669
1670        // Fallback: relax permissions on safe subdirectories only.
1671        // proxy/ is never touched.
1672        chmod_safe_subtrees(state_dir);
1673        debug!(
1674            "relaxed permissions on safe subtrees at {}",
1675            state_dir.display()
1676        );
1677    }
1678}
1679
1680#[cfg(unix)]
1681pub(crate) fn state_owner_ids() -> Option<(u32, u32)> {
1682    if !nix::unistd::Uid::effective().is_root() {
1683        return None;
1684    }
1685
1686    let s = settings();
1687    let user = s.supervisor.user.trim();
1688    if !user.is_empty() {
1689        return resolve_supervisor_user_ids(user).or_else(|| {
1690            warn!(
1691                "failed to resolve supervisor.user '{user}' for state ownership; falling back to the invoking user"
1692            );
1693            env::invoking_user_ids()
1694        });
1695    }
1696
1697    env::invoking_user_ids()
1698}
1699
1700#[cfg(unix)]
1701fn resolve_supervisor_user_ids(user: &str) -> Option<(u32, u32)> {
1702    let user_record = if user.chars().all(|c| c.is_ascii_digit()) {
1703        let uid = user.parse::<u32>().ok()?;
1704        nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1705            .ok()
1706            .flatten()
1707    } else {
1708        nix::unistd::User::from_name(user).ok().flatten()
1709    }?;
1710
1711    Some((user_record.uid.as_raw(), user_record.gid.as_raw()))
1712}
1713
1714/// Recursively `chown` a directory tree. If `skip_proxy` is true, the `proxy/`
1715/// subdirectory is skipped entirely to protect the CA private key.
1716#[cfg(unix)]
1717fn chown_recursive(dir: &std::path::Path, uid: u32, gid: u32, skip_proxy: bool) {
1718    // chown the directory itself
1719    let _ = chown_path(dir, uid, gid);
1720
1721    let entries = match std::fs::read_dir(dir) {
1722        Ok(e) => e,
1723        Err(_) => return,
1724    };
1725    for entry in entries.flatten() {
1726        let path = entry.path();
1727        if path.is_dir() {
1728            // Skip proxy/ at the top level of the state directory
1729            if skip_proxy
1730                && let Some(name) = path.file_name().and_then(|n| n.to_str())
1731                && name == "proxy"
1732            {
1733                continue;
1734            }
1735            chown_recursive(&path, uid, gid, false);
1736        } else {
1737            let _ = chown_path(&path, uid, gid);
1738        }
1739    }
1740}
1741
1742/// `chown` a single path using libc. Returns Ok(()) on success.
1743#[cfg(unix)]
1744fn chown_path(path: &std::path::Path, uid: u32, gid: u32) -> std::io::Result<()> {
1745    use std::ffi::CString;
1746    use std::os::unix::ffi::OsStrExt;
1747    let c_path = CString::new(path.as_os_str().as_bytes())
1748        .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
1749    let ret = unsafe { libc::chown(c_path.as_ptr(), uid, gid) };
1750    if ret == 0 {
1751        Ok(())
1752    } else {
1753        Err(std::io::Error::last_os_error())
1754    }
1755}
1756
1757/// Fallback: relax permissions on safe subdirectories only (sock/, logs/, and
1758/// state.toml). The proxy/ subtree is never touched.
1759#[cfg(unix)]
1760fn chmod_safe_subtrees(state_dir: &std::path::Path) {
1761    // The state directory itself needs to be traversable
1762    let _ = fs::set_permissions(state_dir, fs::Permissions::from_mode(0o755));
1763
1764    // state.toml — needs to be readable by CLI clients
1765    let state_file = state_dir.join("state.toml");
1766    if state_file.exists() {
1767        let _ = fs::set_permissions(&state_file, fs::Permissions::from_mode(0o644));
1768    }
1769
1770    // Safe subdirectories: sock/ and logs/
1771    for subdir_name in &["sock", "logs"] {
1772        let subdir = state_dir.join(subdir_name);
1773        if subdir.is_dir() {
1774            chmod_recursive(&subdir);
1775        }
1776    }
1777}
1778
1779/// On startup, reconcile daemon processes left behind by a previous supervisor
1780/// that was terminated unexpectedly (e.g. `kill -9`).
1781///
1782/// This iterates the state file for daemon entries with a recorded PID. If the
1783/// PID is still alive and its current identity matches the recorded start time
1784/// (or the recorded title for older state files), it is assumed to be an orphan
1785/// from the previous supervisor session and `supervisor.orphan_policy` decides
1786/// its fate: `adopt` (default) resumes supervision via a poll monitor and keeps
1787/// the daemon's state intact; `kill` terminates it and resets its state to
1788/// `Stopped` with no PID. If a matching live process cannot be terminated
1789/// securely, its running state is retained to prevent a duplicate instance
1790/// from being started.
1791///
1792/// Missing or mismatched identity data fails closed so a PID recycled by an
1793/// unrelated process is never adopted or killed. On Unix platforms without
1794/// durable process handles, orphan termination also fails closed because the
1795/// PID/PGID cannot be pinned between identity validation and signaling.
1796///
1797/// This is gated by the `supervisor.cleanup_orphans` setting (default: true).
1798///
1799/// `unclean` is whether the previous supervisor exited uncleanly (see
1800/// [`supervisor_exited_uncleanly`]). It is read in `start()` before the
1801/// starting supervisor records itself in the state file — this function runs
1802/// in the background after that record exists, so reading it here would
1803/// always report unclean.
1804async fn cleanup_orphaned_daemons(supervisor: &Supervisor, unclean: bool) {
1805    if !settings().supervisor.cleanup_orphans {
1806        return;
1807    }
1808
1809    let candidates: Vec<_> = {
1810        let state = supervisor.state_file.lock().await;
1811        state
1812            .daemons
1813            .values()
1814            .filter(|d| d.id != DaemonId::pitchfork() && d.pid.is_some())
1815            .cloned()
1816            .collect()
1817    };
1818
1819    if candidates.is_empty() {
1820        return;
1821    }
1822
1823    info!(
1824        "checking {} daemon(s) for orphaned processes",
1825        candidates.len()
1826    );
1827
1828    let policy = orphan_policy();
1829    let boot_time = PROCS.boot_time();
1830
1831    // Reconcile orphans in parallel — a kill waits for the daemon's whole
1832    // process group to exit, bounded by that daemon's stop budget, so
1833    // sequential processing would make total cleanup time the sum of the
1834    // budgets.
1835    let tasks: Vec<_> = candidates
1836        .into_iter()
1837        .map(|daemon| {
1838            let policy = policy.clone();
1839            tokio::spawn(cleanup_orphaned_daemon(daemon, policy, boot_time, unclean))
1840        })
1841        .collect();
1842    for task in tasks {
1843        let _ = task.await;
1844    }
1845}
1846
1847/// Reconcile a single orphan candidate: adopt it, kill it, or reset its state,
1848/// per the policy and identity checks described on [`cleanup_orphaned_daemons`].
1849///
1850/// Holds the daemon's stop lock so a concurrent Run/Stop request for the same
1851/// daemon (cleanup runs in the background) serializes with the orphan kill,
1852/// and re-checks the recorded PID under the lock: if it changed, another path
1853/// already replaced or cleaned up this record and the snapshot is stale.
1854async fn cleanup_orphaned_daemon(
1855    daemon: crate::daemon::Daemon,
1856    policy: String,
1857    boot_time: u64,
1858    unclean: bool,
1859) {
1860    let supervisor: &Supervisor = &SUPERVISOR;
1861    let Some(pid) = daemon.pid else { return };
1862
1863    let lock = supervisor.stop_lock(&daemon.id).await;
1864    let _guard = lock.lock().await;
1865    let current_pid = {
1866        let state = supervisor.state_file.lock().await;
1867        state.daemons.get(&daemon.id).and_then(|d| d.pid)
1868    };
1869    if current_pid != Some(pid) {
1870        debug!(
1871            "orphan cleanup: daemon {} pid changed (recorded {pid}, now {current_pid:?}), skipping",
1872            daemon.id
1873        );
1874        return;
1875    }
1876
1877    // Refresh the candidate immediately before checking it: waiting for the
1878    // stop lock can await another path's stop timeout, during which this PID
1879    // may exit and be recycled.
1880    PROCS.refresh_pids(&[pid]);
1881
1882    if !PROCS.is_running(pid) {
1883        // PID already dead — the daemon exited while unsupervised, so
1884        // record a terminal status that reflects whether it died under a
1885        // crashed supervisor (retryable) or with the machine.
1886        let status = unobserved_exit_status(
1887            &daemon.status,
1888            daemon.boot_time,
1889            boot_time,
1890            unclean,
1891            daemon.oneshot,
1892        );
1893        reset_daemon_state(supervisor, &daemon.id, status, ExitObservation::Unobserved).await;
1894        return;
1895    }
1896
1897    // Safety check: verify the live process really is the daemon we
1898    // recorded, not an unrelated process that received a recycled PID.
1899    // The kernel start time is a stable identity for the lifetime of a
1900    // process, and is the only thing accepted as one.
1901    let current_start_time = PROCS.start_time(pid);
1902    let matches = process_identity_matches(daemon.start_time, current_start_time);
1903
1904    if !matches {
1905        // Either side missing means the identity cannot be checked at all,
1906        // which is different from checking it and finding a stranger: retain
1907        // the running state rather than resetting a record whose process may
1908        // well still be the daemon.
1909        if daemon.start_time.is_none() || current_start_time.is_none() {
1910            warn!(
1911                "could not verify the identity of live pid {pid} recorded for daemon {}; retaining running state",
1912                daemon.id,
1913            );
1914            return;
1915        }
1916        warn!(
1917            "pid {pid} recorded for daemon {} belongs to a different process now (PID recycled); resetting state without killing",
1918            daemon.id,
1919        );
1920        // The daemon died at some unknown point and the OS handed its PID
1921        // to something else — same unobserved exit as a dead PID.
1922        let status = unobserved_exit_status(
1923            &daemon.status,
1924            daemon.boot_time,
1925            boot_time,
1926            unclean,
1927            daemon.oneshot,
1928        );
1929        reset_daemon_state(supervisor, &daemon.id, status, ExitObservation::Unobserved).await;
1930        return;
1931    }
1932
1933    // Both policies need a verified start time: killing revalidates it
1934    // while pinned to the process, and adoption anchors its poll monitor
1935    // to it so a later PID recycle is never mistaken for the daemon.
1936    let Some(expected_start_time) = current_start_time else {
1937        warn!(
1938            "could not read start time for live pid {pid} recorded for daemon {}; retaining running state",
1939            daemon.id,
1940        );
1941        return;
1942    };
1943
1944    // Identity verified — the process really is our orphaned daemon.
1945    // The policy decides whether supervision resumes or the slate is
1946    // wiped clean.
1947    if policy == "adopt" {
1948        supervisor
1949            .adopt_daemon(&daemon, pid, expected_start_time)
1950            .await;
1951        return;
1952    }
1953
1954    info!("terminating orphaned daemon {} (pid {pid})", daemon.id);
1955
1956    let stop_cfg = daemon.stop_signal.unwrap_or_default();
1957    let termination_result = PROCS
1958        .kill_process_group_if_start_time_matches_async(
1959            pid,
1960            Some(expected_start_time),
1961            stop_cfg.signal.into(),
1962            stop_cfg.timeout,
1963        )
1964        .await;
1965
1966    match termination_result {
1967        Ok(true) => {}
1968        Ok(false) => {
1969            warn!(
1970                "could not securely terminate orphaned daemon {} (pid {pid}); retaining running state",
1971                daemon.id
1972            );
1973            return;
1974        }
1975        Err(err) => {
1976            warn!(
1977                "failed to terminate orphaned daemon {} (pid {pid}): {err}; retaining running state",
1978                daemon.id
1979            );
1980            return;
1981        }
1982    }
1983
1984    // We terminated the orphan ourselves, so this is an observed,
1985    // intentional stop rather than an unobserved exit.
1986    reset_daemon_state(
1987        supervisor,
1988        &daemon.id,
1989        DaemonStatus::Stopped,
1990        ExitObservation::Terminated,
1991    )
1992    .await;
1993}
1994
1995/// Effective `supervisor.orphan_policy`, warning on an unrecognized value
1996/// (which falls back to the default of adopting).
1997pub(crate) fn orphan_policy() -> String {
1998    let policy = settings().supervisor.orphan_policy.clone();
1999    match policy.as_str() {
2000        "adopt" | "kill" => policy,
2001        other => {
2002            warn!("unknown supervisor.orphan_policy '{other}', defaulting to 'adopt'");
2003            "adopt".to_string()
2004        }
2005    }
2006}
2007
2008/// Verify that live process identity matches the persisted daemon identity.
2009///
2010/// Both start times are required. A process name was once accepted in place of
2011/// a recorded start time, for state written before start times existed, but a
2012/// name is not an identity: a recycled PID belonging to another copy of the same
2013/// program matches it, and adopting or killing on that basis acts on the wrong
2014/// process. Missing identity, on either side, means unverifiable — and
2015/// unverifiable must never authorize acting on a process.
2016fn process_identity_matches(
2017    recorded_start_time: Option<u64>,
2018    current_start_time: Option<u64>,
2019) -> bool {
2020    match (recorded_start_time, current_start_time) {
2021        (Some(recorded), Some(current)) => recorded == current,
2022        _ => false,
2023    }
2024}
2025
2026/// Whether a PID read from persisted state may be signalled.
2027///
2028/// Stopping a daemon signals its whole process *group*, so acting on a PID that
2029/// has been recycled since it was recorded takes down an unrelated process tree.
2030/// Records are refused only when their identity is positively contradicted: if
2031/// either start time is unknown the PID stays as signallable as it was before
2032/// identities were recorded, so a daemon whose record predates the field can
2033/// still be stopped rather than becoming permanently unstoppable.
2034///
2035/// This is deliberately weaker than [`process_identity_matches`], which decides
2036/// whether to adopt or kill a process nobody asked about. Here the user has
2037/// named the daemon and asked for it to stop; the check exists to catch the
2038/// case where the answer is provably the wrong process.
2039pub(crate) fn signalling_pid_is_authorized(
2040    recorded_start_time: Option<u64>,
2041    current_start_time: Option<u64>,
2042) -> bool {
2043    !matches!(
2044        (recorded_start_time, current_start_time),
2045        (Some(recorded), Some(current)) if recorded != current
2046    )
2047}
2048
2049/// How a daemon's run ended, which decides what happens to the recorded
2050/// `last_exit_success` that cron `retrigger = "success" | "fail"` reads.
2051#[derive(Clone, Copy, PartialEq, Eq)]
2052pub(crate) enum ExitObservation {
2053    /// Nobody saw how the run ended, because the monitor that would have
2054    /// observed it died with a previous supervisor. The recorded outcome is
2055    /// cleared to `None`.
2056    ///
2057    /// Every option here is imperfect, so this picks the one that asserts
2058    /// nothing false. `Some(false)` would fabricate a failure, silently
2059    /// breaking a `retrigger = "success"` chain whose run may well have
2060    /// succeeded; `Some(true)` fabricates the opposite; keeping the previous
2061    /// value attributes an earlier run's outcome to this one. `None` says
2062    /// "unknown", reusing the reading the cron watcher already applies to a
2063    /// daemon that has never run.
2064    ///
2065    /// The tradeoff is that `None` satisfies both `retrigger = "success"`
2066    /// (`unwrap_or(true)`) and `retrigger = "fail"` (`!unwrap_or(false)`), so
2067    /// such a daemon fires once at its next scheduled time regardless of which
2068    /// it configured. That is schedule-gated rather than a loop, and it biases
2069    /// toward running the daemon over leaving it permanently untriggered.
2070    /// Distinguishing "unknown" from "never ran" would require a third cron
2071    /// state and is deliberately left out of scope here.
2072    Unobserved,
2073    /// We terminated the process ourselves, so the outcome is not a mystery:
2074    /// it stopped because we asked it to. Recorded as a success, matching the
2075    /// convention `Supervisor::stop` already uses for a deliberate stop.
2076    Terminated,
2077}
2078
2079impl ExitObservation {
2080    /// The `last_exit_success` value this observation implies.
2081    pub(crate) fn last_exit_success(self) -> Option<bool> {
2082        match self {
2083            ExitObservation::Unobserved => None,
2084            ExitObservation::Terminated => Some(true),
2085        }
2086    }
2087}
2088
2089/// Clear a daemon's runtime state (pid, process identity, active port) after
2090/// its process is gone or is no longer ours to manage.
2091///
2092/// Config fields are preserved by cloning the existing record, so a reset can
2093/// never drop a daemon's command, retry policy, or schedule.
2094async fn reset_daemon_state(
2095    supervisor: &Supervisor,
2096    id: &DaemonId,
2097    status: DaemonStatus,
2098    observation: ExitObservation,
2099) {
2100    let mut state_file = supervisor.state_file.lock().await;
2101    let Some(existing) = state_file.daemons.get(id) else {
2102        return;
2103    };
2104    let mut daemon = existing.clone();
2105    daemon.pid = None;
2106    daemon.title = None;
2107    daemon.start_time = None;
2108    daemon.boot_time = None;
2109    daemon.status = status;
2110    daemon.last_exit_success = observation.last_exit_success();
2111    daemon.active_port = None;
2112    state_file.clear_active_port(id);
2113    state_file.insert_daemon(id, daemon);
2114}
2115
2116/// Boot times this far apart are treated as different boots.
2117///
2118/// Sized to the only platform that reports a jittery value: Windows derives
2119/// boot time as `now - GetTickCount64()`, sampling two clocks independently,
2120/// so consecutive calls within one boot can differ by about a second. Linux
2121/// (`/proc/stat` btime) and macOS (`kern.boottime`) report stable values.
2122///
2123/// Deliberately kept this tight so a genuine reboot can never fall inside it:
2124/// a prior session would have to boot, start the supervisor, spawn a daemon,
2125/// have that daemon die, and complete a reboot inside two seconds, which no
2126/// real boot cycle reaches. A larger window would misread a short-lived
2127/// previous boot (e.g. a device in a reboot loop) as the current one and
2128/// resurrect daemons a reboot should have left stopped.
2129const BOOT_TIME_TOLERANCE_SECS: u64 = 2;
2130
2131/// Terminal status for a daemon whose process is gone and whose exit was
2132/// never observed, because the monitor that would have seen it died with a
2133/// previous supervisor.
2134///
2135/// A daemon recorded `Running` was expected to still be alive, so it died
2136/// under the crashed supervisor: `Errored(-1)` ("unknown exit code") makes it
2137/// eligible for its configured retries. Two cases stay `Stopped` instead:
2138///
2139/// - records from an earlier boot, whose processes died with the machine —
2140///   auto-restarting those is what `boot_start` is for, and reviving every
2141///   retry-configured daemon after a reboot would be a surprise
2142/// - any other status (in practice `Stopping`), i.e. an intentional stop that
2143///   completed while the supervisor was gone
2144pub(crate) fn unobserved_exit_status(
2145    status: &DaemonStatus,
2146    recorded_boot_time: Option<u64>,
2147    current_boot_time: u64,
2148    supervisor_exited_uncleanly: bool,
2149    oneshot: bool,
2150) -> DaemonStatus {
2151    let same_boot = recorded_boot_time
2152        .is_some_and(|recorded| recorded.abs_diff(current_boot_time) <= BOOT_TIME_TOLERANCE_SECS);
2153    // A task gets `stopped` rather than `errored` for the same reason the
2154    // adopted path does: `errored` is what `check_retry` looks for, and
2155    // nobody saw how this run ended, so retrying it would re-run a migration
2156    // or a seed that may well have succeeded. Leave re-running to an explicit
2157    // start.
2158    if status.is_running() && same_boot && supervisor_exited_uncleanly && !oneshot {
2159        DaemonStatus::Errored(-1)
2160    } else {
2161        DaemonStatus::Stopped
2162    }
2163}
2164
2165/// Whether the supervisor that owned this state file failed to shut down
2166/// cleanly, meaning any daemon it left behind stopped for reasons nobody
2167/// recorded.
2168///
2169/// A clean shutdown removes the supervisor's own entry: `close()` does it on
2170/// Unix, where the stop signal is delivered and handled, and the
2171/// `supervisor stop` command does it on Windows, which has no POSIX signals
2172/// and force-terminates the process instead. A crash, an external `kill -9`,
2173/// or a `--force` replacement all leave the entry behind.
2174///
2175/// This must be read before the starting supervisor records itself, which is
2176/// why `cleanup_orphaned_daemons` runs first in `start()`.
2177async fn supervisor_exited_uncleanly(supervisor: &Supervisor) -> bool {
2178    supervisor
2179        .state_file
2180        .lock()
2181        .await
2182        .daemons
2183        .contains_key(&DaemonId::pitchfork())
2184}
2185
2186/// Recursively chmod: directories → 0o755, files → 0o644.
2187#[cfg(unix)]
2188fn chmod_recursive(dir: &std::path::Path) {
2189    let _ = fs::set_permissions(dir, fs::Permissions::from_mode(0o755));
2190    let entries = match fs::read_dir(dir) {
2191        Ok(e) => e,
2192        Err(_) => return,
2193    };
2194    for entry in entries.flatten() {
2195        let path = entry.path();
2196        if path.is_dir() {
2197            chmod_recursive(&path);
2198        } else {
2199            let _ = fs::set_permissions(&path, fs::Permissions::from_mode(0o644));
2200        }
2201    }
2202}
2203
2204#[cfg(test)]
2205mod tests {
2206    use super::{
2207        BOOT_TIME_TOLERANCE_SECS, legacy_supervisor_title_matches, process_identity_matches,
2208        should_remove_liveness_session, signalling_pid_is_authorized, supervisor_identity_matches,
2209        unobserved_exit_status,
2210    };
2211    use crate::daemon_status::DaemonStatus;
2212    use crate::state_file::ProjectSession;
2213
2214    const BOOT: u64 = 1_700_000_000;
2215
2216    #[test]
2217    fn unobserved_running_death_in_current_boot_is_retryable() {
2218        // Died under a crashed supervisor during this boot: Errored(-1) makes
2219        // the daemon eligible for its configured retries.
2220        assert!(matches!(
2221            unobserved_exit_status(&DaemonStatus::Running, Some(BOOT), BOOT, true, false),
2222            DaemonStatus::Errored(-1)
2223        ));
2224    }
2225
2226    #[test]
2227    fn unobserved_task_death_is_stopped_not_retried() {
2228        // Nobody saw how the run ended, so `errored` would hand a migration or
2229        // a seed to check_retry on a guess. Matches what the adopted path
2230        // records, and what the guide promises.
2231        assert!(matches!(
2232            unobserved_exit_status(&DaemonStatus::Running, Some(BOOT), BOOT, true, true),
2233            DaemonStatus::Stopped
2234        ));
2235    }
2236
2237    #[test]
2238    fn unobserved_running_death_from_previous_boot_is_stopped() {
2239        // The process died with the machine; reviving every retry-configured
2240        // daemon after a reboot is what boot_start is for.
2241        assert!(matches!(
2242            unobserved_exit_status(
2243                &DaemonStatus::Running,
2244                Some(BOOT - 86_400),
2245                BOOT,
2246                true,
2247                false
2248            ),
2249            DaemonStatus::Stopped
2250        ));
2251    }
2252
2253    #[test]
2254    fn unobserved_exit_tolerates_boot_time_jitter() {
2255        // Windows recomputes boot time as now - GetTickCount64(), which can
2256        // drift about a second between samples within one boot.
2257        let within = BOOT + BOOT_TIME_TOLERANCE_SECS;
2258        assert!(matches!(
2259            unobserved_exit_status(&DaemonStatus::Running, Some(within), BOOT, true, false),
2260            DaemonStatus::Errored(-1)
2261        ));
2262        let beyond = BOOT + BOOT_TIME_TOLERANCE_SECS + 1;
2263        assert!(matches!(
2264            unobserved_exit_status(&DaemonStatus::Running, Some(beyond), BOOT, true, false),
2265            DaemonStatus::Stopped
2266        ));
2267    }
2268
2269    #[test]
2270    fn unobserved_exit_after_clean_shutdown_is_stopped() {
2271        // A deliberate `supervisor stop` can leave running records behind on
2272        // platforms where the supervisor cannot handle the stop signal. Those
2273        // daemons were stopped on purpose, so they must not be reported as
2274        // failures or resurrected by the retry checker.
2275        assert!(matches!(
2276            unobserved_exit_status(&DaemonStatus::Running, Some(BOOT), BOOT, false, false),
2277            DaemonStatus::Stopped
2278        ));
2279    }
2280
2281    #[test]
2282    fn unobserved_exit_treats_short_previous_boot_as_previous() {
2283        // A device in a reboot loop can produce consecutive boots seconds
2284        // apart. The jitter window must stay far below that so those records
2285        // are still recognised as belonging to an earlier boot.
2286        for gap in [5, 30, 59, 60] {
2287            assert!(
2288                matches!(
2289                    unobserved_exit_status(
2290                        &DaemonStatus::Running,
2291                        Some(BOOT - gap),
2292                        BOOT,
2293                        true,
2294                        false
2295                    ),
2296                    DaemonStatus::Stopped
2297                ),
2298                "boot {gap}s earlier should be treated as a previous boot"
2299            );
2300        }
2301    }
2302
2303    #[test]
2304    fn unobserved_exit_without_recorded_boot_time_is_stopped() {
2305        // Legacy state files predating the field fail closed to today's
2306        // behavior rather than triggering surprise retries.
2307        assert!(matches!(
2308            unobserved_exit_status(&DaemonStatus::Running, None, BOOT, true, false),
2309            DaemonStatus::Stopped
2310        ));
2311    }
2312
2313    #[test]
2314    fn supervisor_identity_matches_same_generation() {
2315        assert!(supervisor_identity_matches(
2316            Some(100),
2317            Some(100),
2318            Some(BOOT),
2319            BOOT
2320        ));
2321    }
2322
2323    #[test]
2324    fn supervisor_identity_survives_clock_steps_when_tokens_match() {
2325        // An NTP step or sleep/resume moves the realtime-derived boot time by
2326        // far more than the tolerance while the supervisor keeps running.
2327        // Matching start tokens prove it is the same process; declaring it
2328        // stale here would start a second supervisor.
2329        assert!(supervisor_identity_matches(
2330            Some(100),
2331            Some(100),
2332            Some(BOOT),
2333            BOOT + 3600
2334        ));
2335        assert!(supervisor_identity_matches(
2336            Some(100),
2337            Some(100),
2338            Some(BOOT + 3600),
2339            BOOT
2340        ));
2341    }
2342
2343    #[test]
2344    fn supervisor_identity_rejects_recycled_pid() {
2345        // The supervisor died and something else got its PID, within this
2346        // boot or across a reboot: the token differs either way.
2347        assert!(!supervisor_identity_matches(
2348            Some(100),
2349            Some(200),
2350            Some(BOOT),
2351            BOOT
2352        ));
2353        assert!(!supervisor_identity_matches(
2354            Some(100),
2355            Some(200),
2356            Some(BOOT),
2357            BOOT + 3600
2358        ));
2359    }
2360
2361    #[test]
2362    fn supervisor_identity_uses_boot_time_when_a_token_is_missing() {
2363        // Discussion #877: the record survived a reboot and an early system
2364        // daemon now owns the PID. With no live token to compare, the boot
2365        // time is what proves the record stale.
2366        assert!(!supervisor_identity_matches(
2367            Some(100),
2368            None,
2369            Some(BOOT),
2370            BOOT + 3600
2371        ));
2372        // Same boot (within Windows boot-time jitter) and no contradiction.
2373        assert!(supervisor_identity_matches(
2374            Some(100),
2375            None,
2376            Some(BOOT),
2377            BOOT + BOOT_TIME_TOLERANCE_SECS
2378        ));
2379    }
2380
2381    #[test]
2382    fn supervisor_identity_tolerates_legacy_records() {
2383        // A record written before either field existed is not contradicted by
2384        // anything here; `supervisor_record_is_live` applies the process-name
2385        // check to those instead.
2386        assert!(supervisor_identity_matches(None, Some(100), None, BOOT));
2387        // A record with only a boot time is still rejected across a reboot.
2388        assert!(!supervisor_identity_matches(
2389            None,
2390            Some(100),
2391            Some(BOOT),
2392            BOOT + 3600
2393        ));
2394        assert!(supervisor_identity_matches(
2395            None,
2396            Some(100),
2397            Some(BOOT),
2398            BOOT
2399        ));
2400    }
2401
2402    #[test]
2403    fn legacy_supervisor_title_requires_a_pitchfork_process() {
2404        assert!(legacy_supervisor_title_matches(Some("pitchfork")));
2405        assert!(legacy_supervisor_title_matches(Some("pitchfork.exe")));
2406        assert!(legacy_supervisor_title_matches(Some("Pitchfork")));
2407        // Discussion #877: an Apple LaunchAgent inherited the PID after a reboot.
2408        assert!(!legacy_supervisor_title_matches(Some(
2409            "AMPDeviceDiscoveryAgent"
2410        )));
2411        assert!(!legacy_supervisor_title_matches(Some("sleep")));
2412        assert!(!legacy_supervisor_title_matches(None));
2413    }
2414
2415    #[test]
2416    fn unobserved_exit_of_stopping_daemon_is_stopped() {
2417        // An intentional stop that completed while the supervisor was gone is
2418        // not a failure, even within the same boot.
2419        assert!(matches!(
2420            unobserved_exit_status(&DaemonStatus::Stopping, Some(BOOT), BOOT, true, false),
2421            DaemonStatus::Stopped
2422        ));
2423    }
2424
2425    #[test]
2426    fn orphan_identity_requires_both_start_times() {
2427        assert!(process_identity_matches(Some(123), Some(123)));
2428        assert!(!process_identity_matches(Some(123), Some(456)));
2429        // Unreadable current identity: unverifiable, so not a match.
2430        assert!(!process_identity_matches(Some(123), None));
2431    }
2432
2433    #[test]
2434    fn signalling_is_refused_only_for_a_contradicted_identity() {
2435        // Provably someone else's process group: refuse.
2436        assert!(!signalling_pid_is_authorized(Some(123), Some(456)));
2437        // Verified as the daemon's own.
2438        assert!(signalling_pid_is_authorized(Some(123), Some(123)));
2439        // Unknown on either side. Stopping stays possible, because the user has
2440        // named this daemon and a record that cannot be verified must not become
2441        // one that can never be stopped.
2442        assert!(signalling_pid_is_authorized(None, Some(123)));
2443        assert!(signalling_pid_is_authorized(Some(123), None));
2444        assert!(signalling_pid_is_authorized(None, None));
2445    }
2446
2447    #[test]
2448    fn orphan_identity_rejects_records_without_a_start_time() {
2449        // State written before start times were recorded. A process name used
2450        // to stand in here, but another copy of the same program on a recycled
2451        // PID matches a name, so such records are no longer verifiable and must
2452        // not authorize adopting or killing anything.
2453        assert!(!process_identity_matches(None, Some(123)));
2454        assert!(!process_identity_matches(None, None));
2455    }
2456
2457    #[test]
2458    fn should_not_remove_when_state_title_differs_from_snapshot() {
2459        // The session was re-entered after the snapshot was taken, producing a
2460        // new title in state. The snapshot title is stale; skip removal.
2461        let session = ProjectSession {
2462            liveness_title: Some("new_title".to_string()),
2463        };
2464        let recorded_title = Some("old_title".to_string());
2465
2466        assert!(!should_remove_liveness_session(
2467            &session,
2468            &recorded_title,
2469            Some("new_title"),
2470            true,
2471        ));
2472    }
2473
2474    #[test]
2475    fn should_remove_when_running_title_mismatches() {
2476        let session = ProjectSession {
2477            liveness_title: Some("recorded_title".to_string()),
2478        };
2479        let recorded_title = Some("recorded_title".to_string());
2480
2481        assert!(should_remove_liveness_session(
2482            &session,
2483            &recorded_title,
2484            Some("different_title"),
2485            true,
2486        ));
2487    }
2488
2489    #[test]
2490    fn should_remove_when_dead() {
2491        let session = ProjectSession {
2492            liveness_title: Some("recorded_title".to_string()),
2493        };
2494        let recorded_title = Some("recorded_title".to_string());
2495
2496        assert!(should_remove_liveness_session(
2497            &session,
2498            &recorded_title,
2499            Some("recorded_title"),
2500            false,
2501        ));
2502    }
2503
2504    #[test]
2505    fn should_not_remove_when_alive_and_title_matches() {
2506        let session = ProjectSession {
2507            liveness_title: Some("recorded_title".to_string()),
2508        };
2509        let recorded_title = Some("recorded_title".to_string());
2510
2511        assert!(!should_remove_liveness_session(
2512            &session,
2513            &recorded_title,
2514            Some("recorded_title"),
2515            true,
2516        ));
2517    }
2518
2519    #[cfg(unix)]
2520    #[test]
2521    fn fd_scan_end_reaches_open_fds_above_the_fd_limit() {
2522        // A caller opened fd 900, then lowered RLIMIT_NOFILE to 256.
2523        assert_eq!(super::fd_scan_end(256, [0, 1, 2, 900]), 901);
2524        assert_eq!(super::fd_scan_end(256, [0, 1, 2, 10]), 256);
2525        assert_eq!(super::fd_scan_end(-1, []), 1 << 16);
2526        assert_eq!(super::fd_scan_end(libc::c_long::MAX, []), 1 << 20);
2527    }
2528
2529    /// Exercises the fcntl fallback directly: on Linux >= 5.11 the spawn path
2530    /// never reaches it, but it is the only implementation on macOS.
2531    #[cfg(unix)]
2532    #[test]
2533    fn cloexec_fd_scan_hides_inherited_fds_from_children() {
2534        use std::os::unix::process::CommandExt;
2535        use std::process::{Command, Stdio};
2536
2537        // A pipe without O_CLOEXEC, like bats' fd 3 or a wrapper script's pipe.
2538        let mut fds = [0; 2];
2539        assert_eq!(unsafe { libc::pipe(fds.as_mut_ptr()) }, 0);
2540        let fd = fds[1];
2541        let child_sees_fd = |scan: bool| {
2542            let mut cmd = Command::new("sh");
2543            cmd.arg("-c")
2544                .arg(format!("[ -e /dev/fd/{fd} ]"))
2545                .stdin(Stdio::null())
2546                .stdout(Stdio::null())
2547                .stderr(Stdio::null());
2548            if scan {
2549                let max_fd = super::max_inherited_fd();
2550                unsafe {
2551                    cmd.pre_exec(move || {
2552                        super::cloexec_fd_scan(max_fd);
2553                        Ok(())
2554                    });
2555                }
2556            }
2557            cmd.status().unwrap().success()
2558        };
2559        let without_scan = child_sees_fd(false);
2560        let with_scan = child_sees_fd(true);
2561        unsafe {
2562            libc::close(fds[0]);
2563            libc::close(fds[1]);
2564        }
2565        assert!(without_scan, "control: child should inherit fd {fd}");
2566        assert!(!with_scan, "fd {fd} leaked past cloexec_fd_scan");
2567    }
2568}