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