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