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