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