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