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