1mod adopt;
13mod autostop;
14mod hooks;
15mod ipc_handlers;
16mod lifecycle;
17#[cfg(unix)]
18mod pty;
19mod retry;
20mod state;
21mod watchers;
22
23use crate::daemon_id::DaemonId;
24use crate::daemon_status::DaemonStatus;
25use crate::deps::compute_reverse_stop_order;
26use crate::ipc::server::{IpcServer, IpcServerHandle};
27
28use crate::procs::PROCS;
29use crate::settings::settings;
30use crate::state_file::StateFile;
31use crate::{Result, env};
32use duct::cmd;
33use miette::IntoDiagnostic;
34use once_cell::sync::Lazy;
35use std::collections::HashMap;
36#[cfg(unix)]
37use std::collections::HashSet;
38use std::fs;
39#[cfg(unix)]
40use std::os::unix::fs::PermissionsExt;
41use std::path::PathBuf;
42use std::process::exit;
43use std::sync::atomic;
44use std::sync::atomic::{AtomicBool, AtomicU32};
45use std::time::Duration;
46#[cfg(unix)]
47use tokio::signal::unix::SignalKind;
48use tokio::sync::{Mutex, Notify};
49use tokio::task::JoinHandle;
50use tokio::{signal, time};
51
52#[cfg(all(unix, not(target_os = "linux")))]
61pub(crate) static REAPED_STATUSES: Lazy<Mutex<HashMap<u32, i32>>> =
62 Lazy::new(|| Mutex::new(HashMap::new()));
63
64pub(crate) use state::UpsertDaemonOpts;
66
67pub struct Supervisor {
68 pub(crate) state_file: Mutex<StateFile>,
69 pub(crate) pending_notifications: Mutex<Vec<(log::LevelFilter, String)>>,
70 pub(crate) last_refreshed_at: Mutex<time::Instant>,
71 pub(crate) pending_autostops: Mutex<HashMap<DaemonId, time::Instant>>,
73 pub(crate) ipc_shutdown: Mutex<Option<IpcServerHandle>>,
75 pub(crate) hook_tasks: Mutex<Vec<JoinHandle<()>>>,
77 pub(crate) active_monitors: AtomicU32,
81 pub(crate) monitor_done: Notify,
84 pub(crate) proxy_cancel: Mutex<Option<tokio_util::sync::CancellationToken>>,
87 pub(crate) proxy_task: Mutex<Option<JoinHandle<()>>>,
89 pub(crate) mdns_publisher:
92 Mutex<Option<std::sync::Arc<tokio::sync::Mutex<crate::proxy::mdns::MdnsPublisher>>>>,
93 pub(crate) lan_monitor_task: Mutex<Option<JoinHandle<()>>>,
95 pub(crate) flush_cancel: std::sync::Mutex<Option<tokio_util::sync::CancellationToken>>,
97 pub(crate) monitored: std::sync::Mutex<HashMap<DaemonId, adopt::MonitorEntry>>,
103}
104
105pub(crate) fn interval_duration() -> Duration {
106 settings().general_interval()
107}
108
109pub static SUPERVISOR: Lazy<Supervisor> =
110 Lazy::new(|| Supervisor::new().expect("Error creating supervisor"));
111
112pub fn start_if_not_running() -> Result<()> {
113 let sf = StateFile::get();
114 if let Some(d) = sf.daemons.get(&DaemonId::pitchfork())
115 && let Some(pid) = d.pid
116 && PROCS.is_running(pid)
117 {
118 return Ok(());
119 }
120 start_in_background()
121}
122
123pub fn start_in_background() -> Result<()> {
124 debug!("starting supervisor in background");
125 let log_file = &*env::PITCHFORK_LOG_FILE;
129 if let Some(parent) = log_file.parent() {
130 let _ = fs::create_dir_all(parent);
131 }
132 #[cfg(unix)]
133 fix_state_dir_permissions();
134
135 #[cfg(unix)]
137 {
138 let stderr_file = fs::OpenOptions::new()
139 .create(true)
140 .append(true)
141 .open(log_file)
142 .into_diagnostic()?;
143 cmd!(&*env::PITCHFORK_BIN, "supervisor", "run")
144 .stdin_null()
145 .stdout_null()
146 .stderr_file(stderr_file)
147 .start()
148 .into_diagnostic()?;
149 }
150
151 #[cfg(windows)]
162 {
163 use windows_sys::Win32::Foundation::{CloseHandle, FALSE};
164 use windows_sys::Win32::System::Threading::{
165 CREATE_NO_WINDOW, CreateProcessW, DETACHED_PROCESS, PROCESS_INFORMATION, STARTUPINFOW,
166 };
167
168 let mut si: STARTUPINFOW = unsafe { std::mem::zeroed() };
176 si.cb = std::mem::size_of::<STARTUPINFOW>() as u32;
177
178 let bin_path = &*env::PITCHFORK_BIN;
179 let mut cmd_line: Vec<u16> = format!("\"{}\" supervisor run\0", bin_path.to_string_lossy())
180 .encode_utf16()
181 .collect();
182
183 let mut pi: PROCESS_INFORMATION = unsafe { std::mem::zeroed() };
184 let ok = unsafe {
185 CreateProcessW(
186 std::ptr::null(),
187 cmd_line.as_mut_ptr(),
188 std::ptr::null(),
189 std::ptr::null(),
190 FALSE, DETACHED_PROCESS | CREATE_NO_WINDOW,
192 std::ptr::null(),
193 std::ptr::null(),
194 &si,
195 &mut pi,
196 )
197 };
198
199 if ok == 0 {
200 return Err(miette::miette!(
201 "CreateProcessW failed for supervisor: {}",
202 std::io::Error::last_os_error()
203 ));
204 }
205
206 unsafe {
208 CloseHandle(pi.hProcess);
209 CloseHandle(pi.hThread);
210 }
211 }
212
213 Ok(())
214}
215
216fn should_remove_liveness_session(
225 session: &crate::state_file::ProjectSession,
226 recorded_title: &Option<String>,
227 current_title: Option<&str>,
228 is_running: bool,
229) -> bool {
230 if session.liveness_title.as_ref() != recorded_title.as_ref() {
234 return false;
235 }
236 if !is_running {
238 return true;
239 }
240 match (recorded_title.as_deref(), current_title) {
245 (Some(recorded), Some(current)) => recorded != current,
246 _ => false,
247 }
248}
249
250impl Supervisor {
251 pub fn new() -> Result<Self> {
252 Ok(Self {
253 state_file: Mutex::new(StateFile::read(&*env::PITCHFORK_STATE_FILE).unwrap_or_else(
254 |e| {
255 warn!("failed to read state file, starting with empty state: {e}");
256 StateFile::new(env::PITCHFORK_STATE_FILE.clone())
257 },
258 )),
259 last_refreshed_at: Mutex::new(time::Instant::now()),
260 pending_notifications: Mutex::new(vec![]),
261 pending_autostops: Mutex::new(HashMap::new()),
262 ipc_shutdown: Mutex::new(None),
263 hook_tasks: Mutex::new(Vec::new()),
264 active_monitors: AtomicU32::new(0),
265 monitor_done: Notify::new(),
266 proxy_cancel: Mutex::new(None),
267 proxy_task: Mutex::new(None),
268 mdns_publisher: Mutex::new(None),
269 lan_monitor_task: Mutex::new(None),
270 flush_cancel: std::sync::Mutex::new(None),
271 monitored: std::sync::Mutex::new(HashMap::new()),
272 })
273 }
274
275 pub async fn start(
276 &self,
277 is_boot: bool,
278 container: bool,
279 web_port: Option<u16>,
280 web_path: Option<String>,
281 ) -> Result<()> {
282 #[cfg(unix)]
287 fix_state_dir_permissions();
288
289 let pid = std::process::id();
290 PROCS.refresh_pids(&[pid]);
292 let container_mode = container || settings().supervisor.container;
294 if container_mode {
295 info!("Starting supervisor in container/PID1 mode with pid {pid}");
296 } else {
297 info!("Starting supervisor with pid {pid}");
298 }
299
300 cleanup_orphaned_daemons(self).await;
305
306 self.upsert_daemon(
307 UpsertDaemonOpts::builder(DaemonId::pitchfork())
308 .set(|o| {
309 o.pid = Some(pid);
310 o.status = DaemonStatus::Running;
311 })
312 .build(),
313 )
314 .await?;
315 #[cfg(unix)]
316 fix_state_dir_permissions();
317
318 if is_boot {
320 info!("Boot start mode enabled, starting boot_start daemons");
321 self.start_boot_daemons().await?;
322 }
323
324 self.interval_watch()?;
325
326 if let Err(e) = self.check_cron_schedules().await {
333 error!("failed to check cron schedules on startup: {e}");
334 }
335
336 self.cron_watch()?;
337 self.signals()?;
338 self.daemon_file_watch()?;
339
340 #[cfg(unix)]
342 if container_mode {
343 self.reap_zombies()?;
344 }
345
346 let s = settings();
348 let effective_port = web_port.or_else(|| {
349 if s.web.auto_start {
350 match u16::try_from(s.web.bind_port).ok().filter(|&p| p > 0) {
351 Some(p) => Some(p),
352 None => {
353 error!(
354 "web.bind_port {} is out of valid port range (1-65535), web UI disabled",
355 s.web.bind_port
356 );
357 None
358 }
359 }
360 } else {
361 None
362 }
363 });
364 let effective_path = web_path.or_else(|| {
366 let bp = s.web.base_path.clone();
367 if bp.is_empty() { None } else { Some(bp) }
368 });
369 if let Some(port) = effective_port {
370 tokio::spawn(async move {
371 if let Err(e) = crate::web::serve(port, effective_path).await {
372 error!("Web server error: {e}");
373 }
374 });
375 }
376
377 let api_port = if s.api.auto_start {
379 match u16::try_from(s.api.bind_port).ok().filter(|&p| p > 0) {
380 Some(p) => Some(p),
381 None => {
382 error!(
383 "api.bind_port {} is out of valid port range (1-65535), API server disabled",
384 s.api.bind_port
385 );
386 None
387 }
388 }
389 } else {
390 None
391 };
392 if let Some(port) = api_port {
393 tokio::spawn(async move {
394 if let Err(e) = crate::web::serve_api(port, None).await {
395 error!("API server error: {e}");
396 }
397 });
398 }
399
400 if s.proxy.enable {
402 #[cfg(feature = "proxy-tls")]
406 if s.proxy.https {
407 let proxy_dir = crate::env::PITCHFORK_STATE_DIR.join("proxy");
408 let ca_cert_path = proxy_dir.join("ca.pem");
409 let ca_key_path = proxy_dir.join("ca-key.pem");
410 if !ca_cert_path.exists() || !ca_key_path.exists() {
411 match crate::proxy::server::generate_ca(&ca_cert_path, &ca_key_path) {
412 Ok(()) => {
413 info!(
414 "Generated local CA certificate at {}",
415 ca_cert_path.display()
416 );
417 }
418 Err(e) => {
419 error!("Failed to generate CA certificate: {e}");
420 }
421 }
422 }
423
424 if s.proxy.auto_trust && ca_cert_path.exists() {
428 use crate::proxy::trust::{AutoTrustResult, auto_trust};
429 match auto_trust(&ca_cert_path) {
430 AutoTrustResult::AlreadyTrusted => {}
431 AutoTrustResult::Trusted => {
432 info!("CA certificate auto-trusted in system store");
433 }
434 AutoTrustResult::NotTrusted { reason } => {
435 warn!("Auto-trust skipped: {reason}");
436 warn!("Run `pitchfork proxy trust` to install manually");
437 }
438 }
439 }
440 }
441 let (bind_tx, bind_rx) = tokio::sync::oneshot::channel();
445 let proxy_cancel = tokio_util::sync::CancellationToken::new();
446 let proxy_cancel_clone = proxy_cancel.clone();
447 *self.proxy_cancel.lock().await = Some(proxy_cancel);
448 let proxy_task = tokio::spawn(async move {
449 if let Err(e) = crate::proxy::server::serve(bind_tx, proxy_cancel_clone).await {
450 error!("Proxy server error: {e}");
451 }
452 });
453 *self.proxy_task.lock().await = Some(proxy_task);
454 match bind_rx.await {
455 Ok(Ok(())) => {
456 info!("Proxy server bound successfully");
457 self.start_mdns().await;
458 }
459 Ok(Err(msg)) => {
460 error!("{msg}");
461 self.add_notification(log::LevelFilter::Error, msg).await;
462 }
463 Err(_) => {
464 }
468 }
469 }
470
471 tokio::spawn(async {
474 crate::proxy::server::get_cached_slugs().await;
475 });
476
477 let (ipc, ipc_handle) = IpcServer::new()?;
478 *self.ipc_shutdown.lock().await = Some(ipc_handle);
479 self.start_state_flush_task();
480 self.conn_watch(ipc).await
481 }
482
483 async fn start_mdns(&self) {
485 let s = crate::settings::settings();
486 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
487 if !s.proxy.enable || !lan_enabled {
488 return;
489 }
490
491 let lan_ip = if !s.proxy.lan_ip.is_empty() {
492 match s.proxy.lan_ip.parse::<std::net::Ipv4Addr>() {
493 Ok(ip) => Some(ip),
494 Err(e) => {
495 error!(
496 "proxy.lan_ip {:?} is not a valid IPv4 address: {e}",
497 s.proxy.lan_ip
498 );
499 return;
500 }
501 }
502 } else {
503 match crate::proxy::lan_ip::detect_lan_ip().await {
504 Some(ip) => Some(ip),
505 None => {
506 error!(
507 "LAN mode is enabled but no LAN IP address could be detected. \
508 Set proxy.lan_ip to a specific address, or ensure you are connected to a network."
509 );
510 return;
511 }
512 }
513 };
514
515 let Some(lan_ip) = lan_ip else { return };
516 let port = u16::try_from(s.proxy.port).unwrap_or(443);
517
518 let Some(mut publisher) = crate::proxy::mdns::MdnsPublisher::new(lan_ip) else {
519 error!("Failed to start mDNS publisher. Is Avahi (Linux) or Bonjour (macOS) running?");
520 return;
521 };
522
523 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
525 for slug in slugs.keys() {
526 let hostname = format!("{slug}.local");
527 publisher.publish(&hostname, port);
528 }
529
530 log::info!(
531 "LAN mode: mDNS publishing on {lan_ip}, {} slug(s) registered",
532 slugs.len()
533 );
534
535 let publisher = std::sync::Arc::new(tokio::sync::Mutex::new(publisher));
536
537 let ip_pinned = !s.proxy.lan_ip.is_empty();
539 if !ip_pinned {
540 let monitor_cancel = self.proxy_cancel.lock().await.clone();
541 let publisher_clone = publisher.clone();
542 let task = tokio::spawn(async move {
543 let mut last_ip = lan_ip;
544 let interval = std::time::Duration::from_secs(5);
545 let mut ticker = tokio::time::interval(interval);
546 ticker.tick().await; loop {
548 ticker.tick().await;
549 if let Some(cancel) = monitor_cancel.as_ref()
550 && cancel.is_cancelled()
551 {
552 break;
553 }
554 if let Some(new_ip) =
555 crate::proxy::lan_ip::detect_lan_ip_if_changed(last_ip).await
556 {
557 log::info!("LAN IP changed: {last_ip} → {new_ip}");
558 last_ip = new_ip;
559 let mut pub_guard = publisher_clone.lock().await;
560 pub_guard.republish_all(new_ip, port);
561 }
562 }
563 });
564 *self.lan_monitor_task.lock().await = Some(task);
565 }
566
567 *self.mdns_publisher.lock().await = Some(publisher);
568 }
569
570 async fn sync_mdns(&self) {
575 let publisher = {
578 let guard = self.mdns_publisher.lock().await;
579 match guard.as_ref() {
580 Some(p) => p.clone(),
581 None => {
582 debug!("sync_mdns: mDNS publisher not active, skipping");
583 return;
584 }
585 }
586 };
587
588 let s = crate::settings::settings();
589 let port = u16::try_from(s.proxy.port).unwrap_or(443);
590
591 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
592 let mut pub_guard = publisher.lock().await;
593
594 let current_keys: Vec<&String> = slugs.keys().collect();
596 let registered: Vec<String> = pub_guard.registered_hostnames();
597 for hostname in ®istered {
598 let slug = hostname.strip_suffix(".local").unwrap_or(hostname);
600 if !current_keys.iter().any(|k| k.as_str() == slug) {
601 log::info!("mDNS: unpublishing removed slug {slug}");
602 pub_guard.unpublish(hostname);
603 }
604 }
605
606 for slug in slugs.keys() {
608 let hostname = format!("{slug}.local");
609 if !pub_guard.is_published(&hostname) {
610 log::info!("mDNS: publishing new slug {slug}");
611 pub_guard.publish(&hostname, port);
612 }
613 }
614 }
615
616 fn start_state_flush_task(&self) {
620 let cancel = tokio_util::sync::CancellationToken::new();
621 *self.flush_cancel.lock().unwrap() = Some(cancel.clone());
622 tokio::spawn(async move {
623 let mut interval = time::interval(Duration::from_secs(1));
624 interval.set_missed_tick_behavior(time::MissedTickBehavior::Skip);
625 loop {
626 tokio::select! {
627 _ = interval.tick() => {}
628 _ = cancel.cancelled() => {
629 debug!("state flush task received shutdown signal");
630 break;
631 }
632 }
633 let state = SUPERVISOR.state_file.lock().await;
634 if state.is_dirty()
635 && let Err(e) = state.write()
636 {
637 warn!("failed to flush state file: {e}");
638 }
639 }
640 debug!("state flush task exiting");
641 });
642 }
643
644 pub(crate) async fn flush_state(&self) {
645 let state = self.state_file.lock().await;
646 if state.is_dirty()
647 && let Err(e) = state.write()
648 {
649 warn!("failed to flush state file: {e}");
650 }
651 }
652
653 pub(crate) async fn refresh(&self) -> Result<()> {
654 trace!("refreshing");
655
656 let dirs_with_pids = self.get_dirs_with_shell_pids().await;
659 let liveness_sessions = self.get_liveness_sessions().await;
660 let pids_to_check: Vec<u32> = dirs_with_pids
661 .values()
662 .flatten()
663 .copied()
664 .chain(liveness_sessions.iter().map(|(pid, _, _)| *pid))
665 .collect::<std::collections::HashSet<_>>()
666 .into_iter()
667 .collect();
668
669 if pids_to_check.is_empty() {
670 trace!("no tracked PIDs to check, skipping process refresh");
672 } else {
673 debug!("refreshing PIDs: {pids_to_check:?}");
674 PROCS.refresh_pids(&pids_to_check);
675 }
676
677 let mut last_refreshed_at = self.last_refreshed_at.lock().await;
678 *last_refreshed_at = time::Instant::now();
679
680 let mut dirs_to_leave: Vec<PathBuf> = Vec::new();
681
682 #[cfg(unix)]
692 for (dir, pids) in dirs_with_pids {
693 let to_remove = pids
694 .iter()
695 .filter(|pid| !PROCS.is_running(**pid))
696 .collect::<Vec<_>>();
697 for pid in &to_remove {
698 self.remove_shell_pid(**pid).await?
699 }
700 if to_remove.len() == pids.len() {
701 dirs_to_leave.push(dir);
702 }
703 }
704
705 #[cfg(unix)]
718 {
719 let mut state = self.state_file.lock().await;
720 for (pid, dir, recorded_title) in liveness_sessions {
721 let Some(session) = state.get_project_session(pid, &dir) else {
722 continue;
723 };
724 let current_title = PROCS.title(pid);
725 let is_running = PROCS.is_running(pid);
726 debug!(
727 "refresh liveness session pid {pid} dir {} recorded_title={recorded_title:?} current_title={current_title:?} is_running={is_running}",
728 dir.display()
729 );
730 if should_remove_liveness_session(
731 session,
732 &recorded_title,
733 current_title.as_deref(),
734 is_running,
735 ) {
736 warn!(
737 "removing project session pid {pid} dir {} (liveness pid title mismatch or dead)",
738 dir.display()
739 );
740 if state.remove_project_session(pid, &dir).is_some() {
741 dirs_to_leave.push(dir);
742 }
743 }
744 }
745 }
746
747 for dir in dirs_to_leave {
748 self.leave_dir(&dir).await?;
749 }
750
751 self.reconcile_unmonitored_daemons().await;
756
757 self.check_retry().await?;
758 self.process_pending_autostops().await?;
759
760 Ok(())
761 }
762
763 #[cfg(unix)]
788 fn reap_zombies(&self) -> Result<()> {
789 let mut stream = signal::unix::signal(SignalKind::child())
790 .map_err(|e| miette::miette!("Failed to register SIGCHLD handler: {e}"))?;
791 tokio::spawn(async move {
792 loop {
793 stream.recv().await;
794 let managed_pids: HashSet<u32> = SUPERVISOR
796 .state_file
797 .lock()
798 .await
799 .daemons
800 .values()
801 .filter_map(|d| d.pid)
802 .collect();
803 Self::reap_unmanaged_zombies(&managed_pids).await;
805 }
806 });
807 info!("container mode: SIGCHLD zombie reaper installed");
808 Ok(())
809 }
810
811 #[cfg(target_os = "linux")]
817 async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
818 use nix::sys::wait::{Id, WaitPidFlag, WaitStatus, waitid, waitpid};
819 use nix::unistd::Pid;
820
821 loop {
822 let peek_flags = WaitPidFlag::WNOHANG | WaitPidFlag::WNOWAIT | WaitPidFlag::WEXITED;
824 match waitid(Id::All, peek_flags) {
825 Ok(WaitStatus::StillAlive) => break,
826 Ok(status) => {
827 let Some(pid_raw) = status.pid().map(|p| p.as_raw() as u32) else {
828 break;
829 };
830 if managed_pids.contains(&pid_raw) {
831 trace!(
835 "zombie reaper: skipping managed daemon pid {pid_raw}, \
836 leaving for Tokio to reap"
837 );
838 break;
839 }
840 match waitpid(Pid::from_raw(pid_raw as i32), Some(WaitPidFlag::WNOHANG)) {
842 Ok(s) => trace!("reaped orphaned zombie child: {s:?}"),
843 Err(nix::errno::Errno::ECHILD) => break,
844 Err(e) => {
845 trace!("waitpid error reaping pid {pid_raw}: {e}");
846 break;
847 }
848 }
849 }
850 Err(nix::errno::Errno::ECHILD) => break, Err(e) => {
852 trace!("waitid error in zombie reaper: {e}");
853 break;
854 }
855 }
856 }
857 }
858
859 #[cfg(all(unix, not(target_os = "linux")))]
865 async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
866 use nix::sys::wait::{WaitPidFlag, WaitStatus, waitpid};
867
868 loop {
869 match waitpid(None, Some(WaitPidFlag::WNOHANG)) {
870 Ok(WaitStatus::StillAlive) => break,
871 Ok(status) => {
872 let Some(pid) = status.pid().map(|p| p.as_raw() as u32) else {
873 continue;
874 };
875 if managed_pids.contains(&pid) {
876 let exit_code = match status {
878 WaitStatus::Exited(_, code) => code,
879 WaitStatus::Signaled(_, sig, _) => -(sig as i32),
880 _ => -1,
881 };
882 warn!(
883 "zombie reaper reaped managed daemon pid {pid} \
884 (exit_code={exit_code}); stashing status for recovery"
885 );
886 REAPED_STATUSES.lock().await.insert(pid, exit_code);
887 } else {
888 trace!("reaped orphaned zombie child: {status:?}");
889 }
890 }
891 Err(nix::errno::Errno::ECHILD) => break, Err(e) => {
893 trace!("waitpid error in zombie reaper: {e}");
894 break;
895 }
896 }
897 }
898 }
899
900 #[cfg(unix)]
901 fn signals(&self) -> Result<()> {
902 let signals = [
903 SignalKind::terminate(),
904 SignalKind::alarm(),
905 SignalKind::interrupt(),
906 SignalKind::quit(),
907 SignalKind::hangup(),
908 SignalKind::user_defined1(),
909 SignalKind::user_defined2(),
910 ];
911 static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
912 for signal in signals {
913 let stream = match signal::unix::signal(signal) {
914 Ok(s) => s,
915 Err(e) => {
916 warn!("Failed to register signal handler for {signal:?}: {e}");
917 continue;
918 }
919 };
920 tokio::spawn(async move {
921 let mut stream = stream;
922 loop {
923 stream.recv().await;
924 if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
925 exit(1);
926 } else {
927 SUPERVISOR.handle_signal().await;
928 }
929 }
930 });
931 }
932 Ok(())
933 }
934
935 #[cfg(windows)]
936 fn signals(&self) -> Result<()> {
937 tokio::spawn(async move {
938 static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
939 loop {
940 if let Err(e) = signal::ctrl_c().await {
941 error!("Failed to wait for ctrl-c: {}", e);
942 return;
943 }
944 if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
945 exit(1);
946 } else {
947 SUPERVISOR.handle_signal().await;
948 }
949 }
950 });
951 Ok(())
952 }
953
954 async fn handle_signal(&self) {
955 info!("received signal, stopping");
956 self.close().await;
957 exit(0)
958 }
959
960 pub(crate) async fn close(&self) {
961 if let Some(cancel) = self.proxy_cancel.lock().await.take() {
965 cancel.cancel();
966 }
967
968 if let Some(monitor_task) = self.lan_monitor_task.lock().await.take() {
970 monitor_task.abort();
971 }
972
973 if let Some(publisher) = self.mdns_publisher.lock().await.take() {
975 publisher.lock().await.shutdown();
976 }
977
978 if let Some(proxy_task) = self.proxy_task.lock().await.take() {
979 let _ = tokio::time::timeout(Duration::from_secs(12), proxy_task).await;
980 }
981
982 let s = settings();
984 if s.proxy.enable && s.proxy.sync_hosts {
985 crate::proxy::hosts::clean_hosts_file();
986 }
987
988 let pitchfork_id = DaemonId::pitchfork();
989 let active = self.active_daemons().await;
990 let active_ids: Vec<DaemonId> = active
991 .iter()
992 .filter(|d| d.id != pitchfork_id)
993 .map(|d| d.id.clone())
994 .collect();
995
996 let stop_levels = compute_reverse_stop_order(&active_ids);
1001 for level in &stop_levels {
1002 let mut tasks = Vec::new();
1003 for id in level {
1004 let id = id.clone();
1005 tasks.push(tokio::spawn(async move {
1006 if let Err(err) = SUPERVISOR.stop(&id).await {
1007 error!("failed to stop daemon {id}: {err}");
1008 }
1009 }));
1010 }
1011 for task in tasks {
1012 let _ = task.await;
1013 }
1014 }
1015 let _ = self.remove_daemon(&pitchfork_id).await;
1016
1017 if let Some(cancel) = self.flush_cancel.lock().unwrap().take() {
1020 cancel.cancel();
1021 }
1022
1023 {
1026 let state = self.state_file.lock().await;
1027 if state.is_dirty()
1028 && let Err(e) = state.write()
1029 {
1030 warn!("failed to flush state file during shutdown: {e}");
1031 }
1032 }
1033
1034 if let Some(mut handle) = self.ipc_shutdown.lock().await.take() {
1036 handle.shutdown();
1037 }
1038
1039 let drain_timeout = time::sleep(Duration::from_secs(5));
1045 tokio::pin!(drain_timeout);
1046 loop {
1047 if self.active_monitors.load(atomic::Ordering::Acquire) == 0 {
1048 break;
1049 }
1050 tokio::select! {
1051 _ = self.monitor_done.notified() => {}
1052 _ = &mut drain_timeout => {
1053 warn!("timed out waiting for monitoring tasks to register hooks, proceeding with shutdown");
1054 break;
1055 }
1056 }
1057 }
1058 let handles: Vec<JoinHandle<()>> = std::mem::take(&mut *self.hook_tasks.lock().await);
1059 let hook_timeout = Duration::from_secs(30);
1060 for handle in handles {
1061 match time::timeout(hook_timeout, handle).await {
1062 Ok(_) => {} Err(_) => {
1064 warn!(
1065 "hook task did not complete within {hook_timeout:?} during shutdown, skipping"
1066 );
1067 }
1068 }
1069 }
1070
1071 #[cfg(unix)]
1073 let _ = fs::remove_dir_all(&*env::IPC_SOCK_DIR);
1074 }
1075
1076 pub(crate) async fn add_notification(&self, level: log::LevelFilter, message: String) {
1077 self.pending_notifications
1078 .lock()
1079 .await
1080 .push((level, message));
1081 }
1082}
1083
1084#[cfg(unix)]
1102fn fix_state_dir_permissions() {
1103 let state_dir = &*env::PITCHFORK_STATE_DIR;
1104 if let Some((uid, gid)) = state_owner_ids() {
1105 if !state_dir.exists()
1106 && let Err(err) = fs::create_dir_all(state_dir)
1107 {
1108 warn!(
1109 "failed to create state directory for ownership fix at {}: {err}",
1110 state_dir.display()
1111 );
1112 return;
1113 }
1114
1115 chown_recursive(state_dir, uid, gid, true);
1117 debug!(
1118 "chowned state directory to uid={uid} gid={gid} at {}",
1119 state_dir.display()
1120 );
1121 } else {
1122 if !state_dir.exists() {
1123 return;
1124 }
1125
1126 chmod_safe_subtrees(state_dir);
1129 debug!(
1130 "relaxed permissions on safe subtrees at {}",
1131 state_dir.display()
1132 );
1133 }
1134}
1135
1136#[cfg(unix)]
1137pub(crate) fn state_owner_ids() -> Option<(u32, u32)> {
1138 if !nix::unistd::Uid::effective().is_root() {
1139 return None;
1140 }
1141
1142 let s = settings();
1143 let user = s.supervisor.user.trim();
1144 if !user.is_empty() {
1145 return resolve_supervisor_user_ids(user).or_else(|| {
1146 warn!(
1147 "failed to resolve supervisor.user '{user}' for state ownership; falling back to SUDO_UID/SUDO_GID"
1148 );
1149 parse_sudo_ids()
1150 });
1151 }
1152
1153 parse_sudo_ids()
1154}
1155
1156#[cfg(unix)]
1157fn resolve_supervisor_user_ids(user: &str) -> Option<(u32, u32)> {
1158 let user_record = if user.chars().all(|c| c.is_ascii_digit()) {
1159 let uid = user.parse::<u32>().ok()?;
1160 nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1161 .ok()
1162 .flatten()
1163 } else {
1164 nix::unistd::User::from_name(user).ok().flatten()
1165 }?;
1166
1167 Some((user_record.uid.as_raw(), user_record.gid.as_raw()))
1168}
1169
1170#[cfg(unix)]
1176fn parse_sudo_ids() -> Option<(u32, u32)> {
1177 if !nix::unistd::Uid::effective().is_root() {
1178 return None;
1179 }
1180 let uid: u32 = std::env::var("SUDO_UID").ok()?.parse().ok()?;
1181 let gid: u32 = std::env::var("SUDO_GID").ok()?.parse().ok()?;
1182 Some((uid, gid))
1183}
1184
1185#[cfg(unix)]
1188fn chown_recursive(dir: &std::path::Path, uid: u32, gid: u32, skip_proxy: bool) {
1189 let _ = chown_path(dir, uid, gid);
1191
1192 let entries = match std::fs::read_dir(dir) {
1193 Ok(e) => e,
1194 Err(_) => return,
1195 };
1196 for entry in entries.flatten() {
1197 let path = entry.path();
1198 if path.is_dir() {
1199 if skip_proxy
1201 && let Some(name) = path.file_name().and_then(|n| n.to_str())
1202 && name == "proxy"
1203 {
1204 continue;
1205 }
1206 chown_recursive(&path, uid, gid, false);
1207 } else {
1208 let _ = chown_path(&path, uid, gid);
1209 }
1210 }
1211}
1212
1213#[cfg(unix)]
1215fn chown_path(path: &std::path::Path, uid: u32, gid: u32) -> std::io::Result<()> {
1216 use std::ffi::CString;
1217 use std::os::unix::ffi::OsStrExt;
1218 let c_path = CString::new(path.as_os_str().as_bytes())
1219 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
1220 let ret = unsafe { libc::chown(c_path.as_ptr(), uid, gid) };
1221 if ret == 0 {
1222 Ok(())
1223 } else {
1224 Err(std::io::Error::last_os_error())
1225 }
1226}
1227
1228#[cfg(unix)]
1231fn chmod_safe_subtrees(state_dir: &std::path::Path) {
1232 let _ = fs::set_permissions(state_dir, fs::Permissions::from_mode(0o755));
1234
1235 let state_file = state_dir.join("state.toml");
1237 if state_file.exists() {
1238 let _ = fs::set_permissions(&state_file, fs::Permissions::from_mode(0o644));
1239 }
1240
1241 for subdir_name in &["sock", "logs"] {
1243 let subdir = state_dir.join(subdir_name);
1244 if subdir.is_dir() {
1245 chmod_recursive(&subdir);
1246 }
1247 }
1248}
1249
1250async fn cleanup_orphaned_daemons(supervisor: &Supervisor) {
1270 if !settings().supervisor.cleanup_orphans {
1271 return;
1272 }
1273
1274 let candidates: Vec<_> = {
1275 let state = supervisor.state_file.lock().await;
1276 state
1277 .daemons
1278 .values()
1279 .filter(|d| d.id != DaemonId::pitchfork() && d.pid.is_some())
1280 .cloned()
1281 .collect()
1282 };
1283
1284 if candidates.is_empty() {
1285 return;
1286 }
1287
1288 info!(
1289 "checking {} daemon(s) for orphaned processes",
1290 candidates.len()
1291 );
1292
1293 let policy = orphan_policy();
1294 let boot_time = PROCS.boot_time();
1295 let unclean = supervisor_exited_uncleanly(supervisor).await;
1296
1297 for daemon in candidates {
1298 let Some(pid) = daemon.pid else { continue };
1299
1300 PROCS.refresh_pids(&[pid]);
1304
1305 if !PROCS.is_running(pid) {
1306 let status =
1310 unobserved_exit_status(&daemon.status, daemon.boot_time, boot_time, unclean);
1311 reset_daemon_state(supervisor, &daemon.id, status, ExitObservation::Unobserved).await;
1312 continue;
1313 }
1314
1315 let current_start_time = PROCS.start_time(pid);
1321 let current_title = PROCS.title(pid);
1322 let matches = process_identity_matches(
1323 daemon.start_time,
1324 daemon.title.as_deref(),
1325 current_start_time,
1326 current_title.as_deref(),
1327 );
1328
1329 if !matches {
1330 if daemon.start_time.is_some() && current_start_time.is_none() {
1331 warn!(
1332 "could not verify start time for live pid {pid} recorded for daemon {}; retaining running state",
1333 daemon.id,
1334 );
1335 continue;
1336 }
1337 warn!(
1338 "pid {pid} recorded for daemon {} belongs to a different process now (PID recycled); resetting state without killing",
1339 daemon.id,
1340 );
1341 let status =
1344 unobserved_exit_status(&daemon.status, daemon.boot_time, boot_time, unclean);
1345 reset_daemon_state(supervisor, &daemon.id, status, ExitObservation::Unobserved).await;
1346 continue;
1347 }
1348
1349 let Some(expected_start_time) = current_start_time else {
1353 warn!(
1354 "could not read start time for live pid {pid} recorded for daemon {}; retaining running state",
1355 daemon.id,
1356 );
1357 continue;
1358 };
1359
1360 if policy == "adopt" {
1364 supervisor
1365 .adopt_daemon(&daemon, pid, expected_start_time)
1366 .await;
1367 continue;
1368 }
1369
1370 info!("terminating orphaned daemon {} (pid {pid})", daemon.id);
1371
1372 let stop_cfg = daemon.stop_signal.unwrap_or_default();
1373 let termination_result = PROCS
1374 .kill_process_group_if_start_time_matches_async(
1375 pid,
1376 expected_start_time,
1377 stop_cfg.signal.into(),
1378 stop_cfg.timeout,
1379 )
1380 .await;
1381
1382 match termination_result {
1383 Ok(true) => {}
1384 Ok(false) => {
1385 warn!(
1386 "could not securely terminate orphaned daemon {} (pid {pid}); retaining running state",
1387 daemon.id
1388 );
1389 continue;
1390 }
1391 Err(err) => {
1392 warn!(
1393 "failed to terminate orphaned daemon {} (pid {pid}): {err}; retaining running state",
1394 daemon.id
1395 );
1396 continue;
1397 }
1398 }
1399
1400 reset_daemon_state(
1403 supervisor,
1404 &daemon.id,
1405 DaemonStatus::Stopped,
1406 ExitObservation::Terminated,
1407 )
1408 .await;
1409 }
1410}
1411
1412pub(crate) fn orphan_policy() -> String {
1415 let policy = settings().supervisor.orphan_policy.clone();
1416 match policy.as_str() {
1417 "adopt" | "kill" => policy,
1418 other => {
1419 warn!("unknown supervisor.orphan_policy '{other}', defaulting to 'adopt'");
1420 "adopt".to_string()
1421 }
1422 }
1423}
1424
1425fn process_identity_matches(
1430 recorded_start_time: Option<u64>,
1431 recorded_title: Option<&str>,
1432 current_start_time: Option<u64>,
1433 current_title: Option<&str>,
1434) -> bool {
1435 match recorded_start_time {
1436 Some(recorded) => current_start_time == Some(recorded),
1437 None => matches!(
1438 (recorded_title, current_title),
1439 (Some(recorded), Some(current)) if recorded == current
1440 ),
1441 }
1442}
1443
1444#[derive(Clone, Copy, PartialEq, Eq)]
1447pub(crate) enum ExitObservation {
1448 Unobserved,
1468 Terminated,
1472}
1473
1474impl ExitObservation {
1475 pub(crate) fn last_exit_success(self) -> Option<bool> {
1477 match self {
1478 ExitObservation::Unobserved => None,
1479 ExitObservation::Terminated => Some(true),
1480 }
1481 }
1482}
1483
1484async fn reset_daemon_state(
1490 supervisor: &Supervisor,
1491 id: &DaemonId,
1492 status: DaemonStatus,
1493 observation: ExitObservation,
1494) {
1495 let mut state_file = supervisor.state_file.lock().await;
1496 let Some(existing) = state_file.daemons.get(id) else {
1497 return;
1498 };
1499 let mut daemon = existing.clone();
1500 daemon.pid = None;
1501 daemon.title = None;
1502 daemon.start_time = None;
1503 daemon.boot_time = None;
1504 daemon.status = status;
1505 daemon.last_exit_success = observation.last_exit_success();
1506 daemon.active_port = None;
1507 state_file.clear_active_port(id);
1508 state_file.insert_daemon(id, daemon);
1509}
1510
1511const BOOT_TIME_TOLERANCE_SECS: u64 = 2;
1525
1526pub(crate) fn unobserved_exit_status(
1540 status: &DaemonStatus,
1541 recorded_boot_time: Option<u64>,
1542 current_boot_time: u64,
1543 supervisor_exited_uncleanly: bool,
1544) -> DaemonStatus {
1545 let same_boot = recorded_boot_time
1546 .is_some_and(|recorded| recorded.abs_diff(current_boot_time) <= BOOT_TIME_TOLERANCE_SECS);
1547 if status.is_running() && same_boot && supervisor_exited_uncleanly {
1548 DaemonStatus::Errored(-1)
1549 } else {
1550 DaemonStatus::Stopped
1551 }
1552}
1553
1554async fn supervisor_exited_uncleanly(supervisor: &Supervisor) -> bool {
1567 supervisor
1568 .state_file
1569 .lock()
1570 .await
1571 .daemons
1572 .contains_key(&DaemonId::pitchfork())
1573}
1574
1575#[cfg(unix)]
1577fn chmod_recursive(dir: &std::path::Path) {
1578 let _ = fs::set_permissions(dir, fs::Permissions::from_mode(0o755));
1579 let entries = match fs::read_dir(dir) {
1580 Ok(e) => e,
1581 Err(_) => return,
1582 };
1583 for entry in entries.flatten() {
1584 let path = entry.path();
1585 if path.is_dir() {
1586 chmod_recursive(&path);
1587 } else {
1588 let _ = fs::set_permissions(&path, fs::Permissions::from_mode(0o644));
1589 }
1590 }
1591}
1592
1593#[cfg(test)]
1594mod tests {
1595 use super::{
1596 BOOT_TIME_TOLERANCE_SECS, process_identity_matches, should_remove_liveness_session,
1597 unobserved_exit_status,
1598 };
1599 use crate::daemon_status::DaemonStatus;
1600 use crate::state_file::ProjectSession;
1601
1602 const BOOT: u64 = 1_700_000_000;
1603
1604 #[test]
1605 fn unobserved_running_death_in_current_boot_is_retryable() {
1606 assert!(matches!(
1609 unobserved_exit_status(&DaemonStatus::Running, Some(BOOT), BOOT, true),
1610 DaemonStatus::Errored(-1)
1611 ));
1612 }
1613
1614 #[test]
1615 fn unobserved_running_death_from_previous_boot_is_stopped() {
1616 assert!(matches!(
1619 unobserved_exit_status(&DaemonStatus::Running, Some(BOOT - 86_400), BOOT, true),
1620 DaemonStatus::Stopped
1621 ));
1622 }
1623
1624 #[test]
1625 fn unobserved_exit_tolerates_boot_time_jitter() {
1626 let within = BOOT + BOOT_TIME_TOLERANCE_SECS;
1629 assert!(matches!(
1630 unobserved_exit_status(&DaemonStatus::Running, Some(within), BOOT, true),
1631 DaemonStatus::Errored(-1)
1632 ));
1633 let beyond = BOOT + BOOT_TIME_TOLERANCE_SECS + 1;
1634 assert!(matches!(
1635 unobserved_exit_status(&DaemonStatus::Running, Some(beyond), BOOT, true),
1636 DaemonStatus::Stopped
1637 ));
1638 }
1639
1640 #[test]
1641 fn unobserved_exit_after_clean_shutdown_is_stopped() {
1642 assert!(matches!(
1647 unobserved_exit_status(&DaemonStatus::Running, Some(BOOT), BOOT, false),
1648 DaemonStatus::Stopped
1649 ));
1650 }
1651
1652 #[test]
1653 fn unobserved_exit_treats_short_previous_boot_as_previous() {
1654 for gap in [5, 30, 59, 60] {
1658 assert!(
1659 matches!(
1660 unobserved_exit_status(&DaemonStatus::Running, Some(BOOT - gap), BOOT, true),
1661 DaemonStatus::Stopped
1662 ),
1663 "boot {gap}s earlier should be treated as a previous boot"
1664 );
1665 }
1666 }
1667
1668 #[test]
1669 fn unobserved_exit_without_recorded_boot_time_is_stopped() {
1670 assert!(matches!(
1673 unobserved_exit_status(&DaemonStatus::Running, None, BOOT, true),
1674 DaemonStatus::Stopped
1675 ));
1676 }
1677
1678 #[test]
1679 fn unobserved_exit_of_stopping_daemon_is_stopped() {
1680 assert!(matches!(
1683 unobserved_exit_status(&DaemonStatus::Stopping, Some(BOOT), BOOT, true),
1684 DaemonStatus::Stopped
1685 ));
1686 }
1687
1688 #[test]
1689 fn orphan_identity_does_not_match_when_current_identity_is_missing() {
1690 assert!(!process_identity_matches(
1691 Some(123),
1692 Some("daemon"),
1693 None,
1694 None,
1695 ));
1696 assert!(!process_identity_matches(None, Some("daemon"), None, None,));
1697 }
1698
1699 #[test]
1700 fn orphan_identity_requires_recorded_start_time_when_available() {
1701 assert!(process_identity_matches(
1702 Some(123),
1703 Some("old-title"),
1704 Some(123),
1705 Some("new-title"),
1706 ));
1707 assert!(!process_identity_matches(
1708 Some(123),
1709 Some("same-title"),
1710 Some(456),
1711 Some("same-title"),
1712 ));
1713 }
1714
1715 #[test]
1716 fn orphan_identity_falls_back_to_title_for_legacy_state() {
1717 assert!(process_identity_matches(
1718 None,
1719 Some("daemon"),
1720 Some(123),
1721 Some("daemon"),
1722 ));
1723 assert!(!process_identity_matches(
1724 None,
1725 Some("daemon"),
1726 Some(123),
1727 Some("unrelated"),
1728 ));
1729 }
1730
1731 #[test]
1732 fn should_not_remove_when_state_title_differs_from_snapshot() {
1733 let session = ProjectSession {
1736 liveness_title: Some("new_title".to_string()),
1737 };
1738 let recorded_title = Some("old_title".to_string());
1739
1740 assert!(!should_remove_liveness_session(
1741 &session,
1742 &recorded_title,
1743 Some("new_title"),
1744 true,
1745 ));
1746 }
1747
1748 #[test]
1749 fn should_remove_when_running_title_mismatches() {
1750 let session = ProjectSession {
1751 liveness_title: Some("recorded_title".to_string()),
1752 };
1753 let recorded_title = Some("recorded_title".to_string());
1754
1755 assert!(should_remove_liveness_session(
1756 &session,
1757 &recorded_title,
1758 Some("different_title"),
1759 true,
1760 ));
1761 }
1762
1763 #[test]
1764 fn should_remove_when_dead() {
1765 let session = ProjectSession {
1766 liveness_title: Some("recorded_title".to_string()),
1767 };
1768 let recorded_title = Some("recorded_title".to_string());
1769
1770 assert!(should_remove_liveness_session(
1771 &session,
1772 &recorded_title,
1773 Some("recorded_title"),
1774 false,
1775 ));
1776 }
1777
1778 #[test]
1779 fn should_not_remove_when_alive_and_title_matches() {
1780 let session = ProjectSession {
1781 liveness_title: Some("recorded_title".to_string()),
1782 };
1783 let recorded_title = Some("recorded_title".to_string());
1784
1785 assert!(!should_remove_liveness_session(
1786 &session,
1787 &recorded_title,
1788 Some("recorded_title"),
1789 true,
1790 ));
1791 }
1792}