Skip to main content

pitchfork_cli/supervisor/
state.rs

1//! State access layer for the supervisor
2//!
3//! All state getter/setter operations for daemons, shell directories, and notifications.
4
5use super::Supervisor;
6use crate::Result;
7use crate::daemon::Daemon;
8use crate::daemon::RunOptions;
9use crate::daemon_id::DaemonId;
10use crate::daemon_status::DaemonStatus;
11use crate::error::FileError;
12use crate::pitchfork_toml::CpuLimit;
13use crate::pitchfork_toml::CronRetrigger;
14use crate::pitchfork_toml::MemoryLimit;
15use crate::pitchfork_toml::PitchforkToml;
16use crate::pitchfork_toml::PortConfig;
17use crate::pitchfork_toml::ReadyCmd;
18use crate::pitchfork_toml::ReadyHttp;
19use crate::pitchfork_toml::ReadyOutput;
20use crate::pitchfork_toml::ReadyPort;
21use crate::pitchfork_toml::Retry;
22use crate::pitchfork_toml::StopConfig;
23use crate::pitchfork_toml::WatchMode;
24use crate::procs::PROCS;
25use indexmap::IndexMap;
26use std::collections::{HashMap, HashSet};
27use std::path::PathBuf;
28
29fn should_clean_daemon(
30    id: &DaemonId,
31    daemon: &Daemon,
32    namespaces: &[String],
33    daemons: &[DaemonId],
34) -> bool {
35    let namespace_matches =
36        namespaces.is_empty() || namespaces.iter().any(|ns| ns == id.namespace());
37    let daemon_matches = daemons.is_empty() || daemons.contains(id);
38    daemon.pid.is_none() && namespace_matches && daemon_matches
39}
40
41fn prune_candidate(
42    id: &DaemonId,
43    daemon: &Daemon,
44    namespaces: &[String],
45    daemons: &[DaemonId],
46) -> Option<(DaemonId, PathBuf)> {
47    should_clean_daemon(id, daemon, namespaces, daemons)
48        .then(|| daemon.dir.clone())
49        .flatten()
50        .map(|dir| (id.clone(), dir))
51}
52
53/// Options for upserting a daemon's state.
54///
55/// Use `UpsertDaemonOpts::builder(id)` to create, then set fields directly and call `.build()`.
56#[derive(Debug, Default)]
57pub(crate) struct UpsertDaemonOpts {
58    pub id: DaemonId,
59    pub pid: Option<u32>,
60    pub status: DaemonStatus,
61    pub shell_pid: Option<u32>,
62    pub dir: Option<PathBuf>,
63    pub cmd: Option<Vec<String>>,
64    pub run: Option<String>,
65    pub autostop: bool,
66    pub cron_schedule: Option<String>,
67    pub cron_retrigger: Option<CronRetrigger>,
68    pub cron_immediate: Option<bool>,
69    pub last_exit_success: Option<bool>,
70    pub retry: Option<Retry>,
71    pub retry_count: Option<u32>,
72    pub ready_delay: Option<u64>,
73    pub ready_output: Option<ReadyOutput>,
74    pub ready_http: Option<ReadyHttp>,
75    pub ready_port: Option<ReadyPort>,
76    pub ready_cmd: Option<ReadyCmd>,
77    /// Port configuration
78    pub port: Option<PortConfig>,
79    /// Resolved ports actually used after auto-bump (may differ from expected).
80    /// `None` inherits the existing record's value; `Some` sets it explicitly,
81    /// including an empty vec, which clears ports left by a previous run.
82    pub resolved_port: Option<Vec<u16>>,
83    /// The first port the process is actually listening on (detected at runtime).
84    pub active_port: Option<u16>,
85    /// Optional stable slug alias for this daemon.
86    pub slug: Option<String>,
87    /// Whether to proxy this daemon (None = use global proxy.enable setting).
88    pub proxy: Option<bool>,
89    pub depends: Option<Vec<DaemonId>>,
90    pub env: Option<IndexMap<String, String>>,
91    pub watch: Option<Vec<String>>,
92    pub watch_mode: Option<WatchMode>,
93    pub watch_base_dir: Option<PathBuf>,
94    pub mise: Option<bool>,
95    /// Unix user to run this daemon as
96    pub user: Option<String>,
97    /// Memory limit for the daemon process
98    pub memory_limit: Option<MemoryLimit>,
99    /// CPU usage limit as a percentage
100    pub cpu_limit: Option<CpuLimit>,
101    /// Unix signal to send for graceful shutdown
102    pub stop_signal: Option<StopConfig>,
103    /// Archive hook command invoked before retention prunes this daemon's logs.
104    pub archive_hook: Option<String>,
105    /// Log format for this daemon.
106    pub log_format: Option<String>,
107    /// Allocate a pseudo-terminal for the daemon process.
108    pub pty: Option<bool>,
109    /// True for config-only cron daemons auto-registered into state.
110    pub config_registered: bool,
111}
112
113/// Builder for UpsertDaemonOpts - ensures daemon ID is always provided.
114///
115/// # Example
116/// ```ignore
117/// let opts = UpsertDaemonOpts::builder(daemon_id)
118///     .set(|o| {
119///         o.pid = Some(pid);
120///         o.status = DaemonStatus::Running;
121///     })
122///     .build();
123/// ```
124#[derive(Debug)]
125pub(crate) struct UpsertDaemonOptsBuilder {
126    pub opts: UpsertDaemonOpts,
127}
128
129impl UpsertDaemonOpts {
130    /// Create a builder with the required daemon ID.
131    pub fn builder(id: DaemonId) -> UpsertDaemonOptsBuilder {
132        UpsertDaemonOptsBuilder {
133            opts: UpsertDaemonOpts {
134                id,
135                ..Default::default()
136            },
137        }
138    }
139
140    /// Build an `UpsertDaemonOptsBuilder` from `RunOptions`, mapping all
141    /// config-carried fields. Callers chain `.set()` to add runtime-specific
142    /// fields (pid, resolved_port, etc.) and then call `.build()`.
143    pub(crate) fn from_run_options(
144        opts: &RunOptions,
145        status: DaemonStatus,
146    ) -> UpsertDaemonOptsBuilder {
147        UpsertDaemonOpts::builder(opts.id.clone()).set(|o| {
148            o.status = status;
149            o.shell_pid = opts.shell_pid;
150            o.dir = Some(opts.dir.0.clone());
151            o.cmd = Some(opts.cmd.clone());
152            o.run = opts.run.clone();
153            o.autostop = opts.autostop;
154            o.cron_schedule = opts.cron_schedule.clone();
155            o.cron_retrigger = opts.cron_retrigger;
156            o.cron_immediate = opts.cron_immediate;
157            o.retry = Some(opts.retry);
158            o.retry_count = Some(opts.retry_count);
159            o.ready_delay = opts.ready_delay;
160            o.ready_output = opts.ready_output.clone();
161            o.ready_http = opts.ready_http.clone();
162            o.ready_port = opts.ready_port.clone();
163            o.ready_cmd = opts.ready_cmd.clone();
164            o.port = opts.port.clone();
165            o.depends = Some(opts.depends.clone());
166            o.env = opts.env.clone();
167            o.watch = Some(opts.watch.clone());
168            o.watch_mode = Some(opts.watch_mode);
169            o.watch_base_dir = opts.watch_base_dir.clone();
170            o.mise = opts.mise;
171            o.user = opts.user.clone();
172            o.memory_limit = opts.memory_limit;
173            o.cpu_limit = opts.cpu_limit;
174            o.stop_signal = opts.stop_signal;
175            o.pty = opts.pty;
176            o.archive_hook = opts.archive_hook.clone();
177            o.log_format = opts.log_format.clone();
178        })
179    }
180}
181
182impl UpsertDaemonOptsBuilder {
183    /// Modify opts fields with a closure.
184    pub fn set<F: FnOnce(&mut UpsertDaemonOpts)>(mut self, f: F) -> Self {
185        f(&mut self.opts);
186        self
187    }
188
189    /// Build the UpsertDaemonOpts.
190    pub fn build(self) -> UpsertDaemonOpts {
191        self.opts
192    }
193}
194
195impl Supervisor {
196    /// Upsert a daemon's state, merging with existing values
197    pub(crate) async fn upsert_daemon(&self, opts: UpsertDaemonOpts) -> Result<Daemon> {
198        info!(
199            "upserting daemon: {} pid: {} status: {}",
200            opts.id,
201            opts.pid.unwrap_or(0),
202            opts.status
203        );
204        let mut state_file = self.state_file.lock().await;
205        let existing = state_file.daemons.get(&opts.id);
206        let daemon = Daemon {
207            id: opts.id.clone(),
208            // title/start_time identify the process for orphan cleanup after a
209            // supervisor crash. They are looked up from the process cache; if
210            // the cache has no entry (e.g. an upsert between refreshes) fall
211            // back to the existing values as long as the PID is unchanged, so
212            // the recorded identity is not wiped mid-lifetime.
213            title: opts.pid.and_then(|pid| {
214                PROCS.title(pid).or_else(|| {
215                    existing
216                        .filter(|d| d.pid == Some(pid))
217                        .and_then(|d| d.title.clone())
218                })
219            }),
220            start_time: opts.pid.and_then(|pid| {
221                PROCS.start_time(pid).or_else(|| {
222                    existing
223                        .filter(|d| d.pid == Some(pid))
224                        .and_then(|d| d.start_time)
225                })
226            }),
227            // Boot time is a system-wide constant, so it needs no cache
228            // fallback — it is recorded for any daemon with a live PID and
229            // dropped when the PID is cleared.
230            boot_time: opts.pid.map(|_| PROCS.boot_time()),
231            pid: opts.pid,
232            status: opts.status,
233            shell_pid: opts.shell_pid,
234            autostop: opts.autostop || existing.is_some_and(|d| d.autostop),
235            dir: opts.dir.or(existing.and_then(|d| d.dir.clone())),
236            cmd: opts.cmd.or(existing.and_then(|d| d.cmd.clone())),
237            run: opts.run.or(existing.and_then(|d| d.run.clone())),
238            cron_schedule: opts
239                .cron_schedule
240                .or(existing.and_then(|d| d.cron_schedule.clone())),
241            cron_retrigger: opts
242                .cron_retrigger
243                .or(existing.and_then(|d| d.cron_retrigger)),
244            cron_immediate: opts
245                .cron_immediate
246                .or(existing.and_then(|d| d.cron_immediate)),
247            last_cron_triggered: existing.and_then(|d| d.last_cron_triggered),
248            last_exit_success: opts
249                .last_exit_success
250                .or(existing.and_then(|d| d.last_exit_success)),
251            retry: opts
252                .retry
253                .unwrap_or_else(|| existing.map(|d| d.retry).unwrap_or_default()),
254            retry_count: opts
255                .retry_count
256                .unwrap_or(existing.map(|d| d.retry_count).unwrap_or(0)),
257            ready_delay: opts.ready_delay.or(existing.and_then(|d| d.ready_delay)),
258            ready_output: opts
259                .ready_output
260                .or(existing.and_then(|d| d.ready_output.clone())),
261            ready_http: opts
262                .ready_http
263                .or(existing.and_then(|d| d.ready_http.clone())),
264            ready_port: opts
265                .ready_port
266                .or(existing.and_then(|d| d.ready_port.clone())),
267            ready_cmd: opts
268                .ready_cmd
269                .or(existing.and_then(|d| d.ready_cmd.clone())),
270            port: opts.port.or_else(|| existing.and_then(|d| d.port.clone())),
271            resolved_port: match opts.resolved_port {
272                Some(ports) => ports,
273                None => existing
274                    .map(|d| d.resolved_port.clone())
275                    .unwrap_or_default(),
276            },
277            depends: opts
278                .depends
279                .unwrap_or_else(|| existing.map(|d| d.depends.clone()).unwrap_or_default()),
280            env: opts.env.or(existing.and_then(|d| d.env.clone())),
281            watch: opts
282                .watch
283                .unwrap_or_else(|| existing.map(|d| d.watch.clone()).unwrap_or_default()),
284            watch_mode: opts
285                .watch_mode
286                .unwrap_or_else(|| existing.map(|d| d.watch_mode).unwrap_or_default()),
287            watch_base_dir: opts
288                .watch_base_dir
289                .or(existing.and_then(|d| d.watch_base_dir.clone())),
290            mise: opts.mise.or(existing.and_then(|d| d.mise)),
291            user: opts.user.or(existing.and_then(|d| d.user.clone())),
292            proxy: opts.proxy.or(existing.and_then(|d| d.proxy)),
293            // active_port is intentionally NOT inherited from the existing daemon.
294            // When a daemon restarts, the new process has not yet bound a port, so
295            // carrying over the old process's active_port would cause the proxy to
296            // route to a port that is no longer listening.  The port will be
297            // re-detected by detect_and_store_active_port once the new process is ready.
298            active_port: opts.active_port,
299            slug: opts.slug.or(existing.and_then(|d| d.slug.clone())),
300            memory_limit: opts.memory_limit.or(existing.and_then(|d| d.memory_limit)),
301            cpu_limit: opts.cpu_limit.or(existing.and_then(|d| d.cpu_limit)),
302            stop_signal: opts.stop_signal.or(existing.and_then(|d| d.stop_signal)),
303            archive_hook: opts
304                .archive_hook
305                .or(existing.and_then(|d| d.archive_hook.clone())),
306            log_format: opts
307                .log_format
308                .or(existing.and_then(|d| d.log_format.clone())),
309            pty: opts.pty.or(existing.and_then(|d| d.pty)),
310            config_registered: opts.config_registered,
311        };
312        state_file.insert_daemon(&opts.id, daemon.clone());
313        Ok(daemon)
314    }
315
316    /// Enable a daemon (remove from disabled set)
317    pub async fn enable(&self, id: &DaemonId) -> Result<bool> {
318        info!("enabling daemon: {id}");
319        let config = PitchforkToml::all_merged_all_namespaces()?;
320        let mut state_file = self.state_file.lock().await;
321        let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
322        if !exists {
323            return Err(miette::miette!("daemon '{}' not found", id));
324        }
325        let result = state_file.enable_daemon(id);
326        Ok(result)
327    }
328
329    /// Disable a daemon (add to disabled set)
330    pub async fn disable(&self, id: &DaemonId) -> Result<bool> {
331        info!("disabling daemon: {id}");
332        let config = PitchforkToml::all_merged_all_namespaces()?;
333        let mut state_file = self.state_file.lock().await;
334        let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
335        if !exists {
336            return Err(miette::miette!("daemon '{}' not found", id));
337        }
338        let result = state_file.disable_daemon(id);
339        Ok(result)
340    }
341
342    /// Get a daemon by ID
343    pub(crate) async fn get_daemon(&self, id: &DaemonId) -> Option<Daemon> {
344        self.state_file.lock().await.daemons.get(id).cloned()
345    }
346
347    /// Get all active daemons (those with PIDs, excluding pitchfork itself)
348    pub(crate) async fn active_daemons(&self) -> Vec<Daemon> {
349        let pitchfork_id = DaemonId::pitchfork();
350        self.state_file
351            .lock()
352            .await
353            .daemons
354            .values()
355            .filter(|d| d.pid.is_some() && d.id != pitchfork_id)
356            .cloned()
357            .collect()
358    }
359
360    /// Remove a daemon from state
361    pub(crate) async fn remove_daemon(&self, id: &DaemonId) -> Result<()> {
362        let mut state_file = self.state_file.lock().await;
363        state_file.remove_daemon(id);
364        Ok(())
365    }
366
367    /// Set the shell's working directory
368    pub(crate) async fn set_shell_dir(&self, shell_pid: u32, dir: PathBuf) -> Result<()> {
369        let mut state_file = self.state_file.lock().await;
370        state_file.set_shell_dir(shell_pid, dir);
371        Ok(())
372    }
373
374    /// Get the shell's working directory
375    pub(crate) async fn get_shell_dir(&self, shell_pid: u32) -> Option<PathBuf> {
376        self.state_file
377            .lock()
378            .await
379            .shell_dirs
380            .get(&shell_pid.to_string())
381            .cloned()
382    }
383
384    /// Remove a shell PID from tracking
385    pub(crate) async fn remove_shell_pid(&self, shell_pid: u32) -> Result<()> {
386        let mut state_file = self.state_file.lock().await;
387        state_file.remove_shell_dir(shell_pid);
388        Ok(())
389    }
390
391    /// Get all directories with their associated shell PIDs
392    pub(crate) async fn get_dirs_with_shell_pids(&self) -> HashMap<PathBuf, Vec<u32>> {
393        self.state_file.lock().await.shell_dirs.iter().fold(
394            HashMap::new(),
395            |mut acc, (pid, dir)| {
396                if let Ok(pid) = pid.parse() {
397                    acc.entry(dir.clone()).or_default().push(pid);
398                }
399                acc
400            },
401        )
402    }
403
404    /// Get pending notifications and clear the queue
405    pub(crate) async fn get_notifications(&self) -> Vec<(log::LevelFilter, String)> {
406        self.pending_notifications.lock().await.drain(..).collect()
407    }
408
409    /// Clean up daemons that have no PID
410    pub(crate) async fn clean(&self) -> Result<()> {
411        self.clean_filtered(&[], &[], false).await?;
412        Ok(())
413    }
414
415    /// Clean up PID-less daemon registrations matching all supplied filters.
416    pub(crate) async fn clean_filtered(
417        &self,
418        namespaces: &[String],
419        daemons: &[DaemonId],
420        prune: bool,
421    ) -> Result<u64> {
422        if prune {
423            let candidates: Vec<(DaemonId, PathBuf)> = {
424                let state_file = self.state_file.lock().await;
425                state_file
426                    .daemons
427                    .iter()
428                    .filter_map(|(id, daemon)| prune_candidate(id, daemon, namespaces, daemons))
429                    .collect()
430            };
431
432            let mut missing = HashMap::new();
433            for (id, dir) in candidates {
434                if !tokio::fs::try_exists(&dir)
435                    .await
436                    .map_err(|source| FileError::ReadError {
437                        path: dir.clone(),
438                        source,
439                    })?
440                {
441                    missing.insert(id, dir);
442                }
443            }
444
445            let mut removed = 0;
446            let mut state_file = self.state_file.lock().await;
447            state_file.retain_daemons(|id, daemon| {
448                let remove = should_clean_daemon(id, daemon, namespaces, daemons)
449                    && missing
450                        .get(id)
451                        .is_some_and(|missing_dir| daemon.dir.as_ref() == Some(missing_dir));
452                removed += u64::from(remove);
453                !remove
454            });
455            return Ok(removed);
456        }
457
458        let mut removed = 0;
459        let mut state_file = self.state_file.lock().await;
460        state_file.retain_daemons(|id, daemon| {
461            let remove = should_clean_daemon(id, daemon, namespaces, daemons);
462            removed += u64::from(remove);
463            !remove
464        });
465        Ok(removed)
466    }
467
468    /// Return the union of active directories from shell tracking and project
469    /// sessions. These are the directories that should keep auto-stop daemons
470    /// alive.
471    pub(crate) async fn get_active_directories(&self) -> Vec<PathBuf> {
472        let state = self.state_file.lock().await;
473        let mut dirs: HashSet<PathBuf> = state.shell_dirs.values().cloned().collect();
474        for (_, dir, _) in state.iter_project_sessions() {
475            dirs.insert(dir.clone());
476        }
477        dirs.into_iter().collect()
478    }
479
480    /// Collect all project sessions as `(pid, dir, liveness_title)`. Every
481    /// project session carries a host PID in its key, so every session is a
482    /// liveness session. Used by the refresh loop to clean up stale sessions.
483    pub(crate) async fn get_liveness_sessions(&self) -> Vec<(u32, PathBuf, Option<String>)> {
484        self.state_file
485            .lock()
486            .await
487            .iter_project_sessions()
488            .into_iter()
489            .filter_map(|(pid_str, dir, session)| {
490                pid_str
491                    .parse()
492                    .ok()
493                    .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
494            })
495            .collect()
496    }
497
498    /// Build a snapshot of all project sessions with live liveness status
499    /// filled in from `PROCS`. Used to answer `GetProjectSessions` IPC
500    /// requests.
501    pub(crate) async fn get_project_sessions_info(&self) -> Vec<crate::ipc::ProjectSessionInfo> {
502        let sessions: Vec<(u32, PathBuf, Option<String>)> = self
503            .state_file
504            .lock()
505            .await
506            .iter_project_sessions()
507            .into_iter()
508            .filter_map(|(pid_str, dir, session)| {
509                pid_str
510                    .parse()
511                    .ok()
512                    .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
513            })
514            .collect();
515        let pids: Vec<u32> = sessions.iter().map(|(pid, _, _)| *pid).collect();
516        if !pids.is_empty() {
517            PROCS.refresh_pids(&pids);
518        }
519        sessions
520            .into_iter()
521            .map(
522                |(pid, directory, liveness_title)| crate::ipc::ProjectSessionInfo {
523                    pid,
524                    directory,
525                    liveness_title,
526                    alive: PROCS.is_running(pid),
527                    current_title: PROCS.title(pid),
528                },
529            )
530            .collect()
531    }
532
533    /// Atomically enter (or replace) a project session for `(pid, dir)`.
534    /// Returns the previous session, if any, so the caller can evaluate the
535    /// previous entry for autostop.
536    pub(crate) async fn enter_project_session(
537        &self,
538        pid: u32,
539        dir: PathBuf,
540    ) -> Result<Option<crate::state_file::ProjectSession>> {
541        if pid == 0 {
542            return Err(miette::miette!("invalid host PID 0"));
543        }
544        if pid > i32::MAX as u32 {
545            return Err(miette::miette!("host PID {pid} exceeds i32::MAX"));
546        }
547        PROCS.refresh_pids(&[pid]);
548        let liveness_title = PROCS.title(pid);
549        // On Unix, reject dead host PIDs up front so we don't register a
550        // session the refresh loop would immediately evict. On Windows, Git
551        // Bash `$$` is a Cygwin-internal PID invisible to sysinfo, so the
552        // liveness check is skipped and sessions rely on explicit
553        // `project leave`.
554        #[cfg(unix)]
555        if !PROCS.is_running(pid) {
556            return Err(miette::miette!("host PID {pid} is not running"));
557        }
558        let mut state_file = self.state_file.lock().await;
559        let previous = state_file.set_project_session(
560            pid,
561            dir,
562            crate::state_file::ProjectSession { liveness_title },
563        );
564        Ok(previous)
565    }
566
567    /// Atomically remove a project session for `(pid, dir)`. Returns the
568    /// directory of the removed session so the caller can evaluate it for
569    /// autostop.
570    pub(crate) async fn leave_project_session(
571        &self,
572        pid: u32,
573        dir: &std::path::Path,
574    ) -> Result<Option<PathBuf>> {
575        let mut state_file = self.state_file.lock().await;
576        if state_file.remove_project_session(pid, dir).is_some() {
577            Ok(Some(dir.to_path_buf()))
578        } else {
579            Ok(None)
580        }
581    }
582}
583
584#[cfg(test)]
585mod tests {
586    use super::*;
587
588    fn daemon(id: &DaemonId, pid: Option<u32>, dir: Option<PathBuf>) -> Daemon {
589        Daemon {
590            id: id.clone(),
591            pid,
592            dir,
593            ..Daemon::default()
594        }
595    }
596
597    #[test]
598    fn clean_filters_intersect_and_preserve_running_daemons() {
599        let api = DaemonId::new("project-a", "api");
600        let worker = DaemonId::new("project-a", "worker");
601        let namespaces = vec!["project-a".to_string()];
602        let daemons = vec![api.clone()];
603
604        assert!(should_clean_daemon(
605            &api,
606            &daemon(&api, None, None),
607            &namespaces,
608            &daemons,
609        ));
610        assert!(!should_clean_daemon(
611            &worker,
612            &daemon(&worker, None, None),
613            &namespaces,
614            &daemons,
615        ));
616        assert!(!should_clean_daemon(
617            &api,
618            &daemon(&api, Some(42), None),
619            &namespaces,
620            &daemons,
621        ));
622    }
623
624    #[test]
625    fn prune_candidates_require_a_recorded_directory() {
626        let id = DaemonId::new("project-a", "api");
627        assert!(prune_candidate(&id, &daemon(&id, None, None), &[], &[]).is_none());
628        assert_eq!(
629            prune_candidate(
630                &id,
631                &daemon(&id, None, Some(PathBuf::from("missing"))),
632                &[],
633                &[]
634            ),
635            Some((id, PathBuf::from("missing")))
636        );
637    }
638}