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//! - `autostop`: Autostop logic and boot daemon startup
7//! - `retry`: Retry logic with backoff
8//! - `watchers`: Background tasks (interval, cron, file watching)
9//! - `ipc_handlers`: IPC request dispatch
10
11mod autostop;
12mod hooks;
13mod ipc_handlers;
14mod lifecycle;
15#[cfg(unix)]
16mod pty;
17mod retry;
18mod state;
19mod watchers;
20
21use crate::daemon_id::DaemonId;
22use crate::daemon_status::DaemonStatus;
23use crate::deps::compute_reverse_stop_order;
24use crate::ipc::server::{IpcServer, IpcServerHandle};
25
26use crate::procs::PROCS;
27use crate::settings::settings;
28use crate::state_file::StateFile;
29use crate::{Result, env};
30use duct::cmd;
31use miette::IntoDiagnostic;
32use once_cell::sync::Lazy;
33use std::collections::HashMap;
34#[cfg(unix)]
35use std::collections::HashSet;
36use std::fs;
37#[cfg(unix)]
38use std::os::unix::fs::PermissionsExt;
39use std::process::exit;
40use std::sync::atomic;
41use std::sync::atomic::{AtomicBool, AtomicU32};
42use std::time::Duration;
43#[cfg(unix)]
44use tokio::signal::unix::SignalKind;
45use tokio::sync::{Mutex, Notify};
46use tokio::task::JoinHandle;
47use tokio::{signal, time};
48
49/// Exit statuses reaped by the container-mode zombie reaper for managed daemon
50/// PIDs. On non-Linux Unix platforms where `waitid(WNOWAIT)` is unavailable,
51/// `waitpid(None, WNOHANG)` may race with Tokio's `child.wait()`. When the
52/// zombie reaper wins, the exit status is stashed here so the monitoring task
53/// in lifecycle.rs can recover it instead of treating the ECHILD as a failure.
54///
55/// On Linux this map is unused because the reaper uses `waitid` with `WNOWAIT`
56/// to peek before reaping, which avoids the race entirely.
57#[cfg(all(unix, not(target_os = "linux")))]
58pub(crate) static REAPED_STATUSES: Lazy<Mutex<HashMap<u32, i32>>> =
59    Lazy::new(|| Mutex::new(HashMap::new()));
60
61// Re-export types needed by other modules
62pub(crate) use state::UpsertDaemonOpts;
63
64pub struct Supervisor {
65    pub(crate) state_file: Mutex<StateFile>,
66    pub(crate) pending_notifications: Mutex<Vec<(log::LevelFilter, String)>>,
67    pub(crate) last_refreshed_at: Mutex<time::Instant>,
68    /// Map of daemon ID to scheduled autostop time
69    pub(crate) pending_autostops: Mutex<HashMap<DaemonId, time::Instant>>,
70    /// Handle for graceful IPC server shutdown
71    pub(crate) ipc_shutdown: Mutex<Option<IpcServerHandle>>,
72    /// Tracks in-flight hook tasks so shutdown can wait for them to complete
73    pub(crate) hook_tasks: Mutex<Vec<JoinHandle<()>>>,
74    /// Number of monitoring tasks that are still running (between process exit
75    /// and hook registration completion). Used by `close()` to know when it is
76    /// safe to drain `hook_tasks`.
77    pub(crate) active_monitors: AtomicU32,
78    /// Signalled by each monitoring task after it finishes registering hooks
79    /// (or decides it has nothing to register). `close()` waits on this.
80    pub(crate) monitor_done: Notify,
81    /// Cancellation token for the proxy server — cancelled on shutdown to
82    /// stop accepting new connections and drain in-flight ones.
83    pub(crate) proxy_cancel: Mutex<Option<tokio_util::sync::CancellationToken>>,
84    /// Join handle for the proxy task so shutdown can wait for cleanup.
85    pub(crate) proxy_task: Mutex<Option<JoinHandle<()>>>,
86    /// mDNS publisher for LAN mode (None if LAN mode is disabled).
87    /// Shared with the LAN IP monitor task so it can re-publish on IP change.
88    pub(crate) mdns_publisher:
89        Mutex<Option<std::sync::Arc<tokio::sync::Mutex<crate::proxy::mdns::MdnsPublisher>>>>,
90    /// Join handle for the LAN IP monitor task.
91    pub(crate) lan_monitor_task: Mutex<Option<JoinHandle<()>>>,
92    /// Cancellation token for the background state flush task.
93    pub(crate) flush_cancel: std::sync::Mutex<Option<tokio_util::sync::CancellationToken>>,
94}
95
96pub(crate) fn interval_duration() -> Duration {
97    settings().general_interval()
98}
99
100pub static SUPERVISOR: Lazy<Supervisor> =
101    Lazy::new(|| Supervisor::new().expect("Error creating supervisor"));
102
103pub fn start_if_not_running() -> Result<()> {
104    let sf = StateFile::get();
105    if let Some(d) = sf.daemons.get(&DaemonId::pitchfork())
106        && let Some(pid) = d.pid
107        && PROCS.is_running(pid)
108    {
109        return Ok(());
110    }
111    start_in_background()
112}
113
114pub fn start_in_background() -> Result<()> {
115    debug!("starting supervisor in background");
116    // Ensure the log directory exists so we can redirect stderr there.
117    // Panics and other fatal errors from the background supervisor process
118    // would otherwise be silently swallowed.
119    let log_file = &*env::PITCHFORK_LOG_FILE;
120    if let Some(parent) = log_file.parent() {
121        let _ = fs::create_dir_all(parent);
122    }
123    let stderr_file = fs::OpenOptions::new()
124        .create(true)
125        .append(true)
126        .open(log_file)
127        .into_diagnostic()?;
128    #[cfg(unix)]
129    fix_state_dir_permissions();
130    cmd!(&*env::PITCHFORK_BIN, "supervisor", "run")
131        .stdout_null()
132        .stderr_file(stderr_file)
133        .start()
134        .into_diagnostic()?;
135    Ok(())
136}
137
138impl Supervisor {
139    pub fn new() -> Result<Self> {
140        Ok(Self {
141            state_file: Mutex::new(StateFile::read(&*env::PITCHFORK_STATE_FILE).unwrap_or_else(
142                |e| {
143                    warn!("failed to read state file, starting with empty state: {e}");
144                    StateFile::new(env::PITCHFORK_STATE_FILE.clone())
145                },
146            )),
147            last_refreshed_at: Mutex::new(time::Instant::now()),
148            pending_notifications: Mutex::new(vec![]),
149            pending_autostops: Mutex::new(HashMap::new()),
150            ipc_shutdown: Mutex::new(None),
151            hook_tasks: Mutex::new(Vec::new()),
152            active_monitors: AtomicU32::new(0),
153            monitor_done: Notify::new(),
154            proxy_cancel: Mutex::new(None),
155            proxy_task: Mutex::new(None),
156            mdns_publisher: Mutex::new(None),
157            lan_monitor_task: Mutex::new(None),
158            flush_cancel: std::sync::Mutex::new(None),
159        })
160    }
161
162    pub async fn start(
163        &self,
164        is_boot: bool,
165        container: bool,
166        web_port: Option<u16>,
167        web_path: Option<String>,
168    ) -> Result<()> {
169        // Ensure the state directory and its contents are accessible by non-root
170        // users. This is needed when the supervisor is started with `sudo` — all
171        // files it creates are owned by root, which prevents normal CLI clients
172        // from reading/writing state or connecting to the IPC socket.
173        #[cfg(unix)]
174        fix_state_dir_permissions();
175
176        let pid = std::process::id();
177        // Ensure PROCS has data for the supervisor PID before upsert_daemon reads title()
178        PROCS.refresh_pids(&[pid]);
179        // Determine container mode: CLI flag takes priority, then settings
180        let container_mode = container || settings().supervisor.container;
181        if container_mode {
182            info!("Starting supervisor in container/PID1 mode with pid {pid}");
183        } else {
184            info!("Starting supervisor with pid {pid}");
185        }
186
187        // If the previous supervisor died uncleanly, its daemon child processes
188        // may still be alive (orphaned, re-parented to init).  Terminate them
189        // before we take over so we don't end up with duplicate processes
190        // holding the same ports.
191        cleanup_orphaned_daemons(self).await;
192
193        self.upsert_daemon(
194            UpsertDaemonOpts::builder(DaemonId::pitchfork())
195                .set(|o| {
196                    o.pid = Some(pid);
197                    o.status = DaemonStatus::Running;
198                })
199                .build(),
200        )
201        .await?;
202        #[cfg(unix)]
203        fix_state_dir_permissions();
204
205        // If this is a boot start, automatically start boot_start daemons
206        if is_boot {
207            info!("Boot start mode enabled, starting boot_start daemons");
208            self.start_boot_daemons().await?;
209        }
210
211        self.interval_watch()?;
212        self.cron_watch()?;
213        self.signals()?;
214        self.daemon_file_watch()?;
215
216        // In container mode, install SIGCHLD handler to reap orphaned/zombie processes
217        #[cfg(unix)]
218        if container_mode {
219            self.reap_zombies()?;
220        }
221
222        // Start web server: CLI --web-port takes priority, then settings.web.auto_start + bind_port
223        let s = settings();
224        let effective_port = web_port.or_else(|| {
225            if s.web.auto_start {
226                match u16::try_from(s.web.bind_port).ok().filter(|&p| p > 0) {
227                    Some(p) => Some(p),
228                    None => {
229                        error!(
230                            "web.bind_port {} is out of valid port range (1-65535), web UI disabled",
231                            s.web.bind_port
232                        );
233                        None
234                    }
235                }
236            } else {
237                None
238            }
239        });
240        // CLI --web-path takes priority, then settings.web.base_path
241        let effective_path = web_path.or_else(|| {
242            let bp = s.web.base_path.clone();
243            if bp.is_empty() { None } else { Some(bp) }
244        });
245        if let Some(port) = effective_port {
246            tokio::spawn(async move {
247                if let Err(e) = crate::web::serve(port, effective_path).await {
248                    error!("Web server error: {e}");
249                }
250            });
251        }
252
253        // Start standalone API server if configured
254        let api_port = if s.api.auto_start {
255            match u16::try_from(s.api.bind_port).ok().filter(|&p| p > 0) {
256                Some(p) => Some(p),
257                None => {
258                    error!(
259                        "api.bind_port {} is out of valid port range (1-65535), API server disabled",
260                        s.api.bind_port
261                    );
262                    None
263                }
264            }
265        } else {
266            None
267        };
268        if let Some(port) = api_port {
269            tokio::spawn(async move {
270                if let Err(e) = crate::web::serve_api(port, None).await {
271                    error!("API server error: {e}");
272                }
273            });
274        }
275
276        // Start reverse proxy server if enabled
277        if s.proxy.enable {
278            // Pre-generate the TLS certificate synchronously before spawning the proxy
279            // task. This ensures the cert exists immediately after `sup start` returns,
280            // so `proxy trust` can be run right away without waiting for the async task.
281            #[cfg(feature = "proxy-tls")]
282            if s.proxy.https {
283                let proxy_dir = crate::env::PITCHFORK_STATE_DIR.join("proxy");
284                let ca_cert_path = proxy_dir.join("ca.pem");
285                let ca_key_path = proxy_dir.join("ca-key.pem");
286                if !ca_cert_path.exists() || !ca_key_path.exists() {
287                    match crate::proxy::server::generate_ca(&ca_cert_path, &ca_key_path) {
288                        Ok(()) => {
289                            info!(
290                                "Generated local CA certificate at {}",
291                                ca_cert_path.display()
292                            );
293                        }
294                        Err(e) => {
295                            error!("Failed to generate CA certificate: {e}");
296                        }
297                    }
298                }
299
300                // Auto-trust: attempt to install the CA certificate into the
301                // system trust store. May fail silently due to permissions;
302                // user can run `pitchfork proxy trust` manually.
303                if s.proxy.auto_trust && ca_cert_path.exists() {
304                    use crate::proxy::trust::{AutoTrustResult, auto_trust};
305                    match auto_trust(&ca_cert_path) {
306                        AutoTrustResult::AlreadyTrusted => {}
307                        AutoTrustResult::Trusted => {
308                            info!("CA certificate auto-trusted in system store");
309                        }
310                        AutoTrustResult::NotTrusted { reason } => {
311                            warn!("Auto-trust skipped: {reason}");
312                            warn!("Run `pitchfork proxy trust` to install manually");
313                        }
314                    }
315                }
316            }
317            // Spawn the proxy server and wait for its bind result via a oneshot
318            // channel.  This avoids the TOCTOU race of a pre-flight bind check
319            // while still surfacing binding failures immediately.
320            let (bind_tx, bind_rx) = tokio::sync::oneshot::channel();
321            let proxy_cancel = tokio_util::sync::CancellationToken::new();
322            let proxy_cancel_clone = proxy_cancel.clone();
323            *self.proxy_cancel.lock().await = Some(proxy_cancel);
324            let proxy_task = tokio::spawn(async move {
325                if let Err(e) = crate::proxy::server::serve(bind_tx, proxy_cancel_clone).await {
326                    error!("Proxy server error: {e}");
327                }
328            });
329            *self.proxy_task.lock().await = Some(proxy_task);
330            match bind_rx.await {
331                Ok(Ok(())) => {
332                    info!("Proxy server bound successfully");
333                    self.start_mdns().await;
334                }
335                Ok(Err(msg)) => {
336                    error!("{msg}");
337                    self.add_notification(log::LevelFilter::Error, msg).await;
338                }
339                Err(_) => {
340                    // Sender dropped without sending — serve() panicked or
341                    // returned before signalling.  Already logged by the
342                    // spawn error handler above.
343                }
344            }
345        }
346
347        // Pre-warm slug cache so the first /api/proxies request is fast.
348        // Spawned as a background task so it does not block startup.
349        tokio::spawn(async {
350            crate::proxy::server::get_cached_slugs().await;
351        });
352
353        let (ipc, ipc_handle) = IpcServer::new()?;
354        *self.ipc_shutdown.lock().await = Some(ipc_handle);
355        self.start_state_flush_task();
356        self.conn_watch(ipc).await
357    }
358
359    /// Start mDNS publishing for LAN mode (called after the proxy binds successfully).
360    async fn start_mdns(&self) {
361        let s = crate::settings::settings();
362        let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
363        if !s.proxy.enable || !lan_enabled {
364            return;
365        }
366
367        let lan_ip = if !s.proxy.lan_ip.is_empty() {
368            match s.proxy.lan_ip.parse::<std::net::Ipv4Addr>() {
369                Ok(ip) => Some(ip),
370                Err(e) => {
371                    error!(
372                        "proxy.lan_ip {:?} is not a valid IPv4 address: {e}",
373                        s.proxy.lan_ip
374                    );
375                    return;
376                }
377            }
378        } else {
379            match crate::proxy::lan_ip::detect_lan_ip().await {
380                Some(ip) => Some(ip),
381                None => {
382                    error!(
383                        "LAN mode is enabled but no LAN IP address could be detected. \
384                         Set proxy.lan_ip to a specific address, or ensure you are connected to a network."
385                    );
386                    return;
387                }
388            }
389        };
390
391        let Some(lan_ip) = lan_ip else { return };
392        let port = u16::try_from(s.proxy.port).unwrap_or(443);
393
394        let Some(mut publisher) = crate::proxy::mdns::MdnsPublisher::new(lan_ip) else {
395            error!("Failed to start mDNS publisher. Is Avahi (Linux) or Bonjour (macOS) running?");
396            return;
397        };
398
399        // Publish all registered slugs.
400        let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
401        for slug in slugs.keys() {
402            let hostname = format!("{slug}.local");
403            publisher.publish(&hostname, port);
404        }
405
406        log::info!(
407            "LAN mode: mDNS publishing on {lan_ip}, {} slug(s) registered",
408            slugs.len()
409        );
410
411        let publisher = std::sync::Arc::new(tokio::sync::Mutex::new(publisher));
412
413        // Start the IP monitor (only when IP is auto-detected, not pinned).
414        let ip_pinned = !s.proxy.lan_ip.is_empty();
415        if !ip_pinned {
416            let monitor_cancel = self.proxy_cancel.lock().await.clone();
417            let publisher_clone = publisher.clone();
418            let task = tokio::spawn(async move {
419                let mut last_ip = lan_ip;
420                let interval = std::time::Duration::from_secs(5);
421                let mut ticker = tokio::time::interval(interval);
422                ticker.tick().await; // first tick is immediate
423                loop {
424                    ticker.tick().await;
425                    if let Some(cancel) = monitor_cancel.as_ref() {
426                        if cancel.is_cancelled() {
427                            break;
428                        }
429                    }
430                    if let Some(new_ip) =
431                        crate::proxy::lan_ip::detect_lan_ip_if_changed(last_ip).await
432                    {
433                        log::info!("LAN IP changed: {last_ip} → {new_ip}");
434                        last_ip = new_ip;
435                        let mut pub_guard = publisher_clone.lock().await;
436                        pub_guard.republish_all(new_ip, port);
437                    }
438                }
439            });
440            *self.lan_monitor_task.lock().await = Some(task);
441        }
442
443        *self.mdns_publisher.lock().await = Some(publisher);
444    }
445
446    /// Re-read slugs from config and update mDNS records.
447    ///
448    /// Publishes new slugs and unpublishes removed ones. Called via IPC when
449    /// `proxy add` or `proxy remove` modifies the slug registry.
450    async fn sync_mdns(&self) {
451        // Clone the Arc and release the outer lock immediately so we don't
452        // block close() from taking the publisher during shutdown.
453        let publisher = {
454            let guard = self.mdns_publisher.lock().await;
455            match guard.as_ref() {
456                Some(p) => p.clone(),
457                None => {
458                    debug!("sync_mdns: mDNS publisher not active, skipping");
459                    return;
460                }
461            }
462        };
463
464        let s = crate::settings::settings();
465        let port = u16::try_from(s.proxy.port).unwrap_or(443);
466
467        let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
468        let mut pub_guard = publisher.lock().await;
469
470        // Unpublish slugs that no longer exist in config.
471        let current_keys: Vec<&String> = slugs.keys().collect();
472        let registered: Vec<String> = pub_guard.registered_hostnames();
473        for hostname in &registered {
474            // hostname is "slug.local" — extract slug part.
475            let slug = hostname.strip_suffix(".local").unwrap_or(hostname);
476            if !current_keys.iter().any(|k| k.as_str() == slug) {
477                log::info!("mDNS: unpublishing removed slug {slug}");
478                pub_guard.unpublish(hostname);
479            }
480        }
481
482        // Publish new slugs that aren't yet registered.
483        for slug in slugs.keys() {
484            let hostname = format!("{slug}.local");
485            if !pub_guard.is_published(&hostname) {
486                log::info!("mDNS: publishing new slug {slug}");
487                pub_guard.publish(&hostname, port);
488            }
489        }
490    }
491
492    /// Spawn a background task that periodically flushes the state file to
493    /// disk if it has been marked dirty.  Uses debouncing (1s interval) to
494    /// batch rapid state changes.
495    fn start_state_flush_task(&self) {
496        let cancel = tokio_util::sync::CancellationToken::new();
497        *self.flush_cancel.lock().unwrap() = Some(cancel.clone());
498        tokio::spawn(async move {
499            let mut interval = time::interval(Duration::from_secs(1));
500            interval.set_missed_tick_behavior(time::MissedTickBehavior::Skip);
501            loop {
502                tokio::select! {
503                    _ = interval.tick() => {}
504                    _ = cancel.cancelled() => {
505                        debug!("state flush task received shutdown signal");
506                        break;
507                    }
508                }
509                let state = SUPERVISOR.state_file.lock().await;
510                if state.is_dirty() {
511                    if let Err(e) = state.write() {
512                        warn!("failed to flush state file: {e}");
513                    }
514                }
515            }
516            debug!("state flush task exiting");
517        });
518    }
519
520    pub(crate) async fn flush_state(&self) {
521        let state = self.state_file.lock().await;
522        if state.is_dirty() {
523            if let Err(e) = state.write() {
524                warn!("failed to flush state file: {e}");
525            }
526        }
527    }
528
529    pub(crate) async fn refresh(&self) -> Result<()> {
530        trace!("refreshing");
531
532        let dirs_with_pids = self.get_dirs_with_shell_pids().await;
533
534        let mut last_refreshed_at = self.last_refreshed_at.lock().await;
535        *last_refreshed_at = time::Instant::now();
536
537        for (dir, pids) in dirs_with_pids {
538            let to_remove = pids
539                .iter()
540                .filter(|pid| !PROCS.is_running(**pid))
541                .collect::<Vec<_>>();
542            for pid in &to_remove {
543                self.remove_shell_pid(**pid).await?
544            }
545            if to_remove.len() == pids.len() {
546                self.leave_dir(&dir).await?;
547            }
548        }
549
550        self.check_retry().await?;
551        self.process_pending_autostops().await?;
552
553        Ok(())
554    }
555
556    /// Install a SIGCHLD handler that reaps orphaned zombie child processes.
557    ///
558    /// When running as PID 1 inside a container, orphaned processes are
559    /// re-parented to PID 1. Without explicit reaping, they accumulate
560    /// as zombies in the process table indefinitely.
561    ///
562    /// Only reaps processes that are NOT managed by the supervisor (i.e.
563    /// not tracked in the state file). Managed daemon processes are reaped
564    /// by their monitoring tasks via `child.wait()`.
565    ///
566    /// ## Strategy
567    ///
568    /// **Linux**: Uses `waitid(Id::All, WNOHANG | WNOWAIT | WEXITED)` to
569    /// *peek* at the next zombie without consuming its status. If the PID
570    /// belongs to a managed daemon, the reaper skips it so Tokio's
571    /// `child.wait()` can collect the status normally. Only unmanaged
572    /// orphans are actually reaped (via `waitpid(Pid, WNOHANG)`). This
573    /// eliminates the race entirely.
574    ///
575    /// **Non-Linux Unix** (e.g. macOS — mainly for local development;
576    /// container mode targets Linux): `waitid` is unavailable, so we fall
577    /// back to `waitpid(None, WNOHANG)`. If the reaper accidentally
578    /// consumes a managed PID's status, it stashes the exit code in
579    /// [`REAPED_STATUSES`] for the monitoring task to recover.
580    #[cfg(unix)]
581    fn reap_zombies(&self) -> Result<()> {
582        let mut stream = signal::unix::signal(SignalKind::child())
583            .map_err(|e| miette::miette!("Failed to register SIGCHLD handler: {e}"))?;
584        tokio::spawn(async move {
585            loop {
586                stream.recv().await;
587                // Collect PIDs of managed daemons so we don't steal their exit status
588                let managed_pids: HashSet<u32> = SUPERVISOR
589                    .state_file
590                    .lock()
591                    .await
592                    .daemons
593                    .values()
594                    .filter_map(|d| d.pid)
595                    .collect();
596                // Reap all available zombie children that are NOT managed
597                Self::reap_unmanaged_zombies(&managed_pids).await;
598            }
599        });
600        info!("container mode: SIGCHLD zombie reaper installed");
601        Ok(())
602    }
603
604    /// Linux implementation: peek with `waitid(WNOWAIT)` then selectively reap.
605    ///
606    /// `WNOWAIT` leaves the zombie in the table so we can inspect its PID
607    /// without consuming the exit status. Only if the PID is *not* managed
608    /// do we call `waitpid(Pid, WNOHANG)` to actually reap it.
609    #[cfg(target_os = "linux")]
610    async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
611        use nix::sys::wait::{Id, WaitPidFlag, WaitStatus, waitid, waitpid};
612        use nix::unistd::Pid;
613
614        loop {
615            // Peek at the next zombie without consuming it
616            let peek_flags = WaitPidFlag::WNOHANG | WaitPidFlag::WNOWAIT | WaitPidFlag::WEXITED;
617            match waitid(Id::All, peek_flags) {
618                Ok(WaitStatus::StillAlive) => break,
619                Ok(status) => {
620                    let Some(pid_raw) = status.pid().map(|p| p.as_raw() as u32) else {
621                        break;
622                    };
623                    if managed_pids.contains(&pid_raw) {
624                        // This is a managed daemon — leave it for Tokio's child.wait().
625                        // We must break out of the loop because waitid(Id::All) would
626                        // keep returning the same zombie if we don't consume it.
627                        trace!(
628                            "zombie reaper: skipping managed daemon pid {pid_raw}, \
629                             leaving for Tokio to reap"
630                        );
631                        break;
632                    }
633                    // Not managed — actually reap it
634                    match waitpid(Pid::from_raw(pid_raw as i32), Some(WaitPidFlag::WNOHANG)) {
635                        Ok(s) => trace!("reaped orphaned zombie child: {s:?}"),
636                        Err(nix::errno::Errno::ECHILD) => break,
637                        Err(e) => {
638                            trace!("waitpid error reaping pid {pid_raw}: {e}");
639                            break;
640                        }
641                    }
642                }
643                Err(nix::errno::Errno::ECHILD) => break, // no children at all
644                Err(e) => {
645                    trace!("waitid error in zombie reaper: {e}");
646                    break;
647                }
648            }
649        }
650    }
651
652    /// Non-Linux fallback: blind `waitpid(None, WNOHANG)` with stash recovery.
653    ///
654    /// Since `waitid(WNOWAIT)` is not available, we cannot peek. If we
655    /// accidentally reap a managed PID, we stash the exit code in
656    /// [`REAPED_STATUSES`] so the monitoring task can recover it.
657    #[cfg(all(unix, not(target_os = "linux")))]
658    async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
659        use nix::sys::wait::{WaitPidFlag, WaitStatus, waitpid};
660
661        loop {
662            match waitpid(None, Some(WaitPidFlag::WNOHANG)) {
663                Ok(WaitStatus::StillAlive) => break,
664                Ok(status) => {
665                    let Some(pid) = status.pid().map(|p| p.as_raw() as u32) else {
666                        continue;
667                    };
668                    if managed_pids.contains(&pid) {
669                        // Race lost — stash the exit code for lifecycle recovery
670                        let exit_code = match status {
671                            WaitStatus::Exited(_, code) => code,
672                            WaitStatus::Signaled(_, sig, _) => -(sig as i32),
673                            _ => -1,
674                        };
675                        warn!(
676                            "zombie reaper reaped managed daemon pid {pid} \
677                             (exit_code={exit_code}); stashing status for recovery"
678                        );
679                        REAPED_STATUSES.lock().await.insert(pid, exit_code);
680                    } else {
681                        trace!("reaped orphaned zombie child: {status:?}");
682                    }
683                }
684                Err(nix::errno::Errno::ECHILD) => break, // no more children
685                Err(e) => {
686                    trace!("waitpid error in zombie reaper: {e}");
687                    break;
688                }
689            }
690        }
691    }
692
693    #[cfg(unix)]
694    fn signals(&self) -> Result<()> {
695        let signals = [
696            SignalKind::terminate(),
697            SignalKind::alarm(),
698            SignalKind::interrupt(),
699            SignalKind::quit(),
700            SignalKind::hangup(),
701            SignalKind::user_defined1(),
702            SignalKind::user_defined2(),
703        ];
704        static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
705        for signal in signals {
706            let stream = match signal::unix::signal(signal) {
707                Ok(s) => s,
708                Err(e) => {
709                    warn!("Failed to register signal handler for {signal:?}: {e}");
710                    continue;
711                }
712            };
713            tokio::spawn(async move {
714                let mut stream = stream;
715                loop {
716                    stream.recv().await;
717                    if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
718                        exit(1);
719                    } else {
720                        SUPERVISOR.handle_signal().await;
721                    }
722                }
723            });
724        }
725        Ok(())
726    }
727
728    #[cfg(windows)]
729    fn signals(&self) -> Result<()> {
730        tokio::spawn(async move {
731            static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
732            loop {
733                if let Err(e) = signal::ctrl_c().await {
734                    error!("Failed to wait for ctrl-c: {}", e);
735                    return;
736                }
737                if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
738                    exit(1);
739                } else {
740                    SUPERVISOR.handle_signal().await;
741                }
742            }
743        });
744        Ok(())
745    }
746
747    async fn handle_signal(&self) {
748        info!("received signal, stopping");
749        self.close().await;
750        exit(0)
751    }
752
753    pub(crate) async fn close(&self) {
754        // Signal the proxy server to stop accepting new connections
755        // and drain in-flight ones, *before* stopping daemons so the
756        // proxy has time to finish forwarding active requests.
757        if let Some(cancel) = self.proxy_cancel.lock().await.take() {
758            cancel.cancel();
759        }
760
761        // Stop the LAN IP monitor task.
762        if let Some(monitor_task) = self.lan_monitor_task.lock().await.take() {
763            monitor_task.abort();
764        }
765
766        // Shutdown the mDNS publisher (sends goodbye packets).
767        if let Some(publisher) = self.mdns_publisher.lock().await.take() {
768            publisher.lock().await.shutdown();
769        }
770
771        if let Some(proxy_task) = self.proxy_task.lock().await.take() {
772            let _ = tokio::time::timeout(Duration::from_secs(12), proxy_task).await;
773        }
774
775        // Clean up /etc/hosts entries managed by pitchfork
776        let s = settings();
777        if s.proxy.enable && s.proxy.sync_hosts {
778            crate::proxy::hosts::clean_hosts_file();
779        }
780
781        let pitchfork_id = DaemonId::pitchfork();
782        let active = self.active_daemons().await;
783        let active_ids: Vec<DaemonId> = active
784            .iter()
785            .filter(|d| d.id != pitchfork_id)
786            .map(|d| d.id.clone())
787            .collect();
788
789        // Stop daemons in reverse dependency order.
790        // If dependency resolution fails (e.g. config changed), fall back to
791        // stopping in arbitrary order so we still shut down cleanly.
792        // Daemons within the same level are stopped concurrently.
793        let stop_levels = compute_reverse_stop_order(&active_ids);
794        for level in &stop_levels {
795            let mut tasks = Vec::new();
796            for id in level {
797                let id = id.clone();
798                tasks.push(tokio::spawn(async move {
799                    if let Err(err) = SUPERVISOR.stop(&id).await {
800                        error!("failed to stop daemon {id}: {err}");
801                    }
802                }));
803            }
804            for task in tasks {
805                let _ = task.await;
806            }
807        }
808        let _ = self.remove_daemon(&pitchfork_id).await;
809
810        // Signal the background state flush task to exit so it doesn't
811        // keep waking up and acquiring the state mutex after shutdown.
812        if let Some(cancel) = self.flush_cancel.lock().unwrap().take() {
813            cancel.cancel();
814        }
815
816        // Force-flush state to disk before shutting down IPC so no
817        // in-memory-only changes are lost.
818        {
819            let state = self.state_file.lock().await;
820            if state.is_dirty() {
821                if let Err(e) = state.write() {
822                    warn!("failed to flush state file during shutdown: {e}");
823                }
824            }
825        }
826
827        // Signal IPC server to shut down gracefully
828        if let Some(mut handle) = self.ipc_shutdown.lock().await.take() {
829            handle.shutdown();
830        }
831
832        // Wait for all in-flight monitoring tasks to finish registering their
833        // hook handles. Each monitoring task increments `active_monitors` when
834        // its process exits, and decrements it (+ notifies `monitor_done`)
835        // after all fire_hook() calls complete. This replaces the old
836        // yield_now() approach which had a race window.
837        let drain_timeout = time::sleep(Duration::from_secs(5));
838        tokio::pin!(drain_timeout);
839        loop {
840            if self.active_monitors.load(atomic::Ordering::Acquire) == 0 {
841                break;
842            }
843            tokio::select! {
844                _ = self.monitor_done.notified() => {}
845                _ = &mut drain_timeout => {
846                    warn!("timed out waiting for monitoring tasks to register hooks, proceeding with shutdown");
847                    break;
848                }
849            }
850        }
851        let handles: Vec<JoinHandle<()>> = std::mem::take(&mut *self.hook_tasks.lock().await);
852        let hook_timeout = Duration::from_secs(30);
853        for handle in handles {
854            match time::timeout(hook_timeout, handle).await {
855                Ok(_) => {} // Hook completed (success or error, doesn't matter)
856                Err(_) => {
857                    warn!(
858                        "hook task did not complete within {hook_timeout:?} during shutdown, skipping"
859                    );
860                }
861            }
862        }
863
864        let _ = fs::remove_dir_all(&*env::IPC_SOCK_DIR);
865    }
866
867    pub(crate) async fn add_notification(&self, level: log::LevelFilter, message: String) {
868        self.pending_notifications
869            .lock()
870            .await
871            .push((level, message));
872    }
873}
874
875/// Fix ownership on the state directory so non-root users can access files
876/// created by a `sudo`-started supervisor.
877///
878/// When `[settings.supervisor] user` or `SUDO_UID`/`SUDO_GID` are set, we
879/// `chown` the state directory and safe subdirectories back to that non-root
880/// runtime user. This is strictly better than `chmod 0o666` because it does not
881/// widen the permission bits — the files stay owner-only (0o600/0o700) but the
882/// *owner* is the user that daemon processes and CLI clients need to share.
883///
884/// **Security**: The `proxy/` subtree is intentionally skipped. It contains
885/// `ca-key.pem` which must remain `0o600` and owned by the process that
886/// generated it. Changing its ownership or permissions would expose the CA
887/// private key to other local users.
888///
889/// If neither `user` nor `SUDO_UID`/`SUDO_GID` are available (e.g. direct
890/// root login), we fall back to relaxing permissions on only the `sock/` and
891/// `logs/` subdirectories (plus `state.toml`) so CLI clients can still function.
892#[cfg(unix)]
893fn fix_state_dir_permissions() {
894    let state_dir = &*env::PITCHFORK_STATE_DIR;
895    if let Some((uid, gid)) = state_owner_ids() {
896        if !state_dir.exists()
897            && let Err(err) = fs::create_dir_all(state_dir)
898        {
899            warn!(
900                "failed to create state directory for ownership fix at {}: {err}",
901                state_dir.display()
902            );
903            return;
904        }
905
906        // Best path: chown back to the runtime user. Permissions stay tight.
907        chown_recursive(state_dir, uid, gid, true);
908        debug!(
909            "chowned state directory to uid={uid} gid={gid} at {}",
910            state_dir.display()
911        );
912    } else {
913        if !state_dir.exists() {
914            return;
915        }
916
917        // Fallback: relax permissions on safe subdirectories only.
918        // proxy/ is never touched.
919        chmod_safe_subtrees(state_dir);
920        debug!(
921            "relaxed permissions on safe subtrees at {}",
922            state_dir.display()
923        );
924    }
925}
926
927#[cfg(unix)]
928pub(crate) fn state_owner_ids() -> Option<(u32, u32)> {
929    if !nix::unistd::Uid::effective().is_root() {
930        return None;
931    }
932
933    let s = settings();
934    let user = s.supervisor.user.trim();
935    if !user.is_empty() {
936        return resolve_supervisor_user_ids(user).or_else(|| {
937            warn!(
938                "failed to resolve supervisor.user '{user}' for state ownership; falling back to SUDO_UID/SUDO_GID"
939            );
940            parse_sudo_ids()
941        });
942    }
943
944    parse_sudo_ids()
945}
946
947#[cfg(unix)]
948fn resolve_supervisor_user_ids(user: &str) -> Option<(u32, u32)> {
949    let user_record = if user.chars().all(|c| c.is_ascii_digit()) {
950        let uid = user.parse::<u32>().ok()?;
951        nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
952            .ok()
953            .flatten()
954    } else {
955        nix::unistd::User::from_name(user).ok().flatten()
956    }?;
957
958    Some((user_record.uid.as_raw(), user_record.gid.as_raw()))
959}
960
961/// Parse `SUDO_UID` and `SUDO_GID` environment variables into numeric IDs.
962///
963/// Returns `None` unless the effective UID is 0 (root). This prevents stale
964/// `SUDO_UID`/`SUDO_GID` values inherited into non-sudo environments from
965/// triggering incorrect `chown` operations.
966#[cfg(unix)]
967fn parse_sudo_ids() -> Option<(u32, u32)> {
968    if !nix::unistd::Uid::effective().is_root() {
969        return None;
970    }
971    let uid: u32 = std::env::var("SUDO_UID").ok()?.parse().ok()?;
972    let gid: u32 = std::env::var("SUDO_GID").ok()?.parse().ok()?;
973    Some((uid, gid))
974}
975
976/// Recursively `chown` a directory tree. If `skip_proxy` is true, the `proxy/`
977/// subdirectory is skipped entirely to protect the CA private key.
978#[cfg(unix)]
979fn chown_recursive(dir: &std::path::Path, uid: u32, gid: u32, skip_proxy: bool) {
980    // chown the directory itself
981    let _ = chown_path(dir, uid, gid);
982
983    let entries = match std::fs::read_dir(dir) {
984        Ok(e) => e,
985        Err(_) => return,
986    };
987    for entry in entries.flatten() {
988        let path = entry.path();
989        if path.is_dir() {
990            // Skip proxy/ at the top level of the state directory
991            if skip_proxy {
992                if let Some(name) = path.file_name().and_then(|n| n.to_str()) {
993                    if name == "proxy" {
994                        continue;
995                    }
996                }
997            }
998            chown_recursive(&path, uid, gid, false);
999        } else {
1000            let _ = chown_path(&path, uid, gid);
1001        }
1002    }
1003}
1004
1005/// `chown` a single path using libc. Returns Ok(()) on success.
1006#[cfg(unix)]
1007fn chown_path(path: &std::path::Path, uid: u32, gid: u32) -> std::io::Result<()> {
1008    use std::ffi::CString;
1009    use std::os::unix::ffi::OsStrExt;
1010    let c_path = CString::new(path.as_os_str().as_bytes())
1011        .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
1012    let ret = unsafe { libc::chown(c_path.as_ptr(), uid, gid) };
1013    if ret == 0 {
1014        Ok(())
1015    } else {
1016        Err(std::io::Error::last_os_error())
1017    }
1018}
1019
1020/// Fallback: relax permissions on safe subdirectories only (sock/, logs/, and
1021/// state.toml). The proxy/ subtree is never touched.
1022#[cfg(unix)]
1023fn chmod_safe_subtrees(state_dir: &std::path::Path) {
1024    // The state directory itself needs to be traversable
1025    let _ = fs::set_permissions(state_dir, fs::Permissions::from_mode(0o755));
1026
1027    // state.toml — needs to be readable by CLI clients
1028    let state_file = state_dir.join("state.toml");
1029    if state_file.exists() {
1030        let _ = fs::set_permissions(&state_file, fs::Permissions::from_mode(0o644));
1031    }
1032
1033    // Safe subdirectories: sock/ and logs/
1034    for subdir_name in &["sock", "logs"] {
1035        let subdir = state_dir.join(subdir_name);
1036        if subdir.is_dir() {
1037            chmod_recursive(&subdir);
1038        }
1039    }
1040}
1041
1042/// On startup, kill any daemon processes left behind by a previous supervisor
1043/// that was terminated unexpectedly (e.g. `kill -9`).
1044///
1045/// This iterates the state file for daemon entries with a recorded PID.  If the
1046/// PID is still alive and the process name matches the recorded `title`, it is
1047/// assumed to be an orphan from the previous supervisor session and is killed.
1048/// The daemon's state is then reset to `Stopped` with no PID.
1049///
1050/// The process-name check protects against the rare case where a PID has been
1051/// recycled by an unrelated process since the state file was written.
1052///
1053/// This is gated by the `supervisor.cleanup_orphans` setting (default: true).
1054async fn cleanup_orphaned_daemons(supervisor: &Supervisor) {
1055    if !settings().supervisor.cleanup_orphans {
1056        return;
1057    }
1058
1059    let candidates: Vec<_> = {
1060        let state = supervisor.state_file.lock().await;
1061        state
1062            .daemons
1063            .values()
1064            .filter(|d| d.id != DaemonId::pitchfork() && d.pid.is_some())
1065            .cloned()
1066            .collect()
1067    };
1068
1069    if candidates.is_empty() {
1070        return;
1071    }
1072
1073    info!(
1074        "checking {} daemon(s) for orphaned processes",
1075        candidates.len()
1076    );
1077
1078    for daemon in candidates {
1079        let Some(pid) = daemon.pid else { continue };
1080
1081        if !PROCS.is_running(pid) {
1082            // PID already dead — just reset its state
1083            let _ = supervisor
1084                .upsert_daemon(
1085                    UpsertDaemonOpts::builder(daemon.id.clone())
1086                        .set(|o| {
1087                            o.pid = None;
1088                            o.status = DaemonStatus::Stopped;
1089                            o.active_port = None;
1090                        })
1091                        .build(),
1092                )
1093                .await;
1094            continue;
1095        }
1096
1097        // Safety check: verify the process name matches what we recorded.
1098        // This avoids accidentally killing an unrelated process if the PID
1099        // was recycled by the kernel between state-file writes.
1100        let current_title = PROCS.title(pid);
1101        let matches = match (&current_title, &daemon.title) {
1102            (Some(current), Some(expected)) => current == expected,
1103            // If we don't have a recorded title, fall back to allowing the kill
1104            // (this is a degraded but functional state — the process was probably
1105            // started before title tracking was added, or the state was reset).
1106            _ => true,
1107        };
1108
1109        if !matches {
1110            warn!(
1111                "pid {pid} for daemon {} has changed name (expected '{}', found '{}'); skipping orphan cleanup",
1112                daemon.id,
1113                daemon.title.as_deref().unwrap_or("?"),
1114                current_title.as_deref().unwrap_or("?")
1115            );
1116            // Still reset state so we don't track a stale PID
1117            let _ = supervisor
1118                .upsert_daemon(
1119                    UpsertDaemonOpts::builder(daemon.id.clone())
1120                        .set(|o| {
1121                            o.pid = None;
1122                            o.status = DaemonStatus::Stopped;
1123                            o.active_port = None;
1124                        })
1125                        .build(),
1126                )
1127                .await;
1128            continue;
1129        }
1130
1131        info!("terminating orphaned daemon {} (pid {pid})", daemon.id);
1132
1133        let stop_cfg = daemon.stop_signal.unwrap_or_default();
1134        let _ = PROCS
1135            .kill_process_group_async(pid, stop_cfg.signal.into(), stop_cfg.timeout)
1136            .await;
1137
1138        let _ = supervisor
1139            .upsert_daemon(
1140                UpsertDaemonOpts::builder(daemon.id.clone())
1141                    .set(|o| {
1142                        o.pid = None;
1143                        o.status = DaemonStatus::Stopped;
1144                        o.active_port = None;
1145                    })
1146                    .build(),
1147            )
1148            .await;
1149    }
1150}
1151
1152/// Recursively chmod: directories → 0o755, files → 0o644.
1153#[cfg(unix)]
1154fn chmod_recursive(dir: &std::path::Path) {
1155    let _ = fs::set_permissions(dir, fs::Permissions::from_mode(0o755));
1156    let entries = match fs::read_dir(dir) {
1157        Ok(e) => e,
1158        Err(_) => return,
1159    };
1160    for entry in entries.flatten() {
1161        let path = entry.path();
1162        if path.is_dir() {
1163            chmod_recursive(&path);
1164        } else {
1165            let _ = fs::set_permissions(&path, fs::Permissions::from_mode(0o644));
1166        }
1167    }
1168}