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