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 mut command = Shell::default_for_platform().command(cmd);
85 command
86 .current_dir(dir)
87 .stdout(std::process::Stdio::null())
88 .stderr(std::process::Stdio::null())
89 .kill_on_drop(true);
90 let mut child = match command.spawn() {
91 Ok(child) => child,
92 Err(e) => {
93 warn!("daemon {id}: failed to spawn readiness command probe: {e}");
94 let (cancel_tx, _) = tokio::sync::oneshot::channel();
98 let (_, result_rx) = tokio::sync::oneshot::channel();
99 return CmdProbe {
100 cancel_tx,
101 result_rx,
102 };
103 }
104 };
105
106 let (cancel_tx, mut cancel_rx) = tokio::sync::oneshot::channel();
107 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
108
109 tokio::spawn(async move {
110 let status = tokio::select! {
111 status = child.wait() => status,
112 _ = &mut cancel_rx => {
113 let mut child = child;
114 let _ = child.kill().await;
115 child.wait().await
116 }
117 };
118 let _ = result_tx.send(status);
119 });
120
121 CmdProbe {
122 cancel_tx,
123 result_rx,
124 }
125}
126
127fn stop_cmd_probe_state(probe: &mut Option<CmdProbe>) {
129 if let Some(p) = probe.take() {
130 let _ = p.cancel_tx.send(());
131 }
132}
133
134#[allow(clippy::too_many_arguments)]
139fn any_ready_check_remaining(
140 ready_output: Option<&ReadyOutput>,
141 output_exhausted: bool,
142 ready_port: Option<&ReadyPort>,
143 port_exhausted: bool,
144 ready_http: Option<&ReadyHttp>,
145 http_exhausted: bool,
146 ready_cmd: Option<&ReadyCmd>,
147 cmd_exhausted: bool,
148) -> bool {
149 ready_output.is_some_and(|o| o.timeout.is_none() || !output_exhausted)
150 || ready_port.is_some_and(|p| p.timeout.is_none() || !port_exhausted)
151 || ready_http.is_some_and(|h| h.timeout.is_none() || !http_exhausted)
152 || ready_cmd.is_some_and(|c| c.timeout.is_none() || !cmd_exhausted)
153}
154
155impl Supervisor {
156 pub async fn run(&self, opts: RunOptions) -> Result<IpcResponse> {
158 let id = &opts.id;
159 let cmd = opts.cmd.clone();
160
161 {
163 let mut pending = self.pending_autostops.lock().await;
164 if pending.remove(id).is_some() {
165 info!("cleared pending autostop for {id} (daemon starting)");
166 }
167 }
168
169 let daemon = self.get_daemon(id).await;
170 if let Some(daemon) = daemon {
171 if !daemon.status.is_stopping()
174 && !daemon.status.is_stopped()
175 && let Some(pid) = daemon.pid
176 {
177 if opts.force {
178 self.stop(id).await?;
179 info!("run: stop completed for daemon {id}");
180 } else {
181 warn!("daemon {id} already running with pid {pid}");
182 return Ok(IpcResponse::DaemonAlreadyRunning);
183 }
184 }
185 }
186
187 if opts.wait_ready && opts.retry.count() > 0 {
189 let max_attempts = opts.retry.count().saturating_add(1);
191 for attempt in 0..max_attempts {
192 let mut retry_opts = opts.clone();
193 retry_opts.retry_count = attempt;
194 retry_opts.cmd = cmd.clone();
195
196 let result = self.run_once(retry_opts).await?;
197
198 match result {
199 IpcResponse::DaemonReady { daemon } => {
200 return Ok(IpcResponse::DaemonReady { daemon });
201 }
202 IpcResponse::DaemonFailedWithCode { exit_code } => {
203 if attempt < opts.retry.count() {
204 let backoff_secs = 2u64.saturating_pow(attempt).min(3600);
205 info!(
206 "daemon {id} failed (attempt {}/{}), retrying in {}s",
207 attempt + 1,
208 max_attempts,
209 backoff_secs
210 );
211 fire_hook(
212 HookType::OnRetry,
213 id.clone(),
214 opts.dir.0.clone(),
215 attempt + 1,
216 opts.env.clone(),
217 vec![],
218 )
219 .await;
220 time::sleep(Duration::from_secs(backoff_secs)).await;
221 continue;
222 } else {
223 info!("daemon {id} failed after {max_attempts} attempts");
224 return Ok(IpcResponse::DaemonFailedWithCode { exit_code });
225 }
226 }
227 other => return Ok(other),
228 }
229 }
230 }
231
232 self.run_once(opts).await
234 }
235
236 pub(crate) async fn run_once(&self, opts: RunOptions) -> Result<IpcResponse> {
238 let id = &opts.id;
239 let original_cmd = opts.cmd.clone(); let (ready_tx, ready_rx) = if opts.wait_ready {
243 let (tx, rx) = oneshot::channel();
244 (Some(tx), Some(rx))
245 } else {
246 (None, None)
247 };
248
249 let expected_ports = opts
251 .port
252 .as_ref()
253 .map(|p| p.expect.clone())
254 .unwrap_or_default();
255 let (resolved_ports, effective_ready_port) = if !expected_ports.is_empty() {
256 let port_cfg = opts.port.as_ref().unwrap();
257 match check_ports_available(
258 &expected_ports,
259 port_cfg.auto_bump(),
260 port_cfg.max_bump_attempts(),
261 )
262 .await
263 {
264 Ok(resolved) => {
265 let ready_port = if let Some(configured_port) =
266 opts.ready_port.as_ref().and_then(|p| p.as_port())
267 {
268 let bump_offset = resolved
270 .first()
271 .unwrap_or(&0)
272 .saturating_sub(*expected_ports.first().unwrap_or(&0));
273 if expected_ports.contains(&configured_port) && bump_offset > 0 {
274 configured_port
275 .checked_add(bump_offset)
276 .or(Some(configured_port))
277 } else {
278 Some(configured_port)
279 }
280 } else if opts.ready_output.is_none()
281 && opts.ready_http.is_none()
282 && opts.ready_cmd.is_none()
283 && opts.ready_delay.is_none()
284 {
285 resolved.first().copied().filter(|&p| p != 0)
289 } else {
290 None
294 };
295 info!("daemon {id}: ports {expected_ports:?} resolved to {resolved:?}");
296 (resolved, ready_port)
297 }
298 Err(e) => {
299 error!("daemon {id}: port check failed: {e}");
300 if let Some(port_error) = e.downcast_ref::<PortError>() {
302 match port_error {
303 PortError::InUse { port, process, pid } => {
304 return Ok(IpcResponse::PortConflict {
305 port: *port,
306 process: process.clone(),
307 pid: *pid,
308 });
309 }
310 PortError::NoAvailablePort {
311 start_port,
312 attempts,
313 } => {
314 return Ok(IpcResponse::NoAvailablePort {
315 start_port: *start_port,
316 attempts: *attempts,
317 });
318 }
319 }
320 }
321 return Ok(IpcResponse::DaemonFailed {
322 error: e.to_string(),
323 });
324 }
325 }
326 } else {
327 if let Some(port) = opts.ready_port.as_ref().and_then(|p| p.as_port()) {
333 if port > 0 {
334 if let Some((pid, process)) = detect_port_conflict(port).await {
335 return Ok(IpcResponse::PortConflict { port, process, pid });
336 }
337 }
338 }
339 (
340 Vec::new(),
341 opts.ready_port.as_ref().and_then(|p| p.as_port()),
342 )
343 };
344
345 let shell_setting = settings().general.shell.clone();
349 let shell_parts = match shell_words::split(&shell_setting) {
350 Ok(parts) if !parts.is_empty() => parts,
351 Ok(_) => {
352 return Ok(IpcResponse::DaemonFailed {
353 error: "general.shell setting is empty".to_string(),
354 });
355 }
356 Err(e) => {
357 return Ok(IpcResponse::DaemonFailed {
358 error: format!("failed to parse general.shell setting {shell_setting:?}: {e}"),
359 });
360 }
361 };
362 let (shell_program, shell_args) = shell_parts.split_first().unwrap();
363
364 let run_script = opts
369 .run
370 .clone()
371 .unwrap_or_else(|| shell_words::join(&original_cmd));
372
373 let (program, args) = if opts.mise.unwrap_or(settings().general.mise) {
374 match settings().resolve_mise_bin() {
375 Some(mise_bin) => {
376 let mise_bin_str = mise_bin.to_string_lossy().to_string();
377 info!("daemon {id}: wrapping command with mise ({mise_bin_str})");
378 let mut args = vec!["x".to_string(), "--".to_string()];
379 args.push(shell_program.clone());
380 args.extend(shell_args.iter().cloned());
381 args.push(run_script);
382 (mise_bin_str, args)
383 }
384 None => {
385 warn!("daemon {id}: mise=true but mise binary not found, running without mise");
386 let mut args: Vec<String> = shell_args.to_vec();
387 args.push(run_script);
388 (shell_program.clone(), args)
389 }
390 }
391 } else {
392 let mut args: Vec<String> = shell_args.to_vec();
393 args.push(run_script);
394 (shell_program.clone(), args)
395 };
396 #[cfg(unix)]
397 let run_identity = match resolve_effective_run_identity(opts.user.as_deref()) {
398 Ok(identity) => identity,
399 Err(e) => {
400 return Ok(IpcResponse::DaemonFailed {
401 error: e.to_string(),
402 });
403 }
404 };
405 info!("run: spawning daemon {id} with {program} {args:?}");
406
407 #[cfg(unix)]
409 let pty_pair = if opts.pty.unwrap_or(false) {
410 match super::pty::openpty() {
411 Ok(pair) => {
412 info!("daemon {id}: allocated PTY (pty = true)");
413 Some(pair)
414 }
415 Err(e) => {
416 warn!("daemon {id}: failed to allocate PTY, falling back to pipes: {e}");
417 None
418 }
419 }
420 } else {
421 None
422 };
423
424 let mut cmd = tokio::process::Command::new(&program);
425
426 #[cfg(unix)]
427 if let Some(ref pair) = pty_pair {
428 let slave_file = std::fs::File::from(
432 pair.slave
433 .try_clone()
434 .map_err(|e| miette::miette!("failed to dup slave PTY fd: {e}"))?,
435 );
436 cmd.stdin(std::process::Stdio::from(slave_file.try_clone().map_err(
437 |e| miette::miette!("failed to clone slave PTY fd for stdin: {e}"),
438 )?));
439 cmd.stdout(std::process::Stdio::from(slave_file.try_clone().map_err(
440 |e| miette::miette!("failed to clone slave PTY fd for stdout: {e}"),
441 )?));
442 cmd.stderr(std::process::Stdio::from(slave_file));
443 } else {
444 cmd.stdout(std::process::Stdio::piped())
445 .stderr(std::process::Stdio::piped());
446 }
447
448 #[cfg(not(unix))]
449 {
450 cmd.stdout(std::process::Stdio::piped())
451 .stderr(std::process::Stdio::piped());
452 }
453
454 cmd.args(&args).current_dir(&opts.dir);
455
456 #[cfg(unix)]
457 if pty_pair.is_none() {
458 cmd.stdin(std::process::Stdio::null());
459 }
460
461 #[cfg(not(unix))]
462 cmd.stdin(std::process::Stdio::null());
463
464 if let Some(ref path) = *env::ORIGINAL_PATH {
466 cmd.env("PATH", path);
467 }
468
469 if let Some(ref env_vars) = opts.env {
471 cmd.envs(env_vars);
472 }
473
474 cmd.env("PITCHFORK_DAEMON_ID", id.qualified());
476 cmd.env("PITCHFORK_DAEMON_NAMESPACE", id.namespace());
477 cmd.env("PITCHFORK_RETRY_COUNT", opts.retry_count.to_string());
478
479 if !resolved_ports.is_empty() {
481 cmd.env("PORT", resolved_ports[0].to_string());
485 for (i, port) in resolved_ports.iter().enumerate() {
487 cmd.env(format!("PORT{i}"), port.to_string());
488 }
489 }
490
491 inject_proxy_env(&mut cmd, &opts.slug);
493
494 #[cfg(unix)]
495 {
496 let run_identity = run_identity.clone();
497 let use_pty = pty_pair.is_some();
498 unsafe {
499 cmd.pre_exec(move || {
500 nix::unistd::setsid().map_err(nix_to_io_error)?;
501
502 if use_pty {
506 let ret = libc::ioctl(0, libc::TIOCSCTTY as libc::c_ulong, 0);
507 if ret < 0 {
508 #[cfg(target_os = "linux")]
511 eprintln!(
512 "pitchfork: TIOCSCTTY failed: {}",
513 std::io::Error::last_os_error()
514 );
515 }
516 }
517
518 apply_run_identity(&run_identity)?;
519 Ok(())
520 });
521 }
522 }
523
524 let mut child = cmd.spawn().into_diagnostic()?;
525 let pid = match child.id() {
526 Some(p) => p,
527 None => {
528 warn!("Daemon {id} exited before PID could be captured");
529 return Ok(IpcResponse::DaemonFailed {
530 error: "Process exited immediately".to_string(),
531 });
532 }
533 };
534 info!("started daemon {id} with pid {pid}");
535 PROCS.refresh_pids(&[pid]);
536 let daemon = self
537 .upsert_daemon(
538 UpsertDaemonOpts::from_run_options(&opts, DaemonStatus::Running)
539 .set(|o| {
540 o.pid = Some(pid);
541 o.cmd = Some(original_cmd);
542 o.ready_port = effective_ready_port.map(|p| ReadyPort {
543 port: Some(p),
544 template: None,
545 timeout: opts.ready_port.as_ref().and_then(|rp| rp.timeout),
546 });
547 o.port = crate::config_types::PortConfig::from_parts(
548 expected_ports,
549 opts.port.as_ref().map(|p| p.bump).unwrap_or_default(),
550 );
551 o.resolved_port = resolved_ports;
552 })
553 .build(),
554 )
555 .await?;
556
557 let id_clone = id.clone();
558 let ready_delay = opts.ready_delay;
559 let ready_output = opts.ready_output.clone();
560 let ready_http = opts.ready_http.clone();
561 let ready_port = effective_ready_port;
562 let implicit_ready_port = ready_port.map(|p| ReadyPort {
563 port: Some(p),
564 template: None,
565 timeout: None,
566 });
567 let ready_port_config = opts.ready_port.clone().or(implicit_ready_port);
568 let ready_cmd = opts.ready_cmd.clone();
569 let daemon_dir = opts.dir.0.clone();
570 let hook_retry_count = opts.retry_count;
571 let hook_retry = opts.retry;
572 let hook_daemon_env = opts.env.clone();
573 let on_output_hook = opts.on_output_hook.clone();
574 let has_port_config = opts.port.as_ref().is_some_and(|p| !p.expect.is_empty())
580 || (settings().proxy.enable && is_daemon_slug_target(id));
581 let daemon_pid = pid;
582
583 #[cfg(unix)]
587 let pty_reader = pty_pair.map(|p| {
588 tokio::io::BufReader::new(tokio::fs::File::from_std(std::fs::File::from(p.master)))
589 .lines()
590 });
591 #[cfg(not(unix))]
592 let pty_reader: Option<tokio::io::Lines<tokio::io::BufReader<tokio::fs::File>>> = None;
593 let stdout_reader = if pty_reader.is_none() {
594 child
595 .stdout
596 .take()
597 .map(|s| tokio::io::BufReader::new(s).lines())
598 } else {
599 None
600 };
601 let stderr_reader = if pty_reader.is_none() {
602 child
603 .stderr
604 .take()
605 .map(|s| tokio::io::BufReader::new(s).lines())
606 } else {
607 None
608 };
609
610 if pty_reader.is_none() && (stdout_reader.is_none() || stderr_reader.is_none()) {
611 error!("Failed to capture stdout/stderr for daemon {id}");
612 }
613
614 tokio::spawn(async move {
615 let id = id_clone;
616
617 let (output_tx, mut output_rx) = tokio::sync::mpsc::channel::<String>(256);
619
620 if let Some(mut reader) = pty_reader {
621 tokio::spawn(async move {
625 while let Ok(Some(mut line)) = reader.next_line().await {
626 if line.ends_with('\r') {
628 line.pop();
629 }
630 if output_tx.send(line).await.is_err() {
631 break;
632 }
633 }
634 });
635 } else {
636 if let Some(mut stdout) = stdout_reader {
641 let tx = output_tx.clone();
642 tokio::spawn(async move {
643 while let Ok(Some(line)) = stdout.next_line().await {
644 if tx.send(line).await.is_err() {
645 break;
646 }
647 }
648 });
649 }
650 if let Some(mut stderr) = stderr_reader {
651 let tx = output_tx.clone();
652 tokio::spawn(async move {
653 while let Ok(Some(line)) = stderr.next_line().await {
654 if tx.send(line).await.is_err() {
655 break;
656 }
657 }
658 });
659 }
660 drop(output_tx);
662 }
663 let log_store = Arc::clone(&LOG_STORE);
664 let log_format = opts
665 .log_format
666 .clone()
667 .unwrap_or_else(|| crate::settings::settings().logs.log_format.clone());
668 let parse_line = move |line: &str| crate::log_parse::parse(line, &log_format);
669
670 const LOG_BATCH_SIZE: usize = 100;
671 const LOG_FLUSH_INTERVAL: Duration = Duration::from_millis(100);
672 let mut log_buffer: Vec<crate::log_parse::ParsedLog> =
673 Vec::with_capacity(LOG_BATCH_SIZE);
674 let mut log_flush_interval = tokio::time::interval(LOG_FLUSH_INTERVAL);
675 log_flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
676
677 let flush_logs =
678 |buffer: &mut Vec<crate::log_parse::ParsedLog>| -> Option<tokio::task::JoinHandle<()>> {
679 if buffer.is_empty() {
680 return None;
681 }
682 let store = Arc::clone(&log_store);
683 let id = id.clone();
684 let batch = std::mem::take(buffer);
685 Some(tokio::task::spawn_blocking(move || {
686 if let Err(e) = store.append_structured_batch(&id, &batch) {
687 error!("Failed to write batch to log for daemon {id}: {e}");
688 }
689 }))
690 };
691
692 let mut ready_notified = false;
696 let mut ready_tx = ready_tx;
697 let ready_pattern = ready_output
698 .as_ref()
699 .and_then(|o| get_or_compile_regex(&o.pattern));
700 let mut active_port_spawned = false;
702
703 let on_output_hook = match on_output_hook {
707 Some(ref hook) => match hook.validate(id.name()) {
708 Ok(()) => on_output_hook,
709 Err(e) => {
710 error!("{e}");
711 None
712 }
713 },
714 None => None,
715 };
716
717 let on_output_pattern: Option<regex::Regex> = on_output_hook
720 .as_ref()
721 .and_then(|h| h.regex.as_deref().and_then(get_or_compile_regex));
722 let on_output_debounce = on_output_hook
723 .as_ref()
724 .map(|h| h.debounce_duration())
725 .unwrap_or(Duration::from_millis(1000));
726 let mut on_output_last_fired: Option<std::time::Instant> = None;
728
729 let mut delay_timer =
730 ready_delay.map(|secs| Box::pin(time::sleep(Duration::from_secs(secs))));
731
732 let mut http_exhausted = false;
734 let mut cmd_exhausted = false;
735 let mut port_exhausted = false;
736 let mut output_exhausted = false;
737
738 let s = settings();
740 let ready_check_interval = s.supervisor_ready_check_interval();
741 let http_client_timeout = s.supervisor_http_client_timeout();
742
743 let mut output_deadline = ready_output
745 .as_ref()
746 .and_then(|o| o.timeout)
747 .map(|d| Box::pin(time::sleep(d)));
748
749 let mut http_check_interval = ready_http
751 .as_ref()
752 .map(|_| tokio::time::interval(ready_check_interval));
753 let mut http_deadline = ready_http
754 .as_ref()
755 .and_then(|h| h.timeout)
756 .map(|d| Box::pin(time::sleep(d)));
757 let http_client = ready_http.as_ref().map(|_| {
758 reqwest::Client::builder()
759 .timeout(http_client_timeout)
760 .build()
761 .unwrap_or_default()
762 });
763
764 let mut port_check_interval =
766 ready_port.map(|_| tokio::time::interval(ready_check_interval));
767 let mut port_deadline = ready_port_config
768 .as_ref()
769 .and_then(|p| p.timeout)
770 .map(|d| Box::pin(time::sleep(d)));
771
772 let mut cmd_probe: Option<CmdProbe> = None;
775 let mut cmd_respawn_delay: Option<_> = None;
776 let mut cmd_deadline = ready_cmd
777 .as_ref()
778 .and_then(|c| c.timeout)
779 .map(|d| Box::pin(time::sleep(d)));
780 if let Some(ref cmd) = ready_cmd {
781 cmd_probe = Some(spawn_cmd_probe(&id, &cmd.run, daemon_dir.as_path()));
782 }
783
784 let (exit_tx, mut exit_rx) =
786 tokio::sync::mpsc::channel::<std::io::Result<std::process::ExitStatus>>(1);
787
788 let child_pid = child.id().unwrap_or(0);
790 tokio::spawn(async move {
791 let result = child.wait().await;
792 #[cfg(all(unix, not(target_os = "linux")))]
802 let result = match &result {
803 Err(e) if e.raw_os_error() == Some(nix::libc::ECHILD) => {
804 if let Some(code) = super::REAPED_STATUSES.lock().await.remove(&child_pid) {
805 warn!(
806 "daemon pid {child_pid} wait() got ECHILD; \
807 recovered exit code {code} from zombie reaper"
808 );
809 use std::os::unix::process::ExitStatusExt;
814 if code >= 0 {
815 Ok(std::process::ExitStatus::from_raw(code << 8))
816 } else {
817 Ok(std::process::ExitStatus::from_raw((-code) & 0x7f))
819 }
820 } else {
821 warn!(
822 "daemon pid {child_pid} wait() got ECHILD but no \
823 stashed status found; reporting as error"
824 );
825 result
826 }
827 }
828 _ => result,
829 };
830 debug!("daemon pid {child_pid} wait() completed with result: {result:?}");
831 let _ = exit_tx.send(result).await;
832 });
833
834 #[allow(unused_assignments)]
835 let mut exit_status = None;
837
838 if has_port_config
844 && ready_pattern.is_none()
845 && ready_http.is_none()
846 && ready_port.is_none()
847 && ready_cmd.is_none()
848 && delay_timer.is_none()
849 {
850 active_port_spawned = true;
851 detect_and_store_active_port(id.clone(), daemon_pid);
852 }
853
854 loop {
855 select! {
856 Some(line) = output_rx.recv() => {
857 let parsed = parse_line(&line);
858 log_buffer.push(parsed);
859 if log_buffer.len() >= LOG_BATCH_SIZE {
860 let _ = flush_logs(&mut log_buffer);
861 }
862 trace!("output: {id} {line}");
863
864 let line_clean = console::strip_ansi_codes(&line).to_string();
867
868 if !ready_notified
870 && !output_exhausted
871 && let Some(ref pattern) = ready_pattern
872 && pattern.is_match(&line_clean)
873 {
874 if let Some(handle) = flush_logs(&mut log_buffer) {
879 let _ = handle.await;
880 }
881 info!("daemon {id} ready: output matched pattern");
882 ready_notified = true;
883 if let Some(tx) = ready_tx.take() {
884 let _ = tx.send(Ok(()));
885 }
886 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
887 stop_cmd_probe_state(&mut cmd_probe);
888 http_deadline = None;
889 cmd_deadline = None;
890 port_deadline = None;
891 output_deadline = None;
892 if !active_port_spawned && has_port_config {
893 active_port_spawned = true;
894 detect_and_store_active_port(id.clone(), daemon_pid);
895 }
896 }
897
898 if let Some(ref hook) = on_output_hook {
900 let matched = match (&hook.filter, &on_output_pattern) {
901 (Some(substr), _) => line_clean.contains(substr.as_str()),
902 (None, Some(re)) => re.is_match(&line_clean),
903 (None, None) => true,
904 };
905 if matched {
906 let now = std::time::Instant::now();
907 let elapsed = on_output_last_fired.map(|t| now.duration_since(t));
908 if elapsed.is_none_or(|e| e >= on_output_debounce) {
909 on_output_last_fired = Some(now);
910 hooks::fire_output_hook(id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook.run.clone(), line_clean.clone()).await;
911 }
912 }
913 }
914 tokio::task::yield_now().await;
917 }
918 Some(result) = exit_rx.recv() => {
919 exit_status = Some(result);
921 debug!("daemon {id} process exited, exit_status: {exit_status:?}");
922 if !ready_notified {
923 if let Some(tx) = ready_tx.take() {
924 let is_success = exit_status.as_ref()
926 .and_then(|r| r.as_ref().ok())
927 .map(|s| s.success())
928 .unwrap_or(false);
929
930 if is_success {
931 debug!("daemon {id} exited successfully before ready check, sending success notification");
932 let _ = tx.send(Ok(()));
933 } else {
934 let exit_code = exit_status.as_ref()
935 .and_then(|r| r.as_ref().ok())
936 .and_then(|s| s.code());
937 debug!("daemon {id} exited with failure before ready check, sending failure notification with exit_code: {exit_code:?}");
938 let _ = tx.send(Err(exit_code));
939 }
940 }
941 } else {
942 debug!("daemon {id} was already marked ready, not sending notification");
943 }
944 break;
945 },
946 _ = async {
947 if let Some(ref mut interval) = http_check_interval {
948 interval.tick().await;
949 } else {
950 std::future::pending::<()>().await;
951 }
952 }, if !ready_notified && ready_http.is_some() && !http_exhausted => {
953 if let (Some(http), Some(client)) = (&ready_http, &http_client) {
954 match client.get(&http.url).send().await {
955 Ok(response) if http.accepts_status(response.status().as_u16()) => {
956 info!("daemon {id} ready: HTTP check passed (status {})", response.status());
957 ready_notified = true;
958 if let Some(tx) = ready_tx.take() {
959 let _ = tx.send(Ok(()));
960 }
961 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
962 http_check_interval = None;
963 http_deadline = None;
964 stop_cmd_probe_state(&mut cmd_probe);
965 cmd_deadline = None;
966 port_deadline = None;
967 output_deadline = None;
968 if !active_port_spawned && has_port_config {
969 active_port_spawned = true;
970 detect_and_store_active_port(id.clone(), daemon_pid);
971 }
972 }
973 Ok(response) => {
974 trace!("daemon {id} HTTP check: status {} (not ready)", response.status());
975 }
976 Err(e) => {
977 trace!("daemon {id} HTTP check failed: {e}");
978 }
979 }
980 }
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) = port_check_interval {
1047 interval.tick().await;
1048 } else {
1049 std::future::pending::<()>().await;
1050 }
1051 }, if !ready_notified && ready_port.is_some() && !port_exhausted => {
1052 if let Some(port) = ready_port {
1053 match tokio::net::TcpStream::connect(("127.0.0.1", port)).await {
1054 Ok(_) => {
1055 info!("daemon {id} ready: TCP port {port} is listening");
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 port_check_interval = None;
1063 port_deadline = None;
1064 stop_cmd_probe_state(&mut cmd_probe);
1065 http_deadline = None;
1066 cmd_deadline = None;
1067 output_deadline = None;
1068 if !active_port_spawned && has_port_config {
1069 active_port_spawned = true;
1070 detect_and_store_active_port(id.clone(), daemon_pid);
1071 }
1072 }
1073 Err(_) => {
1074 trace!("daemon {id} port check: port {port} not listening yet");
1075 }
1076 }
1077 }
1078 }
1079 _ = async {
1080 if let Some(ref mut deadline) = port_deadline {
1081 deadline.await;
1082 } else {
1083 std::future::pending::<()>().await;
1084 }
1085 }, if !ready_notified && ready_port.is_some() => {
1086 port_exhausted = true;
1087 port_deadline = None;
1088 port_check_interval = None;
1089 warn!("daemon {id}: TCP port readiness check timed out");
1090 let any_remaining = any_ready_check_remaining(
1091 ready_output.as_ref(),
1092 output_exhausted,
1093 ready_port_config.as_ref(),
1094 port_exhausted,
1095 ready_http.as_ref(),
1096 http_exhausted,
1097 ready_cmd.as_ref(),
1098 cmd_exhausted,
1099 );
1100 if !any_remaining {
1101 error!("daemon {id}: all readiness checks exhausted, failing");
1102 stop_cmd_probe_state(&mut cmd_probe);
1103 if let Some(tx) = ready_tx.take() {
1104 let _ = tx.send(Err(Some(124)));
1105 }
1106 let stop_cfg = opts.stop_signal.unwrap_or_default();
1107 let _ = PROCS.kill_process_group_async(daemon_pid, stop_cfg.signal.into(), stop_cfg.timeout).await;
1108 break;
1109 }
1110 }
1111 _ = async {
1112 if let Some(ref mut delay) = cmd_respawn_delay {
1113 delay.await;
1114 } else {
1115 std::future::pending::<()>().await;
1116 }
1117 }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted && cmd_probe.is_none() => {
1118 if let Some(ref cmd) = ready_cmd {
1119 cmd_probe = Some(spawn_cmd_probe(&id, &cmd.run, daemon_dir.as_path()));
1120 }
1121 cmd_respawn_delay = None;
1122 }
1123 result = async {
1124 if let Some(probe) = cmd_probe.as_mut() {
1125 std::pin::Pin::new(&mut probe.result_rx).await
1126 } else {
1127 std::future::pending::<Result<Result<std::process::ExitStatus, std::io::Error>, tokio::sync::oneshot::error::RecvError>>().await
1128 }
1129 }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted => {
1130 let _ = cmd_probe.take();
1134 match result {
1135 Ok(Ok(status)) if status.success() => {
1136 info!("daemon {id} ready: readiness command succeeded");
1137 ready_notified = true;
1138 if let Some(tx) = ready_tx.take() {
1139 let _ = tx.send(Ok(()));
1140 }
1141 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1142 cmd_respawn_delay = None;
1143 cmd_deadline = None;
1144 http_deadline = None;
1145 port_deadline = None;
1146 output_deadline = None;
1147 if !active_port_spawned && has_port_config {
1148 active_port_spawned = true;
1149 detect_and_store_active_port(id.clone(), daemon_pid);
1150 }
1151 }
1152 Ok(Ok(_)) | Ok(Err(_)) | Err(_) => {
1153 trace!("daemon {id} cmd check: command not ready, will respawn");
1154 cmd_respawn_delay = Some(Box::pin(time::sleep(ready_check_interval)));
1155 }
1156 }
1157 }
1158 _ = async {
1159 if let Some(ref mut deadline) = cmd_deadline {
1160 deadline.await;
1161 } else {
1162 std::future::pending::<()>().await;
1163 }
1164 }, if !ready_notified && ready_cmd.is_some() => {
1165 cmd_exhausted = true;
1166 cmd_deadline = None;
1167 stop_cmd_probe_state(&mut cmd_probe);
1168 cmd_respawn_delay = None;
1169 warn!("daemon {id}: command readiness check timed out");
1170 let any_remaining = any_ready_check_remaining(
1171 ready_output.as_ref(),
1172 output_exhausted,
1173 ready_port_config.as_ref(),
1174 port_exhausted,
1175 ready_http.as_ref(),
1176 http_exhausted,
1177 ready_cmd.as_ref(),
1178 cmd_exhausted,
1179 );
1180 if !any_remaining {
1181 error!("daemon {id}: all readiness checks exhausted, failing");
1182 if let Some(tx) = ready_tx.take() {
1183 let _ = tx.send(Err(Some(124)));
1184 }
1185 let stop_cfg = opts.stop_signal.unwrap_or_default();
1186 let _ = PROCS.kill_process_group_async(daemon_pid, stop_cfg.signal.into(), stop_cfg.timeout).await;
1187 break;
1188 }
1189 }
1190 _ = async {
1191 if let Some(ref mut timer) = delay_timer {
1192 timer.await;
1193 } else {
1194 std::future::pending::<()>().await;
1195 }
1196 } => {
1197 if !ready_notified && ready_pattern.is_none() && ready_http.is_none() && ready_port.is_none() && ready_cmd.is_none() {
1198 info!("daemon {id} ready: delay elapsed");
1199 ready_notified = true;
1200 if let Some(tx) = ready_tx.take() {
1201 let _ = tx.send(Ok(()));
1202 }
1203 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1204 output_deadline = None;
1207 http_deadline = None;
1208 cmd_deadline = None;
1209 port_deadline = None;
1210 stop_cmd_probe_state(&mut cmd_probe);
1211 }
1212 delay_timer = None;
1214 if !active_port_spawned && has_port_config {
1215 active_port_spawned = true;
1216 detect_and_store_active_port(id.clone(), daemon_pid);
1217 }
1218 }
1219 _ = log_flush_interval.tick() => {
1220 let _ = flush_logs(&mut log_buffer);
1221 }
1222 }
1223 }
1224
1225 let drain_deadline = tokio::time::Instant::now() + Duration::from_secs(5);
1233 loop {
1234 let now = tokio::time::Instant::now();
1235 if now >= drain_deadline {
1236 break;
1237 }
1238 let Ok(Some(line)) =
1239 tokio::time::timeout(drain_deadline - now, output_rx.recv()).await
1240 else {
1241 break;
1242 };
1243 log_buffer.push(parse_line(&line));
1244 }
1245 if let Some(handle) = flush_logs(&mut log_buffer) {
1248 let _ = handle.await;
1249 }
1250
1251 {
1253 let mut state_file = SUPERVISOR.state_file.lock().await;
1254 state_file.clear_active_port(&id);
1255 }
1256
1257 let exit_status = if let Some(status) = exit_status {
1259 status
1260 } else {
1261 match exit_rx.recv().await {
1263 Some(status) => status,
1264 None => {
1265 warn!("daemon {id} exit channel closed without receiving status");
1266 Err(std::io::Error::other("exit channel closed"))
1267 }
1268 }
1269 };
1270 let current_daemon = SUPERVISOR.get_daemon(&id).await;
1271
1272 SUPERVISOR
1277 .active_monitors
1278 .fetch_add(1, atomic::Ordering::Release);
1279 struct MonitorGuard;
1280 impl Drop for MonitorGuard {
1281 fn drop(&mut self) {
1282 SUPERVISOR
1283 .active_monitors
1284 .fetch_sub(1, atomic::Ordering::Release);
1285 SUPERVISOR.monitor_done.notify_waiters();
1286 }
1287 }
1288 let _monitor_guard = MonitorGuard;
1289 if current_daemon.is_none()
1294 || current_daemon.as_ref().is_some_and(|d| {
1295 d.pid != Some(pid) && !d.status.is_stopped() && !d.status.is_stopping()
1296 })
1297 {
1298 return;
1300 }
1301 let already_stopped = current_daemon
1306 .as_ref()
1307 .is_some_and(|d| d.status.is_stopped());
1308 let is_stopping = already_stopped
1309 || current_daemon
1310 .as_ref()
1311 .is_some_and(|d| d.status.is_stopping());
1312
1313 let (exit_code, exit_reason) = match (&exit_status, is_stopping) {
1315 (Ok(status), true) => {
1316 (status.code().unwrap_or(-1), "stop")
1320 }
1321 (Ok(status), false) if status.success() => (status.code().unwrap_or(-1), "exit"),
1322 (Ok(status), false) => (status.code().unwrap_or(-1), "fail"),
1323 (Err(_), true) => {
1324 (-1, "stop")
1326 }
1327 (Err(_), false) => (-1, "fail"),
1328 };
1329
1330 if !already_stopped {
1332 if let Ok(status) = &exit_status {
1333 info!("daemon {id} exited with status {status}");
1334 }
1335 let (new_status, last_exit_success) = match exit_reason {
1336 "stop" | "exit" => (
1337 DaemonStatus::Stopped,
1338 exit_status.as_ref().map(|s| s.success()).unwrap_or(true),
1339 ),
1340 _ => (DaemonStatus::Errored(exit_code), false),
1341 };
1342 if let Err(e) = SUPERVISOR
1343 .upsert_daemon(
1344 UpsertDaemonOpts::builder(id.clone())
1345 .set(|o| {
1346 o.pid = None;
1347 o.status = new_status;
1348 o.last_exit_success = Some(last_exit_success);
1349 })
1350 .build(),
1351 )
1352 .await
1353 {
1354 error!("Failed to update daemon state for {id}: {e}");
1355 }
1356 }
1357
1358 let hook_extra_env = vec![
1360 ("PITCHFORK_EXIT_CODE".to_string(), exit_code.to_string()),
1361 ("PITCHFORK_EXIT_REASON".to_string(), exit_reason.to_string()),
1362 ];
1363
1364 let hooks_to_fire: Vec<HookType> = match exit_reason {
1366 "stop" => vec![HookType::OnStop, HookType::OnExit],
1367 "exit" => vec![HookType::OnExit],
1368 _ if hook_retry_count >= hook_retry.count() => {
1370 vec![HookType::OnFail, HookType::OnExit]
1371 }
1372 _ => vec![],
1373 };
1374
1375 for hook_type in hooks_to_fire {
1376 fire_hook(
1377 hook_type,
1378 id.clone(),
1379 daemon_dir.clone(),
1380 hook_retry_count,
1381 hook_daemon_env.clone(),
1382 hook_extra_env.clone(),
1383 )
1384 .await;
1385 }
1386 });
1387
1388 if let Some(ready_rx) = ready_rx {
1390 match ready_rx.await {
1391 Ok(Ok(())) => {
1392 info!("daemon {id} is ready");
1393 Ok(IpcResponse::DaemonReady { daemon })
1394 }
1395 Ok(Err(exit_code)) => {
1396 error!("daemon {id} failed before becoming ready");
1397 Ok(IpcResponse::DaemonFailedWithCode { exit_code })
1398 }
1399 Err(_) => {
1400 error!("readiness channel closed unexpectedly for daemon {id}");
1401 Ok(IpcResponse::DaemonStart { daemon })
1402 }
1403 }
1404 } else {
1405 Ok(IpcResponse::DaemonStart { daemon })
1406 }
1407 }
1408
1409 pub async fn stop(&self, id: &DaemonId) -> Result<IpcResponse> {
1411 let pitchfork_id = DaemonId::pitchfork();
1412 if *id == pitchfork_id {
1413 return Ok(IpcResponse::Error(
1414 "Cannot stop supervisor via stop command".into(),
1415 ));
1416 }
1417 info!("stopping daemon: {id}");
1418 if let Some(daemon) = self.get_daemon(id).await {
1419 trace!("daemon to stop: {daemon}");
1420 if let Some(pid) = daemon.pid {
1421 trace!("killing pid: {pid}");
1422 if PROCS.is_running(pid) {
1423 self.upsert_daemon(
1425 UpsertDaemonOpts::builder(id.clone())
1426 .set(|o| {
1427 o.pid = Some(pid);
1428 o.status = DaemonStatus::Stopping;
1429 })
1430 .build(),
1431 )
1432 .await?;
1433
1434 let stop_cfg = daemon.stop_signal.unwrap_or_default();
1437 let stop_signal: i32 = stop_cfg.signal.into();
1438 if let Err(e) = PROCS
1439 .kill_process_group_async(pid, stop_signal, stop_cfg.timeout)
1440 .await
1441 {
1442 debug!("failed to kill pid {pid}: {e}");
1443 if PROCS.is_running(pid) {
1445 debug!("failed to stop pid {pid}: process still running after kill");
1447 self.upsert_daemon(
1448 UpsertDaemonOpts::builder(id.clone())
1449 .set(|o| {
1450 o.pid = Some(pid); o.status = DaemonStatus::Running;
1452 })
1453 .build(),
1454 )
1455 .await?;
1456 return Ok(IpcResponse::DaemonStopFailed {
1457 error: format!(
1458 "process {pid} still running after kill attempt: {e}"
1459 ),
1460 });
1461 }
1462 }
1463
1464 self.upsert_daemon(
1469 UpsertDaemonOpts::builder(id.clone())
1470 .set(|o| {
1471 o.pid = None;
1472 o.status = DaemonStatus::Stopped;
1473 o.last_exit_success = Some(true); })
1475 .build(),
1476 )
1477 .await?;
1478 } else {
1479 debug!("pid {pid} not running, process may have exited unexpectedly");
1480 self.upsert_daemon(
1483 UpsertDaemonOpts::builder(id.clone())
1484 .set(|o| {
1485 o.pid = None;
1486 o.status = DaemonStatus::Stopped;
1487 })
1488 .build(),
1489 )
1490 .await?;
1491 return Ok(IpcResponse::DaemonWasNotRunning);
1492 }
1493 Ok(IpcResponse::Ok)
1494 } else {
1495 debug!("daemon {id} not running");
1496 Ok(IpcResponse::DaemonNotRunning)
1497 }
1498 } else {
1499 debug!("daemon {id} not found");
1500 Ok(IpcResponse::DaemonNotFound)
1501 }
1502 }
1503}
1504
1505#[cfg(unix)]
1506fn resolve_effective_run_identity(daemon_user: Option<&str>) -> Result<RunIdentity> {
1507 let s = settings();
1508 let settings_user = s.supervisor.user.trim();
1509 let daemon_user = daemon_user.map(str::trim).filter(|user| !user.is_empty());
1510 let settings_user = (!settings_user.is_empty()).then_some(settings_user);
1511 let configured = daemon_user.or(settings_user);
1512 let current_uid = nix::unistd::Uid::effective().as_raw();
1513 let current_gid = nix::unistd::Gid::effective().as_raw();
1514 resolve_run_identity(
1515 configured,
1516 current_uid,
1517 current_gid,
1518 std::env::var("SUDO_UID").ok().as_deref(),
1519 std::env::var("SUDO_GID").ok().as_deref(),
1520 )
1521}
1522
1523#[cfg(unix)]
1524fn resolve_run_identity(
1525 configured: Option<&str>,
1526 current_uid: u32,
1527 current_gid: u32,
1528 sudo_uid: Option<&str>,
1529 sudo_gid: Option<&str>,
1530) -> Result<RunIdentity> {
1531 let current_uid = nix::unistd::Uid::from_raw(current_uid);
1532 let current_gid = nix::unistd::Gid::from_raw(current_gid);
1533 if let Some(user) = configured {
1534 let identity = resolve_configured_user(user)?;
1535 ensure_can_use_identity(user, &identity, current_uid, current_gid)?;
1536 if identity.matches(current_uid, current_gid) {
1537 return Ok(RunIdentity::Inherit);
1538 }
1539 return Ok(identity);
1540 }
1541
1542 if current_uid.is_root()
1543 && let Some(identity) = resolve_sudo_identity(sudo_uid, sudo_gid)
1544 {
1545 return Ok(identity);
1546 }
1547
1548 Ok(RunIdentity::Inherit)
1549}
1550
1551#[cfg(unix)]
1552fn resolve_configured_user(user: &str) -> Result<RunIdentity> {
1553 if user.chars().all(|c| c.is_ascii_digit()) {
1554 let uid = user
1555 .parse::<u32>()
1556 .map_err(|e| miette::miette!("invalid run user UID '{}': {}", user, e))?;
1557 let user_record = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1558 .into_diagnostic()?
1559 .ok_or_else(|| miette::miette!("run user UID '{}' does not exist", user))?;
1560 return run_identity_from_user_record(user_record);
1561 }
1562
1563 let user_record = nix::unistd::User::from_name(user)
1564 .into_diagnostic()?
1565 .ok_or_else(|| miette::miette!("run user '{}' does not exist", user))?;
1566 run_identity_from_user_record(user_record)
1567}
1568
1569#[cfg(unix)]
1570fn run_identity_from_user_record(user: nix::unistd::User) -> Result<RunIdentity> {
1571 let username = CString::new(user.name)
1572 .map_err(|e| miette::miette!("run user name contains an interior nul byte: {}", e))?;
1573 Ok(RunIdentity::Switch {
1574 uid: user.uid,
1575 gid: user.gid,
1576 username: Some(username),
1577 })
1578}
1579
1580#[cfg(unix)]
1581fn run_identity_from_raw_ids(uid: u32, gid: u32, username: Option<CString>) -> RunIdentity {
1582 RunIdentity::Switch {
1583 uid: nix::unistd::Uid::from_raw(uid),
1584 gid: nix::unistd::Gid::from_raw(gid),
1585 username,
1586 }
1587}
1588
1589#[cfg(unix)]
1590fn resolve_sudo_identity(sudo_uid: Option<&str>, sudo_gid: Option<&str>) -> Option<RunIdentity> {
1591 let uid = sudo_uid?.parse::<u32>().ok()?;
1592 let gid = sudo_gid?.parse::<u32>().ok()?;
1593 let username = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1594 .ok()
1595 .flatten()
1596 .and_then(|u| CString::new(u.name).ok());
1597 Some(run_identity_from_raw_ids(uid, gid, username))
1598}
1599
1600#[cfg(unix)]
1601fn ensure_can_use_identity(
1602 configured_user: &str,
1603 identity: &RunIdentity,
1604 current_uid: nix::unistd::Uid,
1605 current_gid: nix::unistd::Gid,
1606) -> Result<()> {
1607 let RunIdentity::Switch { uid, gid, .. } = identity else {
1608 return Ok(());
1609 };
1610 if *uid == current_uid && *gid == current_gid {
1611 return Ok(());
1612 }
1613 if current_uid.is_root() {
1614 return Ok(());
1615 }
1616 Err(miette::miette!(
1617 "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.",
1618 configured_user,
1619 current_uid.as_raw(),
1620 current_gid.as_raw(),
1621 uid.as_raw(),
1622 gid.as_raw()
1623 ))
1624}
1625
1626#[cfg(unix)]
1627fn apply_run_identity(identity: &RunIdentity) -> std::io::Result<()> {
1628 let RunIdentity::Switch { uid, gid, username } = identity else {
1629 return Ok(());
1630 };
1631 if let Some(username) = username {
1632 initgroups_for_user(username, *gid)?;
1633 } else {
1634 setgroups_to_primary(*gid)?;
1635 }
1636 nix::unistd::setgid(*gid).map_err(nix_to_io_error)?;
1637 nix::unistd::setuid(*uid).map_err(nix_to_io_error)?;
1638 Ok(())
1639}
1640
1641#[cfg(unix)]
1642impl RunIdentity {
1643 fn matches(&self, uid: nix::unistd::Uid, gid: nix::unistd::Gid) -> bool {
1644 matches!(self, RunIdentity::Switch { uid: u, gid: g, .. } if *u == uid && *g == gid)
1645 }
1646}
1647
1648#[cfg(unix)]
1649fn setgroups_to_primary(gid: nix::unistd::Gid) -> std::io::Result<()> {
1650 let groups = [gid.as_raw() as libc::gid_t];
1651 #[cfg(any(target_os = "linux", target_os = "android"))]
1652 let group_count = groups.len();
1653 #[cfg(not(any(target_os = "linux", target_os = "android")))]
1654 let group_count = groups.len() as libc::c_int;
1655 let rc = unsafe { libc::setgroups(group_count, groups.as_ptr()) };
1656 if rc == -1 {
1657 Err(std::io::Error::last_os_error())
1658 } else {
1659 Ok(())
1660 }
1661}
1662
1663#[cfg(unix)]
1664fn initgroups_for_user(username: &CString, gid: nix::unistd::Gid) -> std::io::Result<()> {
1665 let gid = gid.as_raw();
1666 #[cfg(any(
1667 target_os = "macos",
1668 target_os = "ios",
1669 target_os = "tvos",
1670 target_os = "watchos"
1671 ))]
1672 let base_gid = i32::try_from(gid)
1673 .map_err(|_| std::io::Error::other(format!("gid {gid} is out of range")))?;
1674
1675 #[cfg(not(any(
1676 target_os = "macos",
1677 target_os = "ios",
1678 target_os = "tvos",
1679 target_os = "watchos"
1680 )))]
1681 let base_gid = gid as libc::gid_t;
1682
1683 let rc = unsafe { libc::initgroups(username.as_ptr(), base_gid) };
1686 if rc == -1 {
1687 Err(std::io::Error::last_os_error())
1688 } else {
1689 Ok(())
1690 }
1691}
1692
1693#[cfg(unix)]
1694fn nix_to_io_error(err: nix::errno::Errno) -> std::io::Error {
1695 std::io::Error::from_raw_os_error(err as i32)
1696}
1697
1698async fn check_ports_available(
1705 expected_ports: &[u16],
1706 auto_bump: bool,
1707 max_attempts: u32,
1708) -> Result<Vec<u16>> {
1709 if expected_ports.is_empty() {
1710 return Ok(Vec::new());
1711 }
1712
1713 for bump_offset in 0..=max_attempts {
1714 let candidate_ports: Vec<u16> = expected_ports
1716 .iter()
1717 .map(|&p| p.wrapping_add(bump_offset as u16))
1718 .collect();
1719
1720 let mut all_available = true;
1722 let mut conflicting_port = None;
1723
1724 for &port in &candidate_ports {
1725 if port == 0 {
1728 continue;
1729 }
1730
1731 if is_port_in_use(port).await {
1745 all_available = false;
1746 conflicting_port = Some(port);
1747 break;
1748 }
1749 }
1750
1751 if all_available {
1752 if candidate_ports.contains(&0) && !expected_ports.contains(&0) {
1756 return Err(PortError::NoAvailablePort {
1757 start_port: expected_ports[0],
1758 attempts: bump_offset + 1,
1759 }
1760 .into());
1761 }
1762 if bump_offset > 0 {
1763 info!("ports {expected_ports:?} bumped by {bump_offset} to {candidate_ports:?}");
1764 }
1765 return Ok(candidate_ports);
1766 }
1767
1768 if bump_offset == 0 && !auto_bump {
1770 if let Some(port) = conflicting_port {
1771 let (pid, process) = identify_port_owner(port).await;
1772 return Err(PortError::InUse { port, process, pid }.into());
1773 }
1774 }
1775 }
1776
1777 Err(PortError::NoAvailablePort {
1779 start_port: expected_ports[0],
1780 attempts: max_attempts + 1,
1781 }
1782 .into())
1783}
1784
1785async fn is_port_in_use(port: u16) -> bool {
1791 tokio::task::spawn_blocking(move || {
1792 for &addr in &["0.0.0.0", "127.0.0.1", "::1"] {
1793 match std::net::TcpListener::bind((addr, port)) {
1794 Ok(listener) => drop(listener),
1795 Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => return true,
1796 Err(_) => continue,
1797 }
1798 }
1799 false
1800 })
1801 .await
1802 .unwrap_or(false)
1803}
1804
1805async fn identify_port_owner(port: u16) -> (u32, String) {
1810 tokio::task::spawn_blocking(move || {
1811 listeners::get_all()
1812 .ok()
1813 .and_then(|list| {
1814 list.into_iter()
1815 .find(|l| l.socket.port() == port)
1816 .map(|l| (l.process.pid, l.process.name))
1817 })
1818 .unwrap_or((0, "unknown".to_string()))
1819 })
1820 .await
1821 .unwrap_or((0, "unknown".to_string()))
1822}
1823
1824async fn detect_port_conflict(port: u16) -> Option<(u32, String)> {
1829 if !is_port_in_use(port).await {
1830 return None;
1831 }
1832 Some(identify_port_owner(port).await)
1833}
1834
1835fn detect_and_store_active_port(id: DaemonId, pid: u32) {
1850 tokio::spawn(async move {
1851 for delay_ms in [500u64, 1000, 2000, 4000] {
1855 tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
1856
1857 let expected_port: Option<u16> = {
1860 let state_file = SUPERVISOR.state_file.lock().await;
1861 match state_file.daemons.get(&id) {
1862 Some(d) if d.pid.is_none() => {
1863 debug!("daemon {id}: aborting active_port detection — process exited");
1864 return;
1865 }
1866 Some(d) => d
1867 .port
1868 .as_ref()
1869 .and_then(|p| p.expect.first().copied())
1870 .filter(|&p| p > 0),
1871 None => None,
1872 }
1873 };
1874
1875 let active_port = tokio::task::spawn_blocking(move || {
1876 let listeners = listeners::get_all().ok()?;
1877
1878 PROCS.refresh_processes();
1880
1881 let descendant_pids: std::collections::HashSet<u32> = PROCS
1882 .all_children(pid)
1883 .into_iter()
1884 .chain(std::iter::once(pid))
1885 .collect();
1886
1887 let process_ports: Vec<u16> = listeners
1888 .into_iter()
1889 .filter(|listener| descendant_pids.contains(&listener.process.pid))
1890 .map(|listener| listener.socket.port())
1891 .filter(|&port| port > 0)
1892 .collect();
1893
1894 if process_ports.is_empty() {
1895 return None;
1896 }
1897
1898 if let Some(ep) = expected_port {
1901 if process_ports.contains(&ep) {
1902 return Some(ep);
1903 }
1904 }
1905
1906 process_ports.into_iter().next()
1912 })
1913 .await
1914 .ok()
1915 .flatten();
1916
1917 if let Some(port) = active_port {
1918 debug!("daemon {id} active_port detected: {port}");
1919 let mut state_file = SUPERVISOR.state_file.lock().await;
1920 if let Some(d) = state_file.daemons.get(&id) {
1921 if d.pid == Some(pid) {
1925 state_file.set_active_port(&id, port);
1926 } else {
1927 debug!(
1928 "daemon {id}: skipping active_port write — PID mismatch \
1929 (expected {pid}, current {:?})",
1930 d.pid
1931 );
1932 return;
1933 }
1934 }
1935 return;
1936 }
1937
1938 debug!(
1939 "daemon {id}: no active port detected for pid {pid} or its descendants (will retry)"
1940 );
1941 }
1942
1943 debug!(
1944 "daemon {id}: active port detection exhausted all retries for pid {pid} and its descendants"
1945 );
1946 });
1947}
1948
1949fn is_daemon_slug_target(id: &DaemonId) -> bool {
1957 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
1961 slugs.iter().any(|(slug, entry)| {
1962 let daemon_name = entry.daemon.as_deref().unwrap_or(slug);
1963 id.name() == daemon_name
1964 })
1965}
1966
1967#[cfg(all(test, unix))]
1968mod tests {
1969 use super::*;
1970
1971 #[test]
1972 fn test_resolve_run_identity_empty_without_sudo() {
1973 let identity = resolve_run_identity(None, 501, 20, None, None).unwrap();
1974 assert_eq!(identity, RunIdentity::Inherit);
1975 }
1976
1977 #[test]
1978 fn test_resolve_run_identity_sudo_fallback() {
1979 let identity = resolve_run_identity(None, 0, 0, Some("501"), Some("20")).unwrap();
1980 let RunIdentity::Switch { uid, gid, .. } = identity else {
1981 panic!("expected identity switch");
1982 };
1983 assert_eq!(uid.as_raw(), 501);
1984 assert_eq!(gid.as_raw(), 20);
1985 }
1986
1987 #[test]
1988 fn test_resolve_run_identity_ignores_stale_sudo_when_not_root() {
1989 let identity = resolve_run_identity(None, 501, 20, Some("0"), Some("0")).unwrap();
1990 assert_eq!(identity, RunIdentity::Inherit);
1991 }
1992
1993 #[test]
1994 fn test_resolve_configured_user_root_name() {
1995 let identity = resolve_configured_user("root").unwrap();
1996 let RunIdentity::Switch { uid, username, .. } = identity else {
1997 panic!("expected identity switch");
1998 };
1999 assert_eq!(uid.as_raw(), 0);
2000 assert_eq!(
2001 username.as_deref().and_then(|s| s.to_str().ok()),
2002 Some("root")
2003 );
2004 }
2005
2006 #[test]
2007 fn test_resolve_configured_user_root_uid() {
2008 let identity = resolve_configured_user("0").unwrap();
2009 let RunIdentity::Switch { uid, username, .. } = identity else {
2010 panic!("expected identity switch");
2011 };
2012 assert_eq!(uid.as_raw(), 0);
2013 assert_eq!(
2014 username.as_deref().and_then(|s| s.to_str().ok()),
2015 Some("root")
2016 );
2017 }
2018
2019 #[test]
2020 fn test_resolve_configured_user_missing_user_fails() {
2021 let err = resolve_configured_user("pitchfork-user-that-should-not-exist")
2022 .unwrap_err()
2023 .to_string();
2024 assert!(err.contains("does not exist"));
2025 }
2026
2027 #[test]
2028 fn test_resolve_run_identity_requires_root_for_user_switch() {
2029 let err = resolve_run_identity(Some("root"), 501, 20, None, None)
2030 .unwrap_err()
2031 .to_string();
2032 assert!(err.contains("Restart the supervisor with sudo"));
2033 }
2034
2035 #[test]
2036 fn test_resolve_run_identity_same_user_is_noop() {
2037 let identity = resolve_run_identity(Some("root"), 0, 0, Some("501"), Some("20")).unwrap();
2038 assert_eq!(identity, RunIdentity::Inherit);
2039 }
2040}
2041
2042fn inject_proxy_env(cmd: &mut tokio::process::Command, slug: &Option<String>) {
2051 let s = crate::settings::settings();
2052 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
2053
2054 if should_force_loopback_host(slug) && !lan_enabled {
2055 cmd.env("HOST", "127.0.0.1");
2058 }
2059
2060 if let Some(url) = build_pitchfork_url(slug, &s) {
2062 cmd.env("PITCHFORK_URL", &url);
2063 }
2064
2065 if s.proxy.enable && s.proxy.https {
2067 let ca_path = if s.proxy.tls_cert.is_empty() {
2068 crate::env::PITCHFORK_STATE_DIR.join("proxy").join("ca.pem")
2069 } else {
2070 std::path::PathBuf::from(&s.proxy.tls_cert)
2071 };
2072 if ca_path.exists() {
2073 cmd.env("NODE_EXTRA_CA_CERTS", ca_path.to_string_lossy().to_string());
2074 }
2075 }
2076
2077 if s.proxy.enable {
2079 let tld = if lan_enabled { "local" } else { &s.proxy.tld };
2080 cmd.env("__VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS", format!(".{tld}"));
2081 }
2082
2083 if lan_enabled {
2085 cmd.env("PITCHFORK_LAN", "1");
2086 }
2087}
2088
2089fn should_force_loopback_host(slug: &Option<String>) -> bool {
2090 let Some(slug) = slug.as_deref() else {
2091 return false;
2092 };
2093
2094 let s = crate::settings::settings();
2095 if !s.proxy.enable {
2096 return false;
2097 }
2098
2099 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
2100 slugs.contains_key(slug)
2101}
2102
2103fn build_pitchfork_url(slug: &Option<String>, s: &crate::settings::Settings) -> Option<String> {
2107 let slug = slug.as_ref()?;
2108 if !s.proxy.enable {
2109 return None;
2110 }
2111 let scheme = if s.proxy.https { "https" } else { "http" };
2112 let port = u16::try_from(s.proxy.port).ok().filter(|&p| p > 0)?;
2113 let port_suffix = if (scheme == "https" && port == 443) || (scheme == "http" && port == 80) {
2114 String::new()
2115 } else {
2116 format!(":{port}")
2117 };
2118 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
2119 let tld = if lan_enabled { "local" } else { &s.proxy.tld };
2120 Some(format!("{scheme}://{slug}.{tld}{port_suffix}",))
2121}
2122
2123#[cfg(test)]
2124mod ready_check_tests {
2125 use super::*;
2126 use std::time::Duration;
2127
2128 #[test]
2129 fn any_ready_check_remaining_prefers_unbounded_checks() {
2130 let http = ReadyHttp::new("http://localhost/health");
2131 let cmd = ReadyCmd::new("true");
2132
2133 assert!(any_ready_check_remaining(
2134 None,
2135 false,
2136 None,
2137 false,
2138 Some(&http),
2139 false,
2140 None,
2141 false
2142 ));
2143 assert!(any_ready_check_remaining(
2144 None,
2145 false,
2146 None,
2147 false,
2148 None,
2149 false,
2150 Some(&cmd),
2151 false
2152 ));
2153 assert!(any_ready_check_remaining(
2154 None,
2155 false,
2156 Some(&ReadyPort::new(8080)),
2157 false,
2158 Some(&http),
2159 true,
2160 Some(&cmd),
2161 true
2162 ));
2163 }
2164
2165 #[test]
2166 fn any_ready_check_remaining_exhausted_timed_checks() {
2167 let http = ReadyHttp {
2168 url: "http://localhost/health".to_string(),
2169 status: vec![],
2170 timeout: Some(Duration::from_secs(5)),
2171 };
2172 let cmd = ReadyCmd {
2173 run: "true".to_string(),
2174 timeout: Some(Duration::from_secs(5)),
2175 };
2176
2177 assert!(any_ready_check_remaining(
2178 None,
2179 false,
2180 None,
2181 false,
2182 Some(&http),
2183 false,
2184 Some(&cmd),
2185 false
2186 ));
2187 assert!(!any_ready_check_remaining(
2188 None,
2189 false,
2190 None,
2191 false,
2192 Some(&http),
2193 true,
2194 Some(&cmd),
2195 true
2196 ));
2197 }
2198
2199 #[tokio::test]
2200 async fn spawn_cmd_probe_reports_success() {
2201 let id = DaemonId::new("global", "probe-test");
2202 let probe = spawn_cmd_probe(&id, "true", &std::env::temp_dir());
2203 let status = probe.result_rx.await.unwrap().unwrap();
2204 assert!(status.success());
2205 }
2206
2207 #[tokio::test]
2208 async fn spawn_cmd_probe_stops_on_request() {
2209 let id = DaemonId::new("global", "probe-test");
2210 let probe = spawn_cmd_probe(&id, "sleep 30", &std::env::temp_dir());
2211 let CmdProbe {
2212 cancel_tx,
2213 result_rx,
2214 } = probe;
2215 let _ = cancel_tx.send(());
2216 let status = result_rx.await.unwrap().unwrap();
2217 assert!(!status.success());
2218 }
2219}