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::ReadyHttp;
17use crate::pitchfork_toml::Retry;
18use crate::pitchfork_toml::StopConfig;
19use crate::pitchfork_toml::WatchMode;
20use crate::procs::PROCS;
21use indexmap::IndexMap;
22use std::collections::HashMap;
23use std::path::PathBuf;
24
25/// Options for upserting a daemon's state.
26///
27/// Use `UpsertDaemonOpts::builder(id)` to create, then set fields directly and call `.build()`.
28#[derive(Debug, Default)]
29pub(crate) struct UpsertDaemonOpts {
30    pub id: DaemonId,
31    pub pid: Option<u32>,
32    pub status: DaemonStatus,
33    pub shell_pid: Option<u32>,
34    pub dir: Option<PathBuf>,
35    pub cmd: Option<Vec<String>>,
36    pub run: Option<String>,
37    pub autostop: bool,
38    pub cron_schedule: Option<String>,
39    pub cron_retrigger: Option<CronRetrigger>,
40    pub cron_immediate: Option<bool>,
41    pub last_exit_success: Option<bool>,
42    pub retry: Option<Retry>,
43    pub retry_count: Option<u32>,
44    pub ready_delay: Option<u64>,
45    pub ready_output: Option<String>,
46    pub ready_http: Option<ReadyHttp>,
47    pub ready_port: Option<u16>,
48    pub ready_cmd: Option<String>,
49    /// Port configuration
50    pub port: Option<PortConfig>,
51    /// Resolved ports actually used after auto-bump (may differ from expected)
52    pub resolved_port: Vec<u16>,
53    /// The first port the process is actually listening on (detected at runtime).
54    pub active_port: Option<u16>,
55    /// Optional stable slug alias for this daemon.
56    pub slug: Option<String>,
57    /// Whether to proxy this daemon (None = use global proxy.enable setting).
58    pub proxy: Option<bool>,
59    pub depends: Option<Vec<DaemonId>>,
60    pub env: Option<IndexMap<String, String>>,
61    pub watch: Option<Vec<String>>,
62    pub watch_mode: Option<WatchMode>,
63    pub watch_base_dir: Option<PathBuf>,
64    pub mise: Option<bool>,
65    /// Unix user to run this daemon as
66    pub user: Option<String>,
67    /// Memory limit for the daemon process
68    pub memory_limit: Option<MemoryLimit>,
69    /// CPU usage limit as a percentage
70    pub cpu_limit: Option<CpuLimit>,
71    /// Unix signal to send for graceful shutdown
72    pub stop_signal: Option<StopConfig>,
73    /// Archive hook command invoked before retention prunes this daemon's logs.
74    pub archive_hook: Option<String>,
75    /// Allocate a pseudo-terminal for the daemon process.
76    pub pty: Option<bool>,
77    /// True for config-only cron daemons auto-registered into state.
78    pub config_registered: bool,
79}
80
81/// Builder for UpsertDaemonOpts - ensures daemon ID is always provided.
82///
83/// # Example
84/// ```ignore
85/// let opts = UpsertDaemonOpts::builder(daemon_id)
86///     .set(|o| {
87///         o.pid = Some(pid);
88///         o.status = DaemonStatus::Running;
89///     })
90///     .build();
91/// ```
92#[derive(Debug)]
93pub(crate) struct UpsertDaemonOptsBuilder {
94    pub opts: UpsertDaemonOpts,
95}
96
97impl UpsertDaemonOpts {
98    /// Create a builder with the required daemon ID.
99    pub fn builder(id: DaemonId) -> UpsertDaemonOptsBuilder {
100        UpsertDaemonOptsBuilder {
101            opts: UpsertDaemonOpts {
102                id,
103                ..Default::default()
104            },
105        }
106    }
107
108    /// Build an `UpsertDaemonOptsBuilder` from `RunOptions`, mapping all
109    /// config-carried fields. Callers chain `.set()` to add runtime-specific
110    /// fields (pid, resolved_port, etc.) and then call `.build()`.
111    pub(crate) fn from_run_options(
112        opts: &RunOptions,
113        status: DaemonStatus,
114    ) -> UpsertDaemonOptsBuilder {
115        UpsertDaemonOpts::builder(opts.id.clone()).set(|o| {
116            o.status = status;
117            o.shell_pid = opts.shell_pid;
118            o.dir = Some(opts.dir.0.clone());
119            o.cmd = Some(opts.cmd.clone());
120            o.run = opts.run.clone();
121            o.autostop = opts.autostop;
122            o.cron_schedule = opts.cron_schedule.clone();
123            o.cron_retrigger = opts.cron_retrigger;
124            o.cron_immediate = opts.cron_immediate;
125            o.retry = Some(opts.retry);
126            o.retry_count = Some(opts.retry_count);
127            o.ready_delay = opts.ready_delay;
128            o.ready_output = opts.ready_output.clone();
129            o.ready_http = opts.ready_http.clone();
130            o.ready_port = opts.ready_port;
131            o.ready_cmd = opts.ready_cmd.clone();
132            o.port = opts.port.clone();
133            o.depends = Some(opts.depends.clone());
134            o.env = opts.env.clone();
135            o.watch = Some(opts.watch.clone());
136            o.watch_mode = Some(opts.watch_mode);
137            o.watch_base_dir = opts.watch_base_dir.clone();
138            o.mise = opts.mise;
139            o.user = opts.user.clone();
140            o.memory_limit = opts.memory_limit;
141            o.cpu_limit = opts.cpu_limit;
142            o.stop_signal = opts.stop_signal;
143            o.pty = opts.pty;
144            o.archive_hook = opts.archive_hook.clone();
145        })
146    }
147}
148
149impl UpsertDaemonOptsBuilder {
150    /// Modify opts fields with a closure.
151    pub fn set<F: FnOnce(&mut UpsertDaemonOpts)>(mut self, f: F) -> Self {
152        f(&mut self.opts);
153        self
154    }
155
156    /// Build the UpsertDaemonOpts.
157    pub fn build(self) -> UpsertDaemonOpts {
158        self.opts
159    }
160}
161
162impl Supervisor {
163    /// Upsert a daemon's state, merging with existing values
164    pub(crate) async fn upsert_daemon(&self, opts: UpsertDaemonOpts) -> Result<Daemon> {
165        info!(
166            "upserting daemon: {} pid: {} status: {}",
167            opts.id,
168            opts.pid.unwrap_or(0),
169            opts.status
170        );
171        let mut state_file = self.state_file.lock().await;
172        let existing = state_file.daemons.get(&opts.id);
173        let daemon = Daemon {
174            id: opts.id.clone(),
175            title: opts.pid.and_then(|pid| PROCS.title(pid)),
176            pid: opts.pid,
177            status: opts.status,
178            shell_pid: opts.shell_pid,
179            autostop: opts.autostop || existing.is_some_and(|d| d.autostop),
180            dir: opts.dir.or(existing.and_then(|d| d.dir.clone())),
181            cmd: opts.cmd.or(existing.and_then(|d| d.cmd.clone())),
182            run: opts.run.or(existing.and_then(|d| d.run.clone())),
183            cron_schedule: opts
184                .cron_schedule
185                .or(existing.and_then(|d| d.cron_schedule.clone())),
186            cron_retrigger: opts
187                .cron_retrigger
188                .or(existing.and_then(|d| d.cron_retrigger)),
189            cron_immediate: opts
190                .cron_immediate
191                .or(existing.and_then(|d| d.cron_immediate)),
192            last_cron_triggered: existing.and_then(|d| d.last_cron_triggered),
193            last_exit_success: opts
194                .last_exit_success
195                .or(existing.and_then(|d| d.last_exit_success)),
196            retry: opts
197                .retry
198                .unwrap_or_else(|| existing.map(|d| d.retry).unwrap_or_default()),
199            retry_count: opts
200                .retry_count
201                .unwrap_or(existing.map(|d| d.retry_count).unwrap_or(0)),
202            ready_delay: opts.ready_delay.or(existing.and_then(|d| d.ready_delay)),
203            ready_output: opts
204                .ready_output
205                .or(existing.and_then(|d| d.ready_output.clone())),
206            ready_http: opts
207                .ready_http
208                .or(existing.and_then(|d| d.ready_http.clone())),
209            ready_port: opts.ready_port.or(existing.and_then(|d| d.ready_port)),
210            ready_cmd: opts
211                .ready_cmd
212                .or(existing.and_then(|d| d.ready_cmd.clone())),
213            port: opts.port.or_else(|| existing.and_then(|d| d.port.clone())),
214            resolved_port: if opts.resolved_port.is_empty() {
215                existing
216                    .map(|d| d.resolved_port.clone())
217                    .unwrap_or_default()
218            } else {
219                opts.resolved_port
220            },
221            depends: opts
222                .depends
223                .unwrap_or_else(|| existing.map(|d| d.depends.clone()).unwrap_or_default()),
224            env: opts.env.or(existing.and_then(|d| d.env.clone())),
225            watch: opts
226                .watch
227                .unwrap_or_else(|| existing.map(|d| d.watch.clone()).unwrap_or_default()),
228            watch_mode: opts
229                .watch_mode
230                .unwrap_or_else(|| existing.map(|d| d.watch_mode).unwrap_or_default()),
231            watch_base_dir: opts
232                .watch_base_dir
233                .or(existing.and_then(|d| d.watch_base_dir.clone())),
234            mise: opts.mise.or(existing.and_then(|d| d.mise)),
235            user: opts.user.or(existing.and_then(|d| d.user.clone())),
236            proxy: opts.proxy.or(existing.and_then(|d| d.proxy)),
237            // active_port is intentionally NOT inherited from the existing daemon.
238            // When a daemon restarts, the new process has not yet bound a port, so
239            // carrying over the old process's active_port would cause the proxy to
240            // route to a port that is no longer listening.  The port will be
241            // re-detected by detect_and_store_active_port once the new process is ready.
242            active_port: opts.active_port,
243            slug: opts.slug.or(existing.and_then(|d| d.slug.clone())),
244            memory_limit: opts.memory_limit.or(existing.and_then(|d| d.memory_limit)),
245            cpu_limit: opts.cpu_limit.or(existing.and_then(|d| d.cpu_limit)),
246            stop_signal: opts.stop_signal.or(existing.and_then(|d| d.stop_signal)),
247            archive_hook: opts
248                .archive_hook
249                .or(existing.and_then(|d| d.archive_hook.clone())),
250            pty: opts.pty.or(existing.and_then(|d| d.pty)),
251            config_registered: opts.config_registered,
252        };
253        state_file.insert_daemon(&opts.id, daemon.clone());
254        Ok(daemon)
255    }
256
257    /// Enable a daemon (remove from disabled set)
258    pub async fn enable(&self, id: &DaemonId) -> Result<bool> {
259        info!("enabling daemon: {id}");
260        let config = PitchforkToml::all_merged_all_namespaces()?;
261        let mut state_file = self.state_file.lock().await;
262        let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
263        if !exists {
264            return Err(miette::miette!("daemon '{}' not found", id));
265        }
266        let result = state_file.enable_daemon(id);
267        Ok(result)
268    }
269
270    /// Disable a daemon (add to disabled set)
271    pub async fn disable(&self, id: &DaemonId) -> Result<bool> {
272        info!("disabling daemon: {id}");
273        let config = PitchforkToml::all_merged_all_namespaces()?;
274        let mut state_file = self.state_file.lock().await;
275        let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
276        if !exists {
277            return Err(miette::miette!("daemon '{}' not found", id));
278        }
279        let result = state_file.disable_daemon(id);
280        Ok(result)
281    }
282
283    /// Get a daemon by ID
284    pub(crate) async fn get_daemon(&self, id: &DaemonId) -> Option<Daemon> {
285        self.state_file.lock().await.daemons.get(id).cloned()
286    }
287
288    /// Get all active daemons (those with PIDs, excluding pitchfork itself)
289    pub(crate) async fn active_daemons(&self) -> Vec<Daemon> {
290        let pitchfork_id = DaemonId::pitchfork();
291        self.state_file
292            .lock()
293            .await
294            .daemons
295            .values()
296            .filter(|d| d.pid.is_some() && d.id != pitchfork_id)
297            .cloned()
298            .collect()
299    }
300
301    /// Remove a daemon from state
302    pub(crate) async fn remove_daemon(&self, id: &DaemonId) -> Result<()> {
303        let mut state_file = self.state_file.lock().await;
304        state_file.remove_daemon(id);
305        Ok(())
306    }
307
308    /// Set the shell's working directory
309    pub(crate) async fn set_shell_dir(&self, shell_pid: u32, dir: PathBuf) -> Result<()> {
310        let mut state_file = self.state_file.lock().await;
311        state_file.set_shell_dir(shell_pid, dir);
312        Ok(())
313    }
314
315    /// Get the shell's working directory
316    pub(crate) async fn get_shell_dir(&self, shell_pid: u32) -> Option<PathBuf> {
317        self.state_file
318            .lock()
319            .await
320            .shell_dirs
321            .get(&shell_pid.to_string())
322            .cloned()
323    }
324
325    /// Remove a shell PID from tracking
326    pub(crate) async fn remove_shell_pid(&self, shell_pid: u32) -> Result<()> {
327        let mut state_file = self.state_file.lock().await;
328        state_file.remove_shell_dir(shell_pid);
329        Ok(())
330    }
331
332    /// Get all directories with their associated shell PIDs
333    pub(crate) async fn get_dirs_with_shell_pids(&self) -> HashMap<PathBuf, Vec<u32>> {
334        self.state_file.lock().await.shell_dirs.iter().fold(
335            HashMap::new(),
336            |mut acc, (pid, dir)| {
337                if let Ok(pid) = pid.parse() {
338                    acc.entry(dir.clone()).or_default().push(pid);
339                }
340                acc
341            },
342        )
343    }
344
345    /// Get pending notifications and clear the queue
346    pub(crate) async fn get_notifications(&self) -> Vec<(log::LevelFilter, String)> {
347        self.pending_notifications.lock().await.drain(..).collect()
348    }
349
350    /// Clean up daemons that have no PID
351    pub(crate) async fn clean(&self) -> Result<()> {
352        let mut state_file = self.state_file.lock().await;
353        state_file.retain_daemons(|_id, d| d.pid.is_some());
354        Ok(())
355    }
356}