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