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 miette::IntoDiagnostic;
21use once_cell::sync::Lazy;
22use regex::Regex;
23use std::collections::HashMap;
24#[cfg(unix)]
25use std::ffi::CString;
26use std::sync::{Arc, atomic};
27use std::time::Duration;
28use tokio::io::AsyncBufReadExt;
29use tokio::select;
30use tokio::sync::oneshot;
31use tokio::time;
32
33static REGEX_CACHE: Lazy<std::sync::Mutex<HashMap<String, Regex>>> =
35 Lazy::new(|| std::sync::Mutex::new(HashMap::new()));
36
37#[cfg(unix)]
38#[derive(Clone, Debug, PartialEq, Eq)]
39enum RunIdentity {
40 Inherit,
41 Switch {
42 uid: nix::unistd::Uid,
43 gid: nix::unistd::Gid,
44 username: Option<CString>,
45 },
46}
47
48pub(crate) fn get_or_compile_regex(pattern: &str) -> Option<Regex> {
50 let mut cache = REGEX_CACHE.lock().unwrap_or_else(|e| e.into_inner());
51 if let Some(re) = cache.get(pattern) {
52 return Some(re.clone());
53 }
54 match Regex::new(pattern) {
55 Ok(re) => {
56 cache.insert(pattern.to_string(), re.clone());
57 Some(re)
58 }
59 Err(e) => {
60 error!("invalid regex pattern '{pattern}': {e}");
61 None
62 }
63 }
64}
65
66struct CmdProbe {
73 cancel_tx: tokio::sync::oneshot::Sender<()>,
74 result_rx: tokio::sync::oneshot::Receiver<std::io::Result<std::process::ExitStatus>>,
75}
76
77fn spawn_cmd_probe(id: &DaemonId, cmd: &str, dir: &std::path::Path) -> CmdProbe {
84 let shell_setting = settings().general.shell.clone();
90 let mut command = match shell_words::split(&shell_setting) {
91 Ok(parts) if !parts.is_empty() => {
92 let (program, args) = parts.split_first().unwrap();
93 let mut c = tokio::process::Command::new(program);
94 c.args(args);
95 c.arg(cmd);
96 c
97 }
98 _ => Shell::default_for_platform().command(cmd),
99 };
100 command
101 .current_dir(dir)
102 .stdout(std::process::Stdio::null())
103 .stderr(std::process::Stdio::null())
104 .kill_on_drop(true);
105 let mut child = match command.spawn() {
106 Ok(child) => child,
107 Err(e) => {
108 warn!("daemon {id}: failed to spawn readiness command probe: {e}");
109 let (cancel_tx, _) = tokio::sync::oneshot::channel();
113 let (_, result_rx) = tokio::sync::oneshot::channel();
114 return CmdProbe {
115 cancel_tx,
116 result_rx,
117 };
118 }
119 };
120
121 let (cancel_tx, mut cancel_rx) = tokio::sync::oneshot::channel();
122 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
123
124 tokio::spawn(async move {
125 let status = tokio::select! {
126 status = child.wait() => status,
127 _ = &mut cancel_rx => {
128 let mut child = child;
129 let _ = child.kill().await;
130 child.wait().await
131 }
132 };
133 let _ = result_tx.send(status);
134 });
135
136 CmdProbe {
137 cancel_tx,
138 result_rx,
139 }
140}
141
142fn stop_cmd_probe_state(probe: &mut Option<CmdProbe>) {
144 if let Some(p) = probe.take() {
145 let _ = p.cancel_tx.send(());
146 }
147}
148
149#[allow(clippy::too_many_arguments)]
154fn any_ready_check_remaining(
155 ready_output: Option<&ReadyOutput>,
156 output_exhausted: bool,
157 ready_port: Option<&ReadyPort>,
158 port_exhausted: bool,
159 ready_http: Option<&ReadyHttp>,
160 http_exhausted: bool,
161 ready_cmd: Option<&ReadyCmd>,
162 cmd_exhausted: bool,
163) -> bool {
164 ready_output.is_some_and(|o| o.timeout.is_none() || !output_exhausted)
165 || ready_port.is_some_and(|p| p.timeout.is_none() || !port_exhausted)
166 || ready_http.is_some_and(|h| h.timeout.is_none() || !http_exhausted)
167 || ready_cmd.is_some_and(|c| c.timeout.is_none() || !cmd_exhausted)
168}
169
170impl Supervisor {
171 pub async fn run(&self, opts: RunOptions) -> Result<IpcResponse> {
173 let id = &opts.id;
174 let cmd = opts.cmd.clone();
175
176 {
178 let mut pending = self.pending_autostops.lock().await;
179 if pending.remove(id).is_some() {
180 info!("cleared pending autostop for {id} (daemon starting)");
181 }
182 }
183
184 let daemon = self.get_daemon(id).await;
185 if let Some(daemon) = daemon {
186 if !daemon.status.is_stopping()
189 && !daemon.status.is_stopped()
190 && let Some(pid) = daemon.pid
191 {
192 if opts.force {
193 self.stop(id).await?;
194 info!("run: stop completed for daemon {id}");
195 } else {
196 warn!("daemon {id} already running with pid {pid}");
197 return Ok(IpcResponse::DaemonAlreadyRunning);
198 }
199 }
200 }
201
202 if opts.wait_ready && opts.retry.count() > 0 {
204 let max_attempts = opts.retry.count().saturating_add(1);
206 for attempt in 0..max_attempts {
207 let mut retry_opts = opts.clone();
208 retry_opts.retry_count = attempt;
209 retry_opts.cmd = cmd.clone();
210
211 let result = self.run_once(retry_opts).await?;
212
213 match result {
214 IpcResponse::DaemonReady { daemon } => {
215 return Ok(IpcResponse::DaemonReady { daemon });
216 }
217 IpcResponse::DaemonFailedWithCode { exit_code } => {
218 if attempt < opts.retry.count() {
219 let backoff_secs = 2u64.saturating_pow(attempt).min(3600);
220 info!(
221 "daemon {id} failed (attempt {}/{}), retrying in {}s",
222 attempt + 1,
223 max_attempts,
224 backoff_secs
225 );
226 fire_hook(
227 HookType::OnRetry,
228 id.clone(),
229 opts.dir.0.clone(),
230 attempt + 1,
231 opts.env.clone(),
232 vec![],
233 )
234 .await;
235 time::sleep(Duration::from_secs(backoff_secs)).await;
236 continue;
237 } else {
238 info!("daemon {id} failed after {max_attempts} attempts");
239 return Ok(IpcResponse::DaemonFailedWithCode { exit_code });
240 }
241 }
242 other => return Ok(other),
243 }
244 }
245 }
246
247 self.run_once(opts).await
249 }
250
251 pub(crate) async fn run_once(&self, opts: RunOptions) -> Result<IpcResponse> {
253 let id = &opts.id;
254 let original_cmd = opts.cmd.clone(); let (ready_tx, ready_rx) = if opts.wait_ready {
258 let (tx, rx) = oneshot::channel();
259 (Some(tx), Some(rx))
260 } else {
261 (None, None)
262 };
263
264 let expected_ports = opts
266 .port
267 .as_ref()
268 .map(|p| p.expect.clone())
269 .unwrap_or_default();
270 let (resolved_ports, effective_ready_port) = if !expected_ports.is_empty() {
271 let port_cfg = opts.port.as_ref().unwrap();
272 match check_ports_available(
273 &expected_ports,
274 port_cfg.auto_bump(),
275 port_cfg.max_bump_attempts(),
276 )
277 .await
278 {
279 Ok(resolved) => {
280 let ready_port = if let Some(configured_port) =
281 opts.ready_port.as_ref().and_then(|p| p.as_port())
282 {
283 let bump_offset = resolved
285 .first()
286 .unwrap_or(&0)
287 .saturating_sub(*expected_ports.first().unwrap_or(&0));
288 if expected_ports.contains(&configured_port) && bump_offset > 0 {
289 configured_port
290 .checked_add(bump_offset)
291 .or(Some(configured_port))
292 } else {
293 Some(configured_port)
294 }
295 } else if opts.ready_output.is_none()
296 && opts.ready_http.is_none()
297 && opts.ready_cmd.is_none()
298 && opts.ready_delay.is_none()
299 {
300 resolved.first().copied().filter(|&p| p != 0)
304 } else {
305 None
309 };
310 info!("daemon {id}: ports {expected_ports:?} resolved to {resolved:?}");
311 (resolved, ready_port)
312 }
313 Err(e) => {
314 error!("daemon {id}: port check failed: {e}");
315 if let Some(port_error) = e.downcast_ref::<PortError>() {
317 match port_error {
318 PortError::InUse { port, process, pid } => {
319 return Ok(IpcResponse::PortConflict {
320 port: *port,
321 process: process.clone(),
322 pid: *pid,
323 });
324 }
325 PortError::NoAvailablePort {
326 start_port,
327 attempts,
328 } => {
329 return Ok(IpcResponse::NoAvailablePort {
330 start_port: *start_port,
331 attempts: *attempts,
332 });
333 }
334 }
335 }
336 return Ok(IpcResponse::DaemonFailed {
337 error: e.to_string(),
338 });
339 }
340 }
341 } else {
342 if let Some(port) = opts.ready_port.as_ref().and_then(|p| p.as_port())
348 && port > 0
349 && let Some((pid, process)) = detect_port_conflict(port).await
350 {
351 return Ok(IpcResponse::PortConflict { port, process, pid });
352 }
353 (
354 Vec::new(),
355 opts.ready_port.as_ref().and_then(|p| p.as_port()),
356 )
357 };
358
359 let shell_setting = settings().general.shell.clone();
363 let shell_parts = match shell_words::split(&shell_setting) {
364 Ok(parts) if !parts.is_empty() => parts,
365 Ok(_) => {
366 return Ok(IpcResponse::DaemonFailed {
367 error: "general.shell setting is empty".to_string(),
368 });
369 }
370 Err(e) => {
371 return Ok(IpcResponse::DaemonFailed {
372 error: format!("failed to parse general.shell setting {shell_setting:?}: {e}"),
373 });
374 }
375 };
376 let (shell_program, shell_args) = shell_parts.split_first().unwrap();
377
378 let run_script = opts
383 .run
384 .clone()
385 .unwrap_or_else(|| shell_words::join(&original_cmd));
386
387 let (program, args) = if opts.mise.unwrap_or(settings().general.mise) {
388 match settings().resolve_mise_bin() {
389 Some(mise_bin) => {
390 let mise_bin_str = mise_bin.to_string_lossy().to_string();
391 info!("daemon {id}: wrapping command with mise ({mise_bin_str})");
392 let mut args = vec!["x".to_string(), "--".to_string()];
393 args.push(shell_program.clone());
394 args.extend(shell_args.iter().cloned());
395 args.push(run_script);
396 (mise_bin_str, args)
397 }
398 None => {
399 warn!("daemon {id}: mise=true but mise binary not found, running without mise");
400 let mut args: Vec<String> = shell_args.to_vec();
401 args.push(run_script);
402 (shell_program.clone(), args)
403 }
404 }
405 } else {
406 let mut args: Vec<String> = shell_args.to_vec();
407 args.push(run_script);
408 (shell_program.clone(), args)
409 };
410 #[cfg(unix)]
411 let run_identity = match resolve_effective_run_identity(opts.user.as_deref()) {
412 Ok(identity) => identity,
413 Err(e) => {
414 return Ok(IpcResponse::DaemonFailed {
415 error: e.to_string(),
416 });
417 }
418 };
419 info!("run: spawning daemon {id} with {program} {args:?}");
420
421 #[cfg(unix)]
423 let pty_pair = if opts.pty.unwrap_or(false) {
424 match super::pty::openpty() {
425 Ok(pair) => {
426 info!("daemon {id}: allocated PTY (pty = true)");
427 Some(pair)
428 }
429 Err(e) => {
430 warn!("daemon {id}: failed to allocate PTY, falling back to pipes: {e}");
431 None
432 }
433 }
434 } else {
435 None
436 };
437
438 let mut cmd = tokio::process::Command::new(&program);
439
440 #[cfg(unix)]
441 if let Some(ref pair) = pty_pair {
442 let slave_file = std::fs::File::from(
446 pair.slave
447 .try_clone()
448 .map_err(|e| miette::miette!("failed to dup slave PTY fd: {e}"))?,
449 );
450 cmd.stdin(std::process::Stdio::from(slave_file.try_clone().map_err(
451 |e| miette::miette!("failed to clone slave PTY fd for stdin: {e}"),
452 )?));
453 cmd.stdout(std::process::Stdio::from(slave_file.try_clone().map_err(
454 |e| miette::miette!("failed to clone slave PTY fd for stdout: {e}"),
455 )?));
456 cmd.stderr(std::process::Stdio::from(slave_file));
457 } else {
458 cmd.stdout(std::process::Stdio::piped())
459 .stderr(std::process::Stdio::piped());
460 }
461
462 #[cfg(not(unix))]
463 {
464 cmd.stdout(std::process::Stdio::piped())
465 .stderr(std::process::Stdio::piped());
466 }
467
468 cmd.args(&args).current_dir(&opts.dir);
469
470 #[cfg(unix)]
471 if pty_pair.is_none() {
472 cmd.stdin(std::process::Stdio::null());
473 }
474
475 #[cfg(not(unix))]
476 cmd.stdin(std::process::Stdio::null());
477
478 if let Some(ref path) = *env::ORIGINAL_PATH {
480 cmd.env("PATH", path);
481 }
482
483 if let Some(ref env_vars) = opts.env {
485 cmd.envs(env_vars);
486 }
487
488 cmd.env("PITCHFORK_DAEMON_ID", id.qualified());
490 cmd.env("PITCHFORK_DAEMON_NAMESPACE", id.namespace());
491 cmd.env("PITCHFORK_RETRY_COUNT", opts.retry_count.to_string());
492
493 if !resolved_ports.is_empty() {
495 cmd.env("PORT", resolved_ports[0].to_string());
499 for (i, port) in resolved_ports.iter().enumerate() {
501 cmd.env(format!("PORT{i}"), port.to_string());
502 }
503 }
504
505 inject_proxy_env(&mut cmd, &opts.slug);
507
508 #[cfg(unix)]
509 {
510 let run_identity = run_identity.clone();
511 let use_pty = pty_pair.is_some();
512 unsafe {
513 cmd.pre_exec(move || {
514 nix::unistd::setsid().map_err(nix_to_io_error)?;
515
516 if use_pty {
520 let ret = libc::ioctl(0, libc::TIOCSCTTY as libc::c_ulong, 0);
521 if ret < 0 {
522 #[cfg(target_os = "linux")]
525 eprintln!(
526 "pitchfork: TIOCSCTTY failed: {}",
527 std::io::Error::last_os_error()
528 );
529 }
530 }
531
532 apply_run_identity(&run_identity)?;
533 Ok(())
534 });
535 }
536 }
537
538 let mut child = cmd.spawn().into_diagnostic()?;
539 let pid = match child.id() {
540 Some(p) => p,
541 None => {
542 warn!("Daemon {id} exited before PID could be captured");
543 return Ok(IpcResponse::DaemonFailed {
544 error: "Process exited immediately".to_string(),
545 });
546 }
547 };
548 info!("started daemon {id} with pid {pid}");
549 PROCS.refresh_pids(&[pid]);
550 let monitored_guard = super::adopt::MonitoredGuard::register(id.clone(), pid);
558 let monitor_token = monitored_guard.token();
559 let daemon = self
560 .upsert_daemon(
561 UpsertDaemonOpts::from_run_options(&opts, DaemonStatus::Running)
562 .set(|o| {
563 o.pid = Some(pid);
564 o.cmd = Some(original_cmd);
565 o.ready_port = effective_ready_port.map(|p| ReadyPort {
566 port: Some(p),
567 template: None,
568 timeout: opts.ready_port.as_ref().and_then(|rp| rp.timeout),
569 });
570 o.port = crate::config_types::PortConfig::from_parts(
571 expected_ports,
572 opts.port.as_ref().map(|p| p.bump).unwrap_or_default(),
573 );
574 o.resolved_port = resolved_ports;
575 })
576 .build(),
577 )
578 .await?;
579
580 let id_clone = id.clone();
581 let ready_delay = opts.ready_delay;
582 let ready_output = opts.ready_output.clone();
583 let ready_http = opts.ready_http.clone();
584 let ready_port = effective_ready_port;
585 let implicit_ready_port = ready_port.map(|p| ReadyPort {
586 port: Some(p),
587 template: None,
588 timeout: None,
589 });
590 let ready_port_config = opts.ready_port.clone().or(implicit_ready_port);
591 let ready_cmd = opts.ready_cmd.clone();
592 let daemon_dir = opts.dir.0.clone();
593 let hook_retry_count = opts.retry_count;
594 let hook_retry = opts.retry;
595 let hook_daemon_env = opts.env.clone();
596 let on_output_hook = opts.on_output_hook.clone();
597 let has_port_config = opts.port.as_ref().is_some_and(|p| !p.expect.is_empty())
603 || (settings().proxy.enable && is_daemon_slug_target(id));
604 let expected_port: Option<u16> = opts.port.as_ref().and_then(|p| p.expect.first().copied());
610 let daemon_pid = pid;
611
612 #[cfg(unix)]
616 let pty_reader = pty_pair.map(|p| {
617 tokio::io::BufReader::new(tokio::fs::File::from_std(std::fs::File::from(p.master)))
618 .lines()
619 });
620 #[cfg(not(unix))]
621 let pty_reader: Option<tokio::io::Lines<tokio::io::BufReader<tokio::fs::File>>> = None;
622 let stdout_reader = if pty_reader.is_none() {
623 child
624 .stdout
625 .take()
626 .map(|s| tokio::io::BufReader::new(s).lines())
627 } else {
628 None
629 };
630 let stderr_reader = if pty_reader.is_none() {
631 child
632 .stderr
633 .take()
634 .map(|s| tokio::io::BufReader::new(s).lines())
635 } else {
636 None
637 };
638
639 if pty_reader.is_none() && (stdout_reader.is_none() || stderr_reader.is_none()) {
640 error!("Failed to capture stdout/stderr for daemon {id}");
641 }
642
643 tokio::spawn(async move {
644 let id = id_clone;
645 let _monitored_guard = monitored_guard;
648
649 let (output_tx, mut output_rx) = tokio::sync::mpsc::channel::<String>(256);
651
652 if let Some(mut reader) = pty_reader {
653 tokio::spawn(async move {
657 while let Ok(Some(mut line)) = reader.next_line().await {
658 if line.ends_with('\r') {
660 line.pop();
661 }
662 if output_tx.send(line).await.is_err() {
663 break;
664 }
665 }
666 });
667 } else {
668 if let Some(mut stdout) = stdout_reader {
673 let tx = output_tx.clone();
674 tokio::spawn(async move {
675 while let Ok(Some(line)) = stdout.next_line().await {
676 if tx.send(line).await.is_err() {
677 break;
678 }
679 }
680 });
681 }
682 if let Some(mut stderr) = stderr_reader {
683 let tx = output_tx.clone();
684 tokio::spawn(async move {
685 while let Ok(Some(line)) = stderr.next_line().await {
686 if tx.send(line).await.is_err() {
687 break;
688 }
689 }
690 });
691 }
692 drop(output_tx);
694 }
695 let log_store = Arc::clone(&LOG_STORE);
696 let log_format = opts
697 .log_format
698 .clone()
699 .unwrap_or_else(|| crate::settings::settings().logs.log_format.clone());
700 let parse_line = move |line: &str| crate::log_parse::parse(line, &log_format);
701
702 const LOG_BATCH_SIZE: usize = 100;
703 const LOG_FLUSH_INTERVAL: Duration = Duration::from_millis(100);
704 let mut log_buffer: Vec<crate::log_parse::ParsedLog> =
705 Vec::with_capacity(LOG_BATCH_SIZE);
706 let mut log_flush_interval = tokio::time::interval(LOG_FLUSH_INTERVAL);
707 log_flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
708
709 let flush_logs =
710 |buffer: &mut Vec<crate::log_parse::ParsedLog>| -> Option<tokio::task::JoinHandle<()>> {
711 if buffer.is_empty() {
712 return None;
713 }
714 let store = Arc::clone(&log_store);
715 let id = id.clone();
716 let batch = std::mem::take(buffer);
717 Some(tokio::task::spawn_blocking(move || {
718 if let Err(e) = store.append_structured_batch(&id, &batch) {
719 error!("Failed to write batch to log for daemon {id}: {e}");
720 }
721 }))
722 };
723
724 let mut ready_notified = false;
728 let mut ready_tx = ready_tx;
729 let ready_pattern = ready_output
730 .as_ref()
731 .and_then(|o| get_or_compile_regex(&o.pattern));
732 let mut active_port_spawned = false;
734
735 let on_output_hook = match on_output_hook {
739 Some(ref hook) => match hook.validate(id.name()) {
740 Ok(()) => on_output_hook,
741 Err(e) => {
742 error!("{e}");
743 None
744 }
745 },
746 None => None,
747 };
748
749 let on_output_pattern: Option<regex::Regex> = on_output_hook
752 .as_ref()
753 .and_then(|h| h.regex.as_deref().and_then(get_or_compile_regex));
754 let on_output_debounce = on_output_hook
755 .as_ref()
756 .map(|h| h.debounce_duration())
757 .unwrap_or(Duration::from_millis(1000));
758 let mut on_output_last_fired: Option<std::time::Instant> = None;
760
761 let mut delay_timer =
762 ready_delay.map(|secs| Box::pin(time::sleep(Duration::from_secs(secs))));
763
764 let mut http_exhausted = false;
766 let mut cmd_exhausted = false;
767 let mut port_exhausted = false;
768 let mut output_exhausted = false;
769
770 let s = settings();
772 let ready_check_interval = s.supervisor_ready_check_interval();
773 let http_client_timeout = s.supervisor_http_client_timeout();
774
775 let mut output_deadline = ready_output
777 .as_ref()
778 .and_then(|o| o.timeout)
779 .map(|d| Box::pin(time::sleep(d)));
780
781 let mut http_check_interval = ready_http
783 .as_ref()
784 .map(|_| tokio::time::interval(ready_check_interval));
785 let mut http_deadline = ready_http
786 .as_ref()
787 .and_then(|h| h.timeout)
788 .map(|d| Box::pin(time::sleep(d)));
789 let http_client = ready_http.as_ref().map(|_| {
790 reqwest::Client::builder()
791 .timeout(http_client_timeout)
792 .build()
793 .unwrap_or_default()
794 });
795
796 let mut port_check_interval =
798 ready_port.map(|_| tokio::time::interval(ready_check_interval));
799 let mut port_deadline = ready_port_config
800 .as_ref()
801 .and_then(|p| p.timeout)
802 .map(|d| Box::pin(time::sleep(d)));
803
804 let mut cmd_probe: Option<CmdProbe> = None;
807 let mut cmd_respawn_delay: Option<_> = None;
808 let mut cmd_deadline = ready_cmd
809 .as_ref()
810 .and_then(|c| c.timeout)
811 .map(|d| Box::pin(time::sleep(d)));
812 if let Some(ref cmd) = ready_cmd {
813 cmd_probe = Some(spawn_cmd_probe(&id, &cmd.run, daemon_dir.as_path()));
814 }
815
816 let (exit_tx, mut exit_rx) =
818 tokio::sync::mpsc::channel::<std::io::Result<std::process::ExitStatus>>(1);
819
820 let child_pid = child.id().unwrap_or(0);
822 tokio::spawn(async move {
823 let result = child.wait().await;
824 #[cfg(all(unix, not(target_os = "linux")))]
834 let result = match &result {
835 Err(e) if e.raw_os_error() == Some(nix::libc::ECHILD) => {
836 if let Some(code) = super::REAPED_STATUSES.lock().await.remove(&child_pid) {
837 warn!(
838 "daemon pid {child_pid} wait() got ECHILD; \
839 recovered exit code {code} from zombie reaper"
840 );
841 use std::os::unix::process::ExitStatusExt;
846 if code >= 0 {
847 Ok(std::process::ExitStatus::from_raw(code << 8))
848 } else {
849 Ok(std::process::ExitStatus::from_raw((-code) & 0x7f))
851 }
852 } else {
853 warn!(
854 "daemon pid {child_pid} wait() got ECHILD but no \
855 stashed status found; reporting as error"
856 );
857 result
858 }
859 }
860 _ => result,
861 };
862 debug!("daemon pid {child_pid} wait() completed with result: {result:?}");
863 let _ = exit_tx.send(result).await;
864 });
865
866 #[allow(unused_assignments)]
867 let mut exit_status = None;
869
870 if has_port_config
876 && ready_pattern.is_none()
877 && ready_http.is_none()
878 && ready_port.is_none()
879 && ready_cmd.is_none()
880 && delay_timer.is_none()
881 {
882 active_port_spawned = true;
883 detect_and_store_active_port(id.clone(), daemon_pid);
884 }
885
886 loop {
887 select! {
891 biased;
892 Some(result) = exit_rx.recv() => {
893 exit_status = Some(result);
895 debug!("daemon {id} process exited, exit_status: {exit_status:?}");
896 if !ready_notified {
897 if let Some(tx) = ready_tx.take() {
898 let is_success = exit_status.as_ref()
900 .and_then(|r| r.as_ref().ok())
901 .map(|s| s.success())
902 .unwrap_or(false);
903
904 if is_success {
905 debug!("daemon {id} exited successfully before ready check, sending success notification");
906 let _ = tx.send(Ok(()));
907 } else {
908 let exit_code = exit_status.as_ref()
909 .and_then(|r| r.as_ref().ok())
910 .and_then(|s| s.code());
911 debug!("daemon {id} exited with failure before ready check, sending failure notification with exit_code: {exit_code:?}");
912 let _ = tx.send(Err(exit_code));
913 }
914 }
915 } else {
916 debug!("daemon {id} was already marked ready, not sending notification");
917 }
918 break;
919 },
920 Some(line) = output_rx.recv() => {
921 let parsed = parse_line(&line);
922 log_buffer.push(parsed);
923 if log_buffer.len() >= LOG_BATCH_SIZE {
924 let _ = flush_logs(&mut log_buffer);
925 }
926 trace!("output: {id} {line}");
927
928 let line_clean = console::strip_ansi_codes(&line).to_string();
931
932 if !ready_notified
934 && !output_exhausted
935 && let Some(ref pattern) = ready_pattern
936 && pattern.is_match(&line_clean)
937 {
938 if let Some(handle) = flush_logs(&mut log_buffer) {
943 let _ = handle.await;
944 }
945 info!("daemon {id} ready: output matched pattern");
946 ready_notified = true;
947 if let Some(tx) = ready_tx.take() {
948 let _ = tx.send(Ok(()));
949 }
950 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
951 stop_cmd_probe_state(&mut cmd_probe);
952 http_deadline = None;
953 cmd_deadline = None;
954 port_deadline = None;
955 output_deadline = None;
956 if !active_port_spawned && has_port_config {
957 active_port_spawned = true;
958 detect_and_store_active_port(id.clone(), daemon_pid);
959 }
960 }
961
962 if let Some(ref hook) = on_output_hook {
964 let matched = match (&hook.filter, &on_output_pattern) {
965 (Some(substr), _) => line_clean.contains(substr.as_str()),
966 (None, Some(re)) => re.is_match(&line_clean),
967 (None, None) => true,
968 };
969 if matched {
970 let now = std::time::Instant::now();
971 let elapsed = on_output_last_fired.map(|t| now.duration_since(t));
972 if elapsed.is_none_or(|e| e >= on_output_debounce) {
973 on_output_last_fired = Some(now);
974 hooks::fire_output_hook(id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook.run.clone(), line_clean.clone()).await;
975 }
976 }
977 }
978 tokio::task::yield_now().await;
981 }
982 _ = async {
983 if let Some(ref mut deadline) = http_deadline {
984 deadline.await;
985 } else {
986 std::future::pending::<()>().await;
987 }
988 }, if !ready_notified && ready_http.is_some() => {
989 http_exhausted = true;
990 http_deadline = None;
991 http_check_interval = None;
992 warn!("daemon {id}: HTTP readiness check timed out");
993 let any_remaining = any_ready_check_remaining(
994 ready_output.as_ref(),
995 output_exhausted,
996 ready_port_config.as_ref(),
997 port_exhausted,
998 ready_http.as_ref(),
999 http_exhausted,
1000 ready_cmd.as_ref(),
1001 cmd_exhausted,
1002 );
1003 if !any_remaining {
1004 error!("daemon {id}: all readiness checks exhausted, failing");
1005 stop_cmd_probe_state(&mut cmd_probe);
1006 if let Some(tx) = ready_tx.take() {
1007 let _ = tx.send(Err(Some(124)));
1008 }
1009 let stop_cfg = opts.stop_signal.unwrap_or_default();
1010 let _ = PROCS.kill_process_group_async(daemon_pid, stop_cfg.signal.into(), stop_cfg.timeout).await;
1011 break;
1012 }
1013 }
1014 _ = async {
1015 if let Some(ref mut deadline) = output_deadline {
1016 deadline.await;
1017 } else {
1018 std::future::pending::<()>().await;
1019 }
1020 }, if !ready_notified && ready_output.is_some() => {
1021 output_exhausted = true;
1022 output_deadline = None;
1023 warn!("daemon {id}: output readiness check timed out");
1024 let any_remaining = any_ready_check_remaining(
1025 ready_output.as_ref(),
1026 output_exhausted,
1027 ready_port_config.as_ref(),
1028 port_exhausted,
1029 ready_http.as_ref(),
1030 http_exhausted,
1031 ready_cmd.as_ref(),
1032 cmd_exhausted,
1033 );
1034 if !any_remaining {
1035 error!("daemon {id}: all readiness checks exhausted, failing");
1036 stop_cmd_probe_state(&mut cmd_probe);
1037 if let Some(tx) = ready_tx.take() {
1038 let _ = tx.send(Err(Some(124)));
1039 }
1040 let stop_cfg = opts.stop_signal.unwrap_or_default();
1041 let _ = PROCS.kill_process_group_async(daemon_pid, stop_cfg.signal.into(), stop_cfg.timeout).await;
1042 break;
1043 }
1044 }
1045 _ = async {
1046 if let Some(ref mut interval) = http_check_interval {
1047 interval.tick().await;
1048 } else {
1049 std::future::pending::<()>().await;
1050 }
1051 }, if !ready_notified && ready_http.is_some() && !http_exhausted => {
1052 if let (Some(http), Some(client)) = (&ready_http, &http_client) {
1053 match client.get(&http.url).send().await {
1054 Ok(response) if http.accepts_status(response.status().as_u16()) => {
1055 info!("daemon {id} ready: HTTP check passed (status {})", response.status());
1056 ready_notified = true;
1057 if let Some(tx) = ready_tx.take() {
1058 let _ = tx.send(Ok(()));
1059 }
1060 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1061 http_check_interval = None;
1062 http_deadline = None;
1063 stop_cmd_probe_state(&mut cmd_probe);
1064 cmd_deadline = None;
1065 port_deadline = None;
1066 output_deadline = None;
1067 if !active_port_spawned && has_port_config {
1068 active_port_spawned = true;
1069 detect_and_store_active_port(id.clone(), daemon_pid);
1070 }
1071 }
1072 Ok(response) => {
1073 trace!("daemon {id} HTTP check: status {} (not ready)", response.status());
1074 }
1075 Err(e) => {
1076 trace!("daemon {id} HTTP check failed: {e}");
1077 }
1078 }
1079 }
1080 }
1081 _ = async {
1082 if let Some(ref mut deadline) = port_deadline {
1083 deadline.await;
1084 } else {
1085 std::future::pending::<()>().await;
1086 }
1087 }, if !ready_notified && ready_port.is_some() => {
1088 port_exhausted = true;
1089 port_deadline = None;
1090 port_check_interval = None;
1091 warn!("daemon {id}: TCP port readiness check timed out");
1092 let any_remaining = any_ready_check_remaining(
1093 ready_output.as_ref(),
1094 output_exhausted,
1095 ready_port_config.as_ref(),
1096 port_exhausted,
1097 ready_http.as_ref(),
1098 http_exhausted,
1099 ready_cmd.as_ref(),
1100 cmd_exhausted,
1101 );
1102 if !any_remaining {
1103 error!("daemon {id}: all readiness checks exhausted, failing");
1104 stop_cmd_probe_state(&mut cmd_probe);
1105 if let Some(tx) = ready_tx.take() {
1106 let _ = tx.send(Err(Some(124)));
1107 }
1108 let stop_cfg = opts.stop_signal.unwrap_or_default();
1109 let _ = PROCS.kill_process_group_async(daemon_pid, stop_cfg.signal.into(), stop_cfg.timeout).await;
1110 break;
1111 }
1112 }
1113 _ = async {
1114 if let Some(ref mut interval) = port_check_interval {
1115 interval.tick().await;
1116 } else {
1117 std::future::pending::<()>().await;
1118 }
1119 }, if !ready_notified && ready_port.is_some() && !port_exhausted => {
1120 if let Some(port) = ready_port {
1121 match tokio::net::TcpStream::connect(("127.0.0.1", port)).await {
1122 Ok(_) => {
1123 info!("daemon {id} ready: TCP port {port} is listening");
1124 ready_notified = true;
1125 if let Some(tx) = ready_tx.take() {
1126 let _ = tx.send(Ok(()));
1127 }
1128 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1129 port_check_interval = None;
1131 port_deadline = None;
1132 stop_cmd_probe_state(&mut cmd_probe);
1133 http_deadline = None;
1134 cmd_deadline = None;
1135 output_deadline = None;
1136 if !active_port_spawned && has_port_config {
1137 active_port_spawned = true;
1138 if expected_port == Some(port) {
1148 let mut state_file =
1149 SUPERVISOR.state_file.lock().await;
1150 if let Some(d) = state_file.daemons.get(&id)
1151 && d.pid == Some(daemon_pid)
1152 {
1153 state_file.set_active_port(&id, port);
1154 }
1155 } else {
1156 detect_and_store_active_port(
1157 id.clone(),
1158 daemon_pid,
1159 );
1160 }
1161 }
1162 }
1163 Err(_) => {
1164 trace!("daemon {id} port check: port {port} not listening yet");
1165 }
1166 }
1167 }
1168 }
1169 _ = async {
1170 if let Some(ref mut delay) = cmd_respawn_delay {
1171 delay.await;
1172 } else {
1173 std::future::pending::<()>().await;
1174 }
1175 }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted && cmd_probe.is_none() => {
1176 if let Some(ref cmd) = ready_cmd {
1177 cmd_probe = Some(spawn_cmd_probe(&id, &cmd.run, daemon_dir.as_path()));
1178 }
1179 cmd_respawn_delay = None;
1180 }
1181 result = async {
1182 if let Some(probe) = cmd_probe.as_mut() {
1183 std::pin::Pin::new(&mut probe.result_rx).await
1184 } else {
1185 std::future::pending::<Result<Result<std::process::ExitStatus, std::io::Error>, tokio::sync::oneshot::error::RecvError>>().await
1186 }
1187 }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted => {
1188 let _ = cmd_probe.take();
1192 match result {
1193 Ok(Ok(status)) if status.success() => {
1194 info!("daemon {id} ready: readiness command succeeded");
1195 ready_notified = true;
1196 if let Some(tx) = ready_tx.take() {
1197 let _ = tx.send(Ok(()));
1198 }
1199 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1200 cmd_respawn_delay = None;
1201 cmd_deadline = None;
1202 http_deadline = None;
1203 port_deadline = None;
1204 output_deadline = None;
1205 if !active_port_spawned && has_port_config {
1206 active_port_spawned = true;
1207 detect_and_store_active_port(id.clone(), daemon_pid);
1208 }
1209 }
1210 Ok(Ok(_)) | Ok(Err(_)) | Err(_) => {
1211 trace!("daemon {id} cmd check: command not ready, will respawn");
1212 cmd_respawn_delay = Some(Box::pin(time::sleep(ready_check_interval)));
1213 }
1214 }
1215 }
1216 _ = async {
1217 if let Some(ref mut deadline) = cmd_deadline {
1218 deadline.await;
1219 } else {
1220 std::future::pending::<()>().await;
1221 }
1222 }, if !ready_notified && ready_cmd.is_some() => {
1223 cmd_exhausted = true;
1224 cmd_deadline = None;
1225 stop_cmd_probe_state(&mut cmd_probe);
1226 cmd_respawn_delay = None;
1227 warn!("daemon {id}: command readiness check timed out");
1228 let any_remaining = any_ready_check_remaining(
1229 ready_output.as_ref(),
1230 output_exhausted,
1231 ready_port_config.as_ref(),
1232 port_exhausted,
1233 ready_http.as_ref(),
1234 http_exhausted,
1235 ready_cmd.as_ref(),
1236 cmd_exhausted,
1237 );
1238 if !any_remaining {
1239 error!("daemon {id}: all readiness checks exhausted, failing");
1240 if let Some(tx) = ready_tx.take() {
1241 let _ = tx.send(Err(Some(124)));
1242 }
1243 let stop_cfg = opts.stop_signal.unwrap_or_default();
1244 let _ = PROCS.kill_process_group_async(daemon_pid, stop_cfg.signal.into(), stop_cfg.timeout).await;
1245 break;
1246 }
1247 }
1248 _ = async {
1249 if let Some(ref mut timer) = delay_timer {
1250 timer.await;
1251 } else {
1252 std::future::pending::<()>().await;
1253 }
1254 } => {
1255 if !ready_notified && ready_pattern.is_none() && ready_http.is_none() && ready_port.is_none() && ready_cmd.is_none() {
1256 if exit_status.is_some() {
1261 debug!("daemon {id} exited during ready_delay, not marking as ready");
1262 } else {
1263 PROCS.refresh_pids(&[daemon_pid]);
1266 if !PROCS.is_running(daemon_pid) {
1267 debug!("daemon {id} pid {daemon_pid} not running during ready_delay, deferring to exit handler");
1268 } else {
1269 info!("daemon {id} ready: delay elapsed");
1270 ready_notified = true;
1271 if let Some(tx) = ready_tx.take() {
1272 let _ = tx.send(Ok(()));
1273 }
1274 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1275 }
1276 }
1277 output_deadline = None;
1280 http_deadline = None;
1281 cmd_deadline = None;
1282 port_deadline = None;
1283 stop_cmd_probe_state(&mut cmd_probe);
1284 }
1285 delay_timer = None;
1287 if !active_port_spawned && has_port_config {
1288 active_port_spawned = true;
1289 detect_and_store_active_port(id.clone(), daemon_pid);
1290 }
1291 }
1292 _ = log_flush_interval.tick() => {
1293 let _ = flush_logs(&mut log_buffer);
1294 }
1295 }
1296 }
1297
1298 let pre_drain_daemon = SUPERVISOR.get_daemon(&id).await;
1312 let pre_drain_is_stopping = pre_drain_daemon
1313 .as_ref()
1314 .is_some_and(|d| d.status.is_stopped() || d.status.is_stopping());
1315
1316 let drain_deadline = tokio::time::Instant::now() + Duration::from_secs(5);
1324 loop {
1325 let now = tokio::time::Instant::now();
1326 if now >= drain_deadline {
1327 break;
1328 }
1329 let Ok(Some(line)) =
1330 tokio::time::timeout(drain_deadline - now, output_rx.recv()).await
1331 else {
1332 break;
1333 };
1334 log_buffer.push(parse_line(&line));
1335 }
1336 if let Some(handle) = flush_logs(&mut log_buffer) {
1339 let _ = handle.await;
1340 }
1341
1342 {
1344 let mut state_file = SUPERVISOR.state_file.lock().await;
1345 state_file.clear_active_port(&id);
1346 }
1347
1348 let exit_status = if let Some(status) = exit_status {
1350 status
1351 } else {
1352 match exit_rx.recv().await {
1354 Some(status) => status,
1355 None => {
1356 warn!("daemon {id} exit channel closed without receiving status");
1357 Err(std::io::Error::other("exit channel closed"))
1358 }
1359 }
1360 };
1361 let current_daemon = SUPERVISOR.get_daemon(&id).await;
1362
1363 SUPERVISOR
1368 .active_monitors
1369 .fetch_add(1, atomic::Ordering::Release);
1370 struct MonitorGuard;
1371 impl Drop for MonitorGuard {
1372 fn drop(&mut self) {
1373 SUPERVISOR
1374 .active_monitors
1375 .fetch_sub(1, atomic::Ordering::Release);
1376 SUPERVISOR.monitor_done.notify_waiters();
1377 }
1378 }
1379 let _monitor_guard = MonitorGuard;
1380 if !pre_drain_is_stopping
1385 && (current_daemon.is_none()
1386 || current_daemon.as_ref().is_some_and(|d| {
1387 d.pid != Some(pid) && !d.status.is_stopped() && !d.status.is_stopping()
1388 }))
1389 {
1390 return;
1392 }
1393 let already_stopped = current_daemon
1398 .as_ref()
1399 .is_some_and(|d| d.status.is_stopped());
1400 let is_stopping = already_stopped
1401 || pre_drain_is_stopping
1402 || current_daemon
1403 .as_ref()
1404 .is_some_and(|d| d.status.is_stopping());
1405
1406 let (exit_code, exit_reason) = match (&exit_status, is_stopping) {
1408 (Ok(status), true) => {
1409 (status.code().unwrap_or(-1), "stop")
1413 }
1414 (Ok(status), false) if status.success() => (status.code().unwrap_or(-1), "exit"),
1415 (Ok(status), false) => (status.code().unwrap_or(-1), "fail"),
1416 (Err(_), true) => {
1417 (-1, "stop")
1419 }
1420 (Err(_), false) => (-1, "fail"),
1421 };
1422
1423 if !already_stopped && !pre_drain_is_stopping {
1429 if let Ok(status) = &exit_status {
1430 info!("daemon {id} exited with status {status}");
1431 }
1432 let (new_status, last_exit_success) = match exit_reason {
1433 "stop" | "exit" => (
1434 DaemonStatus::Stopped,
1435 exit_status.as_ref().map(|s| s.success()).unwrap_or(true),
1436 ),
1437 _ => (DaemonStatus::Errored(exit_code), false),
1438 };
1439 if !SUPERVISOR
1445 .finalize_monitored_exit(
1446 &id,
1447 pid,
1448 monitor_token,
1449 new_status,
1450 Some(last_exit_success),
1451 )
1452 .await
1453 {
1454 debug!("daemon {id} exit state was not written; a successor owns the record");
1455 }
1456 }
1457
1458 let hook_extra_env = vec![
1460 ("PITCHFORK_EXIT_CODE".to_string(), exit_code.to_string()),
1461 ("PITCHFORK_EXIT_REASON".to_string(), exit_reason.to_string()),
1462 ];
1463
1464 let hooks_to_fire: Vec<HookType> = match exit_reason {
1466 "stop" => vec![HookType::OnStop, HookType::OnExit],
1467 "exit" => vec![HookType::OnExit],
1468 _ if hook_retry_count >= hook_retry.count() => {
1470 vec![HookType::OnFail, HookType::OnExit]
1471 }
1472 _ => vec![],
1473 };
1474
1475 for hook_type in hooks_to_fire {
1476 fire_hook(
1477 hook_type,
1478 id.clone(),
1479 daemon_dir.clone(),
1480 hook_retry_count,
1481 hook_daemon_env.clone(),
1482 hook_extra_env.clone(),
1483 )
1484 .await;
1485 }
1486 });
1487
1488 if let Some(ready_rx) = ready_rx {
1490 match ready_rx.await {
1491 Ok(Ok(())) => {
1492 info!("daemon {id} is ready");
1493 Ok(IpcResponse::DaemonReady { daemon })
1494 }
1495 Ok(Err(exit_code)) => {
1496 error!("daemon {id} failed before becoming ready");
1497 Ok(IpcResponse::DaemonFailedWithCode { exit_code })
1498 }
1499 Err(_) => {
1500 error!("readiness channel closed unexpectedly for daemon {id}");
1501 Ok(IpcResponse::DaemonStart { daemon })
1502 }
1503 }
1504 } else {
1505 Ok(IpcResponse::DaemonStart { daemon })
1506 }
1507 }
1508
1509 pub async fn stop(&self, id: &DaemonId) -> Result<IpcResponse> {
1511 let pitchfork_id = DaemonId::pitchfork();
1512 if *id == pitchfork_id {
1513 return Ok(IpcResponse::Error(
1514 "Cannot stop supervisor via stop command".into(),
1515 ));
1516 }
1517 info!("stopping daemon: {id}");
1518 if let Some(daemon) = self.get_daemon(id).await {
1519 trace!("daemon to stop: {daemon}");
1520 if let Some(pid) = daemon.pid {
1521 trace!("killing pid: {pid}");
1522 if PROCS.is_running(pid) {
1523 self.upsert_daemon(
1525 UpsertDaemonOpts::builder(id.clone())
1526 .set(|o| {
1527 o.pid = Some(pid);
1528 o.status = DaemonStatus::Stopping;
1529 })
1530 .build(),
1531 )
1532 .await?;
1533
1534 let stop_cfg = daemon.stop_signal.unwrap_or_default();
1537 let stop_signal: i32 = stop_cfg.signal.into();
1538 if let Err(e) = PROCS
1539 .kill_process_group_async(pid, stop_signal, stop_cfg.timeout)
1540 .await
1541 {
1542 debug!("failed to kill pid {pid}: {e}");
1543 if PROCS.is_running(pid) {
1545 debug!("failed to stop pid {pid}: process still running after kill");
1547 self.upsert_daemon(
1548 UpsertDaemonOpts::builder(id.clone())
1549 .set(|o| {
1550 o.pid = Some(pid); o.status = DaemonStatus::Running;
1552 })
1553 .build(),
1554 )
1555 .await?;
1556 return Ok(IpcResponse::DaemonStopFailed {
1557 error: format!(
1558 "process {pid} still running after kill attempt: {e}"
1559 ),
1560 });
1561 }
1562 }
1563
1564 self.upsert_daemon(
1569 UpsertDaemonOpts::builder(id.clone())
1570 .set(|o| {
1571 o.pid = None;
1572 o.status = DaemonStatus::Stopped;
1573 o.last_exit_success = Some(true);
1574 })
1575 .build(),
1576 )
1577 .await?;
1578 } else {
1579 debug!("pid {pid} not running, process may have exited unexpectedly");
1580 self.upsert_daemon(
1586 UpsertDaemonOpts::builder(id.clone())
1587 .set(|o| {
1588 o.pid = None;
1589 o.status = DaemonStatus::Stopped;
1590 })
1591 .build(),
1592 )
1593 .await?;
1594 return Ok(IpcResponse::DaemonWasNotRunning);
1595 }
1596 Ok(IpcResponse::Ok)
1597 } else {
1598 debug!("daemon {id} not running");
1599 Ok(IpcResponse::DaemonNotRunning)
1600 }
1601 } else {
1602 debug!("daemon {id} not found");
1603 Ok(IpcResponse::DaemonNotFound)
1604 }
1605 }
1606}
1607
1608#[cfg(unix)]
1609fn resolve_effective_run_identity(daemon_user: Option<&str>) -> Result<RunIdentity> {
1610 let s = settings();
1611 let settings_user = s.supervisor.user.trim();
1612 let daemon_user = daemon_user.map(str::trim).filter(|user| !user.is_empty());
1613 let settings_user = (!settings_user.is_empty()).then_some(settings_user);
1614 let configured = daemon_user.or(settings_user);
1615 let current_uid = nix::unistd::Uid::effective().as_raw();
1616 let current_gid = nix::unistd::Gid::effective().as_raw();
1617 resolve_run_identity(
1618 configured,
1619 current_uid,
1620 current_gid,
1621 std::env::var("SUDO_UID").ok().as_deref(),
1622 std::env::var("SUDO_GID").ok().as_deref(),
1623 )
1624}
1625
1626#[cfg(unix)]
1627fn resolve_run_identity(
1628 configured: Option<&str>,
1629 current_uid: u32,
1630 current_gid: u32,
1631 sudo_uid: Option<&str>,
1632 sudo_gid: Option<&str>,
1633) -> Result<RunIdentity> {
1634 let current_uid = nix::unistd::Uid::from_raw(current_uid);
1635 let current_gid = nix::unistd::Gid::from_raw(current_gid);
1636 if let Some(user) = configured {
1637 let identity = resolve_configured_user(user)?;
1638 ensure_can_use_identity(user, &identity, current_uid, current_gid)?;
1639 if identity.matches(current_uid, current_gid) {
1640 return Ok(RunIdentity::Inherit);
1641 }
1642 return Ok(identity);
1643 }
1644
1645 if current_uid.is_root()
1646 && let Some(identity) = resolve_sudo_identity(sudo_uid, sudo_gid)
1647 {
1648 return Ok(identity);
1649 }
1650
1651 Ok(RunIdentity::Inherit)
1652}
1653
1654#[cfg(unix)]
1655fn resolve_configured_user(user: &str) -> Result<RunIdentity> {
1656 if user.chars().all(|c| c.is_ascii_digit()) {
1657 let uid = user
1658 .parse::<u32>()
1659 .map_err(|e| miette::miette!("invalid run user UID '{}': {}", user, e))?;
1660 let user_record = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1661 .into_diagnostic()?
1662 .ok_or_else(|| miette::miette!("run user UID '{}' does not exist", user))?;
1663 return run_identity_from_user_record(user_record);
1664 }
1665
1666 let user_record = nix::unistd::User::from_name(user)
1667 .into_diagnostic()?
1668 .ok_or_else(|| miette::miette!("run user '{}' does not exist", user))?;
1669 run_identity_from_user_record(user_record)
1670}
1671
1672#[cfg(unix)]
1673fn run_identity_from_user_record(user: nix::unistd::User) -> Result<RunIdentity> {
1674 let username = CString::new(user.name)
1675 .map_err(|e| miette::miette!("run user name contains an interior nul byte: {}", e))?;
1676 Ok(RunIdentity::Switch {
1677 uid: user.uid,
1678 gid: user.gid,
1679 username: Some(username),
1680 })
1681}
1682
1683#[cfg(unix)]
1684fn run_identity_from_raw_ids(uid: u32, gid: u32, username: Option<CString>) -> RunIdentity {
1685 RunIdentity::Switch {
1686 uid: nix::unistd::Uid::from_raw(uid),
1687 gid: nix::unistd::Gid::from_raw(gid),
1688 username,
1689 }
1690}
1691
1692#[cfg(unix)]
1693fn resolve_sudo_identity(sudo_uid: Option<&str>, sudo_gid: Option<&str>) -> Option<RunIdentity> {
1694 let uid = sudo_uid?.parse::<u32>().ok()?;
1695 let gid = sudo_gid?.parse::<u32>().ok()?;
1696 let username = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1697 .ok()
1698 .flatten()
1699 .and_then(|u| CString::new(u.name).ok());
1700 Some(run_identity_from_raw_ids(uid, gid, username))
1701}
1702
1703#[cfg(unix)]
1704fn ensure_can_use_identity(
1705 configured_user: &str,
1706 identity: &RunIdentity,
1707 current_uid: nix::unistd::Uid,
1708 current_gid: nix::unistd::Gid,
1709) -> Result<()> {
1710 let RunIdentity::Switch { uid, gid, .. } = identity else {
1711 return Ok(());
1712 };
1713 if *uid == current_uid && *gid == current_gid {
1714 return Ok(());
1715 }
1716 if current_uid.is_root() {
1717 return Ok(());
1718 }
1719 Err(miette::miette!(
1720 "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.",
1721 configured_user,
1722 current_uid.as_raw(),
1723 current_gid.as_raw(),
1724 uid.as_raw(),
1725 gid.as_raw()
1726 ))
1727}
1728
1729#[cfg(unix)]
1730fn apply_run_identity(identity: &RunIdentity) -> std::io::Result<()> {
1731 let RunIdentity::Switch { uid, gid, username } = identity else {
1732 return Ok(());
1733 };
1734 if let Some(username) = username {
1735 initgroups_for_user(username, *gid)?;
1736 } else {
1737 setgroups_to_primary(*gid)?;
1738 }
1739 nix::unistd::setgid(*gid).map_err(nix_to_io_error)?;
1740 nix::unistd::setuid(*uid).map_err(nix_to_io_error)?;
1741 Ok(())
1742}
1743
1744#[cfg(unix)]
1745impl RunIdentity {
1746 fn matches(&self, uid: nix::unistd::Uid, gid: nix::unistd::Gid) -> bool {
1747 matches!(self, RunIdentity::Switch { uid: u, gid: g, .. } if *u == uid && *g == gid)
1748 }
1749}
1750
1751#[cfg(unix)]
1752fn setgroups_to_primary(gid: nix::unistd::Gid) -> std::io::Result<()> {
1753 let groups = [gid.as_raw() as libc::gid_t];
1754 #[cfg(any(target_os = "linux", target_os = "android"))]
1755 let group_count = groups.len();
1756 #[cfg(not(any(target_os = "linux", target_os = "android")))]
1757 let group_count = groups.len() as libc::c_int;
1758 let rc = unsafe { libc::setgroups(group_count, groups.as_ptr()) };
1759 if rc == -1 {
1760 Err(std::io::Error::last_os_error())
1761 } else {
1762 Ok(())
1763 }
1764}
1765
1766#[cfg(unix)]
1767fn initgroups_for_user(username: &CString, gid: nix::unistd::Gid) -> std::io::Result<()> {
1768 let gid = gid.as_raw();
1769 #[cfg(any(
1770 target_os = "macos",
1771 target_os = "ios",
1772 target_os = "tvos",
1773 target_os = "watchos"
1774 ))]
1775 let base_gid = i32::try_from(gid)
1776 .map_err(|_| std::io::Error::other(format!("gid {gid} is out of range")))?;
1777
1778 #[cfg(not(any(
1779 target_os = "macos",
1780 target_os = "ios",
1781 target_os = "tvos",
1782 target_os = "watchos"
1783 )))]
1784 let base_gid = gid as libc::gid_t;
1785
1786 let rc = unsafe { libc::initgroups(username.as_ptr(), base_gid) };
1789 if rc == -1 {
1790 Err(std::io::Error::last_os_error())
1791 } else {
1792 Ok(())
1793 }
1794}
1795
1796#[cfg(unix)]
1797fn nix_to_io_error(err: nix::errno::Errno) -> std::io::Error {
1798 std::io::Error::from_raw_os_error(err as i32)
1799}
1800
1801async fn check_ports_available(
1808 expected_ports: &[u16],
1809 auto_bump: bool,
1810 max_attempts: u32,
1811) -> Result<Vec<u16>> {
1812 if expected_ports.is_empty() {
1813 return Ok(Vec::new());
1814 }
1815
1816 for bump_offset in 0..=max_attempts {
1817 let candidate_ports: Vec<u16> = expected_ports
1819 .iter()
1820 .map(|&p| p.wrapping_add(bump_offset as u16))
1821 .collect();
1822
1823 let mut all_available = true;
1825 let mut conflicting_port = None;
1826
1827 for &port in &candidate_ports {
1828 if port == 0 {
1831 continue;
1832 }
1833
1834 if is_port_in_use(port).await {
1848 all_available = false;
1849 conflicting_port = Some(port);
1850 break;
1851 }
1852 }
1853
1854 if all_available {
1855 if candidate_ports.contains(&0) && !expected_ports.contains(&0) {
1859 return Err(PortError::NoAvailablePort {
1860 start_port: expected_ports[0],
1861 attempts: bump_offset + 1,
1862 }
1863 .into());
1864 }
1865 if bump_offset > 0 {
1866 info!("ports {expected_ports:?} bumped by {bump_offset} to {candidate_ports:?}");
1867 }
1868 return Ok(candidate_ports);
1869 }
1870
1871 if bump_offset == 0
1873 && !auto_bump
1874 && let Some(port) = conflicting_port
1875 {
1876 let (pid, process) = identify_port_owner(port).await;
1877 return Err(PortError::InUse { port, process, pid }.into());
1878 }
1879 }
1880
1881 Err(PortError::NoAvailablePort {
1883 start_port: expected_ports[0],
1884 attempts: max_attempts + 1,
1885 }
1886 .into())
1887}
1888
1889async fn is_port_in_use(port: u16) -> bool {
1895 tokio::task::spawn_blocking(move || {
1896 for &addr in &["0.0.0.0", "127.0.0.1", "::1"] {
1897 match std::net::TcpListener::bind((addr, port)) {
1898 Ok(listener) => drop(listener),
1899 Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => return true,
1900 Err(_) => continue,
1901 }
1902 }
1903 false
1904 })
1905 .await
1906 .unwrap_or(false)
1907}
1908
1909async fn identify_port_owner(port: u16) -> (u32, String) {
1914 tokio::task::spawn_blocking(move || {
1915 listeners::get_all()
1916 .ok()
1917 .and_then(|list| {
1918 list.into_iter()
1919 .find(|l| l.socket.port() == port)
1920 .map(|l| (l.process.pid, l.process.name))
1921 })
1922 .unwrap_or((0, "unknown".to_string()))
1923 })
1924 .await
1925 .unwrap_or((0, "unknown".to_string()))
1926}
1927
1928async fn detect_port_conflict(port: u16) -> Option<(u32, String)> {
1933 if !is_port_in_use(port).await {
1934 return None;
1935 }
1936 Some(identify_port_owner(port).await)
1937}
1938
1939fn detect_and_store_active_port(id: DaemonId, pid: u32) {
1954 tokio::spawn(async move {
1955 for delay_ms in [500u64, 1000, 2000, 4000] {
1959 tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
1960
1961 let expected_port: Option<u16> = {
1964 let state_file = SUPERVISOR.state_file.lock().await;
1965 match state_file.daemons.get(&id) {
1966 Some(d) if d.pid.is_none() => {
1967 debug!("daemon {id}: aborting active_port detection — process exited");
1968 return;
1969 }
1970 Some(d) => d
1971 .port
1972 .as_ref()
1973 .and_then(|p| p.expect.first().copied())
1974 .filter(|&p| p > 0),
1975 None => None,
1976 }
1977 };
1978
1979 let active_port = tokio::task::spawn_blocking(move || {
1980 let listeners = listeners::get_all().ok()?;
1981
1982 PROCS.refresh_processes();
1984
1985 let descendant_pids: std::collections::HashSet<u32> = PROCS
1986 .all_children(pid)
1987 .into_iter()
1988 .chain(std::iter::once(pid))
1989 .collect();
1990
1991 let process_ports: Vec<u16> = listeners
1992 .into_iter()
1993 .filter(|listener| descendant_pids.contains(&listener.process.pid))
1994 .map(|listener| listener.socket.port())
1995 .filter(|&port| port > 0)
1996 .collect();
1997
1998 if process_ports.is_empty() {
1999 return None;
2000 }
2001
2002 if let Some(ep) = expected_port
2005 && process_ports.contains(&ep)
2006 {
2007 return Some(ep);
2008 }
2009
2010 process_ports.into_iter().next()
2016 })
2017 .await
2018 .ok()
2019 .flatten();
2020
2021 if let Some(port) = active_port {
2022 debug!("daemon {id} active_port detected: {port}");
2023 let mut state_file = SUPERVISOR.state_file.lock().await;
2024 if let Some(d) = state_file.daemons.get(&id) {
2025 if d.pid == Some(pid) {
2029 state_file.set_active_port(&id, port);
2030 } else {
2031 debug!(
2032 "daemon {id}: skipping active_port write — PID mismatch \
2033 (expected {pid}, current {:?})",
2034 d.pid
2035 );
2036 return;
2037 }
2038 }
2039 return;
2040 }
2041
2042 debug!(
2043 "daemon {id}: no active port detected for pid {pid} or its descendants (will retry)"
2044 );
2045 }
2046
2047 debug!(
2048 "daemon {id}: active port detection exhausted all retries for pid {pid} and its descendants"
2049 );
2050 });
2051}
2052
2053fn is_daemon_slug_target(id: &DaemonId) -> bool {
2061 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
2065 slugs.iter().any(|(slug, entry)| {
2066 let daemon_name = entry.daemon.as_deref().unwrap_or(slug);
2067 id.name() == daemon_name
2068 })
2069}
2070
2071#[cfg(all(test, unix))]
2072mod tests {
2073 use super::*;
2074
2075 #[test]
2076 fn test_resolve_run_identity_empty_without_sudo() {
2077 let identity = resolve_run_identity(None, 501, 20, None, None).unwrap();
2078 assert_eq!(identity, RunIdentity::Inherit);
2079 }
2080
2081 #[test]
2082 fn test_resolve_run_identity_sudo_fallback() {
2083 let identity = resolve_run_identity(None, 0, 0, Some("501"), Some("20")).unwrap();
2084 let RunIdentity::Switch { uid, gid, .. } = identity else {
2085 panic!("expected identity switch");
2086 };
2087 assert_eq!(uid.as_raw(), 501);
2088 assert_eq!(gid.as_raw(), 20);
2089 }
2090
2091 #[test]
2092 fn test_resolve_run_identity_ignores_stale_sudo_when_not_root() {
2093 let identity = resolve_run_identity(None, 501, 20, Some("0"), Some("0")).unwrap();
2094 assert_eq!(identity, RunIdentity::Inherit);
2095 }
2096
2097 #[test]
2098 fn test_resolve_configured_user_root_name() {
2099 let identity = resolve_configured_user("root").unwrap();
2100 let RunIdentity::Switch { uid, username, .. } = identity else {
2101 panic!("expected identity switch");
2102 };
2103 assert_eq!(uid.as_raw(), 0);
2104 assert_eq!(
2105 username.as_deref().and_then(|s| s.to_str().ok()),
2106 Some("root")
2107 );
2108 }
2109
2110 #[test]
2111 fn test_resolve_configured_user_root_uid() {
2112 let identity = resolve_configured_user("0").unwrap();
2113 let RunIdentity::Switch { uid, username, .. } = identity else {
2114 panic!("expected identity switch");
2115 };
2116 assert_eq!(uid.as_raw(), 0);
2117 assert_eq!(
2118 username.as_deref().and_then(|s| s.to_str().ok()),
2119 Some("root")
2120 );
2121 }
2122
2123 #[test]
2124 fn test_resolve_configured_user_missing_user_fails() {
2125 let err = resolve_configured_user("pitchfork-user-that-should-not-exist")
2126 .unwrap_err()
2127 .to_string();
2128 assert!(err.contains("does not exist"));
2129 }
2130
2131 #[test]
2132 fn test_resolve_run_identity_requires_root_for_user_switch() {
2133 let err = resolve_run_identity(Some("root"), 501, 20, None, None)
2134 .unwrap_err()
2135 .to_string();
2136 assert!(err.contains("Restart the supervisor with sudo"));
2137 }
2138
2139 #[test]
2140 fn test_resolve_run_identity_same_user_is_noop() {
2141 let identity = resolve_run_identity(Some("root"), 0, 0, Some("501"), Some("20")).unwrap();
2142 assert_eq!(identity, RunIdentity::Inherit);
2143 }
2144}
2145
2146fn inject_proxy_env(cmd: &mut tokio::process::Command, slug: &Option<String>) {
2155 let s = crate::settings::settings();
2156 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
2157
2158 if should_force_loopback_host(slug) && !lan_enabled {
2159 cmd.env("HOST", "127.0.0.1");
2162 }
2163
2164 if let Some(url) = build_pitchfork_url(slug, &s) {
2166 cmd.env("PITCHFORK_URL", &url);
2167 }
2168
2169 if s.proxy.enable && s.proxy.https {
2171 let ca_path = if s.proxy.tls_cert.is_empty() {
2172 crate::env::PITCHFORK_STATE_DIR.join("proxy").join("ca.pem")
2173 } else {
2174 std::path::PathBuf::from(&s.proxy.tls_cert)
2175 };
2176 if ca_path.exists() {
2177 cmd.env("NODE_EXTRA_CA_CERTS", ca_path.to_string_lossy().to_string());
2178 }
2179 }
2180
2181 if s.proxy.enable {
2183 let tld = if lan_enabled { "local" } else { &s.proxy.tld };
2184 cmd.env("__VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS", format!(".{tld}"));
2185 }
2186
2187 if lan_enabled {
2189 cmd.env("PITCHFORK_LAN", "1");
2190 }
2191}
2192
2193fn should_force_loopback_host(slug: &Option<String>) -> bool {
2194 let Some(slug) = slug.as_deref() else {
2195 return false;
2196 };
2197
2198 let s = crate::settings::settings();
2199 if !s.proxy.enable {
2200 return false;
2201 }
2202
2203 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
2204 slugs.contains_key(slug)
2205}
2206
2207fn build_pitchfork_url(slug: &Option<String>, s: &crate::settings::Settings) -> Option<String> {
2211 let slug = slug.as_ref()?;
2212 if !s.proxy.enable {
2213 return None;
2214 }
2215 let scheme = if s.proxy.https { "https" } else { "http" };
2216 let port = u16::try_from(s.proxy.port).ok().filter(|&p| p > 0)?;
2217 let port_suffix = if (scheme == "https" && port == 443) || (scheme == "http" && port == 80) {
2218 String::new()
2219 } else {
2220 format!(":{port}")
2221 };
2222 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
2223 let tld = if lan_enabled { "local" } else { &s.proxy.tld };
2224 Some(format!("{scheme}://{slug}.{tld}{port_suffix}",))
2225}
2226
2227#[cfg(test)]
2228mod ready_check_tests {
2229 use super::*;
2230 use std::time::Duration;
2231
2232 #[test]
2233 fn any_ready_check_remaining_prefers_unbounded_checks() {
2234 let http = ReadyHttp::new("http://localhost/health");
2235 let cmd = ReadyCmd::new("true");
2236
2237 assert!(any_ready_check_remaining(
2238 None,
2239 false,
2240 None,
2241 false,
2242 Some(&http),
2243 false,
2244 None,
2245 false
2246 ));
2247 assert!(any_ready_check_remaining(
2248 None,
2249 false,
2250 None,
2251 false,
2252 None,
2253 false,
2254 Some(&cmd),
2255 false
2256 ));
2257 assert!(any_ready_check_remaining(
2258 None,
2259 false,
2260 Some(&ReadyPort::new(8080)),
2261 false,
2262 Some(&http),
2263 true,
2264 Some(&cmd),
2265 true
2266 ));
2267 }
2268
2269 #[test]
2270 fn any_ready_check_remaining_exhausted_timed_checks() {
2271 let http = ReadyHttp {
2272 url: "http://localhost/health".to_string(),
2273 status: vec![],
2274 timeout: Some(Duration::from_secs(5)),
2275 };
2276 let cmd = ReadyCmd {
2277 run: "true".to_string(),
2278 timeout: Some(Duration::from_secs(5)),
2279 };
2280
2281 assert!(any_ready_check_remaining(
2282 None,
2283 false,
2284 None,
2285 false,
2286 Some(&http),
2287 false,
2288 Some(&cmd),
2289 false
2290 ));
2291 assert!(!any_ready_check_remaining(
2292 None,
2293 false,
2294 None,
2295 false,
2296 Some(&http),
2297 true,
2298 Some(&cmd),
2299 true
2300 ));
2301 }
2302
2303 #[tokio::test]
2304 async fn spawn_cmd_probe_reports_success() {
2305 let id = DaemonId::new("global", "probe-test");
2306 let probe = spawn_cmd_probe(&id, "true", &std::env::temp_dir());
2307 let status = probe.result_rx.await.unwrap().unwrap();
2308 assert!(status.success());
2309 }
2310
2311 #[tokio::test]
2312 async fn spawn_cmd_probe_stops_on_request() {
2313 let id = DaemonId::new("global", "probe-test");
2314 let probe = spawn_cmd_probe(&id, "sleep 30", &std::env::temp_dir());
2315 let CmdProbe {
2316 cancel_tx,
2317 result_rx,
2318 } = probe;
2319 let _ = cancel_tx.send(());
2320 let status = result_rx.await.unwrap().unwrap();
2321 assert!(!status.success());
2322 }
2323}