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
170const SINK_OUTPUT_TIMEOUT: Duration = Duration::from_millis(400);
175
176impl Supervisor {
177 pub async fn run(&self, opts: RunOptions) -> Result<IpcResponse> {
179 let id = &opts.id;
180 let cmd = opts.cmd.clone();
181
182 {
184 let mut pending = self.pending_autostops.lock().await;
185 if pending.remove(id).is_some() {
186 info!("cleared pending autostop for {id} (daemon starting)");
187 }
188 }
189
190 let daemon = self.get_daemon(id).await;
191 if let Some(daemon) = daemon {
192 if !daemon.status.is_stopping()
195 && !daemon.status.is_stopped()
196 && let Some(pid) = daemon.pid
197 {
198 if opts.force {
199 self.stop(id).await?;
200 info!("run: stop completed for daemon {id}");
201 } else {
202 warn!("daemon {id} already running with pid {pid}");
203 return Ok(IpcResponse::DaemonAlreadyRunning);
204 }
205 }
206 }
207
208 if opts.wait_ready && opts.retry.count() > 0 {
210 let max_attempts = opts.retry.count().saturating_add(1);
212 for attempt in 0..max_attempts {
213 let mut retry_opts = opts.clone();
214 retry_opts.retry_count = attempt;
215 retry_opts.cmd = cmd.clone();
216
217 let result = self.run_once(retry_opts).await?;
218
219 match result {
220 IpcResponse::DaemonReady { daemon } => {
221 return Ok(IpcResponse::DaemonReady { daemon });
222 }
223 IpcResponse::DaemonFailedWithCode { exit_code } => {
224 if attempt < opts.retry.count() {
225 let backoff_secs = 2u64.saturating_pow(attempt).min(3600);
226 info!(
227 "daemon {id} failed (attempt {}/{}), retrying in {}s",
228 attempt + 1,
229 max_attempts,
230 backoff_secs
231 );
232 fire_hook(
233 HookType::OnRetry,
234 id.clone(),
235 opts.dir.0.clone(),
236 attempt + 1,
237 opts.env.clone(),
238 vec![],
239 )
240 .await;
241 time::sleep(Duration::from_secs(backoff_secs)).await;
242 continue;
243 } else {
244 info!("daemon {id} failed after {max_attempts} attempts");
245 return Ok(IpcResponse::DaemonFailedWithCode { exit_code });
246 }
247 }
248 other => return Ok(other),
249 }
250 }
251 }
252
253 self.run_once(opts).await
255 }
256
257 pub(crate) async fn run_once(&self, opts: RunOptions) -> Result<IpcResponse> {
259 let id = &opts.id;
260 let original_cmd = opts.cmd.clone(); let (ready_tx, ready_rx) = if opts.wait_ready {
264 let (tx, rx) = oneshot::channel();
265 (Some(tx), Some(rx))
266 } else {
267 (None, None)
268 };
269
270 let expected_ports = opts
272 .port
273 .as_ref()
274 .map(|p| p.expect.clone())
275 .unwrap_or_default();
276 let (resolved_ports, effective_ready_port) = if !expected_ports.is_empty() {
277 let port_cfg = opts.port.as_ref().unwrap();
278 match check_ports_available(
279 &expected_ports,
280 port_cfg.auto_bump(),
281 port_cfg.max_bump_attempts(),
282 )
283 .await
284 {
285 Ok(resolved) => {
286 let ready_port = if let Some(configured_port) =
287 opts.ready_port.as_ref().and_then(|p| p.as_port())
288 {
289 let bump_offset = resolved
291 .first()
292 .unwrap_or(&0)
293 .saturating_sub(*expected_ports.first().unwrap_or(&0));
294 if expected_ports.contains(&configured_port) && bump_offset > 0 {
295 configured_port
296 .checked_add(bump_offset)
297 .or(Some(configured_port))
298 } else {
299 Some(configured_port)
300 }
301 } else if opts.ready_output.is_none()
302 && opts.ready_http.is_none()
303 && opts.ready_cmd.is_none()
304 && opts.ready_delay.is_none()
305 {
306 resolved.first().copied().filter(|&p| p != 0)
310 } else {
311 None
315 };
316 info!("daemon {id}: ports {expected_ports:?} resolved to {resolved:?}");
317 (resolved, ready_port)
318 }
319 Err(e) => {
320 error!("daemon {id}: port check failed: {e}");
321 if let Some(port_error) = e.downcast_ref::<PortError>() {
323 match port_error {
324 PortError::InUse { port, process, pid } => {
325 return Ok(IpcResponse::PortConflict {
326 port: *port,
327 process: process.clone(),
328 pid: *pid,
329 });
330 }
331 PortError::NoAvailablePort {
332 start_port,
333 attempts,
334 } => {
335 return Ok(IpcResponse::NoAvailablePort {
336 start_port: *start_port,
337 attempts: *attempts,
338 });
339 }
340 }
341 }
342 return Ok(IpcResponse::DaemonFailed {
343 error: e.to_string(),
344 });
345 }
346 }
347 } else {
348 if let Some(port) = opts.ready_port.as_ref().and_then(|p| p.as_port())
354 && port > 0
355 && let Some((pid, process)) = detect_port_conflict(port).await
356 {
357 return Ok(IpcResponse::PortConflict { port, process, pid });
358 }
359 (
360 Vec::new(),
361 opts.ready_port.as_ref().and_then(|p| p.as_port()),
362 )
363 };
364
365 let shell_setting = settings().general.shell.clone();
369 let shell_parts = match shell_words::split(&shell_setting) {
370 Ok(parts) if !parts.is_empty() => parts,
371 Ok(_) => {
372 return Ok(IpcResponse::DaemonFailed {
373 error: "general.shell setting is empty".to_string(),
374 });
375 }
376 Err(e) => {
377 return Ok(IpcResponse::DaemonFailed {
378 error: format!("failed to parse general.shell setting {shell_setting:?}: {e}"),
379 });
380 }
381 };
382 let (shell_program, shell_args) = shell_parts.split_first().unwrap();
383
384 let run_script = opts
389 .run
390 .clone()
391 .unwrap_or_else(|| shell_words::join(&original_cmd));
392
393 let (program, args) = if opts.mise.unwrap_or(settings().general.mise) {
394 match settings().resolve_mise_bin() {
395 Some(mise_bin) => {
396 let mise_bin_str = mise_bin.to_string_lossy().to_string();
397 info!("daemon {id}: wrapping command with mise ({mise_bin_str})");
398 let mut args = vec!["x".to_string(), "--".to_string()];
399 args.push(shell_program.clone());
400 args.extend(shell_args.iter().cloned());
401 args.push(run_script);
402 (mise_bin_str, args)
403 }
404 None => {
405 warn!("daemon {id}: mise=true but mise binary not found, running without mise");
406 let mut args: Vec<String> = shell_args.to_vec();
407 args.push(run_script);
408 (shell_program.clone(), args)
409 }
410 }
411 } else {
412 let mut args: Vec<String> = shell_args.to_vec();
413 args.push(run_script);
414 (shell_program.clone(), args)
415 };
416 #[cfg(unix)]
417 let run_identity = match resolve_effective_run_identity(opts.user.as_deref()) {
418 Ok(identity) => identity,
419 Err(e) => {
420 return Ok(IpcResponse::DaemonFailed {
421 error: e.to_string(),
422 });
423 }
424 };
425 info!("run: spawning daemon {id} with {program} {args:?}");
426
427 #[cfg(unix)]
429 let pty_pair = if opts.pty.unwrap_or(false) {
430 match super::pty::openpty() {
431 Ok(pair) => {
432 info!("daemon {id}: allocated PTY (pty = true)");
433 Some(pair)
434 }
435 Err(e) => {
436 warn!("daemon {id}: failed to allocate PTY, falling back to pipes: {e}");
437 None
438 }
439 }
440 } else {
441 None
442 };
443
444 let mut sink_pipe = None;
447 let mut sink_writer = None;
448 let mut sink_child = None;
449 if super::log_sink::is_supported(&opts) {
450 let log_format = opts
451 .log_format
452 .clone()
453 .unwrap_or_else(|| settings().logs.log_format.clone());
454 match super::log_sink::SinkPipe::new(log_format) {
455 Ok((pipe, writer)) => match pipe.start(id) {
456 Ok(child) => {
457 sink_child = Some(super::log_sink::PendingSink::new(child));
458 sink_pipe = Some(pipe);
459 sink_writer = Some(writer);
460 }
461 Err(e) => {
462 warn!("could not start log sink for {id}, capturing in-process: {e}");
463 }
464 },
465 Err(e) => {
466 warn!("could not create log pipe for {id}, capturing in-process: {e}");
469 }
470 }
471 }
472
473 let mut cmd = tokio::process::Command::new(&program);
474
475 #[cfg(unix)]
476 if let Some(ref pair) = pty_pair {
477 let slave_file = std::fs::File::from(
481 pair.slave
482 .try_clone()
483 .map_err(|e| miette::miette!("failed to dup slave PTY fd: {e}"))?,
484 );
485 cmd.stdin(std::process::Stdio::from(slave_file.try_clone().map_err(
486 |e| miette::miette!("failed to clone slave PTY fd for stdin: {e}"),
487 )?));
488 cmd.stdout(std::process::Stdio::from(slave_file.try_clone().map_err(
489 |e| miette::miette!("failed to clone slave PTY fd for stdout: {e}"),
490 )?));
491 cmd.stderr(std::process::Stdio::from(slave_file));
492 } else if let Some(writer) = sink_writer.take() {
493 let dup = writer
496 .try_clone()
497 .map_err(|e| miette::miette!("failed to dup log pipe for stderr: {e}"))?;
498 cmd.stdout(std::process::Stdio::from(writer))
499 .stderr(std::process::Stdio::from(dup));
500 } else {
501 cmd.stdout(std::process::Stdio::piped())
502 .stderr(std::process::Stdio::piped());
503 }
504
505 #[cfg(not(unix))]
506 if let Some(writer) = sink_writer.take() {
507 let dup = writer
508 .try_clone()
509 .map_err(|e| miette::miette!("failed to dup log pipe for stderr: {e}"))?;
510 cmd.stdout(std::process::Stdio::from(writer))
511 .stderr(std::process::Stdio::from(dup));
512 } else {
513 cmd.stdout(std::process::Stdio::piped())
514 .stderr(std::process::Stdio::piped());
515 }
516
517 cmd.args(&args).current_dir(&opts.dir);
518
519 #[cfg(unix)]
520 if pty_pair.is_none() {
521 cmd.stdin(std::process::Stdio::null());
522 }
523
524 #[cfg(not(unix))]
525 cmd.stdin(std::process::Stdio::null());
526
527 if let Some(ref path) = *env::ORIGINAL_PATH {
529 cmd.env("PATH", path);
530 }
531
532 if let Some(ref env_vars) = opts.env {
534 cmd.envs(env_vars);
535 }
536
537 cmd.env("PITCHFORK_DAEMON_ID", id.qualified());
539 cmd.env("PITCHFORK_DAEMON_NAMESPACE", id.namespace());
540 cmd.env("PITCHFORK_RETRY_COUNT", opts.retry_count.to_string());
541
542 if !resolved_ports.is_empty() {
544 cmd.env("PORT", resolved_ports[0].to_string());
548 for (i, port) in resolved_ports.iter().enumerate() {
550 cmd.env(format!("PORT{i}"), port.to_string());
551 }
552 }
553
554 inject_proxy_env(&mut cmd, &opts.slug);
556
557 #[cfg(unix)]
558 {
559 let run_identity = run_identity.clone();
560 let use_pty = pty_pair.is_some();
561 unsafe {
562 cmd.pre_exec(move || {
563 nix::unistd::setsid().map_err(nix_to_io_error)?;
564
565 if use_pty {
569 let ret = libc::ioctl(0, libc::TIOCSCTTY as libc::c_ulong, 0);
570 if ret < 0 {
571 #[cfg(target_os = "linux")]
574 eprintln!(
575 "pitchfork: TIOCSCTTY failed: {}",
576 std::io::Error::last_os_error()
577 );
578 }
579 }
580
581 apply_run_identity(&run_identity)?;
582 Ok(())
583 });
584 }
585 }
586
587 let spawn_time = chrono::Local::now();
590 let mut child = cmd.spawn().into_diagnostic()?;
596 let pid = match child.id() {
597 Some(p) => p,
598 None => {
599 warn!("Daemon {id} exited before PID could be captured");
600 if sink_child.is_some() {
606 super::log_sink::wait_for_output(id, spawn_time, SINK_OUTPUT_TIMEOUT).await;
607 }
608 return Ok(IpcResponse::DaemonFailed {
609 error: "Process exited immediately".to_string(),
610 });
611 }
612 };
613 info!("started daemon {id} with pid {pid}");
614 PROCS.refresh_pids(&[pid]);
615 let monitored_guard = super::adopt::MonitoredGuard::register(id.clone(), pid);
623 let monitor_token = monitored_guard.token();
624
625 let using_sink = sink_pipe.is_some();
628 if let Some(pipe) = sink_pipe.take()
631 && let Some(child) = sink_child.as_mut().and_then(|pending| pending.take())
632 {
633 pipe.supervise(id.clone(), monitor_token, child);
634 }
635 let daemon = self
636 .upsert_daemon(
637 UpsertDaemonOpts::from_run_options(&opts, DaemonStatus::Running)
638 .set(|o| {
639 o.pid = Some(pid);
640 o.cmd = Some(original_cmd);
641 o.ready_port = effective_ready_port.map(|p| ReadyPort {
642 port: Some(p),
643 template: None,
644 timeout: opts.ready_port.as_ref().and_then(|rp| rp.timeout),
645 });
646 o.port = crate::config_types::PortConfig::from_parts(
647 expected_ports,
648 opts.port.as_ref().map(|p| p.bump).unwrap_or_default(),
649 );
650 o.resolved_port = resolved_ports;
651 })
652 .build(),
653 )
654 .await?;
655
656 let id_clone = id.clone();
657 let ready_delay = opts.ready_delay;
658 let ready_output = opts.ready_output.clone();
659 let ready_http = opts.ready_http.clone();
660 let ready_port = effective_ready_port;
661 let implicit_ready_port = ready_port.map(|p| ReadyPort {
662 port: Some(p),
663 template: None,
664 timeout: None,
665 });
666 let ready_port_config = opts.ready_port.clone().or(implicit_ready_port);
667 let ready_cmd = opts.ready_cmd.clone();
668 let daemon_dir = opts.dir.0.clone();
669 let hook_retry_count = opts.retry_count;
670 let hook_retry = opts.retry;
671 let hook_daemon_env = opts.env.clone();
672 let on_output_hook = opts.on_output_hook.clone();
673 let has_port_config = opts.port.as_ref().is_some_and(|p| !p.expect.is_empty())
679 || (settings().proxy.enable && is_daemon_slug_target(id));
680 let expected_port: Option<u16> = opts.port.as_ref().and_then(|p| p.expect.first().copied());
686 let daemon_pid = pid;
687
688 #[cfg(unix)]
692 let pty_reader = pty_pair.map(|p| {
693 tokio::io::BufReader::new(tokio::fs::File::from_std(std::fs::File::from(p.master)))
694 .lines()
695 });
696 #[cfg(not(unix))]
697 let pty_reader: Option<tokio::io::Lines<tokio::io::BufReader<tokio::fs::File>>> = None;
698 let stdout_reader = if pty_reader.is_none() {
699 child
700 .stdout
701 .take()
702 .map(|s| tokio::io::BufReader::new(s).lines())
703 } else {
704 None
705 };
706 let stderr_reader = if pty_reader.is_none() {
707 child
708 .stderr
709 .take()
710 .map(|s| tokio::io::BufReader::new(s).lines())
711 } else {
712 None
713 };
714
715 if !using_sink
716 && pty_reader.is_none()
717 && (stdout_reader.is_none() || stderr_reader.is_none())
718 {
719 error!("Failed to capture stdout/stderr for daemon {id}");
720 }
721
722 tokio::spawn(async move {
723 let id = id_clone;
724 let _monitored_guard = monitored_guard;
727
728 let (output_tx, mut output_rx) = tokio::sync::mpsc::channel::<String>(256);
730
731 if let Some(mut reader) = pty_reader {
732 tokio::spawn(async move {
736 while let Ok(Some(mut line)) = reader.next_line().await {
737 if line.ends_with('\r') {
739 line.pop();
740 }
741 if output_tx.send(line).await.is_err() {
742 break;
743 }
744 }
745 });
746 } else {
747 if let Some(mut stdout) = stdout_reader {
752 let tx = output_tx.clone();
753 tokio::spawn(async move {
754 while let Ok(Some(line)) = stdout.next_line().await {
755 if tx.send(line).await.is_err() {
756 break;
757 }
758 }
759 });
760 }
761 if let Some(mut stderr) = stderr_reader {
762 let tx = output_tx.clone();
763 tokio::spawn(async move {
764 while let Ok(Some(line)) = stderr.next_line().await {
765 if tx.send(line).await.is_err() {
766 break;
767 }
768 }
769 });
770 }
771 drop(output_tx);
773 }
774 let log_store = Arc::clone(&LOG_STORE);
775 let log_format = opts
776 .log_format
777 .clone()
778 .unwrap_or_else(|| crate::settings::settings().logs.log_format.clone());
779 let parse_line = move |line: &str| crate::log_parse::parse(line, &log_format);
780
781 const LOG_BATCH_SIZE: usize = 100;
782 const LOG_FLUSH_INTERVAL: Duration = Duration::from_millis(100);
783 let mut log_buffer: Vec<crate::log_parse::ParsedLog> =
784 Vec::with_capacity(LOG_BATCH_SIZE);
785 let mut log_flush_interval = tokio::time::interval(LOG_FLUSH_INTERVAL);
786 log_flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
787
788 let flush_logs =
789 |buffer: &mut Vec<crate::log_parse::ParsedLog>| -> Option<tokio::task::JoinHandle<()>> {
790 if buffer.is_empty() {
791 return None;
792 }
793 let store = Arc::clone(&log_store);
794 let id = id.clone();
795 let batch = std::mem::take(buffer);
796 Some(tokio::task::spawn_blocking(move || {
797 if let Err(e) = store.append_structured_batch(&id, &batch) {
798 error!("Failed to write batch to log for daemon {id}: {e}");
799 }
800 }))
801 };
802
803 let mut ready_notified = false;
807 let mut ready_tx = ready_tx;
808 let ready_pattern = ready_output
809 .as_ref()
810 .and_then(|o| get_or_compile_regex(&o.pattern));
811 let mut active_port_spawned = false;
813
814 let on_output_hook = match on_output_hook {
818 Some(ref hook) => match hook.validate(id.name()) {
819 Ok(()) => on_output_hook,
820 Err(e) => {
821 error!("{e}");
822 None
823 }
824 },
825 None => None,
826 };
827
828 let on_output_pattern: Option<regex::Regex> = on_output_hook
831 .as_ref()
832 .and_then(|h| h.regex.as_deref().and_then(get_or_compile_regex));
833 let on_output_debounce = on_output_hook
834 .as_ref()
835 .map(|h| h.debounce_duration())
836 .unwrap_or(Duration::from_millis(1000));
837 let mut on_output_last_fired: Option<std::time::Instant> = None;
839
840 let mut delay_timer =
841 ready_delay.map(|secs| Box::pin(time::sleep(Duration::from_secs(secs))));
842
843 let mut http_exhausted = false;
845 let mut cmd_exhausted = false;
846 let mut port_exhausted = false;
847 let mut output_exhausted = false;
848
849 let s = settings();
851 let ready_check_interval = s.supervisor_ready_check_interval();
852 let http_client_timeout = s.supervisor_http_client_timeout();
853
854 let mut output_deadline = ready_output
856 .as_ref()
857 .and_then(|o| o.timeout)
858 .map(|d| Box::pin(time::sleep(d)));
859
860 let mut http_check_interval = ready_http
862 .as_ref()
863 .map(|_| tokio::time::interval(ready_check_interval));
864 let mut http_deadline = ready_http
865 .as_ref()
866 .and_then(|h| h.timeout)
867 .map(|d| Box::pin(time::sleep(d)));
868 let http_client = ready_http.as_ref().map(|_| {
869 reqwest::Client::builder()
870 .timeout(http_client_timeout)
871 .build()
872 .unwrap_or_default()
873 });
874
875 let mut port_check_interval =
877 ready_port.map(|_| tokio::time::interval(ready_check_interval));
878 let mut port_deadline = ready_port_config
879 .as_ref()
880 .and_then(|p| p.timeout)
881 .map(|d| Box::pin(time::sleep(d)));
882
883 let mut cmd_probe: Option<CmdProbe> = None;
886 let mut cmd_respawn_delay: Option<_> = None;
887 let mut cmd_deadline = ready_cmd
888 .as_ref()
889 .and_then(|c| c.timeout)
890 .map(|d| Box::pin(time::sleep(d)));
891 if let Some(ref cmd) = ready_cmd {
892 cmd_probe = Some(spawn_cmd_probe(&id, &cmd.run, daemon_dir.as_path()));
893 }
894
895 let (exit_tx, mut exit_rx) =
897 tokio::sync::mpsc::channel::<std::io::Result<std::process::ExitStatus>>(1);
898
899 let child_pid = child.id().unwrap_or(0);
901 tokio::spawn(async move {
902 let result = child.wait().await;
903 #[cfg(all(unix, not(target_os = "linux")))]
913 let result = match &result {
914 Err(e) if e.raw_os_error() == Some(nix::libc::ECHILD) => {
915 if let Some(code) = super::REAPED_STATUSES.lock().await.remove(&child_pid) {
916 warn!(
917 "daemon pid {child_pid} wait() got ECHILD; \
918 recovered exit code {code} from zombie reaper"
919 );
920 use std::os::unix::process::ExitStatusExt;
925 if code >= 0 {
926 Ok(std::process::ExitStatus::from_raw(code << 8))
927 } else {
928 Ok(std::process::ExitStatus::from_raw((-code) & 0x7f))
930 }
931 } else {
932 warn!(
933 "daemon pid {child_pid} wait() got ECHILD but no \
934 stashed status found; reporting as error"
935 );
936 result
937 }
938 }
939 _ => result,
940 };
941 debug!("daemon pid {child_pid} wait() completed with result: {result:?}");
942 let _ = exit_tx.send(result).await;
943 });
944
945 #[allow(unused_assignments)]
946 let mut exit_status = None;
948
949 if has_port_config
955 && ready_pattern.is_none()
956 && ready_http.is_none()
957 && ready_port.is_none()
958 && ready_cmd.is_none()
959 && delay_timer.is_none()
960 {
961 active_port_spawned = true;
962 detect_and_store_active_port(id.clone(), daemon_pid);
963 }
964
965 loop {
966 select! {
970 biased;
971 Some(result) = exit_rx.recv() => {
972 exit_status = Some(result);
974 debug!("daemon {id} process exited, exit_status: {exit_status:?}");
975 if !ready_notified {
976 if let Some(tx) = ready_tx.take() {
977 let is_success = exit_status.as_ref()
979 .and_then(|r| r.as_ref().ok())
980 .map(|s| s.success())
981 .unwrap_or(false);
982
983 if is_success {
984 debug!("daemon {id} exited successfully before ready check, sending success notification");
985 let _ = tx.send(Ok(()));
986 } else {
987 let exit_code = exit_status.as_ref()
988 .and_then(|r| r.as_ref().ok())
989 .and_then(|s| s.code());
990 debug!("daemon {id} exited with failure before ready check, sending failure notification with exit_code: {exit_code:?}");
991 let _ = tx.send(Err(exit_code));
992 }
993 }
994 } else {
995 debug!("daemon {id} was already marked ready, not sending notification");
996 }
997 break;
998 },
999 Some(line) = output_rx.recv() => {
1000 let parsed = parse_line(&line);
1001 log_buffer.push(parsed);
1002 if log_buffer.len() >= LOG_BATCH_SIZE {
1003 let _ = flush_logs(&mut log_buffer);
1004 }
1005 trace!("output: {id} {line}");
1006
1007 let line_clean = console::strip_ansi_codes(&line).to_string();
1010
1011 if !ready_notified
1013 && !output_exhausted
1014 && let Some(ref pattern) = ready_pattern
1015 && pattern.is_match(&line_clean)
1016 {
1017 if let Some(handle) = flush_logs(&mut log_buffer) {
1022 let _ = handle.await;
1023 }
1024 info!("daemon {id} ready: output matched pattern");
1025 ready_notified = true;
1026 if let Some(tx) = ready_tx.take() {
1027 let _ = tx.send(Ok(()));
1028 }
1029 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1030 stop_cmd_probe_state(&mut cmd_probe);
1031 http_deadline = None;
1032 cmd_deadline = None;
1033 port_deadline = None;
1034 output_deadline = None;
1035 if !active_port_spawned && has_port_config {
1036 active_port_spawned = true;
1037 detect_and_store_active_port(id.clone(), daemon_pid);
1038 }
1039 }
1040
1041 if let Some(ref hook) = on_output_hook {
1043 let matched = match (&hook.filter, &on_output_pattern) {
1044 (Some(substr), _) => line_clean.contains(substr.as_str()),
1045 (None, Some(re)) => re.is_match(&line_clean),
1046 (None, None) => true,
1047 };
1048 if matched {
1049 let now = std::time::Instant::now();
1050 let elapsed = on_output_last_fired.map(|t| now.duration_since(t));
1051 if elapsed.is_none_or(|e| e >= on_output_debounce) {
1052 on_output_last_fired = Some(now);
1053 hooks::fire_output_hook(id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook.run.clone(), line_clean.clone()).await;
1054 }
1055 }
1056 }
1057 tokio::task::yield_now().await;
1060 }
1061 _ = async {
1062 if let Some(ref mut deadline) = http_deadline {
1063 deadline.await;
1064 } else {
1065 std::future::pending::<()>().await;
1066 }
1067 }, if !ready_notified && ready_http.is_some() => {
1068 http_exhausted = true;
1069 http_deadline = None;
1070 http_check_interval = None;
1071 warn!("daemon {id}: HTTP readiness check timed out");
1072 let any_remaining = any_ready_check_remaining(
1073 ready_output.as_ref(),
1074 output_exhausted,
1075 ready_port_config.as_ref(),
1076 port_exhausted,
1077 ready_http.as_ref(),
1078 http_exhausted,
1079 ready_cmd.as_ref(),
1080 cmd_exhausted,
1081 );
1082 if !any_remaining {
1083 error!("daemon {id}: all readiness checks exhausted, failing");
1084 stop_cmd_probe_state(&mut cmd_probe);
1085 if let Some(tx) = ready_tx.take() {
1086 let _ = tx.send(Err(Some(124)));
1087 }
1088 let stop_cfg = opts.stop_signal.unwrap_or_default();
1089 let _ = PROCS.kill_process_group_async(daemon_pid, stop_cfg.signal.into(), stop_cfg.timeout).await;
1090 break;
1091 }
1092 }
1093 _ = async {
1094 if let Some(ref mut deadline) = output_deadline {
1095 deadline.await;
1096 } else {
1097 std::future::pending::<()>().await;
1098 }
1099 }, if !ready_notified && ready_output.is_some() => {
1100 output_exhausted = true;
1101 output_deadline = None;
1102 warn!("daemon {id}: output readiness check timed out");
1103 let any_remaining = any_ready_check_remaining(
1104 ready_output.as_ref(),
1105 output_exhausted,
1106 ready_port_config.as_ref(),
1107 port_exhausted,
1108 ready_http.as_ref(),
1109 http_exhausted,
1110 ready_cmd.as_ref(),
1111 cmd_exhausted,
1112 );
1113 if !any_remaining {
1114 error!("daemon {id}: all readiness checks exhausted, failing");
1115 stop_cmd_probe_state(&mut cmd_probe);
1116 if let Some(tx) = ready_tx.take() {
1117 let _ = tx.send(Err(Some(124)));
1118 }
1119 let stop_cfg = opts.stop_signal.unwrap_or_default();
1120 let _ = PROCS.kill_process_group_async(daemon_pid, stop_cfg.signal.into(), stop_cfg.timeout).await;
1121 break;
1122 }
1123 }
1124 _ = async {
1125 if let Some(ref mut interval) = http_check_interval {
1126 interval.tick().await;
1127 } else {
1128 std::future::pending::<()>().await;
1129 }
1130 }, if !ready_notified && ready_http.is_some() && !http_exhausted => {
1131 if let (Some(http), Some(client)) = (&ready_http, &http_client) {
1132 match client.get(&http.url).send().await {
1133 Ok(response) if http.accepts_status(response.status().as_u16()) => {
1134 info!("daemon {id} ready: HTTP check passed (status {})", response.status());
1135 ready_notified = true;
1136 if let Some(tx) = ready_tx.take() {
1137 let _ = tx.send(Ok(()));
1138 }
1139 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1140 http_check_interval = None;
1141 http_deadline = None;
1142 stop_cmd_probe_state(&mut cmd_probe);
1143 cmd_deadline = None;
1144 port_deadline = None;
1145 output_deadline = None;
1146 if !active_port_spawned && has_port_config {
1147 active_port_spawned = true;
1148 detect_and_store_active_port(id.clone(), daemon_pid);
1149 }
1150 }
1151 Ok(response) => {
1152 trace!("daemon {id} HTTP check: status {} (not ready)", response.status());
1153 }
1154 Err(e) => {
1155 trace!("daemon {id} HTTP check failed: {e}");
1156 }
1157 }
1158 }
1159 }
1160 _ = async {
1161 if let Some(ref mut deadline) = port_deadline {
1162 deadline.await;
1163 } else {
1164 std::future::pending::<()>().await;
1165 }
1166 }, if !ready_notified && ready_port.is_some() => {
1167 port_exhausted = true;
1168 port_deadline = None;
1169 port_check_interval = None;
1170 warn!("daemon {id}: TCP port readiness check timed out");
1171 let any_remaining = any_ready_check_remaining(
1172 ready_output.as_ref(),
1173 output_exhausted,
1174 ready_port_config.as_ref(),
1175 port_exhausted,
1176 ready_http.as_ref(),
1177 http_exhausted,
1178 ready_cmd.as_ref(),
1179 cmd_exhausted,
1180 );
1181 if !any_remaining {
1182 error!("daemon {id}: all readiness checks exhausted, failing");
1183 stop_cmd_probe_state(&mut cmd_probe);
1184 if let Some(tx) = ready_tx.take() {
1185 let _ = tx.send(Err(Some(124)));
1186 }
1187 let stop_cfg = opts.stop_signal.unwrap_or_default();
1188 let _ = PROCS.kill_process_group_async(daemon_pid, stop_cfg.signal.into(), stop_cfg.timeout).await;
1189 break;
1190 }
1191 }
1192 _ = async {
1193 if let Some(ref mut interval) = port_check_interval {
1194 interval.tick().await;
1195 } else {
1196 std::future::pending::<()>().await;
1197 }
1198 }, if !ready_notified && ready_port.is_some() && !port_exhausted => {
1199 if let Some(port) = ready_port {
1200 match tokio::net::TcpStream::connect(("127.0.0.1", port)).await {
1201 Ok(_) => {
1202 info!("daemon {id} ready: TCP port {port} is listening");
1203 ready_notified = true;
1204 if let Some(tx) = ready_tx.take() {
1205 let _ = tx.send(Ok(()));
1206 }
1207 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1208 port_check_interval = None;
1210 port_deadline = None;
1211 stop_cmd_probe_state(&mut cmd_probe);
1212 http_deadline = None;
1213 cmd_deadline = None;
1214 output_deadline = None;
1215 if !active_port_spawned && has_port_config {
1216 active_port_spawned = true;
1217 if expected_port == Some(port) {
1227 let mut state_file =
1228 SUPERVISOR.state_file.lock().await;
1229 if let Some(d) = state_file.daemons.get(&id)
1230 && d.pid == Some(daemon_pid)
1231 {
1232 state_file.set_active_port(&id, port);
1233 }
1234 } else {
1235 detect_and_store_active_port(
1236 id.clone(),
1237 daemon_pid,
1238 );
1239 }
1240 }
1241 }
1242 Err(_) => {
1243 trace!("daemon {id} port check: port {port} not listening yet");
1244 }
1245 }
1246 }
1247 }
1248 _ = async {
1249 if let Some(ref mut delay) = cmd_respawn_delay {
1250 delay.await;
1251 } else {
1252 std::future::pending::<()>().await;
1253 }
1254 }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted && cmd_probe.is_none() => {
1255 if let Some(ref cmd) = ready_cmd {
1256 cmd_probe = Some(spawn_cmd_probe(&id, &cmd.run, daemon_dir.as_path()));
1257 }
1258 cmd_respawn_delay = None;
1259 }
1260 result = async {
1261 if let Some(probe) = cmd_probe.as_mut() {
1262 std::pin::Pin::new(&mut probe.result_rx).await
1263 } else {
1264 std::future::pending::<Result<Result<std::process::ExitStatus, std::io::Error>, tokio::sync::oneshot::error::RecvError>>().await
1265 }
1266 }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted => {
1267 let _ = cmd_probe.take();
1271 match result {
1272 Ok(Ok(status)) if status.success() => {
1273 info!("daemon {id} ready: readiness command succeeded");
1274 ready_notified = true;
1275 if let Some(tx) = ready_tx.take() {
1276 let _ = tx.send(Ok(()));
1277 }
1278 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1279 cmd_respawn_delay = None;
1280 cmd_deadline = None;
1281 http_deadline = None;
1282 port_deadline = None;
1283 output_deadline = None;
1284 if !active_port_spawned && has_port_config {
1285 active_port_spawned = true;
1286 detect_and_store_active_port(id.clone(), daemon_pid);
1287 }
1288 }
1289 Ok(Ok(_)) | Ok(Err(_)) | Err(_) => {
1290 trace!("daemon {id} cmd check: command not ready, will respawn");
1291 cmd_respawn_delay = Some(Box::pin(time::sleep(ready_check_interval)));
1292 }
1293 }
1294 }
1295 _ = async {
1296 if let Some(ref mut deadline) = cmd_deadline {
1297 deadline.await;
1298 } else {
1299 std::future::pending::<()>().await;
1300 }
1301 }, if !ready_notified && ready_cmd.is_some() => {
1302 cmd_exhausted = true;
1303 cmd_deadline = None;
1304 stop_cmd_probe_state(&mut cmd_probe);
1305 cmd_respawn_delay = None;
1306 warn!("daemon {id}: command readiness check timed out");
1307 let any_remaining = any_ready_check_remaining(
1308 ready_output.as_ref(),
1309 output_exhausted,
1310 ready_port_config.as_ref(),
1311 port_exhausted,
1312 ready_http.as_ref(),
1313 http_exhausted,
1314 ready_cmd.as_ref(),
1315 cmd_exhausted,
1316 );
1317 if !any_remaining {
1318 error!("daemon {id}: all readiness checks exhausted, failing");
1319 if let Some(tx) = ready_tx.take() {
1320 let _ = tx.send(Err(Some(124)));
1321 }
1322 let stop_cfg = opts.stop_signal.unwrap_or_default();
1323 let _ = PROCS.kill_process_group_async(daemon_pid, stop_cfg.signal.into(), stop_cfg.timeout).await;
1324 break;
1325 }
1326 }
1327 _ = async {
1328 if let Some(ref mut timer) = delay_timer {
1329 timer.await;
1330 } else {
1331 std::future::pending::<()>().await;
1332 }
1333 } => {
1334 if !ready_notified && ready_pattern.is_none() && ready_http.is_none() && ready_port.is_none() && ready_cmd.is_none() {
1335 if exit_status.is_some() {
1340 debug!("daemon {id} exited during ready_delay, not marking as ready");
1341 } else {
1342 PROCS.refresh_pids(&[daemon_pid]);
1345 if !PROCS.is_running(daemon_pid) {
1346 debug!("daemon {id} pid {daemon_pid} not running during ready_delay, deferring to exit handler");
1347 } else {
1348 info!("daemon {id} ready: delay elapsed");
1349 ready_notified = true;
1350 if let Some(tx) = ready_tx.take() {
1351 let _ = tx.send(Ok(()));
1352 }
1353 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1354 }
1355 }
1356 output_deadline = None;
1359 http_deadline = None;
1360 cmd_deadline = None;
1361 port_deadline = None;
1362 stop_cmd_probe_state(&mut cmd_probe);
1363 }
1364 delay_timer = None;
1366 if !active_port_spawned && has_port_config {
1367 active_port_spawned = true;
1368 detect_and_store_active_port(id.clone(), daemon_pid);
1369 }
1370 }
1371 _ = log_flush_interval.tick() => {
1372 let _ = flush_logs(&mut log_buffer);
1373 }
1374 }
1375 }
1376
1377 let pre_drain_daemon = SUPERVISOR.get_daemon(&id).await;
1391 let pre_drain_is_stopping = pre_drain_daemon
1392 .as_ref()
1393 .is_some_and(|d| d.status.is_stopped() || d.status.is_stopping());
1394
1395 let drain_deadline = tokio::time::Instant::now() + Duration::from_secs(5);
1403 loop {
1404 let now = tokio::time::Instant::now();
1405 if now >= drain_deadline {
1406 break;
1407 }
1408 let Ok(Some(line)) =
1409 tokio::time::timeout(drain_deadline - now, output_rx.recv()).await
1410 else {
1411 break;
1412 };
1413 log_buffer.push(parse_line(&line));
1414 }
1415 if let Some(handle) = flush_logs(&mut log_buffer) {
1418 let _ = handle.await;
1419 }
1420
1421 {
1423 let mut state_file = SUPERVISOR.state_file.lock().await;
1424 state_file.clear_active_port(&id);
1425 }
1426
1427 let exit_status = if let Some(status) = exit_status {
1429 status
1430 } else {
1431 match exit_rx.recv().await {
1433 Some(status) => status,
1434 None => {
1435 warn!("daemon {id} exit channel closed without receiving status");
1436 Err(std::io::Error::other("exit channel closed"))
1437 }
1438 }
1439 };
1440 let current_daemon = SUPERVISOR.get_daemon(&id).await;
1441
1442 SUPERVISOR
1447 .active_monitors
1448 .fetch_add(1, atomic::Ordering::Release);
1449 struct MonitorGuard;
1450 impl Drop for MonitorGuard {
1451 fn drop(&mut self) {
1452 SUPERVISOR
1453 .active_monitors
1454 .fetch_sub(1, atomic::Ordering::Release);
1455 SUPERVISOR.monitor_done.notify_waiters();
1456 }
1457 }
1458 let _monitor_guard = MonitorGuard;
1459 if !pre_drain_is_stopping
1464 && (current_daemon.is_none()
1465 || current_daemon.as_ref().is_some_and(|d| {
1466 d.pid != Some(pid) && !d.status.is_stopped() && !d.status.is_stopping()
1467 }))
1468 {
1469 return;
1471 }
1472 let already_stopped = current_daemon
1477 .as_ref()
1478 .is_some_and(|d| d.status.is_stopped());
1479 let is_stopping = already_stopped
1480 || pre_drain_is_stopping
1481 || current_daemon
1482 .as_ref()
1483 .is_some_and(|d| d.status.is_stopping());
1484
1485 let (exit_code, exit_reason) = match (&exit_status, is_stopping) {
1487 (Ok(status), true) => {
1488 (status.code().unwrap_or(-1), "stop")
1492 }
1493 (Ok(status), false) if status.success() => (status.code().unwrap_or(-1), "exit"),
1494 (Ok(status), false) => (status.code().unwrap_or(-1), "fail"),
1495 (Err(_), true) => {
1496 (-1, "stop")
1498 }
1499 (Err(_), false) => (-1, "fail"),
1500 };
1501
1502 if !already_stopped && !pre_drain_is_stopping {
1508 if let Ok(status) = &exit_status {
1509 info!("daemon {id} exited with status {status}");
1510 }
1511 let (new_status, last_exit_success) = match exit_reason {
1512 "stop" | "exit" => (
1513 DaemonStatus::Stopped,
1514 exit_status.as_ref().map(|s| s.success()).unwrap_or(true),
1515 ),
1516 _ => (DaemonStatus::Errored(exit_code), false),
1517 };
1518 if !SUPERVISOR
1524 .finalize_monitored_exit(
1525 &id,
1526 pid,
1527 monitor_token,
1528 new_status,
1529 Some(last_exit_success),
1530 )
1531 .await
1532 {
1533 debug!("daemon {id} exit state was not written; a successor owns the record");
1534 }
1535 }
1536
1537 let hook_extra_env = vec![
1539 ("PITCHFORK_EXIT_CODE".to_string(), exit_code.to_string()),
1540 ("PITCHFORK_EXIT_REASON".to_string(), exit_reason.to_string()),
1541 ];
1542
1543 let hooks_to_fire: Vec<HookType> = match exit_reason {
1545 "stop" => vec![HookType::OnStop, HookType::OnExit],
1546 "exit" => vec![HookType::OnExit],
1547 _ if hook_retry_count >= hook_retry.count() => {
1549 vec![HookType::OnFail, HookType::OnExit]
1550 }
1551 _ => vec![],
1552 };
1553
1554 for hook_type in hooks_to_fire {
1555 fire_hook(
1556 hook_type,
1557 id.clone(),
1558 daemon_dir.clone(),
1559 hook_retry_count,
1560 hook_daemon_env.clone(),
1561 hook_extra_env.clone(),
1562 )
1563 .await;
1564 }
1565 });
1566
1567 if let Some(ready_rx) = ready_rx {
1569 match ready_rx.await {
1570 Ok(Ok(())) => {
1571 info!("daemon {id} is ready");
1572 Ok(IpcResponse::DaemonReady { daemon })
1573 }
1574 Ok(Err(exit_code)) => {
1575 error!("daemon {id} failed before becoming ready");
1576 let last_attempt = opts.retry_count >= opts.retry.count();
1587 if using_sink && last_attempt {
1588 super::log_sink::wait_for_output(id, spawn_time, SINK_OUTPUT_TIMEOUT).await;
1589 }
1590 Ok(IpcResponse::DaemonFailedWithCode { exit_code })
1591 }
1592 Err(_) => {
1593 error!("readiness channel closed unexpectedly for daemon {id}");
1594 Ok(IpcResponse::DaemonStart { daemon })
1595 }
1596 }
1597 } else {
1598 Ok(IpcResponse::DaemonStart { daemon })
1599 }
1600 }
1601
1602 pub async fn stop(&self, id: &DaemonId) -> Result<IpcResponse> {
1604 let pitchfork_id = DaemonId::pitchfork();
1605 if *id == pitchfork_id {
1606 return Ok(IpcResponse::Error(
1607 "Cannot stop supervisor via stop command".into(),
1608 ));
1609 }
1610 info!("stopping daemon: {id}");
1611 if let Some(daemon) = self.get_daemon(id).await {
1612 trace!("daemon to stop: {daemon}");
1613 if let Some(pid) = daemon.pid {
1614 trace!("killing pid: {pid}");
1615 if PROCS.is_running(pid) {
1616 if !super::signalling_pid_is_authorized(
1622 daemon.start_time,
1623 PROCS.start_time(pid),
1624 ) {
1625 warn!(
1626 "pid {pid} recorded for daemon {id} belongs to another process now; not signalling it"
1627 );
1628 self.upsert_daemon(
1629 UpsertDaemonOpts::builder(id.clone())
1630 .set(|o| {
1631 o.pid = None;
1632 o.status = DaemonStatus::Stopped;
1633 })
1634 .build(),
1635 )
1636 .await?;
1637 return Ok(IpcResponse::DaemonWasNotRunning);
1638 }
1639
1640 self.upsert_daemon(
1642 UpsertDaemonOpts::builder(id.clone())
1643 .set(|o| {
1644 o.pid = Some(pid);
1645 o.status = DaemonStatus::Stopping;
1646 })
1647 .build(),
1648 )
1649 .await?;
1650
1651 let stop_cfg = daemon.stop_signal.unwrap_or_default();
1654 let stop_signal: i32 = stop_cfg.signal.into();
1655 if let Err(e) = PROCS
1656 .kill_process_group_async(pid, stop_signal, stop_cfg.timeout)
1657 .await
1658 {
1659 debug!("failed to kill pid {pid}: {e}");
1660 if PROCS.is_running(pid) {
1662 debug!("failed to stop pid {pid}: process still running after kill");
1664 self.upsert_daemon(
1665 UpsertDaemonOpts::builder(id.clone())
1666 .set(|o| {
1667 o.pid = Some(pid); o.status = DaemonStatus::Running;
1669 })
1670 .build(),
1671 )
1672 .await?;
1673 return Ok(IpcResponse::DaemonStopFailed {
1674 error: format!(
1675 "process {pid} still running after kill attempt: {e}"
1676 ),
1677 });
1678 }
1679 }
1680
1681 self.upsert_daemon(
1686 UpsertDaemonOpts::builder(id.clone())
1687 .set(|o| {
1688 o.pid = None;
1689 o.status = DaemonStatus::Stopped;
1690 o.last_exit_success = Some(true);
1691 })
1692 .build(),
1693 )
1694 .await?;
1695 } else {
1696 debug!("pid {pid} not running, process may have exited unexpectedly");
1697 self.upsert_daemon(
1703 UpsertDaemonOpts::builder(id.clone())
1704 .set(|o| {
1705 o.pid = None;
1706 o.status = DaemonStatus::Stopped;
1707 })
1708 .build(),
1709 )
1710 .await?;
1711 return Ok(IpcResponse::DaemonWasNotRunning);
1712 }
1713 Ok(IpcResponse::Ok)
1714 } else {
1715 debug!("daemon {id} not running");
1716 Ok(IpcResponse::DaemonNotRunning)
1717 }
1718 } else {
1719 debug!("daemon {id} not found");
1720 Ok(IpcResponse::DaemonNotFound)
1721 }
1722 }
1723}
1724
1725#[cfg(unix)]
1726fn resolve_effective_run_identity(daemon_user: Option<&str>) -> Result<RunIdentity> {
1727 let s = settings();
1728 let settings_user = s.supervisor.user.trim();
1729 let daemon_user = daemon_user.map(str::trim).filter(|user| !user.is_empty());
1730 let settings_user = (!settings_user.is_empty()).then_some(settings_user);
1731 let configured = daemon_user.or(settings_user);
1732 let current_uid = nix::unistd::Uid::effective().as_raw();
1733 let current_gid = nix::unistd::Gid::effective().as_raw();
1734 resolve_run_identity(
1735 configured,
1736 current_uid,
1737 current_gid,
1738 std::env::var("SUDO_UID").ok().as_deref(),
1739 std::env::var("SUDO_GID").ok().as_deref(),
1740 )
1741}
1742
1743#[cfg(unix)]
1744fn resolve_run_identity(
1745 configured: Option<&str>,
1746 current_uid: u32,
1747 current_gid: u32,
1748 sudo_uid: Option<&str>,
1749 sudo_gid: Option<&str>,
1750) -> Result<RunIdentity> {
1751 let current_uid = nix::unistd::Uid::from_raw(current_uid);
1752 let current_gid = nix::unistd::Gid::from_raw(current_gid);
1753 if let Some(user) = configured {
1754 let identity = resolve_configured_user(user)?;
1755 ensure_can_use_identity(user, &identity, current_uid, current_gid)?;
1756 if identity.matches(current_uid, current_gid) {
1757 return Ok(RunIdentity::Inherit);
1758 }
1759 return Ok(identity);
1760 }
1761
1762 if current_uid.is_root()
1763 && let Some(identity) = resolve_sudo_identity(sudo_uid, sudo_gid)
1764 {
1765 return Ok(identity);
1766 }
1767
1768 Ok(RunIdentity::Inherit)
1769}
1770
1771#[cfg(unix)]
1772fn resolve_configured_user(user: &str) -> Result<RunIdentity> {
1773 if user.chars().all(|c| c.is_ascii_digit()) {
1774 let uid = user
1775 .parse::<u32>()
1776 .map_err(|e| miette::miette!("invalid run user UID '{}': {}", user, e))?;
1777 let user_record = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1778 .into_diagnostic()?
1779 .ok_or_else(|| miette::miette!("run user UID '{}' does not exist", user))?;
1780 return run_identity_from_user_record(user_record);
1781 }
1782
1783 let user_record = nix::unistd::User::from_name(user)
1784 .into_diagnostic()?
1785 .ok_or_else(|| miette::miette!("run user '{}' does not exist", user))?;
1786 run_identity_from_user_record(user_record)
1787}
1788
1789#[cfg(unix)]
1790fn run_identity_from_user_record(user: nix::unistd::User) -> Result<RunIdentity> {
1791 let username = CString::new(user.name)
1792 .map_err(|e| miette::miette!("run user name contains an interior nul byte: {}", e))?;
1793 Ok(RunIdentity::Switch {
1794 uid: user.uid,
1795 gid: user.gid,
1796 username: Some(username),
1797 })
1798}
1799
1800#[cfg(unix)]
1801fn run_identity_from_raw_ids(uid: u32, gid: u32, username: Option<CString>) -> RunIdentity {
1802 RunIdentity::Switch {
1803 uid: nix::unistd::Uid::from_raw(uid),
1804 gid: nix::unistd::Gid::from_raw(gid),
1805 username,
1806 }
1807}
1808
1809#[cfg(unix)]
1810fn resolve_sudo_identity(sudo_uid: Option<&str>, sudo_gid: Option<&str>) -> Option<RunIdentity> {
1811 let uid = sudo_uid?.parse::<u32>().ok()?;
1812 let gid = sudo_gid?.parse::<u32>().ok()?;
1813 let username = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1814 .ok()
1815 .flatten()
1816 .and_then(|u| CString::new(u.name).ok());
1817 Some(run_identity_from_raw_ids(uid, gid, username))
1818}
1819
1820#[cfg(unix)]
1821fn ensure_can_use_identity(
1822 configured_user: &str,
1823 identity: &RunIdentity,
1824 current_uid: nix::unistd::Uid,
1825 current_gid: nix::unistd::Gid,
1826) -> Result<()> {
1827 let RunIdentity::Switch { uid, gid, .. } = identity else {
1828 return Ok(());
1829 };
1830 if *uid == current_uid && *gid == current_gid {
1831 return Ok(());
1832 }
1833 if current_uid.is_root() {
1834 return Ok(());
1835 }
1836 Err(miette::miette!(
1837 "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.",
1838 configured_user,
1839 current_uid.as_raw(),
1840 current_gid.as_raw(),
1841 uid.as_raw(),
1842 gid.as_raw()
1843 ))
1844}
1845
1846#[cfg(unix)]
1847fn apply_run_identity(identity: &RunIdentity) -> std::io::Result<()> {
1848 let RunIdentity::Switch { uid, gid, username } = identity else {
1849 return Ok(());
1850 };
1851 if let Some(username) = username {
1852 initgroups_for_user(username, *gid)?;
1853 } else {
1854 setgroups_to_primary(*gid)?;
1855 }
1856 nix::unistd::setgid(*gid).map_err(nix_to_io_error)?;
1857 nix::unistd::setuid(*uid).map_err(nix_to_io_error)?;
1858 Ok(())
1859}
1860
1861#[cfg(unix)]
1862impl RunIdentity {
1863 fn matches(&self, uid: nix::unistd::Uid, gid: nix::unistd::Gid) -> bool {
1864 matches!(self, RunIdentity::Switch { uid: u, gid: g, .. } if *u == uid && *g == gid)
1865 }
1866}
1867
1868#[cfg(unix)]
1869fn setgroups_to_primary(gid: nix::unistd::Gid) -> std::io::Result<()> {
1870 let groups = [gid.as_raw() as libc::gid_t];
1871 #[cfg(any(target_os = "linux", target_os = "android"))]
1872 let group_count = groups.len();
1873 #[cfg(not(any(target_os = "linux", target_os = "android")))]
1874 let group_count = groups.len() as libc::c_int;
1875 let rc = unsafe { libc::setgroups(group_count, groups.as_ptr()) };
1876 if rc == -1 {
1877 Err(std::io::Error::last_os_error())
1878 } else {
1879 Ok(())
1880 }
1881}
1882
1883#[cfg(unix)]
1884fn initgroups_for_user(username: &CString, gid: nix::unistd::Gid) -> std::io::Result<()> {
1885 let gid = gid.as_raw();
1886 #[cfg(any(
1887 target_os = "macos",
1888 target_os = "ios",
1889 target_os = "tvos",
1890 target_os = "watchos"
1891 ))]
1892 let base_gid = i32::try_from(gid)
1893 .map_err(|_| std::io::Error::other(format!("gid {gid} is out of range")))?;
1894
1895 #[cfg(not(any(
1896 target_os = "macos",
1897 target_os = "ios",
1898 target_os = "tvos",
1899 target_os = "watchos"
1900 )))]
1901 let base_gid = gid as libc::gid_t;
1902
1903 let rc = unsafe { libc::initgroups(username.as_ptr(), base_gid) };
1906 if rc == -1 {
1907 Err(std::io::Error::last_os_error())
1908 } else {
1909 Ok(())
1910 }
1911}
1912
1913#[cfg(unix)]
1914fn nix_to_io_error(err: nix::errno::Errno) -> std::io::Error {
1915 std::io::Error::from_raw_os_error(err as i32)
1916}
1917
1918async fn check_ports_available(
1925 expected_ports: &[u16],
1926 auto_bump: bool,
1927 max_attempts: u32,
1928) -> Result<Vec<u16>> {
1929 if expected_ports.is_empty() {
1930 return Ok(Vec::new());
1931 }
1932
1933 for bump_offset in 0..=max_attempts {
1934 let candidate_ports: Vec<u16> = expected_ports
1936 .iter()
1937 .map(|&p| p.wrapping_add(bump_offset as u16))
1938 .collect();
1939
1940 let mut all_available = true;
1942 let mut conflicting_port = None;
1943
1944 for &port in &candidate_ports {
1945 if port == 0 {
1948 continue;
1949 }
1950
1951 if is_port_in_use(port).await {
1965 all_available = false;
1966 conflicting_port = Some(port);
1967 break;
1968 }
1969 }
1970
1971 if all_available {
1972 if candidate_ports.contains(&0) && !expected_ports.contains(&0) {
1976 return Err(PortError::NoAvailablePort {
1977 start_port: expected_ports[0],
1978 attempts: bump_offset + 1,
1979 }
1980 .into());
1981 }
1982 if bump_offset > 0 {
1983 info!("ports {expected_ports:?} bumped by {bump_offset} to {candidate_ports:?}");
1984 }
1985 return Ok(candidate_ports);
1986 }
1987
1988 if bump_offset == 0
1990 && !auto_bump
1991 && let Some(port) = conflicting_port
1992 {
1993 let (pid, process) = identify_port_owner(port).await;
1994 return Err(PortError::InUse { port, process, pid }.into());
1995 }
1996 }
1997
1998 Err(PortError::NoAvailablePort {
2000 start_port: expected_ports[0],
2001 attempts: max_attempts + 1,
2002 }
2003 .into())
2004}
2005
2006async fn is_port_in_use(port: u16) -> bool {
2012 tokio::task::spawn_blocking(move || {
2013 for &addr in &["0.0.0.0", "127.0.0.1", "::1"] {
2014 match std::net::TcpListener::bind((addr, port)) {
2015 Ok(listener) => drop(listener),
2016 Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => return true,
2017 Err(_) => continue,
2018 }
2019 }
2020 false
2021 })
2022 .await
2023 .unwrap_or(false)
2024}
2025
2026async fn identify_port_owner(port: u16) -> (u32, String) {
2031 tokio::task::spawn_blocking(move || {
2032 listeners::get_all()
2033 .ok()
2034 .and_then(|list| {
2035 list.into_iter()
2036 .find(|l| l.socket.port() == port)
2037 .map(|l| (l.process.pid, l.process.name))
2038 })
2039 .unwrap_or((0, "unknown".to_string()))
2040 })
2041 .await
2042 .unwrap_or((0, "unknown".to_string()))
2043}
2044
2045async fn detect_port_conflict(port: u16) -> Option<(u32, String)> {
2050 if !is_port_in_use(port).await {
2051 return None;
2052 }
2053 Some(identify_port_owner(port).await)
2054}
2055
2056fn detect_and_store_active_port(id: DaemonId, pid: u32) {
2071 tokio::spawn(async move {
2072 for delay_ms in [500u64, 1000, 2000, 4000] {
2076 tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
2077
2078 let expected_port: Option<u16> = {
2081 let state_file = SUPERVISOR.state_file.lock().await;
2082 match state_file.daemons.get(&id) {
2083 Some(d) if d.pid.is_none() => {
2084 debug!("daemon {id}: aborting active_port detection — process exited");
2085 return;
2086 }
2087 Some(d) => d
2088 .port
2089 .as_ref()
2090 .and_then(|p| p.expect.first().copied())
2091 .filter(|&p| p > 0),
2092 None => None,
2093 }
2094 };
2095
2096 let active_port = tokio::task::spawn_blocking(move || {
2097 let listeners = listeners::get_all().ok()?;
2098
2099 PROCS.refresh_processes();
2101
2102 let descendant_pids: std::collections::HashSet<u32> = PROCS
2103 .all_children(pid)
2104 .into_iter()
2105 .chain(std::iter::once(pid))
2106 .collect();
2107
2108 let process_ports: Vec<u16> = listeners
2109 .into_iter()
2110 .filter(|listener| descendant_pids.contains(&listener.process.pid))
2111 .map(|listener| listener.socket.port())
2112 .filter(|&port| port > 0)
2113 .collect();
2114
2115 if process_ports.is_empty() {
2116 return None;
2117 }
2118
2119 if let Some(ep) = expected_port
2122 && process_ports.contains(&ep)
2123 {
2124 return Some(ep);
2125 }
2126
2127 process_ports.into_iter().next()
2133 })
2134 .await
2135 .ok()
2136 .flatten();
2137
2138 if let Some(port) = active_port {
2139 debug!("daemon {id} active_port detected: {port}");
2140 let mut state_file = SUPERVISOR.state_file.lock().await;
2141 if let Some(d) = state_file.daemons.get(&id) {
2142 if d.pid == Some(pid) {
2146 state_file.set_active_port(&id, port);
2147 } else {
2148 debug!(
2149 "daemon {id}: skipping active_port write — PID mismatch \
2150 (expected {pid}, current {:?})",
2151 d.pid
2152 );
2153 return;
2154 }
2155 }
2156 return;
2157 }
2158
2159 debug!(
2160 "daemon {id}: no active port detected for pid {pid} or its descendants (will retry)"
2161 );
2162 }
2163
2164 debug!(
2165 "daemon {id}: active port detection exhausted all retries for pid {pid} and its descendants"
2166 );
2167 });
2168}
2169
2170fn is_daemon_slug_target(id: &DaemonId) -> bool {
2178 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
2182 slugs.iter().any(|(slug, entry)| {
2183 let daemon_name = entry.daemon.as_deref().unwrap_or(slug);
2184 id.name() == daemon_name
2185 })
2186}
2187
2188#[cfg(all(test, unix))]
2189mod tests {
2190 use super::*;
2191
2192 #[test]
2193 fn test_resolve_run_identity_empty_without_sudo() {
2194 let identity = resolve_run_identity(None, 501, 20, None, None).unwrap();
2195 assert_eq!(identity, RunIdentity::Inherit);
2196 }
2197
2198 #[test]
2199 fn test_resolve_run_identity_sudo_fallback() {
2200 let identity = resolve_run_identity(None, 0, 0, Some("501"), Some("20")).unwrap();
2201 let RunIdentity::Switch { uid, gid, .. } = identity else {
2202 panic!("expected identity switch");
2203 };
2204 assert_eq!(uid.as_raw(), 501);
2205 assert_eq!(gid.as_raw(), 20);
2206 }
2207
2208 #[test]
2209 fn test_resolve_run_identity_ignores_stale_sudo_when_not_root() {
2210 let identity = resolve_run_identity(None, 501, 20, Some("0"), Some("0")).unwrap();
2211 assert_eq!(identity, RunIdentity::Inherit);
2212 }
2213
2214 #[test]
2215 fn test_resolve_configured_user_root_name() {
2216 let identity = resolve_configured_user("root").unwrap();
2217 let RunIdentity::Switch { uid, username, .. } = identity else {
2218 panic!("expected identity switch");
2219 };
2220 assert_eq!(uid.as_raw(), 0);
2221 assert_eq!(
2222 username.as_deref().and_then(|s| s.to_str().ok()),
2223 Some("root")
2224 );
2225 }
2226
2227 #[test]
2228 fn test_resolve_configured_user_root_uid() {
2229 let identity = resolve_configured_user("0").unwrap();
2230 let RunIdentity::Switch { uid, username, .. } = identity else {
2231 panic!("expected identity switch");
2232 };
2233 assert_eq!(uid.as_raw(), 0);
2234 assert_eq!(
2235 username.as_deref().and_then(|s| s.to_str().ok()),
2236 Some("root")
2237 );
2238 }
2239
2240 #[test]
2241 fn test_resolve_configured_user_missing_user_fails() {
2242 let err = resolve_configured_user("pitchfork-user-that-should-not-exist")
2243 .unwrap_err()
2244 .to_string();
2245 assert!(err.contains("does not exist"));
2246 }
2247
2248 #[test]
2249 fn test_resolve_run_identity_requires_root_for_user_switch() {
2250 let err = resolve_run_identity(Some("root"), 501, 20, None, None)
2251 .unwrap_err()
2252 .to_string();
2253 assert!(err.contains("Restart the supervisor with sudo"));
2254 }
2255
2256 #[test]
2257 fn test_resolve_run_identity_same_user_is_noop() {
2258 let identity = resolve_run_identity(Some("root"), 0, 0, Some("501"), Some("20")).unwrap();
2259 assert_eq!(identity, RunIdentity::Inherit);
2260 }
2261}
2262
2263fn inject_proxy_env(cmd: &mut tokio::process::Command, slug: &Option<String>) {
2272 let s = crate::settings::settings();
2273 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
2274
2275 if should_force_loopback_host(slug) && !lan_enabled {
2276 cmd.env("HOST", "127.0.0.1");
2279 }
2280
2281 if let Some(url) = build_pitchfork_url(slug, &s) {
2283 cmd.env("PITCHFORK_URL", &url);
2284 }
2285
2286 if s.proxy.enable && s.proxy.https {
2288 let ca_path = if s.proxy.tls_cert.is_empty() {
2289 crate::env::PITCHFORK_STATE_DIR.join("proxy").join("ca.pem")
2290 } else {
2291 std::path::PathBuf::from(&s.proxy.tls_cert)
2292 };
2293 if ca_path.exists() {
2294 cmd.env("NODE_EXTRA_CA_CERTS", ca_path.to_string_lossy().to_string());
2295 }
2296 }
2297
2298 if s.proxy.enable {
2300 let tld = if lan_enabled { "local" } else { &s.proxy.tld };
2301 cmd.env("__VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS", format!(".{tld}"));
2302 }
2303
2304 if lan_enabled {
2306 cmd.env("PITCHFORK_LAN", "1");
2307 }
2308}
2309
2310fn should_force_loopback_host(slug: &Option<String>) -> bool {
2311 let Some(slug) = slug.as_deref() else {
2312 return false;
2313 };
2314
2315 let s = crate::settings::settings();
2316 if !s.proxy.enable {
2317 return false;
2318 }
2319
2320 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
2321 slugs.contains_key(slug)
2322}
2323
2324fn build_pitchfork_url(slug: &Option<String>, s: &crate::settings::Settings) -> Option<String> {
2328 let slug = slug.as_ref()?;
2329 if !s.proxy.enable {
2330 return None;
2331 }
2332 let scheme = if s.proxy.https { "https" } else { "http" };
2333 let port = u16::try_from(s.proxy.port).ok().filter(|&p| p > 0)?;
2334 let port_suffix = if (scheme == "https" && port == 443) || (scheme == "http" && port == 80) {
2335 String::new()
2336 } else {
2337 format!(":{port}")
2338 };
2339 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
2340 let tld = if lan_enabled { "local" } else { &s.proxy.tld };
2341 Some(format!("{scheme}://{slug}.{tld}{port_suffix}",))
2342}
2343
2344#[cfg(test)]
2345mod ready_check_tests {
2346 use super::*;
2347 use std::time::Duration;
2348
2349 #[test]
2350 fn any_ready_check_remaining_prefers_unbounded_checks() {
2351 let http = ReadyHttp::new("http://localhost/health");
2352 let cmd = ReadyCmd::new("true");
2353
2354 assert!(any_ready_check_remaining(
2355 None,
2356 false,
2357 None,
2358 false,
2359 Some(&http),
2360 false,
2361 None,
2362 false
2363 ));
2364 assert!(any_ready_check_remaining(
2365 None,
2366 false,
2367 None,
2368 false,
2369 None,
2370 false,
2371 Some(&cmd),
2372 false
2373 ));
2374 assert!(any_ready_check_remaining(
2375 None,
2376 false,
2377 Some(&ReadyPort::new(8080)),
2378 false,
2379 Some(&http),
2380 true,
2381 Some(&cmd),
2382 true
2383 ));
2384 }
2385
2386 #[test]
2387 fn any_ready_check_remaining_exhausted_timed_checks() {
2388 let http = ReadyHttp {
2389 url: "http://localhost/health".to_string(),
2390 status: vec![],
2391 timeout: Some(Duration::from_secs(5)),
2392 };
2393 let cmd = ReadyCmd {
2394 run: "true".to_string(),
2395 timeout: Some(Duration::from_secs(5)),
2396 };
2397
2398 assert!(any_ready_check_remaining(
2399 None,
2400 false,
2401 None,
2402 false,
2403 Some(&http),
2404 false,
2405 Some(&cmd),
2406 false
2407 ));
2408 assert!(!any_ready_check_remaining(
2409 None,
2410 false,
2411 None,
2412 false,
2413 Some(&http),
2414 true,
2415 Some(&cmd),
2416 true
2417 ));
2418 }
2419
2420 #[tokio::test]
2421 async fn spawn_cmd_probe_reports_success() {
2422 let id = DaemonId::new("global", "probe-test");
2423 let probe = spawn_cmd_probe(&id, "true", &std::env::temp_dir());
2424 let status = probe.result_rx.await.unwrap().unwrap();
2425 assert!(status.success());
2426 }
2427
2428 #[tokio::test]
2429 async fn spawn_cmd_probe_stops_on_request() {
2430 let id = DaemonId::new("global", "probe-test");
2431 let probe = spawn_cmd_probe(&id, "sleep 30", &std::env::temp_dir());
2432 let CmdProbe {
2433 cancel_tx,
2434 result_rx,
2435 } = probe;
2436 let _ = cancel_tx.send(());
2437 let status = result_rx.await.unwrap().unwrap();
2438 assert!(!status.success());
2439 }
2440}