1use 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#[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 pub port: Option<PortConfig>,
79 pub resolved_port: Option<Vec<u16>>,
83 pub active_port: Option<u16>,
85 pub slug: Option<String>,
87 pub proxy: Option<bool>,
89 pub depends: Option<Vec<DaemonId>>,
90 pub env: Option<IndexMap<String, String>>,
91 pub watch: Option<Vec<String>>,
92 pub watch_mode: Option<WatchMode>,
93 pub watch_base_dir: Option<PathBuf>,
94 pub mise: Option<bool>,
95 pub user: Option<String>,
97 pub memory_limit: Option<MemoryLimit>,
99 pub cpu_limit: Option<CpuLimit>,
101 pub stop_signal: Option<StopConfig>,
103 pub archive_hook: Option<String>,
105 pub log_format: Option<String>,
107 pub pty: Option<bool>,
109 pub config_registered: bool,
111}
112
113#[derive(Debug)]
125pub(crate) struct UpsertDaemonOptsBuilder {
126 pub opts: UpsertDaemonOpts,
127}
128
129impl UpsertDaemonOpts {
130 pub fn builder(id: DaemonId) -> UpsertDaemonOptsBuilder {
132 UpsertDaemonOptsBuilder {
133 opts: UpsertDaemonOpts {
134 id,
135 ..Default::default()
136 },
137 }
138 }
139
140 pub(crate) fn from_run_options(
144 opts: &RunOptions,
145 status: DaemonStatus,
146 ) -> UpsertDaemonOptsBuilder {
147 UpsertDaemonOpts::builder(opts.id.clone()).set(|o| {
148 o.status = status;
149 o.shell_pid = opts.shell_pid;
150 o.dir = Some(opts.dir.0.clone());
151 o.cmd = Some(opts.cmd.clone());
152 o.run = opts.run.clone();
153 o.autostop = opts.autostop;
154 o.cron_schedule = opts.cron_schedule.clone();
155 o.cron_retrigger = opts.cron_retrigger;
156 o.cron_immediate = opts.cron_immediate;
157 o.retry = Some(opts.retry);
158 o.retry_count = Some(opts.retry_count);
159 o.ready_delay = opts.ready_delay;
160 o.ready_output = opts.ready_output.clone();
161 o.ready_http = opts.ready_http.clone();
162 o.ready_port = opts.ready_port.clone();
163 o.ready_cmd = opts.ready_cmd.clone();
164 o.port = opts.port.clone();
165 o.depends = Some(opts.depends.clone());
166 o.env = opts.env.clone();
167 o.watch = Some(opts.watch.clone());
168 o.watch_mode = Some(opts.watch_mode);
169 o.watch_base_dir = opts.watch_base_dir.clone();
170 o.mise = opts.mise;
171 o.user = opts.user.clone();
172 o.memory_limit = opts.memory_limit;
173 o.cpu_limit = opts.cpu_limit;
174 o.stop_signal = opts.stop_signal;
175 o.pty = opts.pty;
176 o.archive_hook = opts.archive_hook.clone();
177 o.log_format = opts.log_format.clone();
178 })
179 }
180}
181
182impl UpsertDaemonOptsBuilder {
183 pub fn set<F: FnOnce(&mut UpsertDaemonOpts)>(mut self, f: F) -> Self {
185 f(&mut self.opts);
186 self
187 }
188
189 pub fn build(self) -> UpsertDaemonOpts {
191 self.opts
192 }
193}
194
195impl Supervisor {
196 pub(crate) async fn upsert_daemon(&self, opts: UpsertDaemonOpts) -> Result<Daemon> {
198 info!(
199 "upserting daemon: {} pid: {} status: {}",
200 opts.id,
201 opts.pid.unwrap_or(0),
202 opts.status
203 );
204 let mut state_file = self.state_file.lock().await;
205 let existing = state_file.daemons.get(&opts.id);
206 let daemon = Daemon {
207 id: opts.id.clone(),
208 title: opts.pid.and_then(|pid| {
214 PROCS.title(pid).or_else(|| {
215 existing
216 .filter(|d| d.pid == Some(pid))
217 .and_then(|d| d.title.clone())
218 })
219 }),
220 start_time: opts.pid.and_then(|pid| {
221 PROCS.start_time(pid).or_else(|| {
222 existing
223 .filter(|d| d.pid == Some(pid))
224 .and_then(|d| d.start_time)
225 })
226 }),
227 boot_time: opts.pid.map(|_| PROCS.boot_time()),
231 pid: opts.pid,
232 status: opts.status,
233 shell_pid: opts.shell_pid,
234 autostop: opts.autostop || existing.is_some_and(|d| d.autostop),
235 dir: opts.dir.or(existing.and_then(|d| d.dir.clone())),
236 cmd: opts.cmd.or(existing.and_then(|d| d.cmd.clone())),
237 run: opts.run.or(existing.and_then(|d| d.run.clone())),
238 cron_schedule: opts
239 .cron_schedule
240 .or(existing.and_then(|d| d.cron_schedule.clone())),
241 cron_retrigger: opts
242 .cron_retrigger
243 .or(existing.and_then(|d| d.cron_retrigger)),
244 cron_immediate: opts
245 .cron_immediate
246 .or(existing.and_then(|d| d.cron_immediate)),
247 last_cron_triggered: existing.and_then(|d| d.last_cron_triggered),
248 last_exit_success: opts
249 .last_exit_success
250 .or(existing.and_then(|d| d.last_exit_success)),
251 retry: opts
252 .retry
253 .unwrap_or_else(|| existing.map(|d| d.retry).unwrap_or_default()),
254 retry_count: opts
255 .retry_count
256 .unwrap_or(existing.map(|d| d.retry_count).unwrap_or(0)),
257 ready_delay: opts.ready_delay.or(existing.and_then(|d| d.ready_delay)),
258 ready_output: opts
259 .ready_output
260 .or(existing.and_then(|d| d.ready_output.clone())),
261 ready_http: opts
262 .ready_http
263 .or(existing.and_then(|d| d.ready_http.clone())),
264 ready_port: opts
265 .ready_port
266 .or(existing.and_then(|d| d.ready_port.clone())),
267 ready_cmd: opts
268 .ready_cmd
269 .or(existing.and_then(|d| d.ready_cmd.clone())),
270 port: opts.port.or_else(|| existing.and_then(|d| d.port.clone())),
271 resolved_port: match opts.resolved_port {
272 Some(ports) => ports,
273 None => existing
274 .map(|d| d.resolved_port.clone())
275 .unwrap_or_default(),
276 },
277 depends: opts
278 .depends
279 .unwrap_or_else(|| existing.map(|d| d.depends.clone()).unwrap_or_default()),
280 env: opts.env.or(existing.and_then(|d| d.env.clone())),
281 watch: opts
282 .watch
283 .unwrap_or_else(|| existing.map(|d| d.watch.clone()).unwrap_or_default()),
284 watch_mode: opts
285 .watch_mode
286 .unwrap_or_else(|| existing.map(|d| d.watch_mode).unwrap_or_default()),
287 watch_base_dir: opts
288 .watch_base_dir
289 .or(existing.and_then(|d| d.watch_base_dir.clone())),
290 mise: opts.mise.or(existing.and_then(|d| d.mise)),
291 user: opts.user.or(existing.and_then(|d| d.user.clone())),
292 proxy: opts.proxy.or(existing.and_then(|d| d.proxy)),
293 active_port: opts.active_port,
299 slug: opts.slug.or(existing.and_then(|d| d.slug.clone())),
300 memory_limit: opts.memory_limit.or(existing.and_then(|d| d.memory_limit)),
301 cpu_limit: opts.cpu_limit.or(existing.and_then(|d| d.cpu_limit)),
302 stop_signal: opts.stop_signal.or(existing.and_then(|d| d.stop_signal)),
303 archive_hook: opts
304 .archive_hook
305 .or(existing.and_then(|d| d.archive_hook.clone())),
306 log_format: opts
307 .log_format
308 .or(existing.and_then(|d| d.log_format.clone())),
309 pty: opts.pty.or(existing.and_then(|d| d.pty)),
310 config_registered: opts.config_registered,
311 };
312 state_file.insert_daemon(&opts.id, daemon.clone());
313 Ok(daemon)
314 }
315
316 pub async fn enable(&self, id: &DaemonId) -> Result<bool> {
318 info!("enabling daemon: {id}");
319 let config = PitchforkToml::all_merged_all_namespaces()?;
320 let mut state_file = self.state_file.lock().await;
321 let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
322 if !exists {
323 return Err(miette::miette!("daemon '{}' not found", id));
324 }
325 let result = state_file.enable_daemon(id);
326 Ok(result)
327 }
328
329 pub async fn disable(&self, id: &DaemonId) -> Result<bool> {
331 info!("disabling daemon: {id}");
332 let config = PitchforkToml::all_merged_all_namespaces()?;
333 let mut state_file = self.state_file.lock().await;
334 let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
335 if !exists {
336 return Err(miette::miette!("daemon '{}' not found", id));
337 }
338 let result = state_file.disable_daemon(id);
339 Ok(result)
340 }
341
342 pub(crate) async fn get_daemon(&self, id: &DaemonId) -> Option<Daemon> {
344 self.state_file.lock().await.daemons.get(id).cloned()
345 }
346
347 pub(crate) async fn active_daemons(&self) -> Vec<Daemon> {
349 let pitchfork_id = DaemonId::pitchfork();
350 self.state_file
351 .lock()
352 .await
353 .daemons
354 .values()
355 .filter(|d| d.pid.is_some() && d.id != pitchfork_id)
356 .cloned()
357 .collect()
358 }
359
360 pub(crate) async fn remove_daemon(&self, id: &DaemonId) -> Result<()> {
362 let mut state_file = self.state_file.lock().await;
363 state_file.remove_daemon(id);
364 Ok(())
365 }
366
367 pub(crate) async fn set_shell_dir(&self, shell_pid: u32, dir: PathBuf) -> Result<()> {
369 let mut state_file = self.state_file.lock().await;
370 state_file.set_shell_dir(shell_pid, dir);
371 Ok(())
372 }
373
374 pub(crate) async fn get_shell_dir(&self, shell_pid: u32) -> Option<PathBuf> {
376 self.state_file
377 .lock()
378 .await
379 .shell_dirs
380 .get(&shell_pid.to_string())
381 .cloned()
382 }
383
384 pub(crate) async fn remove_shell_pid(&self, shell_pid: u32) -> Result<()> {
386 let mut state_file = self.state_file.lock().await;
387 state_file.remove_shell_dir(shell_pid);
388 Ok(())
389 }
390
391 pub(crate) async fn get_dirs_with_shell_pids(&self) -> HashMap<PathBuf, Vec<u32>> {
393 self.state_file.lock().await.shell_dirs.iter().fold(
394 HashMap::new(),
395 |mut acc, (pid, dir)| {
396 if let Ok(pid) = pid.parse() {
397 acc.entry(dir.clone()).or_default().push(pid);
398 }
399 acc
400 },
401 )
402 }
403
404 pub(crate) async fn get_notifications(&self) -> Vec<(log::LevelFilter, String)> {
406 self.pending_notifications.lock().await.drain(..).collect()
407 }
408
409 pub(crate) async fn clean(&self) -> Result<()> {
411 self.clean_filtered(&[], &[], false).await?;
412 Ok(())
413 }
414
415 pub(crate) async fn clean_filtered(
417 &self,
418 namespaces: &[String],
419 daemons: &[DaemonId],
420 prune: bool,
421 ) -> Result<u64> {
422 if prune {
423 let candidates: Vec<(DaemonId, PathBuf)> = {
424 let state_file = self.state_file.lock().await;
425 state_file
426 .daemons
427 .iter()
428 .filter_map(|(id, daemon)| prune_candidate(id, daemon, namespaces, daemons))
429 .collect()
430 };
431
432 let mut missing = HashMap::new();
433 for (id, dir) in candidates {
434 if !tokio::fs::try_exists(&dir)
435 .await
436 .map_err(|source| FileError::ReadError {
437 path: dir.clone(),
438 source,
439 })?
440 {
441 missing.insert(id, dir);
442 }
443 }
444
445 let mut removed = 0;
446 let mut state_file = self.state_file.lock().await;
447 state_file.retain_daemons(|id, daemon| {
448 let remove = should_clean_daemon(id, daemon, namespaces, daemons)
449 && missing
450 .get(id)
451 .is_some_and(|missing_dir| daemon.dir.as_ref() == Some(missing_dir));
452 removed += u64::from(remove);
453 !remove
454 });
455 return Ok(removed);
456 }
457
458 let mut removed = 0;
459 let mut state_file = self.state_file.lock().await;
460 state_file.retain_daemons(|id, daemon| {
461 let remove = should_clean_daemon(id, daemon, namespaces, daemons);
462 removed += u64::from(remove);
463 !remove
464 });
465 Ok(removed)
466 }
467
468 pub(crate) async fn get_active_directories(&self) -> Vec<PathBuf> {
472 let state = self.state_file.lock().await;
473 let mut dirs: HashSet<PathBuf> = state.shell_dirs.values().cloned().collect();
474 for (_, dir, _) in state.iter_project_sessions() {
475 dirs.insert(dir.clone());
476 }
477 dirs.into_iter().collect()
478 }
479
480 pub(crate) async fn get_liveness_sessions(&self) -> Vec<(u32, PathBuf, Option<String>)> {
484 self.state_file
485 .lock()
486 .await
487 .iter_project_sessions()
488 .into_iter()
489 .filter_map(|(pid_str, dir, session)| {
490 pid_str
491 .parse()
492 .ok()
493 .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
494 })
495 .collect()
496 }
497
498 pub(crate) async fn get_project_sessions_info(&self) -> Vec<crate::ipc::ProjectSessionInfo> {
502 let sessions: Vec<(u32, PathBuf, Option<String>)> = self
503 .state_file
504 .lock()
505 .await
506 .iter_project_sessions()
507 .into_iter()
508 .filter_map(|(pid_str, dir, session)| {
509 pid_str
510 .parse()
511 .ok()
512 .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
513 })
514 .collect();
515 let pids: Vec<u32> = sessions.iter().map(|(pid, _, _)| *pid).collect();
516 if !pids.is_empty() {
517 PROCS.refresh_pids(&pids);
518 }
519 sessions
520 .into_iter()
521 .map(
522 |(pid, directory, liveness_title)| crate::ipc::ProjectSessionInfo {
523 pid,
524 directory,
525 liveness_title,
526 alive: PROCS.is_running(pid),
527 current_title: PROCS.title(pid),
528 },
529 )
530 .collect()
531 }
532
533 pub(crate) async fn enter_project_session(
537 &self,
538 pid: u32,
539 dir: PathBuf,
540 ) -> Result<Option<crate::state_file::ProjectSession>> {
541 if pid == 0 {
542 return Err(miette::miette!("invalid host PID 0"));
543 }
544 if pid > i32::MAX as u32 {
545 return Err(miette::miette!("host PID {pid} exceeds i32::MAX"));
546 }
547 PROCS.refresh_pids(&[pid]);
548 let liveness_title = PROCS.title(pid);
549 #[cfg(unix)]
555 if !PROCS.is_running(pid) {
556 return Err(miette::miette!("host PID {pid} is not running"));
557 }
558 let mut state_file = self.state_file.lock().await;
559 let previous = state_file.set_project_session(
560 pid,
561 dir,
562 crate::state_file::ProjectSession { liveness_title },
563 );
564 Ok(previous)
565 }
566
567 pub(crate) async fn leave_project_session(
571 &self,
572 pid: u32,
573 dir: &std::path::Path,
574 ) -> Result<Option<PathBuf>> {
575 let mut state_file = self.state_file.lock().await;
576 if state_file.remove_project_session(pid, dir).is_some() {
577 Ok(Some(dir.to_path_buf()))
578 } else {
579 Ok(None)
580 }
581 }
582}
583
584#[cfg(test)]
585mod tests {
586 use super::*;
587
588 fn daemon(id: &DaemonId, pid: Option<u32>, dir: Option<PathBuf>) -> Daemon {
589 Daemon {
590 id: id.clone(),
591 pid,
592 dir,
593 ..Daemon::default()
594 }
595 }
596
597 #[test]
598 fn clean_filters_intersect_and_preserve_running_daemons() {
599 let api = DaemonId::new("project-a", "api");
600 let worker = DaemonId::new("project-a", "worker");
601 let namespaces = vec!["project-a".to_string()];
602 let daemons = vec![api.clone()];
603
604 assert!(should_clean_daemon(
605 &api,
606 &daemon(&api, None, None),
607 &namespaces,
608 &daemons,
609 ));
610 assert!(!should_clean_daemon(
611 &worker,
612 &daemon(&worker, None, None),
613 &namespaces,
614 &daemons,
615 ));
616 assert!(!should_clean_daemon(
617 &api,
618 &daemon(&api, Some(42), None),
619 &namespaces,
620 &daemons,
621 ));
622 }
623
624 #[test]
625 fn prune_candidates_require_a_recorded_directory() {
626 let id = DaemonId::new("project-a", "api");
627 assert!(prune_candidate(&id, &daemon(&id, None, None), &[], &[]).is_none());
628 assert_eq!(
629 prune_candidate(
630 &id,
631 &daemon(&id, None, Some(PathBuf::from("missing"))),
632 &[],
633 &[]
634 ),
635 Some((id, PathBuf::from("missing")))
636 );
637 }
638}