pitchfork_cli/supervisor/lifecycle.rs
1//! Daemon lifecycle management - start/stop operations
2//!
3//! Contains the core `run()`, `run_once()`, and `stop()` methods for daemon process management.
4
5use super::hooks::{self, HookType, fire_hook};
6use super::{SUPERVISOR, Supervisor};
7use crate::config_types::OneshotWait;
8use crate::daemon::RunOptions;
9use crate::daemon_id::DaemonId;
10use crate::daemon_status::DaemonStatus;
11use crate::error::PortError;
12use crate::ipc::IpcResponse;
13use crate::log_store::LogStore;
14use crate::log_store::sqlite::LOG_STORE;
15use crate::pitchfork_toml::{ReadyCmd, ReadyHttp, ReadyOutput, ReadyPort};
16use crate::procs::PROCS;
17use crate::settings::{resolve_shell, settings};
18use crate::shell::{HideConsoleWindow, Shell, ShellScript};
19use crate::supervisor::state::UpsertDaemonOpts;
20use crate::{Result, env};
21use indexmap::IndexMap;
22use miette::IntoDiagnostic;
23use once_cell::sync::Lazy;
24use regex::Regex;
25use std::collections::HashMap;
26#[cfg(unix)]
27use std::ffi::CString;
28use std::sync::{Arc, atomic};
29use std::time::Duration;
30use tokio::select;
31use tokio::sync::oneshot;
32use tokio::time;
33
34/// Cache for compiled regex patterns to avoid recompilation on daemon restarts
35static REGEX_CACHE: Lazy<std::sync::Mutex<HashMap<String, Regex>>> =
36 Lazy::new(|| std::sync::Mutex::new(HashMap::new()));
37
38fn resolve_configured_ready_port(
39 configured_port: u16,
40 expected_ports: &[u16],
41 resolved_ports: &[u16],
42) -> u16 {
43 let bump_offset = resolved_ports
44 .first()
45 .unwrap_or(&0)
46 .saturating_sub(*expected_ports.first().unwrap_or(&0));
47 if expected_ports.contains(&configured_port) && bump_offset > 0 {
48 configured_port
49 .checked_add(bump_offset)
50 .unwrap_or(configured_port)
51 } else {
52 configured_port
53 }
54}
55
56fn active_port_from_ready_port(ready_port: u16, resolved_ports: &[u16]) -> Option<u16> {
57 resolved_ports
58 .first()
59 .copied()
60 .filter(|&primary_port| primary_port == ready_port)
61}
62
63#[cfg(unix)]
64#[derive(Clone, Debug, PartialEq, Eq)]
65enum RunIdentity {
66 Inherit,
67 Switch {
68 uid: nix::unistd::Uid,
69 gid: nix::unistd::Gid,
70 username: Option<CString>,
71 /// Home directory from the user's passwd entry.
72 home: Option<std::path::PathBuf>,
73 },
74}
75
76/// Get or compile a regex pattern, caching the result for future use
77pub(crate) fn get_or_compile_regex(pattern: &str) -> Option<Regex> {
78 let mut cache = REGEX_CACHE.lock().unwrap_or_else(|e| e.into_inner());
79 if let Some(re) = cache.get(pattern) {
80 return Some(re.clone());
81 }
82 match Regex::new(pattern) {
83 Ok(re) => {
84 cache.insert(pattern.to_string(), re.clone());
85 Some(re)
86 }
87 Err(e) => {
88 error!("invalid regex pattern '{pattern}': {e}");
89 None
90 }
91 }
92}
93
94/// Handle for an in-flight readiness command probe.
95///
96/// The spawned task owns the `tokio::process::Child` and waits for either the
97/// process to exit or the cancel signal. Dropping the handle without cancelling
98/// leaves the task running, but the child is started with `kill_on_drop(true)`
99/// so it will still be terminated when the task ends.
100pub(crate) struct CmdProbe {
101 pub(crate) cancel_tx: tokio::sync::oneshot::Sender<()>,
102 pub(crate) result_rx: tokio::sync::oneshot::Receiver<std::io::Result<std::process::ExitStatus>>,
103}
104
105/// Spawn a readiness command probe and return a handle that can be used to wait
106/// for the exit status or cancel the probe.
107///
108/// The probe is started with `kill_on_drop(true)` as a cancellation fallback. The
109/// spawned task waits for the process to exit; if cancellation is requested, it
110/// kills the child and waits for it to reap before reporting the result.
111fn apply_runtime_env(
112 command: &mut tokio::process::Command,
113 id: &DaemonId,
114 retry_count: u32,
115 daemon_env: Option<&IndexMap<String, String>>,
116 resolved_ports: &[u16],
117) {
118 if let Some(ref path) = *env::ORIGINAL_PATH {
119 command.env("PATH", path);
120 }
121 if let Some(env_vars) = daemon_env {
122 command.envs(env_vars);
123 }
124 command
125 .env("PITCHFORK_DAEMON_ID", id.qualified())
126 .env("PITCHFORK_DAEMON_NAMESPACE", id.namespace())
127 .env("PITCHFORK_RETRY_COUNT", retry_count.to_string());
128 if let Some(port) = resolved_ports.first() {
129 command.env("PORT", port.to_string());
130 for (index, port) in resolved_ports.iter().enumerate() {
131 command.env(format!("PORT{index}"), port.to_string());
132 }
133 }
134}
135
136pub(crate) fn spawn_cmd_probe(
137 id: &DaemonId,
138 cmd: &str,
139 dir: &std::path::Path,
140 retry_count: u32,
141 daemon_env: Option<&IndexMap<String, String>>,
142 resolved_ports: &[u16],
143) -> CmdProbe {
144 // Use the same shell as daemon run and hooks. A probe is not worth failing
145 // the daemon over, so an unparseable setting degrades to the platform's own
146 // shell here rather than propagating; run_once has already rejected the
147 // start by then, so this only fires for a daemon whose settings changed
148 // under it.
149 let mut command = match resolve_shell() {
150 Ok(parts) => {
151 let (program, args) = parts.split_first().unwrap();
152 let mut c = tokio::process::Command::new(program);
153 c.shell_script(program, args, cmd);
154 c
155 }
156 Err(e) => {
157 warn!("daemon {id}: {e}; using the platform shell for this probe");
158 Shell::default_for_platform().command(cmd)
159 }
160 };
161 command
162 .current_dir(dir)
163 .stdout(std::process::Stdio::null())
164 .stderr(std::process::Stdio::null())
165 .kill_on_drop(true)
166 .hide_console_window();
167 apply_runtime_env(&mut command, id, retry_count, daemon_env, resolved_ports);
168 let mut child = match command.spawn() {
169 Ok(child) => child,
170 Err(e) => {
171 warn!("daemon {id}: failed to spawn command probe: {e}");
172 // Return a probe whose result channel is already closed. The caller will
173 // treat this the same as a probe that exited non-zero and respawn after
174 // the ready_check_interval, preserving the existing retry behaviour.
175 let (cancel_tx, _) = tokio::sync::oneshot::channel();
176 let (_, result_rx) = tokio::sync::oneshot::channel();
177 return CmdProbe {
178 cancel_tx,
179 result_rx,
180 };
181 }
182 };
183
184 let (cancel_tx, mut cancel_rx) = tokio::sync::oneshot::channel();
185 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
186
187 tokio::spawn(async move {
188 let status = tokio::select! {
189 status = child.wait() => status,
190 _ = &mut cancel_rx => {
191 let mut child = child;
192 let _ = child.kill().await;
193 child.wait().await
194 }
195 };
196 let _ = result_tx.send(status);
197 });
198
199 CmdProbe {
200 cancel_tx,
201 result_rx,
202 }
203}
204
205/// Cancel an active command probe and clear its handle.
206fn stop_cmd_probe_state(probe: &mut Option<CmdProbe>) {
207 if let Some(p) = probe.take() {
208 let _ = p.cancel_tx.send(());
209 }
210}
211
212/// Spawn a detached task that kills a daemon's process group after its
213/// readiness checks are exhausted, logging a failed kill instead of
214/// discarding it. The returned handle is awaited before the readiness
215/// failure is reported so the process group is down by then.
216fn spawn_ready_fail_kill(
217 id: DaemonId,
218 pid: u32,
219 stop_cfg: crate::config_types::StopConfig,
220) -> tokio::task::JoinHandle<()> {
221 tokio::spawn(async move {
222 if let Err(e) = PROCS
223 .kill_process_group_async(pid, stop_cfg.signal.into(), stop_cfg.timeout)
224 .await
225 {
226 error!("daemon {id}: failed to kill pid {pid} after readiness failure: {e}");
227 }
228 })
229}
230
231/// Returns true if any configured readiness check can still succeed.
232/// A check with no timeout is unbounded; a timed check can still succeed until its
233/// deadline fires. `ready_delay` is only used as a fallback when no other check is
234/// configured, so it is not counted here.
235#[allow(clippy::too_many_arguments)]
236fn any_ready_check_remaining(
237 ready_output: Option<&ReadyOutput>,
238 output_exhausted: bool,
239 ready_port: Option<&ReadyPort>,
240 port_exhausted: bool,
241 ready_http: Option<&ReadyHttp>,
242 http_exhausted: bool,
243 ready_cmd: Option<&ReadyCmd>,
244 cmd_exhausted: bool,
245) -> bool {
246 ready_output.is_some_and(|o| o.timeout.is_none() || !output_exhausted)
247 || ready_port.is_some_and(|p| p.timeout.is_none() || !port_exhausted)
248 || ready_http.is_some_and(|h| h.timeout.is_none() || !http_exhausted)
249 || ready_cmd.is_some_and(|c| c.timeout.is_none() || !cmd_exhausted)
250}
251
252fn delay_readiness_succeeded(
253 ready_notified: bool,
254 has_other_ready_check: bool,
255 process_exited: bool,
256 process_running: bool,
257) -> bool {
258 !ready_notified && !has_other_ready_check && !process_exited && process_running
259}
260
261/// Terminal state recorded for a daemon run that has ended, and whether that
262/// ending counts as a successful exit.
263///
264/// A `oneshot` daemon's whole job is to finish, so a clean exit of its own
265/// accord is `Completed` rather than `Stopped` — that is what makes it
266/// distinguishable from a service that is merely not running, and what lets
267/// `depends` treat it as satisfied. An explicit stop is still a stop: the task
268/// was interrupted, not completed.
269fn terminal_exit_state(
270 exit_reason: &str,
271 oneshot: bool,
272 exit_code: i32,
273 exited_cleanly: bool,
274) -> (DaemonStatus, bool) {
275 match exit_reason {
276 "exit" if oneshot => (DaemonStatus::Completed, true),
277 "stop" | "exit" => (DaemonStatus::Stopped, exited_cleanly),
278 _ => (DaemonStatus::Errored(exit_code), false),
279 }
280}
281
282/// Whether a stop that arrived after a run's process was already gone should
283/// leave the status its monitor settled on alone.
284///
285/// Only a completed task is left alone: it had already done its work, so there
286/// was nothing for the stop to interrupt, and overwriting it would report a
287/// failure to anyone waiting on it. Anything else — a failure above all — is
288/// replaced by the stop, so the retry checker does not carry on with a task
289/// the user has stopped.
290/// Why the argv form of `run` cannot be started, if it cannot.
291///
292/// Config load checks the array as written, but a template can still render
293/// the program to nothing, or to `exec`.
294fn invalid_argv_program(id: &DaemonId, argv: &[String]) -> Option<String> {
295 match argv.first().map(String::as_str) {
296 None | Some("") => Some(format!(
297 "daemon {id} has no program to run: its run array starts with an empty value"
298 )),
299 Some("exec") => Some(format!(
300 "daemon {id} starts its run array with \"exec\"; a run array starts the program directly, so remove \"exec\""
301 )),
302 Some(_) => None,
303 }
304}
305
306/// Program and arguments to spawn for `words` — the daemon's program followed
307/// by its arguments — run through `mise x` when `mise_bin` is given.
308///
309/// Each word stays one argument: mise hands everything after `--` to the
310/// program as it received it.
311fn launch_command(words: Vec<String>, mise_bin: Option<&std::path::Path>) -> (String, Vec<String>) {
312 match mise_bin {
313 Some(mise_bin) => {
314 let mut args = vec!["x".to_string(), "--".to_string()];
315 args.extend(words);
316 (mise_bin.to_string_lossy().to_string(), args)
317 }
318 None => {
319 let mut words = words.into_iter();
320 // Never empty: a shell resolves to at least its program, and an
321 // empty argv is refused before this is reached.
322 let program = words.next().unwrap_or_default();
323 (program, words.collect())
324 }
325 }
326}
327
328/// Text of one line read from a daemon's output, its line ending removed.
329///
330/// Decoded as the log sink decodes lines, so output that is not UTF-8 is
331/// logged rather than ending the read. A PTY slave's ONLCR turns `\n` into
332/// `\r\n`, so every trailing `\r` goes, as the sink also strips them.
333fn output_line_text(line: &[u8]) -> String {
334 let line = line.strip_suffix(b"\n").unwrap_or(line);
335 crate::cli::log_sink::decode_line(line)
336 .trim_end_matches('\r')
337 .to_string()
338}
339
340/// Send each line of `reader` to `tx` until the output ends or nobody is
341/// listening.
342///
343/// Reading stops only there, never on a line's content: a reader that gave up
344/// would stop draining the daemon's pipe or PTY, and the daemon would block,
345/// or fail, on its next write.
346async fn forward_output_lines<R>(mut reader: R, tx: tokio::sync::mpsc::Sender<super::OutputLine>)
347where
348 R: tokio::io::AsyncBufRead + Unpin,
349{
350 use tokio::io::AsyncBufReadExt;
351 let mut line = Vec::new();
352 loop {
353 line.clear();
354 // A read error ends the output too: a Linux PTY master reports the
355 // slave closing, once the daemon is gone, as EIO rather than end of
356 // file. Whatever of a last, unterminated line arrived before it is
357 // still sent.
358 let last = match reader.read_until(b'\n', &mut line).await {
359 Ok(0) => break,
360 Ok(_) => false,
361 Err(_) if line.is_empty() => break,
362 Err(_) => true,
363 };
364 let sent = tx
365 .send(super::OutputLine {
366 text: output_line_text(&line),
367 source: super::OutputSource::Local,
368 })
369 .await;
370 if sent.is_err() || last {
371 break;
372 }
373 }
374}
375
376fn stop_keeps_finalized_status(status: &DaemonStatus) -> bool {
377 status.is_completed()
378}
379
380/// How long a failed start waits for the daemon's output to become queryable
381/// before reporting. Typically satisfied in a few dozen milliseconds; a daemon
382/// that failed without printing anything waits the whole of it, so keep it
383/// short.
384const SINK_OUTPUT_TIMEOUT: Duration = Duration::from_millis(400);
385
386/// Marks a daemon as having its retries managed by a foreground `run` for as
387/// long as this value lives, so the background checker does not start an
388/// attempt out from under it. Released on every exit from the retry loop,
389/// including the early returns.
390/// Counts a stop of this daemon once the stop is done, while its lock is
391/// still held. See `Supervisor::stop_epochs`.
392struct StopEpochGuard(DaemonId);
393
394impl Drop for StopEpochGuard {
395 fn drop(&mut self) {
396 SUPERVISOR.bump_stop_epoch(&self.0);
397 }
398}
399
400pub(crate) struct RetryingGuard {
401 id: DaemonId,
402 cancel: std::sync::Arc<std::sync::atomic::AtomicBool>,
403}
404
405impl RetryingGuard {
406 /// Whether a `stop` has asked this retry sequence to end.
407 fn is_cancelled(&self) -> bool {
408 self.cancel.load(std::sync::atomic::Ordering::Acquire)
409 }
410}
411
412impl Drop for RetryingGuard {
413 fn drop(&mut self) {
414 let mut retrying = SUPERVISOR
415 .retrying
416 .lock()
417 .unwrap_or_else(|e| e.into_inner());
418 // Drop this claim's flag only. Another sequence for the same daemon
419 // may still be running, and it has to stay both protected from the
420 // retry checker and reachable by a stop.
421 if let Some(claims) = retrying.get_mut(&self.id) {
422 claims.retain(|flag| !std::sync::Arc::ptr_eq(flag, &self.cancel));
423 if claims.is_empty() {
424 retrying.remove(&self.id);
425 }
426 }
427 }
428}
429
430impl Supervisor {
431 /// Whether a foreground `run` is already working through this daemon's
432 /// retries.
433 pub(crate) fn is_retrying(&self, id: &DaemonId) -> bool {
434 self.retrying
435 .lock()
436 .unwrap_or_else(|e| e.into_inner())
437 .get(id)
438 .is_some_and(|claims| !claims.is_empty())
439 }
440
441 /// How many times this daemon has been stopped so far.
442 pub(crate) fn stop_epoch(&self, id: &DaemonId) -> u64 {
443 self.stop_epochs
444 .lock()
445 .unwrap_or_else(|e| e.into_inner())
446 .get(id)
447 .copied()
448 .unwrap_or(0)
449 }
450
451 fn bump_stop_epoch(&self, id: &DaemonId) {
452 *self
453 .stop_epochs
454 .lock()
455 .unwrap_or_else(|e| e.into_inner())
456 .entry(id.clone())
457 .or_default() += 1;
458 }
459
460 fn mark_retrying(&self, id: &DaemonId) -> RetryingGuard {
461 let cancel = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
462 self.retrying
463 .lock()
464 .unwrap_or_else(|e| e.into_inner())
465 .entry(id.clone())
466 .or_default()
467 .push(cancel.clone());
468 RetryingGuard {
469 id: id.clone(),
470 cancel,
471 }
472 }
473
474 /// Ask a foreground retry sequence for this daemon, if there is one, to
475 /// end. A stop is a decision about the daemon, not about one of its
476 /// attempts, so the attempts left must not go ahead behind it.
477 pub(crate) fn cancel_retrying(&self, id: &DaemonId) {
478 if let Some(claims) = self
479 .retrying
480 .lock()
481 .unwrap_or_else(|e| e.into_inner())
482 .get(id)
483 {
484 // Every claim, not just the newest: a start that is sleeping out a
485 // backoff is as much a sequence the stop has to end as the one that
486 // claimed the daemon last.
487 for cancel in claims {
488 cancel.store(true, std::sync::atomic::Ordering::Release);
489 }
490 }
491 }
492
493 /// Run a daemon, handling retries if configured
494 pub async fn run(&self, opts: RunOptions) -> Result<IpcResponse> {
495 self.run_inner(opts, None).await
496 }
497
498 /// Run an attempt the retry checker decided on while `stop_epoch` read
499 /// `approved_at`. If the daemon has been stopped since, the attempt is
500 /// abandoned instead of started.
501 pub(crate) async fn run_retry(
502 &self,
503 opts: RunOptions,
504 approved_at: u64,
505 ) -> Result<IpcResponse> {
506 self.run_inner(opts, Some(approved_at)).await
507 }
508
509 async fn run_inner(&self, opts: RunOptions, approved_at: Option<u64>) -> Result<IpcResponse> {
510 let id = &opts.id;
511 let cmd = opts.cmd.clone();
512
513 // Clear any pending autostop for this daemon since it's being started
514 {
515 let mut pending = self.pending_autostops.lock().await;
516 if pending.remove(id).is_some() {
517 info!("cleared pending autostop for {id} (daemon starting)");
518 }
519 }
520
521 // Serialize against any in-flight stop of this daemon: a stop now
522 // waits for the whole process group to exit, so the Stopping window
523 // can last seconds instead of milliseconds. Starting through that
524 // window would collide with the dying instance (duplicate processes,
525 // port conflicts). Acquiring the stop lock waits the stop out; the
526 // state is re-read afterwards. The guard is owned and handed to
527 // run_once, which holds it until the new daemon's Running state and
528 // PID are persisted — releasing it before that point would let a
529 // concurrent run pass this same check (duplicate processes) or let a
530 // concurrent stop see no PID and return without stopping anything.
531 let mut stop_guard = Some(self.stop_lock(id).await.lock_owned().await);
532 // Checked here, under the daemon's lock, because that is what a stop
533 // takes too: an approval from before the stop cannot slip past it.
534 // Writing `stopped` over the record is not enough on its own, since an
535 // attempt already approved would read that as a daemon free to start.
536 if let Some(approved_at) = approved_at
537 && self.stop_epoch(id) != approved_at
538 {
539 info!("daemon {id} was stopped after this retry was decided on; not starting it");
540 return Ok(IpcResponse::DaemonNotRunning);
541 }
542 if let Some(response) = self.claim_or_defer(&opts, &mut stop_guard).await? {
543 return Ok(response);
544 }
545
546 // If wait_ready is true and retry is configured, implement retry loop
547 if opts.wait_ready && opts.retry.count() > 0 {
548 // Claim this daemon's retries for the duration of the loop. The
549 // backoff between attempts leaves the record errored with no PID,
550 // which is what `check_retry` scans for, and an attempt started
551 // there would leave this call reporting on a run it does not own.
552 let retrying_claim = self.mark_retrying(id);
553 // Use saturating_add to avoid overflow when retry = u32::MAX (infinite)
554 let max_attempts = opts.retry.count().saturating_add(1);
555 for attempt in 0..max_attempts {
556 let mut retry_opts = opts.clone();
557 retry_opts.retry_count = attempt;
558 retry_opts.cmd = cmd.clone();
559
560 // The first attempt starts under the guard held since the
561 // running check above; later attempts re-acquire it so stops
562 // are not locked out during the backoff sleeps.
563 let mut guard = Some(match stop_guard.take() {
564 Some(guard) => guard,
565 None => self.stop_lock(id).await.lock_owned().await,
566 });
567 // Ownership has to be re-checked on every attempt, not just
568 // the first. The backoff leaves the daemon errored with no PID,
569 // which is exactly what `check_retry` looks for, so the
570 // background checker can start the next attempt during the
571 // sleep. Spawning another process here would replace that
572 // attempt's monitor registration and leave its process running
573 // unmonitored.
574 if let Some(response) = self.claim_or_defer(&retry_opts, &mut guard).await? {
575 return Ok(response);
576 }
577 // The background retry checker may have run this attempt for
578 // us and seen it succeed while we slept. Starting again would
579 // repeat a task that has already done its work — for a
580 // migration or a seed, repeating its side effects.
581 //
582 // Only after a backoff, though. A completed record on the first
583 // attempt is the previous run's, and a start is defined to
584 // re-run a completed oneshot; short-circuiting here would make
585 // that true only for oneshots without `retry`.
586 if attempt > 0
587 && let Some(daemon) = self.get_daemon(id).await
588 && daemon.status.is_completed()
589 {
590 info!("daemon {id} completed while waiting to retry; not running it again");
591 return Ok(IpcResponse::DaemonReady { daemon });
592 }
593 // A stop that arrived during the backoff ends the sequence.
594 // Without this the loop would start the next attempt on a
595 // daemon the user has just stopped, and the stop would look
596 // like it had done nothing.
597 if retrying_claim.is_cancelled() {
598 info!("daemon {id} was stopped while waiting to retry; abandoning its retries");
599 return Ok(IpcResponse::DaemonFailed {
600 error: "stopped while retrying".to_string(),
601 });
602 }
603 let Some(guard) = guard else {
604 // Only the deferring paths take the guard, and each of
605 // those returned above.
606 return Ok(IpcResponse::DaemonAlreadyRunning);
607 };
608 let result = self.run_once(retry_opts, guard).await?;
609
610 match result {
611 IpcResponse::DaemonReady { daemon } => {
612 return Ok(IpcResponse::DaemonReady { daemon });
613 }
614 IpcResponse::DaemonFailedWithCode {
615 exit_code,
616 resolved_ports,
617 } => {
618 if attempt < opts.retry.count() {
619 // `run_once` reports failure the moment the process
620 // exits, but its monitor finalizes the record only
621 // after draining the process's remaining output.
622 // Until then the record still names this attempt's
623 // PID, and the next attempt's ownership check would
624 // read its own dead predecessor as a competing run
625 // and abandon the retries that are left.
626 let attempt_pid = self.get_daemon(id).await.and_then(|d| d.pid);
627 self.wait_for_exit_finalized(id, attempt_pid).await;
628 let backoff_secs = 2u64.saturating_pow(attempt).min(3600);
629 info!(
630 "daemon {id} failed (attempt {}/{}), retrying in {}s",
631 attempt + 1,
632 max_attempts,
633 backoff_secs
634 );
635 fire_hook(
636 HookType::OnRetry,
637 id.clone(),
638 opts.dir.0.clone(),
639 attempt + 1,
640 opts.env.clone(),
641 resolved_ports,
642 vec![],
643 )
644 .await;
645 // Slept in slices so a stop arriving during a
646 // long backoff — they grow to an hour — is acted
647 // on when it arrives rather than when the sleep
648 // happens to end.
649 let backoff_deadline =
650 tokio::time::Instant::now() + Duration::from_secs(backoff_secs);
651 while tokio::time::Instant::now() < backoff_deadline
652 && !retrying_claim.is_cancelled()
653 {
654 let remaining = backoff_deadline - tokio::time::Instant::now();
655 time::sleep(remaining.min(Duration::from_millis(200))).await;
656 }
657 continue;
658 } else {
659 info!("daemon {id} failed after {max_attempts} attempts");
660 return Ok(IpcResponse::DaemonFailedWithCode {
661 exit_code,
662 resolved_ports,
663 });
664 }
665 }
666 other => return Ok(other),
667 }
668 }
669 }
670
671 // No retry or wait_ready is false
672 let guard = match stop_guard.take() {
673 Some(guard) => guard,
674 None => self.stop_lock(id).await.lock_owned().await,
675 };
676 self.run_once(opts, guard).await
677 }
678
679 /// Wait for a just-failed attempt's monitor to write its terminal state,
680 /// clearing the PID from the record.
681 ///
682 /// Bounded a little beyond the monitor's own five-second output drain, the
683 /// longest it can hold the record after the process has gone. Giving up
684 /// early is safe: the ownership check that follows simply sees a PID and
685 /// defers, which is what it would have done anyway.
686 ///
687 /// `pid` names the run being waited for, so a record that has moved on to
688 /// another run is not mistaken for this one still finishing.
689 async fn wait_for_exit_finalized(&self, id: &DaemonId, pid: Option<u32>) {
690 let deadline = tokio::time::Instant::now() + Duration::from_secs(8);
691 loop {
692 match self.get_daemon(id).await {
693 Some(daemon) if pid.map_or(daemon.pid.is_some(), |pid| daemon.pid == Some(pid)) => {
694 }
695 _ => return,
696 }
697 if tokio::time::Instant::now() >= deadline {
698 debug!("daemon {id}: previous attempt has not finalized yet; continuing anyway");
699 return;
700 }
701 time::sleep(Duration::from_millis(50)).await;
702 }
703 }
704
705 /// Decide whether this start may take the daemon's record, or must stand
706 /// down because a live run already owns it.
707 ///
708 /// Returns `Some(response)` when the caller must report that run's outcome
709 /// instead of spawning a second process, and `None` when the record is free
710 /// (including after a forced stop of the previous instance).
711 ///
712 /// `stop_guard` is released before an in-flight oneshot is awaited: that
713 /// wait lasts as long as the task does, and holding the lock would block a
714 /// stop of the very run being waited on.
715 async fn claim_or_defer(
716 &self,
717 opts: &RunOptions,
718 stop_guard: &mut Option<tokio::sync::OwnedMutexGuard<()>>,
719 ) -> Result<Option<IpcResponse>> {
720 let id = &opts.id;
721 let Some(daemon) = self.get_daemon(id).await else {
722 return Ok(None);
723 };
724 // Entering a directory does not re-run a finished task, at any level of
725 // the dependency graph — the layout this exists for reaches the task
726 // through a service's `depends`, not by naming it. Decided here rather
727 // than in the client because this is the authoritative state: the state
728 // file lags it by up to the flush interval, which is exactly the window
729 // a second entry lands in after the task completes.
730 // `opts.oneshot` rather than the record's: the request carries what
731 // config says now, while the stored flag is only refreshed by a run, so
732 // a daemon that used to be a task would otherwise stay skipped forever
733 // after being turned into a service.
734 if opts.on_directory_enter && opts.oneshot && daemon.status.is_completed() {
735 debug!("daemon {id} already completed; directory entry leaves it alone");
736 return Ok(Some(IpcResponse::DaemonReady { daemon }));
737 }
738 // Stopping is treated as "not running": the monitoring task will clean
739 // it up. Only a live PID under a non-terminal status blocks a start.
740 if daemon.status.is_stopping() || daemon.status.is_stopped() || daemon.status.is_completed()
741 {
742 return Ok(None);
743 }
744 let Some(pid) = daemon.pid else {
745 return Ok(None);
746 };
747 if opts.force {
748 self.stop_locked(id).await?;
749 info!("run: stop completed for daemon {id}");
750 return Ok(None);
751 }
752 if daemon.oneshot && opts.wait_ready {
753 // An in-flight oneshot has not done its work yet, so reporting
754 // "already running" would let dependents start against the state
755 // the task is still establishing. Wait for the run already under
756 // way instead.
757 info!("daemon {id} is an in-flight oneshot (pid {pid}); waiting for it to finish");
758 drop(stop_guard.take());
759 return Ok(Some(
760 self.await_running_oneshot(id, opts.oneshot_wait, pid).await,
761 ));
762 }
763 // A record can name a PID that has already exited: `stop` leaves the
764 // terminal state to a monitor that still owns the daemon, and that
765 // monitor writes it only after draining the process's output. Rejecting
766 // a start against a dead PID would fail an ordinary stop-then-start for
767 // the length of that drain, so confirm the process is really there
768 // before refusing. The oneshot branch above deliberately comes first: a
769 // task whose process has exited is about to be recorded as completed,
770 // and starting a second copy of it is exactly what waiting prevents.
771 PROCS.refresh_pids(&[pid]);
772 if !PROCS.is_running(pid) {
773 debug!(
774 "daemon {id}: record still names pid {pid}, which has exited; its monitor has not finalized yet"
775 );
776 return Ok(None);
777 }
778 warn!("daemon {id} already running with pid {pid}");
779 Ok(Some(IpcResponse::DaemonAlreadyRunning))
780 }
781
782 /// Wait for a oneshot that is already running to reach a terminal state,
783 /// and report it as if this call had started the task itself.
784 ///
785 /// Polls the state file because the terminal state is written by the
786 /// monitoring task of the *other* run; this call has no readiness channel
787 /// of its own to await.
788 async fn await_running_oneshot(
789 &self,
790 id: &DaemonId,
791 wait: Option<OneshotWait>,
792 watched_pid: u32,
793 ) -> IpcResponse {
794 let interval = settings().supervisor_ready_check_interval();
795 // The caller resolved this from the project's settings and sent it, so
796 // both processes wait exactly as long. Falling back to this process's
797 // own settings would read the directory the supervisor happens to have
798 // started in, where a project's `oneshot_timeout` is not visible — and
799 // the shorter of the two deadlines would silently win, releasing
800 // dependents while the task was still running.
801 //
802 // `None` here means the setting asked for no limit, so there is no
803 // deadline to reach rather than a distant one.
804 // One deadline for the whole wait, retries and backoffs included.
805 // `oneshot_timeout` is documented as the longest `pitchfork start` will
806 // wait, and the client bounds its own request by the same value without
807 // restarting it, so a per-attempt budget here would both break that
808 // promise — unboundedly, with infinite retries — and put the two sides
809 // back to disagreeing about when one task has gone on too long.
810 let deadline = wait
811 .unwrap_or_else(|| settings().supervisor_oneshot_wait())
812 .duration()
813 .map(|d| tokio::time::Instant::now() + d);
814 // Which run this wait is reporting on. A terminal state is only that
815 // run's while the record still names its PID or names none at all; once
816 // another PID appears, something else has started the task and the
817 // outcome that follows belongs to that run, not this one. Following the
818 // handoff keeps the answer useful to a dependent — it still learns
819 // whether the task succeeded — without quietly attributing an unrelated
820 // run's failure to the one it asked about.
821 let mut watched_pid = watched_pid;
822 loop {
823 let Some(daemon) = self.get_daemon(id).await else {
824 return IpcResponse::DaemonNotFound;
825 };
826 if let Some(current) = daemon.pid
827 && current != watched_pid
828 {
829 info!(
830 "daemon {id}: the run being waited on (pid {watched_pid}) was replaced by pid {current}; following it"
831 );
832 watched_pid = current;
833 }
834 match &daemon.status {
835 DaemonStatus::Completed => {
836 info!("daemon {id}: the in-flight oneshot completed");
837 return IpcResponse::DaemonReady { daemon };
838 }
839 DaemonStatus::Errored(code) => {
840 // A failed attempt is persisted before the in-flight `run`
841 // sleeps out its backoff, so an errored record with
842 // attempts left is a gap between tries rather than the
843 // result. Same condition `check_retry` uses to decide
844 // whether another attempt is still owed.
845 if daemon.retry.count() > 0 && daemon.retry_count < daemon.retry.count() {
846 debug!(
847 "daemon {id}: in-flight oneshot failed attempt {} of {}; still waiting",
848 daemon.retry_count + 1,
849 daemon.retry.count() + 1
850 );
851 } else {
852 // -1 records an unobservable exit code; the caller
853 // renders `None` as a plain failure rather than
854 // "exit code -1".
855 let exit_code = Some(*code).filter(|c| *c != -1);
856 return IpcResponse::DaemonFailedWithCode {
857 exit_code,
858 resolved_ports: daemon.resolved_port.clone(),
859 };
860 }
861 }
862 DaemonStatus::Failed(error) => {
863 return IpcResponse::DaemonFailed {
864 error: error.clone(),
865 };
866 }
867 DaemonStatus::Stopped => {
868 // Stopped, not completed: the task was interrupted, so it
869 // never established what its dependents are waiting for.
870 warn!("daemon {id}: the in-flight oneshot was stopped before completing");
871 return IpcResponse::DaemonFailedWithCode {
872 exit_code: None,
873 resolved_ports: daemon.resolved_port.clone(),
874 };
875 }
876 DaemonStatus::Running | DaemonStatus::Waiting | DaemonStatus::Stopping => {}
877 }
878 if deadline.is_some_and(|deadline| tokio::time::Instant::now() >= deadline) {
879 warn!("daemon {id}: gave up waiting for the in-flight oneshot to finish");
880 // Reported as a failure rather than as "already running": the
881 // batch start path only counts a result carrying an exit code
882 // as failed, so anything else would let dependents start
883 // against a task that never finished. 124 is the code a
884 // readiness timeout already uses.
885 return IpcResponse::DaemonFailedWithCode {
886 exit_code: Some(124),
887 resolved_ports: Vec::new(),
888 };
889 }
890 time::sleep(interval).await;
891 }
892 }
893
894 /// Run a daemon once (single attempt).
895 ///
896 /// `stop_guard` is this daemon's stop lock, acquired by `run` before the
897 /// already-running check. It is held through spawning until the Running
898 /// state and PID are persisted (or an early failure returns), then dropped
899 /// before the potentially unbounded readiness wait.
900 pub(crate) async fn run_once(
901 &self,
902 opts: RunOptions,
903 stop_guard: tokio::sync::OwnedMutexGuard<()>,
904 ) -> Result<IpcResponse> {
905 let id = &opts.id;
906 let original_cmd = opts.cmd.clone(); // Save original command for persistence
907
908 // Create channel for readiness notification if wait_ready is true
909 let (ready_tx, ready_rx) = if opts.wait_ready {
910 let (tx, rx) = oneshot::channel();
911 (Some(tx), Some(rx))
912 } else {
913 (None, None)
914 };
915
916 // Check port availability and apply auto-bump if configured
917 let expected_ports = opts
918 .port
919 .as_ref()
920 .map(|p| p.expect.clone())
921 .unwrap_or_default();
922 let (resolved_ports, effective_ready_port) = if !expected_ports.is_empty() {
923 let port_cfg = opts.port.as_ref().unwrap();
924 match check_ports_available(
925 &expected_ports,
926 port_cfg.auto_bump(),
927 port_cfg.max_bump_attempts(),
928 )
929 .await
930 {
931 Ok(resolved) => {
932 let ready_port = if let Some(configured_port) =
933 opts.ready_port.as_ref().and_then(|p| p.as_port())
934 {
935 Some(resolve_configured_ready_port(
936 configured_port,
937 &expected_ports,
938 &resolved,
939 ))
940 } else if opts.ready_output.is_none()
941 && opts.ready_http.is_none()
942 && opts.ready_cmd.is_none()
943 && opts.ready_delay.is_none()
944 {
945 // No other ready check configured — use the first expected port as a
946 // TCP port readiness check so the daemon is considered ready once it
947 // starts listening. Skip port 0 (ephemeral port request).
948 resolved.first().copied().filter(|&p| p != 0)
949 } else {
950 // Another ready check is configured (output/http/cmd/delay).
951 // Don't add an implicit TCP port check — it could race and fire
952 // before the daemon has produced any output.
953 None
954 };
955 info!("daemon {id}: ports {expected_ports:?} resolved to {resolved:?}");
956 (resolved, ready_port)
957 }
958 Err(e) => {
959 error!("daemon {id}: port check failed: {e}");
960 // Convert PortError to structured IPC response
961 if let Some(port_error) = e.downcast_ref::<PortError>() {
962 match port_error {
963 PortError::InUse { port, process, pid } => {
964 return Ok(IpcResponse::PortConflict {
965 port: *port,
966 process: process.clone(),
967 pid: *pid,
968 });
969 }
970 PortError::NoAvailablePort {
971 start_port,
972 attempts,
973 } => {
974 return Ok(IpcResponse::NoAvailablePort {
975 start_port: *start_port,
976 attempts: *attempts,
977 });
978 }
979 }
980 }
981 return Ok(IpcResponse::DaemonFailed {
982 error: e.to_string(),
983 });
984 }
985 }
986 } else {
987 // When ready_port is set without expected_port, check that the port
988 // is not already occupied. If another process is listening on it,
989 // the TCP readiness probe would immediately succeed and pitchfork
990 // would falsely consider the daemon ready — routing proxy traffic to
991 // the wrong process.
992 if let Some(port) = opts.ready_port.as_ref().and_then(|p| p.as_port())
993 && port > 0
994 && let Some((pid, process)) = detect_port_conflict(port).await
995 {
996 return Ok(IpcResponse::PortConflict { port, process, pid });
997 }
998 (
999 Vec::new(),
1000 opts.ready_port.as_ref().and_then(|p| p.as_port()),
1001 )
1002 };
1003
1004 // The program and arguments that start the daemon, before any mise
1005 // wrapping.
1006 // The program and arguments that start the daemon, before any mise
1007 // wrapping, and the script for the shell when `run` is a string.
1008 let (mut words, script) = if opts.no_shell {
1009 // The argv form of `run`: started as written, with no shell to
1010 // reinterpret quotes, `%`, `&` or anything else in the arguments.
1011 if let Some(error) = invalid_argv_program(id, &original_cmd) {
1012 return Ok(IpcResponse::DaemonFailed { error });
1013 }
1014 (original_cmd.clone(), None)
1015 } else {
1016 // Resolve the shell for this platform into program + args. The run
1017 // script is passed verbatim as the final argument, avoiding the lossy
1018 // split->join round-trip that previously mangled $VAR/glob expansion.
1019 let words = match resolve_shell() {
1020 Ok(parts) => parts,
1021 Err(error) => return Ok(IpcResponse::DaemonFailed { error }),
1022 };
1023 // Use the original run string verbatim; fall back to joining cmd for
1024 // ad-hoc commands (e.g. `pitchfork run -- cmd args`) that have no run string.
1025 // We don't prepend `exec` because it breaks compound commands (e.g. `exec a && b`
1026 // silently drops `b`). Users can add `exec` themselves in the run string.
1027 let script = opts
1028 .run
1029 .clone()
1030 .unwrap_or_else(|| shell_words::join(&original_cmd));
1031 (words, Some(script))
1032 };
1033
1034 let mise_bin = if opts.mise.unwrap_or(settings().general.mise) {
1035 let mise_bin = settings().resolve_mise_bin();
1036 if mise_bin.is_none() {
1037 warn!("daemon {id}: mise=true but mise binary not found, running without mise");
1038 }
1039 mise_bin
1040 } else {
1041 None
1042 };
1043 // Started directly, the shell gets its script from `shell_script`.
1044 // Under mise it goes in as an ordinary argument: mise starts the shell
1045 // itself, re-quoting each argument, so the raw command line cmd.exe
1046 // needs for a script with `"` cannot reach it that way.
1047 let script = match &mise_bin {
1048 Some(mise_bin) => {
1049 info!(
1050 "daemon {id}: wrapping command with mise ({})",
1051 mise_bin.display()
1052 );
1053 words.extend(script);
1054 None
1055 }
1056 None => script,
1057 };
1058 let (program, args) = launch_command(words, mise_bin.as_deref());
1059 #[cfg(unix)]
1060 let run_identity = match resolve_effective_run_identity(opts.user.as_deref()) {
1061 Ok(identity) => identity,
1062 Err(e) => {
1063 return Ok(IpcResponse::DaemonFailed {
1064 error: e.to_string(),
1065 });
1066 }
1067 };
1068 info!("run: spawning daemon {id} with {program} {args:?} {script:?}");
1069
1070 // Allocate PTY if configured
1071 #[cfg(unix)]
1072 let pty_pair = if opts.pty.unwrap_or(false) {
1073 match super::pty::openpty() {
1074 Ok(pair) => {
1075 info!("daemon {id}: allocated PTY (pty = true)");
1076 Some(pair)
1077 }
1078 Err(e) => {
1079 warn!("daemon {id}: failed to allocate PTY, falling back to pipes: {e}");
1080 None
1081 }
1082 }
1083 } else {
1084 None
1085 };
1086
1087 // Output reaches the monitoring task either from readers this process
1088 // owns or, when a sink owns the stream, relayed over IPC. The channel is
1089 // created here rather than in that task so it exists before the sink
1090 // starts: a daemon whose very first line matches its readiness pattern
1091 // would otherwise have the match reported with nowhere to deliver it.
1092 let (output_tx, output_rx) = tokio::sync::mpsc::channel::<super::OutputLine>(256);
1093 let mut output_relay = None;
1094
1095 // Set up out-of-process capture before building the command, so the
1096 // daemon can be handed the pipe's write end directly.
1097 let mut sink_pipe = None;
1098 let mut sink_writer = None;
1099 let mut sink_child = None;
1100 if super::log_sink::is_supported(&opts) {
1101 let log_format = opts
1102 .log_format
1103 .clone()
1104 .unwrap_or_else(|| settings().logs.log_format.clone());
1105 let watch_for = super::log_sink::WatchFor::from_opts(id, &opts);
1106 // The token ties this attempt's sink to this attempt's channel, so
1107 // a sink still draining a previous attempt cannot report into it.
1108 let relay_token = if watch_for.is_empty() {
1109 0
1110 } else {
1111 let relay = super::log_sink::OutputRelay::register(id, output_tx.clone());
1112 let token = relay.token();
1113 output_relay = Some(relay);
1114 token
1115 };
1116 match super::log_sink::SinkPipe::new(log_format, watch_for, relay_token) {
1117 Ok((pipe, writer)) => match pipe.start(id) {
1118 Ok(child) => {
1119 sink_child = Some(super::log_sink::PendingSink::new(child));
1120 sink_pipe = Some(pipe);
1121 sink_writer = Some(writer);
1122 }
1123 Err(e) => {
1124 warn!("could not start log sink for {id}, capturing in-process: {e}");
1125 }
1126 },
1127 Err(e) => {
1128 // Fall back to in-process capture rather than refusing to
1129 // start the daemon.
1130 warn!("could not create log pipe for {id}, capturing in-process: {e}");
1131 }
1132 }
1133 }
1134
1135 let mut cmd = tokio::process::Command::new(&program);
1136
1137 #[cfg(unix)]
1138 if let Some(ref pair) = pty_pair {
1139 // PTY mode: connect both stdout and stderr to the slave PTY.
1140 // The child uses the slave for stdin/stdout/stderr, and we read
1141 // output from the master.
1142 let slave_file = std::fs::File::from(
1143 pair.slave
1144 .try_clone()
1145 .map_err(|e| miette::miette!("failed to dup slave PTY fd: {e}"))?,
1146 );
1147 cmd.stdin(std::process::Stdio::from(slave_file.try_clone().map_err(
1148 |e| miette::miette!("failed to clone slave PTY fd for stdin: {e}"),
1149 )?));
1150 cmd.stdout(std::process::Stdio::from(slave_file.try_clone().map_err(
1151 |e| miette::miette!("failed to clone slave PTY fd for stdout: {e}"),
1152 )?));
1153 cmd.stderr(std::process::Stdio::from(slave_file));
1154 } else if let Some(writer) = sink_writer.take() {
1155 // Capture belongs to a sibling sink process, so the daemon writes
1156 // to a pipe this process does not read. See supervisor::log_sink.
1157 let dup = writer
1158 .try_clone()
1159 .map_err(|e| miette::miette!("failed to dup log pipe for stderr: {e}"))?;
1160 cmd.stdout(std::process::Stdio::from(writer))
1161 .stderr(std::process::Stdio::from(dup));
1162 } else {
1163 cmd.stdout(std::process::Stdio::piped())
1164 .stderr(std::process::Stdio::piped());
1165 }
1166
1167 #[cfg(not(unix))]
1168 if let Some(writer) = sink_writer.take() {
1169 let dup = writer
1170 .try_clone()
1171 .map_err(|e| miette::miette!("failed to dup log pipe for stderr: {e}"))?;
1172 cmd.stdout(std::process::Stdio::from(writer))
1173 .stderr(std::process::Stdio::from(dup));
1174 } else {
1175 cmd.stdout(std::process::Stdio::piped())
1176 .stderr(std::process::Stdio::piped());
1177 }
1178
1179 match &script {
1180 Some(script) => cmd.shell_script(&program, &args, script),
1181 None => cmd.args(&args),
1182 };
1183 cmd.current_dir(&opts.dir).hide_console_window();
1184
1185 #[cfg(unix)]
1186 if pty_pair.is_none() {
1187 cmd.stdin(std::process::Stdio::null());
1188 }
1189
1190 #[cfg(not(unix))]
1191 cmd.stdin(std::process::Stdio::null());
1192
1193 // Before the runtime env, so a daemon's own `env` entries win.
1194 #[cfg(unix)]
1195 apply_identity_env(&mut cmd, &run_identity);
1196 apply_runtime_env(
1197 &mut cmd,
1198 id,
1199 opts.retry_count,
1200 opts.env.as_ref(),
1201 &resolved_ports,
1202 );
1203
1204 // Inject proxy-related environment variables
1205 inject_proxy_env(&mut cmd, &daemon_proxy_host(&opts).await);
1206
1207 #[cfg(unix)]
1208 {
1209 let run_identity = run_identity.clone();
1210 let use_pty = pty_pair.is_some();
1211 unsafe {
1212 cmd.pre_exec(move || {
1213 nix::unistd::setsid().map_err(nix_to_io_error)?;
1214
1215 // When using a PTY, set the slave as the controlling terminal.
1216 // The slave FD has already been dup'd onto stdin/stdout/stderr
1217 // by tokio, so we can use stdin (fd 0) for TIOCSCTTY.
1218 if use_pty {
1219 let ret = libc::ioctl(0, libc::TIOCSCTTY as libc::c_ulong, 0);
1220 if ret < 0 {
1221 // Non-fatal: the process can still run without
1222 // a controlling terminal.
1223 #[cfg(target_os = "linux")]
1224 eprintln!(
1225 "pitchfork: TIOCSCTTY failed: {}",
1226 std::io::Error::last_os_error()
1227 );
1228 }
1229 }
1230
1231 apply_run_identity(&run_identity)?;
1232 Ok(())
1233 });
1234 }
1235 }
1236
1237 // Timestamp the run so a failed start can wait for this attempt's output
1238 // specifically, rather than seeing an earlier attempt's.
1239 let spawn_time = chrono::Local::now();
1240 // A sink is already running at this point. Both bail-outs below have to
1241 // reap it explicitly: dropping the handle only reaps on a best-effort
1242 // basis, and run_once runs once per retry attempt, so a daemon that
1243 // consistently fails to spawn would otherwise accumulate sinks.
1244 // A failed spawn returns here; the sink is terminated by PendingSink.
1245 let mut child = cmd.spawn().into_diagnostic()?;
1246 // A process now exists, which is exactly what `last_cron_run` records.
1247 // Written here rather than from the watcher's view of the response
1248 // because that view cannot tell a start that failed before spawning
1249 // from one that spawned and exited before its PID could be read: a
1250 // port conflict, an unresolvable shell and an instant exit all report
1251 // `DaemonFailed`. Only the last of those ran, and a cron job that
1252 // fails that fast is precisely the one whose timing a user needs.
1253 // Ordered before the `Running` upsert below, which inherits it.
1254 //
1255 // Written synchronously rather than left to the background flush, for
1256 // the same reason `last_cron_triggered` is: a supervisor that dies in
1257 // the window between the two would come back with no record that this
1258 // run happened, and for a short-lived job that window is as long as
1259 // the job itself. A scheduled spawn is rare enough -- at most one per
1260 // `cron_check_interval` -- for the extra write to cost nothing.
1261 if opts.cron_started {
1262 let mut state_file = self.state_file.lock().await;
1263 if state_file.set_last_cron_run(id, spawn_time)
1264 && let Err(e) = state_file.write()
1265 {
1266 error!("failed to persist last_cron_run for daemon {id}: {e}");
1267 }
1268 }
1269 let pid = match child.id() {
1270 Some(p) => p,
1271 None => {
1272 warn!("Daemon {id} exited before PID could be captured");
1273 // Unlike a daemon that never started, this one ran and may have
1274 // said why it gave up, and its output is the only diagnosis
1275 // available. Its write end is already closed, so the sink is on
1276 // its way to end of file: let it finish writing before reporting,
1277 // then reap whatever is left of it.
1278 if sink_child.is_some() {
1279 super::log_sink::wait_for_output(id, spawn_time, SINK_OUTPUT_TIMEOUT).await;
1280 }
1281 return Ok(IpcResponse::DaemonFailed {
1282 error: "Process exited immediately".to_string(),
1283 });
1284 }
1285 };
1286 info!("started daemon {id} with pid {pid}");
1287 PROCS.refresh_pids(&[pid]);
1288 // Register the daemon as monitored BEFORE persisting the Running
1289 // state. The orphan reconciler treats any running, unmonitored PID
1290 // as an orphan; if the state became visible first, a concurrent
1291 // reconciliation pass could adopt — or under the kill policy,
1292 // terminate — a daemon that was just legitimately started. The RAII
1293 // guard unregisters on any early-error path below and is otherwise
1294 // handed to the monitoring task.
1295 let monitored_guard = super::adopt::MonitoredGuard::register(id.clone(), pid);
1296 let monitor_token = monitored_guard.token();
1297
1298 // Hand the retained read end to a sink and keep one running for as long
1299 // as this daemon is monitored.
1300 let using_sink = sink_pipe.is_some();
1301 // Take the sink out of the guard only once there is a pipe to supervise
1302 // it with, so it is never left running unsupervised.
1303 if let Some(pipe) = sink_pipe.take()
1304 && let Some(child) = sink_child.as_mut().and_then(|pending| pending.take())
1305 {
1306 pipe.supervise(id.clone(), monitor_token, child);
1307 }
1308 // The attempt's actual resolved ports, captured before the upsert
1309 // moves them into state. Hooks, readiness probes, and the failure
1310 // response must reflect this attempt: the state merge keeps the
1311 // existing resolved_port when an update is empty, so a no-port
1312 // attempt would otherwise inherit a previous run's stale ports
1313 // through the upserted record.
1314 let attempt_resolved_ports = resolved_ports.clone();
1315 let daemon = self
1316 .upsert_daemon(
1317 UpsertDaemonOpts::from_run_options(&opts, DaemonStatus::Running)
1318 .set(|o| {
1319 o.pid = Some(pid);
1320 o.cmd = Some(original_cmd);
1321 o.ready_port = effective_ready_port.map(|p| ReadyPort {
1322 port: Some(p),
1323 template: None,
1324 timeout: opts.ready_port.as_ref().and_then(|rp| rp.timeout),
1325 });
1326 o.port = crate::config_types::PortConfig::from_parts(
1327 expected_ports,
1328 opts.port.as_ref().map(|p| p.bump).unwrap_or_default(),
1329 );
1330 o.resolved_port = Some(resolved_ports);
1331 })
1332 .build(),
1333 )
1334 .await?;
1335
1336 // Running state and PID are now persisted: concurrent run/stop calls
1337 // observe a running daemon and behave correctly, so release the stop
1338 // lock rather than holding it through the readiness wait below, which
1339 // can take arbitrarily long.
1340 drop(stop_guard);
1341
1342 let id_clone = id.clone();
1343 // A oneshot is ready only when its process exits 0, so no readiness
1344 // check may run alongside it — one that fired first would report the
1345 // task ready before it had done its work, and would suppress the
1346 // completion notification entirely. Config load rejects explicit
1347 // `ready_*` fields and the client clears CLI overrides, but the
1348 // implicit port check is derived here from `port.expect`, so the
1349 // suppression has to happen here rather than being trusted to callers.
1350 let ready_delay = (!opts.oneshot).then_some(opts.ready_delay).flatten();
1351 let ready_output = (!opts.oneshot).then(|| opts.ready_output.clone()).flatten();
1352 let ready_http = (!opts.oneshot).then(|| opts.ready_http.clone()).flatten();
1353 let ready_port = (!opts.oneshot).then_some(effective_ready_port).flatten();
1354 let implicit_ready_port = ready_port.map(|p| ReadyPort {
1355 port: Some(p),
1356 template: None,
1357 timeout: None,
1358 });
1359 let ready_port_config = (!opts.oneshot)
1360 .then(|| opts.ready_port.clone())
1361 .flatten()
1362 .or(implicit_ready_port);
1363 let ready_cmd = (!opts.oneshot).then(|| opts.ready_cmd.clone()).flatten();
1364 let daemon_dir = opts.dir.0.clone();
1365 let hook_retry_count = opts.retry_count;
1366 let hook_retry = opts.retry;
1367 let hook_daemon_env = opts.env.clone();
1368 // Ports of THIS attempt, snapshotted before the monitor starts: a retry
1369 // or restart may replace state.resolved_port before a hook task runs.
1370 // Sourced from the attempt's local value, not the upserted record —
1371 // the state merge inherits stale ports for a no-port attempt.
1372 let hook_resolved_ports = attempt_resolved_ports.clone();
1373 let readiness_daemon_env = opts.env.clone();
1374 let readiness_resolved_ports = attempt_resolved_ports.clone();
1375 let on_output_hook = opts.on_output_hook.clone();
1376 // Whether this daemon has any port-related config — used to skip the
1377 // active_port detection task for daemons that never bind a port (e.g. `sleep 60`).
1378 // When the proxy is enabled, only detect active_port for daemons that are
1379 // actually referenced by a registered slug, rather than blanket-polling every
1380 // daemon (which wastes ~7.5 s of listeners::get_all() calls per port-less daemon).
1381 let has_port_config = opts.port.as_ref().is_some_and(|p| !p.expect.is_empty())
1382 || (settings().proxy.enable && is_daemon_slug_target(id));
1383 // When the ready_port check succeeds on the first resolved port we can
1384 // set active_port directly instead
1385 // of spawning detect_and_store_active_port (which relies on
1386 // listeners::get_all() + process-tree traversal and is unreliable on
1387 // Windows where Git Bash PID mapping can break descendant lookups).
1388 let daemon_pid = pid;
1389
1390 // Prepare output readers before spawning the monitoring task.
1391 // In PTY mode, we read from the PTY master FD.
1392 // In pipe mode, we read from separate stdout/stderr pipes.
1393 #[cfg(unix)]
1394 let pty_reader = pty_pair.map(|p| {
1395 tokio::io::BufReader::new(tokio::fs::File::from_std(std::fs::File::from(p.master)))
1396 });
1397 #[cfg(not(unix))]
1398 let pty_reader: Option<tokio::io::BufReader<tokio::fs::File>> = None;
1399 let stdout_reader = if pty_reader.is_none() {
1400 child.stdout.take().map(tokio::io::BufReader::new)
1401 } else {
1402 None
1403 };
1404 let stderr_reader = if pty_reader.is_none() {
1405 child.stderr.take().map(tokio::io::BufReader::new)
1406 } else {
1407 None
1408 };
1409
1410 if !using_sink
1411 && pty_reader.is_none()
1412 && (stdout_reader.is_none() || stderr_reader.is_none())
1413 {
1414 error!("Failed to capture stdout/stderr for daemon {id}");
1415 }
1416
1417 tokio::spawn(async move {
1418 let id = id_clone;
1419 // Registered before the Running upsert above; unregisters when
1420 // this monitoring task ends.
1421 let _monitored_guard = monitored_guard;
1422 // Likewise for sink-relayed output: dropping this stops the
1423 // supervisor delivering into a channel nobody is reading. Dropped
1424 // explicitly once the daemon exits, before the drain below.
1425 let output_relay = output_relay;
1426
1427 // Merge all output sources (PTY master OR stdout+stderr, or a
1428 // sink's IPC reports) into a single channel.
1429 let mut output_rx = output_rx;
1430
1431 if let Some(reader) = pty_reader {
1432 // PTY mode: single merged stream from the master.
1433 // output_tx is moved into the spawn; when the reader ends the
1434 // channel closes automatically.
1435 tokio::spawn(forward_output_lines(reader, output_tx));
1436 } else {
1437 // Pipe mode: stdout and stderr are merged into the same channel.
1438 // Both `ready_output` and `on_output_hook` patterns match against
1439 // lines from either stream, which is the expected behavior (a
1440 // "server ready" message may appear on stderr in some tools).
1441 if let Some(stdout) = stdout_reader {
1442 tokio::spawn(forward_output_lines(stdout, output_tx.clone()));
1443 }
1444 if let Some(stderr) = stderr_reader {
1445 tokio::spawn(forward_output_lines(stderr, output_tx.clone()));
1446 }
1447 // Drop the last sender so the channel closes when all readers
1448 // finish. The relay holds its own clone, so a sink's reports
1449 // still have somewhere to go after these end.
1450 drop(output_tx);
1451 }
1452 let log_store = Arc::clone(&LOG_STORE);
1453 let log_format = opts
1454 .log_format
1455 .clone()
1456 .unwrap_or_else(|| crate::settings::settings().logs.log_format.clone());
1457 let parse_line = move |line: &str| crate::log_parse::parse(line, &log_format);
1458
1459 const LOG_BATCH_SIZE: usize = 100;
1460 const LOG_FLUSH_INTERVAL: Duration = Duration::from_millis(100);
1461 let mut log_buffer: Vec<crate::log_parse::ParsedLog> =
1462 Vec::with_capacity(LOG_BATCH_SIZE);
1463 let mut log_flush_interval = tokio::time::interval(LOG_FLUSH_INTERVAL);
1464 log_flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1465
1466 let flush_logs =
1467 |buffer: &mut Vec<crate::log_parse::ParsedLog>| -> Option<tokio::task::JoinHandle<()>> {
1468 if buffer.is_empty() {
1469 return None;
1470 }
1471 let store = Arc::clone(&log_store);
1472 let id = id.clone();
1473 let batch = std::mem::take(buffer);
1474 Some(tokio::task::spawn_blocking(move || {
1475 if let Err(e) = store.append_structured_batch(&id, &batch) {
1476 error!("Failed to write batch to log for daemon {id}: {e}");
1477 }
1478 }))
1479 };
1480
1481 // SQLite WAL mode provides automatic durability; no explicit flush needed.
1482
1483 // Setup readiness checking
1484 let mut ready_notified = false;
1485 // Set when a oneshot's process exits 0. Its readiness *is* its
1486 // completion, so the notification is held back until the
1487 // `completed` state has been persisted — a caller that returns
1488 // from `pitchfork start` must not still see the daemon running.
1489 let mut oneshot_completion_pending = false;
1490 let mut ready_tx = ready_tx;
1491 let ready_pattern = ready_output
1492 .as_ref()
1493 .and_then(|o| get_or_compile_regex(&o.pattern));
1494 // Track whether we've already spawned the active_port detection task
1495 let mut active_port_spawned = false;
1496
1497 // Validate on_output config early; discard the hook on any error so
1498 // a bad regex does not silently fall through to the (None, None) => true
1499 // match arm and fire on every line.
1500 let on_output_hook = match on_output_hook {
1501 Some(ref hook) => match hook.validate(id.name()) {
1502 Ok(()) => on_output_hook,
1503 Err(e) => {
1504 error!("{e}");
1505 None
1506 }
1507 },
1508 None => None,
1509 };
1510
1511 // Compile the regex pattern after validation so we only attempt this
1512 // when the hook is known-good (validate() already checked the syntax).
1513 let on_output_pattern: Option<regex::Regex> = on_output_hook
1514 .as_ref()
1515 .and_then(|h| h.regex.as_deref().and_then(get_or_compile_regex));
1516 let on_output_debounce = on_output_hook
1517 .as_ref()
1518 .map(|h| h.debounce_duration())
1519 .unwrap_or(Duration::from_millis(1000));
1520 // Last time the on_output hook fired; None means it has never fired.
1521 let mut on_output_last_fired: Option<std::time::Instant> = None;
1522
1523 let mut delay_timer =
1524 ready_delay.map(|secs| Box::pin(time::sleep(Duration::from_secs(secs))));
1525
1526 // Track exhaustion of timed checks
1527 let mut http_exhausted = false;
1528 let mut cmd_exhausted = false;
1529 let mut port_exhausted = false;
1530 let mut output_exhausted = false;
1531
1532 // Get settings for intervals
1533 let s = settings();
1534 let ready_check_interval = s.supervisor_ready_check_interval();
1535 let http_client_timeout = s.supervisor_http_client_timeout();
1536
1537 // Setup output readiness check deadline
1538 let mut output_deadline = ready_output
1539 .as_ref()
1540 .and_then(|o| o.timeout)
1541 .map(|d| Box::pin(time::sleep(d)));
1542
1543 // Setup HTTP readiness check interval and deadline
1544 let mut http_check_interval = ready_http
1545 .as_ref()
1546 .map(|_| tokio::time::interval(ready_check_interval));
1547 let mut http_deadline = ready_http
1548 .as_ref()
1549 .and_then(|h| h.timeout)
1550 .map(|d| Box::pin(time::sleep(d)));
1551 let http_client = ready_http.as_ref().map(|_| {
1552 reqwest::Client::builder()
1553 .timeout(http_client_timeout)
1554 .build()
1555 .unwrap_or_default()
1556 });
1557
1558 // Setup TCP port readiness check interval and deadline
1559 let mut port_check_interval =
1560 ready_port.map(|_| tokio::time::interval(ready_check_interval));
1561 let mut port_deadline = ready_port_config
1562 .as_ref()
1563 .and_then(|p| p.timeout)
1564 .map(|d| Box::pin(time::sleep(d)));
1565
1566 // Setup command readiness check state. Probes are spawned one at a time;
1567 // a non-zero result triggers a respawn delay, and a timeout stops the probe.
1568 let mut cmd_probe: Option<CmdProbe> = None;
1569 let mut cmd_respawn_delay: Option<_> = None;
1570 let mut cmd_deadline = ready_cmd
1571 .as_ref()
1572 .and_then(|c| c.timeout)
1573 .map(|d| Box::pin(time::sleep(d)));
1574 if let Some(ref cmd) = ready_cmd {
1575 cmd_probe = Some(spawn_cmd_probe(
1576 &id,
1577 &cmd.run,
1578 daemon_dir.as_path(),
1579 hook_retry_count,
1580 readiness_daemon_env.as_ref(),
1581 &readiness_resolved_ports,
1582 ));
1583 }
1584
1585 // Use a channel to communicate process exit status
1586 let (exit_tx, mut exit_rx) =
1587 tokio::sync::mpsc::channel::<std::io::Result<std::process::ExitStatus>>(1);
1588
1589 // Spawn a task to wait for process exit
1590 let child_pid = child.id().unwrap_or(0);
1591 tokio::spawn(async move {
1592 let result = child.wait().await;
1593 // On non-Linux Unix (e.g. macOS) the zombie reaper may win the
1594 // race and consume the exit status via waitpid(None, WNOHANG)
1595 // before Tokio's child.wait() gets to it. When that happens,
1596 // Tokio returns an ECHILD io::Error. We recover by checking
1597 // REAPED_STATUSES for the stashed exit code.
1598 //
1599 // On Linux this is unnecessary because the reaper uses
1600 // waitid(WNOWAIT) to peek before reaping, which avoids the
1601 // race entirely.
1602 #[cfg(all(unix, not(target_os = "linux")))]
1603 let result = match &result {
1604 Err(e) if e.raw_os_error() == Some(nix::libc::ECHILD) => {
1605 if let Some(code) = super::REAPED_STATUSES.lock().await.remove(&child_pid) {
1606 warn!(
1607 "daemon pid {child_pid} wait() got ECHILD; \
1608 recovered exit code {code} from zombie reaper"
1609 );
1610 // Synthesize an ExitStatus from the stashed code.
1611 // On Unix we can use `ExitStatus::from_raw()` with
1612 // a wait-style status word (code << 8 for normal
1613 // exit, or raw signal number for signal death).
1614 use std::os::unix::process::ExitStatusExt;
1615 if code >= 0 {
1616 Ok(std::process::ExitStatus::from_raw(code << 8))
1617 } else {
1618 // Negative code means killed by signal (-sig)
1619 Ok(std::process::ExitStatus::from_raw((-code) & 0x7f))
1620 }
1621 } else {
1622 warn!(
1623 "daemon pid {child_pid} wait() got ECHILD but no \
1624 stashed status found; reporting as error"
1625 );
1626 result
1627 }
1628 }
1629 _ => result,
1630 };
1631 debug!("daemon pid {child_pid} wait() completed with result: {result:?}");
1632 let _ = exit_tx.send(result).await;
1633 });
1634
1635 #[allow(unused_assignments)]
1636 // Initial None is a safety net; loop only exits via exit_rx.recv() which sets it
1637 let mut exit_status = None;
1638
1639 // If there is no ready check of any kind and no delay, the daemon is
1640 // considered immediately ready and the active_port detection task would
1641 // never be triggered inside the select loop. Kick it off right away so
1642 // that daemons without any readiness configuration still get their
1643 // active_port populated (needed for proxy routing).
1644 if has_port_config
1645 && ready_pattern.is_none()
1646 && ready_http.is_none()
1647 && ready_port.is_none()
1648 && ready_cmd.is_none()
1649 && delay_timer.is_none()
1650 {
1651 active_port_spawned = true;
1652 detect_and_store_active_port(id.clone(), daemon_pid);
1653 }
1654
1655 // Set when readiness checks exhaust. The group kill runs as a
1656 // separate task so this loop can exit and the post-loop drain
1657 // keeps consuming output — children logging during SIGTERM
1658 // cleanup would otherwise block on a full pipe and never exit.
1659 // The ready failure is only sent once the kill task completes,
1660 // so the retry loop cannot respawn into the dying group.
1661 let mut ready_fail_kill: Option<tokio::task::JoinHandle<()>> = None;
1662
1663 loop {
1664 // biased: evaluate in exit → output → delay order so that
1665 // process exit pre-empts both buffered output and the delay
1666 // timer, preventing a dead daemon from being marked ready.
1667 select! {
1668 biased;
1669 Some(result) = exit_rx.recv() => {
1670 // Process exited - save exit status and notify if not ready yet
1671 exit_status = Some(result);
1672 debug!("daemon {id} process exited, exit_status: {exit_status:?}");
1673 if !ready_notified {
1674 // Check if process exited successfully
1675 let is_success = exit_status.as_ref()
1676 .and_then(|r| r.as_ref().ok())
1677 .map(|s| s.success())
1678 .unwrap_or(false);
1679 if is_success && opts.oneshot {
1680 debug!("daemon {id} completed, deferring success notification until the completed state is persisted");
1681 oneshot_completion_pending = true;
1682 } else if let Some(tx) = ready_tx.take() {
1683 if is_success {
1684 debug!("daemon {id} exited successfully before ready check, sending success notification");
1685 let _ = tx.send(Ok(()));
1686 } else {
1687 let exit_code = exit_status.as_ref()
1688 .and_then(|r| r.as_ref().ok())
1689 .and_then(|s| s.code());
1690 debug!("daemon {id} exited with failure before ready check, sending failure notification with exit_code: {exit_code:?}");
1691 let _ = tx.send(Err(exit_code));
1692 }
1693 }
1694 } else {
1695 debug!("daemon {id} was already marked ready, not sending notification");
1696 }
1697 break;
1698 },
1699 Some(super::OutputLine { text: line, source }) = output_rx.recv() => {
1700 // A line relayed by a sink is already in the store —
1701 // the sink wrote and flushed it before reporting it —
1702 // so it arrives here only to be acted on.
1703 if matches!(source, super::OutputSource::Local) {
1704 let parsed = parse_line(&line);
1705 log_buffer.push(parsed);
1706 if log_buffer.len() >= LOG_BATCH_SIZE {
1707 let _ = flush_logs(&mut log_buffer);
1708 }
1709 }
1710 trace!("output: {id} {line}");
1711
1712 // Strip ANSI for pattern matching so user-written patterns
1713 // work regardless of whether the process emits color codes.
1714 let line_clean = console::strip_ansi_codes(&line).to_string();
1715
1716 // Check if output matches ready pattern
1717 if !ready_notified
1718 && !output_exhausted
1719 && let Some(ref pattern) = ready_pattern
1720 && pattern.is_match(&line_clean)
1721 {
1722 // Flush buffered logs synchronously before signalling
1723 // readiness, so collect_startup_logs sees the line
1724 // that triggered the match (and any co-buffered lines)
1725 // in SQLite.
1726 if let Some(handle) = flush_logs(&mut log_buffer) {
1727 let _ = handle.await;
1728 }
1729 info!("daemon {id} ready: output matched pattern");
1730 ready_notified = true;
1731 if let Some(tx) = ready_tx.take() {
1732 let _ = tx.send(Ok(()));
1733 }
1734 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), vec![]).await;
1735 stop_cmd_probe_state(&mut cmd_probe);
1736 http_deadline = None;
1737 cmd_deadline = None;
1738 port_deadline = None;
1739 output_deadline = None;
1740 if !active_port_spawned && has_port_config {
1741 active_port_spawned = true;
1742 detect_and_store_active_port(id.clone(), daemon_pid);
1743 }
1744 }
1745
1746 // Check on_output hook. A sink has already applied the
1747 // filter, and says so per line: a line reported only
1748 // because it announced readiness must not fire a hook
1749 // that filters for something else.
1750 if let Some(ref hook) = on_output_hook {
1751 let matched = match source {
1752 super::OutputSource::Sink { fires_hook } => fires_hook,
1753 super::OutputSource::Local => match (&hook.filter, &on_output_pattern) {
1754 (Some(substr), _) => line_clean.contains(substr.as_str()),
1755 (None, Some(re)) => re.is_match(&line_clean),
1756 (None, None) => true,
1757 },
1758 };
1759 if matched {
1760 // The debounce is applied here as well as in the
1761 // sink. A replacement sink starts with a fresh
1762 // clock, and would otherwise let the hook fire
1763 // twice inside one configured window.
1764 let now = std::time::Instant::now();
1765 let elapsed = on_output_last_fired.map(|t| now.duration_since(t));
1766 if elapsed.is_none_or(|e| e >= on_output_debounce) {
1767 on_output_last_fired = Some(now);
1768 hooks::fire_output_hook(id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), hook.run.clone(), line_clean.clone()).await;
1769 }
1770 }
1771 }
1772 // Yield briefly so that the output readiness deadline can be
1773 // evaluated even when output is produced continuously.
1774 tokio::task::yield_now().await;
1775 }
1776 _ = async {
1777 if let Some(ref mut deadline) = http_deadline {
1778 deadline.await;
1779 } else {
1780 std::future::pending::<()>().await;
1781 }
1782 }, if !ready_notified && ready_http.is_some() => {
1783 http_exhausted = true;
1784 http_deadline = None;
1785 http_check_interval = None;
1786 warn!("daemon {id}: HTTP readiness check timed out");
1787 let any_remaining = any_ready_check_remaining(
1788 ready_output.as_ref(),
1789 output_exhausted,
1790 ready_port_config.as_ref(),
1791 port_exhausted,
1792 ready_http.as_ref(),
1793 http_exhausted,
1794 ready_cmd.as_ref(),
1795 cmd_exhausted,
1796 );
1797 if !any_remaining {
1798 error!("daemon {id}: all readiness checks exhausted, failing");
1799 stop_cmd_probe_state(&mut cmd_probe);
1800 ready_fail_kill = Some(spawn_ready_fail_kill(
1801 id.clone(),
1802 daemon_pid,
1803 opts.stop_signal.unwrap_or_default(),
1804 ));
1805 break;
1806 }
1807 }
1808 _ = async {
1809 if let Some(ref mut deadline) = output_deadline {
1810 deadline.await;
1811 } else {
1812 std::future::pending::<()>().await;
1813 }
1814 }, if !ready_notified && ready_output.is_some() => {
1815 output_exhausted = true;
1816 output_deadline = None;
1817 warn!("daemon {id}: output readiness check timed out");
1818 let any_remaining = any_ready_check_remaining(
1819 ready_output.as_ref(),
1820 output_exhausted,
1821 ready_port_config.as_ref(),
1822 port_exhausted,
1823 ready_http.as_ref(),
1824 http_exhausted,
1825 ready_cmd.as_ref(),
1826 cmd_exhausted,
1827 );
1828 if !any_remaining {
1829 error!("daemon {id}: all readiness checks exhausted, failing");
1830 stop_cmd_probe_state(&mut cmd_probe);
1831 ready_fail_kill = Some(spawn_ready_fail_kill(
1832 id.clone(),
1833 daemon_pid,
1834 opts.stop_signal.unwrap_or_default(),
1835 ));
1836 break;
1837 }
1838 }
1839 _ = async {
1840 if let Some(ref mut interval) = http_check_interval {
1841 interval.tick().await;
1842 } else {
1843 std::future::pending::<()>().await;
1844 }
1845 }, if !ready_notified && ready_http.is_some() && !http_exhausted => {
1846 if let (Some(http), Some(client)) = (&ready_http, &http_client) {
1847 match client.get(&http.url).send().await {
1848 Ok(response) if http.accepts_status(response.status().as_u16()) => {
1849 info!("daemon {id} ready: HTTP check passed (status {})", response.status());
1850 ready_notified = true;
1851 if let Some(tx) = ready_tx.take() {
1852 let _ = tx.send(Ok(()));
1853 }
1854 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), vec![]).await;
1855 http_check_interval = None;
1856 http_deadline = None;
1857 stop_cmd_probe_state(&mut cmd_probe);
1858 cmd_deadline = None;
1859 port_deadline = None;
1860 output_deadline = None;
1861 if !active_port_spawned && has_port_config {
1862 active_port_spawned = true;
1863 detect_and_store_active_port(id.clone(), daemon_pid);
1864 }
1865 }
1866 Ok(response) => {
1867 trace!("daemon {id} HTTP check: status {} (not ready)", response.status());
1868 }
1869 Err(e) => {
1870 trace!("daemon {id} HTTP check failed: {e}");
1871 }
1872 }
1873 }
1874 }
1875 _ = async {
1876 if let Some(ref mut deadline) = port_deadline {
1877 deadline.await;
1878 } else {
1879 std::future::pending::<()>().await;
1880 }
1881 }, if !ready_notified && ready_port.is_some() => {
1882 port_exhausted = true;
1883 port_deadline = None;
1884 port_check_interval = None;
1885 warn!("daemon {id}: TCP port readiness check timed out");
1886 let any_remaining = any_ready_check_remaining(
1887 ready_output.as_ref(),
1888 output_exhausted,
1889 ready_port_config.as_ref(),
1890 port_exhausted,
1891 ready_http.as_ref(),
1892 http_exhausted,
1893 ready_cmd.as_ref(),
1894 cmd_exhausted,
1895 );
1896 if !any_remaining {
1897 error!("daemon {id}: all readiness checks exhausted, failing");
1898 stop_cmd_probe_state(&mut cmd_probe);
1899 ready_fail_kill = Some(spawn_ready_fail_kill(
1900 id.clone(),
1901 daemon_pid,
1902 opts.stop_signal.unwrap_or_default(),
1903 ));
1904 break;
1905 }
1906 }
1907 _ = async {
1908 if let Some(ref mut interval) = port_check_interval {
1909 interval.tick().await;
1910 } else {
1911 std::future::pending::<()>().await;
1912 }
1913 }, if !ready_notified && ready_port.is_some() && !port_exhausted => {
1914 if let Some(port) = ready_port {
1915 match tokio::net::TcpStream::connect(("127.0.0.1", port)).await {
1916 Ok(_) => {
1917 info!("daemon {id} ready: TCP port {port} is listening");
1918 ready_notified = true;
1919 if let Some(tx) = ready_tx.take() {
1920 let _ = tx.send(Ok(()));
1921 }
1922 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), vec![]).await;
1923 // Stop checking once ready
1924 port_check_interval = None;
1925 port_deadline = None;
1926 stop_cmd_probe_state(&mut cmd_probe);
1927 http_deadline = None;
1928 cmd_deadline = None;
1929 output_deadline = None;
1930 if !active_port_spawned && has_port_config {
1931 active_port_spawned = true;
1932 // ready_port check just TCP-connected to this
1933 // port, so it is definitely listening. If it
1934 // matches the first resolved port, write
1935 // active_port directly instead of spawning
1936 // detect_and_store_active_port, which sleeps
1937 // 500 ms then relies on listeners::get_all()
1938 // + process-tree traversal — unreliable on
1939 // Windows where Git Bash PID mapping can
1940 // break descendant lookups.
1941 if let Some(active_port) = active_port_from_ready_port(
1942 port,
1943 &readiness_resolved_ports,
1944 ) {
1945 let mut state_file =
1946 SUPERVISOR.state_file.lock().await;
1947 if let Some(d) = state_file.daemons.get(&id)
1948 && d.pid == Some(daemon_pid)
1949 {
1950 state_file.set_active_port(&id, active_port);
1951 }
1952 } else {
1953 detect_and_store_active_port(
1954 id.clone(),
1955 daemon_pid,
1956 );
1957 }
1958 }
1959 }
1960 Err(_) => {
1961 trace!("daemon {id} port check: port {port} not listening yet");
1962 }
1963 }
1964 }
1965 }
1966 _ = async {
1967 if let Some(ref mut delay) = cmd_respawn_delay {
1968 delay.await;
1969 } else {
1970 std::future::pending::<()>().await;
1971 }
1972 }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted && cmd_probe.is_none() => {
1973 if let Some(ref cmd) = ready_cmd {
1974 cmd_probe = Some(spawn_cmd_probe(
1975 &id,
1976 &cmd.run,
1977 daemon_dir.as_path(),
1978 hook_retry_count,
1979 readiness_daemon_env.as_ref(),
1980 &readiness_resolved_ports,
1981 ));
1982 }
1983 cmd_respawn_delay = None;
1984 }
1985 result = async {
1986 if let Some(probe) = cmd_probe.as_mut() {
1987 std::pin::Pin::new(&mut probe.result_rx).await
1988 } else {
1989 std::future::pending::<Result<Result<std::process::ExitStatus, std::io::Error>, tokio::sync::oneshot::error::RecvError>>().await
1990 }
1991 }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted => {
1992 // The probe task has finished; remove the handle so it is not
1993 // cancelled or reused. This must happen only after this branch
1994 // actually wins the select, not while constructing the future.
1995 let _ = cmd_probe.take();
1996 match result {
1997 Ok(Ok(status)) if status.success() => {
1998 info!("daemon {id} ready: readiness command succeeded");
1999 ready_notified = true;
2000 if let Some(tx) = ready_tx.take() {
2001 let _ = tx.send(Ok(()));
2002 }
2003 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), vec![]).await;
2004 cmd_respawn_delay = None;
2005 cmd_deadline = None;
2006 http_deadline = None;
2007 port_deadline = None;
2008 output_deadline = None;
2009 if !active_port_spawned && has_port_config {
2010 active_port_spawned = true;
2011 detect_and_store_active_port(id.clone(), daemon_pid);
2012 }
2013 }
2014 Ok(Ok(_)) | Ok(Err(_)) | Err(_) => {
2015 trace!("daemon {id} cmd check: command not ready, will respawn");
2016 cmd_respawn_delay = Some(Box::pin(time::sleep(ready_check_interval)));
2017 }
2018 }
2019 }
2020 _ = async {
2021 if let Some(ref mut deadline) = cmd_deadline {
2022 deadline.await;
2023 } else {
2024 std::future::pending::<()>().await;
2025 }
2026 }, if !ready_notified && ready_cmd.is_some() => {
2027 cmd_exhausted = true;
2028 cmd_deadline = None;
2029 stop_cmd_probe_state(&mut cmd_probe);
2030 cmd_respawn_delay = None;
2031 warn!("daemon {id}: command readiness check timed out");
2032 let any_remaining = any_ready_check_remaining(
2033 ready_output.as_ref(),
2034 output_exhausted,
2035 ready_port_config.as_ref(),
2036 port_exhausted,
2037 ready_http.as_ref(),
2038 http_exhausted,
2039 ready_cmd.as_ref(),
2040 cmd_exhausted,
2041 );
2042 if !any_remaining {
2043 error!("daemon {id}: all readiness checks exhausted, failing");
2044 ready_fail_kill = Some(spawn_ready_fail_kill(
2045 id.clone(),
2046 daemon_pid,
2047 opts.stop_signal.unwrap_or_default(),
2048 ));
2049 break;
2050 }
2051 }
2052 _ = async {
2053 if let Some(ref mut timer) = delay_timer {
2054 timer.await;
2055 } else {
2056 std::future::pending::<()>().await;
2057 }
2058 } => {
2059 let has_other_ready_check = ready_pattern.is_some()
2060 || ready_http.is_some()
2061 || ready_port.is_some()
2062 || ready_cmd.is_some();
2063 let delay_is_only_readiness = !ready_notified && !has_other_ready_check;
2064 let process_exited = exit_status.is_some();
2065 let process_running = if delay_is_only_readiness && !process_exited {
2066 // Force-refresh sysinfo for this PID before checking.
2067 // On Windows, the cached process list may be stale.
2068 PROCS.refresh_pids(&[daemon_pid]);
2069 PROCS.is_running(daemon_pid)
2070 } else {
2071 false
2072 };
2073
2074 if delay_readiness_succeeded(
2075 ready_notified,
2076 has_other_ready_check,
2077 process_exited,
2078 process_running,
2079 ) {
2080 info!("daemon {id} ready: delay elapsed");
2081 ready_notified = true;
2082 if let Some(tx) = ready_tx.take() {
2083 let _ = tx.send(Ok(()));
2084 }
2085 fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook_resolved_ports.clone(), vec![]).await;
2086 if !active_port_spawned && has_port_config {
2087 active_port_spawned = true;
2088 detect_and_store_active_port(id.clone(), daemon_pid);
2089 }
2090 } else if delay_is_only_readiness {
2091 if process_exited {
2092 debug!("daemon {id} exited during ready_delay, not marking as ready");
2093 } else {
2094 debug!("daemon {id} pid {daemon_pid} not running during ready_delay, deferring to exit handler");
2095 }
2096 }
2097
2098 if delay_is_only_readiness {
2099 // Clear all deadlines — no other checks are configured
2100 // when delay fires as readiness, but clear defensively.
2101 output_deadline = None;
2102 http_deadline = None;
2103 cmd_deadline = None;
2104 port_deadline = None;
2105 stop_cmd_probe_state(&mut cmd_probe);
2106 }
2107 // Disable timer after it fires
2108 delay_timer = None;
2109 }
2110 _ = log_flush_interval.tick() => {
2111 let _ = flush_logs(&mut log_buffer);
2112 }
2113 }
2114 }
2115
2116 // Snapshot the daemon state BEFORE draining output.
2117 //
2118 // The drain can take up to 5s (e.g. when child processes keep the
2119 // stdout pipe open). During that time, a subsequent start() call
2120 // (e.g. from `pitchfork restart`) can upsert the daemon with a new
2121 // PID and Running status. If we only checked state AFTER the drain,
2122 // the monitoring task would see d.pid != Some(old_pid) && !is_stopped()
2123 // && !is_stopping() and return early without firing on_stop/on_exit
2124 // hooks.
2125 //
2126 // By snapshotting is_stopping before the drain, we preserve the
2127 // knowledge that stop() was called, so hooks fire correctly even
2128 // if start() has since changed the state.
2129 let pre_drain_daemon = SUPERVISOR.get_daemon(&id).await;
2130 let pre_drain_is_stopping = pre_drain_daemon
2131 .as_ref()
2132 .is_some_and(|d| d.status.is_stopped() || d.status.is_stopping());
2133
2134 // Drain any in-flight output lines that were still in the mpsc
2135 // channel or the OS pipe buffer when the child exited. Without
2136 // this, trailing log lines from short-lived daemons get dropped.
2137 // The reader tasks drop their senders on EOF, so recv() returns
2138 // None when all data has been consumed. A total deadline of 5 s
2139 // guards against a stuck reader (e.g. PTY master FD not closing)
2140 // while ensuring drain doesn't block post-exit cleanup indefinitely.
2141 //
2142 // Stop accepting relayed output first: the relay holds a sender of
2143 // its own, so leaving it registered would keep the channel open and
2144 // make every drain wait out the whole deadline. Readiness is moot
2145 // now anyway — the process has exited.
2146 drop(output_relay);
2147 let drain_deadline = tokio::time::Instant::now() + Duration::from_secs(5);
2148 loop {
2149 let now = tokio::time::Instant::now();
2150 if now >= drain_deadline {
2151 break;
2152 }
2153 let Ok(Some(line)) =
2154 tokio::time::timeout(drain_deadline - now, output_rx.recv()).await
2155 else {
2156 break;
2157 };
2158 // Sink-relayed lines are already stored; see the select loop.
2159 if matches!(line.source, super::OutputSource::Local) {
2160 log_buffer.push(parse_line(&line.text));
2161 }
2162 }
2163 // Flush any remaining log lines (including drained) before the process exits.
2164 // Await the flush to guarantee all buffered logs are persisted before cleanup.
2165 if let Some(handle) = flush_logs(&mut log_buffer) {
2166 let _ = handle.await;
2167 }
2168
2169 // Clear active_port since the process is no longer running
2170 {
2171 let mut state_file = SUPERVISOR.state_file.lock().await;
2172 state_file.clear_active_port(&id);
2173 }
2174
2175 // Get the final exit status
2176 let exit_status = if let Some(status) = exit_status {
2177 status
2178 } else {
2179 // Streams closed but process hasn't exited yet, wait for it
2180 match exit_rx.recv().await {
2181 Some(status) => status,
2182 None => {
2183 warn!("daemon {id} exit channel closed without receiving status");
2184 Err(std::io::Error::other("exit channel closed"))
2185 }
2186 }
2187 };
2188
2189 // If the loop exited via readiness exhaustion, wait for the group
2190 // kill to finish before reporting the failure so the retry loop
2191 // (or a waiting client) cannot start a replacement while the old
2192 // process group is still terminating.
2193 if let Some(kill) = ready_fail_kill {
2194 let _ = kill.await;
2195 if let Some(tx) = ready_tx.take() {
2196 let _ = tx.send(Err(Some(124)));
2197 }
2198 }
2199
2200 let current_daemon = SUPERVISOR.get_daemon(&id).await;
2201
2202 // Signal that this monitoring task is processing its exit path.
2203 // The RAII guard will decrement the counter and notify close()
2204 // when the task finishes (including all fire_hook registrations),
2205 // regardless of which return path is taken.
2206 SUPERVISOR
2207 .active_monitors
2208 .fetch_add(1, atomic::Ordering::Release);
2209 struct MonitorGuard;
2210 impl Drop for MonitorGuard {
2211 fn drop(&mut self) {
2212 SUPERVISOR
2213 .active_monitors
2214 .fetch_sub(1, atomic::Ordering::Release);
2215 SUPERVISOR.monitor_done.notify_waiters();
2216 }
2217 }
2218 let _monitor_guard = MonitorGuard;
2219 // Check if this monitoring task is for the current daemon process.
2220 // If the daemon was intentionally stopped (pre_drain_is_stopping),
2221 // skip this check — we must still fire on_stop/on_exit hooks even
2222 // if start() has since changed the PID and status.
2223 if !pre_drain_is_stopping
2224 && (current_daemon.is_none()
2225 || current_daemon.as_ref().is_some_and(|d| {
2226 d.pid != Some(pid) && !d.status.is_stopped() && !d.status.is_stopping()
2227 }))
2228 {
2229 // Another process has taken over, don't update status. The
2230 // task itself did finish, so a caller waiting on it is still
2231 // told so rather than left to time out.
2232 if oneshot_completion_pending && let Some(tx) = ready_tx.take() {
2233 let _ = tx.send(Ok(()));
2234 }
2235 return;
2236 }
2237 // Capture the intentional-stop flag. Combine pre-drain and
2238 // post-drain state to handle both race orders:
2239 // - stop() set Stopping before drain → pre_drain_is_stopping
2240 // - stop() set Stopped during drain → current_daemon.is_stopped()
2241 let already_stopped = current_daemon
2242 .as_ref()
2243 .is_some_and(|d| d.status.is_stopped());
2244 let is_stopping = already_stopped
2245 || pre_drain_is_stopping
2246 || current_daemon
2247 .as_ref()
2248 .is_some_and(|d| d.status.is_stopping());
2249
2250 // --- Phase 1: Determine exit_code, exit_reason, and update daemon state ---
2251 let (exit_code, exit_reason) = match (&exit_status, is_stopping) {
2252 (Ok(status), true) => {
2253 // Intentional stop (by pitchfork). status.code() returns None
2254 // on Unix when killed by signal (e.g. SIGTERM); use -1 to
2255 // distinguish from a clean exit code 0.
2256 (status.code().unwrap_or(-1), "stop")
2257 }
2258 (Ok(status), false) if status.success() => (status.code().unwrap_or(-1), "exit"),
2259 (Ok(status), false) => (status.code().unwrap_or(-1), "fail"),
2260 (Err(_), true) => {
2261 // child.wait() error while stopping (e.g. sysinfo reaped the process)
2262 (-1, "stop")
2263 }
2264 (Err(_), false) => (-1, "fail"),
2265 };
2266 // A stop that arrived while this monitor was draining did not get
2267 // to write anything, so it is applied here. A run that had already
2268 // succeeded keeps that outcome — there was nothing left to
2269 // interrupt — but a failed one is recorded as stopped, so the
2270 // retry checker leaves it alone.
2271
2272 // Update daemon state unless stop() already did it (won the race),
2273 // OR the daemon was intentionally stopped before the drain
2274 // (pre_drain_is_stopping). In the latter case, start() may have
2275 // upserted Running during the 5s drain, and we must NOT overwrite
2276 // it with Stopped — that would undo the restart.
2277 if !already_stopped && !pre_drain_is_stopping {
2278 if let Ok(status) = &exit_status {
2279 info!("daemon {id} exited with status {status}");
2280 }
2281 let (new_status, last_exit_success) = terminal_exit_state(
2282 exit_reason,
2283 opts.oneshot,
2284 exit_code,
2285 exit_status.as_ref().map(|s| s.success()).unwrap_or(true),
2286 );
2287 // Revalidate ownership inside the same state-lock section that
2288 // performs the write. The snapshot above was taken without
2289 // holding the lock, so a restart running on another thread can
2290 // install a successor in between; overwriting its record would
2291 // clear a live daemon's PID and undo the restart.
2292 if !SUPERVISOR
2293 .finalize_monitored_exit(
2294 &id,
2295 pid,
2296 monitor_token,
2297 new_status,
2298 Some(last_exit_success),
2299 )
2300 .await
2301 {
2302 debug!("daemon {id} exit state was not written; a successor owns the record");
2303 }
2304 }
2305
2306 // The terminal state is now visible, so a caller waiting on this
2307 // oneshot can return and see it. A task that was stopped partway
2308 // never did its work, so it does not satisfy anything waiting on
2309 // it — even when the process caught the signal and exited 0.
2310 if oneshot_completion_pending && let Some(tx) = ready_tx.take() {
2311 if exit_reason == "exit" {
2312 let _ = tx.send(Ok(()));
2313 } else {
2314 warn!("daemon {id}: oneshot was stopped before completing");
2315 let _ = tx.send(Err(None));
2316 }
2317 }
2318
2319 // --- Phase 2: Fire hooks ---
2320 let hook_extra_env = vec![
2321 ("PITCHFORK_EXIT_CODE".to_string(), exit_code.to_string()),
2322 ("PITCHFORK_EXIT_REASON".to_string(), exit_reason.to_string()),
2323 ];
2324
2325 // Determine which hooks to fire based on exit reason
2326 let hooks_to_fire: Vec<HookType> = match exit_reason {
2327 "stop" => vec![HookType::OnStop, HookType::OnExit],
2328 "exit" => vec![HookType::OnExit],
2329 // "fail": fire on_fail + on_exit only when retries are exhausted
2330 _ if hook_retry_count >= hook_retry.count() => {
2331 vec![HookType::OnFail, HookType::OnExit]
2332 }
2333 _ => vec![],
2334 };
2335
2336 for hook_type in hooks_to_fire {
2337 fire_hook(
2338 hook_type,
2339 id.clone(),
2340 daemon_dir.clone(),
2341 hook_retry_count,
2342 hook_daemon_env.clone(),
2343 hook_resolved_ports.clone(),
2344 hook_extra_env.clone(),
2345 )
2346 .await;
2347 }
2348 });
2349
2350 // If wait_ready is true, wait for readiness notification
2351 if let Some(ready_rx) = ready_rx {
2352 match ready_rx.await {
2353 Ok(Ok(())) => {
2354 info!("daemon {id} is ready");
2355 // Re-read rather than returning the snapshot taken at
2356 // spawn: a completed oneshot has since been finalized, and
2357 // the snapshot would tell the caller it is still running
2358 // under a PID that has exited.
2359 //
2360 // Only when the record still describes this run, though. A
2361 // successor that claimed it carries its own PID and start
2362 // time, and reporting those as the outcome of the process
2363 // this call spawned would misattribute them.
2364 let daemon = match self.get_daemon(id).await {
2365 Some(current) if current.pid.is_none() || current.pid == Some(pid) => {
2366 current
2367 }
2368 // A successor owns the record, so neither it nor the
2369 // spawn snapshot describes this run: one carries
2370 // another process's identity, the other still says
2371 // running under a PID that has exited. A oneshot that
2372 // reported ready did finish, so report that outcome
2373 // directly rather than either misleading record.
2374 _ if opts.oneshot => crate::daemon::Daemon {
2375 status: DaemonStatus::Completed,
2376 pid: None,
2377 start_time: None,
2378 boot_time: None,
2379 last_exit_success: Some(true),
2380 ..daemon
2381 },
2382 _ => daemon,
2383 };
2384 Ok(IpcResponse::DaemonReady { daemon })
2385 }
2386 Ok(Err(exit_code)) => {
2387 error!("daemon {id} failed before becoming ready");
2388 // The caller reports this by querying the log store for
2389 // what the daemon printed, so wait for the sink's final
2390 // write first. The in-process path got this ordering by
2391 // flushing synchronously before signalling.
2392 //
2393 // Only on the attempt that gives up: `run` retries inline,
2394 // and waiting after every attempt would both delay the
2395 // backoff and widen the window in which the daemon looks
2396 // errored and idle — long enough for the background retry
2397 // checker to start an attempt of its own alongside it.
2398 let last_attempt = opts.retry_count >= opts.retry.count();
2399 if using_sink && last_attempt {
2400 super::log_sink::wait_for_output(id, spawn_time, SINK_OUTPUT_TIMEOUT).await;
2401 }
2402 Ok(IpcResponse::DaemonFailedWithCode {
2403 exit_code,
2404 resolved_ports: attempt_resolved_ports,
2405 })
2406 }
2407 Err(_) => {
2408 error!("readiness channel closed unexpectedly for daemon {id}");
2409 Ok(IpcResponse::DaemonStart { daemon })
2410 }
2411 }
2412 } else {
2413 Ok(IpcResponse::DaemonStart { daemon })
2414 }
2415 }
2416
2417 /// Stop a running daemon
2418 pub async fn stop(&self, id: &DaemonId) -> Result<IpcResponse> {
2419 // Hold the daemon's stop lock for the whole stop (including the
2420 // whole-group termination wait) so starts and concurrent stops of the
2421 // same daemon serialize against it instead of racing the Stopping window.
2422 let lock = self.stop_lock(id).await;
2423 let _guard = lock.lock().await;
2424 self.stop_locked(id).await
2425 }
2426
2427 /// Stop implementation. Caller must hold the daemon's stop lock.
2428 pub(super) async fn stop_locked(&self, id: &DaemonId) -> Result<IpcResponse> {
2429 let pitchfork_id = DaemonId::pitchfork();
2430 if *id == pitchfork_id {
2431 return Ok(IpcResponse::Error(
2432 "Cannot stop supervisor via stop command".into(),
2433 ));
2434 }
2435 info!("stopping daemon: {id}");
2436 // A foreground `start` may be working through this daemon's retries,
2437 // sleeping out a backoff with no process of its own to kill. Tell it to
2438 // give up, or it would start the next attempt once the stop has
2439 // returned.
2440 self.cancel_retrying(id);
2441 // ...and the retry checker may already have decided on an attempt it
2442 // has not started yet. The count is raised when this stop is done
2443 // rather than now, and while its lock is still held, so a checker that
2444 // reads the count while the stop is still recording itself reads the
2445 // old value and stands down when it reaches the lock. Raising it up
2446 // front would hand that reader a value that still matches once the
2447 // stop has finished.
2448 let _stop_epoch_bump = StopEpochGuard(id.clone());
2449 if let Some(daemon) = self.get_daemon(id).await {
2450 trace!("daemon to stop: {daemon}");
2451 if let Some(pid) = daemon.pid {
2452 trace!("killing pid: {pid}");
2453 if PROCS.is_running(pid) {
2454 // Something is alive on that PID, but the kill below signals
2455 // the entire process group: if the PID was recycled while
2456 // this record sat unsupervised, that group belongs to an
2457 // unrelated process tree. The daemon itself is gone either
2458 // way, so report it as not running and clear the record.
2459 if !super::signalling_pid_is_authorized(
2460 daemon.start_time,
2461 PROCS.start_time(pid),
2462 ) {
2463 warn!(
2464 "pid {pid} recorded for daemon {id} belongs to another process now; not signalling it"
2465 );
2466 self.upsert_daemon(
2467 UpsertDaemonOpts::builder(id.clone())
2468 .set(|o| {
2469 o.pid = None;
2470 o.status = DaemonStatus::Stopped;
2471 })
2472 .build(),
2473 )
2474 .await?;
2475 return Ok(IpcResponse::DaemonWasNotRunning);
2476 }
2477
2478 // First set status to Stopping (preserve PID for monitoring task)
2479 self.upsert_daemon(
2480 UpsertDaemonOpts::builder(id.clone())
2481 .set(|o| {
2482 o.pid = Some(pid);
2483 o.status = DaemonStatus::Stopping;
2484 })
2485 .build(),
2486 )
2487 .await?;
2488
2489 // Kill the entire process group atomically (daemon PID == PGID
2490 // because we called setsid() at spawn time)
2491 let stop_cfg = daemon.stop_signal.unwrap_or_default();
2492 let stop_signal: i32 = stop_cfg.signal.into();
2493 if let Err(e) = PROCS
2494 .kill_process_group_async(pid, stop_signal, stop_cfg.timeout)
2495 .await
2496 {
2497 debug!("failed to kill pid {pid}: {e}");
2498 // Check if the process group is actually gone despite the
2499 // error. Checking only the leader here would mark the daemon
2500 // Stopped while surviving group members (e.g. one stuck in
2501 // uninterruptible sleep) are still alive — letting a restart
2502 // collide with them.
2503 if PROCS.process_group_alive(pid) {
2504 // Group still has live members - set back to Running
2505 debug!(
2506 "failed to stop pid {pid}: process group still alive after kill"
2507 );
2508 self.upsert_daemon(
2509 UpsertDaemonOpts::builder(id.clone())
2510 .set(|o| {
2511 o.pid = Some(pid); // Preserve PID to avoid orphaning the process
2512 o.status = DaemonStatus::Running;
2513 })
2514 .build(),
2515 )
2516 .await?;
2517 return Ok(IpcResponse::DaemonStopFailed {
2518 error: format!(
2519 "process group of {pid} still alive after kill attempt: {e}"
2520 ),
2521 });
2522 }
2523 }
2524
2525 // Process successfully stopped
2526 // Note: kill_process_group_async waits for the ENTIRE process
2527 // group to exit (stop signal -> stop_timeout -> SIGKILL, then a
2528 // bounded verification), so a replacement daemon can be started
2529 // without colliding with a still-terminating instance. The only
2530 // exception is a member stuck in uninterruptible sleep, which is
2531 // logged with a warning.
2532 self.upsert_daemon(
2533 UpsertDaemonOpts::builder(id.clone())
2534 .set(|o| {
2535 o.pid = None;
2536 o.status = DaemonStatus::Stopped;
2537 o.last_exit_success = Some(true);
2538 })
2539 .build(),
2540 )
2541 .await?;
2542 } else if daemon.oneshot && self.is_monitored(id, pid) {
2543 // The task's process is gone but its monitor is still
2544 // running, so the run's real outcome has not been written
2545 // yet — and for a task that finished on its own that
2546 // outcome is `completed`. Writing `stopped` straight over
2547 // it would discard a success the task actually achieved
2548 // and report failure to anyone waiting on it, purely
2549 // because a stop arrived a moment late.
2550 //
2551 // So wait for whoever is monitoring this run — the native
2552 // monitor or an adopted one — to finish, then decide from
2553 // what it wrote. Waiting rather than leaving a note for
2554 // the monitor to find means there is no window in which
2555 // the note lands too late to be read, and nothing left
2556 // behind if it is never read at all. The wait is bounded,
2557 // as is the monitor's own five-second output drain.
2558 //
2559 // Only oneshots take this path. A service has no
2560 // successful exit to preserve, so it falls through to the
2561 // arm below, which records the stop immediately.
2562 debug!(
2563 "pid {pid} not running but daemon {id} is still monitored; waiting for its monitor to settle the outcome"
2564 );
2565 self.wait_for_exit_finalized(id, Some(pid)).await;
2566 let finished = self.get_daemon(id).await;
2567 if finished
2568 .as_ref()
2569 .is_some_and(|d| stop_keeps_finalized_status(&d.status))
2570 {
2571 return Ok(IpcResponse::DaemonWasNotRunning);
2572 }
2573 // The run did not finish its work, so record the stop. A
2574 // failure left in place would be picked up by the retry
2575 // checker, which would start a task the user just stopped.
2576 self.upsert_daemon(
2577 UpsertDaemonOpts::builder(id.clone())
2578 .set(|o| {
2579 o.pid = None;
2580 o.status = DaemonStatus::Stopped;
2581 })
2582 .build(),
2583 )
2584 .await?;
2585 return Ok(IpcResponse::DaemonWasNotRunning);
2586 } else {
2587 debug!("pid {pid} not running, process may have exited unexpectedly");
2588 // Process already dead and unmonitored, so nothing else
2589 // will record an outcome — transition to Stopped so the
2590 // retry checker sees a terminal state and stops
2591 // scheduling new attempts. This is important for an
2592 // explicit `pitchfork stop` on an Errored daemon: the
2593 // user wants to abort retries.
2594 self.upsert_daemon(
2595 UpsertDaemonOpts::builder(id.clone())
2596 .set(|o| {
2597 o.pid = None;
2598 o.status = DaemonStatus::Stopped;
2599 })
2600 .build(),
2601 )
2602 .await?;
2603 return Ok(IpcResponse::DaemonWasNotRunning);
2604 }
2605 Ok(IpcResponse::Ok)
2606 } else {
2607 debug!("daemon {id} not running");
2608 // No process to signal, but a failed record with retries left
2609 // is not inert: `check_retry` starts the next attempt from it,
2610 // whether or not a foreground start is also working through
2611 // them. Record the stop so nothing picks the daemon back up.
2612 if daemon.status.is_errored() && daemon.retry_count < daemon.retry.count() {
2613 self.upsert_daemon(
2614 UpsertDaemonOpts::builder(id.clone())
2615 .set(|o| {
2616 o.pid = None;
2617 o.status = DaemonStatus::Stopped;
2618 })
2619 .build(),
2620 )
2621 .await?;
2622 return Ok(IpcResponse::DaemonWasNotRunning);
2623 }
2624 Ok(IpcResponse::DaemonNotRunning)
2625 }
2626 } else {
2627 debug!("daemon {id} not found");
2628 Ok(IpcResponse::DaemonNotFound)
2629 }
2630 }
2631}
2632
2633#[cfg(unix)]
2634fn resolve_effective_run_identity(daemon_user: Option<&str>) -> Result<RunIdentity> {
2635 let s = settings();
2636 let settings_user = s.supervisor.user.trim();
2637 let daemon_user = daemon_user.map(str::trim).filter(|user| !user.is_empty());
2638 let settings_user = (!settings_user.is_empty()).then_some(settings_user);
2639 let configured = daemon_user.or(settings_user);
2640 let current_uid = nix::unistd::Uid::effective().as_raw();
2641 let current_gid = nix::unistd::Gid::effective().as_raw();
2642 // The recorded invoking user of a boot service stands in for the sudo
2643 // environment that launchd and systemd do not provide.
2644 let invoking = env::invoking_user_ids().map(|(uid, gid)| (uid.to_string(), gid.to_string()));
2645 resolve_run_identity(
2646 configured,
2647 current_uid,
2648 current_gid,
2649 invoking.as_ref().map(|(uid, _)| uid.as_str()),
2650 invoking.as_ref().map(|(_, gid)| gid.as_str()),
2651 )
2652}
2653
2654#[cfg(unix)]
2655fn resolve_run_identity(
2656 configured: Option<&str>,
2657 current_uid: u32,
2658 current_gid: u32,
2659 sudo_uid: Option<&str>,
2660 sudo_gid: Option<&str>,
2661) -> Result<RunIdentity> {
2662 let current_uid = nix::unistd::Uid::from_raw(current_uid);
2663 let current_gid = nix::unistd::Gid::from_raw(current_gid);
2664 if let Some(user) = configured {
2665 let identity = resolve_configured_user(user)?;
2666 ensure_can_use_identity(user, &identity, current_uid, current_gid)?;
2667 if identity.matches(current_uid, current_gid) {
2668 return Ok(RunIdentity::Inherit);
2669 }
2670 return Ok(identity);
2671 }
2672
2673 if current_uid.is_root()
2674 && let Some(identity) = resolve_sudo_identity(sudo_uid, sudo_gid)
2675 {
2676 if identity.matches(current_uid, current_gid) {
2677 return Ok(RunIdentity::Inherit);
2678 }
2679 return Ok(identity);
2680 }
2681
2682 Ok(RunIdentity::Inherit)
2683}
2684
2685#[cfg(unix)]
2686fn resolve_configured_user(user: &str) -> Result<RunIdentity> {
2687 if user.chars().all(|c| c.is_ascii_digit()) {
2688 let uid = user
2689 .parse::<u32>()
2690 .map_err(|e| miette::miette!("invalid run user UID '{}': {}", user, e))?;
2691 let user_record = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
2692 .into_diagnostic()?
2693 .ok_or_else(|| miette::miette!("run user UID '{}' does not exist", user))?;
2694 return run_identity_from_user_record(user_record);
2695 }
2696
2697 let user_record = nix::unistd::User::from_name(user)
2698 .into_diagnostic()?
2699 .ok_or_else(|| miette::miette!("run user '{}' does not exist", user))?;
2700 run_identity_from_user_record(user_record)
2701}
2702
2703#[cfg(unix)]
2704fn run_identity_from_user_record(user: nix::unistd::User) -> Result<RunIdentity> {
2705 let username = CString::new(user.name)
2706 .map_err(|e| miette::miette!("run user name contains an interior nul byte: {}", e))?;
2707 Ok(RunIdentity::Switch {
2708 uid: user.uid,
2709 gid: user.gid,
2710 username: Some(username),
2711 home: Some(user.dir),
2712 })
2713}
2714
2715#[cfg(unix)]
2716fn resolve_sudo_identity(sudo_uid: Option<&str>, sudo_gid: Option<&str>) -> Option<RunIdentity> {
2717 let uid = sudo_uid?.parse::<u32>().ok()?;
2718 let gid = sudo_gid?.parse::<u32>().ok()?;
2719 let user = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
2720 .ok()
2721 .flatten();
2722 let (username, home) = match user {
2723 Some(user) => (CString::new(user.name).ok(), Some(user.dir)),
2724 None => (None, None),
2725 };
2726 Some(RunIdentity::Switch {
2727 uid: nix::unistd::Uid::from_raw(uid),
2728 gid: nix::unistd::Gid::from_raw(gid),
2729 username,
2730 home,
2731 })
2732}
2733
2734#[cfg(unix)]
2735fn ensure_can_use_identity(
2736 configured_user: &str,
2737 identity: &RunIdentity,
2738 current_uid: nix::unistd::Uid,
2739 current_gid: nix::unistd::Gid,
2740) -> Result<()> {
2741 let RunIdentity::Switch { uid, gid, .. } = identity else {
2742 return Ok(());
2743 };
2744 if *uid == current_uid && *gid == current_gid {
2745 return Ok(());
2746 }
2747 if current_uid.is_root() {
2748 return Ok(());
2749 }
2750 Err(miette::miette!(
2751 "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.",
2752 configured_user,
2753 current_uid.as_raw(),
2754 current_gid.as_raw(),
2755 uid.as_raw(),
2756 gid.as_raw()
2757 ))
2758}
2759
2760/// Point HOME, USER and LOGNAME at the user a daemon switches to.
2761///
2762/// The supervisor's environment describes the supervisor's own user (usually
2763/// root), so without this a daemon run as another user would read and write
2764/// root's home. A value the passwd entry cannot supply is removed rather than
2765/// left describing root.
2766#[cfg(unix)]
2767fn apply_identity_env(command: &mut tokio::process::Command, identity: &RunIdentity) {
2768 let RunIdentity::Switch { username, home, .. } = identity else {
2769 return;
2770 };
2771 match home.as_ref().filter(|home| !home.as_os_str().is_empty()) {
2772 Some(home) => command.env("HOME", home),
2773 None => command.env_remove("HOME"),
2774 };
2775 match username.as_ref().and_then(|name| name.to_str().ok()) {
2776 Some(name) => command.env("USER", name).env("LOGNAME", name),
2777 None => command.env_remove("USER").env_remove("LOGNAME"),
2778 };
2779}
2780
2781#[cfg(unix)]
2782fn apply_run_identity(identity: &RunIdentity) -> std::io::Result<()> {
2783 let RunIdentity::Switch {
2784 uid, gid, username, ..
2785 } = identity
2786 else {
2787 return Ok(());
2788 };
2789 if let Some(username) = username {
2790 initgroups_for_user(username, *gid)?;
2791 } else {
2792 setgroups_to_primary(*gid)?;
2793 }
2794 nix::unistd::setgid(*gid).map_err(nix_to_io_error)?;
2795 nix::unistd::setuid(*uid).map_err(nix_to_io_error)?;
2796 Ok(())
2797}
2798
2799#[cfg(unix)]
2800impl RunIdentity {
2801 fn matches(&self, uid: nix::unistd::Uid, gid: nix::unistd::Gid) -> bool {
2802 matches!(self, RunIdentity::Switch { uid: u, gid: g, .. } if *u == uid && *g == gid)
2803 }
2804}
2805
2806#[cfg(unix)]
2807fn setgroups_to_primary(gid: nix::unistd::Gid) -> std::io::Result<()> {
2808 let groups = [gid.as_raw() as libc::gid_t];
2809 #[cfg(any(target_os = "linux", target_os = "android"))]
2810 let group_count = groups.len();
2811 #[cfg(not(any(target_os = "linux", target_os = "android")))]
2812 let group_count = groups.len() as libc::c_int;
2813 let rc = unsafe { libc::setgroups(group_count, groups.as_ptr()) };
2814 if rc == -1 {
2815 Err(std::io::Error::last_os_error())
2816 } else {
2817 Ok(())
2818 }
2819}
2820
2821#[cfg(unix)]
2822fn initgroups_for_user(username: &CString, gid: nix::unistd::Gid) -> std::io::Result<()> {
2823 let gid = gid.as_raw();
2824 #[cfg(any(
2825 target_os = "macos",
2826 target_os = "ios",
2827 target_os = "tvos",
2828 target_os = "watchos"
2829 ))]
2830 let base_gid = i32::try_from(gid)
2831 .map_err(|_| std::io::Error::other(format!("gid {gid} is out of range")))?;
2832
2833 #[cfg(not(any(
2834 target_os = "macos",
2835 target_os = "ios",
2836 target_os = "tvos",
2837 target_os = "watchos"
2838 )))]
2839 let base_gid = gid as libc::gid_t;
2840
2841 // SAFETY: `username` is a valid nul-terminated C string and `base_gid`
2842 // is derived from a resolved system account or sudo-provided gid.
2843 let rc = unsafe { libc::initgroups(username.as_ptr(), base_gid) };
2844 if rc == -1 {
2845 Err(std::io::Error::last_os_error())
2846 } else {
2847 Ok(())
2848 }
2849}
2850
2851#[cfg(unix)]
2852fn nix_to_io_error(err: nix::errno::Errno) -> std::io::Error {
2853 std::io::Error::from_raw_os_error(err as i32)
2854}
2855
2856/// Check if multiple ports are available and optionally auto-bump to find available ports.
2857///
2858/// All ports are bumped by the same offset to maintain relative port spacing.
2859/// Returns the resolved ports (either the original or bumped ones).
2860/// Returns an error if any port is in use and auto_bump is disabled,
2861/// or if no available ports can be found after max attempts.
2862async fn check_ports_available(
2863 expected_ports: &[u16],
2864 auto_bump: bool,
2865 max_attempts: u32,
2866) -> Result<Vec<u16>> {
2867 if expected_ports.is_empty() {
2868 return Ok(Vec::new());
2869 }
2870
2871 for bump_offset in 0..=max_attempts {
2872 // Use wrapping_add to handle overflow correctly - ports wrap around at 65535
2873 let candidate_ports: Vec<u16> = expected_ports
2874 .iter()
2875 .map(|&p| p.wrapping_add(bump_offset as u16))
2876 .collect();
2877
2878 // Check if all ports in this set are available
2879 let mut all_available = true;
2880 let mut conflicting_port = None;
2881
2882 for &port in &candidate_ports {
2883 // Port 0 is a special case - it requests an ephemeral port from the OS.
2884 // Skip the availability check for port 0 since binding to it always succeeds.
2885 if port == 0 {
2886 continue;
2887 }
2888
2889 // Use spawn_blocking to avoid blocking the async runtime during TCP bind checks.
2890 //
2891 // We check multiple addresses to avoid false-negatives caused by SO_REUSEADDR.
2892 // On macOS/BSD, Rust's TcpListener::bind sets SO_REUSEADDR by default, which
2893 // allows binding 0.0.0.0:port even when 127.0.0.1:port is already in use
2894 // (because 0.0.0.0 is technically a different address). Most daemons bind
2895 // to localhost, so checking 127.0.0.1 is essential to detect real conflicts.
2896 // We also check [::1] to cover IPv6 loopback listeners.
2897 //
2898 // NOTE: This check has a time-of-check-to-time-of-use (TOCTOU) race condition.
2899 // Another process could grab the port between our check and the daemon actually
2900 // binding. This is inherent to the approach and acceptable for our use case
2901 // since we're primarily detecting conflicts with already-running daemons.
2902 if is_port_in_use(port).await {
2903 all_available = false;
2904 conflicting_port = Some(port);
2905 break;
2906 }
2907 }
2908
2909 if all_available {
2910 // Check for overflow (port wrapped around to 0 due to wrapping_add)
2911 // If any candidate port is 0 but the original expected port wasn't 0,
2912 // it means we've wrapped around and should stop
2913 if candidate_ports.contains(&0) && !expected_ports.contains(&0) {
2914 return Err(PortError::NoAvailablePort {
2915 start_port: expected_ports[0],
2916 attempts: bump_offset + 1,
2917 }
2918 .into());
2919 }
2920 if bump_offset > 0 {
2921 info!("ports {expected_ports:?} bumped by {bump_offset} to {candidate_ports:?}");
2922 }
2923 return Ok(candidate_ports);
2924 }
2925
2926 // Port is in use
2927 if bump_offset == 0
2928 && !auto_bump
2929 && let Some(port) = conflicting_port
2930 {
2931 let (pid, process) = identify_port_owner(port).await;
2932 return Err(PortError::InUse { port, process, pid }.into());
2933 }
2934 }
2935
2936 // No available ports found after max attempts
2937 Err(PortError::NoAvailablePort {
2938 start_port: expected_ports[0],
2939 attempts: max_attempts + 1,
2940 }
2941 .into())
2942}
2943
2944/// Check whether a port is currently in use by attempting to bind on multiple addresses.
2945///
2946/// Returns `true` when at least one bind attempt gets `AddrInUse`, meaning another
2947/// process is listening. Other errors (e.g. `AddrNotAvailable` on an address family
2948/// the OS doesn't support) are ignored so they don't produce false positives.
2949async fn is_port_in_use(port: u16) -> bool {
2950 tokio::task::spawn_blocking(move || {
2951 for &addr in &["0.0.0.0", "127.0.0.1", "::1"] {
2952 match std::net::TcpListener::bind((addr, port)) {
2953 Ok(listener) => drop(listener),
2954 Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => return true,
2955 Err(_) => continue,
2956 }
2957 }
2958 false
2959 })
2960 .await
2961 .unwrap_or(false)
2962}
2963
2964/// Best-effort lookup of the process occupying a port via `listeners::get_all()`.
2965///
2966/// Returns `(pid, process_name)`. Falls back to `(0, "unknown")` when the
2967/// system call fails (permission error, unsupported OS, etc.).
2968async fn identify_port_owner(port: u16) -> (u32, String) {
2969 tokio::task::spawn_blocking(move || {
2970 listeners::get_all()
2971 .ok()
2972 .and_then(|list| {
2973 list.into_iter()
2974 .find(|l| l.socket.port() == port)
2975 .map(|l| (l.process.pid, l.process.name))
2976 })
2977 .unwrap_or((0, "unknown".to_string()))
2978 })
2979 .await
2980 .unwrap_or((0, "unknown".to_string()))
2981}
2982
2983/// Detect whether a port is in use, and if so, identify the owning process.
2984///
2985/// Combines `is_port_in_use` (reliable bind probe) with `identify_port_owner`
2986/// (best-effort process lookup). Returns `None` when the port is free.
2987async fn detect_port_conflict(port: u16) -> Option<(u32, String)> {
2988 if !is_port_in_use(port).await {
2989 return None;
2990 }
2991 Some(identify_port_owner(port).await)
2992}
2993
2994#[derive(Debug, PartialEq, Eq)]
2995enum ActivePortSelection {
2996 NoCandidates,
2997 Selected(u16),
2998 Ambiguous(Vec<u16>),
2999}
3000
3001fn discovery_preferred_port(daemon: &crate::daemon::Daemon) -> Option<u16> {
3002 daemon
3003 .resolved_port
3004 .first()
3005 .copied()
3006 .or_else(|| {
3007 daemon
3008 .port
3009 .as_ref()
3010 .and_then(|port| port.expect.first().copied())
3011 })
3012 .filter(|&port| port > 0)
3013}
3014
3015fn select_active_port(
3016 listeners: impl IntoIterator<Item = listeners::Listener>,
3017 descendant_pids: &std::collections::HashSet<u32>,
3018 preferred_port: Option<u16>,
3019) -> ActivePortSelection {
3020 let process_ports: std::collections::BTreeSet<u16> = listeners
3021 .into_iter()
3022 .filter(|listener| {
3023 listener.protocol == listeners::Protocol::TCP
3024 && listener.state == listeners::SocketState::Listen
3025 && descendant_pids.contains(&listener.process.pid)
3026 })
3027 .map(|listener| listener.socket.port())
3028 .filter(|&port| port > 0)
3029 .collect();
3030
3031 if let Some(port) = preferred_port
3032 && process_ports.contains(&port)
3033 {
3034 return ActivePortSelection::Selected(port);
3035 }
3036
3037 match process_ports.len() {
3038 0 => ActivePortSelection::NoCandidates,
3039 1 => ActivePortSelection::Selected(*process_ports.first().unwrap()),
3040 _ => ActivePortSelection::Ambiguous(process_ports.into_iter().collect()),
3041 }
3042}
3043
3044/// Spawn a background task that detects the daemon process's active listening port
3045/// and stores it in the state file as `active_port`.
3046///
3047/// This is called once when the daemon becomes ready. The port is cleared when the daemon stops.
3048///
3049/// Port selection strategy:
3050/// 1. Consider only TCP listening sockets owned by the daemon or its descendants.
3051/// 2. Prefer the first resolved port, falling back to the first expected port for
3052/// legacy state without resolved ports.
3053/// 3. Select a sole distinct candidate; leave `active_port` unset when ambiguous.
3054fn detect_and_store_active_port(id: DaemonId, pid: u32) {
3055 tokio::spawn(async move {
3056 // Retry with exponential backoff so that slow-starting daemons (JVM,
3057 // Node.js, Python, etc.) that take more than 500 ms to bind their port
3058 // are still detected. Total wait budget: 500+1000+2000+4000 = 7.5 s.
3059 for delay_ms in [500u64, 1000, 2000, 4000] {
3060 tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
3061
3062 // Read daemon state atomically: check if still alive and get the preferred port
3063 // in a single lock acquisition to avoid TOCTOU and unnecessary lock overhead.
3064 let preferred_port: Option<u16> = {
3065 let state_file = SUPERVISOR.state_file.lock().await;
3066 match state_file.daemons.get(&id) {
3067 Some(d) if d.pid.is_none() => {
3068 debug!("daemon {id}: aborting active_port detection — process exited");
3069 return;
3070 }
3071 Some(d) => discovery_preferred_port(d),
3072 None => None,
3073 }
3074 };
3075
3076 let selection = tokio::task::spawn_blocking(move || {
3077 let listeners = listeners::get_all().ok()?;
3078
3079 // Refresh process tree so all_children sees current descendants.
3080 PROCS.refresh_processes();
3081
3082 let descendant_pids: std::collections::HashSet<u32> = PROCS
3083 .all_children(pid)
3084 .into_iter()
3085 .chain(std::iter::once(pid))
3086 .collect();
3087
3088 Some(select_active_port(
3089 listeners,
3090 &descendant_pids,
3091 preferred_port,
3092 ))
3093 })
3094 .await
3095 .ok()
3096 .flatten()
3097 .unwrap_or(ActivePortSelection::NoCandidates);
3098
3099 let port = match selection {
3100 ActivePortSelection::Selected(port) => port,
3101 ActivePortSelection::Ambiguous(ports) => {
3102 debug!(
3103 "daemon {id}: ambiguous active_port candidates {ports:?} for pid {pid} \
3104 and its descendants; leaving active_port unset (will retry)"
3105 );
3106 continue;
3107 }
3108 ActivePortSelection::NoCandidates => {
3109 debug!(
3110 "daemon {id}: no active port detected for pid {pid} or its descendants \
3111 (will retry)"
3112 );
3113 continue;
3114 }
3115 };
3116
3117 debug!("daemon {id} active_port detected: {port}");
3118 let mut state_file = SUPERVISOR.state_file.lock().await;
3119 if let Some(d) = state_file.daemons.get(&id) {
3120 // Guard against PID reuse: if the original process exited and the OS
3121 // assigned the same PID to an unrelated process that happens to bind
3122 // a port, we must not route proxy traffic to that unrelated service.
3123 if d.pid == Some(pid) {
3124 state_file.set_active_port(&id, port);
3125 } else {
3126 debug!(
3127 "daemon {id}: skipping active_port write — PID mismatch \
3128 (expected {pid}, current {:?})",
3129 d.pid
3130 );
3131 }
3132 }
3133 return;
3134 }
3135
3136 debug!(
3137 "daemon {id}: active port detection exhausted all retries for pid {pid} and its descendants"
3138 );
3139 });
3140}
3141
3142#[cfg(test)]
3143mod active_port_tests {
3144 use super::*;
3145 use crate::config_types::PortConfig;
3146 use listeners::{Listener, Process, Protocol, SocketState};
3147 use std::net::{IpAddr, Ipv4Addr, SocketAddr};
3148
3149 fn listener(pid: u32, port: u16, protocol: Protocol, state: SocketState) -> Listener {
3150 Listener {
3151 process: Process {
3152 pid,
3153 name: "test".to_string(),
3154 path: "/test".to_string(),
3155 },
3156 socket: SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), port),
3157 protocol,
3158 state,
3159 }
3160 }
3161
3162 #[test]
3163 fn active_port_candidates_exclude_outbound_tcp_and_udp_sockets() {
3164 let daemon_pid = 100;
3165 let child_pid = 101;
3166 let descendant_pids = [daemon_pid, child_pid].into_iter().collect();
3167 let listeners = vec![
3168 listener(child_pid, 3004, Protocol::TCP, SocketState::Listen),
3169 listener(daemon_pid, 47082, Protocol::TCP, SocketState::Established),
3170 listener(daemon_pid, 5353, Protocol::UDP, SocketState::Unknown),
3171 listener(999, 9000, Protocol::TCP, SocketState::Listen),
3172 ];
3173
3174 assert_eq!(
3175 select_active_port(listeners, &descendant_pids, None),
3176 ActivePortSelection::Selected(3004)
3177 );
3178 }
3179
3180 #[test]
3181 fn bumped_cmd_readiness_prefers_resolved_primary_port() {
3182 let daemon = crate::daemon::Daemon {
3183 resolved_port: vec![3004],
3184 port: Some(PortConfig {
3185 expect: vec![3000],
3186 ..PortConfig::default()
3187 }),
3188 ..crate::daemon::Daemon::default()
3189 };
3190 let descendant_pids = [100].into_iter().collect();
3191 let listeners = vec![
3192 listener(100, 9000, Protocol::TCP, SocketState::Listen),
3193 listener(100, 3004, Protocol::TCP, SocketState::Listen),
3194 ];
3195
3196 assert_eq!(discovery_preferred_port(&daemon), Some(3004));
3197 assert_eq!(
3198 select_active_port(
3199 listeners,
3200 &descendant_pids,
3201 discovery_preferred_port(&daemon),
3202 ),
3203 ActivePortSelection::Selected(3004)
3204 );
3205 }
3206
3207 #[test]
3208 fn legacy_state_uses_expected_primary_port_for_discovery() {
3209 let daemon = crate::daemon::Daemon {
3210 port: Some(PortConfig {
3211 expect: vec![3000],
3212 ..PortConfig::default()
3213 }),
3214 ..crate::daemon::Daemon::default()
3215 };
3216
3217 assert_eq!(discovery_preferred_port(&daemon), Some(3000));
3218 }
3219
3220 #[test]
3221 fn ambiguous_candidates_leave_active_port_unset() {
3222 let descendant_pids = [100].into_iter().collect();
3223 let listeners = vec![
3224 listener(100, 9000, Protocol::TCP, SocketState::Listen),
3225 listener(100, 3000, Protocol::TCP, SocketState::Listen),
3226 ];
3227
3228 assert_eq!(
3229 select_active_port(listeners, &descendant_pids, None),
3230 ActivePortSelection::Ambiguous(vec![3000, 9000])
3231 );
3232 }
3233
3234 #[test]
3235 fn bumped_ready_port_sets_only_the_resolved_primary_without_scanning() {
3236 assert_eq!(active_port_from_ready_port(3004, &[3004]), Some(3004));
3237 assert_eq!(active_port_from_ready_port(4003, &[3003, 4003]), None);
3238 }
3239
3240 #[test]
3241 fn delay_readiness_requires_delay_only_and_a_running_process() {
3242 assert!(delay_readiness_succeeded(false, false, false, true));
3243 assert!(!delay_readiness_succeeded(false, true, false, true));
3244 assert!(!delay_readiness_succeeded(false, false, true, false));
3245 assert!(!delay_readiness_succeeded(false, false, false, false));
3246 assert!(!delay_readiness_succeeded(true, false, false, true));
3247 }
3248}
3249
3250/// Check whether a daemon (by its qualified ID) is the target of any registered
3251/// slug in the global config. This is used to decide whether to run the
3252/// `detect_and_store_active_port` polling task — only slug-targeted daemons need
3253/// it, avoiding wasted `listeners::get_all()` calls for port-less daemons.
3254///
3255/// Delegates to `proxy::server::is_slug_target()` which uses the same in-memory
3256/// slug cache as the proxy hot path, so this check is cheap.
3257fn is_daemon_slug_target(id: &DaemonId) -> bool {
3258 // read_global_slugs is called once per daemon start — acceptable cost.
3259 // We intentionally avoid making this async to keep has_port_config evaluation
3260 // simple and synchronous in run_once().
3261 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
3262 slugs.iter().any(|(slug, entry)| {
3263 let daemon_name = entry.daemon.as_deref().unwrap_or(slug);
3264 id.name() == daemon_name
3265 })
3266}
3267
3268#[cfg(test)]
3269mod oneshot_tests {
3270 use super::*;
3271
3272 #[test]
3273 fn oneshot_clean_exit_is_completed() {
3274 let (status, success) = terminal_exit_state("exit", true, 0, true);
3275 assert!(matches!(status, DaemonStatus::Completed));
3276 assert!(success);
3277 }
3278
3279 #[test]
3280 fn service_clean_exit_is_still_stopped() {
3281 let (status, success) = terminal_exit_state("exit", false, 0, true);
3282 assert!(matches!(status, DaemonStatus::Stopped));
3283 assert!(success);
3284 }
3285
3286 #[test]
3287 fn oneshot_failure_is_errored_so_retry_applies() {
3288 // check_retry() only picks up errored daemons, so a non-zero exit must
3289 // not be recorded as completed.
3290 let (status, success) = terminal_exit_state("fail", true, 3, false);
3291 assert!(matches!(status, DaemonStatus::Errored(3)));
3292 assert!(!success);
3293 }
3294
3295 #[test]
3296 fn a_stop_cancels_the_retry_sequence_it_finds() {
3297 let id = DaemonId::new("retry-cancel-test", "task");
3298 let claim = SUPERVISOR.mark_retrying(&id);
3299 assert!(!claim.is_cancelled());
3300 assert!(SUPERVISOR.is_retrying(&id));
3301 SUPERVISOR.cancel_retrying(&id);
3302 assert!(claim.is_cancelled());
3303 drop(claim);
3304 assert!(!SUPERVISOR.is_retrying(&id));
3305 }
3306
3307 #[test]
3308 fn a_stop_cancels_every_sequence_for_the_daemon() {
3309 // Two starts can be working through the same daemon's retries: the
3310 // first releases the daemon's lock while it sleeps out a backoff. A
3311 // stop has to end both, not just whichever claimed it last.
3312 let id = DaemonId::new("retry-cancel-test", "concurrent");
3313 let first = SUPERVISOR.mark_retrying(&id);
3314 let second = SUPERVISOR.mark_retrying(&id);
3315 SUPERVISOR.cancel_retrying(&id);
3316 assert!(first.is_cancelled());
3317 assert!(second.is_cancelled());
3318 drop(second);
3319 // The first is still going, so the retry checker must still stand off.
3320 assert!(SUPERVISOR.is_retrying(&id));
3321 drop(first);
3322 assert!(!SUPERVISOR.is_retrying(&id));
3323 }
3324
3325 #[test]
3326 fn a_stop_invalidates_an_attempt_decided_on_before_it() {
3327 // The retry checker reads the epoch when it decides on an attempt and
3328 // `run_retry` compares it under the daemon's lock, so an approval from
3329 // before a stop cannot slip past that stop.
3330 let id = DaemonId::new("stop-epoch-test", "task");
3331 let approved_at = SUPERVISOR.stop_epoch(&id);
3332 assert_eq!(SUPERVISOR.stop_epoch(&id), approved_at);
3333 SUPERVISOR.bump_stop_epoch(&id);
3334 assert_ne!(SUPERVISOR.stop_epoch(&id), approved_at);
3335 // An attempt decided on after the stop is still fine to start.
3336 let approved_after = SUPERVISOR.stop_epoch(&id);
3337 assert_eq!(SUPERVISOR.stop_epoch(&id), approved_after);
3338 }
3339
3340 #[test]
3341 fn stop_epochs_are_tracked_per_daemon() {
3342 let stopped = DaemonId::new("stop-epoch-test", "stopped");
3343 let untouched = DaemonId::new("stop-epoch-test", "untouched");
3344 let approved_at = SUPERVISOR.stop_epoch(&untouched);
3345 SUPERVISOR.bump_stop_epoch(&stopped);
3346 assert_eq!(SUPERVISOR.stop_epoch(&untouched), approved_at);
3347 }
3348
3349 #[test]
3350 fn a_stop_leaves_a_completed_task_alone() {
3351 // It had already done its work, so the stop had nothing to interrupt.
3352 assert!(stop_keeps_finalized_status(&DaemonStatus::Completed));
3353 }
3354
3355 #[test]
3356 fn a_stop_replaces_a_failure_so_retries_do_not_resume() {
3357 // check_retry() picks up errored daemons, so a stop has to overwrite
3358 // one or it will start the task again.
3359 assert!(!stop_keeps_finalized_status(&DaemonStatus::Errored(1)));
3360 assert!(!stop_keeps_finalized_status(&DaemonStatus::Running));
3361 assert!(!stop_keeps_finalized_status(&DaemonStatus::Stopped));
3362 }
3363
3364 #[test]
3365 fn stopped_oneshot_did_not_complete() {
3366 let (status, _) = terminal_exit_state("stop", true, 0, true);
3367 assert!(matches!(status, DaemonStatus::Stopped));
3368 }
3369}
3370
3371#[cfg(all(test, unix))]
3372mod tests {
3373 use super::*;
3374
3375 #[test]
3376 fn test_resolve_run_identity_empty_without_sudo() {
3377 let identity = resolve_run_identity(None, 501, 20, None, None).unwrap();
3378 assert_eq!(identity, RunIdentity::Inherit);
3379 }
3380
3381 #[test]
3382 fn test_resolve_run_identity_sudo_fallback() {
3383 let identity = resolve_run_identity(None, 0, 0, Some("501"), Some("20")).unwrap();
3384 let RunIdentity::Switch { uid, gid, .. } = identity else {
3385 panic!("expected identity switch");
3386 };
3387 assert_eq!(uid.as_raw(), 501);
3388 assert_eq!(gid.as_raw(), 20);
3389 }
3390
3391 #[test]
3392 fn test_resolve_run_identity_ignores_stale_sudo_when_not_root() {
3393 let identity = resolve_run_identity(None, 501, 20, Some("0"), Some("0")).unwrap();
3394 assert_eq!(identity, RunIdentity::Inherit);
3395 }
3396
3397 #[test]
3398 fn test_resolve_configured_user_root_name() {
3399 let identity = resolve_configured_user("root").unwrap();
3400 let RunIdentity::Switch { uid, username, .. } = identity else {
3401 panic!("expected identity switch");
3402 };
3403 assert_eq!(uid.as_raw(), 0);
3404 assert_eq!(
3405 username.as_deref().and_then(|s| s.to_str().ok()),
3406 Some("root")
3407 );
3408 }
3409
3410 #[test]
3411 fn test_resolve_configured_user_root_uid() {
3412 let identity = resolve_configured_user("0").unwrap();
3413 let RunIdentity::Switch { uid, username, .. } = identity else {
3414 panic!("expected identity switch");
3415 };
3416 assert_eq!(uid.as_raw(), 0);
3417 assert_eq!(
3418 username.as_deref().and_then(|s| s.to_str().ok()),
3419 Some("root")
3420 );
3421 }
3422
3423 #[test]
3424 fn test_resolve_configured_user_missing_user_fails() {
3425 let err = resolve_configured_user("pitchfork-user-that-should-not-exist")
3426 .unwrap_err()
3427 .to_string();
3428 assert!(err.contains("does not exist"));
3429 }
3430
3431 #[test]
3432 fn test_resolve_run_identity_requires_root_for_user_switch() {
3433 let err = resolve_run_identity(Some("root"), 501, 20, None, None)
3434 .unwrap_err()
3435 .to_string();
3436 assert!(err.contains("Restart the supervisor with sudo"));
3437 }
3438
3439 #[test]
3440 fn test_resolve_run_identity_same_user_is_noop() {
3441 let identity = resolve_run_identity(Some("root"), 0, 0, Some("501"), Some("20")).unwrap();
3442 assert_eq!(identity, RunIdentity::Inherit);
3443 }
3444
3445 #[test]
3446 fn test_resolve_configured_user_records_home() {
3447 let identity = resolve_configured_user("root").unwrap();
3448 let RunIdentity::Switch { home, .. } = identity else {
3449 panic!("expected identity switch");
3450 };
3451 let expected = nix::unistd::User::from_name("root").unwrap().unwrap().dir;
3452 assert_eq!(home, Some(expected));
3453 }
3454
3455 fn switch_to(name: Option<&str>, home: Option<&str>) -> RunIdentity {
3456 RunIdentity::Switch {
3457 uid: nix::unistd::Uid::from_raw(501),
3458 gid: nix::unistd::Gid::from_raw(20),
3459 username: name.map(|n| CString::new(n).unwrap()),
3460 home: home.map(std::path::PathBuf::from),
3461 }
3462 }
3463
3464 /// The env a command will run with, as set on the command itself:
3465 /// `Some(None)` is an explicit removal, `None` means inherited.
3466 fn command_env(
3467 command: &tokio::process::Command,
3468 key: &str,
3469 ) -> Option<Option<std::ffi::OsString>> {
3470 command
3471 .as_std()
3472 .get_envs()
3473 .find(|(k, _)| *k == key)
3474 .map(|(_, v)| v.map(ToOwned::to_owned))
3475 }
3476
3477 /// Mirrors the order run_once applies them in.
3478 fn daemon_command(
3479 identity: &RunIdentity,
3480 daemon_env: Option<&IndexMap<String, String>>,
3481 ) -> tokio::process::Command {
3482 let mut command = tokio::process::Command::new("true");
3483 apply_identity_env(&mut command, identity);
3484 apply_runtime_env(
3485 &mut command,
3486 &DaemonId::new("identity-env-test", "api"),
3487 0,
3488 daemon_env,
3489 &[],
3490 );
3491 command
3492 }
3493
3494 #[test]
3495 fn test_identity_env_describes_the_switched_user() {
3496 let command = daemon_command(&switch_to(Some("alice"), Some("/home/alice")), None);
3497 assert_eq!(
3498 command_env(&command, "HOME"),
3499 Some(Some("/home/alice".into()))
3500 );
3501 assert_eq!(command_env(&command, "USER"), Some(Some("alice".into())));
3502 assert_eq!(command_env(&command, "LOGNAME"), Some(Some("alice".into())));
3503 }
3504
3505 #[test]
3506 fn test_identity_env_leaves_inherited_env_alone() {
3507 let command = daemon_command(&RunIdentity::Inherit, None);
3508 assert_eq!(command_env(&command, "HOME"), None);
3509 assert_eq!(command_env(&command, "USER"), None);
3510 assert_eq!(command_env(&command, "LOGNAME"), None);
3511 }
3512
3513 #[test]
3514 fn test_root_sudo_to_root_preserves_environment() {
3515 let identity = resolve_run_identity(None, 0, 0, Some("0"), Some("0")).unwrap();
3516 assert_eq!(identity, RunIdentity::Inherit);
3517 let command = daemon_command(&identity, None);
3518 for key in ["HOME", "USER", "LOGNAME"] {
3519 assert_eq!(command_env(&command, key), None);
3520 }
3521 }
3522
3523 #[test]
3524 fn test_daemon_env_overrides_identity_env() {
3525 let mut env = IndexMap::new();
3526 env.insert("HOME".to_string(), "/srv/app".to_string());
3527 env.insert("USER".to_string(), "app".to_string());
3528 let command = daemon_command(&switch_to(Some("alice"), Some("/home/alice")), Some(&env));
3529 assert_eq!(command_env(&command, "HOME"), Some(Some("/srv/app".into())));
3530 assert_eq!(command_env(&command, "USER"), Some(Some("app".into())));
3531 assert_eq!(command_env(&command, "LOGNAME"), Some(Some("alice".into())));
3532 }
3533
3534 #[test]
3535 fn test_identity_env_drops_values_without_a_passwd_entry() {
3536 // A sudo uid with no passwd entry: the supervisor's values would still
3537 // describe root, so they are removed instead.
3538 let command = daemon_command(&switch_to(None, None), None);
3539 assert_eq!(command_env(&command, "HOME"), Some(None));
3540 assert_eq!(command_env(&command, "USER"), Some(None));
3541 assert_eq!(command_env(&command, "LOGNAME"), Some(None));
3542 }
3543
3544 #[test]
3545 fn test_daemon_env_restores_values_without_a_passwd_entry() {
3546 let mut env = IndexMap::new();
3547 env.insert("HOME".to_string(), "/srv/app".to_string());
3548 let command = daemon_command(&switch_to(None, None), Some(&env));
3549 assert_eq!(command_env(&command, "HOME"), Some(Some("/srv/app".into())));
3550 }
3551}
3552
3553/// Inject proxy-related environment variables into a daemon's command.
3554///
3555/// Adds:
3556/// - `HOST` — the address the daemon should bind to (`127.0.0.1`, omitted in LAN mode)
3557/// - `PITCHFORK_URL` — the public proxy URL for this daemon (if it has a slug)
3558/// - `PITCHFORK_CA_FILE` / `NODE_EXTRA_CA_CERTS` — path to the pitchfork CA cert (if HTTPS enabled)
3559/// - `__VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS` — `.<tld>` for Vite host allowlisting
3560/// - `PITCHFORK_LAN` — set to `"1"` when LAN mode is active
3561fn inject_proxy_env(cmd: &mut tokio::process::Command, host: &Option<String>) {
3562 let s = crate::settings::settings();
3563 let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
3564
3565 if s.proxy.enable && host.is_some() && !lan_enabled {
3566 // Only force loopback binding for daemons the proxy actually routes to.
3567 // In LAN mode, daemons need to bind to 0.0.0.0 to be reachable from the network.
3568 cmd.env("HOST", "127.0.0.1");
3569 }
3570
3571 // PITCHFORK_URL: the daemon's public proxy URL (only if it is routed and proxy is enabled)
3572 if let Some(url) = build_pitchfork_url(host, &s) {
3573 cmd.env("PITCHFORK_URL", &url);
3574 }
3575
3576 // PITCHFORK_CA_FILE / NODE_EXTRA_CA_CERTS: let daemons verify TLS to each
3577 // other through the proxy. `PITCHFORK_CA_FILE` is the runtime-agnostic
3578 // name; most TLS libraries take a CA bundle path from configuration, and
3579 // several read one straight out of the environment (for example
3580 // `SSL_CERT_FILE` for OpenSSL or `REQUESTS_CA_BUNDLE` for Python).
3581 if s.proxy.enable && s.proxy.https {
3582 let ca_path = if s.proxy.tls_cert.is_empty() {
3583 crate::env::PITCHFORK_STATE_DIR.join("proxy").join("ca.pem")
3584 } else {
3585 std::path::PathBuf::from(&s.proxy.tls_cert)
3586 };
3587 if ca_path.exists() {
3588 let ca_path = ca_path.to_string_lossy().to_string();
3589 cmd.env("PITCHFORK_CA_FILE", &ca_path);
3590 cmd.env("NODE_EXTRA_CA_CERTS", &ca_path);
3591 }
3592 }
3593
3594 // __VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS: Vite host allowlisting
3595 if s.proxy.enable {
3596 let tld = if lan_enabled { "local" } else { &s.proxy.tld };
3597 cmd.env("__VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS", format!(".{tld}"));
3598 }
3599
3600 // PITCHFORK_LAN: signal to daemons that LAN mode is active
3601 if lan_enabled {
3602 cmd.env("PITCHFORK_LAN", "1");
3603 }
3604}
3605
3606/// The hostname the proxy routes to this daemon, without the TLD.
3607///
3608/// A daemon registered under a legacy `[slugs]` entry keeps that spelling,
3609/// because the proxy resolves slugs first. Otherwise the hostname is derived
3610/// from where the daemon's configuration lives.
3611async fn daemon_proxy_host(opts: &RunOptions) -> Option<String> {
3612 // Nothing consumes a hostname while the proxy is off, and deriving one
3613 // reads configuration and walks the project, so daemon starts skip it.
3614 if !crate::settings::settings().proxy.enable {
3615 return None;
3616 }
3617 // A slug carried on the run options skips the lookup below, so it needs the
3618 // same length check that lookup applies; otherwise the daemon is told a URL
3619 // the proxy refuses to route.
3620 if let Some(slug) = opts.slug.as_deref()
3621 && crate::proxy::hostname::hostname_fits(slug)
3622 {
3623 return opts.slug.clone();
3624 }
3625 // The daemon's own `dir` can point outside its project, so look the config
3626 // up from where it was defined.
3627 let config_dir = opts
3628 .watch_base_dir
3629 .clone()
3630 .unwrap_or_else(|| opts.dir.0.clone());
3631 let id = opts.id.clone();
3632 // Reading the config, the slug registry and the project's checkouts is all
3633 // file I/O, so it happens together on a blocking worker rather than on the
3634 // supervisor's executor. The lookup is the one the CLI and the proxy use,
3635 // so the daemon is told the address they advertise for it — a registered
3636 // slug when it has one, otherwise its automatic hostname.
3637 tokio::task::spawn_blocking(move || {
3638 let pt = crate::pitchfork_toml::PitchforkToml::all_merged_from(&config_dir).ok()?;
3639 let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
3640 crate::proxy::hostname::host_for_daemon(&id, pt.daemons.get(&id), &slugs)
3641 })
3642 .await
3643 .unwrap_or_default()
3644}
3645
3646/// Compute the public proxy URL for a daemon.
3647///
3648/// Returns `None` if the daemon has no hostname or the proxy is not enabled.
3649fn build_pitchfork_url(host: &Option<String>, s: &crate::settings::Settings) -> Option<String> {
3650 crate::proxy::build_proxy_url(host.as_deref(), s)
3651}
3652
3653#[cfg(test)]
3654mod ready_check_tests {
3655 use super::*;
3656 use std::time::Duration;
3657
3658 #[test]
3659 fn any_ready_check_remaining_prefers_unbounded_checks() {
3660 let http = ReadyHttp::new("http://localhost/health");
3661 let cmd = ReadyCmd::new("true");
3662
3663 assert!(any_ready_check_remaining(
3664 None,
3665 false,
3666 None,
3667 false,
3668 Some(&http),
3669 false,
3670 None,
3671 false
3672 ));
3673 assert!(any_ready_check_remaining(
3674 None,
3675 false,
3676 None,
3677 false,
3678 None,
3679 false,
3680 Some(&cmd),
3681 false
3682 ));
3683 assert!(any_ready_check_remaining(
3684 None,
3685 false,
3686 Some(&ReadyPort::new(8080)),
3687 false,
3688 Some(&http),
3689 true,
3690 Some(&cmd),
3691 true
3692 ));
3693 }
3694
3695 #[test]
3696 fn any_ready_check_remaining_exhausted_timed_checks() {
3697 let http = ReadyHttp {
3698 url: "http://localhost/health".to_string(),
3699 status: vec![],
3700 timeout: Some(Duration::from_secs(5)),
3701 };
3702 let cmd = ReadyCmd {
3703 run: "true".into(),
3704 timeout: Some(Duration::from_secs(5)),
3705 };
3706
3707 assert!(any_ready_check_remaining(
3708 None,
3709 false,
3710 None,
3711 false,
3712 Some(&http),
3713 false,
3714 Some(&cmd),
3715 false
3716 ));
3717 assert!(!any_ready_check_remaining(
3718 None,
3719 false,
3720 None,
3721 false,
3722 Some(&http),
3723 true,
3724 Some(&cmd),
3725 true
3726 ));
3727 }
3728
3729 #[tokio::test]
3730 async fn spawn_cmd_probe_reports_success() {
3731 let id = DaemonId::new("global", "probe-test");
3732 let probe = spawn_cmd_probe(&id, "true", &std::env::temp_dir(), 0, None, &[]);
3733 let status = probe.result_rx.await.unwrap().unwrap();
3734 assert!(status.success());
3735 }
3736
3737 #[tokio::test]
3738 async fn spawn_cmd_probe_stops_on_request() {
3739 let id = DaemonId::new("global", "probe-test");
3740 let probe = spawn_cmd_probe(&id, "sleep 30", &std::env::temp_dir(), 0, None, &[]);
3741 let CmdProbe {
3742 cancel_tx,
3743 result_rx,
3744 } = probe;
3745 let _ = cancel_tx.send(());
3746 let status = result_rx.await.unwrap().unwrap();
3747 assert!(!status.success());
3748 }
3749
3750 #[tokio::test]
3751 async fn spawn_cmd_probe_receives_daemon_and_resolved_port_environment() {
3752 let id = DaemonId::new("worktree", "api");
3753 let daemon_env = IndexMap::from([("CUSTOM_VALUE".to_string(), "yes".to_string())]);
3754 // Written for the platform's default shell: cmd.exe reads the probe
3755 // as written, so it needs cmd's own `%VAR%` syntax there.
3756 let check = if cfg!(windows) {
3757 r#"(if "%CUSTOM_VALUE%"=="yes" if "%PORT%"=="4100" if "%PORT0%"=="4100" if "%PORT1%"=="5100" if "%PITCHFORK_DAEMON_ID%"=="worktree/api" if "%PITCHFORK_RETRY_COUNT%"=="2" exit 0) & exit 1"#
3758 } else {
3759 r#"test "$CUSTOM_VALUE" = yes && test "$PORT" = 4100 && test "$PORT0" = 4100 && test "$PORT1" = 5100 && test "$PITCHFORK_DAEMON_ID" = worktree/api && test "$PITCHFORK_RETRY_COUNT" = 2"#
3760 };
3761 let probe = spawn_cmd_probe(
3762 &id,
3763 check,
3764 &std::env::temp_dir(),
3765 2,
3766 Some(&daemon_env),
3767 &[4100, 5100],
3768 );
3769 let status = probe.result_rx.await.unwrap().unwrap();
3770 assert!(status.success());
3771 }
3772
3773 #[test]
3774 fn configured_ready_port_follows_expected_port_bump() {
3775 assert_eq!(resolve_configured_ready_port(3000, &[3000], &[3004]), 3004);
3776 assert_eq!(
3777 resolve_configured_ready_port(4000, &[3000, 4000], &[3003, 4003]),
3778 4003
3779 );
3780 assert_eq!(resolve_configured_ready_port(8080, &[3000], &[3004]), 8080);
3781 }
3782}
3783
3784#[cfg(test)]
3785mod launch_command_tests {
3786 use super::{invalid_argv_program, launch_command};
3787 use crate::daemon_id::DaemonId;
3788
3789 #[test]
3790 fn refuses_an_argv_whose_program_rendered_empty_or_to_exec() {
3791 let id = DaemonId::new("proj", "api");
3792 assert!(invalid_argv_program(&id, &[]).is_some());
3793 assert!(invalid_argv_program(&id, &words(&["", "server.js"])).is_some());
3794 assert!(invalid_argv_program(&id, &words(&["exec", "node"])).is_some());
3795 assert_eq!(invalid_argv_program(&id, &words(&["node", ""])), None);
3796 }
3797
3798 fn words(words: &[&str]) -> Vec<String> {
3799 words.iter().map(|w| w.to_string()).collect()
3800 }
3801
3802 #[test]
3803 fn starts_the_first_word_with_the_rest_as_arguments() {
3804 let argv = words(&["node", "my server.js", "--name=\"a b\"", "&", "%PATH%"]);
3805 assert_eq!(
3806 launch_command(argv, None),
3807 (
3808 "node".to_string(),
3809 words(&["my server.js", "--name=\"a b\"", "&", "%PATH%"])
3810 )
3811 );
3812 }
3813
3814 #[test]
3815 fn mise_receives_every_word_after_the_separator() {
3816 let argv = words(&["node", "my server.js", "'single'"]);
3817 let mise = std::path::Path::new("/opt/mise/bin/mise");
3818 let (program, args) = launch_command(argv, Some(mise));
3819 assert_eq!(program, mise.to_string_lossy());
3820 assert_eq!(
3821 args,
3822 words(&["x", "--", "node", "my server.js", "'single'"])
3823 );
3824 }
3825}
3826
3827#[cfg(test)]
3828mod output_reader_tests {
3829 use super::{forward_output_lines, output_line_text};
3830
3831 #[test]
3832 fn line_endings_are_removed() {
3833 assert_eq!(output_line_text(b"ready\n"), "ready");
3834 // A PTY's ONLCR, and a program that already wrote `\r\n` through it.
3835 assert_eq!(output_line_text(b"ready\r\n"), "ready");
3836 assert_eq!(output_line_text(b"ready\r\r\n"), "ready");
3837 // The last line of the output may have no newline at all.
3838 assert_eq!(output_line_text(b"ready"), "ready");
3839 }
3840
3841 #[tokio::test]
3842 async fn a_line_that_is_not_utf8_does_not_end_the_read() {
3843 // What cut the output off before: `next_line` fails on the second
3844 // line, and nothing after it was read.
3845 let output: &[u8] = b"before\nbad \xff\xfe bytes\nafter 1\nafter 2\nlast";
3846 let (tx, mut rx) = tokio::sync::mpsc::channel(16);
3847 forward_output_lines(tokio::io::BufReader::new(output), tx).await;
3848
3849 let mut lines = Vec::new();
3850 while let Some(line) = rx.recv().await {
3851 lines.push(line.text);
3852 }
3853 assert_eq!(lines.len(), 5, "{lines:?}");
3854 assert_eq!(lines[0], "before");
3855 assert!(lines[1].starts_with("bad "), "{:?}", lines[1]);
3856 assert_eq!(&lines[2..], ["after 1", "after 2", "last"]);
3857 }
3858}
3859
3860#[cfg(test)]
3861mod output_reader_eio_tests {
3862 use super::forward_output_lines;
3863 use std::pin::Pin;
3864 use std::task::{Context, Poll};
3865
3866 /// Output that ends as a Linux PTY master does once the slave closes:
3867 /// the last bytes, with no newline, then `EIO` instead of end of file.
3868 struct EndsWithEio(Option<&'static [u8]>);
3869
3870 impl tokio::io::AsyncRead for EndsWithEio {
3871 fn poll_read(
3872 mut self: Pin<&mut Self>,
3873 _cx: &mut Context<'_>,
3874 buf: &mut tokio::io::ReadBuf<'_>,
3875 ) -> Poll<std::io::Result<()>> {
3876 match self.0.take() {
3877 Some(bytes) => {
3878 buf.put_slice(bytes);
3879 Poll::Ready(Ok(()))
3880 }
3881 None => Poll::Ready(Err(std::io::Error::from_raw_os_error(5))),
3882 }
3883 }
3884 }
3885
3886 #[tokio::test]
3887 async fn a_last_line_cut_off_by_eio_is_still_forwarded() {
3888 let reader = tokio::io::BufReader::new(EndsWithEio(Some(b"first\nlast without newline")));
3889 let (tx, mut rx) = tokio::sync::mpsc::channel(16);
3890 forward_output_lines(reader, tx).await;
3891
3892 let mut lines = Vec::new();
3893 while let Some(line) = rx.recv().await {
3894 lines.push(line.text);
3895 }
3896 assert_eq!(lines, ["first", "last without newline"]);
3897 }
3898}