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 // Ignoring Ctrl+C is inherited, so a supervisor started from a process
628 // that ignores it would neither see Ctrl+C itself nor let its daemons
629 // see it: a daemon with `stop_signal = "SIGINT"` would never get the
630 // Ctrl+C sent to stop it. Handle it again before any daemon starts.
631 #[cfg(windows)]
632 if unsafe { windows_sys::Win32::System::Console::SetConsoleCtrlHandler(None, 0) } == 0 {
633 warn!(
634 "failed to stop ignoring Ctrl+C: {}",
635 std::io::Error::last_os_error()
636 );
637 }
638
639 // Refuse to run beside a supervisor that is already listening, and
640 // do so before recording ourselves in the state file or starting any
641 // daemons: taking over its socket would leave it running but
642 // unreachable. (`--force` has already waited for the one it replaced
643 // to let go of the socket.) The lock is held until our own listener
644 // is bound, so a supervisor starting at the same time waits here and
645 // then finds this one listening.
646 let startup_lock = StartupLock::acquire().await?;
647 if crate::ipc::supervisor_listening().await {
648 return Err(miette::miette!(
649 "another pitchfork supervisor is already listening on {}",
650 crate::ipc::socket_display()
651 ));
652 }
653
654 let pid = std::process::id();
655 // Ensure PROCS has data for the supervisor PID before upsert_daemon reads title()
656 PROCS.refresh_pids(&[pid]);
657 // Determine container mode: CLI flag takes priority, then settings.
658 // Running as PID 1 always enables it: orphaned descendants of daemons
659 // re-parent to us, and without the zombie reaper they would accumulate
660 // as unreaped zombies — which also keep their process group alive,
661 // stalling whole-group stop waits indefinitely.
662 let container_mode =
663 container || settings().supervisor.container || std::process::id() == 1;
664 if container_mode {
665 info!("Starting supervisor in container/PID1 mode with pid {pid}");
666 } else {
667 info!("Starting supervisor with pid {pid}");
668 }
669
670 // Whether the previous supervisor exited uncleanly must be read before
671 // we record ourselves in the state file just below (see
672 // `supervisor_exited_uncleanly`); the background cleanup task runs
673 // after that record exists, so it receives the answer instead of
674 // reading it too late.
675 let unclean = supervisor_exited_uncleanly(self).await;
676
677 self.upsert_daemon(
678 UpsertDaemonOpts::builder(DaemonId::pitchfork())
679 .set(|o| {
680 o.pid = Some(pid);
681 o.status = DaemonStatus::Running;
682 })
683 .build(),
684 )
685 .await?;
686 #[cfg(unix)]
687 fix_state_dir_permissions();
688
689 // Self-heal: if the boot registration points to a stale binary path
690 // (e.g. after a brew/mise upgrade), re-register with the current path.
691 // Runs in the background — must not block or fail supervisor startup.
692 tokio::task::spawn_blocking(|| {
693 if let Ok(boot_manager) = crate::boot_manager::BootManager::new() {
694 boot_manager.check_and_reregister_if_stale();
695 }
696 });
697
698 // If the previous supervisor died uncleanly, its daemon child processes
699 // may still be alive (orphaned, re-parented to init). Terminate them
700 // before starting replacements so we don't end up with duplicate
701 // processes holding the same ports.
702 //
703 // This runs in the background: each orphan kill now waits for its
704 // whole process group to exit (seconds per orphan), and doing that
705 // inline would delay IPC socket creation past the CLI's short connect
706 // budget on autostart. Per-daemon stop locks serialize the cleanup
707 // against any Run/Stop requests that arrive for the same daemon in the
708 // meantime, and boot daemons start after cleanup completes so they
709 // cannot observe an orphan as "already running".
710 let boot_after_cleanup = is_boot;
711 tokio::spawn(async move {
712 cleanup_orphaned_daemons(&SUPERVISOR, unclean).await;
713 if boot_after_cleanup {
714 info!("Boot start mode enabled, starting boot_start daemons");
715 if let Err(e) = SUPERVISOR.start_boot_daemons().await {
716 error!("failed to start boot daemons: {e}");
717 }
718 }
719 });
720
721 self.interval_watch()?;
722
723 // Run the first cron check synchronously before starting the cron
724 // watcher and IPC server. This registers config-only cron daemons and
725 // fires any `immediate=true` triggers in the foreground, so they cannot
726 // race with a concurrent `pitchfork start` IPC. By the time the cron
727 // watcher's first tick runs, `last_cron_triggered` is already anchored
728 // and the immediate daemons are already running.
729 if let Err(e) = self.check_cron_schedules().await {
730 error!("failed to check cron schedules on startup: {e}");
731 }
732
733 self.cron_watch()?;
734 self.signals()?;
735 self.daemon_file_watch()?;
736
737 // In container mode, install SIGCHLD handler to reap orphaned/zombie processes
738 #[cfg(unix)]
739 if container_mode {
740 self.reap_zombies()?;
741 }
742
743 // Start web server: CLI --web-port takes priority, then settings.web.auto_start + bind_port
744 let s = settings();
745 let effective_port = web_port.or_else(|| {
746 if s.web.auto_start {
747 match u16::try_from(s.web.bind_port).ok().filter(|&p| p > 0) {
748 Some(p) => Some(p),
749 None => {
750 error!(
751 "web.bind_port {} is out of valid port range (1-65535), web UI disabled",
752 s.web.bind_port
753 );
754 None
755 }
756 }
757 } else {
758 None
759 }
760 });
761 // CLI --web-path takes priority, then settings.web.base_path
762 let effective_path = web_path.or_else(|| {
763 let bp = s.web.base_path.clone();
764 if bp.is_empty() { None } else { Some(bp) }
765 });
766 if let Some(port) = effective_port {
767 tokio::spawn(async move {
768 if let Err(e) = crate::web::serve(port, effective_path).await {
769 error!("Web server error: {e}");
770 }
771 });
772 }
773
774 // Start standalone API server if configured
775 let api_port = if s.api.auto_start {
776 match u16::try_from(s.api.bind_port).ok().filter(|&p| p > 0) {
777 Some(p) => Some(p),
778 None => {
779 error!(
780 "api.bind_port {} is out of valid port range (1-65535), API server disabled",
781 s.api.bind_port
782 );
783 None
784 }
785 }
786 } else {
787 None
788 };
789 if let Some(port) = api_port {
790 tokio::spawn(async move {
791 if let Err(e) = crate::web::serve_api(port, None).await {
792 error!("API server error: {e}");
793 }
794 });
795 }
796
797 // Start reverse proxy server if enabled
798 if s.proxy.enable {
799 // Pre-generate the TLS certificate synchronously before spawning the proxy
800 // task. This ensures the cert exists immediately after `sup start` returns,
801 // so `proxy trust` can be run right away without waiting for the async task.
802 #[cfg(feature = "proxy-tls")]
803 if s.proxy.https {
804 let proxy_dir = crate::env::PITCHFORK_STATE_DIR.join("proxy");
805 let ca_cert_path = proxy_dir.join("ca.pem");
806 let ca_key_path = proxy_dir.join("ca-key.pem");
807 // Checked and written under the CA lock: `proxy setup` may be
808 // generating the same pair right now.
809 match crate::proxy::server::ensure_ca(&ca_cert_path, &ca_key_path, || {
810 ca_cert_path.exists() && ca_key_path.exists()
811 }) {
812 Ok(true) => {
813 info!(
814 "Generated local CA certificate at {}",
815 ca_cert_path.display()
816 );
817 }
818 Ok(false) => {}
819 Err(e) => {
820 error!("Failed to generate CA certificate: {e}");
821 }
822 }
823
824 // Auto-trust: attempt to install the CA certificate into the
825 // system trust store. May fail silently due to permissions;
826 // user can run `pitchfork proxy trust` manually.
827 if s.proxy.auto_trust && ca_cert_path.exists() {
828 use crate::proxy::trust::{AutoTrustResult, auto_trust};
829 match auto_trust(&ca_cert_path) {
830 AutoTrustResult::AlreadyTrusted => {}
831 AutoTrustResult::Trusted => {
832 info!("CA certificate auto-trusted in system store");
833 }
834 AutoTrustResult::NotTrusted { reason } => {
835 warn!("Auto-trust skipped: {reason}");
836 warn!("Run `pitchfork proxy trust` to install manually");
837 }
838 }
839 }
840 }
841 // Spawn the proxy server and wait for its bind result via a oneshot
842 // channel. This avoids the TOCTOU race of a pre-flight bind check
843 // while still surfacing binding failures immediately.
844 let (bind_tx, bind_rx) = tokio::sync::oneshot::channel();
845 let proxy_cancel = tokio_util::sync::CancellationToken::new();
846 let proxy_cancel_clone = proxy_cancel.clone();
847 *self.proxy_cancel.lock().await = Some(proxy_cancel);
848 let proxy_task = tokio::spawn(async move {
849 if let Err(e) = crate::proxy::server::serve(bind_tx, proxy_cancel_clone).await {
850 error!("Proxy server error: {e}");
851 }
852 });
853 *self.proxy_task.lock().await = Some(proxy_task);
854 match bind_rx.await {
855 Ok(Ok(())) => {
856 info!("Proxy server bound successfully");
857 // Resolved once and shared: mDNS and the DNS resolver
858 // must advertise the same address.
859 let lan_ip = self.resolve_lan_ip().await;
860 self.start_mdns(lan_ip).await;
861 self.start_dns_resolver(lan_ip).await;
862 }
863 Ok(Err(msg)) => {
864 error!("{msg}");
865 self.add_notification(log::LevelFilter::Error, msg).await;
866 }
867 Err(_) => {
868 // Sender dropped without sending — serve() panicked or
869 // returned before signalling. Already logged by the
870 // spawn error handler above.
871 }
872 }
873 }
874
875 // Pre-warm slug cache so the first /api/proxies request is fast.
876 // Spawned as a background task so it does not block startup.
877 tokio::spawn(async {
878 crate::proxy::server::get_cached_slugs().await;
879 });
880
881 let (ipc, ipc_handle) = IpcServer::new(startup_lock).await?;
882 *self.ipc_shutdown.lock().await = Some(ipc_handle);
883 self.start_state_flush_task();
884 self.conn_watch(ipc).await
885 }
886
887 /// Start the loopback DNS resolver for the proxy TLD.
888 ///
889 /// The resolver shares the proxy's cancellation token, so it stops with the
890 /// proxy. A bind failure is a notification rather than a fatal error: the
891 /// proxy still works for anyone who reaches it some other way.
892 async fn start_dns_resolver(&self, lan_ip: Option<std::net::Ipv4Addr>) {
893 let s = crate::settings::settings();
894 if !s.proxy.dns {
895 return;
896 }
897 let cfg = crate::proxy::dns::config_from_settings(&s, lan_ip);
898 let port = crate::proxy::dns::dns_port(&s);
899 // Loopback only: the resolver is for this machine's stub resolver, and
900 // LAN peers are served by mDNS instead.
901 let addr = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, port));
902
903 // The token's guard is held until the handle is stored. `close` takes
904 // the token under this lock before it collects `dns_task`, so it either
905 // runs first — and there is no token to spawn with — or waits until
906 // the handle is in place to be drained. Releasing it earlier left a
907 // window where `close` found no handle and the task outlived shutdown
908 // holding the resolver's sockets.
909 let cancel_guard = self.proxy_cancel.lock().await;
910 let Some(cancel) = cancel_guard.clone() else {
911 return;
912 };
913 let (bind_tx, bind_rx) = tokio::sync::oneshot::channel();
914 let task = tokio::spawn(async move {
915 if let Err(e) = crate::proxy::dns::serve(cfg, addr, bind_tx, cancel).await {
916 error!("DNS resolver error: {e}");
917 }
918 });
919 *self.dns_task.lock().await = Some(task);
920 drop(cancel_guard);
921 match bind_rx.await {
922 Ok(Ok(())) => info!("DNS resolver bound successfully"),
923 Ok(Err(msg)) => {
924 let msg = format!(
925 "{msg}\nProxy host names will not resolve through pitchfork. \
926 Choose another port with proxy.dns_port, or set proxy.dns = false."
927 );
928 error!("{msg}");
929 self.add_notification(log::LevelFilter::Error, msg).await;
930 }
931 Err(_) => {}
932 }
933 }
934
935 /// Watch the LAN address and keep the DNS responder, and mDNS when it is
936 /// running, pointed at the current one.
937 ///
938 /// Started even when the mDNS publisher could not be created: the DNS
939 /// responder serves this machine regardless, and an address it keeps
940 /// answering with after the interface has moved is worse than useless.
941 ///
942 /// Does nothing when `proxy.lan_ip` pins an address. That is a choice to
943 /// respect, not a starting point to drift from — the check lives here so
944 /// neither caller can forget it.
945 async fn start_lan_ip_monitor(
946 &self,
947 initial_ip: std::net::Ipv4Addr,
948 port: u16,
949 publisher: Option<std::sync::Arc<tokio::sync::Mutex<crate::proxy::mdns::MdnsPublisher>>>,
950 ) {
951 if !crate::settings::settings().proxy.lan_ip.is_empty() {
952 return;
953 }
954 // No cancellation token means `close` has already taken it, so shutdown
955 // is under way. Spawning here would leave a task polling with no way to
956 // stop it, and past the point where `close` collects the handle.
957 //
958 // Held until the handle is stored, for the reason `start_dns_resolver`
959 // gives: otherwise `close` can collect `lan_monitor_task` in between.
960 let cancel_guard = self.proxy_cancel.lock().await;
961 let Some(cancel) = cancel_guard.clone() else {
962 debug!("Not starting the LAN IP monitor: the supervisor is shutting down");
963 return;
964 };
965 let monitor_cancel = Some(cancel);
966 let task = tokio::spawn(async move {
967 let mut last_ip = initial_ip;
968 let mut ticker = tokio::time::interval(std::time::Duration::from_secs(5));
969 ticker.tick().await; // first tick is immediate
970 loop {
971 // Cancellation is raced against the tick, not checked after
972 // it. Checking afterwards means the task only notices once the
973 // full interval has elapsed, so a shutdown that waits a second
974 // for it always gives up and aborts instead — the graceful
975 // path would never once be taken.
976 match monitor_cancel.as_ref() {
977 Some(cancel) => {
978 tokio::select! {
979 _ = ticker.tick() => {}
980 _ = cancel.cancelled() => break,
981 }
982 }
983 None => {
984 ticker.tick().await;
985 }
986 }
987 if let Some(new_ip) = crate::proxy::lan_ip::detect_lan_ip_if_changed(last_ip).await
988 {
989 log::info!("LAN IP changed: {last_ip} → {new_ip}");
990 last_ip = new_ip;
991 crate::proxy::dns::update_lan_ip(new_ip);
992 if let Some(publisher) = publisher.as_ref() {
993 publisher.lock().await.republish_all(new_ip, port);
994 }
995 }
996 }
997 });
998 *self.lan_monitor_task.lock().await = Some(task);
999 drop(cancel_guard);
1000 }
1001
1002 /// Start mDNS publishing for LAN mode (called after the proxy binds successfully).
1003 /// The LAN address mDNS publishes and the DNS resolver answers with.
1004 ///
1005 /// Resolved once and handed to both. Detecting separately in each let them
1006 /// disagree when the interface address changed in between, which would
1007 /// advertise one address over mDNS and serve another over DNS, and probed
1008 /// the network twice at startup for one answer.
1009 ///
1010 /// `None` means LAN mode is off, or is on and the address could not be
1011 /// determined; either way the caller has nothing to publish. The reason is
1012 /// reported here so it is said once rather than by each caller.
1013 async fn resolve_lan_ip(&self) -> Option<std::net::Ipv4Addr> {
1014 let s = crate::settings::settings();
1015 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
1016 if !s.proxy.enable || !lan_enabled {
1017 return None;
1018 }
1019 if s.proxy.lan_ip.is_empty() {
1020 let detected = crate::proxy::lan_ip::detect_lan_ip().await;
1021 if detected.is_none() {
1022 error!(
1023 "LAN mode is enabled but no LAN IP address could be detected. \
1024 Set proxy.lan_ip to a specific address, or ensure you are connected to a network."
1025 );
1026 }
1027 return detected;
1028 }
1029 match s.proxy.lan_ip.parse::<std::net::Ipv4Addr>() {
1030 Ok(ip) => Some(ip),
1031 Err(e) => {
1032 let msg = format!(
1033 concat!(
1034 "proxy.lan_ip {:?} is not a valid IPv4 address: {}. ",
1035 "LAN mode will not start; fix the setting or clear it ",
1036 "to auto-detect."
1037 ),
1038 s.proxy.lan_ip, e
1039 );
1040 error!("{msg}");
1041 self.add_notification(log::LevelFilter::Error, msg).await;
1042 None
1043 }
1044 }
1045 }
1046
1047 async fn start_mdns(&self, lan_ip: Option<std::net::Ipv4Addr>) {
1048 let s = crate::settings::settings();
1049 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
1050 if !s.proxy.enable || !lan_enabled {
1051 return;
1052 }
1053
1054 let Some(lan_ip) = lan_ip else { return };
1055 let port = u16::try_from(s.proxy.port).unwrap_or(443);
1056
1057 let Some(mut publisher) = crate::proxy::mdns::MdnsPublisher::new(lan_ip) else {
1058 error!("Failed to start mDNS publisher. Is Avahi (Linux) or Bonjour (macOS) running?");
1059 // The DNS responder hands out this address too, and it is useful on
1060 // this machine whether or not mDNS came up. Keep watching the
1061 // interface so its answers do not go stale.
1062 self.start_lan_ip_monitor(lan_ip, port, None).await;
1063 return;
1064 };
1065
1066 // Publish all registered slugs.
1067 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
1068 for slug in slugs.keys() {
1069 let hostname = format!("{slug}.local");
1070 publisher.publish(&hostname, port);
1071 }
1072
1073 log::info!(
1074 "LAN mode: mDNS publishing on {lan_ip}, {} slug(s) registered",
1075 slugs.len()
1076 );
1077
1078 let publisher = std::sync::Arc::new(tokio::sync::Mutex::new(publisher));
1079
1080 // Start the IP monitor. It declines on its own when the address is
1081 // pinned rather than auto-detected.
1082 self.start_lan_ip_monitor(lan_ip, port, Some(publisher.clone()))
1083 .await;
1084
1085 *self.mdns_publisher.lock().await = Some(publisher);
1086 }
1087
1088 /// Re-read slugs from config and update mDNS records.
1089 ///
1090 /// Publishes new slugs and unpublishes removed ones. Called via IPC when
1091 /// `proxy add` or `proxy remove` modifies the slug registry.
1092 async fn sync_mdns(&self) {
1093 // Clone the Arc and release the outer lock immediately so we don't
1094 // block close() from taking the publisher during shutdown.
1095 let publisher = {
1096 let guard = self.mdns_publisher.lock().await;
1097 match guard.as_ref() {
1098 Some(p) => p.clone(),
1099 None => {
1100 debug!("sync_mdns: mDNS publisher not active, skipping");
1101 return;
1102 }
1103 }
1104 };
1105
1106 let s = crate::settings::settings();
1107 let port = u16::try_from(s.proxy.port).unwrap_or(443);
1108
1109 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
1110 let mut pub_guard = publisher.lock().await;
1111
1112 // Unpublish slugs that no longer exist in config.
1113 let current_keys: Vec<&String> = slugs.keys().collect();
1114 let registered: Vec<String> = pub_guard.registered_hostnames();
1115 for hostname in ®istered {
1116 // hostname is "slug.local" — extract slug part.
1117 let slug = hostname.strip_suffix(".local").unwrap_or(hostname);
1118 if !current_keys.iter().any(|k| k.as_str() == slug) {
1119 log::info!("mDNS: unpublishing removed slug {slug}");
1120 pub_guard.unpublish(hostname);
1121 }
1122 }
1123
1124 // Publish new slugs that aren't yet registered.
1125 for slug in slugs.keys() {
1126 let hostname = format!("{slug}.local");
1127 if !pub_guard.is_published(&hostname) {
1128 log::info!("mDNS: publishing new slug {slug}");
1129 pub_guard.publish(&hostname, port);
1130 }
1131 }
1132 }
1133
1134 /// Spawn a background task that periodically flushes the state file to
1135 /// disk if it has been marked dirty. Uses debouncing (1s interval) to
1136 /// batch rapid state changes.
1137 fn start_state_flush_task(&self) {
1138 let cancel = tokio_util::sync::CancellationToken::new();
1139 *self.flush_cancel.lock().unwrap() = Some(cancel.clone());
1140 tokio::spawn(async move {
1141 let mut interval = time::interval(Duration::from_secs(1));
1142 interval.set_missed_tick_behavior(time::MissedTickBehavior::Skip);
1143 loop {
1144 tokio::select! {
1145 _ = interval.tick() => {}
1146 _ = cancel.cancelled() => {
1147 debug!("state flush task received shutdown signal");
1148 break;
1149 }
1150 }
1151 let state = SUPERVISOR.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 debug!("state flush task exiting");
1159 });
1160 }
1161
1162 pub(crate) async fn flush_state(&self) {
1163 let state = self.state_file.lock().await;
1164 if state.is_dirty()
1165 && let Err(e) = state.write()
1166 {
1167 warn!("failed to flush state file: {e}");
1168 }
1169 }
1170
1171 pub(crate) async fn refresh(&self) -> Result<()> {
1172 trace!("refreshing");
1173
1174 // Collect PIDs we need to check (shell PIDs and liveness PIDs)
1175 // This is more efficient than refreshing all processes on the system
1176 let dirs_with_pids = self.get_dirs_with_shell_pids().await;
1177 let liveness_sessions = self.get_liveness_sessions().await;
1178 let pids_to_check: Vec<u32> = dirs_with_pids
1179 .values()
1180 .flatten()
1181 .copied()
1182 .chain(liveness_sessions.iter().map(|(pid, _, _)| *pid))
1183 .collect::<std::collections::HashSet<_>>()
1184 .into_iter()
1185 .collect();
1186
1187 if pids_to_check.is_empty() {
1188 // No PIDs to check, skip the expensive refresh
1189 trace!("no tracked PIDs to check, skipping process refresh");
1190 } else {
1191 debug!("refreshing PIDs: {pids_to_check:?}");
1192 PROCS.refresh_pids(&pids_to_check);
1193 }
1194
1195 let mut last_refreshed_at = self.last_refreshed_at.lock().await;
1196 *last_refreshed_at = time::Instant::now();
1197
1198 self.restore_own_record().await;
1199
1200 #[cfg_attr(not(unix), allow(unused_mut))]
1201 let mut dirs_to_leave: Vec<PathBuf> = Vec::new();
1202
1203 // Prune shell PIDs that are no longer running. This is essential on
1204 // Unix so that exited shells don't keep daemons alive forever.
1205 //
1206 // On Windows, skip this check: Git Bash (MSYS2) PIDs from `$$` are
1207 // Cygwin-internal PIDs that are invisible to sysinfo (which sees
1208 // Windows PIDs). The is_running check would always return false,
1209 // immediately removing every registered shell and breaking autostop.
1210 // Shell registration/deregistration relies on UpdateShellDir IPC
1211 // messages instead.
1212 #[cfg(unix)]
1213 for (dir, pids) in dirs_with_pids {
1214 let to_remove = pids
1215 .iter()
1216 .filter(|pid| !PROCS.is_running(**pid))
1217 .collect::<Vec<_>>();
1218 for pid in &to_remove {
1219 self.remove_shell_pid(**pid).await?
1220 }
1221 if to_remove.len() == pids.len() {
1222 dirs_to_leave.push(dir);
1223 }
1224 }
1225
1226 // Atomically remove project sessions whose host PID has died or whose
1227 // recorded title no longer matches the current process title. Every
1228 // project session carries a host PID in its key, so we iterate all of
1229 // them. Re-reading the sessions under the lock prevents enter/leave
1230 // interleaving from deleting a session that was just replaced with a
1231 // new title snapshot.
1232 //
1233 // Gated to Unix to mirror the shell-PID pruning above: on Windows,
1234 // Git Bash (MSYS2) `$$` PIDs are Cygwin-internal and invisible to
1235 // sysinfo, so the liveness check would immediately revoke every
1236 // freshly-entered session. Windows relies on explicit `project leave`
1237 // (or shell UpdateShellDir) for deregistration instead.
1238 #[cfg(unix)]
1239 {
1240 let mut state = self.state_file.lock().await;
1241 for (pid, dir, recorded_title) in liveness_sessions {
1242 let Some(session) = state.get_project_session(pid, &dir) else {
1243 continue;
1244 };
1245 let current_title = PROCS.title(pid);
1246 let is_running = PROCS.is_running(pid);
1247 debug!(
1248 "refresh liveness session pid {pid} dir {} recorded_title={recorded_title:?} current_title={current_title:?} is_running={is_running}",
1249 dir.display()
1250 );
1251 if should_remove_liveness_session(
1252 session,
1253 &recorded_title,
1254 current_title.as_deref(),
1255 is_running,
1256 ) {
1257 warn!(
1258 "removing project session pid {pid} dir {} (liveness pid title mismatch or dead)",
1259 dir.display()
1260 );
1261 if state.remove_project_session(pid, &dir).is_some() {
1262 dirs_to_leave.push(dir);
1263 }
1264 }
1265 }
1266 }
1267
1268 for dir in dirs_to_leave {
1269 self.leave_dir(&dir).await?;
1270 }
1271
1272 // Catch state-`running` daemons that lost their monitor (e.g. the
1273 // monitor died with a previous supervisor): mark dead ones errored
1274 // and re-adopt live ones. Runs before check_retry so a daemon marked
1275 // errored here is retried on this same tick.
1276 self.reconcile_unmonitored_daemons().await;
1277
1278 self.check_retry().await?;
1279 self.process_pending_autostops().await?;
1280
1281 Ok(())
1282 }
1283
1284 /// Install a SIGCHLD handler that reaps orphaned zombie child processes.
1285 ///
1286 /// When running as PID 1 inside a container, orphaned processes are
1287 /// re-parented to PID 1. Without explicit reaping, they accumulate
1288 /// as zombies in the process table indefinitely.
1289 ///
1290 /// Only reaps processes that are NOT managed by the supervisor (i.e.
1291 /// not tracked in the state file). Managed daemon processes are reaped
1292 /// by their monitoring tasks via `child.wait()`.
1293 ///
1294 /// ## Strategy
1295 ///
1296 /// **Linux**: Uses `waitid(Id::All, WNOHANG | WNOWAIT | WEXITED)` to
1297 /// *peek* at the next zombie without consuming its status. If the PID
1298 /// belongs to a managed daemon, the reaper skips it so Tokio's
1299 /// `child.wait()` can collect the status normally. Only unmanaged
1300 /// orphans are actually reaped (via `waitpid(Pid, WNOHANG)`). This
1301 /// eliminates the race entirely.
1302 ///
1303 /// **Non-Linux Unix** (e.g. macOS — mainly for local development;
1304 /// container mode targets Linux): `waitid` is unavailable, so we fall
1305 /// back to `waitpid(None, WNOHANG)`. If the reaper accidentally
1306 /// consumes a managed PID's status, it stashes the exit code in
1307 /// [`REAPED_STATUSES`] for the monitoring task to recover.
1308 #[cfg(unix)]
1309 fn reap_zombies(&self) -> Result<()> {
1310 let mut stream = signal::unix::signal(SignalKind::child())
1311 .map_err(|e| miette::miette!("Failed to register SIGCHLD handler: {e}"))?;
1312 tokio::spawn(async move {
1313 loop {
1314 stream.recv().await;
1315 // Collect PIDs of managed daemons so we don't steal their exit status
1316 let managed_pids: HashSet<u32> = SUPERVISOR
1317 .state_file
1318 .lock()
1319 .await
1320 .daemons
1321 .values()
1322 .filter_map(|d| d.pid)
1323 .collect();
1324 // Reap all available zombie children that are NOT managed
1325 Self::reap_unmanaged_zombies(&managed_pids).await;
1326 }
1327 });
1328 info!("container mode: SIGCHLD zombie reaper installed");
1329 Ok(())
1330 }
1331
1332 /// Linux implementation: peek with `waitid(WNOWAIT)` then selectively reap.
1333 ///
1334 /// `WNOWAIT` leaves the zombie in the table so we can inspect its PID
1335 /// without consuming the exit status. Only if the PID is *not* managed
1336 /// do we call `waitpid(Pid, WNOHANG)` to actually reap it.
1337 #[cfg(target_os = "linux")]
1338 async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
1339 use nix::sys::wait::{Id, WaitPidFlag, WaitStatus, waitid, waitpid};
1340 use nix::unistd::Pid;
1341
1342 loop {
1343 // Peek at the next zombie without consuming it
1344 let peek_flags = WaitPidFlag::WNOHANG | WaitPidFlag::WNOWAIT | WaitPidFlag::WEXITED;
1345 match waitid(Id::All, peek_flags) {
1346 Ok(WaitStatus::StillAlive) => break,
1347 Ok(status) => {
1348 let Some(pid_raw) = status.pid().map(|p| p.as_raw() as u32) else {
1349 break;
1350 };
1351 if managed_pids.contains(&pid_raw) {
1352 // This is a managed daemon — leave it for Tokio's child.wait().
1353 // We must break out of the loop because waitid(Id::All) would
1354 // keep returning the same zombie if we don't consume it.
1355 trace!(
1356 "zombie reaper: skipping managed daemon pid {pid_raw}, \
1357 leaving for Tokio to reap"
1358 );
1359 break;
1360 }
1361 // Not managed — actually reap it
1362 match waitpid(Pid::from_raw(pid_raw as i32), Some(WaitPidFlag::WNOHANG)) {
1363 Ok(s) => trace!("reaped orphaned zombie child: {s:?}"),
1364 Err(nix::errno::Errno::ECHILD) => break,
1365 Err(e) => {
1366 trace!("waitpid error reaping pid {pid_raw}: {e}");
1367 break;
1368 }
1369 }
1370 }
1371 Err(nix::errno::Errno::ECHILD) => break, // no children at all
1372 Err(e) => {
1373 trace!("waitid error in zombie reaper: {e}");
1374 break;
1375 }
1376 }
1377 }
1378 }
1379
1380 /// Non-Linux fallback: blind `waitpid(None, WNOHANG)` with stash recovery.
1381 ///
1382 /// Since `waitid(WNOWAIT)` is not available, we cannot peek. If we
1383 /// accidentally reap a managed PID, we stash the exit code in
1384 /// [`REAPED_STATUSES`] so the monitoring task can recover it.
1385 #[cfg(all(unix, not(target_os = "linux")))]
1386 async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
1387 use nix::sys::wait::{WaitPidFlag, WaitStatus, waitpid};
1388
1389 loop {
1390 match waitpid(None, Some(WaitPidFlag::WNOHANG)) {
1391 Ok(WaitStatus::StillAlive) => break,
1392 Ok(status) => {
1393 let Some(pid) = status.pid().map(|p| p.as_raw() as u32) else {
1394 continue;
1395 };
1396 if managed_pids.contains(&pid) {
1397 // Race lost — stash the exit code for lifecycle recovery
1398 let exit_code = match status {
1399 WaitStatus::Exited(_, code) => code,
1400 WaitStatus::Signaled(_, sig, _) => -(sig as i32),
1401 _ => -1,
1402 };
1403 warn!(
1404 "zombie reaper reaped managed daemon pid {pid} \
1405 (exit_code={exit_code}); stashing status for recovery"
1406 );
1407 REAPED_STATUSES.lock().await.insert(pid, exit_code);
1408 } else {
1409 trace!("reaped orphaned zombie child: {status:?}");
1410 }
1411 }
1412 Err(nix::errno::Errno::ECHILD) => break, // no more children
1413 Err(e) => {
1414 trace!("waitpid error in zombie reaper: {e}");
1415 break;
1416 }
1417 }
1418 }
1419 }
1420
1421 #[cfg(unix)]
1422 fn signals(&self) -> Result<()> {
1423 let signals = [
1424 SignalKind::terminate(),
1425 SignalKind::alarm(),
1426 SignalKind::interrupt(),
1427 SignalKind::quit(),
1428 SignalKind::hangup(),
1429 SignalKind::user_defined1(),
1430 SignalKind::user_defined2(),
1431 ];
1432 static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
1433 for signal in signals {
1434 let stream = match signal::unix::signal(signal) {
1435 Ok(s) => s,
1436 Err(e) => {
1437 warn!("Failed to register signal handler for {signal:?}: {e}");
1438 continue;
1439 }
1440 };
1441 tokio::spawn(async move {
1442 let mut stream = stream;
1443 loop {
1444 stream.recv().await;
1445 if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
1446 exit(1);
1447 } else {
1448 SUPERVISOR.handle_signal().await;
1449 }
1450 }
1451 });
1452 }
1453 Ok(())
1454 }
1455
1456 #[cfg(windows)]
1457 fn signals(&self) -> Result<()> {
1458 tokio::spawn(async move {
1459 static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
1460 loop {
1461 if let Err(e) = signal::ctrl_c().await {
1462 error!("Failed to wait for ctrl-c: {}", e);
1463 return;
1464 }
1465 if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
1466 exit(1);
1467 } else {
1468 SUPERVISOR.handle_signal().await;
1469 }
1470 }
1471 });
1472 Ok(())
1473 }
1474
1475 async fn handle_signal(&self) {
1476 info!("received signal, stopping");
1477 self.close().await;
1478 exit(0)
1479 }
1480
1481 pub(crate) async fn close(&self) {
1482 self.shutting_down.store(true, atomic::Ordering::Release);
1483 // Signal the proxy server to stop accepting new connections
1484 // and drain in-flight ones, *before* stopping daemons so the
1485 // proxy has time to finish forwarding active requests.
1486 if let Some(cancel) = self.proxy_cancel.lock().await.take() {
1487 cancel.cancel();
1488 }
1489
1490 // Stop the LAN IP monitor task. It watches the token that was just
1491 // cancelled, so give it a moment to come back on its own rather than
1492 // cutting it off mid-iteration the instant after asking it to stop.
1493 if let Some(mut monitor_task) = self.lan_monitor_task.lock().await.take()
1494 && tokio::time::timeout(Duration::from_secs(1), &mut monitor_task)
1495 .await
1496 .is_err()
1497 {
1498 monitor_task.abort();
1499 }
1500
1501 // Shutdown the mDNS publisher (sends goodbye packets).
1502 if let Some(publisher) = self.mdns_publisher.lock().await.take() {
1503 publisher.lock().await.shutdown();
1504 }
1505
1506 if let Some(mut dns_task) = self.dns_task.lock().await.take()
1507 && tokio::time::timeout(Duration::from_secs(5), &mut dns_task)
1508 .await
1509 .is_err()
1510 {
1511 // Cancelled but still running: aborted, as the LAN monitor is,
1512 // so its sockets are not left bound to the resolver port.
1513 dns_task.abort();
1514 }
1515
1516 if let Some(proxy_task) = self.proxy_task.lock().await.take() {
1517 // Longer than the proxy's own drain budget, so the task finishes on
1518 // its own terms rather than being cut off mid-drain.
1519 let _ = tokio::time::timeout(
1520 crate::proxy::server::SHUTDOWN_DRAIN_BUDGET + Duration::from_secs(2),
1521 proxy_task,
1522 )
1523 .await;
1524 }
1525
1526 // Clean up /etc/hosts entries managed by pitchfork
1527 let s = settings();
1528 if s.proxy.enable && s.proxy.sync_hosts {
1529 crate::proxy::hosts::clean_hosts_file();
1530 }
1531
1532 let pitchfork_id = DaemonId::pitchfork();
1533 let active = self.active_daemons().await;
1534 let active_ids: Vec<DaemonId> = active
1535 .iter()
1536 .filter(|d| d.id != pitchfork_id)
1537 .map(|d| d.id.clone())
1538 .collect();
1539
1540 // Stop daemons in reverse dependency order.
1541 // If dependency resolution fails (e.g. config changed), fall back to
1542 // stopping in arbitrary order so we still shut down cleanly.
1543 // Daemons within the same level are stopped concurrently.
1544 //
1545 // Each stop waits for the daemon's whole process group (bounded by its
1546 // stop budget) and levels are sequential, so total shutdown time is the
1547 // sum of the slowest stop per level. If an external manager (docker,
1548 // systemd) kills us before this completes, cleanup_orphaned_daemons()
1549 // recovers the leftover processes and stale state on the next start.
1550 let stop_levels = compute_reverse_stop_order(&active_ids);
1551 for level in &stop_levels {
1552 let mut tasks = Vec::new();
1553 for id in level {
1554 let id = id.clone();
1555 tasks.push(tokio::spawn(async move {
1556 if let Err(err) = SUPERVISOR.stop(&id).await {
1557 error!("failed to stop daemon {id}: {err}");
1558 }
1559 }));
1560 }
1561 for task in tasks {
1562 let _ = task.await;
1563 }
1564 }
1565 let _ = self.remove_daemon(&pitchfork_id).await;
1566
1567 // Signal the background state flush task to exit so it doesn't
1568 // keep waking up and acquiring the state mutex after shutdown.
1569 if let Some(cancel) = self.flush_cancel.lock().unwrap().take() {
1570 cancel.cancel();
1571 }
1572
1573 // Force-flush state to disk before shutting down IPC so no
1574 // in-memory-only changes are lost.
1575 {
1576 let state = self.state_file.lock().await;
1577 if state.is_dirty()
1578 && let Err(e) = state.write()
1579 {
1580 warn!("failed to flush state file during shutdown: {e}");
1581 }
1582 }
1583
1584 // Signal IPC server to shut down gracefully
1585 if let Some(mut handle) = self.ipc_shutdown.lock().await.take() {
1586 handle.shutdown().await;
1587 }
1588
1589 // Wait for all in-flight monitoring tasks to finish registering their
1590 // hook handles. Each monitoring task increments `active_monitors` when
1591 // its process exits, and decrements it (+ notifies `monitor_done`)
1592 // after all fire_hook() calls complete. This replaces the old
1593 // yield_now() approach which had a race window.
1594 let drain_timeout = time::sleep(Duration::from_secs(5));
1595 tokio::pin!(drain_timeout);
1596 loop {
1597 if self.active_monitors.load(atomic::Ordering::Acquire) == 0 {
1598 break;
1599 }
1600 tokio::select! {
1601 _ = self.monitor_done.notified() => {}
1602 _ = &mut drain_timeout => {
1603 warn!("timed out waiting for monitoring tasks to register hooks, proceeding with shutdown");
1604 break;
1605 }
1606 }
1607 }
1608 let handles: Vec<JoinHandle<()>> = std::mem::take(&mut *self.hook_tasks.lock().await);
1609 let hook_timeout = Duration::from_secs(30);
1610 for handle in handles {
1611 match time::timeout(hook_timeout, handle).await {
1612 Ok(_) => {} // Hook completed (success or error, doesn't matter)
1613 Err(_) => {
1614 warn!(
1615 "hook task did not complete within {hook_timeout:?} during shutdown, skipping"
1616 );
1617 }
1618 }
1619 }
1620
1621 // Unix: remove the socket directory if it is empty. The IPC server
1622 // already removed our socket; anything left belongs to a supervisor
1623 // that replaced this one (e.g. `supervisor run --force`) while we
1624 // were stopping daemons, and must not be deleted.
1625 // Windows: named pipes have no filesystem component.
1626 #[cfg(unix)]
1627 let _ = fs::remove_dir(&*env::IPC_SOCK_DIR);
1628 }
1629
1630 pub(crate) async fn add_notification(&self, level: log::LevelFilter, message: String) {
1631 self.pending_notifications
1632 .lock()
1633 .await
1634 .push((level, message));
1635 }
1636}
1637
1638/// Fix ownership on the state directory so non-root users can access files
1639/// created by a `sudo`-started supervisor.
1640///
1641/// When `[settings.supervisor] user`, a recorded invoking user (see
1642/// [`env::InvokingUser`]), or `SUDO_UID`/`SUDO_GID` are set, we
1643/// `chown` the state directory and safe subdirectories back to that non-root
1644/// runtime user. This is strictly better than `chmod 0o666` because it does not
1645/// widen the permission bits — the files stay owner-only (0o600/0o700) but the
1646/// *owner* is the user that daemon processes and CLI clients need to share.
1647///
1648/// **Security**: The `proxy/` subtree is intentionally skipped. It contains
1649/// `ca-key.pem` which must remain `0o600` and owned by the process that
1650/// generated it. Changing its ownership or permissions would expose the CA
1651/// private key to other local users.
1652///
1653/// If none of these are available (e.g. a direct root login, or a boot service
1654/// installed from a root shell), we fall back to relaxing permissions on only
1655/// the `sock/` and `logs/` subdirectories (plus `state.toml`) so CLI clients
1656/// can still function.
1657#[cfg(unix)]
1658fn fix_state_dir_permissions() {
1659 let state_dir = &*env::PITCHFORK_STATE_DIR;
1660 if let Some((uid, gid)) = state_owner_ids() {
1661 if !state_dir.exists()
1662 && let Err(err) = fs::create_dir_all(state_dir)
1663 {
1664 warn!(
1665 "failed to create state directory for ownership fix at {}: {err}",
1666 state_dir.display()
1667 );
1668 return;
1669 }
1670
1671 // Best path: chown back to the runtime user. Permissions stay tight.
1672 chown_recursive(state_dir, uid, gid, true);
1673 debug!(
1674 "chowned state directory to uid={uid} gid={gid} at {}",
1675 state_dir.display()
1676 );
1677 } else {
1678 if !state_dir.exists() {
1679 return;
1680 }
1681
1682 // Fallback: relax permissions on safe subdirectories only.
1683 // proxy/ is never touched.
1684 chmod_safe_subtrees(state_dir);
1685 debug!(
1686 "relaxed permissions on safe subtrees at {}",
1687 state_dir.display()
1688 );
1689 }
1690}
1691
1692#[cfg(unix)]
1693pub(crate) fn state_owner_ids() -> Option<(u32, u32)> {
1694 if !nix::unistd::Uid::effective().is_root() {
1695 return None;
1696 }
1697
1698 let s = settings();
1699 let user = s.supervisor.user.trim();
1700 if !user.is_empty() {
1701 return resolve_supervisor_user_ids(user).or_else(|| {
1702 warn!(
1703 "failed to resolve supervisor.user '{user}' for state ownership; falling back to the invoking user"
1704 );
1705 env::invoking_user_ids()
1706 });
1707 }
1708
1709 env::invoking_user_ids()
1710}
1711
1712#[cfg(unix)]
1713fn resolve_supervisor_user_ids(user: &str) -> Option<(u32, u32)> {
1714 let user_record = if user.chars().all(|c| c.is_ascii_digit()) {
1715 let uid = user.parse::<u32>().ok()?;
1716 nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1717 .ok()
1718 .flatten()
1719 } else {
1720 nix::unistd::User::from_name(user).ok().flatten()
1721 }?;
1722
1723 Some((user_record.uid.as_raw(), user_record.gid.as_raw()))
1724}
1725
1726/// Recursively `chown` a directory tree. If `skip_proxy` is true, the `proxy/`
1727/// subdirectory is skipped entirely to protect the CA private key.
1728#[cfg(unix)]
1729fn chown_recursive(dir: &std::path::Path, uid: u32, gid: u32, skip_proxy: bool) {
1730 // chown the directory itself
1731 let _ = chown_path(dir, uid, gid);
1732
1733 let entries = match std::fs::read_dir(dir) {
1734 Ok(e) => e,
1735 Err(_) => return,
1736 };
1737 for entry in entries.flatten() {
1738 let path = entry.path();
1739 if path.is_dir() {
1740 // Skip proxy/ at the top level of the state directory
1741 if skip_proxy
1742 && let Some(name) = path.file_name().and_then(|n| n.to_str())
1743 && name == "proxy"
1744 {
1745 continue;
1746 }
1747 chown_recursive(&path, uid, gid, false);
1748 } else {
1749 let _ = chown_path(&path, uid, gid);
1750 }
1751 }
1752}
1753
1754/// `chown` a single path using libc. Returns Ok(()) on success.
1755#[cfg(unix)]
1756fn chown_path(path: &std::path::Path, uid: u32, gid: u32) -> std::io::Result<()> {
1757 use std::ffi::CString;
1758 use std::os::unix::ffi::OsStrExt;
1759 let c_path = CString::new(path.as_os_str().as_bytes())
1760 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
1761 let ret = unsafe { libc::chown(c_path.as_ptr(), uid, gid) };
1762 if ret == 0 {
1763 Ok(())
1764 } else {
1765 Err(std::io::Error::last_os_error())
1766 }
1767}
1768
1769/// Fallback: relax permissions on safe subdirectories only (sock/, logs/, and
1770/// state.toml). The proxy/ subtree is never touched.
1771#[cfg(unix)]
1772fn chmod_safe_subtrees(state_dir: &std::path::Path) {
1773 // The state directory itself needs to be traversable
1774 let _ = fs::set_permissions(state_dir, fs::Permissions::from_mode(0o755));
1775
1776 // state.toml — needs to be readable by CLI clients
1777 let state_file = state_dir.join("state.toml");
1778 if state_file.exists() {
1779 let _ = fs::set_permissions(&state_file, fs::Permissions::from_mode(0o644));
1780 }
1781
1782 // Safe subdirectories: sock/ and logs/
1783 for subdir_name in &["sock", "logs"] {
1784 let subdir = state_dir.join(subdir_name);
1785 if subdir.is_dir() {
1786 chmod_recursive(&subdir);
1787 }
1788 }
1789}
1790
1791/// On startup, reconcile daemon processes left behind by a previous supervisor
1792/// that was terminated unexpectedly (e.g. `kill -9`).
1793///
1794/// This iterates the state file for daemon entries with a recorded PID. If the
1795/// PID is still alive and its current identity matches the recorded start time
1796/// (or the recorded title for older state files), it is assumed to be an orphan
1797/// from the previous supervisor session and `supervisor.orphan_policy` decides
1798/// its fate: `adopt` (default) resumes supervision via a poll monitor and keeps
1799/// the daemon's state intact; `kill` terminates it and resets its state to
1800/// `Stopped` with no PID. If a matching live process cannot be terminated
1801/// securely, its running state is retained to prevent a duplicate instance
1802/// from being started.
1803///
1804/// Missing or mismatched identity data fails closed so a PID recycled by an
1805/// unrelated process is never adopted or killed. On Unix platforms without
1806/// durable process handles, orphan termination also fails closed because the
1807/// PID/PGID cannot be pinned between identity validation and signaling.
1808///
1809/// This is gated by the `supervisor.cleanup_orphans` setting (default: true).
1810///
1811/// `unclean` is whether the previous supervisor exited uncleanly (see
1812/// [`supervisor_exited_uncleanly`]). It is read in `start()` before the
1813/// starting supervisor records itself in the state file — this function runs
1814/// in the background after that record exists, so reading it here would
1815/// always report unclean.
1816async fn cleanup_orphaned_daemons(supervisor: &Supervisor, unclean: bool) {
1817 if !settings().supervisor.cleanup_orphans {
1818 return;
1819 }
1820
1821 let candidates: Vec<_> = {
1822 let state = supervisor.state_file.lock().await;
1823 state
1824 .daemons
1825 .values()
1826 .filter(|d| d.id != DaemonId::pitchfork() && d.pid.is_some())
1827 .cloned()
1828 .collect()
1829 };
1830
1831 if candidates.is_empty() {
1832 return;
1833 }
1834
1835 info!(
1836 "checking {} daemon(s) for orphaned processes",
1837 candidates.len()
1838 );
1839
1840 let policy = orphan_policy();
1841 let boot_time = PROCS.boot_time();
1842
1843 // Reconcile orphans in parallel — a kill waits for the daemon's whole
1844 // process group to exit, bounded by that daemon's stop budget, so
1845 // sequential processing would make total cleanup time the sum of the
1846 // budgets.
1847 let tasks: Vec<_> = candidates
1848 .into_iter()
1849 .map(|daemon| {
1850 let policy = policy.clone();
1851 tokio::spawn(cleanup_orphaned_daemon(daemon, policy, boot_time, unclean))
1852 })
1853 .collect();
1854 for task in tasks {
1855 let _ = task.await;
1856 }
1857}
1858
1859/// Reconcile a single orphan candidate: adopt it, kill it, or reset its state,
1860/// per the policy and identity checks described on [`cleanup_orphaned_daemons`].
1861///
1862/// Holds the daemon's stop lock so a concurrent Run/Stop request for the same
1863/// daemon (cleanup runs in the background) serializes with the orphan kill,
1864/// and re-checks the recorded PID under the lock: if it changed, another path
1865/// already replaced or cleaned up this record and the snapshot is stale.
1866async fn cleanup_orphaned_daemon(
1867 daemon: crate::daemon::Daemon,
1868 policy: String,
1869 boot_time: u64,
1870 unclean: bool,
1871) {
1872 let supervisor: &Supervisor = &SUPERVISOR;
1873 let Some(pid) = daemon.pid else { return };
1874
1875 let lock = supervisor.stop_lock(&daemon.id).await;
1876 let _guard = lock.lock().await;
1877 let current_pid = {
1878 let state = supervisor.state_file.lock().await;
1879 state.daemons.get(&daemon.id).and_then(|d| d.pid)
1880 };
1881 if current_pid != Some(pid) {
1882 debug!(
1883 "orphan cleanup: daemon {} pid changed (recorded {pid}, now {current_pid:?}), skipping",
1884 daemon.id
1885 );
1886 return;
1887 }
1888
1889 // Refresh the candidate immediately before checking it: waiting for the
1890 // stop lock can await another path's stop timeout, during which this PID
1891 // may exit and be recycled.
1892 PROCS.refresh_pids(&[pid]);
1893
1894 if !PROCS.is_running(pid) {
1895 // PID already dead — the daemon exited while unsupervised, so
1896 // record a terminal status that reflects whether it died under a
1897 // crashed supervisor (retryable) or with the machine.
1898 let status = unobserved_exit_status(
1899 &daemon.status,
1900 daemon.boot_time,
1901 boot_time,
1902 unclean,
1903 daemon.oneshot,
1904 );
1905 reset_daemon_state(supervisor, &daemon.id, status, ExitObservation::Unobserved).await;
1906 return;
1907 }
1908
1909 // Safety check: verify the live process really is the daemon we
1910 // recorded, not an unrelated process that received a recycled PID.
1911 // The kernel start time is a stable identity for the lifetime of a
1912 // process, and is the only thing accepted as one.
1913 let current_start_time = PROCS.start_time(pid);
1914 let matches = process_identity_matches(daemon.start_time, current_start_time);
1915
1916 if !matches {
1917 // Either side missing means the identity cannot be checked at all,
1918 // which is different from checking it and finding a stranger: retain
1919 // the running state rather than resetting a record whose process may
1920 // well still be the daemon.
1921 if daemon.start_time.is_none() || current_start_time.is_none() {
1922 warn!(
1923 "could not verify the identity of live pid {pid} recorded for daemon {}; retaining running state",
1924 daemon.id,
1925 );
1926 return;
1927 }
1928 warn!(
1929 "pid {pid} recorded for daemon {} belongs to a different process now (PID recycled); resetting state without killing",
1930 daemon.id,
1931 );
1932 // The daemon died at some unknown point and the OS handed its PID
1933 // to something else — same unobserved exit as a dead PID.
1934 let status = unobserved_exit_status(
1935 &daemon.status,
1936 daemon.boot_time,
1937 boot_time,
1938 unclean,
1939 daemon.oneshot,
1940 );
1941 reset_daemon_state(supervisor, &daemon.id, status, ExitObservation::Unobserved).await;
1942 return;
1943 }
1944
1945 // Both policies need a verified start time: killing revalidates it
1946 // while pinned to the process, and adoption anchors its poll monitor
1947 // to it so a later PID recycle is never mistaken for the daemon.
1948 let Some(expected_start_time) = current_start_time else {
1949 warn!(
1950 "could not read start time for live pid {pid} recorded for daemon {}; retaining running state",
1951 daemon.id,
1952 );
1953 return;
1954 };
1955
1956 // Identity verified — the process really is our orphaned daemon.
1957 // The policy decides whether supervision resumes or the slate is
1958 // wiped clean.
1959 if policy == "adopt" {
1960 supervisor
1961 .adopt_daemon(&daemon, pid, expected_start_time)
1962 .await;
1963 return;
1964 }
1965
1966 info!("terminating orphaned daemon {} (pid {pid})", daemon.id);
1967
1968 let stop_cfg = daemon.stop_signal.unwrap_or_default();
1969 let termination_result = PROCS
1970 .kill_process_group_if_start_time_matches_async(
1971 pid,
1972 Some(expected_start_time),
1973 stop_cfg.signal.into(),
1974 stop_cfg.timeout,
1975 )
1976 .await;
1977
1978 match termination_result {
1979 Ok(true) => {}
1980 Ok(false) => {
1981 warn!(
1982 "could not securely terminate orphaned daemon {} (pid {pid}); retaining running state",
1983 daemon.id
1984 );
1985 return;
1986 }
1987 Err(err) => {
1988 warn!(
1989 "failed to terminate orphaned daemon {} (pid {pid}): {err}; retaining running state",
1990 daemon.id
1991 );
1992 return;
1993 }
1994 }
1995
1996 // We terminated the orphan ourselves, so this is an observed,
1997 // intentional stop rather than an unobserved exit.
1998 reset_daemon_state(
1999 supervisor,
2000 &daemon.id,
2001 DaemonStatus::Stopped,
2002 ExitObservation::Terminated,
2003 )
2004 .await;
2005}
2006
2007/// Effective `supervisor.orphan_policy`, warning on an unrecognized value
2008/// (which falls back to the default of adopting).
2009pub(crate) fn orphan_policy() -> String {
2010 let policy = settings().supervisor.orphan_policy.clone();
2011 match policy.as_str() {
2012 "adopt" | "kill" => policy,
2013 other => {
2014 warn!("unknown supervisor.orphan_policy '{other}', defaulting to 'adopt'");
2015 "adopt".to_string()
2016 }
2017 }
2018}
2019
2020/// Verify that live process identity matches the persisted daemon identity.
2021///
2022/// Both start times are required. A process name was once accepted in place of
2023/// a recorded start time, for state written before start times existed, but a
2024/// name is not an identity: a recycled PID belonging to another copy of the same
2025/// program matches it, and adopting or killing on that basis acts on the wrong
2026/// process. Missing identity, on either side, means unverifiable — and
2027/// unverifiable must never authorize acting on a process.
2028fn process_identity_matches(
2029 recorded_start_time: Option<u64>,
2030 current_start_time: Option<u64>,
2031) -> bool {
2032 match (recorded_start_time, current_start_time) {
2033 (Some(recorded), Some(current)) => recorded == current,
2034 _ => false,
2035 }
2036}
2037
2038/// Whether a PID read from persisted state may be signalled.
2039///
2040/// Stopping a daemon signals its whole process *group*, so acting on a PID that
2041/// has been recycled since it was recorded takes down an unrelated process tree.
2042/// Records are refused only when their identity is positively contradicted: if
2043/// either start time is unknown the PID stays as signallable as it was before
2044/// identities were recorded, so a daemon whose record predates the field can
2045/// still be stopped rather than becoming permanently unstoppable.
2046///
2047/// This is deliberately weaker than [`process_identity_matches`], which decides
2048/// whether to adopt or kill a process nobody asked about. Here the user has
2049/// named the daemon and asked for it to stop; the check exists to catch the
2050/// case where the answer is provably the wrong process.
2051pub(crate) fn signalling_pid_is_authorized(
2052 recorded_start_time: Option<u64>,
2053 current_start_time: Option<u64>,
2054) -> bool {
2055 !matches!(
2056 (recorded_start_time, current_start_time),
2057 (Some(recorded), Some(current)) if recorded != current
2058 )
2059}
2060
2061/// How a daemon's run ended, which decides what happens to the recorded
2062/// `last_exit_success` that cron `retrigger = "success" | "fail"` reads.
2063#[derive(Clone, Copy, PartialEq, Eq)]
2064pub(crate) enum ExitObservation {
2065 /// Nobody saw how the run ended, because the monitor that would have
2066 /// observed it died with a previous supervisor. The recorded outcome is
2067 /// cleared to `None`.
2068 ///
2069 /// Every option here is imperfect, so this picks the one that asserts
2070 /// nothing false. `Some(false)` would fabricate a failure, silently
2071 /// breaking a `retrigger = "success"` chain whose run may well have
2072 /// succeeded; `Some(true)` fabricates the opposite; keeping the previous
2073 /// value attributes an earlier run's outcome to this one. `None` says
2074 /// "unknown", reusing the reading the cron watcher already applies to a
2075 /// daemon that has never run.
2076 ///
2077 /// The tradeoff is that `None` satisfies both `retrigger = "success"`
2078 /// (`unwrap_or(true)`) and `retrigger = "fail"` (`!unwrap_or(false)`), so
2079 /// such a daemon fires once at its next scheduled time regardless of which
2080 /// it configured. That is schedule-gated rather than a loop, and it biases
2081 /// toward running the daemon over leaving it permanently untriggered.
2082 /// Distinguishing "unknown" from "never ran" would require a third cron
2083 /// state and is deliberately left out of scope here.
2084 Unobserved,
2085 /// We terminated the process ourselves, so the outcome is not a mystery:
2086 /// it stopped because we asked it to. Recorded as a success, matching the
2087 /// convention `Supervisor::stop` already uses for a deliberate stop.
2088 Terminated,
2089}
2090
2091impl ExitObservation {
2092 /// The `last_exit_success` value this observation implies.
2093 pub(crate) fn last_exit_success(self) -> Option<bool> {
2094 match self {
2095 ExitObservation::Unobserved => None,
2096 ExitObservation::Terminated => Some(true),
2097 }
2098 }
2099}
2100
2101/// Clear a daemon's runtime state (pid, process identity, active port) after
2102/// its process is gone or is no longer ours to manage.
2103///
2104/// Config fields are preserved by cloning the existing record, so a reset can
2105/// never drop a daemon's command, retry policy, or schedule.
2106async fn reset_daemon_state(
2107 supervisor: &Supervisor,
2108 id: &DaemonId,
2109 status: DaemonStatus,
2110 observation: ExitObservation,
2111) {
2112 let mut state_file = supervisor.state_file.lock().await;
2113 let Some(existing) = state_file.daemons.get(id) else {
2114 return;
2115 };
2116 let mut daemon = existing.clone();
2117 daemon.pid = None;
2118 daemon.title = None;
2119 daemon.start_time = None;
2120 daemon.boot_time = None;
2121 daemon.status = status;
2122 daemon.last_exit_success = observation.last_exit_success();
2123 daemon.active_port = None;
2124 state_file.clear_active_port(id);
2125 state_file.insert_daemon(id, daemon);
2126}
2127
2128/// Boot times this far apart are treated as different boots.
2129///
2130/// Sized to the only platform that reports a jittery value: Windows derives
2131/// boot time as `now - GetTickCount64()`, sampling two clocks independently,
2132/// so consecutive calls within one boot can differ by about a second. Linux
2133/// (`/proc/stat` btime) and macOS (`kern.boottime`) report stable values.
2134///
2135/// Deliberately kept this tight so a genuine reboot can never fall inside it:
2136/// a prior session would have to boot, start the supervisor, spawn a daemon,
2137/// have that daemon die, and complete a reboot inside two seconds, which no
2138/// real boot cycle reaches. A larger window would misread a short-lived
2139/// previous boot (e.g. a device in a reboot loop) as the current one and
2140/// resurrect daemons a reboot should have left stopped.
2141const BOOT_TIME_TOLERANCE_SECS: u64 = 2;
2142
2143/// Terminal status for a daemon whose process is gone and whose exit was
2144/// never observed, because the monitor that would have seen it died with a
2145/// previous supervisor.
2146///
2147/// A daemon recorded `Running` was expected to still be alive, so it died
2148/// under the crashed supervisor: `Errored(-1)` ("unknown exit code") makes it
2149/// eligible for its configured retries. Two cases stay `Stopped` instead:
2150///
2151/// - records from an earlier boot, whose processes died with the machine —
2152/// auto-restarting those is what `boot_start` is for, and reviving every
2153/// retry-configured daemon after a reboot would be a surprise
2154/// - any other status (in practice `Stopping`), i.e. an intentional stop that
2155/// completed while the supervisor was gone
2156pub(crate) fn unobserved_exit_status(
2157 status: &DaemonStatus,
2158 recorded_boot_time: Option<u64>,
2159 current_boot_time: u64,
2160 supervisor_exited_uncleanly: bool,
2161 oneshot: bool,
2162) -> DaemonStatus {
2163 let same_boot = recorded_boot_time
2164 .is_some_and(|recorded| recorded.abs_diff(current_boot_time) <= BOOT_TIME_TOLERANCE_SECS);
2165 // A task gets `stopped` rather than `errored` for the same reason the
2166 // adopted path does: `errored` is what `check_retry` looks for, and
2167 // nobody saw how this run ended, so retrying it would re-run a migration
2168 // or a seed that may well have succeeded. Leave re-running to an explicit
2169 // start.
2170 if status.is_running() && same_boot && supervisor_exited_uncleanly && !oneshot {
2171 DaemonStatus::Errored(-1)
2172 } else {
2173 DaemonStatus::Stopped
2174 }
2175}
2176
2177/// Whether the supervisor that owned this state file failed to shut down
2178/// cleanly, meaning any daemon it left behind stopped for reasons nobody
2179/// recorded.
2180///
2181/// A clean shutdown removes the supervisor's own entry: `close()` does it on
2182/// Unix, where the stop signal is delivered and handled, and the
2183/// `supervisor stop` command does it on Windows, which has no POSIX signals
2184/// and force-terminates the process instead. A crash, an external `kill -9`,
2185/// or a `--force` replacement all leave the entry behind.
2186///
2187/// This must be read before the starting supervisor records itself, which is
2188/// why `cleanup_orphaned_daemons` runs first in `start()`.
2189async fn supervisor_exited_uncleanly(supervisor: &Supervisor) -> bool {
2190 supervisor
2191 .state_file
2192 .lock()
2193 .await
2194 .daemons
2195 .contains_key(&DaemonId::pitchfork())
2196}
2197
2198/// Recursively chmod: directories → 0o755, files → 0o644.
2199#[cfg(unix)]
2200fn chmod_recursive(dir: &std::path::Path) {
2201 let _ = fs::set_permissions(dir, fs::Permissions::from_mode(0o755));
2202 let entries = match fs::read_dir(dir) {
2203 Ok(e) => e,
2204 Err(_) => return,
2205 };
2206 for entry in entries.flatten() {
2207 let path = entry.path();
2208 if path.is_dir() {
2209 chmod_recursive(&path);
2210 } else {
2211 let _ = fs::set_permissions(&path, fs::Permissions::from_mode(0o644));
2212 }
2213 }
2214}
2215
2216#[cfg(test)]
2217mod tests {
2218 use super::{
2219 BOOT_TIME_TOLERANCE_SECS, legacy_supervisor_title_matches, process_identity_matches,
2220 should_remove_liveness_session, signalling_pid_is_authorized, supervisor_identity_matches,
2221 unobserved_exit_status,
2222 };
2223 use crate::daemon_status::DaemonStatus;
2224 use crate::state_file::ProjectSession;
2225
2226 const BOOT: u64 = 1_700_000_000;
2227
2228 #[test]
2229 fn unobserved_running_death_in_current_boot_is_retryable() {
2230 // Died under a crashed supervisor during this boot: Errored(-1) makes
2231 // the daemon eligible for its configured retries.
2232 assert!(matches!(
2233 unobserved_exit_status(&DaemonStatus::Running, Some(BOOT), BOOT, true, false),
2234 DaemonStatus::Errored(-1)
2235 ));
2236 }
2237
2238 #[test]
2239 fn unobserved_task_death_is_stopped_not_retried() {
2240 // Nobody saw how the run ended, so `errored` would hand a migration or
2241 // a seed to check_retry on a guess. Matches what the adopted path
2242 // records, and what the guide promises.
2243 assert!(matches!(
2244 unobserved_exit_status(&DaemonStatus::Running, Some(BOOT), BOOT, true, true),
2245 DaemonStatus::Stopped
2246 ));
2247 }
2248
2249 #[test]
2250 fn unobserved_running_death_from_previous_boot_is_stopped() {
2251 // The process died with the machine; reviving every retry-configured
2252 // daemon after a reboot is what boot_start is for.
2253 assert!(matches!(
2254 unobserved_exit_status(
2255 &DaemonStatus::Running,
2256 Some(BOOT - 86_400),
2257 BOOT,
2258 true,
2259 false
2260 ),
2261 DaemonStatus::Stopped
2262 ));
2263 }
2264
2265 #[test]
2266 fn unobserved_exit_tolerates_boot_time_jitter() {
2267 // Windows recomputes boot time as now - GetTickCount64(), which can
2268 // drift about a second between samples within one boot.
2269 let within = BOOT + BOOT_TIME_TOLERANCE_SECS;
2270 assert!(matches!(
2271 unobserved_exit_status(&DaemonStatus::Running, Some(within), BOOT, true, false),
2272 DaemonStatus::Errored(-1)
2273 ));
2274 let beyond = BOOT + BOOT_TIME_TOLERANCE_SECS + 1;
2275 assert!(matches!(
2276 unobserved_exit_status(&DaemonStatus::Running, Some(beyond), BOOT, true, false),
2277 DaemonStatus::Stopped
2278 ));
2279 }
2280
2281 #[test]
2282 fn unobserved_exit_after_clean_shutdown_is_stopped() {
2283 // A deliberate `supervisor stop` can leave running records behind on
2284 // platforms where the supervisor cannot handle the stop signal. Those
2285 // daemons were stopped on purpose, so they must not be reported as
2286 // failures or resurrected by the retry checker.
2287 assert!(matches!(
2288 unobserved_exit_status(&DaemonStatus::Running, Some(BOOT), BOOT, false, false),
2289 DaemonStatus::Stopped
2290 ));
2291 }
2292
2293 #[test]
2294 fn unobserved_exit_treats_short_previous_boot_as_previous() {
2295 // A device in a reboot loop can produce consecutive boots seconds
2296 // apart. The jitter window must stay far below that so those records
2297 // are still recognised as belonging to an earlier boot.
2298 for gap in [5, 30, 59, 60] {
2299 assert!(
2300 matches!(
2301 unobserved_exit_status(
2302 &DaemonStatus::Running,
2303 Some(BOOT - gap),
2304 BOOT,
2305 true,
2306 false
2307 ),
2308 DaemonStatus::Stopped
2309 ),
2310 "boot {gap}s earlier should be treated as a previous boot"
2311 );
2312 }
2313 }
2314
2315 #[test]
2316 fn unobserved_exit_without_recorded_boot_time_is_stopped() {
2317 // Legacy state files predating the field fail closed to today's
2318 // behavior rather than triggering surprise retries.
2319 assert!(matches!(
2320 unobserved_exit_status(&DaemonStatus::Running, None, BOOT, true, false),
2321 DaemonStatus::Stopped
2322 ));
2323 }
2324
2325 #[test]
2326 fn supervisor_identity_matches_same_generation() {
2327 assert!(supervisor_identity_matches(
2328 Some(100),
2329 Some(100),
2330 Some(BOOT),
2331 BOOT
2332 ));
2333 }
2334
2335 #[test]
2336 fn supervisor_identity_survives_clock_steps_when_tokens_match() {
2337 // An NTP step or sleep/resume moves the realtime-derived boot time by
2338 // far more than the tolerance while the supervisor keeps running.
2339 // Matching start tokens prove it is the same process; declaring it
2340 // stale here would start a second supervisor.
2341 assert!(supervisor_identity_matches(
2342 Some(100),
2343 Some(100),
2344 Some(BOOT),
2345 BOOT + 3600
2346 ));
2347 assert!(supervisor_identity_matches(
2348 Some(100),
2349 Some(100),
2350 Some(BOOT + 3600),
2351 BOOT
2352 ));
2353 }
2354
2355 #[test]
2356 fn supervisor_identity_rejects_recycled_pid() {
2357 // The supervisor died and something else got its PID, within this
2358 // boot or across a reboot: the token differs either way.
2359 assert!(!supervisor_identity_matches(
2360 Some(100),
2361 Some(200),
2362 Some(BOOT),
2363 BOOT
2364 ));
2365 assert!(!supervisor_identity_matches(
2366 Some(100),
2367 Some(200),
2368 Some(BOOT),
2369 BOOT + 3600
2370 ));
2371 }
2372
2373 #[test]
2374 fn supervisor_identity_uses_boot_time_when_a_token_is_missing() {
2375 // Discussion #877: the record survived a reboot and an early system
2376 // daemon now owns the PID. With no live token to compare, the boot
2377 // time is what proves the record stale.
2378 assert!(!supervisor_identity_matches(
2379 Some(100),
2380 None,
2381 Some(BOOT),
2382 BOOT + 3600
2383 ));
2384 // Same boot (within Windows boot-time jitter) and no contradiction.
2385 assert!(supervisor_identity_matches(
2386 Some(100),
2387 None,
2388 Some(BOOT),
2389 BOOT + BOOT_TIME_TOLERANCE_SECS
2390 ));
2391 }
2392
2393 #[test]
2394 fn supervisor_identity_tolerates_legacy_records() {
2395 // A record written before either field existed is not contradicted by
2396 // anything here; `supervisor_record_is_live` applies the process-name
2397 // check to those instead.
2398 assert!(supervisor_identity_matches(None, Some(100), None, BOOT));
2399 // A record with only a boot time is still rejected across a reboot.
2400 assert!(!supervisor_identity_matches(
2401 None,
2402 Some(100),
2403 Some(BOOT),
2404 BOOT + 3600
2405 ));
2406 assert!(supervisor_identity_matches(
2407 None,
2408 Some(100),
2409 Some(BOOT),
2410 BOOT
2411 ));
2412 }
2413
2414 #[test]
2415 fn legacy_supervisor_title_requires_a_pitchfork_process() {
2416 assert!(legacy_supervisor_title_matches(Some("pitchfork")));
2417 assert!(legacy_supervisor_title_matches(Some("pitchfork.exe")));
2418 assert!(legacy_supervisor_title_matches(Some("Pitchfork")));
2419 // Discussion #877: an Apple LaunchAgent inherited the PID after a reboot.
2420 assert!(!legacy_supervisor_title_matches(Some(
2421 "AMPDeviceDiscoveryAgent"
2422 )));
2423 assert!(!legacy_supervisor_title_matches(Some("sleep")));
2424 assert!(!legacy_supervisor_title_matches(None));
2425 }
2426
2427 #[test]
2428 fn unobserved_exit_of_stopping_daemon_is_stopped() {
2429 // An intentional stop that completed while the supervisor was gone is
2430 // not a failure, even within the same boot.
2431 assert!(matches!(
2432 unobserved_exit_status(&DaemonStatus::Stopping, Some(BOOT), BOOT, true, false),
2433 DaemonStatus::Stopped
2434 ));
2435 }
2436
2437 #[test]
2438 fn orphan_identity_requires_both_start_times() {
2439 assert!(process_identity_matches(Some(123), Some(123)));
2440 assert!(!process_identity_matches(Some(123), Some(456)));
2441 // Unreadable current identity: unverifiable, so not a match.
2442 assert!(!process_identity_matches(Some(123), None));
2443 }
2444
2445 #[test]
2446 fn signalling_is_refused_only_for_a_contradicted_identity() {
2447 // Provably someone else's process group: refuse.
2448 assert!(!signalling_pid_is_authorized(Some(123), Some(456)));
2449 // Verified as the daemon's own.
2450 assert!(signalling_pid_is_authorized(Some(123), Some(123)));
2451 // Unknown on either side. Stopping stays possible, because the user has
2452 // named this daemon and a record that cannot be verified must not become
2453 // one that can never be stopped.
2454 assert!(signalling_pid_is_authorized(None, Some(123)));
2455 assert!(signalling_pid_is_authorized(Some(123), None));
2456 assert!(signalling_pid_is_authorized(None, None));
2457 }
2458
2459 #[test]
2460 fn orphan_identity_rejects_records_without_a_start_time() {
2461 // State written before start times were recorded. A process name used
2462 // to stand in here, but another copy of the same program on a recycled
2463 // PID matches a name, so such records are no longer verifiable and must
2464 // not authorize adopting or killing anything.
2465 assert!(!process_identity_matches(None, Some(123)));
2466 assert!(!process_identity_matches(None, None));
2467 }
2468
2469 #[test]
2470 fn should_not_remove_when_state_title_differs_from_snapshot() {
2471 // The session was re-entered after the snapshot was taken, producing a
2472 // new title in state. The snapshot title is stale; skip removal.
2473 let session = ProjectSession {
2474 liveness_title: Some("new_title".to_string()),
2475 };
2476 let recorded_title = Some("old_title".to_string());
2477
2478 assert!(!should_remove_liveness_session(
2479 &session,
2480 &recorded_title,
2481 Some("new_title"),
2482 true,
2483 ));
2484 }
2485
2486 #[test]
2487 fn should_remove_when_running_title_mismatches() {
2488 let session = ProjectSession {
2489 liveness_title: Some("recorded_title".to_string()),
2490 };
2491 let recorded_title = Some("recorded_title".to_string());
2492
2493 assert!(should_remove_liveness_session(
2494 &session,
2495 &recorded_title,
2496 Some("different_title"),
2497 true,
2498 ));
2499 }
2500
2501 #[test]
2502 fn should_remove_when_dead() {
2503 let session = ProjectSession {
2504 liveness_title: Some("recorded_title".to_string()),
2505 };
2506 let recorded_title = Some("recorded_title".to_string());
2507
2508 assert!(should_remove_liveness_session(
2509 &session,
2510 &recorded_title,
2511 Some("recorded_title"),
2512 false,
2513 ));
2514 }
2515
2516 #[test]
2517 fn should_not_remove_when_alive_and_title_matches() {
2518 let session = ProjectSession {
2519 liveness_title: Some("recorded_title".to_string()),
2520 };
2521 let recorded_title = Some("recorded_title".to_string());
2522
2523 assert!(!should_remove_liveness_session(
2524 &session,
2525 &recorded_title,
2526 Some("recorded_title"),
2527 true,
2528 ));
2529 }
2530
2531 #[cfg(unix)]
2532 #[test]
2533 fn fd_scan_end_reaches_open_fds_above_the_fd_limit() {
2534 // A caller opened fd 900, then lowered RLIMIT_NOFILE to 256.
2535 assert_eq!(super::fd_scan_end(256, [0, 1, 2, 900]), 901);
2536 assert_eq!(super::fd_scan_end(256, [0, 1, 2, 10]), 256);
2537 assert_eq!(super::fd_scan_end(-1, []), 1 << 16);
2538 assert_eq!(super::fd_scan_end(libc::c_long::MAX, []), 1 << 20);
2539 }
2540
2541 /// Exercises the fcntl fallback directly: on Linux >= 5.11 the spawn path
2542 /// never reaches it, but it is the only implementation on macOS.
2543 #[cfg(unix)]
2544 #[test]
2545 fn cloexec_fd_scan_hides_inherited_fds_from_children() {
2546 use std::os::unix::process::CommandExt;
2547 use std::process::{Command, Stdio};
2548
2549 // A pipe without O_CLOEXEC, like bats' fd 3 or a wrapper script's pipe.
2550 let mut fds = [0; 2];
2551 assert_eq!(unsafe { libc::pipe(fds.as_mut_ptr()) }, 0);
2552 let fd = fds[1];
2553 let child_sees_fd = |scan: bool| {
2554 let mut cmd = Command::new("sh");
2555 cmd.arg("-c")
2556 .arg(format!("[ -e /dev/fd/{fd} ]"))
2557 .stdin(Stdio::null())
2558 .stdout(Stdio::null())
2559 .stderr(Stdio::null());
2560 if scan {
2561 let max_fd = super::max_inherited_fd();
2562 unsafe {
2563 cmd.pre_exec(move || {
2564 super::cloexec_fd_scan(max_fd);
2565 Ok(())
2566 });
2567 }
2568 }
2569 cmd.status().unwrap().success()
2570 };
2571 let without_scan = child_sees_fd(false);
2572 let with_scan = child_sees_fd(true);
2573 unsafe {
2574 libc::close(fds[0]);
2575 libc::close(fds[1]);
2576 }
2577 assert!(without_scan, "control: child should inherit fd {fd}");
2578 assert!(!with_scan, "fd {fd} leaked past cloexec_fd_scan");
2579 }
2580}