Skip to main content

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}