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 ®istered {
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}