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