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