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