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::daemon::RunOptions;
8use crate::daemon_id::DaemonId;
9use crate::daemon_status::DaemonStatus;
10use crate::error::PortError;
11use crate::ipc::IpcResponse;
12use crate::log_store::LogStore;
13use crate::log_store::sqlite::LOG_STORE;
14use crate::pitchfork_toml::{ReadyCmd, ReadyHttp, ReadyOutput, ReadyPort};
15use crate::procs::PROCS;
16use crate::settings::settings;
17use crate::shell::Shell;
18use crate::supervisor::state::UpsertDaemonOpts;
19use crate::{Result, env};
20use indexmap::IndexMap;
21use miette::IntoDiagnostic;
22use once_cell::sync::Lazy;
23use regex::Regex;
24use std::collections::HashMap;
25#[cfg(unix)]
26use std::ffi::CString;
27use std::sync::{Arc, atomic};
28use std::time::Duration;
29use tokio::io::AsyncBufReadExt;
30use tokio::select;
31use tokio::sync::oneshot;
32use tokio::time;
33
34/// Cache for compiled regex patterns to avoid recompilation on daemon restarts
35static REGEX_CACHE: Lazy<std::sync::Mutex<HashMap<String, Regex>>> =
36    Lazy::new(|| std::sync::Mutex::new(HashMap::new()));
37
38fn resolve_configured_ready_port(
39    configured_port: u16,
40    expected_ports: &[u16],
41    resolved_ports: &[u16],
42) -> u16 {
43    let bump_offset = resolved_ports
44        .first()
45        .unwrap_or(&0)
46        .saturating_sub(*expected_ports.first().unwrap_or(&0));
47    if expected_ports.contains(&configured_port) && bump_offset > 0 {
48        configured_port
49            .checked_add(bump_offset)
50            .unwrap_or(configured_port)
51    } else {
52        configured_port
53    }
54}
55
56#[cfg(unix)]
57#[derive(Clone, Debug, PartialEq, Eq)]
58enum RunIdentity {
59    Inherit,
60    Switch {
61        uid: nix::unistd::Uid,
62        gid: nix::unistd::Gid,
63        username: Option<CString>,
64    },
65}
66
67/// Get or compile a regex pattern, caching the result for future use
68pub(crate) fn get_or_compile_regex(pattern: &str) -> Option<Regex> {
69    let mut cache = REGEX_CACHE.lock().unwrap_or_else(|e| e.into_inner());
70    if let Some(re) = cache.get(pattern) {
71        return Some(re.clone());
72    }
73    match Regex::new(pattern) {
74        Ok(re) => {
75            cache.insert(pattern.to_string(), re.clone());
76            Some(re)
77        }
78        Err(e) => {
79            error!("invalid regex pattern '{pattern}': {e}");
80            None
81        }
82    }
83}
84
85/// Handle for an in-flight readiness command probe.
86///
87/// The spawned task owns the `tokio::process::Child` and waits for either the
88/// process to exit or the cancel signal. Dropping the handle without cancelling
89/// leaves the task running, but the child is started with `kill_on_drop(true)`
90/// so it will still be terminated when the task ends.
91struct CmdProbe {
92    cancel_tx: tokio::sync::oneshot::Sender<()>,
93    result_rx: tokio::sync::oneshot::Receiver<std::io::Result<std::process::ExitStatus>>,
94}
95
96/// Spawn a readiness command probe and return a handle that can be used to wait
97/// for the exit status or cancel the probe.
98///
99/// The probe is started with `kill_on_drop(true)` as a cancellation fallback. The
100/// spawned task waits for the process to exit; if cancellation is requested, it
101/// kills the child and waits for it to reap before reporting the result.
102fn apply_runtime_env(
103    command: &mut tokio::process::Command,
104    id: &DaemonId,
105    retry_count: u32,
106    daemon_env: Option<&IndexMap<String, String>>,
107    resolved_ports: &[u16],
108) {
109    if let Some(ref path) = *env::ORIGINAL_PATH {
110        command.env("PATH", path);
111    }
112    if let Some(env_vars) = daemon_env {
113        command.envs(env_vars);
114    }
115    command
116        .env("PITCHFORK_DAEMON_ID", id.qualified())
117        .env("PITCHFORK_DAEMON_NAMESPACE", id.namespace())
118        .env("PITCHFORK_RETRY_COUNT", retry_count.to_string());
119    if let Some(port) = resolved_ports.first() {
120        command.env("PORT", port.to_string());
121        for (index, port) in resolved_ports.iter().enumerate() {
122            command.env(format!("PORT{index}"), port.to_string());
123        }
124    }
125}
126
127fn spawn_cmd_probe(
128    id: &DaemonId,
129    cmd: &str,
130    dir: &std::path::Path,
131    retry_count: u32,
132    daemon_env: Option<&IndexMap<String, String>>,
133    resolved_ports: &[u16],
134) -> CmdProbe {
135    // Use the configured general.shell setting (same as daemon run and hooks)
136    // instead of default_for_platform(). On Windows, default_for_platform()
137    // returns Shell::Cmd which cannot parse Unix-style commands like
138    // "sleep 1; true". Falls back to default_for_platform() if the setting
139    // is empty or unparseable.
140    let shell_setting = settings().general.shell.clone();
141    let mut command = match shell_words::split(&shell_setting) {
142        Ok(parts) if !parts.is_empty() => {
143            let (program, args) = parts.split_first().unwrap();
144            let mut c = tokio::process::Command::new(program);
145            c.args(args);
146            c.arg(cmd);
147            c
148        }
149        _ => Shell::default_for_platform().command(cmd),
150    };
151    command
152        .current_dir(dir)
153        .stdout(std::process::Stdio::null())
154        .stderr(std::process::Stdio::null())
155        .kill_on_drop(true);
156    apply_runtime_env(&mut command, id, retry_count, daemon_env, resolved_ports);
157    let mut child = match command.spawn() {
158        Ok(child) => child,
159        Err(e) => {
160            warn!("daemon {id}: failed to spawn readiness command probe: {e}");
161            // Return a probe whose result channel is already closed. The caller will
162            // treat this the same as a probe that exited non-zero and respawn after
163            // the ready_check_interval, preserving the existing retry behaviour.
164            let (cancel_tx, _) = tokio::sync::oneshot::channel();
165            let (_, result_rx) = tokio::sync::oneshot::channel();
166            return CmdProbe {
167                cancel_tx,
168                result_rx,
169            };
170        }
171    };
172
173    let (cancel_tx, mut cancel_rx) = tokio::sync::oneshot::channel();
174    let (result_tx, result_rx) = tokio::sync::oneshot::channel();
175
176    tokio::spawn(async move {
177        let status = tokio::select! {
178            status = child.wait() => status,
179            _ = &mut cancel_rx => {
180                let mut child = child;
181                let _ = child.kill().await;
182                child.wait().await
183            }
184        };
185        let _ = result_tx.send(status);
186    });
187
188    CmdProbe {
189        cancel_tx,
190        result_rx,
191    }
192}
193
194/// Cancel an active command probe and clear its handle.
195fn stop_cmd_probe_state(probe: &mut Option<CmdProbe>) {
196    if let Some(p) = probe.take() {
197        let _ = p.cancel_tx.send(());
198    }
199}
200
201/// Spawn a detached task that kills a daemon's process group after its
202/// readiness checks are exhausted, logging a failed kill instead of
203/// discarding it. The returned handle is awaited before the readiness
204/// failure is reported so the process group is down by then.
205fn spawn_ready_fail_kill(
206    id: DaemonId,
207    pid: u32,
208    stop_cfg: crate::config_types::StopConfig,
209) -> tokio::task::JoinHandle<()> {
210    tokio::spawn(async move {
211        if let Err(e) = PROCS
212            .kill_process_group_async(pid, stop_cfg.signal.into(), stop_cfg.timeout)
213            .await
214        {
215            error!("daemon {id}: failed to kill pid {pid} after readiness failure: {e}");
216        }
217    })
218}
219
220/// Returns true if any configured readiness check can still succeed.
221/// A check with no timeout is unbounded; a timed check can still succeed until its
222/// deadline fires. `ready_delay` is only used as a fallback when no other check is
223/// configured, so it is not counted here.
224#[allow(clippy::too_many_arguments)]
225fn any_ready_check_remaining(
226    ready_output: Option<&ReadyOutput>,
227    output_exhausted: bool,
228    ready_port: Option<&ReadyPort>,
229    port_exhausted: bool,
230    ready_http: Option<&ReadyHttp>,
231    http_exhausted: bool,
232    ready_cmd: Option<&ReadyCmd>,
233    cmd_exhausted: bool,
234) -> bool {
235    ready_output.is_some_and(|o| o.timeout.is_none() || !output_exhausted)
236        || ready_port.is_some_and(|p| p.timeout.is_none() || !port_exhausted)
237        || ready_http.is_some_and(|h| h.timeout.is_none() || !http_exhausted)
238        || ready_cmd.is_some_and(|c| c.timeout.is_none() || !cmd_exhausted)
239}
240
241/// How long a failed start waits for the daemon's output to become queryable
242/// before reporting. Typically satisfied in a few dozen milliseconds; a daemon
243/// that failed without printing anything waits the whole of it, so keep it
244/// short.
245const SINK_OUTPUT_TIMEOUT: Duration = Duration::from_millis(400);
246
247impl Supervisor {
248    /// Run a daemon, handling retries if configured
249    pub async fn run(&self, opts: RunOptions) -> Result<IpcResponse> {
250        let id = &opts.id;
251        let cmd = opts.cmd.clone();
252
253        // Clear any pending autostop for this daemon since it's being started
254        {
255            let mut pending = self.pending_autostops.lock().await;
256            if pending.remove(id).is_some() {
257                info!("cleared pending autostop for {id} (daemon starting)");
258            }
259        }
260
261        // Serialize against any in-flight stop of this daemon: a stop now
262        // waits for the whole process group to exit, so the Stopping window
263        // can last seconds instead of milliseconds. Starting through that
264        // window would collide with the dying instance (duplicate processes,
265        // port conflicts). Acquiring the stop lock waits the stop out; the
266        // state is re-read afterwards. The guard is owned and handed to
267        // run_once, which holds it until the new daemon's Running state and
268        // PID are persisted — releasing it before that point would let a
269        // concurrent run pass this same check (duplicate processes) or let a
270        // concurrent stop see no PID and return without stopping anything.
271        let mut stop_guard = Some(self.stop_lock(id).await.lock_owned().await);
272        let daemon = self.get_daemon(id).await;
273        if let Some(daemon) = daemon {
274            // Stopping state is treated as "not running" - the monitoring task will clean it up
275            // Only check for Running state with a valid PID
276            if !daemon.status.is_stopping()
277                && !daemon.status.is_stopped()
278                && let Some(pid) = daemon.pid
279            {
280                if opts.force {
281                    self.stop_locked(id).await?;
282                    info!("run: stop completed for daemon {id}");
283                } else {
284                    warn!("daemon {id} already running with pid {pid}");
285                    return Ok(IpcResponse::DaemonAlreadyRunning);
286                }
287            }
288        }
289
290        // If wait_ready is true and retry is configured, implement retry loop
291        if opts.wait_ready && opts.retry.count() > 0 {
292            // Use saturating_add to avoid overflow when retry = u32::MAX (infinite)
293            let max_attempts = opts.retry.count().saturating_add(1);
294            for attempt in 0..max_attempts {
295                let mut retry_opts = opts.clone();
296                retry_opts.retry_count = attempt;
297                retry_opts.cmd = cmd.clone();
298
299                // The first attempt starts under the guard held since the
300                // running check above; later attempts re-acquire it so stops
301                // are not locked out during the backoff sleeps.
302                let guard = match stop_guard.take() {
303                    Some(guard) => guard,
304                    None => self.stop_lock(id).await.lock_owned().await,
305                };
306                let result = self.run_once(retry_opts, guard).await?;
307
308                match result {
309                    IpcResponse::DaemonReady { daemon } => {
310                        return Ok(IpcResponse::DaemonReady { daemon });
311                    }
312                    IpcResponse::DaemonFailedWithCode { exit_code } => {
313                        if attempt < opts.retry.count() {
314                            let backoff_secs = 2u64.saturating_pow(attempt).min(3600);
315                            info!(
316                                "daemon {id} failed (attempt {}/{}), retrying in {}s",
317                                attempt + 1,
318                                max_attempts,
319                                backoff_secs
320                            );
321                            fire_hook(
322                                HookType::OnRetry,
323                                id.clone(),
324                                opts.dir.0.clone(),
325                                attempt + 1,
326                                opts.env.clone(),
327                                vec![],
328                            )
329                            .await;
330                            time::sleep(Duration::from_secs(backoff_secs)).await;
331                            continue;
332                        } else {
333                            info!("daemon {id} failed after {max_attempts} attempts");
334                            return Ok(IpcResponse::DaemonFailedWithCode { exit_code });
335                        }
336                    }
337                    other => return Ok(other),
338                }
339            }
340        }
341
342        // No retry or wait_ready is false
343        let guard = match stop_guard.take() {
344            Some(guard) => guard,
345            None => self.stop_lock(id).await.lock_owned().await,
346        };
347        self.run_once(opts, guard).await
348    }
349
350    /// Run a daemon once (single attempt).
351    ///
352    /// `stop_guard` is this daemon's stop lock, acquired by `run` before the
353    /// already-running check. It is held through spawning until the Running
354    /// state and PID are persisted (or an early failure returns), then dropped
355    /// before the potentially unbounded readiness wait.
356    pub(crate) async fn run_once(
357        &self,
358        opts: RunOptions,
359        stop_guard: tokio::sync::OwnedMutexGuard<()>,
360    ) -> Result<IpcResponse> {
361        let id = &opts.id;
362        let original_cmd = opts.cmd.clone(); // Save original command for persistence
363
364        // Create channel for readiness notification if wait_ready is true
365        let (ready_tx, ready_rx) = if opts.wait_ready {
366            let (tx, rx) = oneshot::channel();
367            (Some(tx), Some(rx))
368        } else {
369            (None, None)
370        };
371
372        // Check port availability and apply auto-bump if configured
373        let expected_ports = opts
374            .port
375            .as_ref()
376            .map(|p| p.expect.clone())
377            .unwrap_or_default();
378        let (resolved_ports, effective_ready_port) = if !expected_ports.is_empty() {
379            let port_cfg = opts.port.as_ref().unwrap();
380            match check_ports_available(
381                &expected_ports,
382                port_cfg.auto_bump(),
383                port_cfg.max_bump_attempts(),
384            )
385            .await
386            {
387                Ok(resolved) => {
388                    let ready_port = if let Some(configured_port) =
389                        opts.ready_port.as_ref().and_then(|p| p.as_port())
390                    {
391                        Some(resolve_configured_ready_port(
392                            configured_port,
393                            &expected_ports,
394                            &resolved,
395                        ))
396                    } else if opts.ready_output.is_none()
397                        && opts.ready_http.is_none()
398                        && opts.ready_cmd.is_none()
399                        && opts.ready_delay.is_none()
400                    {
401                        // No other ready check configured — use the first expected port as a
402                        // TCP port readiness check so the daemon is considered ready once it
403                        // starts listening.  Skip port 0 (ephemeral port request).
404                        resolved.first().copied().filter(|&p| p != 0)
405                    } else {
406                        // Another ready check is configured (output/http/cmd/delay).
407                        // Don't add an implicit TCP port check — it could race and fire
408                        // before the daemon has produced any output.
409                        None
410                    };
411                    info!("daemon {id}: ports {expected_ports:?} resolved to {resolved:?}");
412                    (resolved, ready_port)
413                }
414                Err(e) => {
415                    error!("daemon {id}: port check failed: {e}");
416                    // Convert PortError to structured IPC response
417                    if let Some(port_error) = e.downcast_ref::<PortError>() {
418                        match port_error {
419                            PortError::InUse { port, process, pid } => {
420                                return Ok(IpcResponse::PortConflict {
421                                    port: *port,
422                                    process: process.clone(),
423                                    pid: *pid,
424                                });
425                            }
426                            PortError::NoAvailablePort {
427                                start_port,
428                                attempts,
429                            } => {
430                                return Ok(IpcResponse::NoAvailablePort {
431                                    start_port: *start_port,
432                                    attempts: *attempts,
433                                });
434                            }
435                        }
436                    }
437                    return Ok(IpcResponse::DaemonFailed {
438                        error: e.to_string(),
439                    });
440                }
441            }
442        } else {
443            // When ready_port is set without expected_port, check that the port
444            // is not already occupied.  If another process is listening on it,
445            // the TCP readiness probe would immediately succeed and pitchfork
446            // would falsely consider the daemon ready — routing proxy traffic to
447            // the wrong process.
448            if let Some(port) = opts.ready_port.as_ref().and_then(|p| p.as_port())
449                && port > 0
450                && let Some((pid, process)) = detect_port_conflict(port).await
451            {
452                return Ok(IpcResponse::PortConflict { port, process, pid });
453            }
454            (
455                Vec::new(),
456                opts.ready_port.as_ref().and_then(|p| p.as_port()),
457            )
458        };
459
460        // Parse the configured shell (default "sh -c") into program + args.
461        // The run script is passed verbatim as the final argument, avoiding the
462        // lossy split->join round-trip that previously mangled $VAR/glob expansion.
463        let shell_setting = settings().general.shell.clone();
464        let shell_parts = match shell_words::split(&shell_setting) {
465            Ok(parts) if !parts.is_empty() => parts,
466            Ok(_) => {
467                return Ok(IpcResponse::DaemonFailed {
468                    error: "general.shell setting is empty".to_string(),
469                });
470            }
471            Err(e) => {
472                return Ok(IpcResponse::DaemonFailed {
473                    error: format!("failed to parse general.shell setting {shell_setting:?}: {e}"),
474                });
475            }
476        };
477        let (shell_program, shell_args) = shell_parts.split_first().unwrap();
478
479        // Use the original run string verbatim; fall back to joining cmd for
480        // ad-hoc commands (e.g. `pitchfork run -- cmd args`) that have no run string.
481        // We don't prepend `exec` because it breaks compound commands (e.g. `exec a && b`
482        // silently drops `b`). Users can add `exec` themselves in the run string.
483        let run_script = opts
484            .run
485            .clone()
486            .unwrap_or_else(|| shell_words::join(&original_cmd));
487
488        let (program, args) = if opts.mise.unwrap_or(settings().general.mise) {
489            match settings().resolve_mise_bin() {
490                Some(mise_bin) => {
491                    let mise_bin_str = mise_bin.to_string_lossy().to_string();
492                    info!("daemon {id}: wrapping command with mise ({mise_bin_str})");
493                    let mut args = vec!["x".to_string(), "--".to_string()];
494                    args.push(shell_program.clone());
495                    args.extend(shell_args.iter().cloned());
496                    args.push(run_script);
497                    (mise_bin_str, args)
498                }
499                None => {
500                    warn!("daemon {id}: mise=true but mise binary not found, running without mise");
501                    let mut args: Vec<String> = shell_args.to_vec();
502                    args.push(run_script);
503                    (shell_program.clone(), args)
504                }
505            }
506        } else {
507            let mut args: Vec<String> = shell_args.to_vec();
508            args.push(run_script);
509            (shell_program.clone(), args)
510        };
511        #[cfg(unix)]
512        let run_identity = match resolve_effective_run_identity(opts.user.as_deref()) {
513            Ok(identity) => identity,
514            Err(e) => {
515                return Ok(IpcResponse::DaemonFailed {
516                    error: e.to_string(),
517                });
518            }
519        };
520        info!("run: spawning daemon {id} with {program} {args:?}");
521
522        // Allocate PTY if configured
523        #[cfg(unix)]
524        let pty_pair = if opts.pty.unwrap_or(false) {
525            match super::pty::openpty() {
526                Ok(pair) => {
527                    info!("daemon {id}: allocated PTY (pty = true)");
528                    Some(pair)
529                }
530                Err(e) => {
531                    warn!("daemon {id}: failed to allocate PTY, falling back to pipes: {e}");
532                    None
533                }
534            }
535        } else {
536            None
537        };
538
539        // Output reaches the monitoring task either from readers this process
540        // owns or, when a sink owns the stream, relayed over IPC. The channel is
541        // created here rather than in that task so it exists before the sink
542        // starts: a daemon whose very first line matches its readiness pattern
543        // would otherwise have the match reported with nowhere to deliver it.
544        let (output_tx, output_rx) = tokio::sync::mpsc::channel::<super::OutputLine>(256);
545        let mut output_relay = None;
546
547        // Set up out-of-process capture before building the command, so the
548        // daemon can be handed the pipe's write end directly.
549        let mut sink_pipe = None;
550        let mut sink_writer = None;
551        let mut sink_child = None;
552        if super::log_sink::is_supported(&opts) {
553            let log_format = opts
554                .log_format
555                .clone()
556                .unwrap_or_else(|| settings().logs.log_format.clone());
557            let watch_for = super::log_sink::WatchFor::from_opts(id, &opts);
558            // The token ties this attempt's sink to this attempt's channel, so
559            // a sink still draining a previous attempt cannot report into it.
560            let relay_token = if watch_for.is_empty() {
561                0
562            } else {
563                let relay = super::log_sink::OutputRelay::register(id, output_tx.clone());
564                let token = relay.token();
565                output_relay = Some(relay);
566                token
567            };
568            match super::log_sink::SinkPipe::new(log_format, watch_for, relay_token) {
569                Ok((pipe, writer)) => match pipe.start(id) {
570                    Ok(child) => {
571                        sink_child = Some(super::log_sink::PendingSink::new(child));
572                        sink_pipe = Some(pipe);
573                        sink_writer = Some(writer);
574                    }
575                    Err(e) => {
576                        warn!("could not start log sink for {id}, capturing in-process: {e}");
577                    }
578                },
579                Err(e) => {
580                    // Fall back to in-process capture rather than refusing to
581                    // start the daemon.
582                    warn!("could not create log pipe for {id}, capturing in-process: {e}");
583                }
584            }
585        }
586
587        let mut cmd = tokio::process::Command::new(&program);
588
589        #[cfg(unix)]
590        if let Some(ref pair) = pty_pair {
591            // PTY mode: connect both stdout and stderr to the slave PTY.
592            // The child uses the slave for stdin/stdout/stderr, and we read
593            // output from the master.
594            let slave_file = std::fs::File::from(
595                pair.slave
596                    .try_clone()
597                    .map_err(|e| miette::miette!("failed to dup slave PTY fd: {e}"))?,
598            );
599            cmd.stdin(std::process::Stdio::from(slave_file.try_clone().map_err(
600                |e| miette::miette!("failed to clone slave PTY fd for stdin: {e}"),
601            )?));
602            cmd.stdout(std::process::Stdio::from(slave_file.try_clone().map_err(
603                |e| miette::miette!("failed to clone slave PTY fd for stdout: {e}"),
604            )?));
605            cmd.stderr(std::process::Stdio::from(slave_file));
606        } else if let Some(writer) = sink_writer.take() {
607            // Capture belongs to a sibling sink process, so the daemon writes
608            // to a pipe this process does not read. See supervisor::log_sink.
609            let dup = writer
610                .try_clone()
611                .map_err(|e| miette::miette!("failed to dup log pipe for stderr: {e}"))?;
612            cmd.stdout(std::process::Stdio::from(writer))
613                .stderr(std::process::Stdio::from(dup));
614        } else {
615            cmd.stdout(std::process::Stdio::piped())
616                .stderr(std::process::Stdio::piped());
617        }
618
619        #[cfg(not(unix))]
620        if let Some(writer) = sink_writer.take() {
621            let dup = writer
622                .try_clone()
623                .map_err(|e| miette::miette!("failed to dup log pipe for stderr: {e}"))?;
624            cmd.stdout(std::process::Stdio::from(writer))
625                .stderr(std::process::Stdio::from(dup));
626        } else {
627            cmd.stdout(std::process::Stdio::piped())
628                .stderr(std::process::Stdio::piped());
629        }
630
631        cmd.args(&args).current_dir(&opts.dir);
632
633        #[cfg(unix)]
634        if pty_pair.is_none() {
635            cmd.stdin(std::process::Stdio::null());
636        }
637
638        #[cfg(not(unix))]
639        cmd.stdin(std::process::Stdio::null());
640
641        apply_runtime_env(
642            &mut cmd,
643            id,
644            opts.retry_count,
645            opts.env.as_ref(),
646            &resolved_ports,
647        );
648
649        // Inject proxy-related environment variables
650        inject_proxy_env(&mut cmd, &opts.slug);
651
652        #[cfg(unix)]
653        {
654            let run_identity = run_identity.clone();
655            let use_pty = pty_pair.is_some();
656            unsafe {
657                cmd.pre_exec(move || {
658                    nix::unistd::setsid().map_err(nix_to_io_error)?;
659
660                    // When using a PTY, set the slave as the controlling terminal.
661                    // The slave FD has already been dup'd onto stdin/stdout/stderr
662                    // by tokio, so we can use stdin (fd 0) for TIOCSCTTY.
663                    if use_pty {
664                        let ret = libc::ioctl(0, libc::TIOCSCTTY as libc::c_ulong, 0);
665                        if ret < 0 {
666                            // Non-fatal: the process can still run without
667                            // a controlling terminal.
668                            #[cfg(target_os = "linux")]
669                            eprintln!(
670                                "pitchfork: TIOCSCTTY failed: {}",
671                                std::io::Error::last_os_error()
672                            );
673                        }
674                    }
675
676                    apply_run_identity(&run_identity)?;
677                    Ok(())
678                });
679            }
680        }
681
682        // Timestamp the run so a failed start can wait for this attempt's output
683        // specifically, rather than seeing an earlier attempt's.
684        let spawn_time = chrono::Local::now();
685        // A sink is already running at this point. Both bail-outs below have to
686        // reap it explicitly: dropping the handle only reaps on a best-effort
687        // basis, and run_once runs once per retry attempt, so a daemon that
688        // consistently fails to spawn would otherwise accumulate sinks.
689        // A failed spawn returns here; the sink is terminated by PendingSink.
690        let mut child = cmd.spawn().into_diagnostic()?;
691        let pid = match child.id() {
692            Some(p) => p,
693            None => {
694                warn!("Daemon {id} exited before PID could be captured");
695                // Unlike a daemon that never started, this one ran and may have
696                // said why it gave up, and its output is the only diagnosis
697                // available. Its write end is already closed, so the sink is on
698                // its way to end of file: let it finish writing before reporting,
699                // then reap whatever is left of it.
700                if sink_child.is_some() {
701                    super::log_sink::wait_for_output(id, spawn_time, SINK_OUTPUT_TIMEOUT).await;
702                }
703                return Ok(IpcResponse::DaemonFailed {
704                    error: "Process exited immediately".to_string(),
705                });
706            }
707        };
708        info!("started daemon {id} with pid {pid}");
709        PROCS.refresh_pids(&[pid]);
710        // Register the daemon as monitored BEFORE persisting the Running
711        // state. The orphan reconciler treats any running, unmonitored PID
712        // as an orphan; if the state became visible first, a concurrent
713        // reconciliation pass could adopt — or under the kill policy,
714        // terminate — a daemon that was just legitimately started. The RAII
715        // guard unregisters on any early-error path below and is otherwise
716        // handed to the monitoring task.
717        let monitored_guard = super::adopt::MonitoredGuard::register(id.clone(), pid);
718        let monitor_token = monitored_guard.token();
719
720        // Hand the retained read end to a sink and keep one running for as long
721        // as this daemon is monitored.
722        let using_sink = sink_pipe.is_some();
723        // Take the sink out of the guard only once there is a pipe to supervise
724        // it with, so it is never left running unsupervised.
725        if let Some(pipe) = sink_pipe.take()
726            && let Some(child) = sink_child.as_mut().and_then(|pending| pending.take())
727        {
728            pipe.supervise(id.clone(), monitor_token, child);
729        }
730        let daemon = self
731            .upsert_daemon(
732                UpsertDaemonOpts::from_run_options(&opts, DaemonStatus::Running)
733                    .set(|o| {
734                        o.pid = Some(pid);
735                        o.cmd = Some(original_cmd);
736                        o.ready_port = effective_ready_port.map(|p| ReadyPort {
737                            port: Some(p),
738                            template: None,
739                            timeout: opts.ready_port.as_ref().and_then(|rp| rp.timeout),
740                        });
741                        o.port = crate::config_types::PortConfig::from_parts(
742                            expected_ports,
743                            opts.port.as_ref().map(|p| p.bump).unwrap_or_default(),
744                        );
745                        o.resolved_port = resolved_ports;
746                    })
747                    .build(),
748            )
749            .await?;
750
751        // Running state and PID are now persisted: concurrent run/stop calls
752        // observe a running daemon and behave correctly, so release the stop
753        // lock rather than holding it through the readiness wait below, which
754        // can take arbitrarily long.
755        drop(stop_guard);
756
757        let id_clone = id.clone();
758        let ready_delay = opts.ready_delay;
759        let ready_output = opts.ready_output.clone();
760        let ready_http = opts.ready_http.clone();
761        let ready_port = effective_ready_port;
762        let implicit_ready_port = ready_port.map(|p| ReadyPort {
763            port: Some(p),
764            template: None,
765            timeout: None,
766        });
767        let ready_port_config = opts.ready_port.clone().or(implicit_ready_port);
768        let ready_cmd = opts.ready_cmd.clone();
769        let daemon_dir = opts.dir.0.clone();
770        let hook_retry_count = opts.retry_count;
771        let hook_retry = opts.retry;
772        let hook_daemon_env = opts.env.clone();
773        let readiness_daemon_env = opts.env.clone();
774        let readiness_resolved_ports = daemon.resolved_port.clone();
775        let on_output_hook = opts.on_output_hook.clone();
776        // Whether this daemon has any port-related config — used to skip the
777        // active_port detection task for daemons that never bind a port (e.g. `sleep 60`).
778        // When the proxy is enabled, only detect active_port for daemons that are
779        // actually referenced by a registered slug, rather than blanket-polling every
780        // daemon (which wastes ~7.5 s of listeners::get_all() calls per port-less daemon).
781        let has_port_config = opts.port.as_ref().is_some_and(|p| !p.expect.is_empty())
782            || (settings().proxy.enable && is_daemon_slug_target(id));
783        // The first expected port, if configured. When the ready_port check
784        // succeeds on this exact port we can set active_port directly instead
785        // of spawning detect_and_store_active_port (which relies on
786        // listeners::get_all() + process-tree traversal and is unreliable on
787        // Windows where Git Bash PID mapping can break descendant lookups).
788        let expected_port: Option<u16> = opts.port.as_ref().and_then(|p| p.expect.first().copied());
789        let daemon_pid = pid;
790
791        // Prepare output readers before spawning the monitoring task.
792        // In PTY mode, we read from the PTY master FD.
793        // In pipe mode, we read from separate stdout/stderr pipes.
794        #[cfg(unix)]
795        let pty_reader = pty_pair.map(|p| {
796            tokio::io::BufReader::new(tokio::fs::File::from_std(std::fs::File::from(p.master)))
797                .lines()
798        });
799        #[cfg(not(unix))]
800        let pty_reader: Option<tokio::io::Lines<tokio::io::BufReader<tokio::fs::File>>> = None;
801        let stdout_reader = if pty_reader.is_none() {
802            child
803                .stdout
804                .take()
805                .map(|s| tokio::io::BufReader::new(s).lines())
806        } else {
807            None
808        };
809        let stderr_reader = if pty_reader.is_none() {
810            child
811                .stderr
812                .take()
813                .map(|s| tokio::io::BufReader::new(s).lines())
814        } else {
815            None
816        };
817
818        if !using_sink
819            && pty_reader.is_none()
820            && (stdout_reader.is_none() || stderr_reader.is_none())
821        {
822            error!("Failed to capture stdout/stderr for daemon {id}");
823        }
824
825        tokio::spawn(async move {
826            let id = id_clone;
827            // Registered before the Running upsert above; unregisters when
828            // this monitoring task ends.
829            let _monitored_guard = monitored_guard;
830            // Likewise for sink-relayed output: dropping this stops the
831            // supervisor delivering into a channel nobody is reading. Dropped
832            // explicitly once the daemon exits, before the drain below.
833            let output_relay = output_relay;
834
835            // Merge all output sources (PTY master OR stdout+stderr, or a
836            // sink's IPC reports) into a single channel.
837            let mut output_rx = output_rx;
838
839            if let Some(mut reader) = pty_reader {
840                // PTY mode: single merged stream from the master.
841                // output_tx is moved into the spawn; when the reader ends the
842                // channel closes automatically.
843                tokio::spawn(async move {
844                    while let Ok(Some(mut line)) = reader.next_line().await {
845                        // PTY slave uses ONLCR: \n → \r\n; strip the trailing \r.
846                        if line.ends_with('\r') {
847                            line.pop();
848                        }
849                        if output_tx
850                            .send(super::OutputLine {
851                                text: line,
852                                source: super::OutputSource::Local,
853                            })
854                            .await
855                            .is_err()
856                        {
857                            break;
858                        }
859                    }
860                });
861            } else {
862                // Pipe mode: stdout and stderr are merged into the same channel.
863                // Both `ready_output` and `on_output_hook` patterns match against
864                // lines from either stream, which is the expected behavior (a
865                // "server ready" message may appear on stderr in some tools).
866                if let Some(mut stdout) = stdout_reader {
867                    let tx = output_tx.clone();
868                    tokio::spawn(async move {
869                        while let Ok(Some(line)) = stdout.next_line().await {
870                            if tx
871                                .send(super::OutputLine {
872                                    text: line,
873                                    source: super::OutputSource::Local,
874                                })
875                                .await
876                                .is_err()
877                            {
878                                break;
879                            }
880                        }
881                    });
882                }
883                if let Some(mut stderr) = stderr_reader {
884                    let tx = output_tx.clone();
885                    tokio::spawn(async move {
886                        while let Ok(Some(line)) = stderr.next_line().await {
887                            if tx
888                                .send(super::OutputLine {
889                                    text: line,
890                                    source: super::OutputSource::Local,
891                                })
892                                .await
893                                .is_err()
894                            {
895                                break;
896                            }
897                        }
898                    });
899                }
900                // Drop the last sender so the channel closes when all readers
901                // finish. The relay holds its own clone, so a sink's reports
902                // still have somewhere to go after these end.
903                drop(output_tx);
904            }
905            let log_store = Arc::clone(&LOG_STORE);
906            let log_format = opts
907                .log_format
908                .clone()
909                .unwrap_or_else(|| crate::settings::settings().logs.log_format.clone());
910            let parse_line = move |line: &str| crate::log_parse::parse(line, &log_format);
911
912            const LOG_BATCH_SIZE: usize = 100;
913            const LOG_FLUSH_INTERVAL: Duration = Duration::from_millis(100);
914            let mut log_buffer: Vec<crate::log_parse::ParsedLog> =
915                Vec::with_capacity(LOG_BATCH_SIZE);
916            let mut log_flush_interval = tokio::time::interval(LOG_FLUSH_INTERVAL);
917            log_flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
918
919            let flush_logs =
920                |buffer: &mut Vec<crate::log_parse::ParsedLog>| -> Option<tokio::task::JoinHandle<()>> {
921                    if buffer.is_empty() {
922                        return None;
923                }
924                let store = Arc::clone(&log_store);
925                let id = id.clone();
926                let batch = std::mem::take(buffer);
927                Some(tokio::task::spawn_blocking(move || {
928                    if let Err(e) = store.append_structured_batch(&id, &batch) {
929                        error!("Failed to write batch to log for daemon {id}: {e}");
930                    }
931                }))
932            };
933
934            // SQLite WAL mode provides automatic durability; no explicit flush needed.
935
936            // Setup readiness checking
937            let mut ready_notified = false;
938            let mut ready_tx = ready_tx;
939            let ready_pattern = ready_output
940                .as_ref()
941                .and_then(|o| get_or_compile_regex(&o.pattern));
942            // Track whether we've already spawned the active_port detection task
943            let mut active_port_spawned = false;
944
945            // Validate on_output config early; discard the hook on any error so
946            // a bad regex does not silently fall through to the (None, None) => true
947            // match arm and fire on every line.
948            let on_output_hook = match on_output_hook {
949                Some(ref hook) => match hook.validate(id.name()) {
950                    Ok(()) => on_output_hook,
951                    Err(e) => {
952                        error!("{e}");
953                        None
954                    }
955                },
956                None => None,
957            };
958
959            // Compile the regex pattern after validation so we only attempt this
960            // when the hook is known-good (validate() already checked the syntax).
961            let on_output_pattern: Option<regex::Regex> = on_output_hook
962                .as_ref()
963                .and_then(|h| h.regex.as_deref().and_then(get_or_compile_regex));
964            let on_output_debounce = on_output_hook
965                .as_ref()
966                .map(|h| h.debounce_duration())
967                .unwrap_or(Duration::from_millis(1000));
968            // Last time the on_output hook fired; None means it has never fired.
969            let mut on_output_last_fired: Option<std::time::Instant> = None;
970
971            let mut delay_timer =
972                ready_delay.map(|secs| Box::pin(time::sleep(Duration::from_secs(secs))));
973
974            // Track exhaustion of timed checks
975            let mut http_exhausted = false;
976            let mut cmd_exhausted = false;
977            let mut port_exhausted = false;
978            let mut output_exhausted = false;
979
980            // Get settings for intervals
981            let s = settings();
982            let ready_check_interval = s.supervisor_ready_check_interval();
983            let http_client_timeout = s.supervisor_http_client_timeout();
984
985            // Setup output readiness check deadline
986            let mut output_deadline = ready_output
987                .as_ref()
988                .and_then(|o| o.timeout)
989                .map(|d| Box::pin(time::sleep(d)));
990
991            // Setup HTTP readiness check interval and deadline
992            let mut http_check_interval = ready_http
993                .as_ref()
994                .map(|_| tokio::time::interval(ready_check_interval));
995            let mut http_deadline = ready_http
996                .as_ref()
997                .and_then(|h| h.timeout)
998                .map(|d| Box::pin(time::sleep(d)));
999            let http_client = ready_http.as_ref().map(|_| {
1000                reqwest::Client::builder()
1001                    .timeout(http_client_timeout)
1002                    .build()
1003                    .unwrap_or_default()
1004            });
1005
1006            // Setup TCP port readiness check interval and deadline
1007            let mut port_check_interval =
1008                ready_port.map(|_| tokio::time::interval(ready_check_interval));
1009            let mut port_deadline = ready_port_config
1010                .as_ref()
1011                .and_then(|p| p.timeout)
1012                .map(|d| Box::pin(time::sleep(d)));
1013
1014            // Setup command readiness check state. Probes are spawned one at a time;
1015            // a non-zero result triggers a respawn delay, and a timeout stops the probe.
1016            let mut cmd_probe: Option<CmdProbe> = None;
1017            let mut cmd_respawn_delay: Option<_> = None;
1018            let mut cmd_deadline = ready_cmd
1019                .as_ref()
1020                .and_then(|c| c.timeout)
1021                .map(|d| Box::pin(time::sleep(d)));
1022            if let Some(ref cmd) = ready_cmd {
1023                cmd_probe = Some(spawn_cmd_probe(
1024                    &id,
1025                    &cmd.run,
1026                    daemon_dir.as_path(),
1027                    hook_retry_count,
1028                    readiness_daemon_env.as_ref(),
1029                    &readiness_resolved_ports,
1030                ));
1031            }
1032
1033            // Use a channel to communicate process exit status
1034            let (exit_tx, mut exit_rx) =
1035                tokio::sync::mpsc::channel::<std::io::Result<std::process::ExitStatus>>(1);
1036
1037            // Spawn a task to wait for process exit
1038            let child_pid = child.id().unwrap_or(0);
1039            tokio::spawn(async move {
1040                let result = child.wait().await;
1041                // On non-Linux Unix (e.g. macOS) the zombie reaper may win the
1042                // race and consume the exit status via waitpid(None, WNOHANG)
1043                // before Tokio's child.wait() gets to it. When that happens,
1044                // Tokio returns an ECHILD io::Error. We recover by checking
1045                // REAPED_STATUSES for the stashed exit code.
1046                //
1047                // On Linux this is unnecessary because the reaper uses
1048                // waitid(WNOWAIT) to peek before reaping, which avoids the
1049                // race entirely.
1050                #[cfg(all(unix, not(target_os = "linux")))]
1051                let result = match &result {
1052                    Err(e) if e.raw_os_error() == Some(nix::libc::ECHILD) => {
1053                        if let Some(code) = super::REAPED_STATUSES.lock().await.remove(&child_pid) {
1054                            warn!(
1055                                "daemon pid {child_pid} wait() got ECHILD; \
1056                                 recovered exit code {code} from zombie reaper"
1057                            );
1058                            // Synthesize an ExitStatus from the stashed code.
1059                            // On Unix we can use `ExitStatus::from_raw()` with
1060                            // a wait-style status word (code << 8 for normal
1061                            // exit, or raw signal number for signal death).
1062                            use std::os::unix::process::ExitStatusExt;
1063                            if code >= 0 {
1064                                Ok(std::process::ExitStatus::from_raw(code << 8))
1065                            } else {
1066                                // Negative code means killed by signal (-sig)
1067                                Ok(std::process::ExitStatus::from_raw((-code) & 0x7f))
1068                            }
1069                        } else {
1070                            warn!(
1071                                "daemon pid {child_pid} wait() got ECHILD but no \
1072                                 stashed status found; reporting as error"
1073                            );
1074                            result
1075                        }
1076                    }
1077                    _ => result,
1078                };
1079                debug!("daemon pid {child_pid} wait() completed with result: {result:?}");
1080                let _ = exit_tx.send(result).await;
1081            });
1082
1083            #[allow(unused_assignments)]
1084            // Initial None is a safety net; loop only exits via exit_rx.recv() which sets it
1085            let mut exit_status = None;
1086
1087            // If there is no ready check of any kind and no delay, the daemon is
1088            // considered immediately ready and the active_port detection task would
1089            // never be triggered inside the select loop.  Kick it off right away so
1090            // that daemons without any readiness configuration still get their
1091            // active_port populated (needed for proxy routing).
1092            if has_port_config
1093                && ready_pattern.is_none()
1094                && ready_http.is_none()
1095                && ready_port.is_none()
1096                && ready_cmd.is_none()
1097                && delay_timer.is_none()
1098            {
1099                active_port_spawned = true;
1100                detect_and_store_active_port(id.clone(), daemon_pid);
1101            }
1102
1103            // Set when readiness checks exhaust. The group kill runs as a
1104            // separate task so this loop can exit and the post-loop drain
1105            // keeps consuming output — children logging during SIGTERM
1106            // cleanup would otherwise block on a full pipe and never exit.
1107            // The ready failure is only sent once the kill task completes,
1108            // so the retry loop cannot respawn into the dying group.
1109            let mut ready_fail_kill: Option<tokio::task::JoinHandle<()>> = None;
1110
1111            loop {
1112                // biased: evaluate in exit → output → delay order so that
1113                // process exit pre-empts both buffered output and the delay
1114                // timer, preventing a dead daemon from being marked ready.
1115                select! {
1116                    biased;
1117                    Some(result) = exit_rx.recv() => {
1118                        // Process exited - save exit status and notify if not ready yet
1119                        exit_status = Some(result);
1120                        debug!("daemon {id} process exited, exit_status: {exit_status:?}");
1121                        if !ready_notified {
1122                            if let Some(tx) = ready_tx.take() {
1123                                // Check if process exited successfully
1124                                let is_success = exit_status.as_ref()
1125                                    .and_then(|r| r.as_ref().ok())
1126                                    .map(|s| s.success())
1127                                    .unwrap_or(false);
1128
1129                                if is_success {
1130                                    debug!("daemon {id} exited successfully before ready check, sending success notification");
1131                                    let _ = tx.send(Ok(()));
1132                                } else {
1133                                    let exit_code = exit_status.as_ref()
1134                                        .and_then(|r| r.as_ref().ok())
1135                                        .and_then(|s| s.code());
1136                                    debug!("daemon {id} exited with failure before ready check, sending failure notification with exit_code: {exit_code:?}");
1137                                    let _ = tx.send(Err(exit_code));
1138                                }
1139                            }
1140                        } else {
1141                            debug!("daemon {id} was already marked ready, not sending notification");
1142                        }
1143                        break;
1144                    },
1145                    Some(super::OutputLine { text: line, source }) = output_rx.recv() => {
1146                        // A line relayed by a sink is already in the store —
1147                        // the sink wrote and flushed it before reporting it —
1148                        // so it arrives here only to be acted on.
1149                        if matches!(source, super::OutputSource::Local) {
1150                            let parsed = parse_line(&line);
1151                            log_buffer.push(parsed);
1152                            if log_buffer.len() >= LOG_BATCH_SIZE {
1153                                let _ = flush_logs(&mut log_buffer);
1154                            }
1155                        }
1156                        trace!("output: {id} {line}");
1157
1158                        // Strip ANSI for pattern matching so user-written patterns
1159                        // work regardless of whether the process emits color codes.
1160                        let line_clean = console::strip_ansi_codes(&line).to_string();
1161
1162                        // Check if output matches ready pattern
1163                        if !ready_notified
1164                            && !output_exhausted
1165                            && let Some(ref pattern) = ready_pattern
1166                            && pattern.is_match(&line_clean)
1167                        {
1168                            // Flush buffered logs synchronously before signalling
1169                            // readiness, so collect_startup_logs sees the line
1170                            // that triggered the match (and any co-buffered lines)
1171                            // in SQLite.
1172                            if let Some(handle) = flush_logs(&mut log_buffer) {
1173                                let _ = handle.await;
1174                            }
1175                            info!("daemon {id} ready: output matched pattern");
1176                            ready_notified = true;
1177                            if let Some(tx) = ready_tx.take() {
1178                                let _ = tx.send(Ok(()));
1179                            }
1180                            fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1181                            stop_cmd_probe_state(&mut cmd_probe);
1182                            http_deadline = None;
1183                            cmd_deadline = None;
1184                            port_deadline = None;
1185                            output_deadline = None;
1186                            if !active_port_spawned && has_port_config {
1187                                active_port_spawned = true;
1188                                detect_and_store_active_port(id.clone(), daemon_pid);
1189                            }
1190                        }
1191
1192                        // Check on_output hook. A sink has already applied the
1193                        // filter, and says so per line: a line reported only
1194                        // because it announced readiness must not fire a hook
1195                        // that filters for something else.
1196                        if let Some(ref hook) = on_output_hook {
1197                            let matched = match source {
1198                                super::OutputSource::Sink { fires_hook } => fires_hook,
1199                                super::OutputSource::Local => match (&hook.filter, &on_output_pattern) {
1200                                    (Some(substr), _) => line_clean.contains(substr.as_str()),
1201                                    (None, Some(re)) => re.is_match(&line_clean),
1202                                    (None, None) => true,
1203                                },
1204                            };
1205                            if matched {
1206                                // The debounce is applied here as well as in the
1207                                // sink. A replacement sink starts with a fresh
1208                                // clock, and would otherwise let the hook fire
1209                                // twice inside one configured window.
1210                                let now = std::time::Instant::now();
1211                                let elapsed = on_output_last_fired.map(|t| now.duration_since(t));
1212                                if elapsed.is_none_or(|e| e >= on_output_debounce) {
1213                                    on_output_last_fired = Some(now);
1214                                    hooks::fire_output_hook(id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook.run.clone(), line_clean.clone()).await;
1215                                }
1216                            }
1217                        }
1218                        // Yield briefly so that the output readiness deadline can be
1219                        // evaluated even when output is produced continuously.
1220                        tokio::task::yield_now().await;
1221                    }
1222                    _ = async {
1223                        if let Some(ref mut deadline) = http_deadline {
1224                            deadline.await;
1225                        } else {
1226                            std::future::pending::<()>().await;
1227                        }
1228                    }, if !ready_notified && ready_http.is_some() => {
1229                        http_exhausted = true;
1230                        http_deadline = None;
1231                        http_check_interval = None;
1232                        warn!("daemon {id}: HTTP readiness check timed out");
1233                        let any_remaining = any_ready_check_remaining(
1234                            ready_output.as_ref(),
1235                            output_exhausted,
1236                            ready_port_config.as_ref(),
1237                            port_exhausted,
1238                            ready_http.as_ref(),
1239                            http_exhausted,
1240                            ready_cmd.as_ref(),
1241                            cmd_exhausted,
1242                        );
1243                        if !any_remaining {
1244                            error!("daemon {id}: all readiness checks exhausted, failing");
1245                            stop_cmd_probe_state(&mut cmd_probe);
1246                            ready_fail_kill = Some(spawn_ready_fail_kill(
1247                                id.clone(),
1248                                daemon_pid,
1249                                opts.stop_signal.unwrap_or_default(),
1250                            ));
1251                            break;
1252                        }
1253                    }
1254                    _ = async {
1255                        if let Some(ref mut deadline) = output_deadline {
1256                            deadline.await;
1257                        } else {
1258                            std::future::pending::<()>().await;
1259                        }
1260                    }, if !ready_notified && ready_output.is_some() => {
1261                        output_exhausted = true;
1262                        output_deadline = None;
1263                        warn!("daemon {id}: output readiness check timed out");
1264                        let any_remaining = any_ready_check_remaining(
1265                            ready_output.as_ref(),
1266                            output_exhausted,
1267                            ready_port_config.as_ref(),
1268                            port_exhausted,
1269                            ready_http.as_ref(),
1270                            http_exhausted,
1271                            ready_cmd.as_ref(),
1272                            cmd_exhausted,
1273                        );
1274                        if !any_remaining {
1275                            error!("daemon {id}: all readiness checks exhausted, failing");
1276                            stop_cmd_probe_state(&mut cmd_probe);
1277                            ready_fail_kill = Some(spawn_ready_fail_kill(
1278                                id.clone(),
1279                                daemon_pid,
1280                                opts.stop_signal.unwrap_or_default(),
1281                            ));
1282                            break;
1283                        }
1284                    }
1285                    _ = async {
1286                        if let Some(ref mut interval) = http_check_interval {
1287                            interval.tick().await;
1288                        } else {
1289                            std::future::pending::<()>().await;
1290                        }
1291                    }, if !ready_notified && ready_http.is_some() && !http_exhausted => {
1292                        if let (Some(http), Some(client)) = (&ready_http, &http_client) {
1293                            match client.get(&http.url).send().await {
1294                                Ok(response) if http.accepts_status(response.status().as_u16()) => {
1295                                    info!("daemon {id} ready: HTTP check passed (status {})", response.status());
1296                                    ready_notified = true;
1297                                    if let Some(tx) = ready_tx.take() {
1298                                        let _ = tx.send(Ok(()));
1299                                    }
1300                                    fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1301                                    http_check_interval = None;
1302                                    http_deadline = None;
1303                                    stop_cmd_probe_state(&mut cmd_probe);
1304                                    cmd_deadline = None;
1305                                    port_deadline = None;
1306                                    output_deadline = None;
1307                                    if !active_port_spawned && has_port_config {
1308                                        active_port_spawned = true;
1309                                        detect_and_store_active_port(id.clone(), daemon_pid);
1310                                    }
1311                                }
1312                                Ok(response) => {
1313                                    trace!("daemon {id} HTTP check: status {} (not ready)", response.status());
1314                                }
1315                                Err(e) => {
1316                                    trace!("daemon {id} HTTP check failed: {e}");
1317                                }
1318                            }
1319                        }
1320                    }
1321                    _ = async {
1322                        if let Some(ref mut deadline) = port_deadline {
1323                            deadline.await;
1324                        } else {
1325                            std::future::pending::<()>().await;
1326                        }
1327                    }, if !ready_notified && ready_port.is_some() => {
1328                        port_exhausted = true;
1329                        port_deadline = None;
1330                        port_check_interval = None;
1331                        warn!("daemon {id}: TCP port readiness check timed out");
1332                        let any_remaining = any_ready_check_remaining(
1333                            ready_output.as_ref(),
1334                            output_exhausted,
1335                            ready_port_config.as_ref(),
1336                            port_exhausted,
1337                            ready_http.as_ref(),
1338                            http_exhausted,
1339                            ready_cmd.as_ref(),
1340                            cmd_exhausted,
1341                        );
1342                        if !any_remaining {
1343                            error!("daemon {id}: all readiness checks exhausted, failing");
1344                            stop_cmd_probe_state(&mut cmd_probe);
1345                            ready_fail_kill = Some(spawn_ready_fail_kill(
1346                                id.clone(),
1347                                daemon_pid,
1348                                opts.stop_signal.unwrap_or_default(),
1349                            ));
1350                            break;
1351                        }
1352                    }
1353                    _ = async {
1354                        if let Some(ref mut interval) = port_check_interval {
1355                            interval.tick().await;
1356                        } else {
1357                            std::future::pending::<()>().await;
1358                        }
1359                    }, if !ready_notified && ready_port.is_some() && !port_exhausted => {
1360                        if let Some(port) = ready_port {
1361                            match tokio::net::TcpStream::connect(("127.0.0.1", port)).await {
1362                                Ok(_) => {
1363                                    info!("daemon {id} ready: TCP port {port} is listening");
1364                                    ready_notified = true;
1365                                    if let Some(tx) = ready_tx.take() {
1366                                        let _ = tx.send(Ok(()));
1367                                    }
1368                                    fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1369                                    // Stop checking once ready
1370                                    port_check_interval = None;
1371                                    port_deadline = None;
1372                                    stop_cmd_probe_state(&mut cmd_probe);
1373                                    http_deadline = None;
1374                                    cmd_deadline = None;
1375                                    output_deadline = None;
1376                                    if !active_port_spawned && has_port_config {
1377                                        active_port_spawned = true;
1378                                        // ready_port check just TCP-connected to this
1379                                        // port, so it is definitely listening. If it
1380                                        // matches the configured expected port, write
1381                                        // active_port directly instead of spawning
1382                                        // detect_and_store_active_port, which sleeps
1383                                        // 500 ms then relies on listeners::get_all()
1384                                        // + process-tree traversal — unreliable on
1385                                        // Windows where Git Bash PID mapping can
1386                                        // break descendant lookups.
1387                                        if expected_port == Some(port) {
1388                                            let mut state_file =
1389                                                SUPERVISOR.state_file.lock().await;
1390                                            if let Some(d) = state_file.daemons.get(&id)
1391                                                && d.pid == Some(daemon_pid)
1392                                            {
1393                                                state_file.set_active_port(&id, port);
1394                                            }
1395                                        } else {
1396                                            detect_and_store_active_port(
1397                                                id.clone(),
1398                                                daemon_pid,
1399                                            );
1400                                        }
1401                                    }
1402                                }
1403                                Err(_) => {
1404                                    trace!("daemon {id} port check: port {port} not listening yet");
1405                                }
1406                            }
1407                        }
1408                    }
1409                    _ = async {
1410                        if let Some(ref mut delay) = cmd_respawn_delay {
1411                            delay.await;
1412                        } else {
1413                            std::future::pending::<()>().await;
1414                        }
1415                    }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted && cmd_probe.is_none() => {
1416                        if let Some(ref cmd) = ready_cmd {
1417                            cmd_probe = Some(spawn_cmd_probe(
1418                                &id,
1419                                &cmd.run,
1420                                daemon_dir.as_path(),
1421                                hook_retry_count,
1422                                readiness_daemon_env.as_ref(),
1423                                &readiness_resolved_ports,
1424                            ));
1425                        }
1426                        cmd_respawn_delay = None;
1427                    }
1428                    result = async {
1429                        if let Some(probe) = cmd_probe.as_mut() {
1430                            std::pin::Pin::new(&mut probe.result_rx).await
1431                        } else {
1432                            std::future::pending::<Result<Result<std::process::ExitStatus, std::io::Error>, tokio::sync::oneshot::error::RecvError>>().await
1433                        }
1434                    }, if !ready_notified && ready_cmd.is_some() && !cmd_exhausted => {
1435                        // The probe task has finished; remove the handle so it is not
1436                        // cancelled or reused. This must happen only after this branch
1437                        // actually wins the select, not while constructing the future.
1438                        let _ = cmd_probe.take();
1439                        match result {
1440                            Ok(Ok(status)) if status.success() => {
1441                                info!("daemon {id} ready: readiness command succeeded");
1442                                ready_notified = true;
1443                                if let Some(tx) = ready_tx.take() {
1444                                    let _ = tx.send(Ok(()));
1445                                }
1446                                fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1447                                cmd_respawn_delay = None;
1448                                cmd_deadline = None;
1449                                http_deadline = None;
1450                                port_deadline = None;
1451                                output_deadline = None;
1452                                if !active_port_spawned && has_port_config {
1453                                    active_port_spawned = true;
1454                                    detect_and_store_active_port(id.clone(), daemon_pid);
1455                                }
1456                            }
1457                            Ok(Ok(_)) | Ok(Err(_)) | Err(_) => {
1458                                trace!("daemon {id} cmd check: command not ready, will respawn");
1459                                cmd_respawn_delay = Some(Box::pin(time::sleep(ready_check_interval)));
1460                            }
1461                        }
1462                    }
1463                    _ = async {
1464                        if let Some(ref mut deadline) = cmd_deadline {
1465                            deadline.await;
1466                        } else {
1467                            std::future::pending::<()>().await;
1468                        }
1469                    }, if !ready_notified && ready_cmd.is_some() => {
1470                        cmd_exhausted = true;
1471                        cmd_deadline = None;
1472                        stop_cmd_probe_state(&mut cmd_probe);
1473                        cmd_respawn_delay = None;
1474                        warn!("daemon {id}: command readiness check timed out");
1475                        let any_remaining = any_ready_check_remaining(
1476                            ready_output.as_ref(),
1477                            output_exhausted,
1478                            ready_port_config.as_ref(),
1479                            port_exhausted,
1480                            ready_http.as_ref(),
1481                            http_exhausted,
1482                            ready_cmd.as_ref(),
1483                            cmd_exhausted,
1484                        );
1485                        if !any_remaining {
1486                            error!("daemon {id}: all readiness checks exhausted, failing");
1487                            ready_fail_kill = Some(spawn_ready_fail_kill(
1488                                id.clone(),
1489                                daemon_pid,
1490                                opts.stop_signal.unwrap_or_default(),
1491                            ));
1492                            break;
1493                        }
1494                    }
1495                    _ = async {
1496                        if let Some(ref mut timer) = delay_timer {
1497                            timer.await;
1498                        } else {
1499                            std::future::pending::<()>().await;
1500                        }
1501                    } => {
1502                        if !ready_notified && ready_pattern.is_none() && ready_http.is_none() && ready_port.is_none() && ready_cmd.is_none() {
1503                            // Check if the process already exited or is exiting before
1504                            // declaring it ready. On Windows, sleep(0) fires before
1505                            // child.wait() detects the exit, causing pitchfork start to
1506                            // return success for a daemon that already failed.
1507                            if exit_status.is_some() {
1508                                debug!("daemon {id} exited during ready_delay, not marking as ready");
1509                            } else {
1510                                // Force-refresh sysinfo for this PID before checking.
1511                                // On Windows, the cached process list may be stale.
1512                                PROCS.refresh_pids(&[daemon_pid]);
1513                                if !PROCS.is_running(daemon_pid) {
1514                                    debug!("daemon {id} pid {daemon_pid} not running during ready_delay, deferring to exit handler");
1515                                } else {
1516                                    info!("daemon {id} ready: delay elapsed");
1517                                    ready_notified = true;
1518                                    if let Some(tx) = ready_tx.take() {
1519                                        let _ = tx.send(Ok(()));
1520                                    }
1521                                    fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
1522                                }
1523                            }
1524                            // Clear all deadlines — no other checks are configured
1525                            // when delay fires as readiness, but clear defensively.
1526                            output_deadline = None;
1527                            http_deadline = None;
1528                            cmd_deadline = None;
1529                            port_deadline = None;
1530                            stop_cmd_probe_state(&mut cmd_probe);
1531                        }
1532                        // Disable timer after it fires
1533                        delay_timer = None;
1534                        if !active_port_spawned && has_port_config {
1535                            active_port_spawned = true;
1536                            detect_and_store_active_port(id.clone(), daemon_pid);
1537                        }
1538                    }
1539                    _ = log_flush_interval.tick() => {
1540                        let _ = flush_logs(&mut log_buffer);
1541                    }
1542                }
1543            }
1544
1545            // Snapshot the daemon state BEFORE draining output.
1546            //
1547            // The drain can take up to 5s (e.g. when child processes keep the
1548            // stdout pipe open). During that time, a subsequent start() call
1549            // (e.g. from `pitchfork restart`) can upsert the daemon with a new
1550            // PID and Running status. If we only checked state AFTER the drain,
1551            // the monitoring task would see d.pid != Some(old_pid) && !is_stopped()
1552            // && !is_stopping() and return early without firing on_stop/on_exit
1553            // hooks.
1554            //
1555            // By snapshotting is_stopping before the drain, we preserve the
1556            // knowledge that stop() was called, so hooks fire correctly even
1557            // if start() has since changed the state.
1558            let pre_drain_daemon = SUPERVISOR.get_daemon(&id).await;
1559            let pre_drain_is_stopping = pre_drain_daemon
1560                .as_ref()
1561                .is_some_and(|d| d.status.is_stopped() || d.status.is_stopping());
1562
1563            // Drain any in-flight output lines that were still in the mpsc
1564            // channel or the OS pipe buffer when the child exited. Without
1565            // this, trailing log lines from short-lived daemons get dropped.
1566            // The reader tasks drop their senders on EOF, so recv() returns
1567            // None when all data has been consumed. A total deadline of 5 s
1568            // guards against a stuck reader (e.g. PTY master FD not closing)
1569            // while ensuring drain doesn't block post-exit cleanup indefinitely.
1570            //
1571            // Stop accepting relayed output first: the relay holds a sender of
1572            // its own, so leaving it registered would keep the channel open and
1573            // make every drain wait out the whole deadline. Readiness is moot
1574            // now anyway — the process has exited.
1575            drop(output_relay);
1576            let drain_deadline = tokio::time::Instant::now() + Duration::from_secs(5);
1577            loop {
1578                let now = tokio::time::Instant::now();
1579                if now >= drain_deadline {
1580                    break;
1581                }
1582                let Ok(Some(line)) =
1583                    tokio::time::timeout(drain_deadline - now, output_rx.recv()).await
1584                else {
1585                    break;
1586                };
1587                // Sink-relayed lines are already stored; see the select loop.
1588                if matches!(line.source, super::OutputSource::Local) {
1589                    log_buffer.push(parse_line(&line.text));
1590                }
1591            }
1592            // Flush any remaining log lines (including drained) before the process exits.
1593            // Await the flush to guarantee all buffered logs are persisted before cleanup.
1594            if let Some(handle) = flush_logs(&mut log_buffer) {
1595                let _ = handle.await;
1596            }
1597
1598            // Clear active_port since the process is no longer running
1599            {
1600                let mut state_file = SUPERVISOR.state_file.lock().await;
1601                state_file.clear_active_port(&id);
1602            }
1603
1604            // Get the final exit status
1605            let exit_status = if let Some(status) = exit_status {
1606                status
1607            } else {
1608                // Streams closed but process hasn't exited yet, wait for it
1609                match exit_rx.recv().await {
1610                    Some(status) => status,
1611                    None => {
1612                        warn!("daemon {id} exit channel closed without receiving status");
1613                        Err(std::io::Error::other("exit channel closed"))
1614                    }
1615                }
1616            };
1617
1618            // If the loop exited via readiness exhaustion, wait for the group
1619            // kill to finish before reporting the failure so the retry loop
1620            // (or a waiting client) cannot start a replacement while the old
1621            // process group is still terminating.
1622            if let Some(kill) = ready_fail_kill {
1623                let _ = kill.await;
1624                if let Some(tx) = ready_tx.take() {
1625                    let _ = tx.send(Err(Some(124)));
1626                }
1627            }
1628
1629            let current_daemon = SUPERVISOR.get_daemon(&id).await;
1630
1631            // Signal that this monitoring task is processing its exit path.
1632            // The RAII guard will decrement the counter and notify close()
1633            // when the task finishes (including all fire_hook registrations),
1634            // regardless of which return path is taken.
1635            SUPERVISOR
1636                .active_monitors
1637                .fetch_add(1, atomic::Ordering::Release);
1638            struct MonitorGuard;
1639            impl Drop for MonitorGuard {
1640                fn drop(&mut self) {
1641                    SUPERVISOR
1642                        .active_monitors
1643                        .fetch_sub(1, atomic::Ordering::Release);
1644                    SUPERVISOR.monitor_done.notify_waiters();
1645                }
1646            }
1647            let _monitor_guard = MonitorGuard;
1648            // Check if this monitoring task is for the current daemon process.
1649            // If the daemon was intentionally stopped (pre_drain_is_stopping),
1650            // skip this check — we must still fire on_stop/on_exit hooks even
1651            // if start() has since changed the PID and status.
1652            if !pre_drain_is_stopping
1653                && (current_daemon.is_none()
1654                    || current_daemon.as_ref().is_some_and(|d| {
1655                        d.pid != Some(pid) && !d.status.is_stopped() && !d.status.is_stopping()
1656                    }))
1657            {
1658                // Another process has taken over, don't update status
1659                return;
1660            }
1661            // Capture the intentional-stop flag. Combine pre-drain and
1662            // post-drain state to handle both race orders:
1663            //  - stop() set Stopping before drain → pre_drain_is_stopping
1664            //  - stop() set Stopped during drain → current_daemon.is_stopped()
1665            let already_stopped = current_daemon
1666                .as_ref()
1667                .is_some_and(|d| d.status.is_stopped());
1668            let is_stopping = already_stopped
1669                || pre_drain_is_stopping
1670                || current_daemon
1671                    .as_ref()
1672                    .is_some_and(|d| d.status.is_stopping());
1673
1674            // --- Phase 1: Determine exit_code, exit_reason, and update daemon state ---
1675            let (exit_code, exit_reason) = match (&exit_status, is_stopping) {
1676                (Ok(status), true) => {
1677                    // Intentional stop (by pitchfork). status.code() returns None
1678                    // on Unix when killed by signal (e.g. SIGTERM); use -1 to
1679                    // distinguish from a clean exit code 0.
1680                    (status.code().unwrap_or(-1), "stop")
1681                }
1682                (Ok(status), false) if status.success() => (status.code().unwrap_or(-1), "exit"),
1683                (Ok(status), false) => (status.code().unwrap_or(-1), "fail"),
1684                (Err(_), true) => {
1685                    // child.wait() error while stopping (e.g. sysinfo reaped the process)
1686                    (-1, "stop")
1687                }
1688                (Err(_), false) => (-1, "fail"),
1689            };
1690
1691            // Update daemon state unless stop() already did it (won the race),
1692            // OR the daemon was intentionally stopped before the drain
1693            // (pre_drain_is_stopping). In the latter case, start() may have
1694            // upserted Running during the 5s drain, and we must NOT overwrite
1695            // it with Stopped — that would undo the restart.
1696            if !already_stopped && !pre_drain_is_stopping {
1697                if let Ok(status) = &exit_status {
1698                    info!("daemon {id} exited with status {status}");
1699                }
1700                let (new_status, last_exit_success) = match exit_reason {
1701                    "stop" | "exit" => (
1702                        DaemonStatus::Stopped,
1703                        exit_status.as_ref().map(|s| s.success()).unwrap_or(true),
1704                    ),
1705                    _ => (DaemonStatus::Errored(exit_code), false),
1706                };
1707                // Revalidate ownership inside the same state-lock section that
1708                // performs the write. The snapshot above was taken without
1709                // holding the lock, so a restart running on another thread can
1710                // install a successor in between; overwriting its record would
1711                // clear a live daemon's PID and undo the restart.
1712                if !SUPERVISOR
1713                    .finalize_monitored_exit(
1714                        &id,
1715                        pid,
1716                        monitor_token,
1717                        new_status,
1718                        Some(last_exit_success),
1719                    )
1720                    .await
1721                {
1722                    debug!("daemon {id} exit state was not written; a successor owns the record");
1723                }
1724            }
1725
1726            // --- Phase 2: Fire hooks ---
1727            let hook_extra_env = vec![
1728                ("PITCHFORK_EXIT_CODE".to_string(), exit_code.to_string()),
1729                ("PITCHFORK_EXIT_REASON".to_string(), exit_reason.to_string()),
1730            ];
1731
1732            // Determine which hooks to fire based on exit reason
1733            let hooks_to_fire: Vec<HookType> = match exit_reason {
1734                "stop" => vec![HookType::OnStop, HookType::OnExit],
1735                "exit" => vec![HookType::OnExit],
1736                // "fail": fire on_fail + on_exit only when retries are exhausted
1737                _ if hook_retry_count >= hook_retry.count() => {
1738                    vec![HookType::OnFail, HookType::OnExit]
1739                }
1740                _ => vec![],
1741            };
1742
1743            for hook_type in hooks_to_fire {
1744                fire_hook(
1745                    hook_type,
1746                    id.clone(),
1747                    daemon_dir.clone(),
1748                    hook_retry_count,
1749                    hook_daemon_env.clone(),
1750                    hook_extra_env.clone(),
1751                )
1752                .await;
1753            }
1754        });
1755
1756        // If wait_ready is true, wait for readiness notification
1757        if let Some(ready_rx) = ready_rx {
1758            match ready_rx.await {
1759                Ok(Ok(())) => {
1760                    info!("daemon {id} is ready");
1761                    Ok(IpcResponse::DaemonReady { daemon })
1762                }
1763                Ok(Err(exit_code)) => {
1764                    error!("daemon {id} failed before becoming ready");
1765                    // The caller reports this by querying the log store for
1766                    // what the daemon printed, so wait for the sink's final
1767                    // write first. The in-process path got this ordering by
1768                    // flushing synchronously before signalling.
1769                    //
1770                    // Only on the attempt that gives up: `run` retries inline,
1771                    // and waiting after every attempt would both delay the
1772                    // backoff and widen the window in which the daemon looks
1773                    // errored and idle — long enough for the background retry
1774                    // checker to start an attempt of its own alongside it.
1775                    let last_attempt = opts.retry_count >= opts.retry.count();
1776                    if using_sink && last_attempt {
1777                        super::log_sink::wait_for_output(id, spawn_time, SINK_OUTPUT_TIMEOUT).await;
1778                    }
1779                    Ok(IpcResponse::DaemonFailedWithCode { exit_code })
1780                }
1781                Err(_) => {
1782                    error!("readiness channel closed unexpectedly for daemon {id}");
1783                    Ok(IpcResponse::DaemonStart { daemon })
1784                }
1785            }
1786        } else {
1787            Ok(IpcResponse::DaemonStart { daemon })
1788        }
1789    }
1790
1791    /// Stop a running daemon
1792    pub async fn stop(&self, id: &DaemonId) -> Result<IpcResponse> {
1793        // Hold the daemon's stop lock for the whole stop (including the
1794        // whole-group termination wait) so starts and concurrent stops of the
1795        // same daemon serialize against it instead of racing the Stopping window.
1796        let lock = self.stop_lock(id).await;
1797        let _guard = lock.lock().await;
1798        self.stop_locked(id).await
1799    }
1800
1801    /// Stop implementation. Caller must hold the daemon's stop lock.
1802    async fn stop_locked(&self, id: &DaemonId) -> Result<IpcResponse> {
1803        let pitchfork_id = DaemonId::pitchfork();
1804        if *id == pitchfork_id {
1805            return Ok(IpcResponse::Error(
1806                "Cannot stop supervisor via stop command".into(),
1807            ));
1808        }
1809        info!("stopping daemon: {id}");
1810        if let Some(daemon) = self.get_daemon(id).await {
1811            trace!("daemon to stop: {daemon}");
1812            if let Some(pid) = daemon.pid {
1813                trace!("killing pid: {pid}");
1814                if PROCS.is_running(pid) {
1815                    // Something is alive on that PID, but the kill below signals
1816                    // the entire process group: if the PID was recycled while
1817                    // this record sat unsupervised, that group belongs to an
1818                    // unrelated process tree. The daemon itself is gone either
1819                    // way, so report it as not running and clear the record.
1820                    if !super::signalling_pid_is_authorized(
1821                        daemon.start_time,
1822                        PROCS.start_time(pid),
1823                    ) {
1824                        warn!(
1825                            "pid {pid} recorded for daemon {id} belongs to another process now; not signalling it"
1826                        );
1827                        self.upsert_daemon(
1828                            UpsertDaemonOpts::builder(id.clone())
1829                                .set(|o| {
1830                                    o.pid = None;
1831                                    o.status = DaemonStatus::Stopped;
1832                                })
1833                                .build(),
1834                        )
1835                        .await?;
1836                        return Ok(IpcResponse::DaemonWasNotRunning);
1837                    }
1838
1839                    // First set status to Stopping (preserve PID for monitoring task)
1840                    self.upsert_daemon(
1841                        UpsertDaemonOpts::builder(id.clone())
1842                            .set(|o| {
1843                                o.pid = Some(pid);
1844                                o.status = DaemonStatus::Stopping;
1845                            })
1846                            .build(),
1847                    )
1848                    .await?;
1849
1850                    // Kill the entire process group atomically (daemon PID == PGID
1851                    // because we called setsid() at spawn time)
1852                    let stop_cfg = daemon.stop_signal.unwrap_or_default();
1853                    let stop_signal: i32 = stop_cfg.signal.into();
1854                    if let Err(e) = PROCS
1855                        .kill_process_group_async(pid, stop_signal, stop_cfg.timeout)
1856                        .await
1857                    {
1858                        debug!("failed to kill pid {pid}: {e}");
1859                        // Check if the process group is actually gone despite the
1860                        // error. Checking only the leader here would mark the daemon
1861                        // Stopped while surviving group members (e.g. one stuck in
1862                        // uninterruptible sleep) are still alive — letting a restart
1863                        // collide with them.
1864                        if PROCS.process_group_alive(pid) {
1865                            // Group still has live members - set back to Running
1866                            debug!(
1867                                "failed to stop pid {pid}: process group still alive after kill"
1868                            );
1869                            self.upsert_daemon(
1870                                UpsertDaemonOpts::builder(id.clone())
1871                                    .set(|o| {
1872                                        o.pid = Some(pid); // Preserve PID to avoid orphaning the process
1873                                        o.status = DaemonStatus::Running;
1874                                    })
1875                                    .build(),
1876                            )
1877                            .await?;
1878                            return Ok(IpcResponse::DaemonStopFailed {
1879                                error: format!(
1880                                    "process group of {pid} still alive after kill attempt: {e}"
1881                                ),
1882                            });
1883                        }
1884                    }
1885
1886                    // Process successfully stopped
1887                    // Note: kill_process_group_async waits for the ENTIRE process
1888                    // group to exit (stop signal -> stop_timeout -> SIGKILL, then a
1889                    // bounded verification), so a replacement daemon can be started
1890                    // without colliding with a still-terminating instance. The only
1891                    // exception is a member stuck in uninterruptible sleep, which is
1892                    // logged with a warning.
1893                    self.upsert_daemon(
1894                        UpsertDaemonOpts::builder(id.clone())
1895                            .set(|o| {
1896                                o.pid = None;
1897                                o.status = DaemonStatus::Stopped;
1898                                o.last_exit_success = Some(true);
1899                            })
1900                            .build(),
1901                    )
1902                    .await?;
1903                } else {
1904                    debug!("pid {pid} not running, process may have exited unexpectedly");
1905                    // Process already dead — transition to Stopped so the
1906                    // retry checker sees a terminal state and stops
1907                    // scheduling new attempts. This is important for an
1908                    // explicit `pitchfork stop` on an Errored daemon: the
1909                    // user wants to abort retries.
1910                    self.upsert_daemon(
1911                        UpsertDaemonOpts::builder(id.clone())
1912                            .set(|o| {
1913                                o.pid = None;
1914                                o.status = DaemonStatus::Stopped;
1915                            })
1916                            .build(),
1917                    )
1918                    .await?;
1919                    return Ok(IpcResponse::DaemonWasNotRunning);
1920                }
1921                Ok(IpcResponse::Ok)
1922            } else {
1923                debug!("daemon {id} not running");
1924                Ok(IpcResponse::DaemonNotRunning)
1925            }
1926        } else {
1927            debug!("daemon {id} not found");
1928            Ok(IpcResponse::DaemonNotFound)
1929        }
1930    }
1931}
1932
1933#[cfg(unix)]
1934fn resolve_effective_run_identity(daemon_user: Option<&str>) -> Result<RunIdentity> {
1935    let s = settings();
1936    let settings_user = s.supervisor.user.trim();
1937    let daemon_user = daemon_user.map(str::trim).filter(|user| !user.is_empty());
1938    let settings_user = (!settings_user.is_empty()).then_some(settings_user);
1939    let configured = daemon_user.or(settings_user);
1940    let current_uid = nix::unistd::Uid::effective().as_raw();
1941    let current_gid = nix::unistd::Gid::effective().as_raw();
1942    resolve_run_identity(
1943        configured,
1944        current_uid,
1945        current_gid,
1946        std::env::var("SUDO_UID").ok().as_deref(),
1947        std::env::var("SUDO_GID").ok().as_deref(),
1948    )
1949}
1950
1951#[cfg(unix)]
1952fn resolve_run_identity(
1953    configured: Option<&str>,
1954    current_uid: u32,
1955    current_gid: u32,
1956    sudo_uid: Option<&str>,
1957    sudo_gid: Option<&str>,
1958) -> Result<RunIdentity> {
1959    let current_uid = nix::unistd::Uid::from_raw(current_uid);
1960    let current_gid = nix::unistd::Gid::from_raw(current_gid);
1961    if let Some(user) = configured {
1962        let identity = resolve_configured_user(user)?;
1963        ensure_can_use_identity(user, &identity, current_uid, current_gid)?;
1964        if identity.matches(current_uid, current_gid) {
1965            return Ok(RunIdentity::Inherit);
1966        }
1967        return Ok(identity);
1968    }
1969
1970    if current_uid.is_root()
1971        && let Some(identity) = resolve_sudo_identity(sudo_uid, sudo_gid)
1972    {
1973        return Ok(identity);
1974    }
1975
1976    Ok(RunIdentity::Inherit)
1977}
1978
1979#[cfg(unix)]
1980fn resolve_configured_user(user: &str) -> Result<RunIdentity> {
1981    if user.chars().all(|c| c.is_ascii_digit()) {
1982        let uid = user
1983            .parse::<u32>()
1984            .map_err(|e| miette::miette!("invalid run user UID '{}': {}", user, e))?;
1985        let user_record = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1986            .into_diagnostic()?
1987            .ok_or_else(|| miette::miette!("run user UID '{}' does not exist", user))?;
1988        return run_identity_from_user_record(user_record);
1989    }
1990
1991    let user_record = nix::unistd::User::from_name(user)
1992        .into_diagnostic()?
1993        .ok_or_else(|| miette::miette!("run user '{}' does not exist", user))?;
1994    run_identity_from_user_record(user_record)
1995}
1996
1997#[cfg(unix)]
1998fn run_identity_from_user_record(user: nix::unistd::User) -> Result<RunIdentity> {
1999    let username = CString::new(user.name)
2000        .map_err(|e| miette::miette!("run user name contains an interior nul byte: {}", e))?;
2001    Ok(RunIdentity::Switch {
2002        uid: user.uid,
2003        gid: user.gid,
2004        username: Some(username),
2005    })
2006}
2007
2008#[cfg(unix)]
2009fn run_identity_from_raw_ids(uid: u32, gid: u32, username: Option<CString>) -> RunIdentity {
2010    RunIdentity::Switch {
2011        uid: nix::unistd::Uid::from_raw(uid),
2012        gid: nix::unistd::Gid::from_raw(gid),
2013        username,
2014    }
2015}
2016
2017#[cfg(unix)]
2018fn resolve_sudo_identity(sudo_uid: Option<&str>, sudo_gid: Option<&str>) -> Option<RunIdentity> {
2019    let uid = sudo_uid?.parse::<u32>().ok()?;
2020    let gid = sudo_gid?.parse::<u32>().ok()?;
2021    let username = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
2022        .ok()
2023        .flatten()
2024        .and_then(|u| CString::new(u.name).ok());
2025    Some(run_identity_from_raw_ids(uid, gid, username))
2026}
2027
2028#[cfg(unix)]
2029fn ensure_can_use_identity(
2030    configured_user: &str,
2031    identity: &RunIdentity,
2032    current_uid: nix::unistd::Uid,
2033    current_gid: nix::unistd::Gid,
2034) -> Result<()> {
2035    let RunIdentity::Switch { uid, gid, .. } = identity else {
2036        return Ok(());
2037    };
2038    if *uid == current_uid && *gid == current_gid {
2039        return Ok(());
2040    }
2041    if current_uid.is_root() {
2042        return Ok(());
2043    }
2044    Err(miette::miette!(
2045        "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.",
2046        configured_user,
2047        current_uid.as_raw(),
2048        current_gid.as_raw(),
2049        uid.as_raw(),
2050        gid.as_raw()
2051    ))
2052}
2053
2054#[cfg(unix)]
2055fn apply_run_identity(identity: &RunIdentity) -> std::io::Result<()> {
2056    let RunIdentity::Switch { uid, gid, username } = identity else {
2057        return Ok(());
2058    };
2059    if let Some(username) = username {
2060        initgroups_for_user(username, *gid)?;
2061    } else {
2062        setgroups_to_primary(*gid)?;
2063    }
2064    nix::unistd::setgid(*gid).map_err(nix_to_io_error)?;
2065    nix::unistd::setuid(*uid).map_err(nix_to_io_error)?;
2066    Ok(())
2067}
2068
2069#[cfg(unix)]
2070impl RunIdentity {
2071    fn matches(&self, uid: nix::unistd::Uid, gid: nix::unistd::Gid) -> bool {
2072        matches!(self, RunIdentity::Switch { uid: u, gid: g, .. } if *u == uid && *g == gid)
2073    }
2074}
2075
2076#[cfg(unix)]
2077fn setgroups_to_primary(gid: nix::unistd::Gid) -> std::io::Result<()> {
2078    let groups = [gid.as_raw() as libc::gid_t];
2079    #[cfg(any(target_os = "linux", target_os = "android"))]
2080    let group_count = groups.len();
2081    #[cfg(not(any(target_os = "linux", target_os = "android")))]
2082    let group_count = groups.len() as libc::c_int;
2083    let rc = unsafe { libc::setgroups(group_count, groups.as_ptr()) };
2084    if rc == -1 {
2085        Err(std::io::Error::last_os_error())
2086    } else {
2087        Ok(())
2088    }
2089}
2090
2091#[cfg(unix)]
2092fn initgroups_for_user(username: &CString, gid: nix::unistd::Gid) -> std::io::Result<()> {
2093    let gid = gid.as_raw();
2094    #[cfg(any(
2095        target_os = "macos",
2096        target_os = "ios",
2097        target_os = "tvos",
2098        target_os = "watchos"
2099    ))]
2100    let base_gid = i32::try_from(gid)
2101        .map_err(|_| std::io::Error::other(format!("gid {gid} is out of range")))?;
2102
2103    #[cfg(not(any(
2104        target_os = "macos",
2105        target_os = "ios",
2106        target_os = "tvos",
2107        target_os = "watchos"
2108    )))]
2109    let base_gid = gid as libc::gid_t;
2110
2111    // SAFETY: `username` is a valid nul-terminated C string and `base_gid`
2112    // is derived from a resolved system account or sudo-provided gid.
2113    let rc = unsafe { libc::initgroups(username.as_ptr(), base_gid) };
2114    if rc == -1 {
2115        Err(std::io::Error::last_os_error())
2116    } else {
2117        Ok(())
2118    }
2119}
2120
2121#[cfg(unix)]
2122fn nix_to_io_error(err: nix::errno::Errno) -> std::io::Error {
2123    std::io::Error::from_raw_os_error(err as i32)
2124}
2125
2126/// Check if multiple ports are available and optionally auto-bump to find available ports.
2127///
2128/// All ports are bumped by the same offset to maintain relative port spacing.
2129/// Returns the resolved ports (either the original or bumped ones).
2130/// Returns an error if any port is in use and auto_bump is disabled,
2131/// or if no available ports can be found after max attempts.
2132async fn check_ports_available(
2133    expected_ports: &[u16],
2134    auto_bump: bool,
2135    max_attempts: u32,
2136) -> Result<Vec<u16>> {
2137    if expected_ports.is_empty() {
2138        return Ok(Vec::new());
2139    }
2140
2141    for bump_offset in 0..=max_attempts {
2142        // Use wrapping_add to handle overflow correctly - ports wrap around at 65535
2143        let candidate_ports: Vec<u16> = expected_ports
2144            .iter()
2145            .map(|&p| p.wrapping_add(bump_offset as u16))
2146            .collect();
2147
2148        // Check if all ports in this set are available
2149        let mut all_available = true;
2150        let mut conflicting_port = None;
2151
2152        for &port in &candidate_ports {
2153            // Port 0 is a special case - it requests an ephemeral port from the OS.
2154            // Skip the availability check for port 0 since binding to it always succeeds.
2155            if port == 0 {
2156                continue;
2157            }
2158
2159            // Use spawn_blocking to avoid blocking the async runtime during TCP bind checks.
2160            //
2161            // We check multiple addresses to avoid false-negatives caused by SO_REUSEADDR.
2162            // On macOS/BSD, Rust's TcpListener::bind sets SO_REUSEADDR by default, which
2163            // allows binding 0.0.0.0:port even when 127.0.0.1:port is already in use
2164            // (because 0.0.0.0 is technically a different address).  Most daemons bind
2165            // to localhost, so checking 127.0.0.1 is essential to detect real conflicts.
2166            // We also check [::1] to cover IPv6 loopback listeners.
2167            //
2168            // NOTE: This check has a time-of-check-to-time-of-use (TOCTOU) race condition.
2169            // Another process could grab the port between our check and the daemon actually
2170            // binding. This is inherent to the approach and acceptable for our use case
2171            // since we're primarily detecting conflicts with already-running daemons.
2172            if is_port_in_use(port).await {
2173                all_available = false;
2174                conflicting_port = Some(port);
2175                break;
2176            }
2177        }
2178
2179        if all_available {
2180            // Check for overflow (port wrapped around to 0 due to wrapping_add)
2181            // If any candidate port is 0 but the original expected port wasn't 0,
2182            // it means we've wrapped around and should stop
2183            if candidate_ports.contains(&0) && !expected_ports.contains(&0) {
2184                return Err(PortError::NoAvailablePort {
2185                    start_port: expected_ports[0],
2186                    attempts: bump_offset + 1,
2187                }
2188                .into());
2189            }
2190            if bump_offset > 0 {
2191                info!("ports {expected_ports:?} bumped by {bump_offset} to {candidate_ports:?}");
2192            }
2193            return Ok(candidate_ports);
2194        }
2195
2196        // Port is in use
2197        if bump_offset == 0
2198            && !auto_bump
2199            && let Some(port) = conflicting_port
2200        {
2201            let (pid, process) = identify_port_owner(port).await;
2202            return Err(PortError::InUse { port, process, pid }.into());
2203        }
2204    }
2205
2206    // No available ports found after max attempts
2207    Err(PortError::NoAvailablePort {
2208        start_port: expected_ports[0],
2209        attempts: max_attempts + 1,
2210    }
2211    .into())
2212}
2213
2214/// Check whether a port is currently in use by attempting to bind on multiple addresses.
2215///
2216/// Returns `true` when at least one bind attempt gets `AddrInUse`, meaning another
2217/// process is listening.  Other errors (e.g. `AddrNotAvailable` on an address family
2218/// the OS doesn't support) are ignored so they don't produce false positives.
2219async fn is_port_in_use(port: u16) -> bool {
2220    tokio::task::spawn_blocking(move || {
2221        for &addr in &["0.0.0.0", "127.0.0.1", "::1"] {
2222            match std::net::TcpListener::bind((addr, port)) {
2223                Ok(listener) => drop(listener),
2224                Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => return true,
2225                Err(_) => continue,
2226            }
2227        }
2228        false
2229    })
2230    .await
2231    .unwrap_or(false)
2232}
2233
2234/// Best-effort lookup of the process occupying a port via `listeners::get_all()`.
2235///
2236/// Returns `(pid, process_name)`.  Falls back to `(0, "unknown")` when the
2237/// system call fails (permission error, unsupported OS, etc.).
2238async fn identify_port_owner(port: u16) -> (u32, String) {
2239    tokio::task::spawn_blocking(move || {
2240        listeners::get_all()
2241            .ok()
2242            .and_then(|list| {
2243                list.into_iter()
2244                    .find(|l| l.socket.port() == port)
2245                    .map(|l| (l.process.pid, l.process.name))
2246            })
2247            .unwrap_or((0, "unknown".to_string()))
2248    })
2249    .await
2250    .unwrap_or((0, "unknown".to_string()))
2251}
2252
2253/// Detect whether a port is in use, and if so, identify the owning process.
2254///
2255/// Combines `is_port_in_use` (reliable bind probe) with `identify_port_owner`
2256/// (best-effort process lookup).  Returns `None` when the port is free.
2257async fn detect_port_conflict(port: u16) -> Option<(u32, String)> {
2258    if !is_port_in_use(port).await {
2259        return None;
2260    }
2261    Some(identify_port_owner(port).await)
2262}
2263
2264/// Spawn a background task that detects the first port the daemon process is listening on
2265/// and stores it in the state file as `active_port`.
2266///
2267/// This is called once when the daemon becomes ready. The port is cleared when the daemon stops.
2268///
2269/// Port selection strategy:
2270/// 1. If the daemon has `expected_port` configured, prefer the first port from that list
2271///    (it is the port the operator explicitly designated as the primary service port).
2272/// 2. Otherwise, take the first port the process is actually listening on (in the order
2273///    returned by the OS), which is typically the port bound earliest.
2274///
2275/// Using `min()` (lowest port number) was previously used here but is incorrect: many
2276/// applications listen on multiple ports (e.g. HTTP + metrics) and the lowest-numbered
2277/// port is not necessarily the primary service port.
2278fn detect_and_store_active_port(id: DaemonId, pid: u32) {
2279    tokio::spawn(async move {
2280        // Retry with exponential backoff so that slow-starting daemons (JVM,
2281        // Node.js, Python, etc.) that take more than 500 ms to bind their port
2282        // are still detected.  Total wait budget: 500+1000+2000+4000 = 7.5 s.
2283        for delay_ms in [500u64, 1000, 2000, 4000] {
2284            tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
2285
2286            // Read daemon state atomically: check if still alive and get expected_port
2287            // in a single lock acquisition to avoid TOCTOU and unnecessary lock overhead.
2288            let expected_port: Option<u16> = {
2289                let state_file = SUPERVISOR.state_file.lock().await;
2290                match state_file.daemons.get(&id) {
2291                    Some(d) if d.pid.is_none() => {
2292                        debug!("daemon {id}: aborting active_port detection — process exited");
2293                        return;
2294                    }
2295                    Some(d) => d
2296                        .port
2297                        .as_ref()
2298                        .and_then(|p| p.expect.first().copied())
2299                        .filter(|&p| p > 0),
2300                    None => None,
2301                }
2302            };
2303
2304            let active_port = tokio::task::spawn_blocking(move || {
2305                let listeners = listeners::get_all().ok()?;
2306
2307                // Refresh process tree so all_children sees current descendants.
2308                PROCS.refresh_processes();
2309
2310                let descendant_pids: std::collections::HashSet<u32> = PROCS
2311                    .all_children(pid)
2312                    .into_iter()
2313                    .chain(std::iter::once(pid))
2314                    .collect();
2315
2316                let process_ports: Vec<u16> = listeners
2317                    .into_iter()
2318                    .filter(|listener| descendant_pids.contains(&listener.process.pid))
2319                    .map(|listener| listener.socket.port())
2320                    .filter(|&port| port > 0)
2321                    .collect();
2322
2323                if process_ports.is_empty() {
2324                    return None;
2325                }
2326
2327                // Prefer the configured expected_port if the process is actually
2328                // listening on it; otherwise fall back to the first port found.
2329                if let Some(ep) = expected_port
2330                    && process_ports.contains(&ep)
2331                {
2332                    return Some(ep);
2333                }
2334
2335                // No expected_port match — return the first port in the list.
2336                // The list order reflects the order the OS reports listeners,
2337                // which is generally the order they were bound (earliest first).
2338                // Do NOT sort: the lowest-numbered port is not necessarily the
2339                // primary service port (e.g. HTTP vs metrics).
2340                process_ports.into_iter().next()
2341            })
2342            .await
2343            .ok()
2344            .flatten();
2345
2346            if let Some(port) = active_port {
2347                debug!("daemon {id} active_port detected: {port}");
2348                let mut state_file = SUPERVISOR.state_file.lock().await;
2349                if let Some(d) = state_file.daemons.get(&id) {
2350                    // Guard against PID reuse: if the original process exited and the OS
2351                    // assigned the same PID to an unrelated process that happens to bind
2352                    // a port, we must not route proxy traffic to that unrelated service.
2353                    if d.pid == Some(pid) {
2354                        state_file.set_active_port(&id, port);
2355                    } else {
2356                        debug!(
2357                            "daemon {id}: skipping active_port write — PID mismatch \
2358                             (expected {pid}, current {:?})",
2359                            d.pid
2360                        );
2361                        return;
2362                    }
2363                }
2364                return;
2365            }
2366
2367            debug!(
2368                "daemon {id}: no active port detected for pid {pid} or its descendants (will retry)"
2369            );
2370        }
2371
2372        debug!(
2373            "daemon {id}: active port detection exhausted all retries for pid {pid} and its descendants"
2374        );
2375    });
2376}
2377
2378/// Check whether a daemon (by its qualified ID) is the target of any registered
2379/// slug in the global config.  This is used to decide whether to run the
2380/// `detect_and_store_active_port` polling task — only slug-targeted daemons need
2381/// it, avoiding wasted `listeners::get_all()` calls for port-less daemons.
2382///
2383/// Delegates to `proxy::server::is_slug_target()` which uses the same in-memory
2384/// slug cache as the proxy hot path, so this check is cheap.
2385fn is_daemon_slug_target(id: &DaemonId) -> bool {
2386    // read_global_slugs is called once per daemon start — acceptable cost.
2387    // We intentionally avoid making this async to keep has_port_config evaluation
2388    // simple and synchronous in run_once().
2389    let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
2390    slugs.iter().any(|(slug, entry)| {
2391        let daemon_name = entry.daemon.as_deref().unwrap_or(slug);
2392        id.name() == daemon_name
2393    })
2394}
2395
2396#[cfg(all(test, unix))]
2397mod tests {
2398    use super::*;
2399
2400    #[test]
2401    fn test_resolve_run_identity_empty_without_sudo() {
2402        let identity = resolve_run_identity(None, 501, 20, None, None).unwrap();
2403        assert_eq!(identity, RunIdentity::Inherit);
2404    }
2405
2406    #[test]
2407    fn test_resolve_run_identity_sudo_fallback() {
2408        let identity = resolve_run_identity(None, 0, 0, Some("501"), Some("20")).unwrap();
2409        let RunIdentity::Switch { uid, gid, .. } = identity else {
2410            panic!("expected identity switch");
2411        };
2412        assert_eq!(uid.as_raw(), 501);
2413        assert_eq!(gid.as_raw(), 20);
2414    }
2415
2416    #[test]
2417    fn test_resolve_run_identity_ignores_stale_sudo_when_not_root() {
2418        let identity = resolve_run_identity(None, 501, 20, Some("0"), Some("0")).unwrap();
2419        assert_eq!(identity, RunIdentity::Inherit);
2420    }
2421
2422    #[test]
2423    fn test_resolve_configured_user_root_name() {
2424        let identity = resolve_configured_user("root").unwrap();
2425        let RunIdentity::Switch { uid, username, .. } = identity else {
2426            panic!("expected identity switch");
2427        };
2428        assert_eq!(uid.as_raw(), 0);
2429        assert_eq!(
2430            username.as_deref().and_then(|s| s.to_str().ok()),
2431            Some("root")
2432        );
2433    }
2434
2435    #[test]
2436    fn test_resolve_configured_user_root_uid() {
2437        let identity = resolve_configured_user("0").unwrap();
2438        let RunIdentity::Switch { uid, username, .. } = identity else {
2439            panic!("expected identity switch");
2440        };
2441        assert_eq!(uid.as_raw(), 0);
2442        assert_eq!(
2443            username.as_deref().and_then(|s| s.to_str().ok()),
2444            Some("root")
2445        );
2446    }
2447
2448    #[test]
2449    fn test_resolve_configured_user_missing_user_fails() {
2450        let err = resolve_configured_user("pitchfork-user-that-should-not-exist")
2451            .unwrap_err()
2452            .to_string();
2453        assert!(err.contains("does not exist"));
2454    }
2455
2456    #[test]
2457    fn test_resolve_run_identity_requires_root_for_user_switch() {
2458        let err = resolve_run_identity(Some("root"), 501, 20, None, None)
2459            .unwrap_err()
2460            .to_string();
2461        assert!(err.contains("Restart the supervisor with sudo"));
2462    }
2463
2464    #[test]
2465    fn test_resolve_run_identity_same_user_is_noop() {
2466        let identity = resolve_run_identity(Some("root"), 0, 0, Some("501"), Some("20")).unwrap();
2467        assert_eq!(identity, RunIdentity::Inherit);
2468    }
2469}
2470
2471/// Inject proxy-related environment variables into a daemon's command.
2472///
2473/// Adds:
2474/// - `HOST` — the address the daemon should bind to (`127.0.0.1`, omitted in LAN mode)
2475/// - `PITCHFORK_URL` — the public proxy URL for this daemon (if it has a slug)
2476/// - `NODE_EXTRA_CA_CERTS` — path to the pitchfork CA cert (if HTTPS enabled)
2477/// - `__VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS` — `.<tld>` for Vite host allowlisting
2478/// - `PITCHFORK_LAN` — set to `"1"` when LAN mode is active
2479fn inject_proxy_env(cmd: &mut tokio::process::Command, slug: &Option<String>) {
2480    let s = crate::settings::settings();
2481    let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
2482
2483    if should_force_loopback_host(slug) && !lan_enabled {
2484        // Only force loopback binding for daemons that are actually routed via a slug.
2485        // In LAN mode, daemons need to bind to 0.0.0.0 to be reachable from the network.
2486        cmd.env("HOST", "127.0.0.1");
2487    }
2488
2489    // PITCHFORK_URL: the daemon's public proxy URL (only if it has a slug and proxy is enabled)
2490    if let Some(url) = build_pitchfork_url(slug, &s) {
2491        cmd.env("PITCHFORK_URL", &url);
2492    }
2493
2494    // NODE_EXTRA_CA_CERTS: let Node.js backends trust the pitchfork CA
2495    if s.proxy.enable && s.proxy.https {
2496        let ca_path = if s.proxy.tls_cert.is_empty() {
2497            crate::env::PITCHFORK_STATE_DIR.join("proxy").join("ca.pem")
2498        } else {
2499            std::path::PathBuf::from(&s.proxy.tls_cert)
2500        };
2501        if ca_path.exists() {
2502            cmd.env("NODE_EXTRA_CA_CERTS", ca_path.to_string_lossy().to_string());
2503        }
2504    }
2505
2506    // __VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS: Vite host allowlisting
2507    if s.proxy.enable {
2508        let tld = if lan_enabled { "local" } else { &s.proxy.tld };
2509        cmd.env("__VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS", format!(".{tld}"));
2510    }
2511
2512    // PITCHFORK_LAN: signal to daemons that LAN mode is active
2513    if lan_enabled {
2514        cmd.env("PITCHFORK_LAN", "1");
2515    }
2516}
2517
2518fn should_force_loopback_host(slug: &Option<String>) -> bool {
2519    let Some(slug) = slug.as_deref() else {
2520        return false;
2521    };
2522
2523    let s = crate::settings::settings();
2524    if !s.proxy.enable {
2525        return false;
2526    }
2527
2528    let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
2529    slugs.contains_key(slug)
2530}
2531
2532/// Compute the public proxy URL for a daemon.
2533///
2534/// Returns `None` if the daemon has no slug or the proxy is not enabled.
2535fn build_pitchfork_url(slug: &Option<String>, s: &crate::settings::Settings) -> Option<String> {
2536    let slug = slug.as_ref()?;
2537    if !s.proxy.enable {
2538        return None;
2539    }
2540    let scheme = if s.proxy.https { "https" } else { "http" };
2541    let port = u16::try_from(s.proxy.port).ok().filter(|&p| p > 0)?;
2542    let port_suffix = if (scheme == "https" && port == 443) || (scheme == "http" && port == 80) {
2543        String::new()
2544    } else {
2545        format!(":{port}")
2546    };
2547    let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
2548    let tld = if lan_enabled { "local" } else { &s.proxy.tld };
2549    Some(format!("{scheme}://{slug}.{tld}{port_suffix}",))
2550}
2551
2552#[cfg(test)]
2553mod ready_check_tests {
2554    use super::*;
2555    use std::time::Duration;
2556
2557    #[test]
2558    fn any_ready_check_remaining_prefers_unbounded_checks() {
2559        let http = ReadyHttp::new("http://localhost/health");
2560        let cmd = ReadyCmd::new("true");
2561
2562        assert!(any_ready_check_remaining(
2563            None,
2564            false,
2565            None,
2566            false,
2567            Some(&http),
2568            false,
2569            None,
2570            false
2571        ));
2572        assert!(any_ready_check_remaining(
2573            None,
2574            false,
2575            None,
2576            false,
2577            None,
2578            false,
2579            Some(&cmd),
2580            false
2581        ));
2582        assert!(any_ready_check_remaining(
2583            None,
2584            false,
2585            Some(&ReadyPort::new(8080)),
2586            false,
2587            Some(&http),
2588            true,
2589            Some(&cmd),
2590            true
2591        ));
2592    }
2593
2594    #[test]
2595    fn any_ready_check_remaining_exhausted_timed_checks() {
2596        let http = ReadyHttp {
2597            url: "http://localhost/health".to_string(),
2598            status: vec![],
2599            timeout: Some(Duration::from_secs(5)),
2600        };
2601        let cmd = ReadyCmd {
2602            run: "true".to_string(),
2603            timeout: Some(Duration::from_secs(5)),
2604        };
2605
2606        assert!(any_ready_check_remaining(
2607            None,
2608            false,
2609            None,
2610            false,
2611            Some(&http),
2612            false,
2613            Some(&cmd),
2614            false
2615        ));
2616        assert!(!any_ready_check_remaining(
2617            None,
2618            false,
2619            None,
2620            false,
2621            Some(&http),
2622            true,
2623            Some(&cmd),
2624            true
2625        ));
2626    }
2627
2628    #[tokio::test]
2629    async fn spawn_cmd_probe_reports_success() {
2630        let id = DaemonId::new("global", "probe-test");
2631        let probe = spawn_cmd_probe(&id, "true", &std::env::temp_dir(), 0, None, &[]);
2632        let status = probe.result_rx.await.unwrap().unwrap();
2633        assert!(status.success());
2634    }
2635
2636    #[tokio::test]
2637    async fn spawn_cmd_probe_stops_on_request() {
2638        let id = DaemonId::new("global", "probe-test");
2639        let probe = spawn_cmd_probe(&id, "sleep 30", &std::env::temp_dir(), 0, None, &[]);
2640        let CmdProbe {
2641            cancel_tx,
2642            result_rx,
2643        } = probe;
2644        let _ = cancel_tx.send(());
2645        let status = result_rx.await.unwrap().unwrap();
2646        assert!(!status.success());
2647    }
2648
2649    #[tokio::test]
2650    async fn spawn_cmd_probe_receives_daemon_and_resolved_port_environment() {
2651        let id = DaemonId::new("worktree", "api");
2652        let daemon_env = IndexMap::from([("CUSTOM_VALUE".to_string(), "yes".to_string())]);
2653        let probe = spawn_cmd_probe(
2654            &id,
2655            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"#,
2656            &std::env::temp_dir(),
2657            2,
2658            Some(&daemon_env),
2659            &[4100, 5100],
2660        );
2661        let status = probe.result_rx.await.unwrap().unwrap();
2662        assert!(status.success());
2663    }
2664
2665    #[test]
2666    fn configured_ready_port_follows_expected_port_bump() {
2667        assert_eq!(resolve_configured_ready_port(3000, &[3000], &[3004]), 3004);
2668        assert_eq!(
2669            resolve_configured_ready_port(4000, &[3000, 4000], &[3003, 4003]),
2670            4003
2671        );
2672        assert_eq!(resolve_configured_ready_port(8080, &[3000], &[3004]), 8080);
2673    }
2674}