pitchfork_cli/supervisor/
mod.rs1mod 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::new(env::PITCHFORK_STATE_FILE.clone())),
142 last_refreshed_at: Mutex::new(time::Instant::now()),
143 pending_notifications: Mutex::new(vec![]),
144 pending_autostops: Mutex::new(HashMap::new()),
145 ipc_shutdown: Mutex::new(None),
146 hook_tasks: Mutex::new(Vec::new()),
147 active_monitors: AtomicU32::new(0),
148 monitor_done: Notify::new(),
149 proxy_cancel: Mutex::new(None),
150 proxy_task: Mutex::new(None),
151 mdns_publisher: Mutex::new(None),
152 lan_monitor_task: Mutex::new(None),
153 flush_cancel: std::sync::Mutex::new(None),
154 })
155 }
156
157 pub async fn start(
158 &self,
159 is_boot: bool,
160 container: bool,
161 web_port: Option<u16>,
162 web_path: Option<String>,
163 ) -> Result<()> {
164 #[cfg(unix)]
169 fix_state_dir_permissions();
170
171 let pid = std::process::id();
172 PROCS.refresh_pids(&[pid]);
174 let container_mode = container || settings().supervisor.container;
176 if container_mode {
177 info!("Starting supervisor in container/PID1 mode with pid {pid}");
178 } else {
179 info!("Starting supervisor with pid {pid}");
180 }
181
182 cleanup_orphaned_daemons(self).await;
187
188 self.upsert_daemon(
189 UpsertDaemonOpts::builder(DaemonId::pitchfork())
190 .set(|o| {
191 o.pid = Some(pid);
192 o.status = DaemonStatus::Running;
193 })
194 .build(),
195 )
196 .await?;
197 #[cfg(unix)]
198 fix_state_dir_permissions();
199
200 if is_boot {
202 info!("Boot start mode enabled, starting boot_start daemons");
203 self.start_boot_daemons().await?;
204 }
205
206 self.interval_watch()?;
207 self.cron_watch()?;
208 self.signals()?;
209 self.daemon_file_watch()?;
210
211 #[cfg(unix)]
213 if container_mode {
214 self.reap_zombies()?;
215 }
216
217 let s = settings();
219 let effective_port = web_port.or_else(|| {
220 if s.web.auto_start {
221 match u16::try_from(s.web.bind_port).ok().filter(|&p| p > 0) {
222 Some(p) => Some(p),
223 None => {
224 error!(
225 "web.bind_port {} is out of valid port range (1-65535), web UI disabled",
226 s.web.bind_port
227 );
228 None
229 }
230 }
231 } else {
232 None
233 }
234 });
235 let effective_path = web_path.or_else(|| {
237 let bp = s.web.base_path.clone();
238 if bp.is_empty() { None } else { Some(bp) }
239 });
240 if let Some(port) = effective_port {
241 tokio::spawn(async move {
242 if let Err(e) = crate::web::serve(port, effective_path).await {
243 error!("Web server error: {e}");
244 }
245 });
246 }
247
248 let api_port = if s.api.auto_start {
250 match u16::try_from(s.api.bind_port).ok().filter(|&p| p > 0) {
251 Some(p) => Some(p),
252 None => {
253 error!(
254 "api.bind_port {} is out of valid port range (1-65535), API server disabled",
255 s.api.bind_port
256 );
257 None
258 }
259 }
260 } else {
261 None
262 };
263 if let Some(port) = api_port {
264 tokio::spawn(async move {
265 if let Err(e) = crate::web::serve_api(port, None).await {
266 error!("API server error: {e}");
267 }
268 });
269 }
270
271 if s.proxy.enable {
273 #[cfg(feature = "proxy-tls")]
277 if s.proxy.https {
278 let proxy_dir = crate::env::PITCHFORK_STATE_DIR.join("proxy");
279 let ca_cert_path = proxy_dir.join("ca.pem");
280 let ca_key_path = proxy_dir.join("ca-key.pem");
281 if !ca_cert_path.exists() || !ca_key_path.exists() {
282 match crate::proxy::server::generate_ca(&ca_cert_path, &ca_key_path) {
283 Ok(()) => {
284 info!(
285 "Generated local CA certificate at {}",
286 ca_cert_path.display()
287 );
288 }
289 Err(e) => {
290 error!("Failed to generate CA certificate: {e}");
291 }
292 }
293 }
294
295 if s.proxy.auto_trust && ca_cert_path.exists() {
299 use crate::proxy::trust::{AutoTrustResult, auto_trust};
300 match auto_trust(&ca_cert_path) {
301 AutoTrustResult::AlreadyTrusted => {}
302 AutoTrustResult::Trusted => {
303 info!("CA certificate auto-trusted in system store");
304 }
305 AutoTrustResult::NotTrusted { reason } => {
306 warn!("Auto-trust skipped: {reason}");
307 warn!("Run `pitchfork proxy trust` to install manually");
308 }
309 }
310 }
311 }
312 let (bind_tx, bind_rx) = tokio::sync::oneshot::channel();
316 let proxy_cancel = tokio_util::sync::CancellationToken::new();
317 let proxy_cancel_clone = proxy_cancel.clone();
318 *self.proxy_cancel.lock().await = Some(proxy_cancel);
319 let proxy_task = tokio::spawn(async move {
320 if let Err(e) = crate::proxy::server::serve(bind_tx, proxy_cancel_clone).await {
321 error!("Proxy server error: {e}");
322 }
323 });
324 *self.proxy_task.lock().await = Some(proxy_task);
325 match bind_rx.await {
326 Ok(Ok(())) => {
327 info!("Proxy server bound successfully");
328 self.start_mdns().await;
329 }
330 Ok(Err(msg)) => {
331 error!("{msg}");
332 self.add_notification(log::LevelFilter::Error, msg).await;
333 }
334 Err(_) => {
335 }
339 }
340 }
341
342 tokio::spawn(async {
345 crate::proxy::server::get_cached_slugs().await;
346 });
347
348 let (ipc, ipc_handle) = IpcServer::new()?;
349 *self.ipc_shutdown.lock().await = Some(ipc_handle);
350 self.start_state_flush_task();
351 self.conn_watch(ipc).await
352 }
353
354 async fn start_mdns(&self) {
356 let s = crate::settings::settings();
357 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
358 if !s.proxy.enable || !lan_enabled {
359 return;
360 }
361
362 let lan_ip = if !s.proxy.lan_ip.is_empty() {
363 match s.proxy.lan_ip.parse::<std::net::Ipv4Addr>() {
364 Ok(ip) => Some(ip),
365 Err(e) => {
366 error!(
367 "proxy.lan_ip {:?} is not a valid IPv4 address: {e}",
368 s.proxy.lan_ip
369 );
370 return;
371 }
372 }
373 } else {
374 match crate::proxy::lan_ip::detect_lan_ip().await {
375 Some(ip) => Some(ip),
376 None => {
377 error!(
378 "LAN mode is enabled but no LAN IP address could be detected. \
379 Set proxy.lan_ip to a specific address, or ensure you are connected to a network."
380 );
381 return;
382 }
383 }
384 };
385
386 let Some(lan_ip) = lan_ip else { return };
387 let port = u16::try_from(s.proxy.port).unwrap_or(443);
388
389 let Some(mut publisher) = crate::proxy::mdns::MdnsPublisher::new(lan_ip) else {
390 error!("Failed to start mDNS publisher. Is Avahi (Linux) or Bonjour (macOS) running?");
391 return;
392 };
393
394 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
396 for slug in slugs.keys() {
397 let hostname = format!("{slug}.local");
398 publisher.publish(&hostname, port);
399 }
400
401 log::info!(
402 "LAN mode: mDNS publishing on {lan_ip}, {} slug(s) registered",
403 slugs.len()
404 );
405
406 let publisher = std::sync::Arc::new(tokio::sync::Mutex::new(publisher));
407
408 let ip_pinned = !s.proxy.lan_ip.is_empty();
410 if !ip_pinned {
411 let monitor_cancel = self.proxy_cancel.lock().await.clone();
412 let publisher_clone = publisher.clone();
413 let task = tokio::spawn(async move {
414 let mut last_ip = lan_ip;
415 let interval = std::time::Duration::from_secs(5);
416 let mut ticker = tokio::time::interval(interval);
417 ticker.tick().await; loop {
419 ticker.tick().await;
420 if let Some(cancel) = monitor_cancel.as_ref() {
421 if cancel.is_cancelled() {
422 break;
423 }
424 }
425 if let Some(new_ip) =
426 crate::proxy::lan_ip::detect_lan_ip_if_changed(last_ip).await
427 {
428 log::info!("LAN IP changed: {last_ip} → {new_ip}");
429 last_ip = new_ip;
430 let mut pub_guard = publisher_clone.lock().await;
431 pub_guard.republish_all(new_ip, port);
432 }
433 }
434 });
435 *self.lan_monitor_task.lock().await = Some(task);
436 }
437
438 *self.mdns_publisher.lock().await = Some(publisher);
439 }
440
441 async fn sync_mdns(&self) {
446 let publisher = {
449 let guard = self.mdns_publisher.lock().await;
450 match guard.as_ref() {
451 Some(p) => p.clone(),
452 None => {
453 debug!("sync_mdns: mDNS publisher not active, skipping");
454 return;
455 }
456 }
457 };
458
459 let s = crate::settings::settings();
460 let port = u16::try_from(s.proxy.port).unwrap_or(443);
461
462 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
463 let mut pub_guard = publisher.lock().await;
464
465 let current_keys: Vec<&String> = slugs.keys().collect();
467 let registered: Vec<String> = pub_guard.registered_hostnames();
468 for hostname in ®istered {
469 let slug = hostname.strip_suffix(".local").unwrap_or(hostname);
471 if !current_keys.iter().any(|k| k.as_str() == slug) {
472 log::info!("mDNS: unpublishing removed slug {slug}");
473 pub_guard.unpublish(hostname);
474 }
475 }
476
477 for slug in slugs.keys() {
479 let hostname = format!("{slug}.local");
480 if !pub_guard.is_published(&hostname) {
481 log::info!("mDNS: publishing new slug {slug}");
482 pub_guard.publish(&hostname, port);
483 }
484 }
485 }
486
487 fn start_state_flush_task(&self) {
491 let cancel = tokio_util::sync::CancellationToken::new();
492 *self.flush_cancel.lock().unwrap() = Some(cancel.clone());
493 tokio::spawn(async move {
494 let mut interval = time::interval(Duration::from_secs(1));
495 interval.set_missed_tick_behavior(time::MissedTickBehavior::Skip);
496 loop {
497 tokio::select! {
498 _ = interval.tick() => {}
499 _ = cancel.cancelled() => {
500 debug!("state flush task received shutdown signal");
501 break;
502 }
503 }
504 let state = SUPERVISOR.state_file.lock().await;
505 if state.is_dirty() {
506 if let Err(e) = state.write() {
507 warn!("failed to flush state file: {e}");
508 }
509 }
510 }
511 debug!("state flush task exiting");
512 });
513 }
514
515 pub(crate) async fn flush_state(&self) {
516 let state = self.state_file.lock().await;
517 if state.is_dirty() {
518 if let Err(e) = state.write() {
519 warn!("failed to flush state file: {e}");
520 }
521 }
522 }
523
524 pub(crate) async fn refresh(&self) -> Result<()> {
525 trace!("refreshing");
526
527 let dirs_with_pids = self.get_dirs_with_shell_pids().await;
528
529 let mut last_refreshed_at = self.last_refreshed_at.lock().await;
530 *last_refreshed_at = time::Instant::now();
531
532 for (dir, pids) in dirs_with_pids {
533 let to_remove = pids
534 .iter()
535 .filter(|pid| !PROCS.is_running(**pid))
536 .collect::<Vec<_>>();
537 for pid in &to_remove {
538 self.remove_shell_pid(**pid).await?
539 }
540 if to_remove.len() == pids.len() {
541 self.leave_dir(&dir).await?;
542 }
543 }
544
545 self.check_retry().await?;
546 self.process_pending_autostops().await?;
547
548 Ok(())
549 }
550
551 #[cfg(unix)]
576 fn reap_zombies(&self) -> Result<()> {
577 let mut stream = signal::unix::signal(SignalKind::child())
578 .map_err(|e| miette::miette!("Failed to register SIGCHLD handler: {e}"))?;
579 tokio::spawn(async move {
580 loop {
581 stream.recv().await;
582 let managed_pids: HashSet<u32> = SUPERVISOR
584 .state_file
585 .lock()
586 .await
587 .daemons
588 .values()
589 .filter_map(|d| d.pid)
590 .collect();
591 Self::reap_unmanaged_zombies(&managed_pids).await;
593 }
594 });
595 info!("container mode: SIGCHLD zombie reaper installed");
596 Ok(())
597 }
598
599 #[cfg(target_os = "linux")]
605 async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
606 use nix::sys::wait::{Id, WaitPidFlag, WaitStatus, waitid, waitpid};
607 use nix::unistd::Pid;
608
609 loop {
610 let peek_flags = WaitPidFlag::WNOHANG | WaitPidFlag::WNOWAIT | WaitPidFlag::WEXITED;
612 match waitid(Id::All, peek_flags) {
613 Ok(WaitStatus::StillAlive) => break,
614 Ok(status) => {
615 let Some(pid_raw) = status.pid().map(|p| p.as_raw() as u32) else {
616 break;
617 };
618 if managed_pids.contains(&pid_raw) {
619 trace!(
623 "zombie reaper: skipping managed daemon pid {pid_raw}, \
624 leaving for Tokio to reap"
625 );
626 break;
627 }
628 match waitpid(Pid::from_raw(pid_raw as i32), Some(WaitPidFlag::WNOHANG)) {
630 Ok(s) => trace!("reaped orphaned zombie child: {s:?}"),
631 Err(nix::errno::Errno::ECHILD) => break,
632 Err(e) => {
633 trace!("waitpid error reaping pid {pid_raw}: {e}");
634 break;
635 }
636 }
637 }
638 Err(nix::errno::Errno::ECHILD) => break, Err(e) => {
640 trace!("waitid error in zombie reaper: {e}");
641 break;
642 }
643 }
644 }
645 }
646
647 #[cfg(all(unix, not(target_os = "linux")))]
653 async fn reap_unmanaged_zombies(managed_pids: &HashSet<u32>) {
654 use nix::sys::wait::{WaitPidFlag, WaitStatus, waitpid};
655
656 loop {
657 match waitpid(None, Some(WaitPidFlag::WNOHANG)) {
658 Ok(WaitStatus::StillAlive) => break,
659 Ok(status) => {
660 let Some(pid) = status.pid().map(|p| p.as_raw() as u32) else {
661 continue;
662 };
663 if managed_pids.contains(&pid) {
664 let exit_code = match status {
666 WaitStatus::Exited(_, code) => code,
667 WaitStatus::Signaled(_, sig, _) => -(sig as i32),
668 _ => -1,
669 };
670 warn!(
671 "zombie reaper reaped managed daemon pid {pid} \
672 (exit_code={exit_code}); stashing status for recovery"
673 );
674 REAPED_STATUSES.lock().await.insert(pid, exit_code);
675 } else {
676 trace!("reaped orphaned zombie child: {status:?}");
677 }
678 }
679 Err(nix::errno::Errno::ECHILD) => break, Err(e) => {
681 trace!("waitpid error in zombie reaper: {e}");
682 break;
683 }
684 }
685 }
686 }
687
688 #[cfg(unix)]
689 fn signals(&self) -> Result<()> {
690 let signals = [
691 SignalKind::terminate(),
692 SignalKind::alarm(),
693 SignalKind::interrupt(),
694 SignalKind::quit(),
695 SignalKind::hangup(),
696 SignalKind::user_defined1(),
697 SignalKind::user_defined2(),
698 ];
699 static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
700 for signal in signals {
701 let stream = match signal::unix::signal(signal) {
702 Ok(s) => s,
703 Err(e) => {
704 warn!("Failed to register signal handler for {signal:?}: {e}");
705 continue;
706 }
707 };
708 tokio::spawn(async move {
709 let mut stream = stream;
710 loop {
711 stream.recv().await;
712 if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
713 exit(1);
714 } else {
715 SUPERVISOR.handle_signal().await;
716 }
717 }
718 });
719 }
720 Ok(())
721 }
722
723 #[cfg(windows)]
724 fn signals(&self) -> Result<()> {
725 tokio::spawn(async move {
726 static RECEIVED_SIGNAL: AtomicBool = AtomicBool::new(false);
727 loop {
728 if let Err(e) = signal::ctrl_c().await {
729 error!("Failed to wait for ctrl-c: {}", e);
730 return;
731 }
732 if RECEIVED_SIGNAL.swap(true, atomic::Ordering::SeqCst) {
733 exit(1);
734 } else {
735 SUPERVISOR.handle_signal().await;
736 }
737 }
738 });
739 Ok(())
740 }
741
742 async fn handle_signal(&self) {
743 info!("received signal, stopping");
744 self.close().await;
745 exit(0)
746 }
747
748 pub(crate) async fn close(&self) {
749 if let Some(cancel) = self.proxy_cancel.lock().await.take() {
753 cancel.cancel();
754 }
755
756 if let Some(monitor_task) = self.lan_monitor_task.lock().await.take() {
758 monitor_task.abort();
759 }
760
761 if let Some(publisher) = self.mdns_publisher.lock().await.take() {
763 publisher.lock().await.shutdown();
764 }
765
766 if let Some(proxy_task) = self.proxy_task.lock().await.take() {
767 let _ = tokio::time::timeout(Duration::from_secs(12), proxy_task).await;
768 }
769
770 let s = settings();
772 if s.proxy.enable && s.proxy.sync_hosts {
773 crate::proxy::hosts::clean_hosts_file();
774 }
775
776 let pitchfork_id = DaemonId::pitchfork();
777 let active = self.active_daemons().await;
778 let active_ids: Vec<DaemonId> = active
779 .iter()
780 .filter(|d| d.id != pitchfork_id)
781 .map(|d| d.id.clone())
782 .collect();
783
784 let stop_levels = compute_reverse_stop_order(&active_ids);
789 for level in &stop_levels {
790 let mut tasks = Vec::new();
791 for id in level {
792 let id = id.clone();
793 tasks.push(tokio::spawn(async move {
794 if let Err(err) = SUPERVISOR.stop(&id).await {
795 error!("failed to stop daemon {id}: {err}");
796 }
797 }));
798 }
799 for task in tasks {
800 let _ = task.await;
801 }
802 }
803 let _ = self.remove_daemon(&pitchfork_id).await;
804
805 if let Some(cancel) = self.flush_cancel.lock().unwrap().take() {
808 cancel.cancel();
809 }
810
811 {
814 let state = self.state_file.lock().await;
815 if state.is_dirty() {
816 if let Err(e) = state.write() {
817 warn!("failed to flush state file during shutdown: {e}");
818 }
819 }
820 }
821
822 if let Some(mut handle) = self.ipc_shutdown.lock().await.take() {
824 handle.shutdown();
825 }
826
827 let drain_timeout = time::sleep(Duration::from_secs(5));
833 tokio::pin!(drain_timeout);
834 loop {
835 if self.active_monitors.load(atomic::Ordering::Acquire) == 0 {
836 break;
837 }
838 tokio::select! {
839 _ = self.monitor_done.notified() => {}
840 _ = &mut drain_timeout => {
841 warn!("timed out waiting for monitoring tasks to register hooks, proceeding with shutdown");
842 break;
843 }
844 }
845 }
846 let handles: Vec<JoinHandle<()>> = std::mem::take(&mut *self.hook_tasks.lock().await);
847 let hook_timeout = Duration::from_secs(30);
848 for handle in handles {
849 match time::timeout(hook_timeout, handle).await {
850 Ok(_) => {} Err(_) => {
852 warn!(
853 "hook task did not complete within {hook_timeout:?} during shutdown, skipping"
854 );
855 }
856 }
857 }
858
859 let _ = fs::remove_dir_all(&*env::IPC_SOCK_DIR);
860 }
861
862 pub(crate) async fn add_notification(&self, level: log::LevelFilter, message: String) {
863 self.pending_notifications
864 .lock()
865 .await
866 .push((level, message));
867 }
868}
869
870#[cfg(unix)]
888fn fix_state_dir_permissions() {
889 let state_dir = &*env::PITCHFORK_STATE_DIR;
890 if let Some((uid, gid)) = state_owner_ids() {
891 if !state_dir.exists()
892 && let Err(err) = fs::create_dir_all(state_dir)
893 {
894 warn!(
895 "failed to create state directory for ownership fix at {}: {err}",
896 state_dir.display()
897 );
898 return;
899 }
900
901 chown_recursive(state_dir, uid, gid, true);
903 debug!(
904 "chowned state directory to uid={uid} gid={gid} at {}",
905 state_dir.display()
906 );
907 } else {
908 if !state_dir.exists() {
909 return;
910 }
911
912 chmod_safe_subtrees(state_dir);
915 debug!(
916 "relaxed permissions on safe subtrees at {}",
917 state_dir.display()
918 );
919 }
920}
921
922#[cfg(unix)]
923pub(crate) fn state_owner_ids() -> Option<(u32, u32)> {
924 if !nix::unistd::Uid::effective().is_root() {
925 return None;
926 }
927
928 let s = settings();
929 let user = s.supervisor.user.trim();
930 if !user.is_empty() {
931 return resolve_supervisor_user_ids(user).or_else(|| {
932 warn!(
933 "failed to resolve supervisor.user '{user}' for state ownership; falling back to SUDO_UID/SUDO_GID"
934 );
935 parse_sudo_ids()
936 });
937 }
938
939 parse_sudo_ids()
940}
941
942#[cfg(unix)]
943fn resolve_supervisor_user_ids(user: &str) -> Option<(u32, u32)> {
944 let user_record = if user.chars().all(|c| c.is_ascii_digit()) {
945 let uid = user.parse::<u32>().ok()?;
946 nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
947 .ok()
948 .flatten()
949 } else {
950 nix::unistd::User::from_name(user).ok().flatten()
951 }?;
952
953 Some((user_record.uid.as_raw(), user_record.gid.as_raw()))
954}
955
956#[cfg(unix)]
962fn parse_sudo_ids() -> Option<(u32, u32)> {
963 if !nix::unistd::Uid::effective().is_root() {
964 return None;
965 }
966 let uid: u32 = std::env::var("SUDO_UID").ok()?.parse().ok()?;
967 let gid: u32 = std::env::var("SUDO_GID").ok()?.parse().ok()?;
968 Some((uid, gid))
969}
970
971#[cfg(unix)]
974fn chown_recursive(dir: &std::path::Path, uid: u32, gid: u32, skip_proxy: bool) {
975 let _ = chown_path(dir, uid, gid);
977
978 let entries = match std::fs::read_dir(dir) {
979 Ok(e) => e,
980 Err(_) => return,
981 };
982 for entry in entries.flatten() {
983 let path = entry.path();
984 if path.is_dir() {
985 if skip_proxy {
987 if let Some(name) = path.file_name().and_then(|n| n.to_str()) {
988 if name == "proxy" {
989 continue;
990 }
991 }
992 }
993 chown_recursive(&path, uid, gid, false);
994 } else {
995 let _ = chown_path(&path, uid, gid);
996 }
997 }
998}
999
1000#[cfg(unix)]
1002fn chown_path(path: &std::path::Path, uid: u32, gid: u32) -> std::io::Result<()> {
1003 use std::ffi::CString;
1004 use std::os::unix::ffi::OsStrExt;
1005 let c_path = CString::new(path.as_os_str().as_bytes())
1006 .map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidInput, e))?;
1007 let ret = unsafe { libc::chown(c_path.as_ptr(), uid, gid) };
1008 if ret == 0 {
1009 Ok(())
1010 } else {
1011 Err(std::io::Error::last_os_error())
1012 }
1013}
1014
1015#[cfg(unix)]
1018fn chmod_safe_subtrees(state_dir: &std::path::Path) {
1019 let _ = fs::set_permissions(state_dir, fs::Permissions::from_mode(0o755));
1021
1022 let state_file = state_dir.join("state.toml");
1024 if state_file.exists() {
1025 let _ = fs::set_permissions(&state_file, fs::Permissions::from_mode(0o644));
1026 }
1027
1028 for subdir_name in &["sock", "logs"] {
1030 let subdir = state_dir.join(subdir_name);
1031 if subdir.is_dir() {
1032 chmod_recursive(&subdir);
1033 }
1034 }
1035}
1036
1037async fn cleanup_orphaned_daemons(supervisor: &Supervisor) {
1050 if !settings().supervisor.cleanup_orphans {
1051 return;
1052 }
1053
1054 let candidates: Vec<_> = {
1055 let state = supervisor.state_file.lock().await;
1056 state
1057 .daemons
1058 .values()
1059 .filter(|d| d.id != DaemonId::pitchfork() && d.pid.is_some())
1060 .cloned()
1061 .collect()
1062 };
1063
1064 if candidates.is_empty() {
1065 return;
1066 }
1067
1068 info!(
1069 "checking {} daemon(s) for orphaned processes",
1070 candidates.len()
1071 );
1072
1073 for daemon in candidates {
1074 let Some(pid) = daemon.pid else { continue };
1075
1076 if !PROCS.is_running(pid) {
1077 let _ = supervisor
1079 .upsert_daemon(
1080 UpsertDaemonOpts::builder(daemon.id.clone())
1081 .set(|o| {
1082 o.pid = None;
1083 o.status = DaemonStatus::Stopped;
1084 o.active_port = None;
1085 })
1086 .build(),
1087 )
1088 .await;
1089 continue;
1090 }
1091
1092 let current_title = PROCS.title(pid);
1096 let matches = match (¤t_title, &daemon.title) {
1097 (Some(current), Some(expected)) => current == expected,
1098 _ => true,
1102 };
1103
1104 if !matches {
1105 warn!(
1106 "pid {pid} for daemon {} has changed name (expected '{}', found '{}'); skipping orphan cleanup",
1107 daemon.id,
1108 daemon.title.as_deref().unwrap_or("?"),
1109 current_title.as_deref().unwrap_or("?")
1110 );
1111 let _ = supervisor
1113 .upsert_daemon(
1114 UpsertDaemonOpts::builder(daemon.id.clone())
1115 .set(|o| {
1116 o.pid = None;
1117 o.status = DaemonStatus::Stopped;
1118 o.active_port = None;
1119 })
1120 .build(),
1121 )
1122 .await;
1123 continue;
1124 }
1125
1126 info!("terminating orphaned daemon {} (pid {pid})", daemon.id);
1127
1128 let stop_cfg = daemon.stop_signal.unwrap_or_default();
1129 let _ = PROCS
1130 .kill_process_group_async(pid, stop_cfg.signal.into(), stop_cfg.timeout)
1131 .await;
1132
1133 let _ = supervisor
1134 .upsert_daemon(
1135 UpsertDaemonOpts::builder(daemon.id.clone())
1136 .set(|o| {
1137 o.pid = None;
1138 o.status = DaemonStatus::Stopped;
1139 o.active_port = None;
1140 })
1141 .build(),
1142 )
1143 .await;
1144 }
1145}
1146
1147#[cfg(unix)]
1149fn chmod_recursive(dir: &std::path::Path) {
1150 let _ = fs::set_permissions(dir, fs::Permissions::from_mode(0o755));
1151 let entries = match fs::read_dir(dir) {
1152 Ok(e) => e,
1153 Err(_) => return,
1154 };
1155 for entry in entries.flatten() {
1156 let path = entry.path();
1157 if path.is_dir() {
1158 chmod_recursive(&path);
1159 } else {
1160 let _ = fs::set_permissions(&path, fs::Permissions::from_mode(0o644));
1161 }
1162 }
1163}