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