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