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