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::procs::PROCS;
15use crate::settings::settings;
16use crate::shell::Shell;
17use crate::supervisor::state::UpsertDaemonOpts;
18use crate::{Result, env};
19use miette::IntoDiagnostic;
20use once_cell::sync::Lazy;
21use regex::Regex;
22use std::collections::HashMap;
23#[cfg(unix)]
24use std::ffi::CString;
25use std::sync::{Arc, atomic};
26use std::time::Duration;
27use tokio::io::AsyncBufReadExt;
28use tokio::select;
29use tokio::sync::oneshot;
30use tokio::time;
31
32/// Cache for compiled regex patterns to avoid recompilation on daemon restarts
33static REGEX_CACHE: Lazy<std::sync::Mutex<HashMap<String, Regex>>> =
34    Lazy::new(|| std::sync::Mutex::new(HashMap::new()));
35
36#[cfg(unix)]
37#[derive(Clone, Debug, PartialEq, Eq)]
38enum RunIdentity {
39    Inherit,
40    Switch {
41        uid: nix::unistd::Uid,
42        gid: nix::unistd::Gid,
43        username: Option<CString>,
44    },
45}
46
47/// Get or compile a regex pattern, caching the result for future use
48pub(crate) fn get_or_compile_regex(pattern: &str) -> Option<Regex> {
49    let mut cache = REGEX_CACHE.lock().unwrap_or_else(|e| e.into_inner());
50    if let Some(re) = cache.get(pattern) {
51        return Some(re.clone());
52    }
53    match Regex::new(pattern) {
54        Ok(re) => {
55            cache.insert(pattern.to_string(), re.clone());
56            Some(re)
57        }
58        Err(e) => {
59            error!("invalid regex pattern '{pattern}': {e}");
60            None
61        }
62    }
63}
64
65impl Supervisor {
66    /// Run a daemon, handling retries if configured
67    pub async fn run(&self, opts: RunOptions) -> Result<IpcResponse> {
68        let id = &opts.id;
69        let cmd = opts.cmd.clone();
70
71        // Clear any pending autostop for this daemon since it's being started
72        {
73            let mut pending = self.pending_autostops.lock().await;
74            if pending.remove(id).is_some() {
75                info!("cleared pending autostop for {id} (daemon starting)");
76            }
77        }
78
79        let daemon = self.get_daemon(id).await;
80        if let Some(daemon) = daemon {
81            // Stopping state is treated as "not running" - the monitoring task will clean it up
82            // Only check for Running state with a valid PID
83            if !daemon.status.is_stopping()
84                && !daemon.status.is_stopped()
85                && let Some(pid) = daemon.pid
86            {
87                if opts.force {
88                    self.stop(id).await?;
89                    info!("run: stop completed for daemon {id}");
90                } else {
91                    warn!("daemon {id} already running with pid {pid}");
92                    return Ok(IpcResponse::DaemonAlreadyRunning);
93                }
94            }
95        }
96
97        // If wait_ready is true and retry is configured, implement retry loop
98        if opts.wait_ready && opts.retry.count() > 0 {
99            // Use saturating_add to avoid overflow when retry = u32::MAX (infinite)
100            let max_attempts = opts.retry.count().saturating_add(1);
101            for attempt in 0..max_attempts {
102                let mut retry_opts = opts.clone();
103                retry_opts.retry_count = attempt;
104                retry_opts.cmd = cmd.clone();
105
106                let result = self.run_once(retry_opts).await?;
107
108                match result {
109                    IpcResponse::DaemonReady { daemon } => {
110                        return Ok(IpcResponse::DaemonReady { daemon });
111                    }
112                    IpcResponse::DaemonFailedWithCode { exit_code } => {
113                        if attempt < opts.retry.count() {
114                            let backoff_secs = 2u64.saturating_pow(attempt).min(3600);
115                            info!(
116                                "daemon {id} failed (attempt {}/{}), retrying in {}s",
117                                attempt + 1,
118                                max_attempts,
119                                backoff_secs
120                            );
121                            fire_hook(
122                                HookType::OnRetry,
123                                id.clone(),
124                                opts.dir.0.clone(),
125                                attempt + 1,
126                                opts.env.clone(),
127                                vec![],
128                            )
129                            .await;
130                            time::sleep(Duration::from_secs(backoff_secs)).await;
131                            continue;
132                        } else {
133                            info!("daemon {id} failed after {max_attempts} attempts");
134                            return Ok(IpcResponse::DaemonFailedWithCode { exit_code });
135                        }
136                    }
137                    other => return Ok(other),
138                }
139            }
140        }
141
142        // No retry or wait_ready is false
143        self.run_once(opts).await
144    }
145
146    /// Run a daemon once (single attempt)
147    pub(crate) async fn run_once(&self, opts: RunOptions) -> Result<IpcResponse> {
148        let id = &opts.id;
149        let original_cmd = opts.cmd.clone(); // Save original command for persistence
150
151        // Create channel for readiness notification if wait_ready is true
152        let (ready_tx, ready_rx) = if opts.wait_ready {
153            let (tx, rx) = oneshot::channel();
154            (Some(tx), Some(rx))
155        } else {
156            (None, None)
157        };
158
159        // Check port availability and apply auto-bump if configured
160        let expected_ports = opts
161            .port
162            .as_ref()
163            .map(|p| p.expect.clone())
164            .unwrap_or_default();
165        let (resolved_ports, effective_ready_port) = if !expected_ports.is_empty() {
166            let port_cfg = opts.port.as_ref().unwrap();
167            match check_ports_available(
168                &expected_ports,
169                port_cfg.auto_bump(),
170                port_cfg.max_bump_attempts(),
171            )
172            .await
173            {
174                Ok(resolved) => {
175                    let ready_port = if let Some(configured_port) = opts.ready_port {
176                        // If ready_port matches one of the expected ports, apply the same bump offset
177                        let bump_offset = resolved
178                            .first()
179                            .unwrap_or(&0)
180                            .saturating_sub(*expected_ports.first().unwrap_or(&0));
181                        if expected_ports.contains(&configured_port) && bump_offset > 0 {
182                            configured_port
183                                .checked_add(bump_offset)
184                                .or(Some(configured_port))
185                        } else {
186                            Some(configured_port)
187                        }
188                    } else if opts.ready_output.is_none()
189                        && opts.ready_http.is_none()
190                        && opts.ready_cmd.is_none()
191                        && opts.ready_delay.is_none()
192                    {
193                        // No other ready check configured — use the first expected port as a
194                        // TCP port readiness check so the daemon is considered ready once it
195                        // starts listening.  Skip port 0 (ephemeral port request).
196                        resolved.first().copied().filter(|&p| p != 0)
197                    } else {
198                        // Another ready check is configured (output/http/cmd/delay).
199                        // Don't add an implicit TCP port check — it could race and fire
200                        // before the daemon has produced any output.
201                        None
202                    };
203                    info!("daemon {id}: ports {expected_ports:?} resolved to {resolved:?}");
204                    (resolved, ready_port)
205                }
206                Err(e) => {
207                    error!("daemon {id}: port check failed: {e}");
208                    // Convert PortError to structured IPC response
209                    if let Some(port_error) = e.downcast_ref::<PortError>() {
210                        match port_error {
211                            PortError::InUse { port, process, pid } => {
212                                return Ok(IpcResponse::PortConflict {
213                                    port: *port,
214                                    process: process.clone(),
215                                    pid: *pid,
216                                });
217                            }
218                            PortError::NoAvailablePort {
219                                start_port,
220                                attempts,
221                            } => {
222                                return Ok(IpcResponse::NoAvailablePort {
223                                    start_port: *start_port,
224                                    attempts: *attempts,
225                                });
226                            }
227                        }
228                    }
229                    return Ok(IpcResponse::DaemonFailed {
230                        error: e.to_string(),
231                    });
232                }
233            }
234        } else {
235            // When ready_port is set without expected_port, check that the port
236            // is not already occupied.  If another process is listening on it,
237            // the TCP readiness probe would immediately succeed and pitchfork
238            // would falsely consider the daemon ready — routing proxy traffic to
239            // the wrong process.
240            if let Some(port) = opts.ready_port {
241                if port > 0 {
242                    if let Some((pid, process)) = detect_port_conflict(port).await {
243                        return Ok(IpcResponse::PortConflict { port, process, pid });
244                    }
245                }
246            }
247            (Vec::new(), opts.ready_port)
248        };
249
250        // Parse the configured shell (default "sh -c") into program + args.
251        // The run script is passed verbatim as the final argument, avoiding the
252        // lossy split->join round-trip that previously mangled $VAR/glob expansion.
253        let shell_setting = settings().general.shell.clone();
254        let shell_parts = match shell_words::split(&shell_setting) {
255            Ok(parts) if !parts.is_empty() => parts,
256            Ok(_) => {
257                return Ok(IpcResponse::DaemonFailed {
258                    error: "general.shell setting is empty".to_string(),
259                });
260            }
261            Err(e) => {
262                return Ok(IpcResponse::DaemonFailed {
263                    error: format!("failed to parse general.shell setting {shell_setting:?}: {e}"),
264                });
265            }
266        };
267        let (shell_program, shell_args) = shell_parts.split_first().unwrap();
268
269        // Use the original run string verbatim; fall back to joining cmd for
270        // ad-hoc commands (e.g. `pitchfork run -- cmd args`) that have no run string.
271        // We don't prepend `exec` because it breaks compound commands (e.g. `exec a && b`
272        // silently drops `b`). Users can add `exec` themselves in the run string.
273        let run_script = opts
274            .run
275            .clone()
276            .unwrap_or_else(|| shell_words::join(&original_cmd));
277
278        let (program, args) = if opts.mise.unwrap_or(settings().general.mise) {
279            match settings().resolve_mise_bin() {
280                Some(mise_bin) => {
281                    let mise_bin_str = mise_bin.to_string_lossy().to_string();
282                    info!("daemon {id}: wrapping command with mise ({mise_bin_str})");
283                    let mut args = vec!["x".to_string(), "--".to_string()];
284                    args.push(shell_program.clone());
285                    args.extend(shell_args.iter().cloned());
286                    args.push(run_script);
287                    (mise_bin_str, args)
288                }
289                None => {
290                    warn!("daemon {id}: mise=true but mise binary not found, running without mise");
291                    let mut args: Vec<String> = shell_args.to_vec();
292                    args.push(run_script);
293                    (shell_program.clone(), args)
294                }
295            }
296        } else {
297            let mut args: Vec<String> = shell_args.to_vec();
298            args.push(run_script);
299            (shell_program.clone(), args)
300        };
301        #[cfg(unix)]
302        let run_identity = match resolve_effective_run_identity(opts.user.as_deref()) {
303            Ok(identity) => identity,
304            Err(e) => {
305                return Ok(IpcResponse::DaemonFailed {
306                    error: e.to_string(),
307                });
308            }
309        };
310        info!("run: spawning daemon {id} with {program} {args:?}");
311
312        // Allocate PTY if configured
313        #[cfg(unix)]
314        let pty_pair = if opts.pty.unwrap_or(false) {
315            match super::pty::openpty() {
316                Ok(pair) => {
317                    info!("daemon {id}: allocated PTY (pty = true)");
318                    Some(pair)
319                }
320                Err(e) => {
321                    warn!("daemon {id}: failed to allocate PTY, falling back to pipes: {e}");
322                    None
323                }
324            }
325        } else {
326            None
327        };
328
329        let mut cmd = tokio::process::Command::new(&program);
330
331        #[cfg(unix)]
332        if let Some(ref pair) = pty_pair {
333            // PTY mode: connect both stdout and stderr to the slave PTY.
334            // The child uses the slave for stdin/stdout/stderr, and we read
335            // output from the master.
336            let slave_file = std::fs::File::from(
337                pair.slave
338                    .try_clone()
339                    .map_err(|e| miette::miette!("failed to dup slave PTY fd: {e}"))?,
340            );
341            cmd.stdin(std::process::Stdio::from(slave_file.try_clone().map_err(
342                |e| miette::miette!("failed to clone slave PTY fd for stdin: {e}"),
343            )?));
344            cmd.stdout(std::process::Stdio::from(slave_file.try_clone().map_err(
345                |e| miette::miette!("failed to clone slave PTY fd for stdout: {e}"),
346            )?));
347            cmd.stderr(std::process::Stdio::from(slave_file));
348        } else {
349            cmd.stdout(std::process::Stdio::piped())
350                .stderr(std::process::Stdio::piped());
351        }
352
353        #[cfg(not(unix))]
354        {
355            cmd.stdout(std::process::Stdio::piped())
356                .stderr(std::process::Stdio::piped());
357        }
358
359        cmd.args(&args).current_dir(&opts.dir);
360
361        #[cfg(unix)]
362        if pty_pair.is_none() {
363            cmd.stdin(std::process::Stdio::null());
364        }
365
366        #[cfg(not(unix))]
367        cmd.stdin(std::process::Stdio::null());
368
369        // Ensure daemon can find user tools by using the original PATH
370        if let Some(ref path) = *env::ORIGINAL_PATH {
371            cmd.env("PATH", path);
372        }
373
374        // Apply custom environment variables from config
375        if let Some(ref env_vars) = opts.env {
376            cmd.envs(env_vars);
377        }
378
379        // Inject pitchfork metadata env vars AFTER user env so they can't be overwritten
380        cmd.env("PITCHFORK_DAEMON_ID", id.qualified());
381        cmd.env("PITCHFORK_DAEMON_NAMESPACE", id.namespace());
382        cmd.env("PITCHFORK_RETRY_COUNT", opts.retry_count.to_string());
383
384        // Inject the resolved ports for the daemon to use
385        if !resolved_ports.is_empty() {
386            // Set PORT to the first port for backward compatibility
387            // When there's only one port, both PORT and PORT0 will be set to the same value.
388            // This follows the convention used by many deployment platforms (Heroku, etc.).
389            cmd.env("PORT", resolved_ports[0].to_string());
390            // Set individual ports as PORT0, PORT1, etc.
391            for (i, port) in resolved_ports.iter().enumerate() {
392                cmd.env(format!("PORT{i}"), port.to_string());
393            }
394        }
395
396        // Inject proxy-related environment variables
397        inject_proxy_env(&mut cmd, &opts.slug);
398
399        #[cfg(unix)]
400        {
401            let run_identity = run_identity.clone();
402            let use_pty = pty_pair.is_some();
403            unsafe {
404                cmd.pre_exec(move || {
405                    nix::unistd::setsid().map_err(nix_to_io_error)?;
406
407                    // When using a PTY, set the slave as the controlling terminal.
408                    // The slave FD has already been dup'd onto stdin/stdout/stderr
409                    // by tokio, so we can use stdin (fd 0) for TIOCSCTTY.
410                    if use_pty {
411                        let ret = libc::ioctl(0, libc::TIOCSCTTY as libc::c_ulong, 0);
412                        if ret < 0 {
413                            // Non-fatal: the process can still run without
414                            // a controlling terminal.
415                            #[cfg(target_os = "linux")]
416                            eprintln!(
417                                "pitchfork: TIOCSCTTY failed: {}",
418                                std::io::Error::last_os_error()
419                            );
420                        }
421                    }
422
423                    apply_run_identity(&run_identity)?;
424                    Ok(())
425                });
426            }
427        }
428
429        let mut child = cmd.spawn().into_diagnostic()?;
430        let pid = match child.id() {
431            Some(p) => p,
432            None => {
433                warn!("Daemon {id} exited before PID could be captured");
434                return Ok(IpcResponse::DaemonFailed {
435                    error: "Process exited immediately".to_string(),
436                });
437            }
438        };
439        info!("started daemon {id} with pid {pid}");
440        PROCS.refresh_pids(&[pid]);
441        let daemon = self
442            .upsert_daemon(
443                UpsertDaemonOpts::from_run_options(&opts, DaemonStatus::Running)
444                    .set(|o| {
445                        o.pid = Some(pid);
446                        o.cmd = Some(original_cmd);
447                        o.ready_port = effective_ready_port;
448                        o.port = crate::config_types::PortConfig::from_parts(
449                            expected_ports,
450                            opts.port.as_ref().map(|p| p.bump).unwrap_or_default(),
451                        );
452                        o.resolved_port = resolved_ports;
453                    })
454                    .build(),
455            )
456            .await?;
457
458        let id_clone = id.clone();
459        let ready_delay = opts.ready_delay;
460        let ready_output = opts.ready_output.clone();
461        let ready_http = opts.ready_http.clone();
462        let ready_port = effective_ready_port;
463        let ready_cmd = opts.ready_cmd.clone();
464        let daemon_dir = opts.dir.0.clone();
465        let hook_retry_count = opts.retry_count;
466        let hook_retry = opts.retry;
467        let hook_daemon_env = opts.env.clone();
468        let on_output_hook = opts.on_output_hook.clone();
469        // Whether this daemon has any port-related config — used to skip the
470        // active_port detection task for daemons that never bind a port (e.g. `sleep 60`).
471        // When the proxy is enabled, only detect active_port for daemons that are
472        // actually referenced by a registered slug, rather than blanket-polling every
473        // daemon (which wastes ~7.5 s of listeners::get_all() calls per port-less daemon).
474        let has_port_config = opts.port.as_ref().is_some_and(|p| !p.expect.is_empty())
475            || (settings().proxy.enable && is_daemon_slug_target(id));
476        let daemon_pid = pid;
477
478        // Prepare output readers before spawning the monitoring task.
479        // In PTY mode, we read from the PTY master FD.
480        // In pipe mode, we read from separate stdout/stderr pipes.
481        #[cfg(unix)]
482        let pty_reader = pty_pair.map(|p| {
483            tokio::io::BufReader::new(tokio::fs::File::from_std(std::fs::File::from(p.master)))
484                .lines()
485        });
486        #[cfg(not(unix))]
487        let pty_reader: Option<tokio::io::Lines<tokio::io::BufReader<tokio::fs::File>>> = None;
488        let stdout_reader = if pty_reader.is_none() {
489            child
490                .stdout
491                .take()
492                .map(|s| tokio::io::BufReader::new(s).lines())
493        } else {
494            None
495        };
496        let stderr_reader = if pty_reader.is_none() {
497            child
498                .stderr
499                .take()
500                .map(|s| tokio::io::BufReader::new(s).lines())
501        } else {
502            None
503        };
504
505        if pty_reader.is_none() && (stdout_reader.is_none() || stderr_reader.is_none()) {
506            error!("Failed to capture stdout/stderr for daemon {id}");
507        }
508
509        tokio::spawn(async move {
510            let id = id_clone;
511
512            // Merge all output sources (PTY master OR stdout+stderr) into a single channel.
513            let (output_tx, mut output_rx) = tokio::sync::mpsc::channel::<String>(256);
514
515            if let Some(mut reader) = pty_reader {
516                // PTY mode: single merged stream from the master.
517                // output_tx is moved into the spawn; when the reader ends the
518                // channel closes automatically.
519                tokio::spawn(async move {
520                    while let Ok(Some(mut line)) = reader.next_line().await {
521                        // PTY slave uses ONLCR: \n → \r\n; strip the trailing \r.
522                        if line.ends_with('\r') {
523                            line.pop();
524                        }
525                        if output_tx.send(line).await.is_err() {
526                            break;
527                        }
528                    }
529                });
530            } else {
531                // Pipe mode: stdout and stderr are merged into the same channel.
532                // Both `ready_output` and `on_output_hook` patterns match against
533                // lines from either stream, which is the expected behavior (a
534                // "server ready" message may appear on stderr in some tools).
535                if let Some(mut stdout) = stdout_reader {
536                    let tx = output_tx.clone();
537                    tokio::spawn(async move {
538                        while let Ok(Some(line)) = stdout.next_line().await {
539                            if tx.send(line).await.is_err() {
540                                break;
541                            }
542                        }
543                    });
544                }
545                if let Some(mut stderr) = stderr_reader {
546                    let tx = output_tx.clone();
547                    tokio::spawn(async move {
548                        while let Ok(Some(line)) = stderr.next_line().await {
549                            if tx.send(line).await.is_err() {
550                                break;
551                            }
552                        }
553                    });
554                }
555                // Drop the last sender so the channel closes when all readers finish.
556                drop(output_tx);
557            }
558            let log_store = Arc::clone(&LOG_STORE);
559            let format_line = |line: String| line;
560
561            const LOG_BATCH_SIZE: usize = 100;
562            const LOG_FLUSH_INTERVAL: Duration = Duration::from_millis(100);
563            let mut log_buffer: Vec<String> = Vec::with_capacity(LOG_BATCH_SIZE);
564            let mut log_flush_interval = tokio::time::interval(LOG_FLUSH_INTERVAL);
565            log_flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
566
567            let flush_logs = |buffer: &mut Vec<String>| -> Option<tokio::task::JoinHandle<()>> {
568                if buffer.is_empty() {
569                    return None;
570                }
571                let store = Arc::clone(&log_store);
572                let id = id.clone();
573                let batch = std::mem::take(buffer);
574                Some(tokio::task::spawn_blocking(move || {
575                    if let Err(e) = store.append_batch(&id, &batch) {
576                        error!("Failed to write batch to log for daemon {id}: {e}");
577                    }
578                }))
579            };
580
581            // SQLite WAL mode provides automatic durability; no explicit flush needed.
582
583            // Setup readiness checking
584            let mut ready_notified = false;
585            let mut ready_tx = ready_tx;
586            let ready_pattern = ready_output.as_ref().and_then(|p| get_or_compile_regex(p));
587            // Track whether we've already spawned the active_port detection task
588            let mut active_port_spawned = false;
589
590            // Validate on_output config early; discard the hook on any error so
591            // a bad regex does not silently fall through to the (None, None) => true
592            // match arm and fire on every line.
593            let on_output_hook = match on_output_hook {
594                Some(ref hook) => match hook.validate(id.name()) {
595                    Ok(()) => on_output_hook,
596                    Err(e) => {
597                        error!("{e}");
598                        None
599                    }
600                },
601                None => None,
602            };
603
604            // Compile the regex pattern after validation so we only attempt this
605            // when the hook is known-good (validate() already checked the syntax).
606            let on_output_pattern: Option<regex::Regex> = on_output_hook
607                .as_ref()
608                .and_then(|h| h.regex.as_deref().and_then(get_or_compile_regex));
609            let on_output_debounce = on_output_hook
610                .as_ref()
611                .map(|h| h.debounce_duration())
612                .unwrap_or(Duration::from_millis(1000));
613            // Last time the on_output hook fired; None means it has never fired.
614            let mut on_output_last_fired: Option<std::time::Instant> = None;
615
616            let mut delay_timer =
617                ready_delay.map(|secs| Box::pin(time::sleep(Duration::from_secs(secs))));
618
619            // Get settings for intervals
620            let s = settings();
621            let ready_check_interval = s.supervisor_ready_check_interval();
622            let http_client_timeout = s.supervisor_http_client_timeout();
623
624            // Setup HTTP readiness check interval
625            let mut http_check_interval = ready_http
626                .as_ref()
627                .map(|_| tokio::time::interval(ready_check_interval));
628            let http_client = ready_http.as_ref().map(|_| {
629                reqwest::Client::builder()
630                    .timeout(http_client_timeout)
631                    .build()
632                    .unwrap_or_default()
633            });
634
635            // Setup TCP port readiness check interval
636            let mut port_check_interval =
637                ready_port.map(|_| tokio::time::interval(ready_check_interval));
638
639            // Setup command readiness check interval
640            let mut cmd_check_interval = ready_cmd
641                .as_ref()
642                .map(|_| tokio::time::interval(ready_check_interval));
643
644            // Use a channel to communicate process exit status
645            let (exit_tx, mut exit_rx) =
646                tokio::sync::mpsc::channel::<std::io::Result<std::process::ExitStatus>>(1);
647
648            // Spawn a task to wait for process exit
649            let child_pid = child.id().unwrap_or(0);
650            tokio::spawn(async move {
651                let result = child.wait().await;
652                // On non-Linux Unix (e.g. macOS) the zombie reaper may win the
653                // race and consume the exit status via waitpid(None, WNOHANG)
654                // before Tokio's child.wait() gets to it. When that happens,
655                // Tokio returns an ECHILD io::Error. We recover by checking
656                // REAPED_STATUSES for the stashed exit code.
657                //
658                // On Linux this is unnecessary because the reaper uses
659                // waitid(WNOWAIT) to peek before reaping, which avoids the
660                // race entirely.
661                #[cfg(all(unix, not(target_os = "linux")))]
662                let result = match &result {
663                    Err(e) if e.raw_os_error() == Some(nix::libc::ECHILD) => {
664                        if let Some(code) = super::REAPED_STATUSES.lock().await.remove(&child_pid) {
665                            warn!(
666                                "daemon pid {child_pid} wait() got ECHILD; \
667                                 recovered exit code {code} from zombie reaper"
668                            );
669                            // Synthesize an ExitStatus from the stashed code.
670                            // On Unix we can use `ExitStatus::from_raw()` with
671                            // a wait-style status word (code << 8 for normal
672                            // exit, or raw signal number for signal death).
673                            use std::os::unix::process::ExitStatusExt;
674                            if code >= 0 {
675                                Ok(std::process::ExitStatus::from_raw(code << 8))
676                            } else {
677                                // Negative code means killed by signal (-sig)
678                                Ok(std::process::ExitStatus::from_raw((-code) & 0x7f))
679                            }
680                        } else {
681                            warn!(
682                                "daemon pid {child_pid} wait() got ECHILD but no \
683                                 stashed status found; reporting as error"
684                            );
685                            result
686                        }
687                    }
688                    _ => result,
689                };
690                debug!("daemon pid {child_pid} wait() completed with result: {result:?}");
691                let _ = exit_tx.send(result).await;
692            });
693
694            #[allow(unused_assignments)]
695            // Initial None is a safety net; loop only exits via exit_rx.recv() which sets it
696            let mut exit_status = None;
697
698            // If there is no ready check of any kind and no delay, the daemon is
699            // considered immediately ready and the active_port detection task would
700            // never be triggered inside the select loop.  Kick it off right away so
701            // that daemons without any readiness configuration still get their
702            // active_port populated (needed for proxy routing).
703            if has_port_config
704                && ready_pattern.is_none()
705                && ready_http.is_none()
706                && ready_port.is_none()
707                && ready_cmd.is_none()
708                && delay_timer.is_none()
709            {
710                active_port_spawned = true;
711                detect_and_store_active_port(id.clone(), daemon_pid);
712            }
713
714            loop {
715                select! {
716                    Some(line) = output_rx.recv() => {
717                        let formatted = format_line(line.clone());
718                        log_buffer.push(formatted.clone());
719                        if log_buffer.len() >= LOG_BATCH_SIZE {
720                            let _ = flush_logs(&mut log_buffer);
721                        }
722                        trace!("output: {id} {formatted}");
723
724                        // Strip ANSI for pattern matching so user-written patterns
725                        // work regardless of whether the process emits color codes.
726                        let line_clean = console::strip_ansi_codes(&line).to_string();
727
728                        // Check if output matches ready pattern
729                        if !ready_notified
730                            && let Some(ref pattern) = ready_pattern
731                            && pattern.is_match(&line_clean)
732                        {
733                            // Flush buffered logs synchronously before signalling
734                            // readiness, so collect_startup_logs sees the line
735                            // that triggered the match (and any co-buffered lines)
736                            // in SQLite.
737                            if let Some(handle) = flush_logs(&mut log_buffer) {
738                                let _ = handle.await;
739                            }
740                            info!("daemon {id} ready: output matched pattern");
741                            ready_notified = true;
742                            if let Some(tx) = ready_tx.take() {
743                                let _ = tx.send(Ok(()));
744                            }
745                            fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
746                            if !active_port_spawned && has_port_config {
747                                active_port_spawned = true;
748                                detect_and_store_active_port(id.clone(), daemon_pid);
749                            }
750                        }
751
752                        // Check on_output hook
753                        if let Some(ref hook) = on_output_hook {
754                            let matched = match (&hook.filter, &on_output_pattern) {
755                                (Some(substr), _) => line_clean.contains(substr.as_str()),
756                                (None, Some(re)) => re.is_match(&line_clean),
757                                (None, None) => true,
758                            };
759                            if matched {
760                                let now = std::time::Instant::now();
761                                let elapsed = on_output_last_fired.map(|t| now.duration_since(t));
762                                if elapsed.is_none_or(|e| e >= on_output_debounce) {
763                                    on_output_last_fired = Some(now);
764                                    hooks::fire_output_hook(id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), hook.run.clone(), line_clean.clone()).await;
765                                }
766                            }
767                        }
768                    }
769                    Some(result) = exit_rx.recv() => {
770                        // Process exited - save exit status and notify if not ready yet
771                        exit_status = Some(result);
772                        debug!("daemon {id} process exited, exit_status: {exit_status:?}");
773                        if !ready_notified {
774                            if let Some(tx) = ready_tx.take() {
775                                // Check if process exited successfully
776                                let is_success = exit_status.as_ref()
777                                    .and_then(|r| r.as_ref().ok())
778                                    .map(|s| s.success())
779                                    .unwrap_or(false);
780
781                                if is_success {
782                                    debug!("daemon {id} exited successfully before ready check, sending success notification");
783                                    let _ = tx.send(Ok(()));
784                                } else {
785                                    let exit_code = exit_status.as_ref()
786                                        .and_then(|r| r.as_ref().ok())
787                                        .and_then(|s| s.code());
788                                    debug!("daemon {id} exited with failure before ready check, sending failure notification with exit_code: {exit_code:?}");
789                                    let _ = tx.send(Err(exit_code));
790                                }
791                            }
792                        } else {
793                            debug!("daemon {id} was already marked ready, not sending notification");
794                        }
795                        break;
796                    },
797                    _ = async {
798                        if let Some(ref mut interval) = http_check_interval {
799                            interval.tick().await;
800                        } else {
801                            std::future::pending::<()>().await;
802                        }
803                    }, if !ready_notified && ready_http.is_some() => {
804                        if let (Some(http), Some(client)) = (&ready_http, &http_client) {
805                            match client.get(&http.url).send().await {
806                                Ok(response) if http.accepts_status(response.status().as_u16()) => {
807                                    info!("daemon {id} ready: HTTP check passed (status {})", response.status());
808                                    ready_notified = true;
809                                    if let Some(tx) = ready_tx.take() {
810                                        let _ = tx.send(Ok(()));
811                                    }
812                                    fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
813                                    http_check_interval = None;
814                                    if !active_port_spawned && has_port_config {
815                                        active_port_spawned = true;
816                                        detect_and_store_active_port(id.clone(), daemon_pid);
817                                    }
818                                }
819                                Ok(response) => {
820                                    trace!("daemon {id} HTTP check: status {} (not ready)", response.status());
821                                }
822                                Err(e) => {
823                                    trace!("daemon {id} HTTP check failed: {e}");
824                                }
825                            }
826                        }
827                    }
828                    _ = async {
829                        if let Some(ref mut interval) = port_check_interval {
830                            interval.tick().await;
831                        } else {
832                            std::future::pending::<()>().await;
833                        }
834                    }, if !ready_notified && ready_port.is_some() => {
835                        if let Some(port) = ready_port {
836                            match tokio::net::TcpStream::connect(("127.0.0.1", port)).await {
837                                Ok(_) => {
838                                    info!("daemon {id} ready: TCP port {port} is listening");
839                                    ready_notified = true;
840                                    if let Some(tx) = ready_tx.take() {
841                                        let _ = tx.send(Ok(()));
842                                    }
843                                    fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
844                                    // Stop checking once ready
845                                    port_check_interval = None;
846                                    if !active_port_spawned && has_port_config {
847                                        active_port_spawned = true;
848                                        detect_and_store_active_port(id.clone(), daemon_pid);
849                                    }
850                                }
851                                Err(_) => {
852                                    trace!("daemon {id} port check: port {port} not listening yet");
853                                }
854                            }
855                        }
856                    }
857                    _ = async {
858                        if let Some(ref mut interval) = cmd_check_interval {
859                            interval.tick().await;
860                        } else {
861                            std::future::pending::<()>().await;
862                        }
863                    }, if !ready_notified && ready_cmd.is_some() => {
864                        if let Some(ref cmd) = ready_cmd {
865                            // Run the readiness check command using the shell abstraction
866                            let mut command = Shell::default_for_platform().command(cmd);
867                            command
868                                .current_dir(&daemon_dir)
869                                .stdout(std::process::Stdio::null())
870                                .stderr(std::process::Stdio::null());
871                            let result: std::io::Result<std::process::ExitStatus> = command.status().await;
872                            match result {
873                                Ok(status) if status.success() => {
874                                    info!("daemon {id} ready: readiness command succeeded");
875                                    ready_notified = true;
876                                    if let Some(tx) = ready_tx.take() {
877                                        let _ = tx.send(Ok(()));
878                                    }
879                                    fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
880                                    // Stop checking once ready
881                                    cmd_check_interval = None;
882                                    if !active_port_spawned && has_port_config {
883                                        active_port_spawned = true;
884                                        detect_and_store_active_port(id.clone(), daemon_pid);
885                                    }
886                                }
887                                Ok(_) => {
888                                    trace!("daemon {id} cmd check: command returned non-zero (not ready)");
889                                }
890                                Err(e) => {
891                                    trace!("daemon {id} cmd check failed: {e}");
892                                }
893                            }
894                        }
895                    }
896                    _ = async {
897                        if let Some(ref mut timer) = delay_timer {
898                            timer.await;
899                        } else {
900                            std::future::pending::<()>().await;
901                        }
902                    } => {
903                        if !ready_notified && ready_pattern.is_none() && ready_http.is_none() && ready_port.is_none() && ready_cmd.is_none() {
904                            info!("daemon {id} ready: delay elapsed");
905                            ready_notified = true;
906                            if let Some(tx) = ready_tx.take() {
907                                let _ = tx.send(Ok(()));
908                            }
909                            fire_hook(HookType::OnReady, id.clone(), daemon_dir.clone(), hook_retry_count, hook_daemon_env.clone(), vec![]).await;
910                        }
911                        // Disable timer after it fires
912                        delay_timer = None;
913                        if !active_port_spawned && has_port_config {
914                            active_port_spawned = true;
915                            detect_and_store_active_port(id.clone(), daemon_pid);
916                        }
917                    }
918                    _ = log_flush_interval.tick() => {
919                        let _ = flush_logs(&mut log_buffer);
920                    }
921                }
922            }
923
924            // Drain any in-flight output lines that were still in the mpsc
925            // channel or the OS pipe buffer when the child exited. Without
926            // this, trailing log lines from short-lived daemons get dropped.
927            // The reader tasks drop their senders on EOF, so recv() returns
928            // None when all data has been consumed. A total deadline of 5 s
929            // guards against a stuck reader (e.g. PTY master FD not closing)
930            // while ensuring drain doesn't block post-exit cleanup indefinitely.
931            let drain_deadline = tokio::time::Instant::now() + Duration::from_secs(5);
932            loop {
933                let now = tokio::time::Instant::now();
934                if now >= drain_deadline {
935                    break;
936                }
937                let Ok(Some(line)) =
938                    tokio::time::timeout(drain_deadline - now, output_rx.recv()).await
939                else {
940                    break;
941                };
942                log_buffer.push(format_line(line));
943            }
944            // Flush any remaining log lines (including drained) before the process exits.
945            // Await the flush to guarantee all buffered logs are persisted before cleanup.
946            if let Some(handle) = flush_logs(&mut log_buffer) {
947                let _ = handle.await;
948            }
949
950            // Clear active_port since the process is no longer running
951            {
952                let mut state_file = SUPERVISOR.state_file.lock().await;
953                state_file.clear_active_port(&id);
954            }
955
956            // Get the final exit status
957            let exit_status = if let Some(status) = exit_status {
958                status
959            } else {
960                // Streams closed but process hasn't exited yet, wait for it
961                match exit_rx.recv().await {
962                    Some(status) => status,
963                    None => {
964                        warn!("daemon {id} exit channel closed without receiving status");
965                        Err(std::io::Error::other("exit channel closed"))
966                    }
967                }
968            };
969            let current_daemon = SUPERVISOR.get_daemon(&id).await;
970
971            // Signal that this monitoring task is processing its exit path.
972            // The RAII guard will decrement the counter and notify close()
973            // when the task finishes (including all fire_hook registrations),
974            // regardless of which return path is taken.
975            SUPERVISOR
976                .active_monitors
977                .fetch_add(1, atomic::Ordering::Release);
978            struct MonitorGuard;
979            impl Drop for MonitorGuard {
980                fn drop(&mut self) {
981                    SUPERVISOR
982                        .active_monitors
983                        .fetch_sub(1, atomic::Ordering::Release);
984                    SUPERVISOR.monitor_done.notify_waiters();
985                }
986            }
987            let _monitor_guard = MonitorGuard;
988            // Check if this monitoring task is for the current daemon process.
989            // Allow Stopped/Stopping daemons through: stop() clears pid atomically,
990            // so d.pid != Some(pid) would be true, but we still need the is_stopped()
991            // branch below to fire on_stop/on_exit hooks.
992            if current_daemon.is_none()
993                || current_daemon.as_ref().is_some_and(|d| {
994                    d.pid != Some(pid) && !d.status.is_stopped() && !d.status.is_stopping()
995                })
996            {
997                // Another process has taken over, don't update status
998                return;
999            }
1000            // Capture the intentional-stop flag BEFORE any state changes.
1001            // stop() transitions Stopping → Stopped and clears pid. If stop() wins the race
1002            // and sets Stopped before this task runs, we still need to fire on_stop/on_exit.
1003            // Treat both Stopping and Stopped as "intentional stop by pitchfork".
1004            let already_stopped = current_daemon
1005                .as_ref()
1006                .is_some_and(|d| d.status.is_stopped());
1007            let is_stopping = already_stopped
1008                || current_daemon
1009                    .as_ref()
1010                    .is_some_and(|d| d.status.is_stopping());
1011
1012            // --- Phase 1: Determine exit_code, exit_reason, and update daemon state ---
1013            let (exit_code, exit_reason) = match (&exit_status, is_stopping) {
1014                (Ok(status), true) => {
1015                    // Intentional stop (by pitchfork). status.code() returns None
1016                    // on Unix when killed by signal (e.g. SIGTERM); use -1 to
1017                    // distinguish from a clean exit code 0.
1018                    (status.code().unwrap_or(-1), "stop")
1019                }
1020                (Ok(status), false) if status.success() => (status.code().unwrap_or(-1), "exit"),
1021                (Ok(status), false) => (status.code().unwrap_or(-1), "fail"),
1022                (Err(_), true) => {
1023                    // child.wait() error while stopping (e.g. sysinfo reaped the process)
1024                    (-1, "stop")
1025                }
1026                (Err(_), false) => (-1, "fail"),
1027            };
1028
1029            // Update daemon state unless stop() already did it (won the race).
1030            if !already_stopped {
1031                if let Ok(status) = &exit_status {
1032                    info!("daemon {id} exited with status {status}");
1033                }
1034                let (new_status, last_exit_success) = match exit_reason {
1035                    "stop" | "exit" => (
1036                        DaemonStatus::Stopped,
1037                        exit_status.as_ref().map(|s| s.success()).unwrap_or(true),
1038                    ),
1039                    _ => (DaemonStatus::Errored(exit_code), false),
1040                };
1041                if let Err(e) = SUPERVISOR
1042                    .upsert_daemon(
1043                        UpsertDaemonOpts::builder(id.clone())
1044                            .set(|o| {
1045                                o.pid = None;
1046                                o.status = new_status;
1047                                o.last_exit_success = Some(last_exit_success);
1048                            })
1049                            .build(),
1050                    )
1051                    .await
1052                {
1053                    error!("Failed to update daemon state for {id}: {e}");
1054                }
1055            }
1056
1057            // --- Phase 2: Fire hooks ---
1058            let hook_extra_env = vec![
1059                ("PITCHFORK_EXIT_CODE".to_string(), exit_code.to_string()),
1060                ("PITCHFORK_EXIT_REASON".to_string(), exit_reason.to_string()),
1061            ];
1062
1063            // Determine which hooks to fire based on exit reason
1064            let hooks_to_fire: Vec<HookType> = match exit_reason {
1065                "stop" => vec![HookType::OnStop, HookType::OnExit],
1066                "exit" => vec![HookType::OnExit],
1067                // "fail": fire on_fail + on_exit only when retries are exhausted
1068                _ if hook_retry_count >= hook_retry.count() => {
1069                    vec![HookType::OnFail, HookType::OnExit]
1070                }
1071                _ => vec![],
1072            };
1073
1074            for hook_type in hooks_to_fire {
1075                fire_hook(
1076                    hook_type,
1077                    id.clone(),
1078                    daemon_dir.clone(),
1079                    hook_retry_count,
1080                    hook_daemon_env.clone(),
1081                    hook_extra_env.clone(),
1082                )
1083                .await;
1084            }
1085        });
1086
1087        // If wait_ready is true, wait for readiness notification
1088        if let Some(ready_rx) = ready_rx {
1089            match ready_rx.await {
1090                Ok(Ok(())) => {
1091                    info!("daemon {id} is ready");
1092                    Ok(IpcResponse::DaemonReady { daemon })
1093                }
1094                Ok(Err(exit_code)) => {
1095                    error!("daemon {id} failed before becoming ready");
1096                    Ok(IpcResponse::DaemonFailedWithCode { exit_code })
1097                }
1098                Err(_) => {
1099                    error!("readiness channel closed unexpectedly for daemon {id}");
1100                    Ok(IpcResponse::DaemonStart { daemon })
1101                }
1102            }
1103        } else {
1104            Ok(IpcResponse::DaemonStart { daemon })
1105        }
1106    }
1107
1108    /// Stop a running daemon
1109    pub async fn stop(&self, id: &DaemonId) -> Result<IpcResponse> {
1110        let pitchfork_id = DaemonId::pitchfork();
1111        if *id == pitchfork_id {
1112            return Ok(IpcResponse::Error(
1113                "Cannot stop supervisor via stop command".into(),
1114            ));
1115        }
1116        info!("stopping daemon: {id}");
1117        if let Some(daemon) = self.get_daemon(id).await {
1118            trace!("daemon to stop: {daemon}");
1119            if let Some(pid) = daemon.pid {
1120                trace!("killing pid: {pid}");
1121                if PROCS.is_running(pid) {
1122                    // First set status to Stopping (preserve PID for monitoring task)
1123                    self.upsert_daemon(
1124                        UpsertDaemonOpts::builder(id.clone())
1125                            .set(|o| {
1126                                o.pid = Some(pid);
1127                                o.status = DaemonStatus::Stopping;
1128                            })
1129                            .build(),
1130                    )
1131                    .await?;
1132
1133                    // Kill the entire process group atomically (daemon PID == PGID
1134                    // because we called setsid() at spawn time)
1135                    let stop_cfg = daemon.stop_signal.unwrap_or_default();
1136                    let stop_signal: i32 = stop_cfg.signal.into();
1137                    if let Err(e) = PROCS
1138                        .kill_process_group_async(pid, stop_signal, stop_cfg.timeout)
1139                        .await
1140                    {
1141                        debug!("failed to kill pid {pid}: {e}");
1142                        // Check if the process is actually stopped despite the error
1143                        if PROCS.is_running(pid) {
1144                            // Process still running after kill attempt - set back to Running
1145                            debug!("failed to stop pid {pid}: process still running after kill");
1146                            self.upsert_daemon(
1147                                UpsertDaemonOpts::builder(id.clone())
1148                                    .set(|o| {
1149                                        o.pid = Some(pid); // Preserve PID to avoid orphaning the process
1150                                        o.status = DaemonStatus::Running;
1151                                    })
1152                                    .build(),
1153                            )
1154                            .await?;
1155                            return Ok(IpcResponse::DaemonStopFailed {
1156                                error: format!(
1157                                    "process {pid} still running after kill attempt: {e}"
1158                                ),
1159                            });
1160                        }
1161                    }
1162
1163                    // Process successfully stopped
1164                    // Note: kill_async uses SIGTERM -> wait ~3s -> SIGKILL strategy,
1165                    // and also detects zombie processes, so by the time it returns,
1166                    // the process should be fully terminated.
1167                    self.upsert_daemon(
1168                        UpsertDaemonOpts::builder(id.clone())
1169                            .set(|o| {
1170                                o.pid = None;
1171                                o.status = DaemonStatus::Stopped;
1172                                o.last_exit_success = Some(true); // Manual stop is considered successful
1173                            })
1174                            .build(),
1175                    )
1176                    .await?;
1177                } else {
1178                    debug!("pid {pid} not running, process may have exited unexpectedly");
1179                    // Process already dead, directly mark as stopped
1180                    // Note that the cleanup logic is handled in monitor task
1181                    self.upsert_daemon(
1182                        UpsertDaemonOpts::builder(id.clone())
1183                            .set(|o| {
1184                                o.pid = None;
1185                                o.status = DaemonStatus::Stopped;
1186                            })
1187                            .build(),
1188                    )
1189                    .await?;
1190                    return Ok(IpcResponse::DaemonWasNotRunning);
1191                }
1192                Ok(IpcResponse::Ok)
1193            } else {
1194                debug!("daemon {id} not running");
1195                Ok(IpcResponse::DaemonNotRunning)
1196            }
1197        } else {
1198            debug!("daemon {id} not found");
1199            Ok(IpcResponse::DaemonNotFound)
1200        }
1201    }
1202}
1203
1204#[cfg(unix)]
1205fn resolve_effective_run_identity(daemon_user: Option<&str>) -> Result<RunIdentity> {
1206    let s = settings();
1207    let settings_user = s.supervisor.user.trim();
1208    let daemon_user = daemon_user.map(str::trim).filter(|user| !user.is_empty());
1209    let settings_user = (!settings_user.is_empty()).then_some(settings_user);
1210    let configured = daemon_user.or(settings_user);
1211    let current_uid = nix::unistd::Uid::effective().as_raw();
1212    let current_gid = nix::unistd::Gid::effective().as_raw();
1213    resolve_run_identity(
1214        configured,
1215        current_uid,
1216        current_gid,
1217        std::env::var("SUDO_UID").ok().as_deref(),
1218        std::env::var("SUDO_GID").ok().as_deref(),
1219    )
1220}
1221
1222#[cfg(unix)]
1223fn resolve_run_identity(
1224    configured: Option<&str>,
1225    current_uid: u32,
1226    current_gid: u32,
1227    sudo_uid: Option<&str>,
1228    sudo_gid: Option<&str>,
1229) -> Result<RunIdentity> {
1230    let current_uid = nix::unistd::Uid::from_raw(current_uid);
1231    let current_gid = nix::unistd::Gid::from_raw(current_gid);
1232    if let Some(user) = configured {
1233        let identity = resolve_configured_user(user)?;
1234        ensure_can_use_identity(user, &identity, current_uid, current_gid)?;
1235        if identity.matches(current_uid, current_gid) {
1236            return Ok(RunIdentity::Inherit);
1237        }
1238        return Ok(identity);
1239    }
1240
1241    if current_uid.is_root()
1242        && let Some(identity) = resolve_sudo_identity(sudo_uid, sudo_gid)
1243    {
1244        return Ok(identity);
1245    }
1246
1247    Ok(RunIdentity::Inherit)
1248}
1249
1250#[cfg(unix)]
1251fn resolve_configured_user(user: &str) -> Result<RunIdentity> {
1252    if user.chars().all(|c| c.is_ascii_digit()) {
1253        let uid = user
1254            .parse::<u32>()
1255            .map_err(|e| miette::miette!("invalid run user UID '{}': {}", user, e))?;
1256        let user_record = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1257            .into_diagnostic()?
1258            .ok_or_else(|| miette::miette!("run user UID '{}' does not exist", user))?;
1259        return run_identity_from_user_record(user_record);
1260    }
1261
1262    let user_record = nix::unistd::User::from_name(user)
1263        .into_diagnostic()?
1264        .ok_or_else(|| miette::miette!("run user '{}' does not exist", user))?;
1265    run_identity_from_user_record(user_record)
1266}
1267
1268#[cfg(unix)]
1269fn run_identity_from_user_record(user: nix::unistd::User) -> Result<RunIdentity> {
1270    let username = CString::new(user.name)
1271        .map_err(|e| miette::miette!("run user name contains an interior nul byte: {}", e))?;
1272    Ok(RunIdentity::Switch {
1273        uid: user.uid,
1274        gid: user.gid,
1275        username: Some(username),
1276    })
1277}
1278
1279#[cfg(unix)]
1280fn run_identity_from_raw_ids(uid: u32, gid: u32, username: Option<CString>) -> RunIdentity {
1281    RunIdentity::Switch {
1282        uid: nix::unistd::Uid::from_raw(uid),
1283        gid: nix::unistd::Gid::from_raw(gid),
1284        username,
1285    }
1286}
1287
1288#[cfg(unix)]
1289fn resolve_sudo_identity(sudo_uid: Option<&str>, sudo_gid: Option<&str>) -> Option<RunIdentity> {
1290    let uid = sudo_uid?.parse::<u32>().ok()?;
1291    let gid = sudo_gid?.parse::<u32>().ok()?;
1292    let username = nix::unistd::User::from_uid(nix::unistd::Uid::from_raw(uid))
1293        .ok()
1294        .flatten()
1295        .and_then(|u| CString::new(u.name).ok());
1296    Some(run_identity_from_raw_ids(uid, gid, username))
1297}
1298
1299#[cfg(unix)]
1300fn ensure_can_use_identity(
1301    configured_user: &str,
1302    identity: &RunIdentity,
1303    current_uid: nix::unistd::Uid,
1304    current_gid: nix::unistd::Gid,
1305) -> Result<()> {
1306    let RunIdentity::Switch { uid, gid, .. } = identity else {
1307        return Ok(());
1308    };
1309    if *uid == current_uid && *gid == current_gid {
1310        return Ok(());
1311    }
1312    if current_uid.is_root() {
1313        return Ok(());
1314    }
1315    Err(miette::miette!(
1316        "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.",
1317        configured_user,
1318        current_uid.as_raw(),
1319        current_gid.as_raw(),
1320        uid.as_raw(),
1321        gid.as_raw()
1322    ))
1323}
1324
1325#[cfg(unix)]
1326fn apply_run_identity(identity: &RunIdentity) -> std::io::Result<()> {
1327    let RunIdentity::Switch { uid, gid, username } = identity else {
1328        return Ok(());
1329    };
1330    if let Some(username) = username {
1331        initgroups_for_user(username, *gid)?;
1332    } else {
1333        setgroups_to_primary(*gid)?;
1334    }
1335    nix::unistd::setgid(*gid).map_err(nix_to_io_error)?;
1336    nix::unistd::setuid(*uid).map_err(nix_to_io_error)?;
1337    Ok(())
1338}
1339
1340#[cfg(unix)]
1341impl RunIdentity {
1342    fn matches(&self, uid: nix::unistd::Uid, gid: nix::unistd::Gid) -> bool {
1343        matches!(self, RunIdentity::Switch { uid: u, gid: g, .. } if *u == uid && *g == gid)
1344    }
1345}
1346
1347#[cfg(unix)]
1348fn setgroups_to_primary(gid: nix::unistd::Gid) -> std::io::Result<()> {
1349    let groups = [gid.as_raw() as libc::gid_t];
1350    #[cfg(any(target_os = "linux", target_os = "android"))]
1351    let group_count = groups.len();
1352    #[cfg(not(any(target_os = "linux", target_os = "android")))]
1353    let group_count = groups.len() as libc::c_int;
1354    let rc = unsafe { libc::setgroups(group_count, groups.as_ptr()) };
1355    if rc == -1 {
1356        Err(std::io::Error::last_os_error())
1357    } else {
1358        Ok(())
1359    }
1360}
1361
1362#[cfg(unix)]
1363fn initgroups_for_user(username: &CString, gid: nix::unistd::Gid) -> std::io::Result<()> {
1364    let gid = gid.as_raw();
1365    #[cfg(any(
1366        target_os = "macos",
1367        target_os = "ios",
1368        target_os = "tvos",
1369        target_os = "watchos"
1370    ))]
1371    let base_gid = i32::try_from(gid)
1372        .map_err(|_| std::io::Error::other(format!("gid {gid} is out of range")))?;
1373
1374    #[cfg(not(any(
1375        target_os = "macos",
1376        target_os = "ios",
1377        target_os = "tvos",
1378        target_os = "watchos"
1379    )))]
1380    let base_gid = gid as libc::gid_t;
1381
1382    // SAFETY: `username` is a valid nul-terminated C string and `base_gid`
1383    // is derived from a resolved system account or sudo-provided gid.
1384    let rc = unsafe { libc::initgroups(username.as_ptr(), base_gid) };
1385    if rc == -1 {
1386        Err(std::io::Error::last_os_error())
1387    } else {
1388        Ok(())
1389    }
1390}
1391
1392#[cfg(unix)]
1393fn nix_to_io_error(err: nix::errno::Errno) -> std::io::Error {
1394    std::io::Error::from_raw_os_error(err as i32)
1395}
1396
1397/// Check if multiple ports are available and optionally auto-bump to find available ports.
1398///
1399/// All ports are bumped by the same offset to maintain relative port spacing.
1400/// Returns the resolved ports (either the original or bumped ones).
1401/// Returns an error if any port is in use and auto_bump is disabled,
1402/// or if no available ports can be found after max attempts.
1403async fn check_ports_available(
1404    expected_ports: &[u16],
1405    auto_bump: bool,
1406    max_attempts: u32,
1407) -> Result<Vec<u16>> {
1408    if expected_ports.is_empty() {
1409        return Ok(Vec::new());
1410    }
1411
1412    for bump_offset in 0..=max_attempts {
1413        // Use wrapping_add to handle overflow correctly - ports wrap around at 65535
1414        let candidate_ports: Vec<u16> = expected_ports
1415            .iter()
1416            .map(|&p| p.wrapping_add(bump_offset as u16))
1417            .collect();
1418
1419        // Check if all ports in this set are available
1420        let mut all_available = true;
1421        let mut conflicting_port = None;
1422
1423        for &port in &candidate_ports {
1424            // Port 0 is a special case - it requests an ephemeral port from the OS.
1425            // Skip the availability check for port 0 since binding to it always succeeds.
1426            if port == 0 {
1427                continue;
1428            }
1429
1430            // Use spawn_blocking to avoid blocking the async runtime during TCP bind checks.
1431            //
1432            // We check multiple addresses to avoid false-negatives caused by SO_REUSEADDR.
1433            // On macOS/BSD, Rust's TcpListener::bind sets SO_REUSEADDR by default, which
1434            // allows binding 0.0.0.0:port even when 127.0.0.1:port is already in use
1435            // (because 0.0.0.0 is technically a different address).  Most daemons bind
1436            // to localhost, so checking 127.0.0.1 is essential to detect real conflicts.
1437            // We also check [::1] to cover IPv6 loopback listeners.
1438            //
1439            // NOTE: This check has a time-of-check-to-time-of-use (TOCTOU) race condition.
1440            // Another process could grab the port between our check and the daemon actually
1441            // binding. This is inherent to the approach and acceptable for our use case
1442            // since we're primarily detecting conflicts with already-running daemons.
1443            if is_port_in_use(port).await {
1444                all_available = false;
1445                conflicting_port = Some(port);
1446                break;
1447            }
1448        }
1449
1450        if all_available {
1451            // Check for overflow (port wrapped around to 0 due to wrapping_add)
1452            // If any candidate port is 0 but the original expected port wasn't 0,
1453            // it means we've wrapped around and should stop
1454            if candidate_ports.contains(&0) && !expected_ports.contains(&0) {
1455                return Err(PortError::NoAvailablePort {
1456                    start_port: expected_ports[0],
1457                    attempts: bump_offset + 1,
1458                }
1459                .into());
1460            }
1461            if bump_offset > 0 {
1462                info!("ports {expected_ports:?} bumped by {bump_offset} to {candidate_ports:?}");
1463            }
1464            return Ok(candidate_ports);
1465        }
1466
1467        // Port is in use
1468        if bump_offset == 0 && !auto_bump {
1469            if let Some(port) = conflicting_port {
1470                let (pid, process) = identify_port_owner(port).await;
1471                return Err(PortError::InUse { port, process, pid }.into());
1472            }
1473        }
1474    }
1475
1476    // No available ports found after max attempts
1477    Err(PortError::NoAvailablePort {
1478        start_port: expected_ports[0],
1479        attempts: max_attempts + 1,
1480    }
1481    .into())
1482}
1483
1484/// Check whether a port is currently in use by attempting to bind on multiple addresses.
1485///
1486/// Returns `true` when at least one bind attempt gets `AddrInUse`, meaning another
1487/// process is listening.  Other errors (e.g. `AddrNotAvailable` on an address family
1488/// the OS doesn't support) are ignored so they don't produce false positives.
1489async fn is_port_in_use(port: u16) -> bool {
1490    tokio::task::spawn_blocking(move || {
1491        for &addr in &["0.0.0.0", "127.0.0.1", "::1"] {
1492            match std::net::TcpListener::bind((addr, port)) {
1493                Ok(listener) => drop(listener),
1494                Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => return true,
1495                Err(_) => continue,
1496            }
1497        }
1498        false
1499    })
1500    .await
1501    .unwrap_or(false)
1502}
1503
1504/// Best-effort lookup of the process occupying a port via `listeners::get_all()`.
1505///
1506/// Returns `(pid, process_name)`.  Falls back to `(0, "unknown")` when the
1507/// system call fails (permission error, unsupported OS, etc.).
1508async fn identify_port_owner(port: u16) -> (u32, String) {
1509    tokio::task::spawn_blocking(move || {
1510        listeners::get_all()
1511            .ok()
1512            .and_then(|list| {
1513                list.into_iter()
1514                    .find(|l| l.socket.port() == port)
1515                    .map(|l| (l.process.pid, l.process.name))
1516            })
1517            .unwrap_or((0, "unknown".to_string()))
1518    })
1519    .await
1520    .unwrap_or((0, "unknown".to_string()))
1521}
1522
1523/// Detect whether a port is in use, and if so, identify the owning process.
1524///
1525/// Combines `is_port_in_use` (reliable bind probe) with `identify_port_owner`
1526/// (best-effort process lookup).  Returns `None` when the port is free.
1527async fn detect_port_conflict(port: u16) -> Option<(u32, String)> {
1528    if !is_port_in_use(port).await {
1529        return None;
1530    }
1531    Some(identify_port_owner(port).await)
1532}
1533
1534/// Spawn a background task that detects the first port the daemon process is listening on
1535/// and stores it in the state file as `active_port`.
1536///
1537/// This is called once when the daemon becomes ready. The port is cleared when the daemon stops.
1538///
1539/// Port selection strategy:
1540/// 1. If the daemon has `expected_port` configured, prefer the first port from that list
1541///    (it is the port the operator explicitly designated as the primary service port).
1542/// 2. Otherwise, take the first port the process is actually listening on (in the order
1543///    returned by the OS), which is typically the port bound earliest.
1544///
1545/// Using `min()` (lowest port number) was previously used here but is incorrect: many
1546/// applications listen on multiple ports (e.g. HTTP + metrics) and the lowest-numbered
1547/// port is not necessarily the primary service port.
1548fn detect_and_store_active_port(id: DaemonId, pid: u32) {
1549    tokio::spawn(async move {
1550        // Retry with exponential backoff so that slow-starting daemons (JVM,
1551        // Node.js, Python, etc.) that take more than 500 ms to bind their port
1552        // are still detected.  Total wait budget: 500+1000+2000+4000 = 7.5 s.
1553        for delay_ms in [500u64, 1000, 2000, 4000] {
1554            tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
1555
1556            // Read daemon state atomically: check if still alive and get expected_port
1557            // in a single lock acquisition to avoid TOCTOU and unnecessary lock overhead.
1558            let expected_port: Option<u16> = {
1559                let state_file = SUPERVISOR.state_file.lock().await;
1560                match state_file.daemons.get(&id) {
1561                    Some(d) if d.pid.is_none() => {
1562                        debug!("daemon {id}: aborting active_port detection — process exited");
1563                        return;
1564                    }
1565                    Some(d) => d
1566                        .port
1567                        .as_ref()
1568                        .and_then(|p| p.expect.first().copied())
1569                        .filter(|&p| p > 0),
1570                    None => None,
1571                }
1572            };
1573
1574            let active_port = tokio::task::spawn_blocking(move || {
1575                let listeners = listeners::get_all().ok()?;
1576
1577                // Refresh process tree so all_children sees current descendants.
1578                PROCS.refresh_processes();
1579
1580                let descendant_pids: std::collections::HashSet<u32> = PROCS
1581                    .all_children(pid)
1582                    .into_iter()
1583                    .chain(std::iter::once(pid))
1584                    .collect();
1585
1586                let process_ports: Vec<u16> = listeners
1587                    .into_iter()
1588                    .filter(|listener| descendant_pids.contains(&listener.process.pid))
1589                    .map(|listener| listener.socket.port())
1590                    .filter(|&port| port > 0)
1591                    .collect();
1592
1593                if process_ports.is_empty() {
1594                    return None;
1595                }
1596
1597                // Prefer the configured expected_port if the process is actually
1598                // listening on it; otherwise fall back to the first port found.
1599                if let Some(ep) = expected_port {
1600                    if process_ports.contains(&ep) {
1601                        return Some(ep);
1602                    }
1603                }
1604
1605                // No expected_port match — return the first port in the list.
1606                // The list order reflects the order the OS reports listeners,
1607                // which is generally the order they were bound (earliest first).
1608                // Do NOT sort: the lowest-numbered port is not necessarily the
1609                // primary service port (e.g. HTTP vs metrics).
1610                process_ports.into_iter().next()
1611            })
1612            .await
1613            .ok()
1614            .flatten();
1615
1616            if let Some(port) = active_port {
1617                debug!("daemon {id} active_port detected: {port}");
1618                let mut state_file = SUPERVISOR.state_file.lock().await;
1619                if let Some(d) = state_file.daemons.get(&id) {
1620                    // Guard against PID reuse: if the original process exited and the OS
1621                    // assigned the same PID to an unrelated process that happens to bind
1622                    // a port, we must not route proxy traffic to that unrelated service.
1623                    if d.pid == Some(pid) {
1624                        state_file.set_active_port(&id, port);
1625                    } else {
1626                        debug!(
1627                            "daemon {id}: skipping active_port write — PID mismatch \
1628                             (expected {pid}, current {:?})",
1629                            d.pid
1630                        );
1631                        return;
1632                    }
1633                }
1634                return;
1635            }
1636
1637            debug!(
1638                "daemon {id}: no active port detected for pid {pid} or its descendants (will retry)"
1639            );
1640        }
1641
1642        debug!(
1643            "daemon {id}: active port detection exhausted all retries for pid {pid} and its descendants"
1644        );
1645    });
1646}
1647
1648/// Check whether a daemon (by its qualified ID) is the target of any registered
1649/// slug in the global config.  This is used to decide whether to run the
1650/// `detect_and_store_active_port` polling task — only slug-targeted daemons need
1651/// it, avoiding wasted `listeners::get_all()` calls for port-less daemons.
1652///
1653/// Delegates to `proxy::server::is_slug_target()` which uses the same in-memory
1654/// slug cache as the proxy hot path, so this check is cheap.
1655fn is_daemon_slug_target(id: &DaemonId) -> bool {
1656    // read_global_slugs is called once per daemon start — acceptable cost.
1657    // We intentionally avoid making this async to keep has_port_config evaluation
1658    // simple and synchronous in run_once().
1659    let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
1660    slugs.iter().any(|(slug, entry)| {
1661        let daemon_name = entry.daemon.as_deref().unwrap_or(slug);
1662        id.name() == daemon_name
1663    })
1664}
1665
1666#[cfg(all(test, unix))]
1667mod tests {
1668    use super::*;
1669
1670    #[test]
1671    fn test_resolve_run_identity_empty_without_sudo() {
1672        let identity = resolve_run_identity(None, 501, 20, None, None).unwrap();
1673        assert_eq!(identity, RunIdentity::Inherit);
1674    }
1675
1676    #[test]
1677    fn test_resolve_run_identity_sudo_fallback() {
1678        let identity = resolve_run_identity(None, 0, 0, Some("501"), Some("20")).unwrap();
1679        let RunIdentity::Switch { uid, gid, .. } = identity else {
1680            panic!("expected identity switch");
1681        };
1682        assert_eq!(uid.as_raw(), 501);
1683        assert_eq!(gid.as_raw(), 20);
1684    }
1685
1686    #[test]
1687    fn test_resolve_run_identity_ignores_stale_sudo_when_not_root() {
1688        let identity = resolve_run_identity(None, 501, 20, Some("0"), Some("0")).unwrap();
1689        assert_eq!(identity, RunIdentity::Inherit);
1690    }
1691
1692    #[test]
1693    fn test_resolve_configured_user_root_name() {
1694        let identity = resolve_configured_user("root").unwrap();
1695        let RunIdentity::Switch { uid, username, .. } = identity else {
1696            panic!("expected identity switch");
1697        };
1698        assert_eq!(uid.as_raw(), 0);
1699        assert_eq!(
1700            username.as_deref().and_then(|s| s.to_str().ok()),
1701            Some("root")
1702        );
1703    }
1704
1705    #[test]
1706    fn test_resolve_configured_user_root_uid() {
1707        let identity = resolve_configured_user("0").unwrap();
1708        let RunIdentity::Switch { uid, username, .. } = identity else {
1709            panic!("expected identity switch");
1710        };
1711        assert_eq!(uid.as_raw(), 0);
1712        assert_eq!(
1713            username.as_deref().and_then(|s| s.to_str().ok()),
1714            Some("root")
1715        );
1716    }
1717
1718    #[test]
1719    fn test_resolve_configured_user_missing_user_fails() {
1720        let err = resolve_configured_user("pitchfork-user-that-should-not-exist")
1721            .unwrap_err()
1722            .to_string();
1723        assert!(err.contains("does not exist"));
1724    }
1725
1726    #[test]
1727    fn test_resolve_run_identity_requires_root_for_user_switch() {
1728        let err = resolve_run_identity(Some("root"), 501, 20, None, None)
1729            .unwrap_err()
1730            .to_string();
1731        assert!(err.contains("Restart the supervisor with sudo"));
1732    }
1733
1734    #[test]
1735    fn test_resolve_run_identity_same_user_is_noop() {
1736        let identity = resolve_run_identity(Some("root"), 0, 0, Some("501"), Some("20")).unwrap();
1737        assert_eq!(identity, RunIdentity::Inherit);
1738    }
1739}
1740
1741/// Inject proxy-related environment variables into a daemon's command.
1742///
1743/// Adds:
1744/// - `HOST` — the address the daemon should bind to (`127.0.0.1`, omitted in LAN mode)
1745/// - `PITCHFORK_URL` — the public proxy URL for this daemon (if it has a slug)
1746/// - `NODE_EXTRA_CA_CERTS` — path to the pitchfork CA cert (if HTTPS enabled)
1747/// - `__VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS` — `.<tld>` for Vite host allowlisting
1748/// - `PITCHFORK_LAN` — set to `"1"` when LAN mode is active
1749fn inject_proxy_env(cmd: &mut tokio::process::Command, slug: &Option<String>) {
1750    let s = crate::settings::settings();
1751    let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
1752
1753    if should_force_loopback_host(slug) && !lan_enabled {
1754        // Only force loopback binding for daemons that are actually routed via a slug.
1755        // In LAN mode, daemons need to bind to 0.0.0.0 to be reachable from the network.
1756        cmd.env("HOST", "127.0.0.1");
1757    }
1758
1759    // PITCHFORK_URL: the daemon's public proxy URL (only if it has a slug and proxy is enabled)
1760    if let Some(url) = build_pitchfork_url(slug, &s) {
1761        cmd.env("PITCHFORK_URL", &url);
1762    }
1763
1764    // NODE_EXTRA_CA_CERTS: let Node.js backends trust the pitchfork CA
1765    if s.proxy.enable && s.proxy.https {
1766        let ca_path = if s.proxy.tls_cert.is_empty() {
1767            crate::env::PITCHFORK_STATE_DIR.join("proxy").join("ca.pem")
1768        } else {
1769            std::path::PathBuf::from(&s.proxy.tls_cert)
1770        };
1771        if ca_path.exists() {
1772            cmd.env("NODE_EXTRA_CA_CERTS", ca_path.to_string_lossy().to_string());
1773        }
1774    }
1775
1776    // __VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS: Vite host allowlisting
1777    if s.proxy.enable {
1778        let tld = if lan_enabled { "local" } else { &s.proxy.tld };
1779        cmd.env("__VITE_ADDITIONAL_SERVER_ALLOWED_HOSTS", format!(".{tld}"));
1780    }
1781
1782    // PITCHFORK_LAN: signal to daemons that LAN mode is active
1783    if lan_enabled {
1784        cmd.env("PITCHFORK_LAN", "1");
1785    }
1786}
1787
1788fn should_force_loopback_host(slug: &Option<String>) -> bool {
1789    let Some(slug) = slug.as_deref() else {
1790        return false;
1791    };
1792
1793    let s = crate::settings::settings();
1794    if !s.proxy.enable {
1795        return false;
1796    }
1797
1798    let slugs = crate::pitchfork_toml::PitchforkToml::read_global_slugs();
1799    slugs.contains_key(slug)
1800}
1801
1802/// Compute the public proxy URL for a daemon.
1803///
1804/// Returns `None` if the daemon has no slug or the proxy is not enabled.
1805fn build_pitchfork_url(slug: &Option<String>, s: &crate::settings::Settings) -> Option<String> {
1806    let slug = slug.as_ref()?;
1807    if !s.proxy.enable {
1808        return None;
1809    }
1810    let scheme = if s.proxy.https { "https" } else { "http" };
1811    let port = u16::try_from(s.proxy.port).ok().filter(|&p| p > 0)?;
1812    let port_suffix = if (scheme == "https" && port == 443) || (scheme == "http" && port == 80) {
1813        String::new()
1814    } else {
1815        format!(":{port}")
1816    };
1817    let lan_enabled = s.proxy.lan || !s.proxy.lan_ip.is_empty();
1818    let tld = if lan_enabled { "local" } else { &s.proxy.tld };
1819    Some(format!("{scheme}://{slug}.{tld}{port_suffix}",))
1820}