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::procs::PROCS;
15use crate::settings::settings;
16use crate::shell::Shell;
17use crate::supervisor::state::UpsertDaemonOpts;
18use crate::{Result, env};
19use miette::IntoDiagnostic;
20use once_cell::sync::Lazy;
21use regex::Regex;
22use std::collections::HashMap;
23#[cfg(unix)]
24use std::ffi::CString;
25use std::sync::{Arc, atomic};
26use std::time::Duration;
27use tokio::io::AsyncBufReadExt;
28use tokio::select;
29use tokio::sync::oneshot;
30use tokio::time;
31
32static REGEX_CACHE: Lazy<std::sync::Mutex<HashMap<String, Regex>>> =
34 Lazy::new(|| std::sync::Mutex::new(HashMap::new()));
35
36#[cfg(unix)]
37#[derive(Clone, Debug, PartialEq, Eq)]
38enum RunIdentity {
39 Inherit,
40 Switch {
41 uid: nix::unistd::Uid,
42 gid: nix::unistd::Gid,
43 username: Option<CString>,
44 },
45}
46
47pub(crate) fn get_or_compile_regex(pattern: &str) -> Option<Regex> {
49 let mut cache = REGEX_CACHE.lock().unwrap_or_else(|e| e.into_inner());
50 if let Some(re) = cache.get(pattern) {
51 return Some(re.clone());
52 }
53 match Regex::new(pattern) {
54 Ok(re) => {
55 cache.insert(pattern.to_string(), re.clone());
56 Some(re)
57 }
58 Err(e) => {
59 error!("invalid regex pattern '{pattern}': {e}");
60 None
61 }
62 }
63}
64
65impl Supervisor {
66 pub async fn run(&self, opts: RunOptions) -> Result<IpcResponse> {
68 let id = &opts.id;
69 let cmd = opts.cmd.clone();
70
71 {
73 let mut pending = self.pending_autostops.lock().await;
74 if pending.remove(id).is_some() {
75 info!("cleared pending autostop for {id} (daemon starting)");
76 }
77 }
78
79 let daemon = self.get_daemon(id).await;
80 if let Some(daemon) = daemon {
81 if !daemon.status.is_stopping()
84 && !daemon.status.is_stopped()
85 && let Some(pid) = daemon.pid
86 {
87 if opts.force {
88 self.stop(id).await?;
89 info!("run: stop completed for daemon {id}");
90 } else {
91 warn!("daemon {id} already running with pid {pid}");
92 return Ok(IpcResponse::DaemonAlreadyRunning);
93 }
94 }
95 }
96
97 if opts.wait_ready && opts.retry.count() > 0 {
99 let max_attempts = opts.retry.count().saturating_add(1);
101 for attempt in 0..max_attempts {
102 let mut retry_opts = opts.clone();
103 retry_opts.retry_count = attempt;
104 retry_opts.cmd = cmd.clone();
105
106 let result = self.run_once(retry_opts).await?;
107
108 match result {
109 IpcResponse::DaemonReady { daemon } => {
110 return Ok(IpcResponse::DaemonReady { daemon });
111 }
112 IpcResponse::DaemonFailedWithCode { exit_code } => {
113 if attempt < opts.retry.count() {
114 let backoff_secs = 2u64.saturating_pow(attempt).min(3600);
115 info!(
116 "daemon {id} failed (attempt {}/{}), retrying in {}s",
117 attempt + 1,
118 max_attempts,
119 backoff_secs
120 );
121 fire_hook(
122 HookType::OnRetry,
123 id.clone(),
124 opts.dir.0.clone(),
125 attempt + 1,
126 opts.env.clone(),
127 vec![],
128 )
129 .await;
130 time::sleep(Duration::from_secs(backoff_secs)).await;
131 continue;
132 } else {
133 info!("daemon {id} failed after {max_attempts} attempts");
134 return Ok(IpcResponse::DaemonFailedWithCode { exit_code });
135 }
136 }
137 other => return Ok(other),
138 }
139 }
140 }
141
142 self.run_once(opts).await
144 }
145
146 pub(crate) async fn run_once(&self, opts: RunOptions) -> Result<IpcResponse> {
148 let id = &opts.id;
149 let original_cmd = opts.cmd.clone(); let (ready_tx, ready_rx) = if opts.wait_ready {
153 let (tx, rx) = oneshot::channel();
154 (Some(tx), Some(rx))
155 } else {
156 (None, None)
157 };
158
159 let expected_ports = opts
161 .port
162 .as_ref()
163 .map(|p| p.expect.clone())
164 .unwrap_or_default();
165 let (resolved_ports, effective_ready_port) = if !expected_ports.is_empty() {
166 let port_cfg = opts.port.as_ref().unwrap();
167 match check_ports_available(
168 &expected_ports,
169 port_cfg.auto_bump(),
170 port_cfg.max_bump_attempts(),
171 )
172 .await
173 {
174 Ok(resolved) => {
175 let ready_port = if let Some(configured_port) = opts.ready_port {
176 let bump_offset = resolved
178 .first()
179 .unwrap_or(&0)
180 .saturating_sub(*expected_ports.first().unwrap_or(&0));
181 if expected_ports.contains(&configured_port) && bump_offset > 0 {
182 configured_port
183 .checked_add(bump_offset)
184 .or(Some(configured_port))
185 } else {
186 Some(configured_port)
187 }
188 } else if opts.ready_output.is_none()
189 && opts.ready_http.is_none()
190 && opts.ready_cmd.is_none()
191 && opts.ready_delay.is_none()
192 {
193 resolved.first().copied().filter(|&p| p != 0)
197 } else {
198 None
202 };
203 info!("daemon {id}: ports {expected_ports:?} resolved to {resolved:?}");
204 (resolved, ready_port)
205 }
206 Err(e) => {
207 error!("daemon {id}: port check failed: {e}");
208 if let Some(port_error) = e.downcast_ref::<PortError>() {
210 match port_error {
211 PortError::InUse { port, process, pid } => {
212 return Ok(IpcResponse::PortConflict {
213 port: *port,
214 process: process.clone(),
215 pid: *pid,
216 });
217 }
218 PortError::NoAvailablePort {
219 start_port,
220 attempts,
221 } => {
222 return Ok(IpcResponse::NoAvailablePort {
223 start_port: *start_port,
224 attempts: *attempts,
225 });
226 }
227 }
228 }
229 return Ok(IpcResponse::DaemonFailed {
230 error: e.to_string(),
231 });
232 }
233 }
234 } else {
235 if let Some(port) = opts.ready_port {
241 if port > 0 {
242 if let Some((pid, process)) = detect_port_conflict(port).await {
243 return Ok(IpcResponse::PortConflict { port, process, pid });
244 }
245 }
246 }
247 (Vec::new(), opts.ready_port)
248 };
249
250 let shell_setting = settings().general.shell.clone();
254 let shell_parts = match shell_words::split(&shell_setting) {
255 Ok(parts) if !parts.is_empty() => parts,
256 Ok(_) => {
257 return Ok(IpcResponse::DaemonFailed {
258 error: "general.shell setting is empty".to_string(),
259 });
260 }
261 Err(e) => {
262 return Ok(IpcResponse::DaemonFailed {
263 error: format!("failed to parse general.shell setting {shell_setting:?}: {e}"),
264 });
265 }
266 };
267 let (shell_program, shell_args) = shell_parts.split_first().unwrap();
268
269 let run_script = opts
274 .run
275 .clone()
276 .unwrap_or_else(|| shell_words::join(&original_cmd));
277
278 let (program, args) = if opts.mise.unwrap_or(settings().general.mise) {
279 match settings().resolve_mise_bin() {
280 Some(mise_bin) => {
281 let mise_bin_str = mise_bin.to_string_lossy().to_string();
282 info!("daemon {id}: wrapping command with mise ({mise_bin_str})");
283 let mut args = vec!["x".to_string(), "--".to_string()];
284 args.push(shell_program.clone());
285 args.extend(shell_args.iter().cloned());
286 args.push(run_script);
287 (mise_bin_str, args)
288 }
289 None => {
290 warn!("daemon {id}: mise=true but mise binary not found, running without mise");
291 let mut args: Vec<String> = shell_args.to_vec();
292 args.push(run_script);
293 (shell_program.clone(), args)
294 }
295 }
296 } else {
297 let mut args: Vec<String> = shell_args.to_vec();
298 args.push(run_script);
299 (shell_program.clone(), args)
300 };
301 #[cfg(unix)]
302 let run_identity = match resolve_effective_run_identity(opts.user.as_deref()) {
303 Ok(identity) => identity,
304 Err(e) => {
305 return Ok(IpcResponse::DaemonFailed {
306 error: e.to_string(),
307 });
308 }
309 };
310 info!("run: spawning daemon {id} with {program} {args:?}");
311
312 #[cfg(unix)]
314 let pty_pair = if opts.pty.unwrap_or(false) {
315 match super::pty::openpty() {
316 Ok(pair) => {
317 info!("daemon {id}: allocated PTY (pty = true)");
318 Some(pair)
319 }
320 Err(e) => {
321 warn!("daemon {id}: failed to allocate PTY, falling back to pipes: {e}");
322 None
323 }
324 }
325 } else {
326 None
327 };
328
329 let mut cmd = tokio::process::Command::new(&program);
330
331 #[cfg(unix)]
332 if let Some(ref pair) = pty_pair {
333 let slave_file = std::fs::File::from(
337 pair.slave
338 .try_clone()
339 .map_err(|e| miette::miette!("failed to dup slave PTY fd: {e}"))?,
340 );
341 cmd.stdin(std::process::Stdio::from(slave_file.try_clone().map_err(
342 |e| miette::miette!("failed to clone slave PTY fd for stdin: {e}"),
343 )?));
344 cmd.stdout(std::process::Stdio::from(slave_file.try_clone().map_err(
345 |e| miette::miette!("failed to clone slave PTY fd for stdout: {e}"),
346 )?));
347 cmd.stderr(std::process::Stdio::from(slave_file));
348 } else {
349 cmd.stdout(std::process::Stdio::piped())
350 .stderr(std::process::Stdio::piped());
351 }
352
353 #[cfg(not(unix))]
354 {
355 cmd.stdout(std::process::Stdio::piped())
356 .stderr(std::process::Stdio::piped());
357 }
358
359 cmd.args(&args).current_dir(&opts.dir);
360
361 #[cfg(unix)]
362 if pty_pair.is_none() {
363 cmd.stdin(std::process::Stdio::null());
364 }
365
366 #[cfg(not(unix))]
367 cmd.stdin(std::process::Stdio::null());
368
369 if let Some(ref path) = *env::ORIGINAL_PATH {
371 cmd.env("PATH", path);
372 }
373
374 if let Some(ref env_vars) = opts.env {
376 cmd.envs(env_vars);
377 }
378
379 cmd.env("PITCHFORK_DAEMON_ID", id.qualified());
381 cmd.env("PITCHFORK_DAEMON_NAMESPACE", id.namespace());
382 cmd.env("PITCHFORK_RETRY_COUNT", opts.retry_count.to_string());
383
384 if !resolved_ports.is_empty() {
386 cmd.env("PORT", resolved_ports[0].to_string());
390 for (i, port) in resolved_ports.iter().enumerate() {
392 cmd.env(format!("PORT{i}"), port.to_string());
393 }
394 }
395
396 inject_proxy_env(&mut cmd, &opts.slug);
398
399 #[cfg(unix)]
400 {
401 let run_identity = run_identity.clone();
402 let use_pty = pty_pair.is_some();
403 unsafe {
404 cmd.pre_exec(move || {
405 nix::unistd::setsid().map_err(nix_to_io_error)?;
406
407 if use_pty {
411 let ret = libc::ioctl(0, libc::TIOCSCTTY as libc::c_ulong, 0);
412 if ret < 0 {
413 #[cfg(target_os = "linux")]
416 eprintln!(
417 "pitchfork: TIOCSCTTY failed: {}",
418 std::io::Error::last_os_error()
419 );
420 }
421 }
422
423 apply_run_identity(&run_identity)?;
424 Ok(())
425 });
426 }
427 }
428
429 let mut child = cmd.spawn().into_diagnostic()?;
430 let pid = match child.id() {
431 Some(p) => p,
432 None => {
433 warn!("Daemon {id} exited before PID could be captured");
434 return Ok(IpcResponse::DaemonFailed {
435 error: "Process exited immediately".to_string(),
436 });
437 }
438 };
439 info!("started daemon {id} with pid {pid}");
440 PROCS.refresh_pids(&[pid]);
441 let daemon = self
442 .upsert_daemon(
443 UpsertDaemonOpts::from_run_options(&opts, DaemonStatus::Running)
444 .set(|o| {
445 o.pid = Some(pid);
446 o.cmd = Some(original_cmd);
447 o.ready_port = effective_ready_port;
448 o.port = crate::config_types::PortConfig::from_parts(
449 expected_ports,
450 opts.port.as_ref().map(|p| p.bump).unwrap_or_default(),
451 );
452 o.resolved_port = resolved_ports;
453 })
454 .build(),
455 )
456 .await?;
457
458 let id_clone = id.clone();
459 let ready_delay = opts.ready_delay;
460 let ready_output = opts.ready_output.clone();
461 let ready_http = opts.ready_http.clone();
462 let ready_port = effective_ready_port;
463 let ready_cmd = opts.ready_cmd.clone();
464 let daemon_dir = opts.dir.0.clone();
465 let hook_retry_count = opts.retry_count;
466 let hook_retry = opts.retry;
467 let hook_daemon_env = opts.env.clone();
468 let on_output_hook = opts.on_output_hook.clone();
469 let has_port_config = opts.port.as_ref().is_some_and(|p| !p.expect.is_empty())
475 || (settings().proxy.enable && is_daemon_slug_target(id));
476 let daemon_pid = pid;
477
478 #[cfg(unix)]
482 let pty_reader = pty_pair.map(|p| {
483 tokio::io::BufReader::new(tokio::fs::File::from_std(std::fs::File::from(p.master)))
484 .lines()
485 });
486 #[cfg(not(unix))]
487 let pty_reader: Option<tokio::io::Lines<tokio::io::BufReader<tokio::fs::File>>> = None;
488 let stdout_reader = if pty_reader.is_none() {
489 child
490 .stdout
491 .take()
492 .map(|s| tokio::io::BufReader::new(s).lines())
493 } else {
494 None
495 };
496 let stderr_reader = if pty_reader.is_none() {
497 child
498 .stderr
499 .take()
500 .map(|s| tokio::io::BufReader::new(s).lines())
501 } else {
502 None
503 };
504
505 if pty_reader.is_none() && (stdout_reader.is_none() || stderr_reader.is_none()) {
506 error!("Failed to capture stdout/stderr for daemon {id}");
507 }
508
509 tokio::spawn(async move {
510 let id = id_clone;
511
512 let (output_tx, mut output_rx) = tokio::sync::mpsc::channel::<String>(256);
514
515 if let Some(mut reader) = pty_reader {
516 tokio::spawn(async move {
520 while let Ok(Some(mut line)) = reader.next_line().await {
521 if line.ends_with('\r') {
523 line.pop();
524 }
525 if output_tx.send(line).await.is_err() {
526 break;
527 }
528 }
529 });
530 } else {
531 if let Some(mut stdout) = stdout_reader {
536 let tx = output_tx.clone();
537 tokio::spawn(async move {
538 while let Ok(Some(line)) = stdout.next_line().await {
539 if tx.send(line).await.is_err() {
540 break;
541 }
542 }
543 });
544 }
545 if let Some(mut stderr) = stderr_reader {
546 let tx = output_tx.clone();
547 tokio::spawn(async move {
548 while let Ok(Some(line)) = stderr.next_line().await {
549 if tx.send(line).await.is_err() {
550 break;
551 }
552 }
553 });
554 }
555 drop(output_tx);
557 }
558 let log_store = Arc::clone(&LOG_STORE);
559 let format_line = |line: String| line;
560
561 const LOG_BATCH_SIZE: usize = 100;
562 const LOG_FLUSH_INTERVAL: Duration = Duration::from_millis(100);
563 let mut log_buffer: Vec<String> = Vec::with_capacity(LOG_BATCH_SIZE);
564 let mut log_flush_interval = tokio::time::interval(LOG_FLUSH_INTERVAL);
565 log_flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
566
567 let flush_logs = |buffer: &mut Vec<String>| -> Option<tokio::task::JoinHandle<()>> {
568 if buffer.is_empty() {
569 return None;
570 }
571 let store = Arc::clone(&log_store);
572 let id = id.clone();
573 let batch = std::mem::take(buffer);
574 Some(tokio::task::spawn_blocking(move || {
575 if let Err(e) = store.append_batch(&id, &batch) {
576 error!("Failed to write batch to log for daemon {id}: {e}");
577 }
578 }))
579 };
580
581 let mut ready_notified = false;
585 let mut ready_tx = ready_tx;
586 let ready_pattern = ready_output.as_ref().and_then(|p| get_or_compile_regex(p));
587 let mut active_port_spawned = false;
589
590 let on_output_hook = match on_output_hook {
594 Some(ref hook) => match hook.validate(id.name()) {
595 Ok(()) => on_output_hook,
596 Err(e) => {
597 error!("{e}");
598 None
599 }
600 },
601 None => None,
602 };
603
604 let on_output_pattern: Option<regex::Regex> = on_output_hook
607 .as_ref()
608 .and_then(|h| h.regex.as_deref().and_then(get_or_compile_regex));
609 let on_output_debounce = on_output_hook
610 .as_ref()
611 .map(|h| h.debounce_duration())
612 .unwrap_or(Duration::from_millis(1000));
613 let mut on_output_last_fired: Option<std::time::Instant> = None;
615
616 let mut delay_timer =
617 ready_delay.map(|secs| Box::pin(time::sleep(Duration::from_secs(secs))));
618
619 let s = settings();
621 let ready_check_interval = s.supervisor_ready_check_interval();
622 let http_client_timeout = s.supervisor_http_client_timeout();
623
624 let mut http_check_interval = ready_http
626 .as_ref()
627 .map(|_| tokio::time::interval(ready_check_interval));
628 let http_client = ready_http.as_ref().map(|_| {
629 reqwest::Client::builder()
630 .timeout(http_client_timeout)
631 .build()
632 .unwrap_or_default()
633 });
634
635 let mut port_check_interval =
637 ready_port.map(|_| tokio::time::interval(ready_check_interval));
638
639 let mut cmd_check_interval = ready_cmd
641 .as_ref()
642 .map(|_| tokio::time::interval(ready_check_interval));
643
644 let (exit_tx, mut exit_rx) =
646 tokio::sync::mpsc::channel::<std::io::Result<std::process::ExitStatus>>(1);
647
648 let child_pid = child.id().unwrap_or(0);
650 tokio::spawn(async move {
651 let result = child.wait().await;
652 #[cfg(all(unix, not(target_os = "linux")))]
662 let result = match &result {
663 Err(e) if e.raw_os_error() == Some(nix::libc::ECHILD) => {
664 if let Some(code) = super::REAPED_STATUSES.lock().await.remove(&child_pid) {
665 warn!(
666 "daemon pid {child_pid} wait() got ECHILD; \
667 recovered exit code {code} from zombie reaper"
668 );
669 use std::os::unix::process::ExitStatusExt;
674 if code >= 0 {
675 Ok(std::process::ExitStatus::from_raw(code << 8))
676 } else {
677 Ok(std::process::ExitStatus::from_raw((-code) & 0x7f))
679 }
680 } else {
681 warn!(
682 "daemon pid {child_pid} wait() got ECHILD but no \
683 stashed status found; reporting as error"
684 );
685 result
686 }
687 }
688 _ => result,
689 };
690 debug!("daemon pid {child_pid} wait() completed with result: {result:?}");
691 let _ = exit_tx.send(result).await;
692 });
693
694 #[allow(unused_assignments)]
695 let mut exit_status = None;
697
698 if has_port_config
704 && ready_pattern.is_none()
705 && ready_http.is_none()
706 && ready_port.is_none()
707 && ready_cmd.is_none()
708 && delay_timer.is_none()
709 {
710 active_port_spawned = true;
711 detect_and_store_active_port(id.clone(), daemon_pid);
712 }
713
714 loop {
715 select! {
716 Some(line) = output_rx.recv() => {
717 let formatted = format_line(line.clone());
718 log_buffer.push(formatted.clone());
719 if log_buffer.len() >= LOG_BATCH_SIZE {
720 let _ = flush_logs(&mut log_buffer);
721 }
722 trace!("output: {id} {formatted}");
723
724 let line_clean = console::strip_ansi_codes(&line).to_string();
727
728 if !ready_notified
730 && let Some(ref pattern) = ready_pattern
731 && pattern.is_match(&line_clean)
732 {
733 if let Some(handle) = flush_logs(&mut log_buffer) {
738 let _ = handle.await;
739 }
740 info!("daemon {id} ready: output matched pattern");
741 ready_notified = true;
742 if let Some(tx) = ready_tx.take() {
743 let _ = tx.send(Ok(()));
744 }
745 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
746 if !active_port_spawned && has_port_config {
747 active_port_spawned = true;
748 detect_and_store_active_port(id.clone(), daemon_pid);
749 }
750 }
751
752 if let Some(ref hook) = on_output_hook {
754 let matched = match (&hook.filter, &on_output_pattern) {
755 (Some(substr), _) => line_clean.contains(substr.as_str()),
756 (None, Some(re)) => re.is_match(&line_clean),
757 (None, None) => true,
758 };
759 if matched {
760 let now = std::time::Instant::now();
761 let elapsed = on_output_last_fired.map(|t| now.duration_since(t));
762 if elapsed.is_none_or(|e| e >= on_output_debounce) {
763 on_output_last_fired = Some(now);
764 hooks::fire_output_hook(id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook.run.clone(), line_clean.clone()).await;
765 }
766 }
767 }
768 }
769 Some(result) = exit_rx.recv() => {
770 exit_status = Some(result);
772 debug!("daemon {id} process exited, exit_status: {exit_status:?}");
773 if !ready_notified {
774 if let Some(tx) = ready_tx.take() {
775 let is_success = exit_status.as_ref()
777 .and_then(|r| r.as_ref().ok())
778 .map(|s| s.success())
779 .unwrap_or(false);
780
781 if is_success {
782 debug!("daemon {id} exited successfully before ready check, sending success notification");
783 let _ = tx.send(Ok(()));
784 } else {
785 let exit_code = exit_status.as_ref()
786 .and_then(|r| r.as_ref().ok())
787 .and_then(|s| s.code());
788 debug!("daemon {id} exited with failure before ready check, sending failure notification with exit_code: {exit_code:?}");
789 let _ = tx.send(Err(exit_code));
790 }
791 }
792 } else {
793 debug!("daemon {id} was already marked ready, not sending notification");
794 }
795 break;
796 },
797 _ = async {
798 if let Some(ref mut interval) = http_check_interval {
799 interval.tick().await;
800 } else {
801 std::future::pending::<()>().await;
802 }
803 }, if !ready_notified && ready_http.is_some() => {
804 if let (Some(http), Some(client)) = (&ready_http, &http_client) {
805 match client.get(&http.url).send().await {
806 Ok(response) if http.accepts_status(response.status().as_u16()) => {
807 info!("daemon {id} ready: HTTP check passed (status {})", response.status());
808 ready_notified = true;
809 if let Some(tx) = ready_tx.take() {
810 let _ = tx.send(Ok(()));
811 }
812 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
813 http_check_interval = None;
814 if !active_port_spawned && has_port_config {
815 active_port_spawned = true;
816 detect_and_store_active_port(id.clone(), daemon_pid);
817 }
818 }
819 Ok(response) => {
820 trace!("daemon {id} HTTP check: status {} (not ready)", response.status());
821 }
822 Err(e) => {
823 trace!("daemon {id} HTTP check failed: {e}");
824 }
825 }
826 }
827 }
828 _ = async {
829 if let Some(ref mut interval) = port_check_interval {
830 interval.tick().await;
831 } else {
832 std::future::pending::<()>().await;
833 }
834 }, if !ready_notified && ready_port.is_some() => {
835 if let Some(port) = ready_port {
836 match tokio::net::TcpStream::connect(("127.0.0.1", port)).await {
837 Ok(_) => {
838 info!("daemon {id} ready: TCP port {port} is listening");
839 ready_notified = true;
840 if let Some(tx) = ready_tx.take() {
841 let _ = tx.send(Ok(()));
842 }
843 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
844 port_check_interval = None;
846 if !active_port_spawned && has_port_config {
847 active_port_spawned = true;
848 detect_and_store_active_port(id.clone(), daemon_pid);
849 }
850 }
851 Err(_) => {
852 trace!("daemon {id} port check: port {port} not listening yet");
853 }
854 }
855 }
856 }
857 _ = async {
858 if let Some(ref mut interval) = cmd_check_interval {
859 interval.tick().await;
860 } else {
861 std::future::pending::<()>().await;
862 }
863 }, if !ready_notified && ready_cmd.is_some() => {
864 if let Some(ref cmd) = ready_cmd {
865 let mut command = Shell::default_for_platform().command(cmd);
867 command
868 .current_dir(&daemon_dir)
869 .stdout(std::process::Stdio::null())
870 .stderr(std::process::Stdio::null());
871 let result: std::io::Result<std::process::ExitStatus> = command.status().await;
872 match result {
873 Ok(status) if status.success() => {
874 info!("daemon {id} ready: readiness command succeeded");
875 ready_notified = true;
876 if let Some(tx) = ready_tx.take() {
877 let _ = tx.send(Ok(()));
878 }
879 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
880 cmd_check_interval = None;
882 if !active_port_spawned && has_port_config {
883 active_port_spawned = true;
884 detect_and_store_active_port(id.clone(), daemon_pid);
885 }
886 }
887 Ok(_) => {
888 trace!("daemon {id} cmd check: command returned non-zero (not ready)");
889 }
890 Err(e) => {
891 trace!("daemon {id} cmd check failed: {e}");
892 }
893 }
894 }
895 }
896 _ = async {
897 if let Some(ref mut timer) = delay_timer {
898 timer.await;
899 } else {
900 std::future::pending::<()>().await;
901 }
902 } => {
903 if !ready_notified && ready_pattern.is_none() && ready_http.is_none() && ready_port.is_none() && ready_cmd.is_none() {
904 info!("daemon {id} ready: delay elapsed");
905 ready_notified = true;
906 if let Some(tx) = ready_tx.take() {
907 let _ = tx.send(Ok(()));
908 }
909 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
910 }
911 delay_timer = None;
913 if !active_port_spawned && has_port_config {
914 active_port_spawned = true;
915 detect_and_store_active_port(id.clone(), daemon_pid);
916 }
917 }
918 _ = log_flush_interval.tick() => {
919 let _ = flush_logs(&mut log_buffer);
920 }
921 }
922 }
923
924 let drain_deadline = tokio::time::Instant::now() + Duration::from_secs(5);
932 loop {
933 let now = tokio::time::Instant::now();
934 if now >= drain_deadline {
935 break;
936 }
937 let Ok(Some(line)) =
938 tokio::time::timeout(drain_deadline - now, output_rx.recv()).await
939 else {
940 break;
941 };
942 log_buffer.push(format_line(line));
943 }
944 if let Some(handle) = flush_logs(&mut log_buffer) {
947 let _ = handle.await;
948 }
949
950 {
952 let mut state_file = SUPERVISOR.state_file.lock().await;
953 state_file.clear_active_port(&id);
954 }
955
956 let exit_status = if let Some(status) = exit_status {
958 status
959 } else {
960 match exit_rx.recv().await {
962 Some(status) => status,
963 None => {
964 warn!("daemon {id} exit channel closed without receiving status");
965 Err(std::io::Error::other("exit channel closed"))
966 }
967 }
968 };
969 let current_daemon = SUPERVISOR.get_daemon(&id).await;
970
971 SUPERVISOR
976 .active_monitors
977 .fetch_add(1, atomic::Ordering::Release);
978 struct MonitorGuard;
979 impl Drop for MonitorGuard {
980 fn drop(&mut self) {
981 SUPERVISOR
982 .active_monitors
983 .fetch_sub(1, atomic::Ordering::Release);
984 SUPERVISOR.monitor_done.notify_waiters();
985 }
986 }
987 let _monitor_guard = MonitorGuard;
988 if current_daemon.is_none()
993 || current_daemon.as_ref().is_some_and(|d| {
994 d.pid != Some(pid) && !d.status.is_stopped() && !d.status.is_stopping()
995 })
996 {
997 return;
999 }
1000 let already_stopped = current_daemon
1005 .as_ref()
1006 .is_some_and(|d| d.status.is_stopped());
1007 let is_stopping = already_stopped
1008 || current_daemon
1009 .as_ref()
1010 .is_some_and(|d| d.status.is_stopping());
1011
1012 let (exit_code, exit_reason) = match (&exit_status, is_stopping) {
1014 (Ok(status), true) => {
1015 (status.code().unwrap_or(-1), "stop")
1019 }
1020 (Ok(status), false) if status.success() => (status.code().unwrap_or(-1), "exit"),
1021 (Ok(status), false) => (status.code().unwrap_or(-1), "fail"),
1022 (Err(_), true) => {
1023 (-1, "stop")
1025 }
1026 (Err(_), false) => (-1, "fail"),
1027 };
1028
1029 if !already_stopped {
1031 if let Ok(status) = &exit_status {
1032 info!("daemon {id} exited with status {status}");
1033 }
1034 let (new_status, last_exit_success) = match exit_reason {
1035 "stop" | "exit" => (
1036 DaemonStatus::Stopped,
1037 exit_status.as_ref().map(|s| s.success()).unwrap_or(true),
1038 ),
1039 _ => (DaemonStatus::Errored(exit_code), false),
1040 };
1041 if let Err(e) = SUPERVISOR
1042 .upsert_daemon(
1043 UpsertDaemonOpts::builder(id.clone())
1044 .set(|o| {
1045 o.pid = None;
1046 o.status = new_status;
1047 o.last_exit_success = Some(last_exit_success);
1048 })
1049 .build(),
1050 )
1051 .await
1052 {
1053 error!("Failed to update daemon state for {id}: {e}");
1054 }
1055 }
1056
1057 let hook_extra_env = vec![
1059 ("PITCHFORK_EXIT_CODE".to_string(), exit_code.to_string()),
1060 ("PITCHFORK_EXIT_REASON".to_string(), exit_reason.to_string()),
1061 ];
1062
1063 let hooks_to_fire: Vec<HookType> = match exit_reason {
1065 "stop" => vec![HookType::OnStop, HookType::OnExit],
1066 "exit" => vec![HookType::OnExit],
1067 _ if hook_retry_count >= hook_retry.count() => {
1069 vec![HookType::OnFail, HookType::OnExit]
1070 }
1071 _ => vec![],
1072 };
1073
1074 for hook_type in hooks_to_fire {
1075 fire_hook(
1076 hook_type,
1077 id.clone(),
1078 daemon_dir.clone(),
1079 hook_retry_count,
1080 hook_daemon_env.clone(),
1081 hook_extra_env.clone(),
1082 )
1083 .await;
1084 }
1085 });
1086
1087 if let Some(ready_rx) = ready_rx {
1089 match ready_rx.await {
1090 Ok(Ok(())) => {
1091 info!("daemon {id} is ready");
1092 Ok(IpcResponse::DaemonReady { daemon })
1093 }
1094 Ok(Err(exit_code)) => {
1095 error!("daemon {id} failed before becoming ready");
1096 Ok(IpcResponse::DaemonFailedWithCode { exit_code })
1097 }
1098 Err(_) => {
1099 error!("readiness channel closed unexpectedly for daemon {id}");
1100 Ok(IpcResponse::DaemonStart { daemon })
1101 }
1102 }
1103 } else {
1104 Ok(IpcResponse::DaemonStart { daemon })
1105 }
1106 }
1107
1108 pub async fn stop(&self, id: &DaemonId) -> Result<IpcResponse> {
1110 let pitchfork_id = DaemonId::pitchfork();
1111 if *id == pitchfork_id {
1112 return Ok(IpcResponse::Error(
1113 "Cannot stop supervisor via stop command".into(),
1114 ));
1115 }
1116 info!("stopping daemon: {id}");
1117 if let Some(daemon) = self.get_daemon(id).await {
1118 trace!("daemon to stop: {daemon}");
1119 if let Some(pid) = daemon.pid {
1120 trace!("killing pid: {pid}");
1121 if PROCS.is_running(pid) {
1122 self.upsert_daemon(
1124 UpsertDaemonOpts::builder(id.clone())
1125 .set(|o| {
1126 o.pid = Some(pid);
1127 o.status = DaemonStatus::Stopping;
1128 })
1129 .build(),
1130 )
1131 .await?;
1132
1133 let stop_cfg = daemon.stop_signal.unwrap_or_default();
1136 let stop_signal: i32 = stop_cfg.signal.into();
1137 if let Err(e) = PROCS
1138 .kill_process_group_async(pid, stop_signal, stop_cfg.timeout)
1139 .await
1140 {
1141 debug!("failed to kill pid {pid}: {e}");
1142 if PROCS.is_running(pid) {
1144 debug!("failed to stop pid {pid}: process still running after kill");
1146 self.upsert_daemon(
1147 UpsertDaemonOpts::builder(id.clone())
1148 .set(|o| {
1149 o.pid = Some(pid); o.status = DaemonStatus::Running;
1151 })
1152 .build(),
1153 )
1154 .await?;
1155 return Ok(IpcResponse::DaemonStopFailed {
1156 error: format!(
1157 "process {pid} still running after kill attempt: {e}"
1158 ),
1159 });
1160 }
1161 }
1162
1163 self.upsert_daemon(
1168 UpsertDaemonOpts::builder(id.clone())
1169 .set(|o| {
1170 o.pid = None;
1171 o.status = DaemonStatus::Stopped;
1172 o.last_exit_success = Some(true); })
1174 .build(),
1175 )
1176 .await?;
1177 } else {
1178 debug!("pid {pid} not running, process may have exited unexpectedly");
1179 self.upsert_daemon(
1182 UpsertDaemonOpts::builder(id.clone())
1183 .set(|o| {
1184 o.pid = None;
1185 o.status = DaemonStatus::Stopped;
1186 })
1187 .build(),
1188 )
1189 .await?;
1190 return Ok(IpcResponse::DaemonWasNotRunning);
1191 }
1192 Ok(IpcResponse::Ok)
1193 } else {
1194 debug!("daemon {id} not running");
1195 Ok(IpcResponse::DaemonNotRunning)
1196 }
1197 } else {
1198 debug!("daemon {id} not found");
1199 Ok(IpcResponse::DaemonNotFound)
1200 }
1201 }
1202}
1203
1204#[cfg(unix)]
1205fn resolve_effective_run_identity(daemon_user: Option<&str>) -> Result<RunIdentity> {
1206 let s = settings();
1207 let settings_user = s.supervisor.user.trim();
1208 let daemon_user = daemon_user.map(str::trim).filter(|user| !user.is_empty());
1209 let settings_user = (!settings_user.is_empty()).then_some(settings_user);
1210 let configured = daemon_user.or(settings_user);
1211 let current_uid = nix::unistd::Uid::effective().as_raw();
1212 let current_gid = nix::unistd::Gid::effective().as_raw();
1213 resolve_run_identity(
1214 configured,
1215 current_uid,
1216 current_gid,
1217 std::env::var("SUDO_UID").ok().as_deref(),
1218 std::env::var("SUDO_GID").ok().as_deref(),
1219 )
1220}
1221
1222#[cfg(unix)]
1223fn resolve_run_identity(
1224 configured: Option<&str>,
1225 current_uid: u32,
1226 current_gid: u32,
1227 sudo_uid: Option<&str>,
1228 sudo_gid: Option<&str>,
1229) -> Result<RunIdentity> {
1230 let current_uid = nix::unistd::Uid::from_raw(current_uid);
1231 let current_gid = nix::unistd::Gid::from_raw(current_gid);
1232 if let Some(user) = configured {
1233 let identity = resolve_configured_user(user)?;
1234 ensure_can_use_identity(user, &identity, current_uid, current_gid)?;
1235 if identity.matches(current_uid, current_gid) {
1236 return Ok(RunIdentity::Inherit);
1237 }
1238 return Ok(identity);
1239 }
1240
1241 if current_uid.is_root()
1242 && let Some(identity) = resolve_sudo_identity(sudo_uid, sudo_gid)
1243 {
1244 return Ok(identity);
1245 }
1246
1247 Ok(RunIdentity::Inherit)
1248}
1249
1250#[cfg(unix)]
1251fn resolve_configured_user(user: &str) -> Result<RunIdentity> {
1252 if user.chars().all(|c| c.is_ascii_digit()) {
1253 let uid = user
1254 .parse::<u32>()
1255 .map_err(|e| miette::miette!("invalid run user UID '{}': {}", user, e))?;
1256 let user_record = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1257 .into_diagnostic()?
1258 .ok_or_else(|| miette::miette!("run user UID '{}' does not exist", user))?;
1259 return run_identity_from_user_record(user_record);
1260 }
1261
1262 let user_record = nix::unistd::User::from_name(user)
1263 .into_diagnostic()?
1264 .ok_or_else(|| miette::miette!("run user '{}' does not exist", user))?;
1265 run_identity_from_user_record(user_record)
1266}
1267
1268#[cfg(unix)]
1269fn run_identity_from_user_record(user: nix::unistd::User) -> Result<RunIdentity> {
1270 let username = CString::new(user.name)
1271 .map_err(|e| miette::miette!("run user name contains an interior nul byte: {}", e))?;
1272 Ok(RunIdentity::Switch {
1273 uid: user.uid,
1274 gid: user.gid,
1275 username: Some(username),
1276 })
1277}
1278
1279#[cfg(unix)]
1280fn run_identity_from_raw_ids(uid: u32, gid: u32, username: Option<CString>) -> RunIdentity {
1281 RunIdentity::Switch {
1282 uid: nix::unistd::Uid::from_raw(uid),
1283 gid: nix::unistd::Gid::from_raw(gid),
1284 username,
1285 }
1286}
1287
1288#[cfg(unix)]
1289fn resolve_sudo_identity(sudo_uid: Option<&str>, sudo_gid: Option<&str>) -> Option<RunIdentity> {
1290 let uid = sudo_uid?.parse::<u32>().ok()?;
1291 let gid = sudo_gid?.parse::<u32>().ok()?;
1292 let username = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1293 .ok()
1294 .flatten()
1295 .and_then(|u| CString::new(u.name).ok());
1296 Some(run_identity_from_raw_ids(uid, gid, username))
1297}
1298
1299#[cfg(unix)]
1300fn ensure_can_use_identity(
1301 configured_user: &str,
1302 identity: &RunIdentity,
1303 current_uid: nix::unistd::Uid,
1304 current_gid: nix::unistd::Gid,
1305) -> Result<()> {
1306 let RunIdentity::Switch { uid, gid, .. } = identity else {
1307 return Ok(());
1308 };
1309 if *uid == current_uid && *gid == current_gid {
1310 return Ok(());
1311 }
1312 if current_uid.is_root() {
1313 return Ok(());
1314 }
1315 Err(miette::miette!(
1316 "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.",
1317 configured_user,
1318 current_uid.as_raw(),
1319 current_gid.as_raw(),
1320 uid.as_raw(),
1321 gid.as_raw()
1322 ))
1323}
1324
1325#[cfg(unix)]
1326fn apply_run_identity(identity: &RunIdentity) -> std::io::Result<()> {
1327 let RunIdentity::Switch { uid, gid, username } = identity else {
1328 return Ok(());
1329 };
1330 if let Some(username) = username {
1331 initgroups_for_user(username, *gid)?;
1332 } else {
1333 setgroups_to_primary(*gid)?;
1334 }
1335 nix::unistd::setgid(*gid).map_err(nix_to_io_error)?;
1336 nix::unistd::setuid(*uid).map_err(nix_to_io_error)?;
1337 Ok(())
1338}
1339
1340#[cfg(unix)]
1341impl RunIdentity {
1342 fn matches(&self, uid: nix::unistd::Uid, gid: nix::unistd::Gid) -> bool {
1343 matches!(self, RunIdentity::Switch { uid: u, gid: g, .. } if *u == uid && *g == gid)
1344 }
1345}
1346
1347#[cfg(unix)]
1348fn setgroups_to_primary(gid: nix::unistd::Gid) -> std::io::Result<()> {
1349 let groups = [gid.as_raw() as libc::gid_t];
1350 #[cfg(any(target_os = "linux", target_os = "android"))]
1351 let group_count = groups.len();
1352 #[cfg(not(any(target_os = "linux", target_os = "android")))]
1353 let group_count = groups.len() as libc::c_int;
1354 let rc = unsafe { libc::setgroups(group_count, groups.as_ptr()) };
1355 if rc == -1 {
1356 Err(std::io::Error::last_os_error())
1357 } else {
1358 Ok(())
1359 }
1360}
1361
1362#[cfg(unix)]
1363fn initgroups_for_user(username: &CString, gid: nix::unistd::Gid) -> std::io::Result<()> {
1364 let gid = gid.as_raw();
1365 #[cfg(any(
1366 target_os = "macos",
1367 target_os = "ios",
1368 target_os = "tvos",
1369 target_os = "watchos"
1370 ))]
1371 let base_gid = i32::try_from(gid)
1372 .map_err(|_| std::io::Error::other(format!("gid {gid} is out of range")))?;
1373
1374 #[cfg(not(any(
1375 target_os = "macos",
1376 target_os = "ios",
1377 target_os = "tvos",
1378 target_os = "watchos"
1379 )))]
1380 let base_gid = gid as libc::gid_t;
1381
1382 let rc = unsafe { libc::initgroups(username.as_ptr(), base_gid) };
1385 if rc == -1 {
1386 Err(std::io::Error::last_os_error())
1387 } else {
1388 Ok(())
1389 }
1390}
1391
1392#[cfg(unix)]
1393fn nix_to_io_error(err: nix::errno::Errno) -> std::io::Error {
1394 std::io::Error::from_raw_os_error(err as i32)
1395}
1396
1397async fn check_ports_available(
1404 expected_ports: &[u16],
1405 auto_bump: bool,
1406 max_attempts: u32,
1407) -> Result<Vec<u16>> {
1408 if expected_ports.is_empty() {
1409 return Ok(Vec::new());
1410 }
1411
1412 for bump_offset in 0..=max_attempts {
1413 let candidate_ports: Vec<u16> = expected_ports
1415 .iter()
1416 .map(|&p| p.wrapping_add(bump_offset as u16))
1417 .collect();
1418
1419 let mut all_available = true;
1421 let mut conflicting_port = None;
1422
1423 for &port in &candidate_ports {
1424 if port == 0 {
1427 continue;
1428 }
1429
1430 if is_port_in_use(port).await {
1444 all_available = false;
1445 conflicting_port = Some(port);
1446 break;
1447 }
1448 }
1449
1450 if all_available {
1451 if candidate_ports.contains(&0) && !expected_ports.contains(&0) {
1455 return Err(PortError::NoAvailablePort {
1456 start_port: expected_ports[0],
1457 attempts: bump_offset + 1,
1458 }
1459 .into());
1460 }
1461 if bump_offset > 0 {
1462 info!("ports {expected_ports:?} bumped by {bump_offset} to {candidate_ports:?}");
1463 }
1464 return Ok(candidate_ports);
1465 }
1466
1467 if bump_offset == 0 && !auto_bump {
1469 if let Some(port) = conflicting_port {
1470 let (pid, process) = identify_port_owner(port).await;
1471 return Err(PortError::InUse { port, process, pid }.into());
1472 }
1473 }
1474 }
1475
1476 Err(PortError::NoAvailablePort {
1478 start_port: expected_ports[0],
1479 attempts: max_attempts + 1,
1480 }
1481 .into())
1482}
1483
1484async fn is_port_in_use(port: u16) -> bool {
1490 tokio::task::spawn_blocking(move || {
1491 for &addr in &["0.0.0.0", "127.0.0.1", "::1"] {
1492 match std::net::TcpListener::bind((addr, port)) {
1493 Ok(listener) => drop(listener),
1494 Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => return true,
1495 Err(_) => continue,
1496 }
1497 }
1498 false
1499 })
1500 .await
1501 .unwrap_or(false)
1502}
1503
1504async fn identify_port_owner(port: u16) -> (u32, String) {
1509 tokio::task::spawn_blocking(move || {
1510 listeners::get_all()
1511 .ok()
1512 .and_then(|list| {
1513 list.into_iter()
1514 .find(|l| l.socket.port() == port)
1515 .map(|l| (l.process.pid, l.process.name))
1516 })
1517 .unwrap_or((0, "unknown".to_string()))
1518 })
1519 .await
1520 .unwrap_or((0, "unknown".to_string()))
1521}
1522
1523async fn detect_port_conflict(port: u16) -> Option<(u32, String)> {
1528 if !is_port_in_use(port).await {
1529 return None;
1530 }
1531 Some(identify_port_owner(port).await)
1532}
1533
1534fn detect_and_store_active_port(id: DaemonId, pid: u32) {
1549 tokio::spawn(async move {
1550 for delay_ms in [500u64, 1000, 2000, 4000] {
1554 tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
1555
1556 let expected_port: Option<u16> = {
1559 let state_file = SUPERVISOR.state_file.lock().await;
1560 match state_file.daemons.get(&id) {
1561 Some(d) if d.pid.is_none() => {
1562 debug!("daemon {id}: aborting active_port detection — process exited");
1563 return;
1564 }
1565 Some(d) => d
1566 .port
1567 .as_ref()
1568 .and_then(|p| p.expect.first().copied())
1569 .filter(|&p| p > 0),
1570 None => None,
1571 }
1572 };
1573
1574 let active_port = tokio::task::spawn_blocking(move || {
1575 let listeners = listeners::get_all().ok()?;
1576
1577 PROCS.refresh_processes();
1579
1580 let descendant_pids: std::collections::HashSet<u32> = PROCS
1581 .all_children(pid)
1582 .into_iter()
1583 .chain(std::iter::once(pid))
1584 .collect();
1585
1586 let process_ports: Vec<u16> = listeners
1587 .into_iter()
1588 .filter(|listener| descendant_pids.contains(&listener.process.pid))
1589 .map(|listener| listener.socket.port())
1590 .filter(|&port| port > 0)
1591 .collect();
1592
1593 if process_ports.is_empty() {
1594 return None;
1595 }
1596
1597 if let Some(ep) = expected_port {
1600 if process_ports.contains(&ep) {
1601 return Some(ep);
1602 }
1603 }
1604
1605 process_ports.into_iter().next()
1611 })
1612 .await
1613 .ok()
1614 .flatten();
1615
1616 if let Some(port) = active_port {
1617 debug!("daemon {id} active_port detected: {port}");
1618 let mut state_file = SUPERVISOR.state_file.lock().await;
1619 if let Some(d) = state_file.daemons.get(&id) {
1620 if d.pid == Some(pid) {
1624 state_file.set_active_port(&id, port);
1625 } else {
1626 debug!(
1627 "daemon {id}: skipping active_port write — PID mismatch \
1628 (expected {pid}, current {:?})",
1629 d.pid
1630 );
1631 return;
1632 }
1633 }
1634 return;
1635 }
1636
1637 debug!(
1638 "daemon {id}: no active port detected for pid {pid} or its descendants (will retry)"
1639 );
1640 }
1641
1642 debug!(
1643 "daemon {id}: active port detection exhausted all retries for pid {pid} and its descendants"
1644 );
1645 });
1646}
1647
1648fn is_daemon_slug_target(id: &DaemonId) -> bool {
1656 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
1660 slugs.iter().any(|(slug, entry)| {
1661 let daemon_name = entry.daemon.as_deref().unwrap_or(slug);
1662 id.name() == daemon_name
1663 })
1664}
1665
1666#[cfg(all(test, unix))]
1667mod tests {
1668 use super::*;
1669
1670 #[test]
1671 fn test_resolve_run_identity_empty_without_sudo() {
1672 let identity = resolve_run_identity(None, 501, 20, None, None).unwrap();
1673 assert_eq!(identity, RunIdentity::Inherit);
1674 }
1675
1676 #[test]
1677 fn test_resolve_run_identity_sudo_fallback() {
1678 let identity = resolve_run_identity(None, 0, 0, Some("501"), Some("20")).unwrap();
1679 let RunIdentity::Switch { uid, gid, .. } = identity else {
1680 panic!("expected identity switch");
1681 };
1682 assert_eq!(uid.as_raw(), 501);
1683 assert_eq!(gid.as_raw(), 20);
1684 }
1685
1686 #[test]
1687 fn test_resolve_run_identity_ignores_stale_sudo_when_not_root() {
1688 let identity = resolve_run_identity(None, 501, 20, Some("0"), Some("0")).unwrap();
1689 assert_eq!(identity, RunIdentity::Inherit);
1690 }
1691
1692 #[test]
1693 fn test_resolve_configured_user_root_name() {
1694 let identity = resolve_configured_user("root").unwrap();
1695 let RunIdentity::Switch { uid, username, .. } = identity else {
1696 panic!("expected identity switch");
1697 };
1698 assert_eq!(uid.as_raw(), 0);
1699 assert_eq!(
1700 username.as_deref().and_then(|s| s.to_str().ok()),
1701 Some("root")
1702 );
1703 }
1704
1705 #[test]
1706 fn test_resolve_configured_user_root_uid() {
1707 let identity = resolve_configured_user("0").unwrap();
1708 let RunIdentity::Switch { uid, username, .. } = identity else {
1709 panic!("expected identity switch");
1710 };
1711 assert_eq!(uid.as_raw(), 0);
1712 assert_eq!(
1713 username.as_deref().and_then(|s| s.to_str().ok()),
1714 Some("root")
1715 );
1716 }
1717
1718 #[test]
1719 fn test_resolve_configured_user_missing_user_fails() {
1720 let err = resolve_configured_user("pitchfork-user-that-should-not-exist")
1721 .unwrap_err()
1722 .to_string();
1723 assert!(err.contains("does not exist"));
1724 }
1725
1726 #[test]
1727 fn test_resolve_run_identity_requires_root_for_user_switch() {
1728 let err = resolve_run_identity(Some("root"), 501, 20, None, None)
1729 .unwrap_err()
1730 .to_string();
1731 assert!(err.contains("Restart the supervisor with sudo"));
1732 }
1733
1734 #[test]
1735 fn test_resolve_run_identity_same_user_is_noop() {
1736 let identity = resolve_run_identity(Some("root"), 0, 0, Some("501"), Some("20")).unwrap();
1737 assert_eq!(identity, RunIdentity::Inherit);
1738 }
1739}
1740
1741fn inject_proxy_env(cmd: &mut tokio::process::Command, slug: &Option<String>) {
1750 let s = crate::settings::settings();
1751 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
1752
1753 if should_force_loopback_host(slug) && !lan_enabled {
1754 cmd.env("HOST", "127.0.0.1");
1757 }
1758
1759 if let Some(url) = build_pitchfork_url(slug, &s) {
1761 cmd.env("PITCHFORK_URL", &url);
1762 }
1763
1764 if s.proxy.enable && s.proxy.https {
1766 let ca_path = if s.proxy.tls_cert.is_empty() {
1767 crate::env::PITCHFORK_STATE_DIR.join("proxy").join("ca.pem")
1768 } else {
1769 std::path::PathBuf::from(&s.proxy.tls_cert)
1770 };
1771 if ca_path.exists() {
1772 cmd.env("NODE_EXTRA_CA_CERTS", ca_path.to_string_lossy().to_string());
1773 }
1774 }
1775
1776 if s.proxy.enable {
1778 let tld = if lan_enabled { "local" } else { &s.proxy.tld };
1779 cmd.env("__VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS", format!(".{tld}"));
1780 }
1781
1782 if lan_enabled {
1784 cmd.env("PITCHFORK_LAN", "1");
1785 }
1786}
1787
1788fn should_force_loopback_host(slug: &Option<String>) -> bool {
1789 let Some(slug) = slug.as_deref() else {
1790 return false;
1791 };
1792
1793 let s = crate::settings::settings();
1794 if !s.proxy.enable {
1795 return false;
1796 }
1797
1798 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
1799 slugs.contains_key(slug)
1800}
1801
1802fn build_pitchfork_url(slug: &Option<String>, s: &crate::settings::Settings) -> Option<String> {
1806 let slug = slug.as_ref()?;
1807 if !s.proxy.enable {
1808 return None;
1809 }
1810 let scheme = if s.proxy.https { "https" } else { "http" };
1811 let port = u16::try_from(s.proxy.port).ok().filter(|&p| p > 0)?;
1812 let port_suffix = if (scheme == "https" && port == 443) || (scheme == "http" && port == 80) {
1813 String::new()
1814 } else {
1815 format!(":{port}")
1816 };
1817 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
1818 let tld = if lan_enabled { "local" } else { &s.proxy.tld };
1819 Some(format!("{scheme}://{slug}.{tld}{port_suffix}",))
1820}