1mod autostop;
12mod hooks;
13mod ipc_handlers;
14mod lifecycle;
15#[cfg(unix)]
16mod pty;
17mod retry;
18mod state;
19mod watchers;
20
21use crate::daemon_id::DaemonId;
22use crate::daemon_status::DaemonStatus;
23use crate::deps::compute_reverse_stop_order;
24use crate::ipc::server::{IpcServer, IpcServerHandle};
25
26use crate::procs::PROCS;
27use crate::settings::settings;
28use crate::state_file::StateFile;
29use crate::{Result, env};
30use duct::cmd;
31use miette::IntoDiagnostic;
32use once_cell::sync::Lazy;
33use std::collections::HashMap;
34#[cfg(unix)]
35use std::collections::HashSet;
36use std::fs;
37#[cfg(unix)]
38use std::os::unix::fs::PermissionsExt;
39use std::process::exit;
40use std::sync::atomic;
41use std::sync::atomic::{AtomicBool, AtomicU32};
42use std::time::Duration;
43#[cfg(unix)]
44use tokio::signal::unix::SignalKind;
45use tokio::sync::{Mutex, Notify};
46use tokio::task::JoinHandle;
47use tokio::{signal, time};
48
49#[cfg(all(unix, not(target_os = "linux")))]
58pub(crate) static REAPED_STATUSES: Lazy<Mutex<HashMap<u32, i32>>> =
59 Lazy::new(|| Mutex::new(HashMap::new()));
60
61pub(crate) use state::UpsertDaemonOpts;
63
64pub struct Supervisor {
65 pub(crate) state_file: Mutex<StateFile>,
66 pub(crate) pending_notifications: Mutex<Vec<(log::LevelFilter, String)>>,
67 pub(crate) last_refreshed_at: Mutex<time::Instant>,
68 pub(crate) pending_autostops: Mutex<HashMap<DaemonId, time::Instant>>,
70 pub(crate) ipc_shutdown: Mutex<Option<IpcServerHandle>>,
72 pub(crate) hook_tasks: Mutex<Vec<JoinHandle<()>>>,
74 pub(crate) active_monitors: AtomicU32,
78 pub(crate) monitor_done: Notify,
81 pub(crate) proxy_cancel: Mutex<Option<tokio_util::sync::CancellationToken>>,
84 pub(crate) proxy_task: Mutex<Option<JoinHandle<()>>>,
86 pub(crate) mdns_publisher:
89 Mutex<Option<std::sync::Arc<tokio::sync::Mutex<crate::proxy::mdns::MdnsPublisher>>>>,
90 pub(crate) lan_monitor_task: Mutex<Option<JoinHandle<()>>>,
92 pub(crate) flush_cancel: std::sync::Mutex<Option<tokio_util::sync::CancellationToken>>,
94}
95
96pub(crate) fn interval_duration() -> Duration {
97 settings().general_interval()
98}
99
100pub static SUPERVISOR: Lazy<Supervisor> =
101 Lazy::new(|| Supervisor::new().expect("Error creating supervisor"));
102
103pub fn start_if_not_running() -> Result<()> {
104 let sf = StateFile::get();
105 if let Some(d) = sf.daemons.get(&DaemonId::pitchfork())
106 && let Some(pid) = d.pid
107 && PROCS.is_running(pid)
108 {
109 return Ok(());
110 }
111 start_in_background()
112}
113
114pub fn start_in_background() -> Result<()> {
115 debug!("starting supervisor in background");
116 let log_file = &*env::PITCHFORK_LOG_FILE;
120 if let Some(parent) = log_file.parent() {
121 let _ = fs::create_dir_all(parent);
122 }
123 let stderr_file = fs::OpenOptions::new()
124 .create(true)
125 .append(true)
126 .open(log_file)
127 .into_diagnostic()?;
128 #[cfg(unix)]
129 fix_state_dir_permissions();
130 cmd!(&*env::PITCHFORK_BIN, "supervisor", "run")
131 .stdout_null()
132 .stderr_file(stderr_file)
133 .start()
134 .into_diagnostic()?;
135 Ok(())
136}
137
138impl Supervisor {
139 pub fn new() -> Result<Self> {
140 Ok(Self {
141 state_file: Mutex::new(StateFile::read(&*env::PITCHFORK_STATE_FILE).unwrap_or_else(
142 |e| {
143 warn!("failed to read state file, starting with empty state: {e}");
144 StateFile::new(env::PITCHFORK_STATE_FILE.clone())
145 },
146 )),
147 last_refreshed_at: Mutex::new(time::Instant::now()),
148 pending_notifications: Mutex::new(vec![]),
149 pending_autostops: Mutex::new(HashMap::new()),
150 ipc_shutdown: Mutex::new(None),
151 hook_tasks: Mutex::new(Vec::new()),
152 active_monitors: AtomicU32::new(0),
153 monitor_done: Notify::new(),
154 proxy_cancel: Mutex::new(None),
155 proxy_task: Mutex::new(None),
156 mdns_publisher: Mutex::new(None),
157 lan_monitor_task: Mutex::new(None),
158 flush_cancel: std::sync::Mutex::new(None),
159 })
160 }
161
162 pub async fn start(
163 &self,
164 is_boot: bool,
165 container: bool,
166 web_port: Option<u16>,
167 web_path: Option<String>,
168 ) -> Result<()> {
169 #[cfg(unix)]
174 fix_state_dir_permissions();
175
176 let pid = std::process::id();
177 PROCS.refresh_pids(&[pid]);
179 let container_mode = container || settings().supervisor.container;
181 if container_mode {
182 info!("Starting supervisor in container/PID1 mode with pid {pid}");
183 } else {
184 info!("Starting supervisor with pid {pid}");
185 }
186
187 cleanup_orphaned_daemons(self).await;
192
193 self.upsert_daemon(
194 UpsertDaemonOpts::builder(DaemonId::pitchfork())
195 .set(|o| {
196 o.pid = Some(pid);
197 o.status = DaemonStatus::Running;
198 })
199 .build(),
200 )
201 .await?;
202 #[cfg(unix)]
203 fix_state_dir_permissions();
204
205 if is_boot {
207 info!("Boot start mode enabled, starting boot_start daemons");
208 self.start_boot_daemons().await?;
209 }
210
211 self.interval_watch()?;
212 self.cron_watch()?;
213 self.signals()?;
214 self.daemon_file_watch()?;
215
216 #[cfg(unix)]
218 if container_mode {
219 self.reap_zombies()?;
220 }
221
222 let s = settings();
224 let effective_port = web_port.or_else(|| {
225 if s.web.auto_start {
226 match u16::try_from(s.web.bind_port).ok().filter(|&p| p > 0) {
227 Some(p) => Some(p),
228 None => {
229 error!(
230 "web.bind_port {} is out of valid port range (1-65535), web UI disabled",
231 s.web.bind_port
232 );
233 None
234 }
235 }
236 } else {
237 None
238 }
239 });
240 let effective_path = web_path.or_else(|| {
242 let bp = s.web.base_path.clone();
243 if bp.is_empty() { None } else { Some(bp) }
244 });
245 if let Some(port) = effective_port {
246 tokio::spawn(async move {
247 if let Err(e) = crate::web::serve(port, effective_path).await {
248 error!("Web server error: {e}");
249 }
250 });
251 }
252
253 let api_port = if s.api.auto_start {
255 match u16::try_from(s.api.bind_port).ok().filter(|&p| p > 0) {
256 Some(p) => Some(p),
257 None => {
258 error!(
259 "api.bind_port {} is out of valid port range (1-65535), API server disabled",
260 s.api.bind_port
261 );
262 None
263 }
264 }
265 } else {
266 None
267 };
268 if let Some(port) = api_port {
269 tokio::spawn(async move {
270 if let Err(e) = crate::web::serve_api(port, None).await {
271 error!("API server error: {e}");
272 }
273 });
274 }
275
276 if s.proxy.enable {
278 #[cfg(feature = "proxy-tls")]
282 if s.proxy.https {
283 let proxy_dir = crate::env::PITCHFORK_STATE_DIR.join("proxy");
284 let ca_cert_path = proxy_dir.join("ca.pem");
285 let ca_key_path = proxy_dir.join("ca-key.pem");
286 if !ca_cert_path.exists() || !ca_key_path.exists() {
287 match crate::proxy::server::generate_ca(&ca_cert_path, &ca_key_path) {
288 Ok(()) => {
289 info!(
290 "Generated local CA certificate at {}",
291 ca_cert_path.display()
292 );
293 }
294 Err(e) => {
295 error!("Failed to generate CA certificate: {e}");
296 }
297 }
298 }
299
300 if s.proxy.auto_trust && ca_cert_path.exists() {
304 use crate::proxy::trust::{AutoTrustResult, auto_trust};
305 match auto_trust(&ca_cert_path) {
306 AutoTrustResult::AlreadyTrusted => {}
307 AutoTrustResult::Trusted => {
308 info!("CA certificate auto-trusted in system store");
309 }
310 AutoTrustResult::NotTrusted { reason } => {
311 warn!("Auto-trust skipped: {reason}");
312 warn!("Run `pitchfork proxy trust` to install manually");
313 }
314 }
315 }
316 }
317 let (bind_tx, bind_rx) = tokio::sync::oneshot::channel();
321 let proxy_cancel = tokio_util::sync::CancellationToken::new();
322 let proxy_cancel_clone = proxy_cancel.clone();
323 *self.proxy_cancel.lock().await = Some(proxy_cancel);
324 let proxy_task = tokio::spawn(async move {
325 if let Err(e) = crate::proxy::server::serve(bind_tx, proxy_cancel_clone).await {
326 error!("Proxy server error: {e}");
327 }
328 });
329 *self.proxy_task.lock().await = Some(proxy_task);
330 match bind_rx.await {
331 Ok(Ok(())) => {
332 info!("Proxy server bound successfully");
333 self.start_mdns().await;
334 }
335 Ok(Err(msg)) => {
336 error!("{msg}");
337 self.add_notification(log::LevelFilter::Error, msg).await;
338 }
339 Err(_) => {
340 }
344 }
345 }
346
347 tokio::spawn(async {
350 crate::proxy::server::get_cached_slugs().await;
351 });
352
353 let (ipc, ipc_handle) = IpcServer::new()?;
354 *self.ipc_shutdown.lock().await = Some(ipc_handle);
355 self.start_state_flush_task();
356 self.conn_watch(ipc).await
357 }
358
359 async fn start_mdns(&self) {
361 let s = crate::settings::settings();
362 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
363 if !s.proxy.enable || !lan_enabled {
364 return;
365 }
366
367 let lan_ip = if !s.proxy.lan_ip.is_empty() {
368 match s.proxy.lan_ip.parse::<std::net::Ipv4Addr>() {
369 Ok(ip) => Some(ip),
370 Err(e) => {
371 error!(
372 "proxy.lan_ip {:?} is not a valid IPv4 address: {e}",
373 s.proxy.lan_ip
374 );
375 return;
376 }
377 }
378 } else {
379 match crate::proxy::lan_ip::detect_lan_ip().await {
380 Some(ip) => Some(ip),
381 None => {
382 error!(
383 "LAN mode is enabled but no LAN IP address could be detected. \
384 Set proxy.lan_ip to a specific address, or ensure you are connected to a network."
385 );
386 return;
387 }
388 }
389 };
390
391 let Some(lan_ip) = lan_ip else { return };
392 let port = u16::try_from(s.proxy.port).unwrap_or(443);
393
394 let Some(mut publisher) = crate::proxy::mdns::MdnsPublisher::new(lan_ip) else {
395 error!("Failed to start mDNS publisher. Is Avahi (Linux) or Bonjour (macOS) running?");
396 return;
397 };
398
399 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
401 for slug in slugs.keys() {
402 let hostname = format!("{slug}.local");
403 publisher.publish(&hostname, port);
404 }
405
406 log::info!(
407 "LAN mode: mDNS publishing on {lan_ip}, {} slug(s) registered",
408 slugs.len()
409 );
410
411 let publisher = std::sync::Arc::new(tokio::sync::Mutex::new(publisher));
412
413 let ip_pinned = !s.proxy.lan_ip.is_empty();
415 if !ip_pinned {
416 let monitor_cancel = self.proxy_cancel.lock().await.clone();
417 let publisher_clone = publisher.clone();
418 let task = tokio::spawn(async move {
419 let mut last_ip = lan_ip;
420 let interval = std::time::Duration::from_secs(5);
421 let mut ticker = tokio::time::interval(interval);
422 ticker.tick().await; loop {
424 ticker.tick().await;
425 if let Some(cancel) = monitor_cancel.as_ref() {
426 if cancel.is_cancelled() {
427 break;
428 }
429 }
430 if let Some(new_ip) =
431 crate::proxy::lan_ip::detect_lan_ip_if_changed(last_ip).await
432 {
433 log::info!("LAN IP changed: {last_ip} → {new_ip}");
434 last_ip = new_ip;
435 let mut pub_guard = publisher_clone.lock().await;
436 pub_guard.republish_all(new_ip, port);
437 }
438 }
439 });
440 *self.lan_monitor_task.lock().await = Some(task);
441 }
442
443 *self.mdns_publisher.lock().await = Some(publisher);
444 }
445
446 async fn sync_mdns(&self) {
451 let publisher = {
454 let guard = self.mdns_publisher.lock().await;
455 match guard.as_ref() {
456 Some(p) => p.clone(),
457 None => {
458 debug!("sync_mdns: mDNS publisher not active, skipping");
459 return;
460 }
461 }
462 };
463
464 let s = crate::settings::settings();
465 let port = u16::try_from(s.proxy.port).unwrap_or(443);
466
467 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
468 let mut pub_guard = publisher.lock().await;
469
470 let current_keys: Vec<&String> = slugs.keys().collect();
472 let registered: Vec<String> = pub_guard.registered_hostnames();
473 for hostname in ®istered {
474 let slug = hostname.strip_suffix(".local").unwrap_or(hostname);
476 if !current_keys.iter().any(|k| k.as_str() == slug) {
477 log::info!("mDNS: unpublishing removed slug {slug}");
478 pub_guard.unpublish(hostname);
479 }
480 }
481
482 for slug in slugs.keys() {
484 let hostname = format!("{slug}.local");
485 if !pub_guard.is_published(&hostname) {
486 log::info!("mDNS: publishing new slug {slug}");
487 pub_guard.publish(&hostname, port);
488 }
489 }
490 }
491
492 fn start_state_flush_task(&self) {
496 let cancel = tokio_util::sync::CancellationToken::new();
497 *self.flush_cancel.lock().unwrap() = Some(cancel.clone());
498 tokio::spawn(async move {
499 let mut interval = time::interval(Duration::from_secs(1));
500 interval.set_missed_tick_behavior(time::MissedTickBehavior::Skip);
501 loop {
502 tokio::select! {
503 _ = interval.tick() => {}
504 _ = cancel.cancelled() => {
505 debug!("state flush task received shutdown signal");
506 break;
507 }
508 }
509 let state = SUPERVISOR.state_file.lock().await;
510 if state.is_dirty() {
511 if let Err(e) = state.write() {
512 warn!("failed to flush state file: {e}");
513 }
514 }
515 }
516 debug!("state flush task exiting");
517 });
518 }
519
520 pub(crate) async fn flush_state(&self) {
521 let state = self.state_file.lock().await;
522 if state.is_dirty() {
523 if let Err(e) = state.write() {
524 warn!("failed to flush state file: {e}");
525 }
526 }
527 }
528
529 pub(crate) async fn refresh(&self) -> Result<()> {
530 trace!("refreshing");
531
532 let dirs_with_pids = self.get_dirs_with_shell_pids().await;
533
534 let mut last_refreshed_at = self.last_refreshed_at.lock().await;
535 *last_refreshed_at = time::Instant::now();
536
537 for (dir, pids) in dirs_with_pids {
538 let to_remove = pids
539 .iter()
540 .filter(|pid| !PROCS.is_running(**pid))
541 .collect::<Vec<_>>();
542 for pid in &to_remove {
543 self.remove_shell_pid(**pid).await?
544 }
545 if to_remove.len() == pids.len() {
546 self.leave_dir(&dir).await?;
547 }
548 }
549
550 self.check_retry().await?;
551 self.process_pending_autostops().await?;
552
553 Ok(())
554 }
555
556 #[cfg(unix)]
581 fn reap_zombies(&self) -> Result<()> {
582 let mut stream = signal::unix::signal(SignalKind::child())
583 .map_err(|e| miette::miette!("Failed to register SIGCHLD handler: {e}"))?;
584 tokio::spawn(async move {
585 loop {
586 stream.recv().await;
587 let managed_pids: HashSet<u32> = SUPERVISOR
589 .state_file
590 .lock()
591 .await
592 .daemons
593 .values()
594 .filter_map(|d| d.pid)
595 .collect();
596 Self::reap_unmanaged_zombies(&managed_pids).await;
598 }
599 });
600 info!("container mode: SIGCHLD zombie reaper installed");
601 Ok(())
602 }
603
604 #[cfg(target_os = "linux")]
610 async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
611 use nix::sys::wait::{Id, WaitPidFlag, WaitStatus, waitid, waitpid};
612 use nix::unistd::Pid;
613
614 loop {
615 let peek_flags = WaitPidFlag::WNOHANG | WaitPidFlag::WNOWAIT | WaitPidFlag::WEXITED;
617 match waitid(Id::All, peek_flags) {
618 Ok(WaitStatus::StillAlive) => break,
619 Ok(status) => {
620 let Some(pid_raw) = status.pid().map(|p| p.as_raw() as u32) else {
621 break;
622 };
623 if managed_pids.contains(&pid_raw) {
624 trace!(
628 "zombie reaper: skipping managed daemon pid {pid_raw}, \
629 leaving for Tokio to reap"
630 );
631 break;
632 }
633 match waitpid(Pid::from_raw(pid_raw as i32), Some(WaitPidFlag::WNOHANG)) {
635 Ok(s) => trace!("reaped orphaned zombie child: {s:?}"),
636 Err(nix::errno::Errno::ECHILD) => break,
637 Err(e) => {
638 trace!("waitpid error reaping pid {pid_raw}: {e}");
639 break;
640 }
641 }
642 }
643 Err(nix::errno::Errno::ECHILD) => break, Err(e) => {
645 trace!("waitid error in zombie reaper: {e}");
646 break;
647 }
648 }
649 }
650 }
651
652 #[cfg(all(unix, not(target_os = "linux")))]
658 async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
659 use nix::sys::wait::{WaitPidFlag, WaitStatus, waitpid};
660
661 loop {
662 match waitpid(None, Some(WaitPidFlag::WNOHANG)) {
663 Ok(WaitStatus::StillAlive) => break,
664 Ok(status) => {
665 let Some(pid) = status.pid().map(|p| p.as_raw() as u32) else {
666 continue;
667 };
668 if managed_pids.contains(&pid) {
669 let exit_code = match status {
671 WaitStatus::Exited(_, code) => code,
672 WaitStatus::Signaled(_, sig, _) => -(sig as i32),
673 _ => -1,
674 };
675 warn!(
676 "zombie reaper reaped managed daemon pid {pid} \
677 (exit_code={exit_code}); stashing status for recovery"
678 );
679 REAPED_STATUSES.lock().await.insert(pid, exit_code);
680 } else {
681 trace!("reaped orphaned zombie child: {status:?}");
682 }
683 }
684 Err(nix::errno::Errno::ECHILD) => break, Err(e) => {
686 trace!("waitpid error in zombie reaper: {e}");
687 break;
688 }
689 }
690 }
691 }
692
693 #[cfg(unix)]
694 fn signals(&self) -> Result<()> {
695 let signals = [
696 SignalKind::terminate(),
697 SignalKind::alarm(),
698 SignalKind::interrupt(),
699 SignalKind::quit(),
700 SignalKind::hangup(),
701 SignalKind::user_defined1(),
702 SignalKind::user_defined2(),
703 ];
704 static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
705 for signal in signals {
706 let stream = match signal::unix::signal(signal) {
707 Ok(s) => s,
708 Err(e) => {
709 warn!("Failed to register signal handler for {signal:?}: {e}");
710 continue;
711 }
712 };
713 tokio::spawn(async move {
714 let mut stream = stream;
715 loop {
716 stream.recv().await;
717 if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
718 exit(1);
719 } else {
720 SUPERVISOR.handle_signal().await;
721 }
722 }
723 });
724 }
725 Ok(())
726 }
727
728 #[cfg(windows)]
729 fn signals(&self) -> Result<()> {
730 tokio::spawn(async move {
731 static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
732 loop {
733 if let Err(e) = signal::ctrl_c().await {
734 error!("Failed to wait for ctrl-c: {}", e);
735 return;
736 }
737 if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
738 exit(1);
739 } else {
740 SUPERVISOR.handle_signal().await;
741 }
742 }
743 });
744 Ok(())
745 }
746
747 async fn handle_signal(&self) {
748 info!("received signal, stopping");
749 self.close().await;
750 exit(0)
751 }
752
753 pub(crate) async fn close(&self) {
754 if let Some(cancel) = self.proxy_cancel.lock().await.take() {
758 cancel.cancel();
759 }
760
761 if let Some(monitor_task) = self.lan_monitor_task.lock().await.take() {
763 monitor_task.abort();
764 }
765
766 if let Some(publisher) = self.mdns_publisher.lock().await.take() {
768 publisher.lock().await.shutdown();
769 }
770
771 if let Some(proxy_task) = self.proxy_task.lock().await.take() {
772 let _ = tokio::time::timeout(Duration::from_secs(12), proxy_task).await;
773 }
774
775 let s = settings();
777 if s.proxy.enable && s.proxy.sync_hosts {
778 crate::proxy::hosts::clean_hosts_file();
779 }
780
781 let pitchfork_id = DaemonId::pitchfork();
782 let active = self.active_daemons().await;
783 let active_ids: Vec<DaemonId> = active
784 .iter()
785 .filter(|d| d.id != pitchfork_id)
786 .map(|d| d.id.clone())
787 .collect();
788
789 let stop_levels = compute_reverse_stop_order(&active_ids);
794 for level in &stop_levels {
795 let mut tasks = Vec::new();
796 for id in level {
797 let id = id.clone();
798 tasks.push(tokio::spawn(async move {
799 if let Err(err) = SUPERVISOR.stop(&id).await {
800 error!("failed to stop daemon {id}: {err}");
801 }
802 }));
803 }
804 for task in tasks {
805 let _ = task.await;
806 }
807 }
808 let _ = self.remove_daemon(&pitchfork_id).await;
809
810 if let Some(cancel) = self.flush_cancel.lock().unwrap().take() {
813 cancel.cancel();
814 }
815
816 {
819 let state = self.state_file.lock().await;
820 if state.is_dirty() {
821 if let Err(e) = state.write() {
822 warn!("failed to flush state file during shutdown: {e}");
823 }
824 }
825 }
826
827 if let Some(mut handle) = self.ipc_shutdown.lock().await.take() {
829 handle.shutdown();
830 }
831
832 let drain_timeout = time::sleep(Duration::from_secs(5));
838 tokio::pin!(drain_timeout);
839 loop {
840 if self.active_monitors.load(atomic::Ordering::Acquire) == 0 {
841 break;
842 }
843 tokio::select! {
844 _ = self.monitor_done.notified() => {}
845 _ = &mut drain_timeout => {
846 warn!("timed out waiting for monitoring tasks to register hooks, proceeding with shutdown");
847 break;
848 }
849 }
850 }
851 let handles: Vec<JoinHandle<()>> = std::mem::take(&mut *self.hook_tasks.lock().await);
852 let hook_timeout = Duration::from_secs(30);
853 for handle in handles {
854 match time::timeout(hook_timeout, handle).await {
855 Ok(_) => {} Err(_) => {
857 warn!(
858 "hook task did not complete within {hook_timeout:?} during shutdown, skipping"
859 );
860 }
861 }
862 }
863
864 let _ = fs::remove_dir_all(&*env::IPC_SOCK_DIR);
865 }
866
867 pub(crate) async fn add_notification(&self, level: log::LevelFilter, message: String) {
868 self.pending_notifications
869 .lock()
870 .await
871 .push((level, message));
872 }
873}
874
875#[cfg(unix)]
893fn fix_state_dir_permissions() {
894 let state_dir = &*env::PITCHFORK_STATE_DIR;
895 if let Some((uid, gid)) = state_owner_ids() {
896 if !state_dir.exists()
897 && let Err(err) = fs::create_dir_all(state_dir)
898 {
899 warn!(
900 "failed to create state directory for ownership fix at {}: {err}",
901 state_dir.display()
902 );
903 return;
904 }
905
906 chown_recursive(state_dir, uid, gid, true);
908 debug!(
909 "chowned state directory to uid={uid} gid={gid} at {}",
910 state_dir.display()
911 );
912 } else {
913 if !state_dir.exists() {
914 return;
915 }
916
917 chmod_safe_subtrees(state_dir);
920 debug!(
921 "relaxed permissions on safe subtrees at {}",
922 state_dir.display()
923 );
924 }
925}
926
927#[cfg(unix)]
928pub(crate) fn state_owner_ids() -> Option<(u32, u32)> {
929 if !nix::unistd::Uid::effective().is_root() {
930 return None;
931 }
932
933 let s = settings();
934 let user = s.supervisor.user.trim();
935 if !user.is_empty() {
936 return resolve_supervisor_user_ids(user).or_else(|| {
937 warn!(
938 "failed to resolve supervisor.user '{user}' for state ownership; falling back to SUDO_UID/SUDO_GID"
939 );
940 parse_sudo_ids()
941 });
942 }
943
944 parse_sudo_ids()
945}
946
947#[cfg(unix)]
948fn resolve_supervisor_user_ids(user: &str) -> Option<(u32, u32)> {
949 let user_record = if user.chars().all(|c| c.is_ascii_digit()) {
950 let uid = user.parse::<u32>().ok()?;
951 nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
952 .ok()
953 .flatten()
954 } else {
955 nix::unistd::User::from_name(user).ok().flatten()
956 }?;
957
958 Some((user_record.uid.as_raw(), user_record.gid.as_raw()))
959}
960
961#[cfg(unix)]
967fn parse_sudo_ids() -> Option<(u32, u32)> {
968 if !nix::unistd::Uid::effective().is_root() {
969 return None;
970 }
971 let uid: u32 = std::env::var("SUDO_UID").ok()?.parse().ok()?;
972 let gid: u32 = std::env::var("SUDO_GID").ok()?.parse().ok()?;
973 Some((uid, gid))
974}
975
976#[cfg(unix)]
979fn chown_recursive(dir: &std::path::Path, uid: u32, gid: u32, skip_proxy: bool) {
980 let _ = chown_path(dir, uid, gid);
982
983 let entries = match std::fs::read_dir(dir) {
984 Ok(e) => e,
985 Err(_) => return,
986 };
987 for entry in entries.flatten() {
988 let path = entry.path();
989 if path.is_dir() {
990 if skip_proxy {
992 if let Some(name) = path.file_name().and_then(|n| n.to_str()) {
993 if name == "proxy" {
994 continue;
995 }
996 }
997 }
998 chown_recursive(&path, uid, gid, false);
999 } else {
1000 let _ = chown_path(&path, uid, gid);
1001 }
1002 }
1003}
1004
1005#[cfg(unix)]
1007fn chown_path(path: &std::path::Path, uid: u32, gid: u32) -> std::io::Result<()> {
1008 use std::ffi::CString;
1009 use std::os::unix::ffi::OsStrExt;
1010 let c_path = CString::new(path.as_os_str().as_bytes())
1011 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
1012 let ret = unsafe { libc::chown(c_path.as_ptr(), uid, gid) };
1013 if ret == 0 {
1014 Ok(())
1015 } else {
1016 Err(std::io::Error::last_os_error())
1017 }
1018}
1019
1020#[cfg(unix)]
1023fn chmod_safe_subtrees(state_dir: &std::path::Path) {
1024 let _ = fs::set_permissions(state_dir, fs::Permissions::from_mode(0o755));
1026
1027 let state_file = state_dir.join("state.toml");
1029 if state_file.exists() {
1030 let _ = fs::set_permissions(&state_file, fs::Permissions::from_mode(0o644));
1031 }
1032
1033 for subdir_name in &["sock", "logs"] {
1035 let subdir = state_dir.join(subdir_name);
1036 if subdir.is_dir() {
1037 chmod_recursive(&subdir);
1038 }
1039 }
1040}
1041
1042async fn cleanup_orphaned_daemons(supervisor: &Supervisor) {
1055 if !settings().supervisor.cleanup_orphans {
1056 return;
1057 }
1058
1059 let candidates: Vec<_> = {
1060 let state = supervisor.state_file.lock().await;
1061 state
1062 .daemons
1063 .values()
1064 .filter(|d| d.id != DaemonId::pitchfork() && d.pid.is_some())
1065 .cloned()
1066 .collect()
1067 };
1068
1069 if candidates.is_empty() {
1070 return;
1071 }
1072
1073 info!(
1074 "checking {} daemon(s) for orphaned processes",
1075 candidates.len()
1076 );
1077
1078 for daemon in candidates {
1079 let Some(pid) = daemon.pid else { continue };
1080
1081 if !PROCS.is_running(pid) {
1082 let _ = supervisor
1084 .upsert_daemon(
1085 UpsertDaemonOpts::builder(daemon.id.clone())
1086 .set(|o| {
1087 o.pid = None;
1088 o.status = DaemonStatus::Stopped;
1089 o.active_port = None;
1090 })
1091 .build(),
1092 )
1093 .await;
1094 continue;
1095 }
1096
1097 let current_title = PROCS.title(pid);
1101 let matches = match (¤t_title, &daemon.title) {
1102 (Some(current), Some(expected)) => current == expected,
1103 _ => true,
1107 };
1108
1109 if !matches {
1110 warn!(
1111 "pid {pid} for daemon {} has changed name (expected '{}', found '{}'); skipping orphan cleanup",
1112 daemon.id,
1113 daemon.title.as_deref().unwrap_or("?"),
1114 current_title.as_deref().unwrap_or("?")
1115 );
1116 let _ = supervisor
1118 .upsert_daemon(
1119 UpsertDaemonOpts::builder(daemon.id.clone())
1120 .set(|o| {
1121 o.pid = None;
1122 o.status = DaemonStatus::Stopped;
1123 o.active_port = None;
1124 })
1125 .build(),
1126 )
1127 .await;
1128 continue;
1129 }
1130
1131 info!("terminating orphaned daemon {} (pid {pid})", daemon.id);
1132
1133 let stop_cfg = daemon.stop_signal.unwrap_or_default();
1134 let _ = PROCS
1135 .kill_process_group_async(pid, stop_cfg.signal.into(), stop_cfg.timeout)
1136 .await;
1137
1138 let _ = supervisor
1139 .upsert_daemon(
1140 UpsertDaemonOpts::builder(daemon.id.clone())
1141 .set(|o| {
1142 o.pid = None;
1143 o.status = DaemonStatus::Stopped;
1144 o.active_port = None;
1145 })
1146 .build(),
1147 )
1148 .await;
1149 }
1150}
1151
1152#[cfg(unix)]
1154fn chmod_recursive(dir: &std::path::Path) {
1155 let _ = fs::set_permissions(dir, fs::Permissions::from_mode(0o755));
1156 let entries = match fs::read_dir(dir) {
1157 Ok(e) => e,
1158 Err(_) => return,
1159 };
1160 for entry in entries.flatten() {
1161 let path = entry.path();
1162 if path.is_dir() {
1163 chmod_recursive(&path);
1164 } else {
1165 let _ = fs::set_permissions(&path, fs::Permissions::from_mode(0o644));
1166 }
1167 }
1168}