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;
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: opts.pid.and_then(|pid| PROCS.title(pid)),
182            pid: opts.pid,
183            status: opts.status,
184            shell_pid: opts.shell_pid,
185            autostop: opts.autostop || existing.is_some_and(|d| d.autostop),
186            dir: opts.dir.or(existing.and_then(|d| d.dir.clone())),
187            cmd: opts.cmd.or(existing.and_then(|d| d.cmd.clone())),
188            run: opts.run.or(existing.and_then(|d| d.run.clone())),
189            cron_schedule: opts
190                .cron_schedule
191                .or(existing.and_then(|d| d.cron_schedule.clone())),
192            cron_retrigger: opts
193                .cron_retrigger
194                .or(existing.and_then(|d| d.cron_retrigger)),
195            cron_immediate: opts
196                .cron_immediate
197                .or(existing.and_then(|d| d.cron_immediate)),
198            last_cron_triggered: existing.and_then(|d| d.last_cron_triggered),
199            last_exit_success: opts
200                .last_exit_success
201                .or(existing.and_then(|d| d.last_exit_success)),
202            retry: opts
203                .retry
204                .unwrap_or_else(|| existing.map(|d| d.retry).unwrap_or_default()),
205            retry_count: opts
206                .retry_count
207                .unwrap_or(existing.map(|d| d.retry_count).unwrap_or(0)),
208            ready_delay: opts.ready_delay.or(existing.and_then(|d| d.ready_delay)),
209            ready_output: opts
210                .ready_output
211                .or(existing.and_then(|d| d.ready_output.clone())),
212            ready_http: opts
213                .ready_http
214                .or(existing.and_then(|d| d.ready_http.clone())),
215            ready_port: opts
216                .ready_port
217                .or(existing.and_then(|d| d.ready_port.clone())),
218            ready_cmd: opts
219                .ready_cmd
220                .or(existing.and_then(|d| d.ready_cmd.clone())),
221            port: opts.port.or_else(|| existing.and_then(|d| d.port.clone())),
222            resolved_port: if opts.resolved_port.is_empty() {
223                existing
224                    .map(|d| d.resolved_port.clone())
225                    .unwrap_or_default()
226            } else {
227                opts.resolved_port
228            },
229            depends: opts
230                .depends
231                .unwrap_or_else(|| existing.map(|d| d.depends.clone()).unwrap_or_default()),
232            env: opts.env.or(existing.and_then(|d| d.env.clone())),
233            watch: opts
234                .watch
235                .unwrap_or_else(|| existing.map(|d| d.watch.clone()).unwrap_or_default()),
236            watch_mode: opts
237                .watch_mode
238                .unwrap_or_else(|| existing.map(|d| d.watch_mode).unwrap_or_default()),
239            watch_base_dir: opts
240                .watch_base_dir
241                .or(existing.and_then(|d| d.watch_base_dir.clone())),
242            mise: opts.mise.or(existing.and_then(|d| d.mise)),
243            user: opts.user.or(existing.and_then(|d| d.user.clone())),
244            proxy: opts.proxy.or(existing.and_then(|d| d.proxy)),
245            // active_port is intentionally NOT inherited from the existing daemon.
246            // When a daemon restarts, the new process has not yet bound a port, so
247            // carrying over the old process's active_port would cause the proxy to
248            // route to a port that is no longer listening.  The port will be
249            // re-detected by detect_and_store_active_port once the new process is ready.
250            active_port: opts.active_port,
251            slug: opts.slug.or(existing.and_then(|d| d.slug.clone())),
252            memory_limit: opts.memory_limit.or(existing.and_then(|d| d.memory_limit)),
253            cpu_limit: opts.cpu_limit.or(existing.and_then(|d| d.cpu_limit)),
254            stop_signal: opts.stop_signal.or(existing.and_then(|d| d.stop_signal)),
255            archive_hook: opts
256                .archive_hook
257                .or(existing.and_then(|d| d.archive_hook.clone())),
258            log_format: opts
259                .log_format
260                .or(existing.and_then(|d| d.log_format.clone())),
261            pty: opts.pty.or(existing.and_then(|d| d.pty)),
262            config_registered: opts.config_registered,
263        };
264        state_file.insert_daemon(&opts.id, daemon.clone());
265        Ok(daemon)
266    }
267
268    /// Enable a daemon (remove from disabled set)
269    pub async fn enable(&self, id: &DaemonId) -> Result<bool> {
270        info!("enabling daemon: {id}");
271        let config = PitchforkToml::all_merged_all_namespaces()?;
272        let mut state_file = self.state_file.lock().await;
273        let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
274        if !exists {
275            return Err(miette::miette!("daemon '{}' not found", id));
276        }
277        let result = state_file.enable_daemon(id);
278        Ok(result)
279    }
280
281    /// Disable a daemon (add to disabled set)
282    pub async fn disable(&self, id: &DaemonId) -> Result<bool> {
283        info!("disabling daemon: {id}");
284        let config = PitchforkToml::all_merged_all_namespaces()?;
285        let mut state_file = self.state_file.lock().await;
286        let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
287        if !exists {
288            return Err(miette::miette!("daemon '{}' not found", id));
289        }
290        let result = state_file.disable_daemon(id);
291        Ok(result)
292    }
293
294    /// Get a daemon by ID
295    pub(crate) async fn get_daemon(&self, id: &DaemonId) -> Option<Daemon> {
296        self.state_file.lock().await.daemons.get(id).cloned()
297    }
298
299    /// Get all active daemons (those with PIDs, excluding pitchfork itself)
300    pub(crate) async fn active_daemons(&self) -> Vec<Daemon> {
301        let pitchfork_id = DaemonId::pitchfork();
302        self.state_file
303            .lock()
304            .await
305            .daemons
306            .values()
307            .filter(|d| d.pid.is_some() && d.id != pitchfork_id)
308            .cloned()
309            .collect()
310    }
311
312    /// Remove a daemon from state
313    pub(crate) async fn remove_daemon(&self, id: &DaemonId) -> Result<()> {
314        let mut state_file = self.state_file.lock().await;
315        state_file.remove_daemon(id);
316        Ok(())
317    }
318
319    /// Set the shell's working directory
320    pub(crate) async fn set_shell_dir(&self, shell_pid: u32, dir: PathBuf) -> Result<()> {
321        let mut state_file = self.state_file.lock().await;
322        state_file.set_shell_dir(shell_pid, dir);
323        Ok(())
324    }
325
326    /// Get the shell's working directory
327    pub(crate) async fn get_shell_dir(&self, shell_pid: u32) -> Option<PathBuf> {
328        self.state_file
329            .lock()
330            .await
331            .shell_dirs
332            .get(&shell_pid.to_string())
333            .cloned()
334    }
335
336    /// Remove a shell PID from tracking
337    pub(crate) async fn remove_shell_pid(&self, shell_pid: u32) -> Result<()> {
338        let mut state_file = self.state_file.lock().await;
339        state_file.remove_shell_dir(shell_pid);
340        Ok(())
341    }
342
343    /// Get all directories with their associated shell PIDs
344    pub(crate) async fn get_dirs_with_shell_pids(&self) -> HashMap<PathBuf, Vec<u32>> {
345        self.state_file.lock().await.shell_dirs.iter().fold(
346            HashMap::new(),
347            |mut acc, (pid, dir)| {
348                if let Ok(pid) = pid.parse() {
349                    acc.entry(dir.clone()).or_default().push(pid);
350                }
351                acc
352            },
353        )
354    }
355
356    /// Get pending notifications and clear the queue
357    pub(crate) async fn get_notifications(&self) -> Vec<(log::LevelFilter, String)> {
358        self.pending_notifications.lock().await.drain(..).collect()
359    }
360
361    /// Clean up daemons that have no PID
362    pub(crate) async fn clean(&self) -> Result<()> {
363        let mut state_file = self.state_file.lock().await;
364        state_file.retain_daemons(|_id, d| d.pid.is_some());
365        Ok(())
366    }
367}