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        // Ignoring Ctrl+C is inherited, so a supervisor started from a process
628        // that ignores it would neither see Ctrl+C itself nor let its daemons
629        // see it: a daemon with `stop_signal = "SIGINT"` would never get the
630        // Ctrl+C sent to stop it. Handle it again before any daemon starts.
631        #[cfg(windows)]
632        if unsafe { windows_sys::Win32::System::Console::SetConsoleCtrlHandler(None, 0) } == 0 {
633            warn!(
634                "failed to stop ignoring Ctrl+C: {}",
635                std::io::Error::last_os_error()
636            );
637        }
638
639        // Refuse to run beside a supervisor that is already listening, and
640        // do so before recording ourselves in the state file or starting any
641        // daemons: taking over its socket would leave it running but
642        // unreachable. (`--force` has already waited for the one it replaced
643        // to let go of the socket.) The lock is held until our own listener
644        // is bound, so a supervisor starting at the same time waits here and
645        // then finds this one listening.
646        let startup_lock = StartupLock::acquire().await?;
647        if crate::ipc::supervisor_listening().await {
648            return Err(miette::miette!(
649                "another pitchfork supervisor is already listening on {}",
650                crate::ipc::socket_display()
651            ));
652        }
653
654        let pid = std::process::id();
655        // Ensure PROCS has data for the supervisor PID before upsert_daemon reads title()
656        PROCS.refresh_pids(&[pid]);
657        // Determine container mode: CLI flag takes priority, then settings.
658        // Running as PID 1 always enables it: orphaned descendants of daemons
659        // re-parent to us, and without the zombie reaper they would accumulate
660        // as unreaped zombies — which also keep their process group alive,
661        // stalling whole-group stop waits indefinitely.
662        let container_mode =
663            container || settings().supervisor.container || std::process::id() == 1;
664        if container_mode {
665            info!("Starting supervisor in container/PID1 mode with pid {pid}");
666        } else {
667            info!("Starting supervisor with pid {pid}");
668        }
669
670        // Whether the previous supervisor exited uncleanly must be read before
671        // we record ourselves in the state file just below (see
672        // `supervisor_exited_uncleanly`); the background cleanup task runs
673        // after that record exists, so it receives the answer instead of
674        // reading it too late.
675        let unclean = supervisor_exited_uncleanly(self).await;
676
677        self.upsert_daemon(
678            UpsertDaemonOpts::builder(DaemonId::pitchfork())
679                .set(|o| {
680                    o.pid = Some(pid);
681                    o.status = DaemonStatus::Running;
682                })
683                .build(),
684        )
685        .await?;
686        #[cfg(unix)]
687        fix_state_dir_permissions();
688
689        // Self-heal: if the boot registration points to a stale binary path
690        // (e.g. after a brew/mise upgrade), re-register with the current path.
691        // Runs in the background — must not block or fail supervisor startup.
692        tokio::task::spawn_blocking(|| {
693            if let Ok(boot_manager) = crate::boot_manager::BootManager::new() {
694                boot_manager.check_and_reregister_if_stale();
695            }
696        });
697
698        // If the previous supervisor died uncleanly, its daemon child processes
699        // may still be alive (orphaned, re-parented to init).  Terminate them
700        // before starting replacements so we don't end up with duplicate
701        // processes holding the same ports.
702        //
703        // This runs in the background: each orphan kill now waits for its
704        // whole process group to exit (seconds per orphan), and doing that
705        // inline would delay IPC socket creation past the CLI's short connect
706        // budget on autostart. Per-daemon stop locks serialize the cleanup
707        // against any Run/Stop requests that arrive for the same daemon in the
708        // meantime, and boot daemons start after cleanup completes so they
709        // cannot observe an orphan as "already running".
710        let boot_after_cleanup = is_boot;
711        tokio::spawn(async move {
712            cleanup_orphaned_daemons(&SUPERVISOR, unclean).await;
713            if boot_after_cleanup {
714                info!("Boot start mode enabled, starting boot_start daemons");
715                if let Err(e) = SUPERVISOR.start_boot_daemons().await {
716                    error!("failed to start boot daemons: {e}");
717                }
718            }
719        });
720
721        self.interval_watch()?;
722
723        // Run the first cron check synchronously before starting the cron
724        // watcher and IPC server. This registers config-only cron daemons and
725        // fires any `immediate=true` triggers in the foreground, so they cannot
726        // race with a concurrent `pitchfork start` IPC. By the time the cron
727        // watcher's first tick runs, `last_cron_triggered` is already anchored
728        // and the immediate daemons are already running.
729        if let Err(e) = self.check_cron_schedules().await {
730            error!("failed to check cron schedules on startup: {e}");
731        }
732
733        self.cron_watch()?;
734        self.signals()?;
735        self.daemon_file_watch()?;
736
737        // In container mode, install SIGCHLD handler to reap orphaned/zombie processes
738        #[cfg(unix)]
739        if container_mode {
740            self.reap_zombies()?;
741        }
742
743        // Start web server: CLI --web-port takes priority, then settings.web.auto_start + bind_port
744        let s = settings();
745        let effective_port = web_port.or_else(|| {
746            if s.web.auto_start {
747                match u16::try_from(s.web.bind_port).ok().filter(|&p| p > 0) {
748                    Some(p) => Some(p),
749                    None => {
750                        error!(
751                            "web.bind_port {} is out of valid port range (1-65535), web UI disabled",
752                            s.web.bind_port
753                        );
754                        None
755                    }
756                }
757            } else {
758                None
759            }
760        });
761        // CLI --web-path takes priority, then settings.web.base_path
762        let effective_path = web_path.or_else(|| {
763            let bp = s.web.base_path.clone();
764            if bp.is_empty() { None } else { Some(bp) }
765        });
766        if let Some(port) = effective_port {
767            tokio::spawn(async move {
768                if let Err(e) = crate::web::serve(port, effective_path).await {
769                    error!("Web server error: {e}");
770                }
771            });
772        }
773
774        // Start standalone API server if configured
775        let api_port = if s.api.auto_start {
776            match u16::try_from(s.api.bind_port).ok().filter(|&p| p > 0) {
777                Some(p) => Some(p),
778                None => {
779                    error!(
780                        "api.bind_port {} is out of valid port range (1-65535), API server disabled",
781                        s.api.bind_port
782                    );
783                    None
784                }
785            }
786        } else {
787            None
788        };
789        if let Some(port) = api_port {
790            tokio::spawn(async move {
791                if let Err(e) = crate::web::serve_api(port, None).await {
792                    error!("API server error: {e}");
793                }
794            });
795        }
796
797        // Start reverse proxy server if enabled
798        if s.proxy.enable {
799            // Pre-generate the TLS certificate synchronously before spawning the proxy
800            // task. This ensures the cert exists immediately after `sup start` returns,
801            // so `proxy trust` can be run right away without waiting for the async task.
802            #[cfg(feature = "proxy-tls")]
803            if s.proxy.https {
804                let proxy_dir = crate::env::PITCHFORK_STATE_DIR.join("proxy");
805                let ca_cert_path = proxy_dir.join("ca.pem");
806                let ca_key_path = proxy_dir.join("ca-key.pem");
807                // Checked and written under the CA lock: `proxy setup` may be
808                // generating the same pair right now.
809                match crate::proxy::server::ensure_ca(&ca_cert_path, &ca_key_path, || {
810                    ca_cert_path.exists() && ca_key_path.exists()
811                }) {
812                    Ok(true) => {
813                        info!(
814                            "Generated local CA certificate at {}",
815                            ca_cert_path.display()
816                        );
817                    }
818                    Ok(false) => {}
819                    Err(e) => {
820                        error!("Failed to generate CA certificate: {e}");
821                    }
822                }
823
824                // Auto-trust: attempt to install the CA certificate into the
825                // system trust store. May fail silently due to permissions;
826                // user can run `pitchfork proxy trust` manually.
827                if s.proxy.auto_trust && ca_cert_path.exists() {
828                    use crate::proxy::trust::{AutoTrustResult, auto_trust};
829                    match auto_trust(&ca_cert_path) {
830                        AutoTrustResult::AlreadyTrusted => {}
831                        AutoTrustResult::Trusted => {
832                            info!("CA certificate auto-trusted in system store");
833                        }
834                        AutoTrustResult::NotTrusted { reason } => {
835                            warn!("Auto-trust skipped: {reason}");
836                            warn!("Run `pitchfork proxy trust` to install manually");
837                        }
838                    }
839                }
840            }
841            // Spawn the proxy server and wait for its bind result via a oneshot
842            // channel.  This avoids the TOCTOU race of a pre-flight bind check
843            // while still surfacing binding failures immediately.
844            let (bind_tx, bind_rx) = tokio::sync::oneshot::channel();
845            let proxy_cancel = tokio_util::sync::CancellationToken::new();
846            let proxy_cancel_clone = proxy_cancel.clone();
847            *self.proxy_cancel.lock().await = Some(proxy_cancel);
848            let proxy_task = tokio::spawn(async move {
849                if let Err(e) = crate::proxy::server::serve(bind_tx, proxy_cancel_clone).await {
850                    error!("Proxy server error: {e}");
851                }
852            });
853            *self.proxy_task.lock().await = Some(proxy_task);
854            match bind_rx.await {
855                Ok(Ok(())) => {
856                    info!("Proxy server bound successfully");
857                    // Resolved once and shared: mDNS and the DNS resolver
858                    // must advertise the same address.
859                    let lan_ip = self.resolve_lan_ip().await;
860                    self.start_mdns(lan_ip).await;
861                    self.start_dns_resolver(lan_ip).await;
862                }
863                Ok(Err(msg)) => {
864                    error!("{msg}");
865                    self.add_notification(log::LevelFilter::Error, msg).await;
866                }
867                Err(_) => {
868                    // Sender dropped without sending — serve() panicked or
869                    // returned before signalling.  Already logged by the
870                    // spawn error handler above.
871                }
872            }
873        }
874
875        // Pre-warm slug cache so the first /api/proxies request is fast.
876        // Spawned as a background task so it does not block startup.
877        tokio::spawn(async {
878            crate::proxy::server::get_cached_slugs().await;
879        });
880
881        let (ipc, ipc_handle) = IpcServer::new(startup_lock).await?;
882        *self.ipc_shutdown.lock().await = Some(ipc_handle);
883        self.start_state_flush_task();
884        self.conn_watch(ipc).await
885    }
886
887    /// Start the loopback DNS resolver for the proxy TLD.
888    ///
889    /// The resolver shares the proxy's cancellation token, so it stops with the
890    /// proxy. A bind failure is a notification rather than a fatal error: the
891    /// proxy still works for anyone who reaches it some other way.
892    async fn start_dns_resolver(&self, lan_ip: Option<std::net::Ipv4Addr>) {
893        let s = crate::settings::settings();
894        if !s.proxy.dns {
895            return;
896        }
897        let cfg = crate::proxy::dns::config_from_settings(&s, lan_ip);
898        let port = crate::proxy::dns::dns_port(&s);
899        // Loopback only: the resolver is for this machine's stub resolver, and
900        // LAN peers are served by mDNS instead.
901        let addr = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, port));
902
903        // The token's guard is held until the handle is stored. `close` takes
904        // the token under this lock before it collects `dns_task`, so it either
905        // runs first — and there is no token to spawn with — or waits until
906        // the handle is in place to be drained. Releasing it earlier left a
907        // window where `close` found no handle and the task outlived shutdown
908        // holding the resolver's sockets.
909        let cancel_guard = self.proxy_cancel.lock().await;
910        let Some(cancel) = cancel_guard.clone() else {
911            return;
912        };
913        let (bind_tx, bind_rx) = tokio::sync::oneshot::channel();
914        let task = tokio::spawn(async move {
915            if let Err(e) = crate::proxy::dns::serve(cfg, addr, bind_tx, cancel).await {
916                error!("DNS resolver error: {e}");
917            }
918        });
919        *self.dns_task.lock().await = Some(task);
920        drop(cancel_guard);
921        match bind_rx.await {
922            Ok(Ok(())) => info!("DNS resolver bound successfully"),
923            Ok(Err(msg)) => {
924                let msg = format!(
925                    "{msg}\nProxy host names will not resolve through pitchfork. \
926                     Choose another port with proxy.dns_port, or set proxy.dns = false."
927                );
928                error!("{msg}");
929                self.add_notification(log::LevelFilter::Error, msg).await;
930            }
931            Err(_) => {}
932        }
933    }
934
935    /// Watch the LAN address and keep the DNS responder, and mDNS when it is
936    /// running, pointed at the current one.
937    ///
938    /// Started even when the mDNS publisher could not be created: the DNS
939    /// responder serves this machine regardless, and an address it keeps
940    /// answering with after the interface has moved is worse than useless.
941    ///
942    /// Does nothing when `proxy.lan_ip` pins an address. That is a choice to
943    /// respect, not a starting point to drift from — the check lives here so
944    /// neither caller can forget it.
945    async fn start_lan_ip_monitor(
946        &self,
947        initial_ip: std::net::Ipv4Addr,
948        port: u16,
949        publisher: Option<std::sync::Arc<tokio::sync::Mutex<crate::proxy::mdns::MdnsPublisher>>>,
950    ) {
951        if !crate::settings::settings().proxy.lan_ip.is_empty() {
952            return;
953        }
954        // No cancellation token means `close` has already taken it, so shutdown
955        // is under way. Spawning here would leave a task polling with no way to
956        // stop it, and past the point where `close` collects the handle.
957        //
958        // Held until the handle is stored, for the reason `start_dns_resolver`
959        // gives: otherwise `close` can collect `lan_monitor_task` in between.
960        let cancel_guard = self.proxy_cancel.lock().await;
961        let Some(cancel) = cancel_guard.clone() else {
962            debug!("Not starting the LAN IP monitor: the supervisor is shutting down");
963            return;
964        };
965        let monitor_cancel = Some(cancel);
966        let task = tokio::spawn(async move {
967            let mut last_ip = initial_ip;
968            let mut ticker = tokio::time::interval(std::time::Duration::from_secs(5));
969            ticker.tick().await; // first tick is immediate
970            loop {
971                // Cancellation is raced against the tick, not checked after
972                // it. Checking afterwards means the task only notices once the
973                // full interval has elapsed, so a shutdown that waits a second
974                // for it always gives up and aborts instead — the graceful
975                // path would never once be taken.
976                match monitor_cancel.as_ref() {
977                    Some(cancel) => {
978                        tokio::select! {
979                            _ = ticker.tick() => {}
980                            _ = cancel.cancelled() => break,
981                        }
982                    }
983                    None => {
984                        ticker.tick().await;
985                    }
986                }
987                if let Some(new_ip) = crate::proxy::lan_ip::detect_lan_ip_if_changed(last_ip).await
988                {
989                    log::info!("LAN IP changed: {last_ip} → {new_ip}");
990                    last_ip = new_ip;
991                    crate::proxy::dns::update_lan_ip(new_ip);
992                    if let Some(publisher) = publisher.as_ref() {
993                        publisher.lock().await.republish_all(new_ip, port);
994                    }
995                }
996            }
997        });
998        *self.lan_monitor_task.lock().await = Some(task);
999        drop(cancel_guard);
1000    }
1001
1002    /// Start mDNS publishing for LAN mode (called after the proxy binds successfully).
1003    /// The LAN address mDNS publishes and the DNS resolver answers with.
1004    ///
1005    /// Resolved once and handed to both. Detecting separately in each let them
1006    /// disagree when the interface address changed in between, which would
1007    /// advertise one address over mDNS and serve another over DNS, and probed
1008    /// the network twice at startup for one answer.
1009    ///
1010    /// `None` means LAN mode is off, or is on and the address could not be
1011    /// determined; either way the caller has nothing to publish. The reason is
1012    /// reported here so it is said once rather than by each caller.
1013    async fn resolve_lan_ip(&self) -> Option<std::net::Ipv4Addr> {
1014        let s = crate::settings::settings();
1015        let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
1016        if !s.proxy.enable || !lan_enabled {
1017            return None;
1018        }
1019        if s.proxy.lan_ip.is_empty() {
1020            let detected = crate::proxy::lan_ip::detect_lan_ip().await;
1021            if detected.is_none() {
1022                error!(
1023                    "LAN mode is enabled but no LAN IP address could be detected. \
1024                     Set proxy.lan_ip to a specific address, or ensure you are connected to a network."
1025                );
1026            }
1027            return detected;
1028        }
1029        match s.proxy.lan_ip.parse::<std::net::Ipv4Addr>() {
1030            Ok(ip) => Some(ip),
1031            Err(e) => {
1032                let msg = format!(
1033                    concat!(
1034                        "proxy.lan_ip {:?} is not a valid IPv4 address: {}. ",
1035                        "LAN mode will not start; fix the setting or clear it ",
1036                        "to auto-detect."
1037                    ),
1038                    s.proxy.lan_ip, e
1039                );
1040                error!("{msg}");
1041                self.add_notification(log::LevelFilter::Error, msg).await;
1042                None
1043            }
1044        }
1045    }
1046
1047    async fn start_mdns(&self, lan_ip: Option<std::net::Ipv4Addr>) {
1048        let s = crate::settings::settings();
1049        let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
1050        if !s.proxy.enable || !lan_enabled {
1051            return;
1052        }
1053
1054        let Some(lan_ip) = lan_ip else { return };
1055        let port = u16::try_from(s.proxy.port).unwrap_or(443);
1056
1057        let Some(mut publisher) = crate::proxy::mdns::MdnsPublisher::new(lan_ip) else {
1058            error!("Failed to start mDNS publisher. Is Avahi (Linux) or Bonjour (macOS) running?");
1059            // The DNS responder hands out this address too, and it is useful on
1060            // this machine whether or not mDNS came up. Keep watching the
1061            // interface so its answers do not go stale.
1062            self.start_lan_ip_monitor(lan_ip, port, None).await;
1063            return;
1064        };
1065
1066        // Publish all registered slugs.
1067        let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
1068        for slug in slugs.keys() {
1069            let hostname = format!("{slug}.local");
1070            publisher.publish(&hostname, port);
1071        }
1072
1073        log::info!(
1074            "LAN mode: mDNS publishing on {lan_ip}, {} slug(s) registered",
1075            slugs.len()
1076        );
1077
1078        let publisher = std::sync::Arc::new(tokio::sync::Mutex::new(publisher));
1079
1080        // Start the IP monitor. It declines on its own when the address is
1081        // pinned rather than auto-detected.
1082        self.start_lan_ip_monitor(lan_ip, port, Some(publisher.clone()))
1083            .await;
1084
1085        *self.mdns_publisher.lock().await = Some(publisher);
1086    }
1087
1088    /// Re-read slugs from config and update mDNS records.
1089    ///
1090    /// Publishes new slugs and unpublishes removed ones. Called via IPC when
1091    /// `proxy add` or `proxy remove` modifies the slug registry.
1092    async fn sync_mdns(&self) {
1093        // Clone the Arc and release the outer lock immediately so we don't
1094        // block close() from taking the publisher during shutdown.
1095        let publisher = {
1096            let guard = self.mdns_publisher.lock().await;
1097            match guard.as_ref() {
1098                Some(p) => p.clone(),
1099                None => {
1100                    debug!("sync_mdns: mDNS publisher not active, skipping");
1101                    return;
1102                }
1103            }
1104        };
1105
1106        let s = crate::settings::settings();
1107        let port = u16::try_from(s.proxy.port).unwrap_or(443);
1108
1109        let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
1110        let mut pub_guard = publisher.lock().await;
1111
1112        // Unpublish slugs that no longer exist in config.
1113        let current_keys: Vec<&String> = slugs.keys().collect();
1114        let registered: Vec<String> = pub_guard.registered_hostnames();
1115        for hostname in &registered {
1116            // hostname is "slug.local" — extract slug part.
1117            let slug = hostname.strip_suffix(".local").unwrap_or(hostname);
1118            if !current_keys.iter().any(|k| k.as_str() == slug) {
1119                log::info!("mDNS: unpublishing removed slug {slug}");
1120                pub_guard.unpublish(hostname);
1121            }
1122        }
1123
1124        // Publish new slugs that aren't yet registered.
1125        for slug in slugs.keys() {
1126            let hostname = format!("{slug}.local");
1127            if !pub_guard.is_published(&hostname) {
1128                log::info!("mDNS: publishing new slug {slug}");
1129                pub_guard.publish(&hostname, port);
1130            }
1131        }
1132    }
1133
1134    /// Spawn a background task that periodically flushes the state file to
1135    /// disk if it has been marked dirty.  Uses debouncing (1s interval) to
1136    /// batch rapid state changes.
1137    fn start_state_flush_task(&self) {
1138        let cancel = tokio_util::sync::CancellationToken::new();
1139        *self.flush_cancel.lock().unwrap() = Some(cancel.clone());
1140        tokio::spawn(async move {
1141            let mut interval = time::interval(Duration::from_secs(1));
1142            interval.set_missed_tick_behavior(time::MissedTickBehavior::Skip);
1143            loop {
1144                tokio::select! {
1145                    _ = interval.tick() => {}
1146                    _ = cancel.cancelled() => {
1147                        debug!("state flush task received shutdown signal");
1148                        break;
1149                    }
1150                }
1151                let state = SUPERVISOR.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            debug!("state flush task exiting");
1159        });
1160    }
1161
1162    pub(crate) async fn flush_state(&self) {
1163        let state = self.state_file.lock().await;
1164        if state.is_dirty()
1165            && let Err(e) = state.write()
1166        {
1167            warn!("failed to flush state file: {e}");
1168        }
1169    }
1170
1171    pub(crate) async fn refresh(&self) -> Result<()> {
1172        trace!("refreshing");
1173
1174        // Collect PIDs we need to check (shell PIDs and liveness PIDs)
1175        // This is more efficient than refreshing all processes on the system
1176        let dirs_with_pids = self.get_dirs_with_shell_pids().await;
1177        let liveness_sessions = self.get_liveness_sessions().await;
1178        let pids_to_check: Vec<u32> = dirs_with_pids
1179            .values()
1180            .flatten()
1181            .copied()
1182            .chain(liveness_sessions.iter().map(|(pid, _, _)| *pid))
1183            .collect::<std::collections::HashSet<_>>()
1184            .into_iter()
1185            .collect();
1186
1187        if pids_to_check.is_empty() {
1188            // No PIDs to check, skip the expensive refresh
1189            trace!("no tracked PIDs to check, skipping process refresh");
1190        } else {
1191            debug!("refreshing PIDs: {pids_to_check:?}");
1192            PROCS.refresh_pids(&pids_to_check);
1193        }
1194
1195        let mut last_refreshed_at = self.last_refreshed_at.lock().await;
1196        *last_refreshed_at = time::Instant::now();
1197
1198        self.restore_own_record().await;
1199
1200        #[cfg_attr(not(unix), allow(unused_mut))]
1201        let mut dirs_to_leave: Vec<PathBuf> = Vec::new();
1202
1203        // Prune shell PIDs that are no longer running. This is essential on
1204        // Unix so that exited shells don't keep daemons alive forever.
1205        //
1206        // On Windows, skip this check: Git Bash (MSYS2) PIDs from `$$` are
1207        // Cygwin-internal PIDs that are invisible to sysinfo (which sees
1208        // Windows PIDs). The is_running check would always return false,
1209        // immediately removing every registered shell and breaking autostop.
1210        // Shell registration/deregistration relies on UpdateShellDir IPC
1211        // messages instead.
1212        #[cfg(unix)]
1213        for (dir, pids) in dirs_with_pids {
1214            let to_remove = pids
1215                .iter()
1216                .filter(|pid| !PROCS.is_running(**pid))
1217                .collect::<Vec<_>>();
1218            for pid in &to_remove {
1219                self.remove_shell_pid(**pid).await?
1220            }
1221            if to_remove.len() == pids.len() {
1222                dirs_to_leave.push(dir);
1223            }
1224        }
1225
1226        // Atomically remove project sessions whose host PID has died or whose
1227        // recorded title no longer matches the current process title. Every
1228        // project session carries a host PID in its key, so we iterate all of
1229        // them. Re-reading the sessions under the lock prevents enter/leave
1230        // interleaving from deleting a session that was just replaced with a
1231        // new title snapshot.
1232        //
1233        // Gated to Unix to mirror the shell-PID pruning above: on Windows,
1234        // Git Bash (MSYS2) `$$` PIDs are Cygwin-internal and invisible to
1235        // sysinfo, so the liveness check would immediately revoke every
1236        // freshly-entered session. Windows relies on explicit `project leave`
1237        // (or shell UpdateShellDir) for deregistration instead.
1238        #[cfg(unix)]
1239        {
1240            let mut state = self.state_file.lock().await;
1241            for (pid, dir, recorded_title) in liveness_sessions {
1242                let Some(session) = state.get_project_session(pid, &dir) else {
1243                    continue;
1244                };
1245                let current_title = PROCS.title(pid);
1246                let is_running = PROCS.is_running(pid);
1247                debug!(
1248                    "refresh liveness session pid {pid} dir {} recorded_title={recorded_title:?} current_title={current_title:?} is_running={is_running}",
1249                    dir.display()
1250                );
1251                if should_remove_liveness_session(
1252                    session,
1253                    &recorded_title,
1254                    current_title.as_deref(),
1255                    is_running,
1256                ) {
1257                    warn!(
1258                        "removing project session pid {pid} dir {} (liveness pid title mismatch or dead)",
1259                        dir.display()
1260                    );
1261                    if state.remove_project_session(pid, &dir).is_some() {
1262                        dirs_to_leave.push(dir);
1263                    }
1264                }
1265            }
1266        }
1267
1268        for dir in dirs_to_leave {
1269            self.leave_dir(&dir).await?;
1270        }
1271
1272        // Catch state-`running` daemons that lost their monitor (e.g. the
1273        // monitor died with a previous supervisor): mark dead ones errored
1274        // and re-adopt live ones. Runs before check_retry so a daemon marked
1275        // errored here is retried on this same tick.
1276        self.reconcile_unmonitored_daemons().await;
1277
1278        self.check_retry().await?;
1279        self.process_pending_autostops().await?;
1280
1281        Ok(())
1282    }
1283
1284    /// Install a SIGCHLD handler that reaps orphaned zombie child processes.
1285    ///
1286    /// When running as PID 1 inside a container, orphaned processes are
1287    /// re-parented to PID 1. Without explicit reaping, they accumulate
1288    /// as zombies in the process table indefinitely.
1289    ///
1290    /// Only reaps processes that are NOT managed by the supervisor (i.e.
1291    /// not tracked in the state file). Managed daemon processes are reaped
1292    /// by their monitoring tasks via `child.wait()`.
1293    ///
1294    /// ## Strategy
1295    ///
1296    /// **Linux**: Uses `waitid(Id::All, WNOHANG | WNOWAIT | WEXITED)` to
1297    /// *peek* at the next zombie without consuming its status. If the PID
1298    /// belongs to a managed daemon, the reaper skips it so Tokio's
1299    /// `child.wait()` can collect the status normally. Only unmanaged
1300    /// orphans are actually reaped (via `waitpid(Pid, WNOHANG)`). This
1301    /// eliminates the race entirely.
1302    ///
1303    /// **Non-Linux Unix** (e.g. macOS — mainly for local development;
1304    /// container mode targets Linux): `waitid` is unavailable, so we fall
1305    /// back to `waitpid(None, WNOHANG)`. If the reaper accidentally
1306    /// consumes a managed PID's status, it stashes the exit code in
1307    /// [`REAPED_STATUSES`] for the monitoring task to recover.
1308    #[cfg(unix)]
1309    fn reap_zombies(&self) -> Result<()> {
1310        let mut stream = signal::unix::signal(SignalKind::child())
1311            .map_err(|e| miette::miette!("Failed to register SIGCHLD handler: {e}"))?;
1312        tokio::spawn(async move {
1313            loop {
1314                stream.recv().await;
1315                // Collect PIDs of managed daemons so we don't steal their exit status
1316                let managed_pids: HashSet<u32> = SUPERVISOR
1317                    .state_file
1318                    .lock()
1319                    .await
1320                    .daemons
1321                    .values()
1322                    .filter_map(|d| d.pid)
1323                    .collect();
1324                // Reap all available zombie children that are NOT managed
1325                Self::reap_unmanaged_zombies(&managed_pids).await;
1326            }
1327        });
1328        info!("container mode: SIGCHLD zombie reaper installed");
1329        Ok(())
1330    }
1331
1332    /// Linux implementation: peek with `waitid(WNOWAIT)` then selectively reap.
1333    ///
1334    /// `WNOWAIT` leaves the zombie in the table so we can inspect its PID
1335    /// without consuming the exit status. Only if the PID is *not* managed
1336    /// do we call `waitpid(Pid, WNOHANG)` to actually reap it.
1337    #[cfg(target_os = "linux")]
1338    async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
1339        use nix::sys::wait::{Id, WaitPidFlag, WaitStatus, waitid, waitpid};
1340        use nix::unistd::Pid;
1341
1342        loop {
1343            // Peek at the next zombie without consuming it
1344            let peek_flags = WaitPidFlag::WNOHANG | WaitPidFlag::WNOWAIT | WaitPidFlag::WEXITED;
1345            match waitid(Id::All, peek_flags) {
1346                Ok(WaitStatus::StillAlive) => break,
1347                Ok(status) => {
1348                    let Some(pid_raw) = status.pid().map(|p| p.as_raw() as u32) else {
1349                        break;
1350                    };
1351                    if managed_pids.contains(&pid_raw) {
1352                        // This is a managed daemon — leave it for Tokio's child.wait().
1353                        // We must break out of the loop because waitid(Id::All) would
1354                        // keep returning the same zombie if we don't consume it.
1355                        trace!(
1356                            "zombie reaper: skipping managed daemon pid {pid_raw}, \
1357                             leaving for Tokio to reap"
1358                        );
1359                        break;
1360                    }
1361                    // Not managed — actually reap it
1362                    match waitpid(Pid::from_raw(pid_raw as i32), Some(WaitPidFlag::WNOHANG)) {
1363                        Ok(s) => trace!("reaped orphaned zombie child: {s:?}"),
1364                        Err(nix::errno::Errno::ECHILD) => break,
1365                        Err(e) => {
1366                            trace!("waitpid error reaping pid {pid_raw}: {e}");
1367                            break;
1368                        }
1369                    }
1370                }
1371                Err(nix::errno::Errno::ECHILD) => break, // no children at all
1372                Err(e) => {
1373                    trace!("waitid error in zombie reaper: {e}");
1374                    break;
1375                }
1376            }
1377        }
1378    }
1379
1380    /// Non-Linux fallback: blind `waitpid(None, WNOHANG)` with stash recovery.
1381    ///
1382    /// Since `waitid(WNOWAIT)` is not available, we cannot peek. If we
1383    /// accidentally reap a managed PID, we stash the exit code in
1384    /// [`REAPED_STATUSES`] so the monitoring task can recover it.
1385    #[cfg(all(unix, not(target_os = "linux")))]
1386    async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
1387        use nix::sys::wait::{WaitPidFlag, WaitStatus, waitpid};
1388
1389        loop {
1390            match waitpid(None, Some(WaitPidFlag::WNOHANG)) {
1391                Ok(WaitStatus::StillAlive) => break,
1392                Ok(status) => {
1393                    let Some(pid) = status.pid().map(|p| p.as_raw() as u32) else {
1394                        continue;
1395                    };
1396                    if managed_pids.contains(&pid) {
1397                        // Race lost — stash the exit code for lifecycle recovery
1398                        let exit_code = match status {
1399                            WaitStatus::Exited(_, code) => code,
1400                            WaitStatus::Signaled(_, sig, _) => -(sig as i32),
1401                            _ => -1,
1402                        };
1403                        warn!(
1404                            "zombie reaper reaped managed daemon pid {pid} \
1405                             (exit_code={exit_code}); stashing status for recovery"
1406                        );
1407                        REAPED_STATUSES.lock().await.insert(pid, exit_code);
1408                    } else {
1409                        trace!("reaped orphaned zombie child: {status:?}");
1410                    }
1411                }
1412                Err(nix::errno::Errno::ECHILD) => break, // no more children
1413                Err(e) => {
1414                    trace!("waitpid error in zombie reaper: {e}");
1415                    break;
1416                }
1417            }
1418        }
1419    }
1420
1421    #[cfg(unix)]
1422    fn signals(&self) -> Result<()> {
1423        let signals = [
1424            SignalKind::terminate(),
1425            SignalKind::alarm(),
1426            SignalKind::interrupt(),
1427            SignalKind::quit(),
1428            SignalKind::hangup(),
1429            SignalKind::user_defined1(),
1430            SignalKind::user_defined2(),
1431        ];
1432        static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
1433        for signal in signals {
1434            let stream = match signal::unix::signal(signal) {
1435                Ok(s) => s,
1436                Err(e) => {
1437                    warn!("Failed to register signal handler for {signal:?}: {e}");
1438                    continue;
1439                }
1440            };
1441            tokio::spawn(async move {
1442                let mut stream = stream;
1443                loop {
1444                    stream.recv().await;
1445                    if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
1446                        exit(1);
1447                    } else {
1448                        SUPERVISOR.handle_signal().await;
1449                    }
1450                }
1451            });
1452        }
1453        Ok(())
1454    }
1455
1456    #[cfg(windows)]
1457    fn signals(&self) -> Result<()> {
1458        tokio::spawn(async move {
1459            static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
1460            loop {
1461                if let Err(e) = signal::ctrl_c().await {
1462                    error!("Failed to wait for ctrl-c: {}", e);
1463                    return;
1464                }
1465                if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
1466                    exit(1);
1467                } else {
1468                    SUPERVISOR.handle_signal().await;
1469                }
1470            }
1471        });
1472        Ok(())
1473    }
1474
1475    async fn handle_signal(&self) {
1476        info!("received signal, stopping");
1477        self.close().await;
1478        exit(0)
1479    }
1480
1481    pub(crate) async fn close(&self) {
1482        self.shutting_down.store(true, atomic::Ordering::Release);
1483        // Signal the proxy server to stop accepting new connections
1484        // and drain in-flight ones, *before* stopping daemons so the
1485        // proxy has time to finish forwarding active requests.
1486        if let Some(cancel) = self.proxy_cancel.lock().await.take() {
1487            cancel.cancel();
1488        }
1489
1490        // Stop the LAN IP monitor task. It watches the token that was just
1491        // cancelled, so give it a moment to come back on its own rather than
1492        // cutting it off mid-iteration the instant after asking it to stop.
1493        if let Some(mut monitor_task) = self.lan_monitor_task.lock().await.take()
1494            && tokio::time::timeout(Duration::from_secs(1), &mut monitor_task)
1495                .await
1496                .is_err()
1497        {
1498            monitor_task.abort();
1499        }
1500
1501        // Shutdown the mDNS publisher (sends goodbye packets).
1502        if let Some(publisher) = self.mdns_publisher.lock().await.take() {
1503            publisher.lock().await.shutdown();
1504        }
1505
1506        if let Some(mut dns_task) = self.dns_task.lock().await.take()
1507            && tokio::time::timeout(Duration::from_secs(5), &mut dns_task)
1508                .await
1509                .is_err()
1510        {
1511            // Cancelled but still running: aborted, as the LAN monitor is,
1512            // so its sockets are not left bound to the resolver port.
1513            dns_task.abort();
1514        }
1515
1516        if let Some(proxy_task) = self.proxy_task.lock().await.take() {
1517            // Longer than the proxy's own drain budget, so the task finishes on
1518            // its own terms rather than being cut off mid-drain.
1519            let _ = tokio::time::timeout(
1520                crate::proxy::server::SHUTDOWN_DRAIN_BUDGET + Duration::from_secs(2),
1521                proxy_task,
1522            )
1523            .await;
1524        }
1525
1526        // Clean up /etc/hosts entries managed by pitchfork
1527        let s = settings();
1528        if s.proxy.enable && s.proxy.sync_hosts {
1529            crate::proxy::hosts::clean_hosts_file();
1530        }
1531
1532        let pitchfork_id = DaemonId::pitchfork();
1533        let active = self.active_daemons().await;
1534        let active_ids: Vec<DaemonId> = active
1535            .iter()
1536            .filter(|d| d.id != pitchfork_id)
1537            .map(|d| d.id.clone())
1538            .collect();
1539
1540        // Stop daemons in reverse dependency order.
1541        // If dependency resolution fails (e.g. config changed), fall back to
1542        // stopping in arbitrary order so we still shut down cleanly.
1543        // Daemons within the same level are stopped concurrently.
1544        //
1545        // Each stop waits for the daemon's whole process group (bounded by its
1546        // stop budget) and levels are sequential, so total shutdown time is the
1547        // sum of the slowest stop per level. If an external manager (docker,
1548        // systemd) kills us before this completes, cleanup_orphaned_daemons()
1549        // recovers the leftover processes and stale state on the next start.
1550        let stop_levels = compute_reverse_stop_order(&active_ids);
1551        for level in &stop_levels {
1552            let mut tasks = Vec::new();
1553            for id in level {
1554                let id = id.clone();
1555                tasks.push(tokio::spawn(async move {
1556                    if let Err(err) = SUPERVISOR.stop(&id).await {
1557                        error!("failed to stop daemon {id}: {err}");
1558                    }
1559                }));
1560            }
1561            for task in tasks {
1562                let _ = task.await;
1563            }
1564        }
1565        let _ = self.remove_daemon(&pitchfork_id).await;
1566
1567        // Signal the background state flush task to exit so it doesn't
1568        // keep waking up and acquiring the state mutex after shutdown.
1569        if let Some(cancel) = self.flush_cancel.lock().unwrap().take() {
1570            cancel.cancel();
1571        }
1572
1573        // Force-flush state to disk before shutting down IPC so no
1574        // in-memory-only changes are lost.
1575        {
1576            let state = self.state_file.lock().await;
1577            if state.is_dirty()
1578                && let Err(e) = state.write()
1579            {
1580                warn!("failed to flush state file during shutdown: {e}");
1581            }
1582        }
1583
1584        // Signal IPC server to shut down gracefully
1585        if let Some(mut handle) = self.ipc_shutdown.lock().await.take() {
1586            handle.shutdown().await;
1587        }
1588
1589        // Wait for all in-flight monitoring tasks to finish registering their
1590        // hook handles. Each monitoring task increments `active_monitors` when
1591        // its process exits, and decrements it (+ notifies `monitor_done`)
1592        // after all fire_hook() calls complete. This replaces the old
1593        // yield_now() approach which had a race window.
1594        let drain_timeout = time::sleep(Duration::from_secs(5));
1595        tokio::pin!(drain_timeout);
1596        loop {
1597            if self.active_monitors.load(atomic::Ordering::Acquire) == 0 {
1598                break;
1599            }
1600            tokio::select! {
1601                _ = self.monitor_done.notified() => {}
1602                _ = &mut drain_timeout => {
1603                    warn!("timed out waiting for monitoring tasks to register hooks, proceeding with shutdown");
1604                    break;
1605                }
1606            }
1607        }
1608        let handles: Vec<JoinHandle<()>> = std::mem::take(&mut *self.hook_tasks.lock().await);
1609        let hook_timeout = Duration::from_secs(30);
1610        for handle in handles {
1611            match time::timeout(hook_timeout, handle).await {
1612                Ok(_) => {} // Hook completed (success or error, doesn't matter)
1613                Err(_) => {
1614                    warn!(
1615                        "hook task did not complete within {hook_timeout:?} during shutdown, skipping"
1616                    );
1617                }
1618            }
1619        }
1620
1621        // Unix: remove the socket directory if it is empty. The IPC server
1622        // already removed our socket; anything left belongs to a supervisor
1623        // that replaced this one (e.g. `supervisor run --force`) while we
1624        // were stopping daemons, and must not be deleted.
1625        // Windows: named pipes have no filesystem component.
1626        #[cfg(unix)]
1627        let _ = fs::remove_dir(&*env::IPC_SOCK_DIR);
1628    }
1629
1630    pub(crate) async fn add_notification(&self, level: log::LevelFilter, message: String) {
1631        self.pending_notifications
1632            .lock()
1633            .await
1634            .push((level, message));
1635    }
1636}
1637
1638/// Fix ownership on the state directory so non-root users can access files
1639/// created by a `sudo`-started supervisor.
1640///
1641/// When `[settings.supervisor] user`, a recorded invoking user (see
1642/// [`env::InvokingUser`]), or `SUDO_UID`/`SUDO_GID` are set, we
1643/// `chown` the state directory and safe subdirectories back to that non-root
1644/// runtime user. This is strictly better than `chmod 0o666` because it does not
1645/// widen the permission bits — the files stay owner-only (0o600/0o700) but the
1646/// *owner* is the user that daemon processes and CLI clients need to share.
1647///
1648/// **Security**: The `proxy/` subtree is intentionally skipped. It contains
1649/// `ca-key.pem` which must remain `0o600` and owned by the process that
1650/// generated it. Changing its ownership or permissions would expose the CA
1651/// private key to other local users.
1652///
1653/// If none of these are available (e.g. a direct root login, or a boot service
1654/// installed from a root shell), we fall back to relaxing permissions on only
1655/// the `sock/` and `logs/` subdirectories (plus `state.toml`) so CLI clients
1656/// can still function.
1657#[cfg(unix)]
1658fn fix_state_dir_permissions() {
1659    let state_dir = &*env::PITCHFORK_STATE_DIR;
1660    if let Some((uid, gid)) = state_owner_ids() {
1661        if !state_dir.exists()
1662            && let Err(err) = fs::create_dir_all(state_dir)
1663        {
1664            warn!(
1665                "failed to create state directory for ownership fix at {}: {err}",
1666                state_dir.display()
1667            );
1668            return;
1669        }
1670
1671        // Best path: chown back to the runtime user. Permissions stay tight.
1672        chown_recursive(state_dir, uid, gid, true);
1673        debug!(
1674            "chowned state directory to uid={uid} gid={gid} at {}",
1675            state_dir.display()
1676        );
1677    } else {
1678        if !state_dir.exists() {
1679            return;
1680        }
1681
1682        // Fallback: relax permissions on safe subdirectories only.
1683        // proxy/ is never touched.
1684        chmod_safe_subtrees(state_dir);
1685        debug!(
1686            "relaxed permissions on safe subtrees at {}",
1687            state_dir.display()
1688        );
1689    }
1690}
1691
1692#[cfg(unix)]
1693pub(crate) fn state_owner_ids() -> Option<(u32, u32)> {
1694    if !nix::unistd::Uid::effective().is_root() {
1695        return None;
1696    }
1697
1698    let s = settings();
1699    let user = s.supervisor.user.trim();
1700    if !user.is_empty() {
1701        return resolve_supervisor_user_ids(user).or_else(|| {
1702            warn!(
1703                "failed to resolve supervisor.user '{user}' for state ownership; falling back to the invoking user"
1704            );
1705            env::invoking_user_ids()
1706        });
1707    }
1708
1709    env::invoking_user_ids()
1710}
1711
1712#[cfg(unix)]
1713fn resolve_supervisor_user_ids(user: &str) -> Option<(u32, u32)> {
1714    let user_record = if user.chars().all(|c| c.is_ascii_digit()) {
1715        let uid = user.parse::<u32>().ok()?;
1716        nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1717            .ok()
1718            .flatten()
1719    } else {
1720        nix::unistd::User::from_name(user).ok().flatten()
1721    }?;
1722
1723    Some((user_record.uid.as_raw(), user_record.gid.as_raw()))
1724}
1725
1726/// Recursively `chown` a directory tree. If `skip_proxy` is true, the `proxy/`
1727/// subdirectory is skipped entirely to protect the CA private key.
1728#[cfg(unix)]
1729fn chown_recursive(dir: &std::path::Path, uid: u32, gid: u32, skip_proxy: bool) {
1730    // chown the directory itself
1731    let _ = chown_path(dir, uid, gid);
1732
1733    let entries = match std::fs::read_dir(dir) {
1734        Ok(e) => e,
1735        Err(_) => return,
1736    };
1737    for entry in entries.flatten() {
1738        let path = entry.path();
1739        if path.is_dir() {
1740            // Skip proxy/ at the top level of the state directory
1741            if skip_proxy
1742                && let Some(name) = path.file_name().and_then(|n| n.to_str())
1743                && name == "proxy"
1744            {
1745                continue;
1746            }
1747            chown_recursive(&path, uid, gid, false);
1748        } else {
1749            let _ = chown_path(&path, uid, gid);
1750        }
1751    }
1752}
1753
1754/// `chown` a single path using libc. Returns Ok(()) on success.
1755#[cfg(unix)]
1756fn chown_path(path: &std::path::Path, uid: u32, gid: u32) -> std::io::Result<()> {
1757    use std::ffi::CString;
1758    use std::os::unix::ffi::OsStrExt;
1759    let c_path = CString::new(path.as_os_str().as_bytes())
1760        .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
1761    let ret = unsafe { libc::chown(c_path.as_ptr(), uid, gid) };
1762    if ret == 0 {
1763        Ok(())
1764    } else {
1765        Err(std::io::Error::last_os_error())
1766    }
1767}
1768
1769/// Fallback: relax permissions on safe subdirectories only (sock/, logs/, and
1770/// state.toml). The proxy/ subtree is never touched.
1771#[cfg(unix)]
1772fn chmod_safe_subtrees(state_dir: &std::path::Path) {
1773    // The state directory itself needs to be traversable
1774    let _ = fs::set_permissions(state_dir, fs::Permissions::from_mode(0o755));
1775
1776    // state.toml — needs to be readable by CLI clients
1777    let state_file = state_dir.join("state.toml");
1778    if state_file.exists() {
1779        let _ = fs::set_permissions(&state_file, fs::Permissions::from_mode(0o644));
1780    }
1781
1782    // Safe subdirectories: sock/ and logs/
1783    for subdir_name in &["sock", "logs"] {
1784        let subdir = state_dir.join(subdir_name);
1785        if subdir.is_dir() {
1786            chmod_recursive(&subdir);
1787        }
1788    }
1789}
1790
1791/// On startup, reconcile daemon processes left behind by a previous supervisor
1792/// that was terminated unexpectedly (e.g. `kill -9`).
1793///
1794/// This iterates the state file for daemon entries with a recorded PID. If the
1795/// PID is still alive and its current identity matches the recorded start time
1796/// (or the recorded title for older state files), it is assumed to be an orphan
1797/// from the previous supervisor session and `supervisor.orphan_policy` decides
1798/// its fate: `adopt` (default) resumes supervision via a poll monitor and keeps
1799/// the daemon's state intact; `kill` terminates it and resets its state to
1800/// `Stopped` with no PID. If a matching live process cannot be terminated
1801/// securely, its running state is retained to prevent a duplicate instance
1802/// from being started.
1803///
1804/// Missing or mismatched identity data fails closed so a PID recycled by an
1805/// unrelated process is never adopted or killed. On Unix platforms without
1806/// durable process handles, orphan termination also fails closed because the
1807/// PID/PGID cannot be pinned between identity validation and signaling.
1808///
1809/// This is gated by the `supervisor.cleanup_orphans` setting (default: true).
1810///
1811/// `unclean` is whether the previous supervisor exited uncleanly (see
1812/// [`supervisor_exited_uncleanly`]). It is read in `start()` before the
1813/// starting supervisor records itself in the state file — this function runs
1814/// in the background after that record exists, so reading it here would
1815/// always report unclean.
1816async fn cleanup_orphaned_daemons(supervisor: &Supervisor, unclean: bool) {
1817    if !settings().supervisor.cleanup_orphans {
1818        return;
1819    }
1820
1821    let candidates: Vec<_> = {
1822        let state = supervisor.state_file.lock().await;
1823        state
1824            .daemons
1825            .values()
1826            .filter(|d| d.id != DaemonId::pitchfork() && d.pid.is_some())
1827            .cloned()
1828            .collect()
1829    };
1830
1831    if candidates.is_empty() {
1832        return;
1833    }
1834
1835    info!(
1836        "checking {} daemon(s) for orphaned processes",
1837        candidates.len()
1838    );
1839
1840    let policy = orphan_policy();
1841    let boot_time = PROCS.boot_time();
1842
1843    // Reconcile orphans in parallel — a kill waits for the daemon's whole
1844    // process group to exit, bounded by that daemon's stop budget, so
1845    // sequential processing would make total cleanup time the sum of the
1846    // budgets.
1847    let tasks: Vec<_> = candidates
1848        .into_iter()
1849        .map(|daemon| {
1850            let policy = policy.clone();
1851            tokio::spawn(cleanup_orphaned_daemon(daemon, policy, boot_time, unclean))
1852        })
1853        .collect();
1854    for task in tasks {
1855        let _ = task.await;
1856    }
1857}
1858
1859/// Reconcile a single orphan candidate: adopt it, kill it, or reset its state,
1860/// per the policy and identity checks described on [`cleanup_orphaned_daemons`].
1861///
1862/// Holds the daemon's stop lock so a concurrent Run/Stop request for the same
1863/// daemon (cleanup runs in the background) serializes with the orphan kill,
1864/// and re-checks the recorded PID under the lock: if it changed, another path
1865/// already replaced or cleaned up this record and the snapshot is stale.
1866async fn cleanup_orphaned_daemon(
1867    daemon: crate::daemon::Daemon,
1868    policy: String,
1869    boot_time: u64,
1870    unclean: bool,
1871) {
1872    let supervisor: &Supervisor = &SUPERVISOR;
1873    let Some(pid) = daemon.pid else { return };
1874
1875    let lock = supervisor.stop_lock(&daemon.id).await;
1876    let _guard = lock.lock().await;
1877    let current_pid = {
1878        let state = supervisor.state_file.lock().await;
1879        state.daemons.get(&daemon.id).and_then(|d| d.pid)
1880    };
1881    if current_pid != Some(pid) {
1882        debug!(
1883            "orphan cleanup: daemon {} pid changed (recorded {pid}, now {current_pid:?}), skipping",
1884            daemon.id
1885        );
1886        return;
1887    }
1888
1889    // Refresh the candidate immediately before checking it: waiting for the
1890    // stop lock can await another path's stop timeout, during which this PID
1891    // may exit and be recycled.
1892    PROCS.refresh_pids(&[pid]);
1893
1894    if !PROCS.is_running(pid) {
1895        // PID already dead — the daemon exited while unsupervised, so
1896        // record a terminal status that reflects whether it died under a
1897        // crashed supervisor (retryable) or with the machine.
1898        let status = unobserved_exit_status(
1899            &daemon.status,
1900            daemon.boot_time,
1901            boot_time,
1902            unclean,
1903            daemon.oneshot,
1904        );
1905        reset_daemon_state(supervisor, &daemon.id, status, ExitObservation::Unobserved).await;
1906        return;
1907    }
1908
1909    // Safety check: verify the live process really is the daemon we
1910    // recorded, not an unrelated process that received a recycled PID.
1911    // The kernel start time is a stable identity for the lifetime of a
1912    // process, and is the only thing accepted as one.
1913    let current_start_time = PROCS.start_time(pid);
1914    let matches = process_identity_matches(daemon.start_time, current_start_time);
1915
1916    if !matches {
1917        // Either side missing means the identity cannot be checked at all,
1918        // which is different from checking it and finding a stranger: retain
1919        // the running state rather than resetting a record whose process may
1920        // well still be the daemon.
1921        if daemon.start_time.is_none() || current_start_time.is_none() {
1922            warn!(
1923                "could not verify the identity of live pid {pid} recorded for daemon {}; retaining running state",
1924                daemon.id,
1925            );
1926            return;
1927        }
1928        warn!(
1929            "pid {pid} recorded for daemon {} belongs to a different process now (PID recycled); resetting state without killing",
1930            daemon.id,
1931        );
1932        // The daemon died at some unknown point and the OS handed its PID
1933        // to something else — same unobserved exit as a dead PID.
1934        let status = unobserved_exit_status(
1935            &daemon.status,
1936            daemon.boot_time,
1937            boot_time,
1938            unclean,
1939            daemon.oneshot,
1940        );
1941        reset_daemon_state(supervisor, &daemon.id, status, ExitObservation::Unobserved).await;
1942        return;
1943    }
1944
1945    // Both policies need a verified start time: killing revalidates it
1946    // while pinned to the process, and adoption anchors its poll monitor
1947    // to it so a later PID recycle is never mistaken for the daemon.
1948    let Some(expected_start_time) = current_start_time else {
1949        warn!(
1950            "could not read start time for live pid {pid} recorded for daemon {}; retaining running state",
1951            daemon.id,
1952        );
1953        return;
1954    };
1955
1956    // Identity verified — the process really is our orphaned daemon.
1957    // The policy decides whether supervision resumes or the slate is
1958    // wiped clean.
1959    if policy == "adopt" {
1960        supervisor
1961            .adopt_daemon(&daemon, pid, expected_start_time)
1962            .await;
1963        return;
1964    }
1965
1966    info!("terminating orphaned daemon {} (pid {pid})", daemon.id);
1967
1968    let stop_cfg = daemon.stop_signal.unwrap_or_default();
1969    let termination_result = PROCS
1970        .kill_process_group_if_start_time_matches_async(
1971            pid,
1972            Some(expected_start_time),
1973            stop_cfg.signal.into(),
1974            stop_cfg.timeout,
1975        )
1976        .await;
1977
1978    match termination_result {
1979        Ok(true) => {}
1980        Ok(false) => {
1981            warn!(
1982                "could not securely terminate orphaned daemon {} (pid {pid}); retaining running state",
1983                daemon.id
1984            );
1985            return;
1986        }
1987        Err(err) => {
1988            warn!(
1989                "failed to terminate orphaned daemon {} (pid {pid}): {err}; retaining running state",
1990                daemon.id
1991            );
1992            return;
1993        }
1994    }
1995
1996    // We terminated the orphan ourselves, so this is an observed,
1997    // intentional stop rather than an unobserved exit.
1998    reset_daemon_state(
1999        supervisor,
2000        &daemon.id,
2001        DaemonStatus::Stopped,
2002        ExitObservation::Terminated,
2003    )
2004    .await;
2005}
2006
2007/// Effective `supervisor.orphan_policy`, warning on an unrecognized value
2008/// (which falls back to the default of adopting).
2009pub(crate) fn orphan_policy() -> String {
2010    let policy = settings().supervisor.orphan_policy.clone();
2011    match policy.as_str() {
2012        "adopt" | "kill" => policy,
2013        other => {
2014            warn!("unknown supervisor.orphan_policy '{other}', defaulting to 'adopt'");
2015            "adopt".to_string()
2016        }
2017    }
2018}
2019
2020/// Verify that live process identity matches the persisted daemon identity.
2021///
2022/// Both start times are required. A process name was once accepted in place of
2023/// a recorded start time, for state written before start times existed, but a
2024/// name is not an identity: a recycled PID belonging to another copy of the same
2025/// program matches it, and adopting or killing on that basis acts on the wrong
2026/// process. Missing identity, on either side, means unverifiable — and
2027/// unverifiable must never authorize acting on a process.
2028fn process_identity_matches(
2029    recorded_start_time: Option<u64>,
2030    current_start_time: Option<u64>,
2031) -> bool {
2032    match (recorded_start_time, current_start_time) {
2033        (Some(recorded), Some(current)) => recorded == current,
2034        _ => false,
2035    }
2036}
2037
2038/// Whether a PID read from persisted state may be signalled.
2039///
2040/// Stopping a daemon signals its whole process *group*, so acting on a PID that
2041/// has been recycled since it was recorded takes down an unrelated process tree.
2042/// Records are refused only when their identity is positively contradicted: if
2043/// either start time is unknown the PID stays as signallable as it was before
2044/// identities were recorded, so a daemon whose record predates the field can
2045/// still be stopped rather than becoming permanently unstoppable.
2046///
2047/// This is deliberately weaker than [`process_identity_matches`], which decides
2048/// whether to adopt or kill a process nobody asked about. Here the user has
2049/// named the daemon and asked for it to stop; the check exists to catch the
2050/// case where the answer is provably the wrong process.
2051pub(crate) fn signalling_pid_is_authorized(
2052    recorded_start_time: Option<u64>,
2053    current_start_time: Option<u64>,
2054) -> bool {
2055    !matches!(
2056        (recorded_start_time, current_start_time),
2057        (Some(recorded), Some(current)) if recorded != current
2058    )
2059}
2060
2061/// How a daemon's run ended, which decides what happens to the recorded
2062/// `last_exit_success` that cron `retrigger = "success" | "fail"` reads.
2063#[derive(Clone, Copy, PartialEq, Eq)]
2064pub(crate) enum ExitObservation {
2065    /// Nobody saw how the run ended, because the monitor that would have
2066    /// observed it died with a previous supervisor. The recorded outcome is
2067    /// cleared to `None`.
2068    ///
2069    /// Every option here is imperfect, so this picks the one that asserts
2070    /// nothing false. `Some(false)` would fabricate a failure, silently
2071    /// breaking a `retrigger = "success"` chain whose run may well have
2072    /// succeeded; `Some(true)` fabricates the opposite; keeping the previous
2073    /// value attributes an earlier run's outcome to this one. `None` says
2074    /// "unknown", reusing the reading the cron watcher already applies to a
2075    /// daemon that has never run.
2076    ///
2077    /// The tradeoff is that `None` satisfies both `retrigger = "success"`
2078    /// (`unwrap_or(true)`) and `retrigger = "fail"` (`!unwrap_or(false)`), so
2079    /// such a daemon fires once at its next scheduled time regardless of which
2080    /// it configured. That is schedule-gated rather than a loop, and it biases
2081    /// toward running the daemon over leaving it permanently untriggered.
2082    /// Distinguishing "unknown" from "never ran" would require a third cron
2083    /// state and is deliberately left out of scope here.
2084    Unobserved,
2085    /// We terminated the process ourselves, so the outcome is not a mystery:
2086    /// it stopped because we asked it to. Recorded as a success, matching the
2087    /// convention `Supervisor::stop` already uses for a deliberate stop.
2088    Terminated,
2089}
2090
2091impl ExitObservation {
2092    /// The `last_exit_success` value this observation implies.
2093    pub(crate) fn last_exit_success(self) -> Option<bool> {
2094        match self {
2095            ExitObservation::Unobserved => None,
2096            ExitObservation::Terminated => Some(true),
2097        }
2098    }
2099}
2100
2101/// Clear a daemon's runtime state (pid, process identity, active port) after
2102/// its process is gone or is no longer ours to manage.
2103///
2104/// Config fields are preserved by cloning the existing record, so a reset can
2105/// never drop a daemon's command, retry policy, or schedule.
2106async fn reset_daemon_state(
2107    supervisor: &Supervisor,
2108    id: &DaemonId,
2109    status: DaemonStatus,
2110    observation: ExitObservation,
2111) {
2112    let mut state_file = supervisor.state_file.lock().await;
2113    let Some(existing) = state_file.daemons.get(id) else {
2114        return;
2115    };
2116    let mut daemon = existing.clone();
2117    daemon.pid = None;
2118    daemon.title = None;
2119    daemon.start_time = None;
2120    daemon.boot_time = None;
2121    daemon.status = status;
2122    daemon.last_exit_success = observation.last_exit_success();
2123    daemon.active_port = None;
2124    state_file.clear_active_port(id);
2125    state_file.insert_daemon(id, daemon);
2126}
2127
2128/// Boot times this far apart are treated as different boots.
2129///
2130/// Sized to the only platform that reports a jittery value: Windows derives
2131/// boot time as `now - GetTickCount64()`, sampling two clocks independently,
2132/// so consecutive calls within one boot can differ by about a second. Linux
2133/// (`/proc/stat` btime) and macOS (`kern.boottime`) report stable values.
2134///
2135/// Deliberately kept this tight so a genuine reboot can never fall inside it:
2136/// a prior session would have to boot, start the supervisor, spawn a daemon,
2137/// have that daemon die, and complete a reboot inside two seconds, which no
2138/// real boot cycle reaches. A larger window would misread a short-lived
2139/// previous boot (e.g. a device in a reboot loop) as the current one and
2140/// resurrect daemons a reboot should have left stopped.
2141const BOOT_TIME_TOLERANCE_SECS: u64 = 2;
2142
2143/// Terminal status for a daemon whose process is gone and whose exit was
2144/// never observed, because the monitor that would have seen it died with a
2145/// previous supervisor.
2146///
2147/// A daemon recorded `Running` was expected to still be alive, so it died
2148/// under the crashed supervisor: `Errored(-1)` ("unknown exit code") makes it
2149/// eligible for its configured retries. Two cases stay `Stopped` instead:
2150///
2151/// - records from an earlier boot, whose processes died with the machine —
2152///   auto-restarting those is what `boot_start` is for, and reviving every
2153///   retry-configured daemon after a reboot would be a surprise
2154/// - any other status (in practice `Stopping`), i.e. an intentional stop that
2155///   completed while the supervisor was gone
2156pub(crate) fn unobserved_exit_status(
2157    status: &DaemonStatus,
2158    recorded_boot_time: Option<u64>,
2159    current_boot_time: u64,
2160    supervisor_exited_uncleanly: bool,
2161    oneshot: bool,
2162) -> DaemonStatus {
2163    let same_boot = recorded_boot_time
2164        .is_some_and(|recorded| recorded.abs_diff(current_boot_time) <= BOOT_TIME_TOLERANCE_SECS);
2165    // A task gets `stopped` rather than `errored` for the same reason the
2166    // adopted path does: `errored` is what `check_retry` looks for, and
2167    // nobody saw how this run ended, so retrying it would re-run a migration
2168    // or a seed that may well have succeeded. Leave re-running to an explicit
2169    // start.
2170    if status.is_running() && same_boot && supervisor_exited_uncleanly && !oneshot {
2171        DaemonStatus::Errored(-1)
2172    } else {
2173        DaemonStatus::Stopped
2174    }
2175}
2176
2177/// Whether the supervisor that owned this state file failed to shut down
2178/// cleanly, meaning any daemon it left behind stopped for reasons nobody
2179/// recorded.
2180///
2181/// A clean shutdown removes the supervisor's own entry: `close()` does it on
2182/// Unix, where the stop signal is delivered and handled, and the
2183/// `supervisor stop` command does it on Windows, which has no POSIX signals
2184/// and force-terminates the process instead. A crash, an external `kill -9`,
2185/// or a `--force` replacement all leave the entry behind.
2186///
2187/// This must be read before the starting supervisor records itself, which is
2188/// why `cleanup_orphaned_daemons` runs first in `start()`.
2189async fn supervisor_exited_uncleanly(supervisor: &Supervisor) -> bool {
2190    supervisor
2191        .state_file
2192        .lock()
2193        .await
2194        .daemons
2195        .contains_key(&DaemonId::pitchfork())
2196}
2197
2198/// Recursively chmod: directories → 0o755, files → 0o644.
2199#[cfg(unix)]
2200fn chmod_recursive(dir: &std::path::Path) {
2201    let _ = fs::set_permissions(dir, fs::Permissions::from_mode(0o755));
2202    let entries = match fs::read_dir(dir) {
2203        Ok(e) => e,
2204        Err(_) => return,
2205    };
2206    for entry in entries.flatten() {
2207        let path = entry.path();
2208        if path.is_dir() {
2209            chmod_recursive(&path);
2210        } else {
2211            let _ = fs::set_permissions(&path, fs::Permissions::from_mode(0o644));
2212        }
2213    }
2214}
2215
2216#[cfg(test)]
2217mod tests {
2218    use super::{
2219        BOOT_TIME_TOLERANCE_SECS, legacy_supervisor_title_matches, process_identity_matches,
2220        should_remove_liveness_session, signalling_pid_is_authorized, supervisor_identity_matches,
2221        unobserved_exit_status,
2222    };
2223    use crate::daemon_status::DaemonStatus;
2224    use crate::state_file::ProjectSession;
2225
2226    const BOOT: u64 = 1_700_000_000;
2227
2228    #[test]
2229    fn unobserved_running_death_in_current_boot_is_retryable() {
2230        // Died under a crashed supervisor during this boot: Errored(-1) makes
2231        // the daemon eligible for its configured retries.
2232        assert!(matches!(
2233            unobserved_exit_status(&DaemonStatus::Running, Some(BOOT), BOOT, true, false),
2234            DaemonStatus::Errored(-1)
2235        ));
2236    }
2237
2238    #[test]
2239    fn unobserved_task_death_is_stopped_not_retried() {
2240        // Nobody saw how the run ended, so `errored` would hand a migration or
2241        // a seed to check_retry on a guess. Matches what the adopted path
2242        // records, and what the guide promises.
2243        assert!(matches!(
2244            unobserved_exit_status(&DaemonStatus::Running, Some(BOOT), BOOT, true, true),
2245            DaemonStatus::Stopped
2246        ));
2247    }
2248
2249    #[test]
2250    fn unobserved_running_death_from_previous_boot_is_stopped() {
2251        // The process died with the machine; reviving every retry-configured
2252        // daemon after a reboot is what boot_start is for.
2253        assert!(matches!(
2254            unobserved_exit_status(
2255                &DaemonStatus::Running,
2256                Some(BOOT - 86_400),
2257                BOOT,
2258                true,
2259                false
2260            ),
2261            DaemonStatus::Stopped
2262        ));
2263    }
2264
2265    #[test]
2266    fn unobserved_exit_tolerates_boot_time_jitter() {
2267        // Windows recomputes boot time as now - GetTickCount64(), which can
2268        // drift about a second between samples within one boot.
2269        let within = BOOT + BOOT_TIME_TOLERANCE_SECS;
2270        assert!(matches!(
2271            unobserved_exit_status(&DaemonStatus::Running, Some(within), BOOT, true, false),
2272            DaemonStatus::Errored(-1)
2273        ));
2274        let beyond = BOOT + BOOT_TIME_TOLERANCE_SECS + 1;
2275        assert!(matches!(
2276            unobserved_exit_status(&DaemonStatus::Running, Some(beyond), BOOT, true, false),
2277            DaemonStatus::Stopped
2278        ));
2279    }
2280
2281    #[test]
2282    fn unobserved_exit_after_clean_shutdown_is_stopped() {
2283        // A deliberate `supervisor stop` can leave running records behind on
2284        // platforms where the supervisor cannot handle the stop signal. Those
2285        // daemons were stopped on purpose, so they must not be reported as
2286        // failures or resurrected by the retry checker.
2287        assert!(matches!(
2288            unobserved_exit_status(&DaemonStatus::Running, Some(BOOT), BOOT, false, false),
2289            DaemonStatus::Stopped
2290        ));
2291    }
2292
2293    #[test]
2294    fn unobserved_exit_treats_short_previous_boot_as_previous() {
2295        // A device in a reboot loop can produce consecutive boots seconds
2296        // apart. The jitter window must stay far below that so those records
2297        // are still recognised as belonging to an earlier boot.
2298        for gap in [5, 30, 59, 60] {
2299            assert!(
2300                matches!(
2301                    unobserved_exit_status(
2302                        &DaemonStatus::Running,
2303                        Some(BOOT - gap),
2304                        BOOT,
2305                        true,
2306                        false
2307                    ),
2308                    DaemonStatus::Stopped
2309                ),
2310                "boot {gap}s earlier should be treated as a previous boot"
2311            );
2312        }
2313    }
2314
2315    #[test]
2316    fn unobserved_exit_without_recorded_boot_time_is_stopped() {
2317        // Legacy state files predating the field fail closed to today's
2318        // behavior rather than triggering surprise retries.
2319        assert!(matches!(
2320            unobserved_exit_status(&DaemonStatus::Running, None, BOOT, true, false),
2321            DaemonStatus::Stopped
2322        ));
2323    }
2324
2325    #[test]
2326    fn supervisor_identity_matches_same_generation() {
2327        assert!(supervisor_identity_matches(
2328            Some(100),
2329            Some(100),
2330            Some(BOOT),
2331            BOOT
2332        ));
2333    }
2334
2335    #[test]
2336    fn supervisor_identity_survives_clock_steps_when_tokens_match() {
2337        // An NTP step or sleep/resume moves the realtime-derived boot time by
2338        // far more than the tolerance while the supervisor keeps running.
2339        // Matching start tokens prove it is the same process; declaring it
2340        // stale here would start a second supervisor.
2341        assert!(supervisor_identity_matches(
2342            Some(100),
2343            Some(100),
2344            Some(BOOT),
2345            BOOT + 3600
2346        ));
2347        assert!(supervisor_identity_matches(
2348            Some(100),
2349            Some(100),
2350            Some(BOOT + 3600),
2351            BOOT
2352        ));
2353    }
2354
2355    #[test]
2356    fn supervisor_identity_rejects_recycled_pid() {
2357        // The supervisor died and something else got its PID, within this
2358        // boot or across a reboot: the token differs either way.
2359        assert!(!supervisor_identity_matches(
2360            Some(100),
2361            Some(200),
2362            Some(BOOT),
2363            BOOT
2364        ));
2365        assert!(!supervisor_identity_matches(
2366            Some(100),
2367            Some(200),
2368            Some(BOOT),
2369            BOOT + 3600
2370        ));
2371    }
2372
2373    #[test]
2374    fn supervisor_identity_uses_boot_time_when_a_token_is_missing() {
2375        // Discussion #877: the record survived a reboot and an early system
2376        // daemon now owns the PID. With no live token to compare, the boot
2377        // time is what proves the record stale.
2378        assert!(!supervisor_identity_matches(
2379            Some(100),
2380            None,
2381            Some(BOOT),
2382            BOOT + 3600
2383        ));
2384        // Same boot (within Windows boot-time jitter) and no contradiction.
2385        assert!(supervisor_identity_matches(
2386            Some(100),
2387            None,
2388            Some(BOOT),
2389            BOOT + BOOT_TIME_TOLERANCE_SECS
2390        ));
2391    }
2392
2393    #[test]
2394    fn supervisor_identity_tolerates_legacy_records() {
2395        // A record written before either field existed is not contradicted by
2396        // anything here; `supervisor_record_is_live` applies the process-name
2397        // check to those instead.
2398        assert!(supervisor_identity_matches(None, Some(100), None, BOOT));
2399        // A record with only a boot time is still rejected across a reboot.
2400        assert!(!supervisor_identity_matches(
2401            None,
2402            Some(100),
2403            Some(BOOT),
2404            BOOT + 3600
2405        ));
2406        assert!(supervisor_identity_matches(
2407            None,
2408            Some(100),
2409            Some(BOOT),
2410            BOOT
2411        ));
2412    }
2413
2414    #[test]
2415    fn legacy_supervisor_title_requires_a_pitchfork_process() {
2416        assert!(legacy_supervisor_title_matches(Some("pitchfork")));
2417        assert!(legacy_supervisor_title_matches(Some("pitchfork.exe")));
2418        assert!(legacy_supervisor_title_matches(Some("Pitchfork")));
2419        // Discussion #877: an Apple LaunchAgent inherited the PID after a reboot.
2420        assert!(!legacy_supervisor_title_matches(Some(
2421            "AMPDeviceDiscoveryAgent"
2422        )));
2423        assert!(!legacy_supervisor_title_matches(Some("sleep")));
2424        assert!(!legacy_supervisor_title_matches(None));
2425    }
2426
2427    #[test]
2428    fn unobserved_exit_of_stopping_daemon_is_stopped() {
2429        // An intentional stop that completed while the supervisor was gone is
2430        // not a failure, even within the same boot.
2431        assert!(matches!(
2432            unobserved_exit_status(&DaemonStatus::Stopping, Some(BOOT), BOOT, true, false),
2433            DaemonStatus::Stopped
2434        ));
2435    }
2436
2437    #[test]
2438    fn orphan_identity_requires_both_start_times() {
2439        assert!(process_identity_matches(Some(123), Some(123)));
2440        assert!(!process_identity_matches(Some(123), Some(456)));
2441        // Unreadable current identity: unverifiable, so not a match.
2442        assert!(!process_identity_matches(Some(123), None));
2443    }
2444
2445    #[test]
2446    fn signalling_is_refused_only_for_a_contradicted_identity() {
2447        // Provably someone else's process group: refuse.
2448        assert!(!signalling_pid_is_authorized(Some(123), Some(456)));
2449        // Verified as the daemon's own.
2450        assert!(signalling_pid_is_authorized(Some(123), Some(123)));
2451        // Unknown on either side. Stopping stays possible, because the user has
2452        // named this daemon and a record that cannot be verified must not become
2453        // one that can never be stopped.
2454        assert!(signalling_pid_is_authorized(None, Some(123)));
2455        assert!(signalling_pid_is_authorized(Some(123), None));
2456        assert!(signalling_pid_is_authorized(None, None));
2457    }
2458
2459    #[test]
2460    fn orphan_identity_rejects_records_without_a_start_time() {
2461        // State written before start times were recorded. A process name used
2462        // to stand in here, but another copy of the same program on a recycled
2463        // PID matches a name, so such records are no longer verifiable and must
2464        // not authorize adopting or killing anything.
2465        assert!(!process_identity_matches(None, Some(123)));
2466        assert!(!process_identity_matches(None, None));
2467    }
2468
2469    #[test]
2470    fn should_not_remove_when_state_title_differs_from_snapshot() {
2471        // The session was re-entered after the snapshot was taken, producing a
2472        // new title in state. The snapshot title is stale; skip removal.
2473        let session = ProjectSession {
2474            liveness_title: Some("new_title".to_string()),
2475        };
2476        let recorded_title = Some("old_title".to_string());
2477
2478        assert!(!should_remove_liveness_session(
2479            &session,
2480            &recorded_title,
2481            Some("new_title"),
2482            true,
2483        ));
2484    }
2485
2486    #[test]
2487    fn should_remove_when_running_title_mismatches() {
2488        let session = ProjectSession {
2489            liveness_title: Some("recorded_title".to_string()),
2490        };
2491        let recorded_title = Some("recorded_title".to_string());
2492
2493        assert!(should_remove_liveness_session(
2494            &session,
2495            &recorded_title,
2496            Some("different_title"),
2497            true,
2498        ));
2499    }
2500
2501    #[test]
2502    fn should_remove_when_dead() {
2503        let session = ProjectSession {
2504            liveness_title: Some("recorded_title".to_string()),
2505        };
2506        let recorded_title = Some("recorded_title".to_string());
2507
2508        assert!(should_remove_liveness_session(
2509            &session,
2510            &recorded_title,
2511            Some("recorded_title"),
2512            false,
2513        ));
2514    }
2515
2516    #[test]
2517    fn should_not_remove_when_alive_and_title_matches() {
2518        let session = ProjectSession {
2519            liveness_title: Some("recorded_title".to_string()),
2520        };
2521        let recorded_title = Some("recorded_title".to_string());
2522
2523        assert!(!should_remove_liveness_session(
2524            &session,
2525            &recorded_title,
2526            Some("recorded_title"),
2527            true,
2528        ));
2529    }
2530
2531    #[cfg(unix)]
2532    #[test]
2533    fn fd_scan_end_reaches_open_fds_above_the_fd_limit() {
2534        // A caller opened fd 900, then lowered RLIMIT_NOFILE to 256.
2535        assert_eq!(super::fd_scan_end(256, [0, 1, 2, 900]), 901);
2536        assert_eq!(super::fd_scan_end(256, [0, 1, 2, 10]), 256);
2537        assert_eq!(super::fd_scan_end(-1, []), 1 << 16);
2538        assert_eq!(super::fd_scan_end(libc::c_long::MAX, []), 1 << 20);
2539    }
2540
2541    /// Exercises the fcntl fallback directly: on Linux >= 5.11 the spawn path
2542    /// never reaches it, but it is the only implementation on macOS.
2543    #[cfg(unix)]
2544    #[test]
2545    fn cloexec_fd_scan_hides_inherited_fds_from_children() {
2546        use std::os::unix::process::CommandExt;
2547        use std::process::{Command, Stdio};
2548
2549        // A pipe without O_CLOEXEC, like bats' fd 3 or a wrapper script's pipe.
2550        let mut fds = [0; 2];
2551        assert_eq!(unsafe { libc::pipe(fds.as_mut_ptr()) }, 0);
2552        let fd = fds[1];
2553        let child_sees_fd = |scan: bool| {
2554            let mut cmd = Command::new("sh");
2555            cmd.arg("-c")
2556                .arg(format!("[ -e /dev/fd/{fd} ]"))
2557                .stdin(Stdio::null())
2558                .stdout(Stdio::null())
2559                .stderr(Stdio::null());
2560            if scan {
2561                let max_fd = super::max_inherited_fd();
2562                unsafe {
2563                    cmd.pre_exec(move || {
2564                        super::cloexec_fd_scan(max_fd);
2565                        Ok(())
2566                    });
2567                }
2568            }
2569            cmd.status().unwrap().success()
2570        };
2571        let without_scan = child_sees_fd(false);
2572        let with_scan = child_sees_fd(true);
2573        unsafe {
2574            libc::close(fds[0]);
2575            libc::close(fds[1]);
2576        }
2577        assert!(without_scan, "control: child should inherit fd {fd}");
2578        assert!(!with_scan, "fd {fd} leaked past cloexec_fd_scan");
2579    }
2580}