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