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::HealthCmd;
15use crate::pitchfork_toml::HealthHttp;
16use crate::pitchfork_toml::HealthPort;
17use crate::pitchfork_toml::MemoryLimit;
18use crate::pitchfork_toml::PitchforkToml;
19use crate::pitchfork_toml::PortConfig;
20use crate::pitchfork_toml::ReadyCmd;
21use crate::pitchfork_toml::ReadyHttp;
22use crate::pitchfork_toml::ReadyOutput;
23use crate::pitchfork_toml::ReadyPort;
24use crate::pitchfork_toml::Retry;
25use crate::pitchfork_toml::StopConfig;
26use crate::pitchfork_toml::WatchMode;
27use crate::procs::PROCS;
28use crate::state_file::DiskRecord;
29use indexmap::IndexMap;
30use std::collections::HashMap;
31use std::path::PathBuf;
32
33fn should_clean_daemon(
34    id: &DaemonId,
35    daemon: &Daemon,
36    namespaces: &[String],
37    daemons: &[DaemonId],
38) -> bool {
39    let namespace_matches =
40        namespaces.is_empty() || namespaces.iter().any(|ns| ns == id.namespace());
41    let daemon_matches = daemons.is_empty() || daemons.contains(id);
42    daemon.pid.is_none() && namespace_matches && daemon_matches
43}
44
45fn prune_candidate(
46    id: &DaemonId,
47    daemon: &Daemon,
48    namespaces: &[String],
49    daemons: &[DaemonId],
50) -> Option<(DaemonId, PathBuf)> {
51    should_clean_daemon(id, daemon, namespaces, daemons)
52        .then(|| daemon.dir.clone())
53        .flatten()
54        .map(|dir| (id.clone(), dir))
55}
56
57/// Options for upserting a daemon's state.
58///
59/// Use `UpsertDaemonOpts::builder(id)` to create, then set fields directly and call `.build()`.
60#[derive(Debug, Default)]
61pub(crate) struct UpsertDaemonOpts {
62    pub id: DaemonId,
63    pub pid: Option<u32>,
64    pub status: DaemonStatus,
65    pub shell_pid: Option<u32>,
66    pub dir: Option<PathBuf>,
67    pub cmd: Option<Vec<String>>,
68    pub run: Option<String>,
69    /// Start `cmd` directly, without a shell. `None` inherits the existing
70    /// record's value, like `oneshot`.
71    pub no_shell: Option<bool>,
72    /// `None` inherits the existing record's value; see
73    /// `Daemon::scheduled_from_config`.
74    pub scheduled_from_config: Option<bool>,
75    pub autostop: bool,
76    /// Run-to-completion task rather than a long-running service.
77    /// `None` inherits the existing record's value, so a status-only upsert
78    /// (stop, exit finalization) does not reclassify the daemon.
79    pub oneshot: Option<bool>,
80    pub cron_schedule: Option<String>,
81    pub cron_retrigger: Option<CronRetrigger>,
82    pub cron_immediate: Option<bool>,
83    pub last_exit_success: Option<bool>,
84    pub retry: Option<Retry>,
85    pub retry_count: Option<u32>,
86    pub ready_delay: Option<u64>,
87    pub ready_output: Option<ReadyOutput>,
88    pub ready_http: Option<ReadyHttp>,
89    pub ready_port: Option<ReadyPort>,
90    pub ready_cmd: Option<ReadyCmd>,
91    pub health_cmd: Option<HealthCmd>,
92    pub health_http: Option<HealthHttp>,
93    pub health_port: Option<HealthPort>,
94    /// Port configuration
95    pub port: Option<PortConfig>,
96    /// Resolved ports actually used after auto-bump (may differ from expected).
97    /// `None` inherits the existing record's value; `Some` sets it explicitly,
98    /// including an empty vec, which clears ports left by a previous run.
99    pub resolved_port: Option<Vec<u16>>,
100    /// The first port the process is actually listening on (detected at runtime).
101    pub active_port: Option<u16>,
102    /// Optional stable slug alias for this daemon.
103    pub slug: Option<String>,
104    /// Whether to proxy this daemon (None = use global proxy.enable setting).
105    pub proxy: Option<bool>,
106    pub depends: Option<Vec<DaemonId>>,
107    pub env: Option<IndexMap<String, String>>,
108    pub watch: Option<Vec<String>>,
109    pub watch_mode: Option<WatchMode>,
110    pub watch_base_dir: Option<PathBuf>,
111    pub mise: Option<bool>,
112    /// Unix user to run this daemon as
113    pub user: Option<String>,
114    /// Memory limit for the daemon process
115    pub memory_limit: Option<MemoryLimit>,
116    /// CPU usage limit as a percentage
117    pub cpu_limit: Option<CpuLimit>,
118    /// Unix signal to send for graceful shutdown
119    pub stop_signal: Option<StopConfig>,
120    /// Archive hook command invoked before retention prunes this daemon's logs.
121    pub archive_hook: Option<String>,
122    /// Log format for this daemon.
123    pub log_format: Option<String>,
124    /// Allocate a pseudo-terminal for the daemon process.
125    pub pty: Option<bool>,
126    /// True for config-only cron daemons auto-registered into state.
127    pub config_registered: bool,
128    /// Idle-shutdown ownership. `None` inherits the existing record's, so a
129    /// status-only upsert (stop, exit finalization) keeps it; a start sets it
130    /// from its `RunOptions`, which is what makes an explicit start explicit.
131    pub proxy_idle_timeout_ms: Option<Option<u64>>,
132}
133
134/// Builder for UpsertDaemonOpts - ensures daemon ID is always provided.
135///
136/// # Example
137/// ```ignore
138/// let opts = UpsertDaemonOpts::builder(daemon_id)
139///     .set(|o| {
140///         o.pid = Some(pid);
141///         o.status = DaemonStatus::Running;
142///     })
143///     .build();
144/// ```
145#[derive(Debug)]
146pub(crate) struct UpsertDaemonOptsBuilder {
147    pub opts: UpsertDaemonOpts,
148}
149
150impl UpsertDaemonOpts {
151    /// Create a builder with the required daemon ID.
152    pub fn builder(id: DaemonId) -> UpsertDaemonOptsBuilder {
153        UpsertDaemonOptsBuilder {
154            opts: UpsertDaemonOpts {
155                id,
156                ..Default::default()
157            },
158        }
159    }
160
161    /// Build an `UpsertDaemonOptsBuilder` from `RunOptions`, mapping all
162    /// config-carried fields. Callers chain `.set()` to add runtime-specific
163    /// fields (pid, resolved_port, etc.) and then call `.build()`.
164    pub(crate) fn from_run_options(
165        opts: &RunOptions,
166        status: DaemonStatus,
167    ) -> UpsertDaemonOptsBuilder {
168        UpsertDaemonOpts::builder(opts.id.clone()).set(|o| {
169            o.status = status;
170            o.shell_pid = opts.shell_pid;
171            o.dir = Some(opts.dir.0.clone());
172            o.cmd = Some(opts.cmd.clone());
173            o.run = opts.run.clone();
174            o.no_shell = Some(opts.no_shell);
175            // Only a client clears it; the supervisor's own starts (the
176            // schedule, retries, file watching) leave it as it was.
177            o.scheduled_from_config = opts.requested_by_client.then_some(false);
178            o.autostop = opts.autostop;
179            o.oneshot = Some(opts.oneshot);
180            o.cron_schedule = opts.cron_schedule.clone();
181            o.cron_retrigger = opts.cron_retrigger;
182            o.cron_immediate = opts.cron_immediate;
183            o.retry = Some(opts.retry);
184            o.retry_count = Some(opts.retry_count);
185            o.ready_delay = opts.ready_delay;
186            o.ready_output = opts.ready_output.clone();
187            o.ready_http = opts.ready_http.clone();
188            o.ready_port = opts.ready_port.clone();
189            o.ready_cmd = opts.ready_cmd.clone();
190            o.health_cmd = opts.health_cmd.clone();
191            o.health_http = opts.health_http.clone();
192            o.health_port = opts.health_port.clone();
193            o.port = opts.port.clone();
194            o.depends = Some(opts.depends.clone());
195            o.env = opts.env.clone();
196            o.watch = Some(opts.watch.clone());
197            o.watch_mode = Some(opts.watch_mode);
198            o.watch_base_dir = opts.watch_base_dir.clone();
199            o.mise = opts.mise;
200            o.user = opts.user.clone();
201            o.memory_limit = opts.memory_limit;
202            o.cpu_limit = opts.cpu_limit;
203            o.stop_signal = opts.stop_signal;
204            o.pty = opts.pty;
205            o.archive_hook = opts.archive_hook.clone();
206            o.log_format = opts.log_format.clone();
207            o.proxy_idle_timeout_ms = Some(opts.proxy_idle_timeout_ms);
208        })
209    }
210}
211
212impl UpsertDaemonOptsBuilder {
213    /// Modify opts fields with a closure.
214    pub fn set<F: FnOnce(&mut UpsertDaemonOpts)>(mut self, f: F) -> Self {
215        f(&mut self.opts);
216        self
217    }
218
219    /// Build the UpsertDaemonOpts.
220    pub fn build(self) -> UpsertDaemonOpts {
221        self.opts
222    }
223}
224
225impl Supervisor {
226    /// Put this supervisor's record back into the state file if the file no
227    /// longer has it.
228    ///
229    /// The CLI decides whether a supervisor is running, and which process
230    /// `supervisor stop` signals, from this record. The supervisor only
231    /// writes the file when its own state changes, so if the file is replaced
232    /// or rewritten by something else the record could stay missing
233    /// indefinitely, leaving this supervisor running but unaccounted for.
234    /// Only the record is restored; the rest of the file is left as it is.
235    pub(crate) async fn restore_own_record(&self) {
236        if self
237            .shutting_down
238            .load(std::sync::atomic::Ordering::Acquire)
239        {
240            return;
241        }
242        let pitchfork_id = DaemonId::pitchfork();
243        let state = self.state_file.lock().await;
244        // Checked again under the lock: `close()` may have begun while this
245        // waited for it. (Its removal of the record happens under this lock
246        // too, so a restore can never outlive it; this just skips the work.)
247        if self
248            .shutting_down
249            .load(std::sync::atomic::Ordering::Acquire)
250        {
251            return;
252        }
253        let Some(own) = state
254            .daemons
255            .get(&pitchfork_id)
256            .filter(|d| d.pid.is_some())
257            .cloned()
258        else {
259            return;
260        };
261        let own_pid = own.pid.unwrap_or_default();
262        let path = state.path.clone();
263        // The state lock stays held across the file work, as it is for a
264        // flush, so the two cannot interleave.
265        let result = tokio::task::spawn_blocking(move || {
266            crate::state_file::StateFile::restore_daemon_in_file(&path, &own)
267        })
268        .await;
269        match result {
270            Ok(Ok(DiskRecord::Present)) => {}
271            Ok(Ok(DiskRecord::Restored)) => {
272                // Not flushed yet (e.g. a client connecting right after
273                // startup): the record was simply written early.
274                if state.is_dirty() {
275                    debug!("recorded this supervisor in the state file ahead of the next flush");
276                } else {
277                    warn!(
278                        "state file {} no longer recorded this supervisor (pid {own_pid}); restored it",
279                        state.path.display()
280                    );
281                }
282                // The file now differs from what this supervisor last wrote,
283                // so its next flush must not be skipped as unchanged.
284                state.forget_written_snapshot();
285            }
286            Ok(Ok(DiskRecord::Unparseable)) => {
287                warn!(
288                    "state file {} cannot be parsed; rewriting it from the supervisor's state",
289                    state.path.display()
290                );
291                state.force_next_write();
292            }
293            Ok(Err(e)) => warn!("failed to restore the supervisor record in the state file: {e}"),
294            Err(e) => warn!("failed to restore the supervisor record in the state file: {e}"),
295        }
296    }
297
298    /// Run options for starting `id` from its config, as `pitchfork start`
299    /// would build them: templates rendered and the top-level `[env]` merged.
300    ///
301    /// Other daemons' ports come from the state held here rather than the file
302    /// on disk, which can trail it — a daemon started a moment ago may not be
303    /// written out yet, and a template referring to its port would fail.
304    pub(crate) async fn run_options_from_config(
305        &self,
306        id: &DaemonId,
307        config: &crate::pitchfork_toml::PitchforkTomlDaemon,
308        pt: &PitchforkToml,
309    ) -> Result<RunOptions> {
310        let resolved_daemons: std::collections::HashMap<DaemonId, Vec<u16>> = {
311            let state = self.state_file.lock().await;
312            state
313                .daemons
314                .iter()
315                .filter(|(_, d)| !d.resolved_port.is_empty())
316                .map(|(id, d)| (id.clone(), d.resolved_port.clone()))
317                .collect()
318        };
319        // On a blocking worker, as for a client's start: building the context
320        // reads configuration and derives hostnames by walking the project's
321        // checkouts, which must not hold up the supervisor's other tasks.
322        let id = id.clone();
323        let mut config = config.clone();
324        let pt = pt.clone();
325        tokio::task::spawn_blocking(move || {
326            crate::ipc::batch::render_daemon_config_with(&id, &mut config, &pt, &resolved_daemons)?;
327            let cmd = config
328                .run
329                .argv()
330                .map_err(|e| miette::miette!("failed to parse command for daemon {id}: {e}"))?;
331            Ok(config.to_run_options(&id, cmd))
332        })
333        .await
334        .map_err(|e| miette::miette!("rendering daemon config panicked: {e}"))?
335    }
336
337    /// Upsert a daemon's state, merging with existing values
338    pub(crate) async fn upsert_daemon(&self, opts: UpsertDaemonOpts) -> Result<Daemon> {
339        info!(
340            "upserting daemon: {} pid: {} status: {}",
341            opts.id,
342            opts.pid.unwrap_or(0),
343            opts.status
344        );
345        let mut state_file = self.state_file.lock().await;
346        let existing = state_file.daemons.get(&opts.id);
347        let daemon = Daemon {
348            id: opts.id.clone(),
349            // title/start_time identify the process for orphan cleanup after a
350            // supervisor crash. They are looked up from the process cache; if
351            // the cache has no entry (e.g. an upsert between refreshes) fall
352            // back to the existing values as long as the PID is unchanged, so
353            // the recorded identity is not wiped mid-lifetime.
354            title: opts.pid.and_then(|pid| {
355                PROCS.title(pid).or_else(|| {
356                    existing
357                        .filter(|d| d.pid == Some(pid))
358                        .and_then(|d| d.title.clone())
359                })
360            }),
361            start_time: opts.pid.and_then(|pid| {
362                PROCS.start_time(pid).or_else(|| {
363                    existing
364                        .filter(|d| d.pid == Some(pid))
365                        .and_then(|d| d.start_time)
366                })
367            }),
368            // Boot time is a system-wide constant, so it needs no cache
369            // fallback — it is recorded for any daemon with a live PID and
370            // dropped when the PID is cleared.
371            boot_time: opts.pid.map(|_| PROCS.boot_time()),
372            pid: opts.pid,
373            status: opts.status,
374            shell_pid: opts.shell_pid,
375            autostop: opts.autostop || existing.is_some_and(|d| d.autostop),
376            // A start carries the current config value (including a removed
377            // `oneshot = true`, which must reclassify the daemon); every other
378            // upsert leaves it None and inherits, so finalizing a completed
379            // oneshot's exit does not forget what it is.
380            oneshot: opts
381                .oneshot
382                .unwrap_or_else(|| existing.is_some_and(|d| d.oneshot)),
383            dir: opts.dir.or(existing.and_then(|d| d.dir.clone())),
384            cmd: opts.cmd.or(existing.and_then(|d| d.cmd.clone())),
385            // A start in the argv form has no command line, and must not
386            // inherit the one a shell-form run left behind.
387            run: if opts.no_shell == Some(true) {
388                None
389            } else {
390                opts.run.or(existing.and_then(|d| d.run.clone()))
391            },
392            no_shell: opts
393                .no_shell
394                .unwrap_or_else(|| existing.is_some_and(|d| d.no_shell)),
395            scheduled_from_config: opts
396                .scheduled_from_config
397                .unwrap_or_else(|| existing.is_some_and(|d| d.scheduled_from_config)),
398            cron_schedule: opts
399                .cron_schedule
400                .or(existing.and_then(|d| d.cron_schedule.clone())),
401            cron_retrigger: opts
402                .cron_retrigger
403                .or(existing.and_then(|d| d.cron_retrigger)),
404            cron_immediate: opts
405                .cron_immediate
406                .or(existing.and_then(|d| d.cron_immediate)),
407            last_cron_triggered: existing.and_then(|d| d.last_cron_triggered),
408            last_cron_run: existing.and_then(|d| d.last_cron_run),
409            last_exit_success: opts
410                .last_exit_success
411                .or(existing.and_then(|d| d.last_exit_success)),
412            retry: opts
413                .retry
414                .unwrap_or_else(|| existing.map(|d| d.retry).unwrap_or_default()),
415            retry_count: opts
416                .retry_count
417                .unwrap_or(existing.map(|d| d.retry_count).unwrap_or(0)),
418            ready_delay: opts.ready_delay.or(existing.and_then(|d| d.ready_delay)),
419            ready_output: opts
420                .ready_output
421                .or(existing.and_then(|d| d.ready_output.clone())),
422            ready_http: opts
423                .ready_http
424                .or(existing.and_then(|d| d.ready_http.clone())),
425            ready_port: opts
426                .ready_port
427                .or(existing.and_then(|d| d.ready_port.clone())),
428            ready_cmd: opts
429                .ready_cmd
430                .or(existing.and_then(|d| d.ready_cmd.clone())),
431            health_cmd: opts
432                .health_cmd
433                .or(existing.and_then(|d| d.health_cmd.clone())),
434            health_http: opts
435                .health_http
436                .or(existing.and_then(|d| d.health_http.clone())),
437            health_port: opts
438                .health_port
439                .or(existing.and_then(|d| d.health_port.clone())),
440            port: opts.port.or_else(|| existing.and_then(|d| d.port.clone())),
441            resolved_port: match opts.resolved_port {
442                Some(ports) => ports,
443                None => existing
444                    .map(|d| d.resolved_port.clone())
445                    .unwrap_or_default(),
446            },
447            depends: opts
448                .depends
449                .unwrap_or_else(|| existing.map(|d| d.depends.clone()).unwrap_or_default()),
450            env: opts.env.or(existing.and_then(|d| d.env.clone())),
451            watch: opts
452                .watch
453                .unwrap_or_else(|| existing.map(|d| d.watch.clone()).unwrap_or_default()),
454            watch_mode: opts
455                .watch_mode
456                .unwrap_or_else(|| existing.map(|d| d.watch_mode).unwrap_or_default()),
457            watch_base_dir: opts
458                .watch_base_dir
459                .or(existing.and_then(|d| d.watch_base_dir.clone())),
460            mise: opts.mise.or(existing.and_then(|d| d.mise)),
461            user: opts.user.or(existing.and_then(|d| d.user.clone())),
462            proxy: opts.proxy.or(existing.and_then(|d| d.proxy)),
463            // active_port is intentionally NOT inherited from the existing daemon.
464            // When a daemon restarts, the new process has not yet bound a port, so
465            // carrying over the old process's active_port would cause the proxy to
466            // route to a port that is no longer listening.  The port will be
467            // re-detected by detect_and_store_active_port once the new process is ready.
468            active_port: opts.active_port,
469            slug: opts.slug.or(existing.and_then(|d| d.slug.clone())),
470            memory_limit: opts.memory_limit.or(existing.and_then(|d| d.memory_limit)),
471            cpu_limit: opts.cpu_limit.or(existing.and_then(|d| d.cpu_limit)),
472            stop_signal: opts.stop_signal.or(existing.and_then(|d| d.stop_signal)),
473            archive_hook: opts
474                .archive_hook
475                .or(existing.and_then(|d| d.archive_hook.clone())),
476            log_format: opts
477                .log_format
478                .or(existing.and_then(|d| d.log_format.clone())),
479            pty: opts.pty.or(existing.and_then(|d| d.pty)),
480            config_registered: opts.config_registered,
481            proxy_idle_timeout_ms: opts
482                .proxy_idle_timeout_ms
483                .unwrap_or_else(|| existing.and_then(|d| d.proxy_idle_timeout_ms)),
484        };
485        state_file.insert_daemon(&opts.id, daemon.clone());
486        Ok(daemon)
487    }
488
489    /// Enable a daemon (remove from disabled set)
490    pub async fn enable(&self, id: &DaemonId) -> Result<bool> {
491        info!("enabling daemon: {id}");
492        let config = PitchforkToml::all_merged_all_namespaces()?;
493        let mut state_file = self.state_file.lock().await;
494        let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
495        if !exists {
496            return Err(miette::miette!("daemon '{}' not found", id));
497        }
498        let result = state_file.enable_daemon(id);
499        Ok(result)
500    }
501
502    /// Disable a daemon (add to disabled set)
503    pub async fn disable(&self, id: &DaemonId) -> Result<bool> {
504        info!("disabling daemon: {id}");
505        let config = PitchforkToml::all_merged_all_namespaces()?;
506        let mut state_file = self.state_file.lock().await;
507        let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
508        if !exists {
509            return Err(miette::miette!("daemon '{}' not found", id));
510        }
511        let result = state_file.disable_daemon(id);
512        Ok(result)
513    }
514
515    /// Get a daemon by ID
516    pub(crate) async fn get_daemon(&self, id: &DaemonId) -> Option<Daemon> {
517        self.state_file.lock().await.daemons.get(id).cloned()
518    }
519
520    /// Get all active daemons (those with PIDs, excluding pitchfork itself)
521    pub(crate) async fn active_daemons(&self) -> Vec<Daemon> {
522        let pitchfork_id = DaemonId::pitchfork();
523        self.state_file
524            .lock()
525            .await
526            .daemons
527            .values()
528            .filter(|d| d.pid.is_some() && d.id != pitchfork_id)
529            .cloned()
530            .collect()
531    }
532
533    /// Remove a daemon from state
534    pub(crate) async fn remove_daemon(&self, id: &DaemonId) -> Result<()> {
535        let mut state_file = self.state_file.lock().await;
536        state_file.remove_daemon(id);
537        Ok(())
538    }
539
540    /// Set the shell's working directory
541    pub(crate) async fn set_shell_dir(&self, shell_pid: u32, dir: PathBuf) -> Result<()> {
542        let mut state_file = self.state_file.lock().await;
543        state_file.set_shell_dir(shell_pid, dir);
544        Ok(())
545    }
546
547    /// Get the shell's working directory
548    pub(crate) async fn get_shell_dir(&self, shell_pid: u32) -> Option<PathBuf> {
549        self.state_file
550            .lock()
551            .await
552            .shell_dirs
553            .get(&shell_pid.to_string())
554            .cloned()
555    }
556
557    /// Remove a shell PID from tracking
558    #[cfg(unix)]
559    pub(crate) async fn remove_shell_pid(&self, shell_pid: u32) -> Result<()> {
560        let mut state_file = self.state_file.lock().await;
561        state_file.remove_shell_dir(shell_pid);
562        Ok(())
563    }
564
565    /// Get all directories with their associated shell PIDs
566    pub(crate) async fn get_dirs_with_shell_pids(&self) -> HashMap<PathBuf, Vec<u32>> {
567        self.state_file.lock().await.shell_dirs.iter().fold(
568            HashMap::new(),
569            |mut acc, (pid, dir)| {
570                if let Ok(pid) = pid.parse() {
571                    acc.entry(dir.clone()).or_default().push(pid);
572                }
573                acc
574            },
575        )
576    }
577
578    /// Get pending notifications and clear the queue
579    pub(crate) async fn get_notifications(&self) -> Vec<(log::LevelFilter, String)> {
580        self.pending_notifications.lock().await.drain(..).collect()
581    }
582
583    /// Clean up daemons that have no PID
584    pub(crate) async fn clean(&self) -> Result<()> {
585        self.clean_filtered(&[], &[], false).await?;
586        Ok(())
587    }
588
589    /// Clean up PID-less daemon registrations matching all supplied filters.
590    pub(crate) async fn clean_filtered(
591        &self,
592        namespaces: &[String],
593        daemons: &[DaemonId],
594        prune: bool,
595    ) -> Result<u64> {
596        if prune {
597            let candidates: Vec<(DaemonId, PathBuf)> = {
598                let state_file = self.state_file.lock().await;
599                state_file
600                    .daemons
601                    .iter()
602                    .filter_map(|(id, daemon)| prune_candidate(id, daemon, namespaces, daemons))
603                    .collect()
604            };
605
606            let mut missing = HashMap::new();
607            for (id, dir) in candidates {
608                if !tokio::fs::try_exists(&dir)
609                    .await
610                    .map_err(|source| FileError::ReadError {
611                        path: dir.clone(),
612                        source,
613                    })?
614                {
615                    missing.insert(id, dir);
616                }
617            }
618
619            let mut removed = 0;
620            let mut state_file = self.state_file.lock().await;
621            state_file.retain_daemons(|id, daemon| {
622                let remove = should_clean_daemon(id, daemon, namespaces, daemons)
623                    && missing
624                        .get(id)
625                        .is_some_and(|missing_dir| daemon.dir.as_ref() == Some(missing_dir));
626                removed += u64::from(remove);
627                !remove
628            });
629            return Ok(removed);
630        }
631
632        let mut removed = 0;
633        let mut state_file = self.state_file.lock().await;
634        state_file.retain_daemons(|id, daemon| {
635            let remove = should_clean_daemon(id, daemon, namespaces, daemons);
636            removed += u64::from(remove);
637            !remove
638        });
639        Ok(removed)
640    }
641
642    /// Return the union of active directories from shell tracking and project
643    /// sessions. These are the directories that should keep auto-stop daemons
644    /// alive.
645    pub(crate) async fn get_active_directories(&self) -> Vec<PathBuf> {
646        self.state_file.lock().await.active_directories()
647    }
648
649    /// Collect all project sessions as `(pid, dir, liveness_title)`. Every
650    /// project session carries a host PID in its key, so every session is a
651    /// liveness session. Used by the refresh loop to clean up stale sessions.
652    pub(crate) async fn get_liveness_sessions(&self) -> Vec<(u32, PathBuf, Option<String>)> {
653        self.state_file
654            .lock()
655            .await
656            .iter_project_sessions()
657            .into_iter()
658            .filter_map(|(pid_str, dir, session)| {
659                pid_str
660                    .parse()
661                    .ok()
662                    .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
663            })
664            .collect()
665    }
666
667    /// Build a snapshot of all project sessions with live liveness status
668    /// filled in from `PROCS`. Used to answer `GetProjectSessions` IPC
669    /// requests.
670    pub(crate) async fn get_project_sessions_info(&self) -> Vec<crate::ipc::ProjectSessionInfo> {
671        let sessions: Vec<(u32, PathBuf, Option<String>)> = self
672            .state_file
673            .lock()
674            .await
675            .iter_project_sessions()
676            .into_iter()
677            .filter_map(|(pid_str, dir, session)| {
678                pid_str
679                    .parse()
680                    .ok()
681                    .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
682            })
683            .collect();
684        let pids: Vec<u32> = sessions.iter().map(|(pid, _, _)| *pid).collect();
685        if !pids.is_empty() {
686            PROCS.refresh_pids(&pids);
687        }
688        sessions
689            .into_iter()
690            .map(
691                |(pid, directory, liveness_title)| crate::ipc::ProjectSessionInfo {
692                    pid,
693                    directory,
694                    liveness_title,
695                    alive: PROCS.is_running(pid),
696                    current_title: PROCS.title(pid),
697                },
698            )
699            .collect()
700    }
701
702    /// Atomically enter (or replace) a project session for `(pid, dir)`.
703    /// Returns the previous session, if any, so the caller can evaluate the
704    /// previous entry for autostop.
705    pub(crate) async fn enter_project_session(
706        &self,
707        pid: u32,
708        dir: PathBuf,
709    ) -> Result<Option<crate::state_file::ProjectSession>> {
710        if pid == 0 {
711            return Err(miette::miette!("invalid host PID 0"));
712        }
713        if pid > i32::MAX as u32 {
714            return Err(miette::miette!("host PID {pid} exceeds i32::MAX"));
715        }
716        PROCS.refresh_pids(&[pid]);
717        let liveness_title = PROCS.title(pid);
718        // On Unix, reject dead host PIDs up front so we don't register a
719        // session the refresh loop would immediately evict. On Windows, Git
720        // Bash `$$` is a Cygwin-internal PID invisible to sysinfo, so the
721        // liveness check is skipped and sessions rely on explicit
722        // `project leave`.
723        #[cfg(unix)]
724        if !PROCS.is_running(pid) {
725            return Err(miette::miette!("host PID {pid} is not running"));
726        }
727        let mut state_file = self.state_file.lock().await;
728        let previous = state_file.set_project_session(
729            pid,
730            dir,
731            crate::state_file::ProjectSession { liveness_title },
732        );
733        Ok(previous)
734    }
735
736    /// Atomically remove a project session for `(pid, dir)`. Returns the
737    /// directory of the removed session so the caller can evaluate it for
738    /// autostop.
739    pub(crate) async fn leave_project_session(
740        &self,
741        pid: u32,
742        dir: &std::path::Path,
743    ) -> Result<Option<PathBuf>> {
744        let mut state_file = self.state_file.lock().await;
745        if state_file.remove_project_session(pid, dir).is_some() {
746            Ok(Some(dir.to_path_buf()))
747        } else {
748            Ok(None)
749        }
750    }
751}
752
753#[cfg(test)]
754mod tests {
755    use super::*;
756
757    fn daemon(id: &DaemonId, pid: Option<u32>, dir: Option<PathBuf>) -> Daemon {
758        Daemon {
759            id: id.clone(),
760            pid,
761            dir,
762            ..Daemon::default()
763        }
764    }
765
766    #[test]
767    fn clean_filters_intersect_and_preserve_running_daemons() {
768        let api = DaemonId::new("project-a", "api");
769        let worker = DaemonId::new("project-a", "worker");
770        let namespaces = vec!["project-a".to_string()];
771        let daemons = vec![api.clone()];
772
773        assert!(should_clean_daemon(
774            &api,
775            &daemon(&api, None, None),
776            &namespaces,
777            &daemons,
778        ));
779        assert!(!should_clean_daemon(
780            &worker,
781            &daemon(&worker, None, None),
782            &namespaces,
783            &daemons,
784        ));
785        assert!(!should_clean_daemon(
786            &api,
787            &daemon(&api, Some(42), None),
788            &namespaces,
789            &daemons,
790        ));
791    }
792
793    #[test]
794    fn prune_candidates_require_a_recorded_directory() {
795        let id = DaemonId::new("project-a", "api");
796        assert!(prune_candidate(&id, &daemon(&id, None, None), &[], &[]).is_none());
797        assert_eq!(
798            prune_candidate(
799                &id,
800                &daemon(&id, None, Some(PathBuf::from("missing"))),
801                &[],
802                &[]
803            ),
804            Some((id, PathBuf::from("missing")))
805        );
806    }
807}