1use super::hooks::{self, HookType, fire_hook};
6use super::{SUPERVISOR, Supervisor};
7use crate::daemon::RunOptions;
8use crate::daemon_id::DaemonId;
9use crate::daemon_status::DaemonStatus;
10use crate::error::PortError;
11use crate::ipc::IpcResponse;
12use crate::log_store::LogStore;
13use crate::log_store::sqlite::LOG_STORE;
14use crate::pitchfork_toml::{ReadyCmd, ReadyHttp, ReadyOutput, ReadyPort};
15use crate::procs::PROCS;
16use crate::settings::settings;
17use crate::shell::{HideConsoleWindow, Shell};
18use crate::supervisor::state::UpsertDaemonOpts;
19use crate::{Result, env};
20use indexmap::IndexMap;
21use miette::IntoDiagnostic;
22use once_cell::sync::Lazy;
23use regex::Regex;
24use std::collections::HashMap;
25#[cfg(unix)]
26use std::ffi::CString;
27use std::sync::{Arc, atomic};
28use std::time::Duration;
29use tokio::io::AsyncBufReadExt;
30use tokio::select;
31use tokio::sync::oneshot;
32use tokio::time;
33
34static REGEX_CACHE: Lazy<std::sync::Mutex<HashMap<String, Regex>>> =
36 Lazy::new(|| std::sync::Mutex::new(HashMap::new()));
37
38fn resolve_configured_ready_port(
39 configured_port: u16,
40 expected_ports: &[u16],
41 resolved_ports: &[u16],
42) -> u16 {
43 let bump_offset = resolved_ports
44 .first()
45 .unwrap_or(&0)
46 .saturating_sub(*expected_ports.first().unwrap_or(&0));
47 if expected_ports.contains(&configured_port) && bump_offset > 0 {
48 configured_port
49 .checked_add(bump_offset)
50 .unwrap_or(configured_port)
51 } else {
52 configured_port
53 }
54}
55
56fn active_port_from_ready_port(ready_port: u16, resolved_ports: &[u16]) -> Option<u16> {
57 resolved_ports
58 .first()
59 .copied()
60 .filter(|&primary_port| primary_port == ready_port)
61}
62
63#[cfg(unix)]
64#[derive(Clone, Debug, PartialEq, Eq)]
65enum RunIdentity {
66 Inherit,
67 Switch {
68 uid: nix::unistd::Uid,
69 gid: nix::unistd::Gid,
70 username: Option<CString>,
71 },
72}
73
74pub(crate) fn get_or_compile_regex(pattern: &str) -> Option<Regex> {
76 let mut cache = REGEX_CACHE.lock().unwrap_or_else(|e| e.into_inner());
77 if let Some(re) = cache.get(pattern) {
78 return Some(re.clone());
79 }
80 match Regex::new(pattern) {
81 Ok(re) => {
82 cache.insert(pattern.to_string(), re.clone());
83 Some(re)
84 }
85 Err(e) => {
86 error!("invalid regex pattern '{pattern}': {e}");
87 None
88 }
89 }
90}
91
92pub(crate) struct CmdProbe {
99 pub(crate) cancel_tx: tokio::sync::oneshot::Sender<()>,
100 pub(crate) result_rx: tokio::sync::oneshot::Receiver<std::io::Result<std::process::ExitStatus>>,
101}
102
103fn apply_runtime_env(
110 command: &mut tokio::process::Command,
111 id: &DaemonId,
112 retry_count: u32,
113 daemon_env: Option<&IndexMap<String, String>>,
114 resolved_ports: &[u16],
115) {
116 if let Some(ref path) = *env::ORIGINAL_PATH {
117 command.env("PATH", path);
118 }
119 if let Some(env_vars) = daemon_env {
120 command.envs(env_vars);
121 }
122 command
123 .env("PITCHFORK_DAEMON_ID", id.qualified())
124 .env("PITCHFORK_DAEMON_NAMESPACE", id.namespace())
125 .env("PITCHFORK_RETRY_COUNT", retry_count.to_string());
126 if let Some(port) = resolved_ports.first() {
127 command.env("PORT", port.to_string());
128 for (index, port) in resolved_ports.iter().enumerate() {
129 command.env(format!("PORT{index}"), port.to_string());
130 }
131 }
132}
133
134pub(crate) fn spawn_cmd_probe(
135 id: &DaemonId,
136 cmd: &str,
137 dir: &std::path::Path,
138 retry_count: u32,
139 daemon_env: Option<&IndexMap<String, String>>,
140 resolved_ports: &[u16],
141) -> CmdProbe {
142 let shell_setting = settings().general.shell.clone();
148 let mut command = match shell_words::split(&shell_setting) {
149 Ok(parts) if !parts.is_empty() => {
150 let (program, args) = parts.split_first().unwrap();
151 let mut c = tokio::process::Command::new(program);
152 c.args(args);
153 c.arg(cmd);
154 c
155 }
156 _ => Shell::default_for_platform().command(cmd),
157 };
158 command
159 .current_dir(dir)
160 .stdout(std::process::Stdio::null())
161 .stderr(std::process::Stdio::null())
162 .kill_on_drop(true)
163 .hide_console_window();
164 apply_runtime_env(&mut command, id, retry_count, daemon_env, resolved_ports);
165 let mut child = match command.spawn() {
166 Ok(child) => child,
167 Err(e) => {
168 warn!("daemon {id}: failed to spawn command probe: {e}");
169 let (cancel_tx, _) = tokio::sync::oneshot::channel();
173 let (_, result_rx) = tokio::sync::oneshot::channel();
174 return CmdProbe {
175 cancel_tx,
176 result_rx,
177 };
178 }
179 };
180
181 let (cancel_tx, mut cancel_rx) = tokio::sync::oneshot::channel();
182 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
183
184 tokio::spawn(async move {
185 let status = tokio::select! {
186 status = child.wait() => status,
187 _ = &mut cancel_rx => {
188 let mut child = child;
189 let _ = child.kill().await;
190 child.wait().await
191 }
192 };
193 let _ = result_tx.send(status);
194 });
195
196 CmdProbe {
197 cancel_tx,
198 result_rx,
199 }
200}
201
202fn stop_cmd_probe_state(probe: &mut Option<CmdProbe>) {
204 if let Some(p) = probe.take() {
205 let _ = p.cancel_tx.send(());
206 }
207}
208
209fn spawn_ready_fail_kill(
214 id: DaemonId,
215 pid: u32,
216 stop_cfg: crate::config_types::StopConfig,
217) -> tokio::task::JoinHandle<()> {
218 tokio::spawn(async move {
219 if let Err(e) = PROCS
220 .kill_process_group_async(pid, stop_cfg.signal.into(), stop_cfg.timeout)
221 .await
222 {
223 error!("daemon {id}: failed to kill pid {pid} after readiness failure: {e}");
224 }
225 })
226}
227
228#[allow(clippy::too_many_arguments)]
233fn any_ready_check_remaining(
234 ready_output: Option<&ReadyOutput>,
235 output_exhausted: bool,
236 ready_port: Option<&ReadyPort>,
237 port_exhausted: bool,
238 ready_http: Option<&ReadyHttp>,
239 http_exhausted: bool,
240 ready_cmd: Option<&ReadyCmd>,
241 cmd_exhausted: bool,
242) -> bool {
243 ready_output.is_some_and(|o| o.timeout.is_none() || !output_exhausted)
244 || ready_port.is_some_and(|p| p.timeout.is_none() || !port_exhausted)
245 || ready_http.is_some_and(|h| h.timeout.is_none() || !http_exhausted)
246 || ready_cmd.is_some_and(|c| c.timeout.is_none() || !cmd_exhausted)
247}
248
249fn delay_readiness_succeeded(
250 ready_notified: bool,
251 has_other_ready_check: bool,
252 process_exited: bool,
253 process_running: bool,
254) -> bool {
255 !ready_notified && !has_other_ready_check && !process_exited && process_running
256}
257
258const SINK_OUTPUT_TIMEOUT: Duration = Duration::from_millis(400);
263
264impl Supervisor {
265 pub async fn run(&self, opts: RunOptions) -> Result<IpcResponse> {
267 let id = &opts.id;
268 let cmd = opts.cmd.clone();
269
270 {
272 let mut pending = self.pending_autostops.lock().await;
273 if pending.remove(id).is_some() {
274 info!("cleared pending autostop for {id} (daemon starting)");
275 }
276 }
277
278 let mut stop_guard = Some(self.stop_lock(id).await.lock_owned().await);
289 let daemon = self.get_daemon(id).await;
290 if let Some(daemon) = daemon {
291 if !daemon.status.is_stopping()
294 && !daemon.status.is_stopped()
295 && let Some(pid) = daemon.pid
296 {
297 if opts.force {
298 self.stop_locked(id).await?;
299 info!("run: stop completed for daemon {id}");
300 } else {
301 warn!("daemon {id} already running with pid {pid}");
302 return Ok(IpcResponse::DaemonAlreadyRunning);
303 }
304 }
305 }
306
307 if opts.wait_ready && opts.retry.count() > 0 {
309 let max_attempts = opts.retry.count().saturating_add(1);
311 for attempt in 0..max_attempts {
312 let mut retry_opts = opts.clone();
313 retry_opts.retry_count = attempt;
314 retry_opts.cmd = cmd.clone();
315
316 let guard = match stop_guard.take() {
320 Some(guard) => guard,
321 None => self.stop_lock(id).await.lock_owned().await,
322 };
323 let result = self.run_once(retry_opts, guard).await?;
324
325 match result {
326 IpcResponse::DaemonReady { daemon } => {
327 return Ok(IpcResponse::DaemonReady { daemon });
328 }
329 IpcResponse::DaemonFailedWithCode {
330 exit_code,
331 resolved_ports,
332 } => {
333 if attempt < opts.retry.count() {
334 let backoff_secs = 2u64.saturating_pow(attempt).min(3600);
335 info!(
336 "daemon {id} failed (attempt {}/{}), retrying in {}s",
337 attempt + 1,
338 max_attempts,
339 backoff_secs
340 );
341 fire_hook(
342 HookType::OnRetry,
343 id.clone(),
344 opts.dir.0.clone(),
345 attempt + 1,
346 opts.env.clone(),
347 resolved_ports,
348 vec![],
349 )
350 .await;
351 time::sleep(Duration::from_secs(backoff_secs)).await;
352 continue;
353 } else {
354 info!("daemon {id} failed after {max_attempts} attempts");
355 return Ok(IpcResponse::DaemonFailedWithCode {
356 exit_code,
357 resolved_ports,
358 });
359 }
360 }
361 other => return Ok(other),
362 }
363 }
364 }
365
366 let guard = match stop_guard.take() {
368 Some(guard) => guard,
369 None => self.stop_lock(id).await.lock_owned().await,
370 };
371 self.run_once(opts, guard).await
372 }
373
374 pub(crate) async fn run_once(
381 &self,
382 opts: RunOptions,
383 stop_guard: tokio::sync::OwnedMutexGuard<()>,
384 ) -> Result<IpcResponse> {
385 let id = &opts.id;
386 let original_cmd = opts.cmd.clone(); let (ready_tx, ready_rx) = if opts.wait_ready {
390 let (tx, rx) = oneshot::channel();
391 (Some(tx), Some(rx))
392 } else {
393 (None, None)
394 };
395
396 let expected_ports = opts
398 .port
399 .as_ref()
400 .map(|p| p.expect.clone())
401 .unwrap_or_default();
402 let (resolved_ports, effective_ready_port) = if !expected_ports.is_empty() {
403 let port_cfg = opts.port.as_ref().unwrap();
404 match check_ports_available(
405 &expected_ports,
406 port_cfg.auto_bump(),
407 port_cfg.max_bump_attempts(),
408 )
409 .await
410 {
411 Ok(resolved) => {
412 let ready_port = if let Some(configured_port) =
413 opts.ready_port.as_ref().and_then(|p| p.as_port())
414 {
415 Some(resolve_configured_ready_port(
416 configured_port,
417 &expected_ports,
418 &resolved,
419 ))
420 } else if opts.ready_output.is_none()
421 && opts.ready_http.is_none()
422 && opts.ready_cmd.is_none()
423 && opts.ready_delay.is_none()
424 {
425 resolved.first().copied().filter(|&p| p != 0)
429 } else {
430 None
434 };
435 info!("daemon {id}: ports {expected_ports:?} resolved to {resolved:?}");
436 (resolved, ready_port)
437 }
438 Err(e) => {
439 error!("daemon {id}: port check failed: {e}");
440 if let Some(port_error) = e.downcast_ref::<PortError>() {
442 match port_error {
443 PortError::InUse { port, process, pid } => {
444 return Ok(IpcResponse::PortConflict {
445 port: *port,
446 process: process.clone(),
447 pid: *pid,
448 });
449 }
450 PortError::NoAvailablePort {
451 start_port,
452 attempts,
453 } => {
454 return Ok(IpcResponse::NoAvailablePort {
455 start_port: *start_port,
456 attempts: *attempts,
457 });
458 }
459 }
460 }
461 return Ok(IpcResponse::DaemonFailed {
462 error: e.to_string(),
463 });
464 }
465 }
466 } else {
467 if let Some(port) = opts.ready_port.as_ref().and_then(|p| p.as_port())
473 && port > 0
474 && let Some((pid, process)) = detect_port_conflict(port).await
475 {
476 return Ok(IpcResponse::PortConflict { port, process, pid });
477 }
478 (
479 Vec::new(),
480 opts.ready_port.as_ref().and_then(|p| p.as_port()),
481 )
482 };
483
484 let shell_setting = settings().general.shell.clone();
488 let shell_parts = match shell_words::split(&shell_setting) {
489 Ok(parts) if !parts.is_empty() => parts,
490 Ok(_) => {
491 return Ok(IpcResponse::DaemonFailed {
492 error: "general.shell setting is empty".to_string(),
493 });
494 }
495 Err(e) => {
496 return Ok(IpcResponse::DaemonFailed {
497 error: format!("failed to parse general.shell setting {shell_setting:?}: {e}"),
498 });
499 }
500 };
501 let (shell_program, shell_args) = shell_parts.split_first().unwrap();
502
503 let run_script = opts
508 .run
509 .clone()
510 .unwrap_or_else(|| shell_words::join(&original_cmd));
511
512 let (program, args) = if opts.mise.unwrap_or(settings().general.mise) {
513 match settings().resolve_mise_bin() {
514 Some(mise_bin) => {
515 let mise_bin_str = mise_bin.to_string_lossy().to_string();
516 info!("daemon {id}: wrapping command with mise ({mise_bin_str})");
517 let mut args = vec!["x".to_string(), "--".to_string()];
518 args.push(shell_program.clone());
519 args.extend(shell_args.iter().cloned());
520 args.push(run_script);
521 (mise_bin_str, args)
522 }
523 None => {
524 warn!("daemon {id}: mise=true but mise binary not found, running without mise");
525 let mut args: Vec<String> = shell_args.to_vec();
526 args.push(run_script);
527 (shell_program.clone(), args)
528 }
529 }
530 } else {
531 let mut args: Vec<String> = shell_args.to_vec();
532 args.push(run_script);
533 (shell_program.clone(), args)
534 };
535 #[cfg(unix)]
536 let run_identity = match resolve_effective_run_identity(opts.user.as_deref()) {
537 Ok(identity) => identity,
538 Err(e) => {
539 return Ok(IpcResponse::DaemonFailed {
540 error: e.to_string(),
541 });
542 }
543 };
544 info!("run: spawning daemon {id} with {program} {args:?}");
545
546 #[cfg(unix)]
548 let pty_pair = if opts.pty.unwrap_or(false) {
549 match super::pty::openpty() {
550 Ok(pair) => {
551 info!("daemon {id}: allocated PTY (pty = true)");
552 Some(pair)
553 }
554 Err(e) => {
555 warn!("daemon {id}: failed to allocate PTY, falling back to pipes: {e}");
556 None
557 }
558 }
559 } else {
560 None
561 };
562
563 let (output_tx, output_rx) = tokio::sync::mpsc::channel::<super::OutputLine>(256);
569 let mut output_relay = None;
570
571 let mut sink_pipe = None;
574 let mut sink_writer = None;
575 let mut sink_child = None;
576 if super::log_sink::is_supported(&opts) {
577 let log_format = opts
578 .log_format
579 .clone()
580 .unwrap_or_else(|| settings().logs.log_format.clone());
581 let watch_for = super::log_sink::WatchFor::from_opts(id, &opts);
582 let relay_token = if watch_for.is_empty() {
585 0
586 } else {
587 let relay = super::log_sink::OutputRelay::register(id, output_tx.clone());
588 let token = relay.token();
589 output_relay = Some(relay);
590 token
591 };
592 match super::log_sink::SinkPipe::new(log_format, watch_for, relay_token) {
593 Ok((pipe, writer)) => match pipe.start(id) {
594 Ok(child) => {
595 sink_child = Some(super::log_sink::PendingSink::new(child));
596 sink_pipe = Some(pipe);
597 sink_writer = Some(writer);
598 }
599 Err(e) => {
600 warn!("could not start log sink for {id}, capturing in-process: {e}");
601 }
602 },
603 Err(e) => {
604 warn!("could not create log pipe for {id}, capturing in-process: {e}");
607 }
608 }
609 }
610
611 let mut cmd = tokio::process::Command::new(&program);
612
613 #[cfg(unix)]
614 if let Some(ref pair) = pty_pair {
615 let slave_file = std::fs::File::from(
619 pair.slave
620 .try_clone()
621 .map_err(|e| miette::miette!("failed to dup slave PTY fd: {e}"))?,
622 );
623 cmd.stdin(std::process::Stdio::from(slave_file.try_clone().map_err(
624 |e| miette::miette!("failed to clone slave PTY fd for stdin: {e}"),
625 )?));
626 cmd.stdout(std::process::Stdio::from(slave_file.try_clone().map_err(
627 |e| miette::miette!("failed to clone slave PTY fd for stdout: {e}"),
628 )?));
629 cmd.stderr(std::process::Stdio::from(slave_file));
630 } else if let Some(writer) = sink_writer.take() {
631 let dup = writer
634 .try_clone()
635 .map_err(|e| miette::miette!("failed to dup log pipe for stderr: {e}"))?;
636 cmd.stdout(std::process::Stdio::from(writer))
637 .stderr(std::process::Stdio::from(dup));
638 } else {
639 cmd.stdout(std::process::Stdio::piped())
640 .stderr(std::process::Stdio::piped());
641 }
642
643 #[cfg(not(unix))]
644 if let Some(writer) = sink_writer.take() {
645 let dup = writer
646 .try_clone()
647 .map_err(|e| miette::miette!("failed to dup log pipe for stderr: {e}"))?;
648 cmd.stdout(std::process::Stdio::from(writer))
649 .stderr(std::process::Stdio::from(dup));
650 } else {
651 cmd.stdout(std::process::Stdio::piped())
652 .stderr(std::process::Stdio::piped());
653 }
654
655 cmd.args(&args).current_dir(&opts.dir).hide_console_window();
656
657 #[cfg(unix)]
658 if pty_pair.is_none() {
659 cmd.stdin(std::process::Stdio::null());
660 }
661
662 #[cfg(not(unix))]
663 cmd.stdin(std::process::Stdio::null());
664
665 apply_runtime_env(
666 &mut cmd,
667 id,
668 opts.retry_count,
669 opts.env.as_ref(),
670 &resolved_ports,
671 );
672
673 inject_proxy_env(&mut cmd, &opts.slug);
675
676 #[cfg(unix)]
677 {
678 let run_identity = run_identity.clone();
679 let use_pty = pty_pair.is_some();
680 unsafe {
681 cmd.pre_exec(move || {
682 nix::unistd::setsid().map_err(nix_to_io_error)?;
683
684 if use_pty {
688 let ret = libc::ioctl(0, libc::TIOCSCTTY as libc::c_ulong, 0);
689 if ret < 0 {
690 #[cfg(target_os = "linux")]
693 eprintln!(
694 "pitchfork: TIOCSCTTY failed: {}",
695 std::io::Error::last_os_error()
696 );
697 }
698 }
699
700 apply_run_identity(&run_identity)?;
701 Ok(())
702 });
703 }
704 }
705
706 let spawn_time = chrono::Local::now();
709 let mut child = cmd.spawn().into_diagnostic()?;
715 let pid = match child.id() {
716 Some(p) => p,
717 None => {
718 warn!("Daemon {id} exited before PID could be captured");
719 if sink_child.is_some() {
725 super::log_sink::wait_for_output(id, spawn_time, SINK_OUTPUT_TIMEOUT).await;
726 }
727 return Ok(IpcResponse::DaemonFailed {
728 error: "Process exited immediately".to_string(),
729 });
730 }
731 };
732 info!("started daemon {id} with pid {pid}");
733 PROCS.refresh_pids(&[pid]);
734 let monitored_guard = super::adopt::MonitoredGuard::register(id.clone(), pid);
742 let monitor_token = monitored_guard.token();
743
744 let using_sink = sink_pipe.is_some();
747 if let Some(pipe) = sink_pipe.take()
750 && let Some(child) = sink_child.as_mut().and_then(|pending| pending.take())
751 {
752 pipe.supervise(id.clone(), monitor_token, child);
753 }
754 let attempt_resolved_ports = resolved_ports.clone();
761 let daemon = self
762 .upsert_daemon(
763 UpsertDaemonOpts::from_run_options(&opts, DaemonStatus::Running)
764 .set(|o| {
765 o.pid = Some(pid);
766 o.cmd = Some(original_cmd);
767 o.ready_port = effective_ready_port.map(|p| ReadyPort {
768 port: Some(p),
769 template: None,
770 timeout: opts.ready_port.as_ref().and_then(|rp| rp.timeout),
771 });
772 o.port = crate::config_types::PortConfig::from_parts(
773 expected_ports,
774 opts.port.as_ref().map(|p| p.bump).unwrap_or_default(),
775 );
776 o.resolved_port = Some(resolved_ports);
777 })
778 .build(),
779 )
780 .await?;
781
782 drop(stop_guard);
787
788 let id_clone = id.clone();
789 let ready_delay = opts.ready_delay;
790 let ready_output = opts.ready_output.clone();
791 let ready_http = opts.ready_http.clone();
792 let ready_port = effective_ready_port;
793 let implicit_ready_port = ready_port.map(|p| ReadyPort {
794 port: Some(p),
795 template: None,
796 timeout: None,
797 });
798 let ready_port_config = opts.ready_port.clone().or(implicit_ready_port);
799 let ready_cmd = opts.ready_cmd.clone();
800 let daemon_dir = opts.dir.0.clone();
801 let hook_retry_count = opts.retry_count;
802 let hook_retry = opts.retry;
803 let hook_daemon_env = opts.env.clone();
804 let hook_resolved_ports = attempt_resolved_ports.clone();
809 let readiness_daemon_env = opts.env.clone();
810 let readiness_resolved_ports = attempt_resolved_ports.clone();
811 let on_output_hook = opts.on_output_hook.clone();
812 let has_port_config = opts.port.as_ref().is_some_and(|p| !p.expect.is_empty())
818 || (settings().proxy.enable && is_daemon_slug_target(id));
819 let daemon_pid = pid;
825
826 #[cfg(unix)]
830 let pty_reader = pty_pair.map(|p| {
831 tokio::io::BufReader::new(tokio::fs::File::from_std(std::fs::File::from(p.master)))
832 .lines()
833 });
834 #[cfg(not(unix))]
835 let pty_reader: Option<tokio::io::Lines<tokio::io::BufReader<tokio::fs::File>>> = None;
836 let stdout_reader = if pty_reader.is_none() {
837 child
838 .stdout
839 .take()
840 .map(|s| tokio::io::BufReader::new(s).lines())
841 } else {
842 None
843 };
844 let stderr_reader = if pty_reader.is_none() {
845 child
846 .stderr
847 .take()
848 .map(|s| tokio::io::BufReader::new(s).lines())
849 } else {
850 None
851 };
852
853 if !using_sink
854 && pty_reader.is_none()
855 && (stdout_reader.is_none() || stderr_reader.is_none())
856 {
857 error!("Failed to capture stdout/stderr for daemon {id}");
858 }
859
860 tokio::spawn(async move {
861 let id = id_clone;
862 let _monitored_guard = monitored_guard;
865 let output_relay = output_relay;
869
870 let mut output_rx = output_rx;
873
874 if let Some(mut reader) = pty_reader {
875 tokio::spawn(async move {
879 while let Ok(Some(mut line)) = reader.next_line().await {
880 if line.ends_with('\r') {
882 line.pop();
883 }
884 if output_tx
885 .send(super::OutputLine {
886 text: line,
887 source: super::OutputSource::Local,
888 })
889 .await
890 .is_err()
891 {
892 break;
893 }
894 }
895 });
896 } else {
897 if let Some(mut stdout) = stdout_reader {
902 let tx = output_tx.clone();
903 tokio::spawn(async move {
904 while let Ok(Some(line)) = stdout.next_line().await {
905 if tx
906 .send(super::OutputLine {
907 text: line,
908 source: super::OutputSource::Local,
909 })
910 .await
911 .is_err()
912 {
913 break;
914 }
915 }
916 });
917 }
918 if let Some(mut stderr) = stderr_reader {
919 let tx = output_tx.clone();
920 tokio::spawn(async move {
921 while let Ok(Some(line)) = stderr.next_line().await {
922 if tx
923 .send(super::OutputLine {
924 text: line,
925 source: super::OutputSource::Local,
926 })
927 .await
928 .is_err()
929 {
930 break;
931 }
932 }
933 });
934 }
935 drop(output_tx);
939 }
940 let log_store = Arc::clone(&LOG_STORE);
941 let log_format = opts
942 .log_format
943 .clone()
944 .unwrap_or_else(|| crate::settings::settings().logs.log_format.clone());
945 let parse_line = move |line: &str| crate::log_parse::parse(line, &log_format);
946
947 const LOG_BATCH_SIZE: usize = 100;
948 const LOG_FLUSH_INTERVAL: Duration = Duration::from_millis(100);
949 let mut log_buffer: Vec<crate::log_parse::ParsedLog> =
950 Vec::with_capacity(LOG_BATCH_SIZE);
951 let mut log_flush_interval = tokio::time::interval(LOG_FLUSH_INTERVAL);
952 log_flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
953
954 let flush_logs =
955 |buffer: &mut Vec<crate::log_parse::ParsedLog>| -> Option<tokio::task::JoinHandle<()>> {
956 if buffer.is_empty() {
957 return None;
958 }
959 let store = Arc::clone(&log_store);
960 let id = id.clone();
961 let batch = std::mem::take(buffer);
962 Some(tokio::task::spawn_blocking(move || {
963 if let Err(e) = store.append_structured_batch(&id, &batch) {
964 error!("Failed to write batch to log for daemon {id}: {e}");
965 }
966 }))
967 };
968
969 let mut ready_notified = false;
973 let mut ready_tx = ready_tx;
974 let ready_pattern = ready_output
975 .as_ref()
976 .and_then(|o| get_or_compile_regex(&o.pattern));
977 let mut active_port_spawned = false;
979
980 let on_output_hook = match on_output_hook {
984 Some(ref hook) => match hook.validate(id.name()) {
985 Ok(()) => on_output_hook,
986 Err(e) => {
987 error!("{e}");
988 None
989 }
990 },
991 None => None,
992 };
993
994 let on_output_pattern: Option<regex::Regex> = on_output_hook
997 .as_ref()
998 .and_then(|h| h.regex.as_deref().and_then(get_or_compile_regex));
999 let on_output_debounce = on_output_hook
1000 .as_ref()
1001 .map(|h| h.debounce_duration())
1002 .unwrap_or(Duration::from_millis(1000));
1003 let mut on_output_last_fired: Option<std::time::Instant> = None;
1005
1006 let mut delay_timer =
1007 ready_delay.map(|secs| Box::pin(time::sleep(Duration::from_secs(secs))));
1008
1009 let mut http_exhausted = false;
1011 let mut cmd_exhausted = false;
1012 let mut port_exhausted = false;
1013 let mut output_exhausted = false;
1014
1015 let s = settings();
1017 let ready_check_interval = s.supervisor_ready_check_interval();
1018 let http_client_timeout = s.supervisor_http_client_timeout();
1019
1020 let mut output_deadline = ready_output
1022 .as_ref()
1023 .and_then(|o| o.timeout)
1024 .map(|d| Box::pin(time::sleep(d)));
1025
1026 let mut http_check_interval = ready_http
1028 .as_ref()
1029 .map(|_| tokio::time::interval(ready_check_interval));
1030 let mut http_deadline = ready_http
1031 .as_ref()
1032 .and_then(|h| h.timeout)
1033 .map(|d| Box::pin(time::sleep(d)));
1034 let http_client = ready_http.as_ref().map(|_| {
1035 reqwest::Client::builder()
1036 .timeout(http_client_timeout)
1037 .build()
1038 .unwrap_or_default()
1039 });
1040
1041 let mut port_check_interval =
1043 ready_port.map(|_| tokio::time::interval(ready_check_interval));
1044 let mut port_deadline = ready_port_config
1045 .as_ref()
1046 .and_then(|p| p.timeout)
1047 .map(|d| Box::pin(time::sleep(d)));
1048
1049 let mut cmd_probe: Option<CmdProbe> = None;
1052 let mut cmd_respawn_delay: Option<_> = None;
1053 let mut cmd_deadline = ready_cmd
1054 .as_ref()
1055 .and_then(|c| c.timeout)
1056 .map(|d| Box::pin(time::sleep(d)));
1057 if let Some(ref cmd) = ready_cmd {
1058 cmd_probe = Some(spawn_cmd_probe(
1059 &id,
1060 &cmd.run,
1061 daemon_dir.as_path(),
1062 hook_retry_count,
1063 readiness_daemon_env.as_ref(),
1064 &readiness_resolved_ports,
1065 ));
1066 }
1067
1068 let (exit_tx, mut exit_rx) =
1070 tokio::sync::mpsc::channel::<std::io::Result<std::process::ExitStatus>>(1);
1071
1072 let child_pid = child.id().unwrap_or(0);
1074 tokio::spawn(async move {
1075 let result = child.wait().await;
1076 #[cfg(all(unix, not(target_os = "linux")))]
1086 let result = match &result {
1087 Err(e) if e.raw_os_error() == Some(nix::libc::ECHILD) => {
1088 if let Some(code) = super::REAPED_STATUSES.lock().await.remove(&child_pid) {
1089 warn!(
1090 "daemon pid {child_pid} wait() got ECHILD; \
1091 recovered exit code {code} from zombie reaper"
1092 );
1093 use std::os::unix::process::ExitStatusExt;
1098 if code >= 0 {
1099 Ok(std::process::ExitStatus::from_raw(code << 8))
1100 } else {
1101 Ok(std::process::ExitStatus::from_raw((-code) & 0x7f))
1103 }
1104 } else {
1105 warn!(
1106 "daemon pid {child_pid} wait() got ECHILD but no \
1107 stashed status found; reporting as error"
1108 );
1109 result
1110 }
1111 }
1112 _ => result,
1113 };
1114 debug!("daemon pid {child_pid} wait() completed with result: {result:?}");
1115 let _ = exit_tx.send(result).await;
1116 });
1117
1118 #[allow(unused_assignments)]
1119 let mut exit_status = None;
1121
1122 if has_port_config
1128 && ready_pattern.is_none()
1129 && ready_http.is_none()
1130 && ready_port.is_none()
1131 && ready_cmd.is_none()
1132 && delay_timer.is_none()
1133 {
1134 active_port_spawned = true;
1135 detect_and_store_active_port(id.clone(), daemon_pid);
1136 }
1137
1138 let mut ready_fail_kill: Option<tokio::task::JoinHandle<()>> = None;
1145
1146 loop {
1147 select! {
1151 biased;
1152 Some(result) = exit_rx.recv() => {
1153 exit_status = Some(result);
1155 debug!("daemon {id} process exited, exit_status: {exit_status:?}");
1156 if !ready_notified {
1157 if let Some(tx) = ready_tx.take() {
1158 let is_success = exit_status.as_ref()
1160 .and_then(|r| r.as_ref().ok())
1161 .map(|s| s.success())
1162 .unwrap_or(false);
1163
1164 if is_success {
1165 debug!("daemon {id} exited successfully before ready check, sending success notification");
1166 let _ = tx.send(Ok(()));
1167 } else {
1168 let exit_code = exit_status.as_ref()
1169 .and_then(|r| r.as_ref().ok())
1170 .and_then(|s| s.code());
1171 debug!("daemon {id} exited with failure before ready check, sending failure notification with exit_code: {exit_code:?}");
1172 let _ = tx.send(Err(exit_code));
1173 }
1174 }
1175 } else {
1176 debug!("daemon {id} was already marked ready, not sending notification");
1177 }
1178 break;
1179 },
1180 Some(super::OutputLine { text: line, source }) = output_rx.recv() => {
1181 if matches!(source, super::OutputSource::Local) {
1185 let parsed = parse_line(&line);
1186 log_buffer.push(parsed);
1187 if log_buffer.len() >= LOG_BATCH_SIZE {
1188 let _ = flush_logs(&mut log_buffer);
1189 }
1190 }
1191 trace!("output: {id} {line}");
1192
1193 let line_clean = console::strip_ansi_codes(&line).to_string();
1196
1197 if !ready_notified
1199 && !output_exhausted
1200 && let Some(ref pattern) = ready_pattern
1201 && pattern.is_match(&line_clean)
1202 {
1203 if let Some(handle) = flush_logs(&mut log_buffer) {
1208 let _ = handle.await;
1209 }
1210 info!("daemon {id} ready: output matched pattern");
1211 ready_notified = true;
1212 if let Some(tx) = ready_tx.take() {
1213 let _ = tx.send(Ok(()));
1214 }
1215 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), vec![]).await;
1216 stop_cmd_probe_state(&mut cmd_probe);
1217 http_deadline = None;
1218 cmd_deadline = None;
1219 port_deadline = None;
1220 output_deadline = None;
1221 if !active_port_spawned && has_port_config {
1222 active_port_spawned = true;
1223 detect_and_store_active_port(id.clone(), daemon_pid);
1224 }
1225 }
1226
1227 if let Some(ref hook) = on_output_hook {
1232 let matched = match source {
1233 super::OutputSource::Sink { fires_hook } => fires_hook,
1234 super::OutputSource::Local => match (&hook.filter, &on_output_pattern) {
1235 (Some(substr), _) => line_clean.contains(substr.as_str()),
1236 (None, Some(re)) => re.is_match(&line_clean),
1237 (None, None) => true,
1238 },
1239 };
1240 if matched {
1241 let now = std::time::Instant::now();
1246 let elapsed = on_output_last_fired.map(|t| now.duration_since(t));
1247 if elapsed.is_none_or(|e| e >= on_output_debounce) {
1248 on_output_last_fired = Some(now);
1249 hooks::fire_output_hook(id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), hook.run.clone(), line_clean.clone()).await;
1250 }
1251 }
1252 }
1253 tokio::task::yield_now().await;
1256 }
1257 _ = async {
1258 if let Some(ref mut deadline) = http_deadline {
1259 deadline.await;
1260 } else {
1261 std::future::pending::<()>().await;
1262 }
1263 }, if !ready_notified && ready_http.is_some() => {
1264 http_exhausted = true;
1265 http_deadline = None;
1266 http_check_interval = None;
1267 warn!("daemon {id}: HTTP readiness check timed out");
1268 let any_remaining = any_ready_check_remaining(
1269 ready_output.as_ref(),
1270 output_exhausted,
1271 ready_port_config.as_ref(),
1272 port_exhausted,
1273 ready_http.as_ref(),
1274 http_exhausted,
1275 ready_cmd.as_ref(),
1276 cmd_exhausted,
1277 );
1278 if !any_remaining {
1279 error!("daemon {id}: all readiness checks exhausted, failing");
1280 stop_cmd_probe_state(&mut cmd_probe);
1281 ready_fail_kill = Some(spawn_ready_fail_kill(
1282 id.clone(),
1283 daemon_pid,
1284 opts.stop_signal.unwrap_or_default(),
1285 ));
1286 break;
1287 }
1288 }
1289 _ = async {
1290 if let Some(ref mut deadline) = output_deadline {
1291 deadline.await;
1292 } else {
1293 std::future::pending::<()>().await;
1294 }
1295 }, if !ready_notified && ready_output.is_some() => {
1296 output_exhausted = true;
1297 output_deadline = None;
1298 warn!("daemon {id}: output readiness check timed out");
1299 let any_remaining = any_ready_check_remaining(
1300 ready_output.as_ref(),
1301 output_exhausted,
1302 ready_port_config.as_ref(),
1303 port_exhausted,
1304 ready_http.as_ref(),
1305 http_exhausted,
1306 ready_cmd.as_ref(),
1307 cmd_exhausted,
1308 );
1309 if !any_remaining {
1310 error!("daemon {id}: all readiness checks exhausted, failing");
1311 stop_cmd_probe_state(&mut cmd_probe);
1312 ready_fail_kill = Some(spawn_ready_fail_kill(
1313 id.clone(),
1314 daemon_pid,
1315 opts.stop_signal.unwrap_or_default(),
1316 ));
1317 break;
1318 }
1319 }
1320 _ = async {
1321 if let Some(ref mut interval) = http_check_interval {
1322 interval.tick().await;
1323 } else {
1324 std::future::pending::<()>().await;
1325 }
1326 }, if !ready_notified && ready_http.is_some() && !http_exhausted => {
1327 if let (Some(http), Some(client)) = (&ready_http, &http_client) {
1328 match client.get(&http.url).send().await {
1329 Ok(response) if http.accepts_status(response.status().as_u16()) => {
1330 info!("daemon {id} ready: HTTP check passed (status {})", response.status());
1331 ready_notified = true;
1332 if let Some(tx) = ready_tx.take() {
1333 let _ = tx.send(Ok(()));
1334 }
1335 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), vec![]).await;
1336 http_check_interval = None;
1337 http_deadline = None;
1338 stop_cmd_probe_state(&mut cmd_probe);
1339 cmd_deadline = None;
1340 port_deadline = None;
1341 output_deadline = None;
1342 if !active_port_spawned && has_port_config {
1343 active_port_spawned = true;
1344 detect_and_store_active_port(id.clone(), daemon_pid);
1345 }
1346 }
1347 Ok(response) => {
1348 trace!("daemon {id} HTTP check: status {} (not ready)", response.status());
1349 }
1350 Err(e) => {
1351 trace!("daemon {id} HTTP check failed: {e}");
1352 }
1353 }
1354 }
1355 }
1356 _ = async {
1357 if let Some(ref mut deadline) = port_deadline {
1358 deadline.await;
1359 } else {
1360 std::future::pending::<()>().await;
1361 }
1362 }, if !ready_notified && ready_port.is_some() => {
1363 port_exhausted = true;
1364 port_deadline = None;
1365 port_check_interval = None;
1366 warn!("daemon {id}: TCP port readiness check timed out");
1367 let any_remaining = any_ready_check_remaining(
1368 ready_output.as_ref(),
1369 output_exhausted,
1370 ready_port_config.as_ref(),
1371 port_exhausted,
1372 ready_http.as_ref(),
1373 http_exhausted,
1374 ready_cmd.as_ref(),
1375 cmd_exhausted,
1376 );
1377 if !any_remaining {
1378 error!("daemon {id}: all readiness checks exhausted, failing");
1379 stop_cmd_probe_state(&mut cmd_probe);
1380 ready_fail_kill = Some(spawn_ready_fail_kill(
1381 id.clone(),
1382 daemon_pid,
1383 opts.stop_signal.unwrap_or_default(),
1384 ));
1385 break;
1386 }
1387 }
1388 _ = async {
1389 if let Some(ref mut interval) = port_check_interval {
1390 interval.tick().await;
1391 } else {
1392 std::future::pending::<()>().await;
1393 }
1394 }, if !ready_notified && ready_port.is_some() && !port_exhausted => {
1395 if let Some(port) = ready_port {
1396 match tokio::net::TcpStream::connect(("127.0.0.1", port)).await {
1397 Ok(_) => {
1398 info!("daemon {id} ready: TCP port {port} is listening");
1399 ready_notified = true;
1400 if let Some(tx) = ready_tx.take() {
1401 let _ = tx.send(Ok(()));
1402 }
1403 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), vec![]).await;
1404 port_check_interval = None;
1406 port_deadline = None;
1407 stop_cmd_probe_state(&mut cmd_probe);
1408 http_deadline = None;
1409 cmd_deadline = None;
1410 output_deadline = None;
1411 if !active_port_spawned && has_port_config {
1412 active_port_spawned = true;
1413 if let Some(active_port) = active_port_from_ready_port(
1423 port,
1424 &readiness_resolved_ports,
1425 ) {
1426 let mut state_file =
1427 SUPERVISOR.state_file.lock().await;
1428 if let Some(d) = state_file.daemons.get(&id)
1429 && d.pid == Some(daemon_pid)
1430 {
1431 state_file.set_active_port(&id, active_port);
1432 }
1433 } else {
1434 detect_and_store_active_port(
1435 id.clone(),
1436 daemon_pid,
1437 );
1438 }
1439 }
1440 }
1441 Err(_) => {
1442 trace!("daemon {id} port check: port {port} not listening yet");
1443 }
1444 }
1445 }
1446 }
1447 _ = async {
1448 if let Some(ref mut delay) = cmd_respawn_delay {
1449 delay.await;
1450 } else {
1451 std::future::pending::<()>().await;
1452 }
1453 }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted && cmd_probe.is_none() => {
1454 if let Some(ref cmd) = ready_cmd {
1455 cmd_probe = Some(spawn_cmd_probe(
1456 &id,
1457 &cmd.run,
1458 daemon_dir.as_path(),
1459 hook_retry_count,
1460 readiness_daemon_env.as_ref(),
1461 &readiness_resolved_ports,
1462 ));
1463 }
1464 cmd_respawn_delay = None;
1465 }
1466 result = async {
1467 if let Some(probe) = cmd_probe.as_mut() {
1468 std::pin::Pin::new(&mut probe.result_rx).await
1469 } else {
1470 std::future::pending::<Result<Result<std::process::ExitStatus, std::io::Error>, tokio::sync::oneshot::error::RecvError>>().await
1471 }
1472 }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted => {
1473 let _ = cmd_probe.take();
1477 match result {
1478 Ok(Ok(status)) if status.success() => {
1479 info!("daemon {id} ready: readiness command succeeded");
1480 ready_notified = true;
1481 if let Some(tx) = ready_tx.take() {
1482 let _ = tx.send(Ok(()));
1483 }
1484 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), vec![]).await;
1485 cmd_respawn_delay = None;
1486 cmd_deadline = None;
1487 http_deadline = None;
1488 port_deadline = None;
1489 output_deadline = None;
1490 if !active_port_spawned && has_port_config {
1491 active_port_spawned = true;
1492 detect_and_store_active_port(id.clone(), daemon_pid);
1493 }
1494 }
1495 Ok(Ok(_)) | Ok(Err(_)) | Err(_) => {
1496 trace!("daemon {id} cmd check: command not ready, will respawn");
1497 cmd_respawn_delay = Some(Box::pin(time::sleep(ready_check_interval)));
1498 }
1499 }
1500 }
1501 _ = async {
1502 if let Some(ref mut deadline) = cmd_deadline {
1503 deadline.await;
1504 } else {
1505 std::future::pending::<()>().await;
1506 }
1507 }, if !ready_notified && ready_cmd.is_some() => {
1508 cmd_exhausted = true;
1509 cmd_deadline = None;
1510 stop_cmd_probe_state(&mut cmd_probe);
1511 cmd_respawn_delay = None;
1512 warn!("daemon {id}: command readiness check timed out");
1513 let any_remaining = any_ready_check_remaining(
1514 ready_output.as_ref(),
1515 output_exhausted,
1516 ready_port_config.as_ref(),
1517 port_exhausted,
1518 ready_http.as_ref(),
1519 http_exhausted,
1520 ready_cmd.as_ref(),
1521 cmd_exhausted,
1522 );
1523 if !any_remaining {
1524 error!("daemon {id}: all readiness checks exhausted, failing");
1525 ready_fail_kill = Some(spawn_ready_fail_kill(
1526 id.clone(),
1527 daemon_pid,
1528 opts.stop_signal.unwrap_or_default(),
1529 ));
1530 break;
1531 }
1532 }
1533 _ = async {
1534 if let Some(ref mut timer) = delay_timer {
1535 timer.await;
1536 } else {
1537 std::future::pending::<()>().await;
1538 }
1539 } => {
1540 let has_other_ready_check = ready_pattern.is_some()
1541 || ready_http.is_some()
1542 || ready_port.is_some()
1543 || ready_cmd.is_some();
1544 let delay_is_only_readiness = !ready_notified && !has_other_ready_check;
1545 let process_exited = exit_status.is_some();
1546 let process_running = if delay_is_only_readiness && !process_exited {
1547 PROCS.refresh_pids(&[daemon_pid]);
1550 PROCS.is_running(daemon_pid)
1551 } else {
1552 false
1553 };
1554
1555 if delay_readiness_succeeded(
1556 ready_notified,
1557 has_other_ready_check,
1558 process_exited,
1559 process_running,
1560 ) {
1561 info!("daemon {id} ready: delay elapsed");
1562 ready_notified = true;
1563 if let Some(tx) = ready_tx.take() {
1564 let _ = tx.send(Ok(()));
1565 }
1566 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), vec![]).await;
1567 if !active_port_spawned && has_port_config {
1568 active_port_spawned = true;
1569 detect_and_store_active_port(id.clone(), daemon_pid);
1570 }
1571 } else if delay_is_only_readiness {
1572 if process_exited {
1573 debug!("daemon {id} exited during ready_delay, not marking as ready");
1574 } else {
1575 debug!("daemon {id} pid {daemon_pid} not running during ready_delay, deferring to exit handler");
1576 }
1577 }
1578
1579 if delay_is_only_readiness {
1580 output_deadline = None;
1583 http_deadline = None;
1584 cmd_deadline = None;
1585 port_deadline = None;
1586 stop_cmd_probe_state(&mut cmd_probe);
1587 }
1588 delay_timer = None;
1590 }
1591 _ = log_flush_interval.tick() => {
1592 let _ = flush_logs(&mut log_buffer);
1593 }
1594 }
1595 }
1596
1597 let pre_drain_daemon = SUPERVISOR.get_daemon(&id).await;
1611 let pre_drain_is_stopping = pre_drain_daemon
1612 .as_ref()
1613 .is_some_and(|d| d.status.is_stopped() || d.status.is_stopping());
1614
1615 drop(output_relay);
1628 let drain_deadline = tokio::time::Instant::now() + Duration::from_secs(5);
1629 loop {
1630 let now = tokio::time::Instant::now();
1631 if now >= drain_deadline {
1632 break;
1633 }
1634 let Ok(Some(line)) =
1635 tokio::time::timeout(drain_deadline - now, output_rx.recv()).await
1636 else {
1637 break;
1638 };
1639 if matches!(line.source, super::OutputSource::Local) {
1641 log_buffer.push(parse_line(&line.text));
1642 }
1643 }
1644 if let Some(handle) = flush_logs(&mut log_buffer) {
1647 let _ = handle.await;
1648 }
1649
1650 {
1652 let mut state_file = SUPERVISOR.state_file.lock().await;
1653 state_file.clear_active_port(&id);
1654 }
1655
1656 let exit_status = if let Some(status) = exit_status {
1658 status
1659 } else {
1660 match exit_rx.recv().await {
1662 Some(status) => status,
1663 None => {
1664 warn!("daemon {id} exit channel closed without receiving status");
1665 Err(std::io::Error::other("exit channel closed"))
1666 }
1667 }
1668 };
1669
1670 if let Some(kill) = ready_fail_kill {
1675 let _ = kill.await;
1676 if let Some(tx) = ready_tx.take() {
1677 let _ = tx.send(Err(Some(124)));
1678 }
1679 }
1680
1681 let current_daemon = SUPERVISOR.get_daemon(&id).await;
1682
1683 SUPERVISOR
1688 .active_monitors
1689 .fetch_add(1, atomic::Ordering::Release);
1690 struct MonitorGuard;
1691 impl Drop for MonitorGuard {
1692 fn drop(&mut self) {
1693 SUPERVISOR
1694 .active_monitors
1695 .fetch_sub(1, atomic::Ordering::Release);
1696 SUPERVISOR.monitor_done.notify_waiters();
1697 }
1698 }
1699 let _monitor_guard = MonitorGuard;
1700 if !pre_drain_is_stopping
1705 && (current_daemon.is_none()
1706 || current_daemon.as_ref().is_some_and(|d| {
1707 d.pid != Some(pid) && !d.status.is_stopped() && !d.status.is_stopping()
1708 }))
1709 {
1710 return;
1712 }
1713 let already_stopped = current_daemon
1718 .as_ref()
1719 .is_some_and(|d| d.status.is_stopped());
1720 let is_stopping = already_stopped
1721 || pre_drain_is_stopping
1722 || current_daemon
1723 .as_ref()
1724 .is_some_and(|d| d.status.is_stopping());
1725
1726 let (exit_code, exit_reason) = match (&exit_status, is_stopping) {
1728 (Ok(status), true) => {
1729 (status.code().unwrap_or(-1), "stop")
1733 }
1734 (Ok(status), false) if status.success() => (status.code().unwrap_or(-1), "exit"),
1735 (Ok(status), false) => (status.code().unwrap_or(-1), "fail"),
1736 (Err(_), true) => {
1737 (-1, "stop")
1739 }
1740 (Err(_), false) => (-1, "fail"),
1741 };
1742
1743 if !already_stopped && !pre_drain_is_stopping {
1749 if let Ok(status) = &exit_status {
1750 info!("daemon {id} exited with status {status}");
1751 }
1752 let (new_status, last_exit_success) = match exit_reason {
1753 "stop" | "exit" => (
1754 DaemonStatus::Stopped,
1755 exit_status.as_ref().map(|s| s.success()).unwrap_or(true),
1756 ),
1757 _ => (DaemonStatus::Errored(exit_code), false),
1758 };
1759 if !SUPERVISOR
1765 .finalize_monitored_exit(
1766 &id,
1767 pid,
1768 monitor_token,
1769 new_status,
1770 Some(last_exit_success),
1771 )
1772 .await
1773 {
1774 debug!("daemon {id} exit state was not written; a successor owns the record");
1775 }
1776 }
1777
1778 let hook_extra_env = vec![
1780 ("PITCHFORK_EXIT_CODE".to_string(), exit_code.to_string()),
1781 ("PITCHFORK_EXIT_REASON".to_string(), exit_reason.to_string()),
1782 ];
1783
1784 let hooks_to_fire: Vec<HookType> = match exit_reason {
1786 "stop" => vec![HookType::OnStop, HookType::OnExit],
1787 "exit" => vec![HookType::OnExit],
1788 _ if hook_retry_count >= hook_retry.count() => {
1790 vec![HookType::OnFail, HookType::OnExit]
1791 }
1792 _ => vec![],
1793 };
1794
1795 for hook_type in hooks_to_fire {
1796 fire_hook(
1797 hook_type,
1798 id.clone(),
1799 daemon_dir.clone(),
1800 hook_retry_count,
1801 hook_daemon_env.clone(),
1802 hook_resolved_ports.clone(),
1803 hook_extra_env.clone(),
1804 )
1805 .await;
1806 }
1807 });
1808
1809 if let Some(ready_rx) = ready_rx {
1811 match ready_rx.await {
1812 Ok(Ok(())) => {
1813 info!("daemon {id} is ready");
1814 Ok(IpcResponse::DaemonReady { daemon })
1815 }
1816 Ok(Err(exit_code)) => {
1817 error!("daemon {id} failed before becoming ready");
1818 let last_attempt = opts.retry_count >= opts.retry.count();
1829 if using_sink && last_attempt {
1830 super::log_sink::wait_for_output(id, spawn_time, SINK_OUTPUT_TIMEOUT).await;
1831 }
1832 Ok(IpcResponse::DaemonFailedWithCode {
1833 exit_code,
1834 resolved_ports: attempt_resolved_ports,
1835 })
1836 }
1837 Err(_) => {
1838 error!("readiness channel closed unexpectedly for daemon {id}");
1839 Ok(IpcResponse::DaemonStart { daemon })
1840 }
1841 }
1842 } else {
1843 Ok(IpcResponse::DaemonStart { daemon })
1844 }
1845 }
1846
1847 pub async fn stop(&self, id: &DaemonId) -> Result<IpcResponse> {
1849 let lock = self.stop_lock(id).await;
1853 let _guard = lock.lock().await;
1854 self.stop_locked(id).await
1855 }
1856
1857 async fn stop_locked(&self, id: &DaemonId) -> Result<IpcResponse> {
1859 let pitchfork_id = DaemonId::pitchfork();
1860 if *id == pitchfork_id {
1861 return Ok(IpcResponse::Error(
1862 "Cannot stop supervisor via stop command".into(),
1863 ));
1864 }
1865 info!("stopping daemon: {id}");
1866 if let Some(daemon) = self.get_daemon(id).await {
1867 trace!("daemon to stop: {daemon}");
1868 if let Some(pid) = daemon.pid {
1869 trace!("killing pid: {pid}");
1870 if PROCS.is_running(pid) {
1871 if !super::signalling_pid_is_authorized(
1877 daemon.start_time,
1878 PROCS.start_time(pid),
1879 ) {
1880 warn!(
1881 "pid {pid} recorded for daemon {id} belongs to another process now; not signalling it"
1882 );
1883 self.upsert_daemon(
1884 UpsertDaemonOpts::builder(id.clone())
1885 .set(|o| {
1886 o.pid = None;
1887 o.status = DaemonStatus::Stopped;
1888 })
1889 .build(),
1890 )
1891 .await?;
1892 return Ok(IpcResponse::DaemonWasNotRunning);
1893 }
1894
1895 self.upsert_daemon(
1897 UpsertDaemonOpts::builder(id.clone())
1898 .set(|o| {
1899 o.pid = Some(pid);
1900 o.status = DaemonStatus::Stopping;
1901 })
1902 .build(),
1903 )
1904 .await?;
1905
1906 let stop_cfg = daemon.stop_signal.unwrap_or_default();
1909 let stop_signal: i32 = stop_cfg.signal.into();
1910 if let Err(e) = PROCS
1911 .kill_process_group_async(pid, stop_signal, stop_cfg.timeout)
1912 .await
1913 {
1914 debug!("failed to kill pid {pid}: {e}");
1915 if PROCS.process_group_alive(pid) {
1921 debug!(
1923 "failed to stop pid {pid}: process group still alive after kill"
1924 );
1925 self.upsert_daemon(
1926 UpsertDaemonOpts::builder(id.clone())
1927 .set(|o| {
1928 o.pid = Some(pid); o.status = DaemonStatus::Running;
1930 })
1931 .build(),
1932 )
1933 .await?;
1934 return Ok(IpcResponse::DaemonStopFailed {
1935 error: format!(
1936 "process group of {pid} still alive after kill attempt: {e}"
1937 ),
1938 });
1939 }
1940 }
1941
1942 self.upsert_daemon(
1950 UpsertDaemonOpts::builder(id.clone())
1951 .set(|o| {
1952 o.pid = None;
1953 o.status = DaemonStatus::Stopped;
1954 o.last_exit_success = Some(true);
1955 })
1956 .build(),
1957 )
1958 .await?;
1959 } else {
1960 debug!("pid {pid} not running, process may have exited unexpectedly");
1961 self.upsert_daemon(
1967 UpsertDaemonOpts::builder(id.clone())
1968 .set(|o| {
1969 o.pid = None;
1970 o.status = DaemonStatus::Stopped;
1971 })
1972 .build(),
1973 )
1974 .await?;
1975 return Ok(IpcResponse::DaemonWasNotRunning);
1976 }
1977 Ok(IpcResponse::Ok)
1978 } else {
1979 debug!("daemon {id} not running");
1980 Ok(IpcResponse::DaemonNotRunning)
1981 }
1982 } else {
1983 debug!("daemon {id} not found");
1984 Ok(IpcResponse::DaemonNotFound)
1985 }
1986 }
1987}
1988
1989#[cfg(unix)]
1990fn resolve_effective_run_identity(daemon_user: Option<&str>) -> Result<RunIdentity> {
1991 let s = settings();
1992 let settings_user = s.supervisor.user.trim();
1993 let daemon_user = daemon_user.map(str::trim).filter(|user| !user.is_empty());
1994 let settings_user = (!settings_user.is_empty()).then_some(settings_user);
1995 let configured = daemon_user.or(settings_user);
1996 let current_uid = nix::unistd::Uid::effective().as_raw();
1997 let current_gid = nix::unistd::Gid::effective().as_raw();
1998 resolve_run_identity(
1999 configured,
2000 current_uid,
2001 current_gid,
2002 std::env::var("SUDO_UID").ok().as_deref(),
2003 std::env::var("SUDO_GID").ok().as_deref(),
2004 )
2005}
2006
2007#[cfg(unix)]
2008fn resolve_run_identity(
2009 configured: Option<&str>,
2010 current_uid: u32,
2011 current_gid: u32,
2012 sudo_uid: Option<&str>,
2013 sudo_gid: Option<&str>,
2014) -> Result<RunIdentity> {
2015 let current_uid = nix::unistd::Uid::from_raw(current_uid);
2016 let current_gid = nix::unistd::Gid::from_raw(current_gid);
2017 if let Some(user) = configured {
2018 let identity = resolve_configured_user(user)?;
2019 ensure_can_use_identity(user, &identity, current_uid, current_gid)?;
2020 if identity.matches(current_uid, current_gid) {
2021 return Ok(RunIdentity::Inherit);
2022 }
2023 return Ok(identity);
2024 }
2025
2026 if current_uid.is_root()
2027 && let Some(identity) = resolve_sudo_identity(sudo_uid, sudo_gid)
2028 {
2029 return Ok(identity);
2030 }
2031
2032 Ok(RunIdentity::Inherit)
2033}
2034
2035#[cfg(unix)]
2036fn resolve_configured_user(user: &str) -> Result<RunIdentity> {
2037 if user.chars().all(|c| c.is_ascii_digit()) {
2038 let uid = user
2039 .parse::<u32>()
2040 .map_err(|e| miette::miette!("invalid run user UID '{}': {}", user, e))?;
2041 let user_record = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
2042 .into_diagnostic()?
2043 .ok_or_else(|| miette::miette!("run user UID '{}' does not exist", user))?;
2044 return run_identity_from_user_record(user_record);
2045 }
2046
2047 let user_record = nix::unistd::User::from_name(user)
2048 .into_diagnostic()?
2049 .ok_or_else(|| miette::miette!("run user '{}' does not exist", user))?;
2050 run_identity_from_user_record(user_record)
2051}
2052
2053#[cfg(unix)]
2054fn run_identity_from_user_record(user: nix::unistd::User) -> Result<RunIdentity> {
2055 let username = CString::new(user.name)
2056 .map_err(|e| miette::miette!("run user name contains an interior nul byte: {}", e))?;
2057 Ok(RunIdentity::Switch {
2058 uid: user.uid,
2059 gid: user.gid,
2060 username: Some(username),
2061 })
2062}
2063
2064#[cfg(unix)]
2065fn run_identity_from_raw_ids(uid: u32, gid: u32, username: Option<CString>) -> RunIdentity {
2066 RunIdentity::Switch {
2067 uid: nix::unistd::Uid::from_raw(uid),
2068 gid: nix::unistd::Gid::from_raw(gid),
2069 username,
2070 }
2071}
2072
2073#[cfg(unix)]
2074fn resolve_sudo_identity(sudo_uid: Option<&str>, sudo_gid: Option<&str>) -> Option<RunIdentity> {
2075 let uid = sudo_uid?.parse::<u32>().ok()?;
2076 let gid = sudo_gid?.parse::<u32>().ok()?;
2077 let username = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
2078 .ok()
2079 .flatten()
2080 .and_then(|u| CString::new(u.name).ok());
2081 Some(run_identity_from_raw_ids(uid, gid, username))
2082}
2083
2084#[cfg(unix)]
2085fn ensure_can_use_identity(
2086 configured_user: &str,
2087 identity: &RunIdentity,
2088 current_uid: nix::unistd::Uid,
2089 current_gid: nix::unistd::Gid,
2090) -> Result<()> {
2091 let RunIdentity::Switch { uid, gid, .. } = identity else {
2092 return Ok(());
2093 };
2094 if *uid == current_uid && *gid == current_gid {
2095 return Ok(());
2096 }
2097 if current_uid.is_root() {
2098 return Ok(());
2099 }
2100 Err(miette::miette!(
2101 "daemon is configured to run as '{}', but the supervisor is running as uid={} gid={}. Restart the supervisor with sudo to switch to uid={} gid={}, or choose a user matching the supervisor.",
2102 configured_user,
2103 current_uid.as_raw(),
2104 current_gid.as_raw(),
2105 uid.as_raw(),
2106 gid.as_raw()
2107 ))
2108}
2109
2110#[cfg(unix)]
2111fn apply_run_identity(identity: &RunIdentity) -> std::io::Result<()> {
2112 let RunIdentity::Switch { uid, gid, username } = identity else {
2113 return Ok(());
2114 };
2115 if let Some(username) = username {
2116 initgroups_for_user(username, *gid)?;
2117 } else {
2118 setgroups_to_primary(*gid)?;
2119 }
2120 nix::unistd::setgid(*gid).map_err(nix_to_io_error)?;
2121 nix::unistd::setuid(*uid).map_err(nix_to_io_error)?;
2122 Ok(())
2123}
2124
2125#[cfg(unix)]
2126impl RunIdentity {
2127 fn matches(&self, uid: nix::unistd::Uid, gid: nix::unistd::Gid) -> bool {
2128 matches!(self, RunIdentity::Switch { uid: u, gid: g, .. } if *u == uid && *g == gid)
2129 }
2130}
2131
2132#[cfg(unix)]
2133fn setgroups_to_primary(gid: nix::unistd::Gid) -> std::io::Result<()> {
2134 let groups = [gid.as_raw() as libc::gid_t];
2135 #[cfg(any(target_os = "linux", target_os = "android"))]
2136 let group_count = groups.len();
2137 #[cfg(not(any(target_os = "linux", target_os = "android")))]
2138 let group_count = groups.len() as libc::c_int;
2139 let rc = unsafe { libc::setgroups(group_count, groups.as_ptr()) };
2140 if rc == -1 {
2141 Err(std::io::Error::last_os_error())
2142 } else {
2143 Ok(())
2144 }
2145}
2146
2147#[cfg(unix)]
2148fn initgroups_for_user(username: &CString, gid: nix::unistd::Gid) -> std::io::Result<()> {
2149 let gid = gid.as_raw();
2150 #[cfg(any(
2151 target_os = "macos",
2152 target_os = "ios",
2153 target_os = "tvos",
2154 target_os = "watchos"
2155 ))]
2156 let base_gid = i32::try_from(gid)
2157 .map_err(|_| std::io::Error::other(format!("gid {gid} is out of range")))?;
2158
2159 #[cfg(not(any(
2160 target_os = "macos",
2161 target_os = "ios",
2162 target_os = "tvos",
2163 target_os = "watchos"
2164 )))]
2165 let base_gid = gid as libc::gid_t;
2166
2167 let rc = unsafe { libc::initgroups(username.as_ptr(), base_gid) };
2170 if rc == -1 {
2171 Err(std::io::Error::last_os_error())
2172 } else {
2173 Ok(())
2174 }
2175}
2176
2177#[cfg(unix)]
2178fn nix_to_io_error(err: nix::errno::Errno) -> std::io::Error {
2179 std::io::Error::from_raw_os_error(err as i32)
2180}
2181
2182async fn check_ports_available(
2189 expected_ports: &[u16],
2190 auto_bump: bool,
2191 max_attempts: u32,
2192) -> Result<Vec<u16>> {
2193 if expected_ports.is_empty() {
2194 return Ok(Vec::new());
2195 }
2196
2197 for bump_offset in 0..=max_attempts {
2198 let candidate_ports: Vec<u16> = expected_ports
2200 .iter()
2201 .map(|&p| p.wrapping_add(bump_offset as u16))
2202 .collect();
2203
2204 let mut all_available = true;
2206 let mut conflicting_port = None;
2207
2208 for &port in &candidate_ports {
2209 if port == 0 {
2212 continue;
2213 }
2214
2215 if is_port_in_use(port).await {
2229 all_available = false;
2230 conflicting_port = Some(port);
2231 break;
2232 }
2233 }
2234
2235 if all_available {
2236 if candidate_ports.contains(&0) && !expected_ports.contains(&0) {
2240 return Err(PortError::NoAvailablePort {
2241 start_port: expected_ports[0],
2242 attempts: bump_offset + 1,
2243 }
2244 .into());
2245 }
2246 if bump_offset > 0 {
2247 info!("ports {expected_ports:?} bumped by {bump_offset} to {candidate_ports:?}");
2248 }
2249 return Ok(candidate_ports);
2250 }
2251
2252 if bump_offset == 0
2254 && !auto_bump
2255 && let Some(port) = conflicting_port
2256 {
2257 let (pid, process) = identify_port_owner(port).await;
2258 return Err(PortError::InUse { port, process, pid }.into());
2259 }
2260 }
2261
2262 Err(PortError::NoAvailablePort {
2264 start_port: expected_ports[0],
2265 attempts: max_attempts + 1,
2266 }
2267 .into())
2268}
2269
2270async fn is_port_in_use(port: u16) -> bool {
2276 tokio::task::spawn_blocking(move || {
2277 for &addr in &["0.0.0.0", "127.0.0.1", "::1"] {
2278 match std::net::TcpListener::bind((addr, port)) {
2279 Ok(listener) => drop(listener),
2280 Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => return true,
2281 Err(_) => continue,
2282 }
2283 }
2284 false
2285 })
2286 .await
2287 .unwrap_or(false)
2288}
2289
2290async fn identify_port_owner(port: u16) -> (u32, String) {
2295 tokio::task::spawn_blocking(move || {
2296 listeners::get_all()
2297 .ok()
2298 .and_then(|list| {
2299 list.into_iter()
2300 .find(|l| l.socket.port() == port)
2301 .map(|l| (l.process.pid, l.process.name))
2302 })
2303 .unwrap_or((0, "unknown".to_string()))
2304 })
2305 .await
2306 .unwrap_or((0, "unknown".to_string()))
2307}
2308
2309async fn detect_port_conflict(port: u16) -> Option<(u32, String)> {
2314 if !is_port_in_use(port).await {
2315 return None;
2316 }
2317 Some(identify_port_owner(port).await)
2318}
2319
2320#[derive(Debug, PartialEq, Eq)]
2321enum ActivePortSelection {
2322 NoCandidates,
2323 Selected(u16),
2324 Ambiguous(Vec<u16>),
2325}
2326
2327fn discovery_preferred_port(daemon: &crate::daemon::Daemon) -> Option<u16> {
2328 daemon
2329 .resolved_port
2330 .first()
2331 .copied()
2332 .or_else(|| {
2333 daemon
2334 .port
2335 .as_ref()
2336 .and_then(|port| port.expect.first().copied())
2337 })
2338 .filter(|&port| port > 0)
2339}
2340
2341fn select_active_port(
2342 listeners: impl IntoIterator<Item = listeners::Listener>,
2343 descendant_pids: &std::collections::HashSet<u32>,
2344 preferred_port: Option<u16>,
2345) -> ActivePortSelection {
2346 let process_ports: std::collections::BTreeSet<u16> = listeners
2347 .into_iter()
2348 .filter(|listener| {
2349 listener.protocol == listeners::Protocol::TCP
2350 && listener.state == listeners::SocketState::Listen
2351 && descendant_pids.contains(&listener.process.pid)
2352 })
2353 .map(|listener| listener.socket.port())
2354 .filter(|&port| port > 0)
2355 .collect();
2356
2357 if let Some(port) = preferred_port
2358 && process_ports.contains(&port)
2359 {
2360 return ActivePortSelection::Selected(port);
2361 }
2362
2363 match process_ports.len() {
2364 0 => ActivePortSelection::NoCandidates,
2365 1 => ActivePortSelection::Selected(*process_ports.first().unwrap()),
2366 _ => ActivePortSelection::Ambiguous(process_ports.into_iter().collect()),
2367 }
2368}
2369
2370fn detect_and_store_active_port(id: DaemonId, pid: u32) {
2381 tokio::spawn(async move {
2382 for delay_ms in [500u64, 1000, 2000, 4000] {
2386 tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
2387
2388 let preferred_port: Option<u16> = {
2391 let state_file = SUPERVISOR.state_file.lock().await;
2392 match state_file.daemons.get(&id) {
2393 Some(d) if d.pid.is_none() => {
2394 debug!("daemon {id}: aborting active_port detection — process exited");
2395 return;
2396 }
2397 Some(d) => discovery_preferred_port(d),
2398 None => None,
2399 }
2400 };
2401
2402 let selection = tokio::task::spawn_blocking(move || {
2403 let listeners = listeners::get_all().ok()?;
2404
2405 PROCS.refresh_processes();
2407
2408 let descendant_pids: std::collections::HashSet<u32> = PROCS
2409 .all_children(pid)
2410 .into_iter()
2411 .chain(std::iter::once(pid))
2412 .collect();
2413
2414 Some(select_active_port(
2415 listeners,
2416 &descendant_pids,
2417 preferred_port,
2418 ))
2419 })
2420 .await
2421 .ok()
2422 .flatten()
2423 .unwrap_or(ActivePortSelection::NoCandidates);
2424
2425 let port = match selection {
2426 ActivePortSelection::Selected(port) => port,
2427 ActivePortSelection::Ambiguous(ports) => {
2428 debug!(
2429 "daemon {id}: ambiguous active_port candidates {ports:?} for pid {pid} \
2430 and its descendants; leaving active_port unset (will retry)"
2431 );
2432 continue;
2433 }
2434 ActivePortSelection::NoCandidates => {
2435 debug!(
2436 "daemon {id}: no active port detected for pid {pid} or its descendants \
2437 (will retry)"
2438 );
2439 continue;
2440 }
2441 };
2442
2443 debug!("daemon {id} active_port detected: {port}");
2444 let mut state_file = SUPERVISOR.state_file.lock().await;
2445 if let Some(d) = state_file.daemons.get(&id) {
2446 if d.pid == Some(pid) {
2450 state_file.set_active_port(&id, port);
2451 } else {
2452 debug!(
2453 "daemon {id}: skipping active_port write — PID mismatch \
2454 (expected {pid}, current {:?})",
2455 d.pid
2456 );
2457 }
2458 }
2459 return;
2460 }
2461
2462 debug!(
2463 "daemon {id}: active port detection exhausted all retries for pid {pid} and its descendants"
2464 );
2465 });
2466}
2467
2468#[cfg(test)]
2469mod active_port_tests {
2470 use super::*;
2471 use crate::config_types::PortConfig;
2472 use listeners::{Listener, Process, Protocol, SocketState};
2473 use std::net::{IpAddr, Ipv4Addr, SocketAddr};
2474
2475 fn listener(pid: u32, port: u16, protocol: Protocol, state: SocketState) -> Listener {
2476 Listener {
2477 process: Process {
2478 pid,
2479 name: "test".to_string(),
2480 path: "/test".to_string(),
2481 },
2482 socket: SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), port),
2483 protocol,
2484 state,
2485 }
2486 }
2487
2488 #[test]
2489 fn active_port_candidates_exclude_outbound_tcp_and_udp_sockets() {
2490 let daemon_pid = 100;
2491 let child_pid = 101;
2492 let descendant_pids = [daemon_pid, child_pid].into_iter().collect();
2493 let listeners = vec![
2494 listener(child_pid, 3004, Protocol::TCP, SocketState::Listen),
2495 listener(daemon_pid, 47082, Protocol::TCP, SocketState::Established),
2496 listener(daemon_pid, 5353, Protocol::UDP, SocketState::Unknown),
2497 listener(999, 9000, Protocol::TCP, SocketState::Listen),
2498 ];
2499
2500 assert_eq!(
2501 select_active_port(listeners, &descendant_pids, None),
2502 ActivePortSelection::Selected(3004)
2503 );
2504 }
2505
2506 #[test]
2507 fn bumped_cmd_readiness_prefers_resolved_primary_port() {
2508 let daemon = crate::daemon::Daemon {
2509 resolved_port: vec![3004],
2510 port: Some(PortConfig {
2511 expect: vec![3000],
2512 ..PortConfig::default()
2513 }),
2514 ..crate::daemon::Daemon::default()
2515 };
2516 let descendant_pids = [100].into_iter().collect();
2517 let listeners = vec![
2518 listener(100, 9000, Protocol::TCP, SocketState::Listen),
2519 listener(100, 3004, Protocol::TCP, SocketState::Listen),
2520 ];
2521
2522 assert_eq!(discovery_preferred_port(&daemon), Some(3004));
2523 assert_eq!(
2524 select_active_port(
2525 listeners,
2526 &descendant_pids,
2527 discovery_preferred_port(&daemon),
2528 ),
2529 ActivePortSelection::Selected(3004)
2530 );
2531 }
2532
2533 #[test]
2534 fn legacy_state_uses_expected_primary_port_for_discovery() {
2535 let daemon = crate::daemon::Daemon {
2536 port: Some(PortConfig {
2537 expect: vec![3000],
2538 ..PortConfig::default()
2539 }),
2540 ..crate::daemon::Daemon::default()
2541 };
2542
2543 assert_eq!(discovery_preferred_port(&daemon), Some(3000));
2544 }
2545
2546 #[test]
2547 fn ambiguous_candidates_leave_active_port_unset() {
2548 let descendant_pids = [100].into_iter().collect();
2549 let listeners = vec![
2550 listener(100, 9000, Protocol::TCP, SocketState::Listen),
2551 listener(100, 3000, Protocol::TCP, SocketState::Listen),
2552 ];
2553
2554 assert_eq!(
2555 select_active_port(listeners, &descendant_pids, None),
2556 ActivePortSelection::Ambiguous(vec![3000, 9000])
2557 );
2558 }
2559
2560 #[test]
2561 fn bumped_ready_port_sets_only_the_resolved_primary_without_scanning() {
2562 assert_eq!(active_port_from_ready_port(3004, &[3004]), Some(3004));
2563 assert_eq!(active_port_from_ready_port(4003, &[3003, 4003]), None);
2564 }
2565
2566 #[test]
2567 fn delay_readiness_requires_delay_only_and_a_running_process() {
2568 assert!(delay_readiness_succeeded(false, false, false, true));
2569 assert!(!delay_readiness_succeeded(false, true, false, true));
2570 assert!(!delay_readiness_succeeded(false, false, true, false));
2571 assert!(!delay_readiness_succeeded(false, false, false, false));
2572 assert!(!delay_readiness_succeeded(true, false, false, true));
2573 }
2574}
2575
2576fn is_daemon_slug_target(id: &DaemonId) -> bool {
2584 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
2588 slugs.iter().any(|(slug, entry)| {
2589 let daemon_name = entry.daemon.as_deref().unwrap_or(slug);
2590 id.name() == daemon_name
2591 })
2592}
2593
2594#[cfg(all(test, unix))]
2595mod tests {
2596 use super::*;
2597
2598 #[test]
2599 fn test_resolve_run_identity_empty_without_sudo() {
2600 let identity = resolve_run_identity(None, 501, 20, None, None).unwrap();
2601 assert_eq!(identity, RunIdentity::Inherit);
2602 }
2603
2604 #[test]
2605 fn test_resolve_run_identity_sudo_fallback() {
2606 let identity = resolve_run_identity(None, 0, 0, Some("501"), Some("20")).unwrap();
2607 let RunIdentity::Switch { uid, gid, .. } = identity else {
2608 panic!("expected identity switch");
2609 };
2610 assert_eq!(uid.as_raw(), 501);
2611 assert_eq!(gid.as_raw(), 20);
2612 }
2613
2614 #[test]
2615 fn test_resolve_run_identity_ignores_stale_sudo_when_not_root() {
2616 let identity = resolve_run_identity(None, 501, 20, Some("0"), Some("0")).unwrap();
2617 assert_eq!(identity, RunIdentity::Inherit);
2618 }
2619
2620 #[test]
2621 fn test_resolve_configured_user_root_name() {
2622 let identity = resolve_configured_user("root").unwrap();
2623 let RunIdentity::Switch { uid, username, .. } = identity else {
2624 panic!("expected identity switch");
2625 };
2626 assert_eq!(uid.as_raw(), 0);
2627 assert_eq!(
2628 username.as_deref().and_then(|s| s.to_str().ok()),
2629 Some("root")
2630 );
2631 }
2632
2633 #[test]
2634 fn test_resolve_configured_user_root_uid() {
2635 let identity = resolve_configured_user("0").unwrap();
2636 let RunIdentity::Switch { uid, username, .. } = identity else {
2637 panic!("expected identity switch");
2638 };
2639 assert_eq!(uid.as_raw(), 0);
2640 assert_eq!(
2641 username.as_deref().and_then(|s| s.to_str().ok()),
2642 Some("root")
2643 );
2644 }
2645
2646 #[test]
2647 fn test_resolve_configured_user_missing_user_fails() {
2648 let err = resolve_configured_user("pitchfork-user-that-should-not-exist")
2649 .unwrap_err()
2650 .to_string();
2651 assert!(err.contains("does not exist"));
2652 }
2653
2654 #[test]
2655 fn test_resolve_run_identity_requires_root_for_user_switch() {
2656 let err = resolve_run_identity(Some("root"), 501, 20, None, None)
2657 .unwrap_err()
2658 .to_string();
2659 assert!(err.contains("Restart the supervisor with sudo"));
2660 }
2661
2662 #[test]
2663 fn test_resolve_run_identity_same_user_is_noop() {
2664 let identity = resolve_run_identity(Some("root"), 0, 0, Some("501"), Some("20")).unwrap();
2665 assert_eq!(identity, RunIdentity::Inherit);
2666 }
2667}
2668
2669fn inject_proxy_env(cmd: &mut tokio::process::Command, slug: &Option<String>) {
2678 let s = crate::settings::settings();
2679 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
2680
2681 if should_force_loopback_host(slug) && !lan_enabled {
2682 cmd.env("HOST", "127.0.0.1");
2685 }
2686
2687 if let Some(url) = build_pitchfork_url(slug, &s) {
2689 cmd.env("PITCHFORK_URL", &url);
2690 }
2691
2692 if s.proxy.enable && s.proxy.https {
2694 let ca_path = if s.proxy.tls_cert.is_empty() {
2695 crate::env::PITCHFORK_STATE_DIR.join("proxy").join("ca.pem")
2696 } else {
2697 std::path::PathBuf::from(&s.proxy.tls_cert)
2698 };
2699 if ca_path.exists() {
2700 cmd.env("NODE_EXTRA_CA_CERTS", ca_path.to_string_lossy().to_string());
2701 }
2702 }
2703
2704 if s.proxy.enable {
2706 let tld = if lan_enabled { "local" } else { &s.proxy.tld };
2707 cmd.env("__VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS", format!(".{tld}"));
2708 }
2709
2710 if lan_enabled {
2712 cmd.env("PITCHFORK_LAN", "1");
2713 }
2714}
2715
2716fn should_force_loopback_host(slug: &Option<String>) -> bool {
2717 let Some(slug) = slug.as_deref() else {
2718 return false;
2719 };
2720
2721 let s = crate::settings::settings();
2722 if !s.proxy.enable {
2723 return false;
2724 }
2725
2726 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
2727 slugs.contains_key(slug)
2728}
2729
2730fn build_pitchfork_url(slug: &Option<String>, s: &crate::settings::Settings) -> Option<String> {
2734 let slug = slug.as_ref()?;
2735 if !s.proxy.enable {
2736 return None;
2737 }
2738 let scheme = if s.proxy.https { "https" } else { "http" };
2739 let port = u16::try_from(s.proxy.port).ok().filter(|&p| p > 0)?;
2740 let port_suffix = if (scheme == "https" && port == 443) || (scheme == "http" && port == 80) {
2741 String::new()
2742 } else {
2743 format!(":{port}")
2744 };
2745 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
2746 let tld = if lan_enabled { "local" } else { &s.proxy.tld };
2747 Some(format!("{scheme}://{slug}.{tld}{port_suffix}",))
2748}
2749
2750#[cfg(test)]
2751mod ready_check_tests {
2752 use super::*;
2753 use std::time::Duration;
2754
2755 #[test]
2756 fn any_ready_check_remaining_prefers_unbounded_checks() {
2757 let http = ReadyHttp::new("http://localhost/health");
2758 let cmd = ReadyCmd::new("true");
2759
2760 assert!(any_ready_check_remaining(
2761 None,
2762 false,
2763 None,
2764 false,
2765 Some(&http),
2766 false,
2767 None,
2768 false
2769 ));
2770 assert!(any_ready_check_remaining(
2771 None,
2772 false,
2773 None,
2774 false,
2775 None,
2776 false,
2777 Some(&cmd),
2778 false
2779 ));
2780 assert!(any_ready_check_remaining(
2781 None,
2782 false,
2783 Some(&ReadyPort::new(8080)),
2784 false,
2785 Some(&http),
2786 true,
2787 Some(&cmd),
2788 true
2789 ));
2790 }
2791
2792 #[test]
2793 fn any_ready_check_remaining_exhausted_timed_checks() {
2794 let http = ReadyHttp {
2795 url: "http://localhost/health".to_string(),
2796 status: vec![],
2797 timeout: Some(Duration::from_secs(5)),
2798 };
2799 let cmd = ReadyCmd {
2800 run: "true".to_string(),
2801 timeout: Some(Duration::from_secs(5)),
2802 };
2803
2804 assert!(any_ready_check_remaining(
2805 None,
2806 false,
2807 None,
2808 false,
2809 Some(&http),
2810 false,
2811 Some(&cmd),
2812 false
2813 ));
2814 assert!(!any_ready_check_remaining(
2815 None,
2816 false,
2817 None,
2818 false,
2819 Some(&http),
2820 true,
2821 Some(&cmd),
2822 true
2823 ));
2824 }
2825
2826 #[tokio::test]
2827 async fn spawn_cmd_probe_reports_success() {
2828 let id = DaemonId::new("global", "probe-test");
2829 let probe = spawn_cmd_probe(&id, "true", &std::env::temp_dir(), 0, None, &[]);
2830 let status = probe.result_rx.await.unwrap().unwrap();
2831 assert!(status.success());
2832 }
2833
2834 #[tokio::test]
2835 async fn spawn_cmd_probe_stops_on_request() {
2836 let id = DaemonId::new("global", "probe-test");
2837 let probe = spawn_cmd_probe(&id, "sleep 30", &std::env::temp_dir(), 0, None, &[]);
2838 let CmdProbe {
2839 cancel_tx,
2840 result_rx,
2841 } = probe;
2842 let _ = cancel_tx.send(());
2843 let status = result_rx.await.unwrap().unwrap();
2844 assert!(!status.success());
2845 }
2846
2847 #[tokio::test]
2848 async fn spawn_cmd_probe_receives_daemon_and_resolved_port_environment() {
2849 let id = DaemonId::new("worktree", "api");
2850 let daemon_env = IndexMap::from([("CUSTOM_VALUE".to_string(), "yes".to_string())]);
2851 let probe = spawn_cmd_probe(
2852 &id,
2853 r#"test "$CUSTOM_VALUE" = yes && test "$PORT" = 4100 && test "$PORT0" = 4100 && test "$PORT1" = 5100 && test "$PITCHFORK_DAEMON_ID" = worktree/api && test "$PITCHFORK_RETRY_COUNT" = 2"#,
2854 &std::env::temp_dir(),
2855 2,
2856 Some(&daemon_env),
2857 &[4100, 5100],
2858 );
2859 let status = probe.result_rx.await.unwrap().unwrap();
2860 assert!(status.success());
2861 }
2862
2863 #[test]
2864 fn configured_ready_port_follows_expected_port_bump() {
2865 assert_eq!(resolve_configured_ready_port(3000, &[3000], &[3004]), 3004);
2866 assert_eq!(
2867 resolve_configured_ready_port(4000, &[3000, 4000], &[3003, 4003]),
2868 4003
2869 );
2870 assert_eq!(resolve_configured_ready_port(8080, &[3000], &[3004]), 8080);
2871 }
2872}