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::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 indexmap::IndexMap;
29use std::collections::{HashMap, HashSet};
30use std::path::PathBuf;
31
32fn should_clean_daemon(
33 id: &DaemonId,
34 daemon: &Daemon,
35 namespaces: &[String],
36 daemons: &[DaemonId],
37) -> bool {
38 let namespace_matches =
39 namespaces.is_empty() || namespaces.iter().any(|ns| ns == id.namespace());
40 let daemon_matches = daemons.is_empty() || daemons.contains(id);
41 daemon.pid.is_none() && namespace_matches && daemon_matches
42}
43
44fn prune_candidate(
45 id: &DaemonId,
46 daemon: &Daemon,
47 namespaces: &[String],
48 daemons: &[DaemonId],
49) -> Option<(DaemonId, PathBuf)> {
50 should_clean_daemon(id, daemon, namespaces, daemons)
51 .then(|| daemon.dir.clone())
52 .flatten()
53 .map(|dir| (id.clone(), dir))
54}
55
56#[derive(Debug, Default)]
60pub(crate) struct UpsertDaemonOpts {
61 pub id: DaemonId,
62 pub pid: Option<u32>,
63 pub status: DaemonStatus,
64 pub shell_pid: Option<u32>,
65 pub dir: Option<PathBuf>,
66 pub cmd: Option<Vec<String>>,
67 pub run: Option<String>,
68 pub autostop: bool,
69 pub oneshot: Option<bool>,
73 pub cron_schedule: Option<String>,
74 pub cron_retrigger: Option<CronRetrigger>,
75 pub cron_immediate: Option<bool>,
76 pub last_exit_success: Option<bool>,
77 pub retry: Option<Retry>,
78 pub retry_count: Option<u32>,
79 pub ready_delay: Option<u64>,
80 pub ready_output: Option<ReadyOutput>,
81 pub ready_http: Option<ReadyHttp>,
82 pub ready_port: Option<ReadyPort>,
83 pub ready_cmd: Option<ReadyCmd>,
84 pub health_cmd: Option<HealthCmd>,
85 pub health_http: Option<HealthHttp>,
86 pub health_port: Option<HealthPort>,
87 pub port: Option<PortConfig>,
89 pub resolved_port: Option<Vec<u16>>,
93 pub active_port: Option<u16>,
95 pub slug: Option<String>,
97 pub proxy: Option<bool>,
99 pub depends: Option<Vec<DaemonId>>,
100 pub env: Option<IndexMap<String, String>>,
101 pub watch: Option<Vec<String>>,
102 pub watch_mode: Option<WatchMode>,
103 pub watch_base_dir: Option<PathBuf>,
104 pub mise: Option<bool>,
105 pub user: Option<String>,
107 pub memory_limit: Option<MemoryLimit>,
109 pub cpu_limit: Option<CpuLimit>,
111 pub stop_signal: Option<StopConfig>,
113 pub archive_hook: Option<String>,
115 pub log_format: Option<String>,
117 pub pty: Option<bool>,
119 pub config_registered: bool,
121}
122
123#[derive(Debug)]
135pub(crate) struct UpsertDaemonOptsBuilder {
136 pub opts: UpsertDaemonOpts,
137}
138
139impl UpsertDaemonOpts {
140 pub fn builder(id: DaemonId) -> UpsertDaemonOptsBuilder {
142 UpsertDaemonOptsBuilder {
143 opts: UpsertDaemonOpts {
144 id,
145 ..Default::default()
146 },
147 }
148 }
149
150 pub(crate) fn from_run_options(
154 opts: &RunOptions,
155 status: DaemonStatus,
156 ) -> UpsertDaemonOptsBuilder {
157 UpsertDaemonOpts::builder(opts.id.clone()).set(|o| {
158 o.status = status;
159 o.shell_pid = opts.shell_pid;
160 o.dir = Some(opts.dir.0.clone());
161 o.cmd = Some(opts.cmd.clone());
162 o.run = opts.run.clone();
163 o.autostop = opts.autostop;
164 o.oneshot = Some(opts.oneshot);
165 o.cron_schedule = opts.cron_schedule.clone();
166 o.cron_retrigger = opts.cron_retrigger;
167 o.cron_immediate = opts.cron_immediate;
168 o.retry = Some(opts.retry);
169 o.retry_count = Some(opts.retry_count);
170 o.ready_delay = opts.ready_delay;
171 o.ready_output = opts.ready_output.clone();
172 o.ready_http = opts.ready_http.clone();
173 o.ready_port = opts.ready_port.clone();
174 o.ready_cmd = opts.ready_cmd.clone();
175 o.health_cmd = opts.health_cmd.clone();
176 o.health_http = opts.health_http.clone();
177 o.health_port = opts.health_port.clone();
178 o.port = opts.port.clone();
179 o.depends = Some(opts.depends.clone());
180 o.env = opts.env.clone();
181 o.watch = Some(opts.watch.clone());
182 o.watch_mode = Some(opts.watch_mode);
183 o.watch_base_dir = opts.watch_base_dir.clone();
184 o.mise = opts.mise;
185 o.user = opts.user.clone();
186 o.memory_limit = opts.memory_limit;
187 o.cpu_limit = opts.cpu_limit;
188 o.stop_signal = opts.stop_signal;
189 o.pty = opts.pty;
190 o.archive_hook = opts.archive_hook.clone();
191 o.log_format = opts.log_format.clone();
192 })
193 }
194}
195
196impl UpsertDaemonOptsBuilder {
197 pub fn set<F: FnOnce(&mut UpsertDaemonOpts)>(mut self, f: F) -> Self {
199 f(&mut self.opts);
200 self
201 }
202
203 pub fn build(self) -> UpsertDaemonOpts {
205 self.opts
206 }
207}
208
209impl Supervisor {
210 pub(crate) async fn upsert_daemon(&self, opts: UpsertDaemonOpts) -> Result<Daemon> {
212 info!(
213 "upserting daemon: {} pid: {} status: {}",
214 opts.id,
215 opts.pid.unwrap_or(0),
216 opts.status
217 );
218 let mut state_file = self.state_file.lock().await;
219 let existing = state_file.daemons.get(&opts.id);
220 let daemon = Daemon {
221 id: opts.id.clone(),
222 title: opts.pid.and_then(|pid| {
228 PROCS.title(pid).or_else(|| {
229 existing
230 .filter(|d| d.pid == Some(pid))
231 .and_then(|d| d.title.clone())
232 })
233 }),
234 start_time: opts.pid.and_then(|pid| {
235 PROCS.start_time(pid).or_else(|| {
236 existing
237 .filter(|d| d.pid == Some(pid))
238 .and_then(|d| d.start_time)
239 })
240 }),
241 boot_time: opts.pid.map(|_| PROCS.boot_time()),
245 pid: opts.pid,
246 status: opts.status,
247 shell_pid: opts.shell_pid,
248 autostop: opts.autostop || existing.is_some_and(|d| d.autostop),
249 oneshot: opts
254 .oneshot
255 .unwrap_or_else(|| existing.is_some_and(|d| d.oneshot)),
256 dir: opts.dir.or(existing.and_then(|d| d.dir.clone())),
257 cmd: opts.cmd.or(existing.and_then(|d| d.cmd.clone())),
258 run: opts.run.or(existing.and_then(|d| d.run.clone())),
259 cron_schedule: opts
260 .cron_schedule
261 .or(existing.and_then(|d| d.cron_schedule.clone())),
262 cron_retrigger: opts
263 .cron_retrigger
264 .or(existing.and_then(|d| d.cron_retrigger)),
265 cron_immediate: opts
266 .cron_immediate
267 .or(existing.and_then(|d| d.cron_immediate)),
268 last_cron_triggered: existing.and_then(|d| d.last_cron_triggered),
269 last_exit_success: opts
270 .last_exit_success
271 .or(existing.and_then(|d| d.last_exit_success)),
272 retry: opts
273 .retry
274 .unwrap_or_else(|| existing.map(|d| d.retry).unwrap_or_default()),
275 retry_count: opts
276 .retry_count
277 .unwrap_or(existing.map(|d| d.retry_count).unwrap_or(0)),
278 ready_delay: opts.ready_delay.or(existing.and_then(|d| d.ready_delay)),
279 ready_output: opts
280 .ready_output
281 .or(existing.and_then(|d| d.ready_output.clone())),
282 ready_http: opts
283 .ready_http
284 .or(existing.and_then(|d| d.ready_http.clone())),
285 ready_port: opts
286 .ready_port
287 .or(existing.and_then(|d| d.ready_port.clone())),
288 ready_cmd: opts
289 .ready_cmd
290 .or(existing.and_then(|d| d.ready_cmd.clone())),
291 health_cmd: opts
292 .health_cmd
293 .or(existing.and_then(|d| d.health_cmd.clone())),
294 health_http: opts
295 .health_http
296 .or(existing.and_then(|d| d.health_http.clone())),
297 health_port: opts
298 .health_port
299 .or(existing.and_then(|d| d.health_port.clone())),
300 port: opts.port.or_else(|| existing.and_then(|d| d.port.clone())),
301 resolved_port: match opts.resolved_port {
302 Some(ports) => ports,
303 None => existing
304 .map(|d| d.resolved_port.clone())
305 .unwrap_or_default(),
306 },
307 depends: opts
308 .depends
309 .unwrap_or_else(|| existing.map(|d| d.depends.clone()).unwrap_or_default()),
310 env: opts.env.or(existing.and_then(|d| d.env.clone())),
311 watch: opts
312 .watch
313 .unwrap_or_else(|| existing.map(|d| d.watch.clone()).unwrap_or_default()),
314 watch_mode: opts
315 .watch_mode
316 .unwrap_or_else(|| existing.map(|d| d.watch_mode).unwrap_or_default()),
317 watch_base_dir: opts
318 .watch_base_dir
319 .or(existing.and_then(|d| d.watch_base_dir.clone())),
320 mise: opts.mise.or(existing.and_then(|d| d.mise)),
321 user: opts.user.or(existing.and_then(|d| d.user.clone())),
322 proxy: opts.proxy.or(existing.and_then(|d| d.proxy)),
323 active_port: opts.active_port,
329 slug: opts.slug.or(existing.and_then(|d| d.slug.clone())),
330 memory_limit: opts.memory_limit.or(existing.and_then(|d| d.memory_limit)),
331 cpu_limit: opts.cpu_limit.or(existing.and_then(|d| d.cpu_limit)),
332 stop_signal: opts.stop_signal.or(existing.and_then(|d| d.stop_signal)),
333 archive_hook: opts
334 .archive_hook
335 .or(existing.and_then(|d| d.archive_hook.clone())),
336 log_format: opts
337 .log_format
338 .or(existing.and_then(|d| d.log_format.clone())),
339 pty: opts.pty.or(existing.and_then(|d| d.pty)),
340 config_registered: opts.config_registered,
341 };
342 state_file.insert_daemon(&opts.id, daemon.clone());
343 Ok(daemon)
344 }
345
346 pub async fn enable(&self, id: &DaemonId) -> Result<bool> {
348 info!("enabling daemon: {id}");
349 let config = PitchforkToml::all_merged_all_namespaces()?;
350 let mut state_file = self.state_file.lock().await;
351 let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
352 if !exists {
353 return Err(miette::miette!("daemon '{}' not found", id));
354 }
355 let result = state_file.enable_daemon(id);
356 Ok(result)
357 }
358
359 pub async fn disable(&self, id: &DaemonId) -> Result<bool> {
361 info!("disabling daemon: {id}");
362 let config = PitchforkToml::all_merged_all_namespaces()?;
363 let mut state_file = self.state_file.lock().await;
364 let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
365 if !exists {
366 return Err(miette::miette!("daemon '{}' not found", id));
367 }
368 let result = state_file.disable_daemon(id);
369 Ok(result)
370 }
371
372 pub(crate) async fn get_daemon(&self, id: &DaemonId) -> Option<Daemon> {
374 self.state_file.lock().await.daemons.get(id).cloned()
375 }
376
377 pub(crate) async fn active_daemons(&self) -> Vec<Daemon> {
379 let pitchfork_id = DaemonId::pitchfork();
380 self.state_file
381 .lock()
382 .await
383 .daemons
384 .values()
385 .filter(|d| d.pid.is_some() && d.id != pitchfork_id)
386 .cloned()
387 .collect()
388 }
389
390 pub(crate) async fn remove_daemon(&self, id: &DaemonId) -> Result<()> {
392 let mut state_file = self.state_file.lock().await;
393 state_file.remove_daemon(id);
394 Ok(())
395 }
396
397 pub(crate) async fn set_shell_dir(&self, shell_pid: u32, dir: PathBuf) -> Result<()> {
399 let mut state_file = self.state_file.lock().await;
400 state_file.set_shell_dir(shell_pid, dir);
401 Ok(())
402 }
403
404 pub(crate) async fn get_shell_dir(&self, shell_pid: u32) -> Option<PathBuf> {
406 self.state_file
407 .lock()
408 .await
409 .shell_dirs
410 .get(&shell_pid.to_string())
411 .cloned()
412 }
413
414 #[cfg(unix)]
416 pub(crate) async fn remove_shell_pid(&self, shell_pid: u32) -> Result<()> {
417 let mut state_file = self.state_file.lock().await;
418 state_file.remove_shell_dir(shell_pid);
419 Ok(())
420 }
421
422 pub(crate) async fn get_dirs_with_shell_pids(&self) -> HashMap<PathBuf, Vec<u32>> {
424 self.state_file.lock().await.shell_dirs.iter().fold(
425 HashMap::new(),
426 |mut acc, (pid, dir)| {
427 if let Ok(pid) = pid.parse() {
428 acc.entry(dir.clone()).or_default().push(pid);
429 }
430 acc
431 },
432 )
433 }
434
435 pub(crate) async fn get_notifications(&self) -> Vec<(log::LevelFilter, String)> {
437 self.pending_notifications.lock().await.drain(..).collect()
438 }
439
440 pub(crate) async fn clean(&self) -> Result<()> {
442 self.clean_filtered(&[], &[], false).await?;
443 Ok(())
444 }
445
446 pub(crate) async fn clean_filtered(
448 &self,
449 namespaces: &[String],
450 daemons: &[DaemonId],
451 prune: bool,
452 ) -> Result<u64> {
453 if prune {
454 let candidates: Vec<(DaemonId, PathBuf)> = {
455 let state_file = self.state_file.lock().await;
456 state_file
457 .daemons
458 .iter()
459 .filter_map(|(id, daemon)| prune_candidate(id, daemon, namespaces, daemons))
460 .collect()
461 };
462
463 let mut missing = HashMap::new();
464 for (id, dir) in candidates {
465 if !tokio::fs::try_exists(&dir)
466 .await
467 .map_err(|source| FileError::ReadError {
468 path: dir.clone(),
469 source,
470 })?
471 {
472 missing.insert(id, dir);
473 }
474 }
475
476 let mut removed = 0;
477 let mut state_file = self.state_file.lock().await;
478 state_file.retain_daemons(|id, daemon| {
479 let remove = should_clean_daemon(id, daemon, namespaces, daemons)
480 && missing
481 .get(id)
482 .is_some_and(|missing_dir| daemon.dir.as_ref() == Some(missing_dir));
483 removed += u64::from(remove);
484 !remove
485 });
486 return Ok(removed);
487 }
488
489 let mut removed = 0;
490 let mut state_file = self.state_file.lock().await;
491 state_file.retain_daemons(|id, daemon| {
492 let remove = should_clean_daemon(id, daemon, namespaces, daemons);
493 removed += u64::from(remove);
494 !remove
495 });
496 Ok(removed)
497 }
498
499 pub(crate) async fn get_active_directories(&self) -> Vec<PathBuf> {
503 let state = self.state_file.lock().await;
504 let mut dirs: HashSet<PathBuf> = state.shell_dirs.values().cloned().collect();
505 for (_, dir, _) in state.iter_project_sessions() {
506 dirs.insert(dir.clone());
507 }
508 dirs.into_iter().collect()
509 }
510
511 pub(crate) async fn get_liveness_sessions(&self) -> Vec<(u32, PathBuf, Option<String>)> {
515 self.state_file
516 .lock()
517 .await
518 .iter_project_sessions()
519 .into_iter()
520 .filter_map(|(pid_str, dir, session)| {
521 pid_str
522 .parse()
523 .ok()
524 .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
525 })
526 .collect()
527 }
528
529 pub(crate) async fn get_project_sessions_info(&self) -> Vec<crate::ipc::ProjectSessionInfo> {
533 let sessions: Vec<(u32, PathBuf, Option<String>)> = self
534 .state_file
535 .lock()
536 .await
537 .iter_project_sessions()
538 .into_iter()
539 .filter_map(|(pid_str, dir, session)| {
540 pid_str
541 .parse()
542 .ok()
543 .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
544 })
545 .collect();
546 let pids: Vec<u32> = sessions.iter().map(|(pid, _, _)| *pid).collect();
547 if !pids.is_empty() {
548 PROCS.refresh_pids(&pids);
549 }
550 sessions
551 .into_iter()
552 .map(
553 |(pid, directory, liveness_title)| crate::ipc::ProjectSessionInfo {
554 pid,
555 directory,
556 liveness_title,
557 alive: PROCS.is_running(pid),
558 current_title: PROCS.title(pid),
559 },
560 )
561 .collect()
562 }
563
564 pub(crate) async fn enter_project_session(
568 &self,
569 pid: u32,
570 dir: PathBuf,
571 ) -> Result<Option<crate::state_file::ProjectSession>> {
572 if pid == 0 {
573 return Err(miette::miette!("invalid host PID 0"));
574 }
575 if pid > i32::MAX as u32 {
576 return Err(miette::miette!("host PID {pid} exceeds i32::MAX"));
577 }
578 PROCS.refresh_pids(&[pid]);
579 let liveness_title = PROCS.title(pid);
580 #[cfg(unix)]
586 if !PROCS.is_running(pid) {
587 return Err(miette::miette!("host PID {pid} is not running"));
588 }
589 let mut state_file = self.state_file.lock().await;
590 let previous = state_file.set_project_session(
591 pid,
592 dir,
593 crate::state_file::ProjectSession { liveness_title },
594 );
595 Ok(previous)
596 }
597
598 pub(crate) async fn leave_project_session(
602 &self,
603 pid: u32,
604 dir: &std::path::Path,
605 ) -> Result<Option<PathBuf>> {
606 let mut state_file = self.state_file.lock().await;
607 if state_file.remove_project_session(pid, dir).is_some() {
608 Ok(Some(dir.to_path_buf()))
609 } else {
610 Ok(None)
611 }
612 }
613}
614
615#[cfg(test)]
616mod tests {
617 use super::*;
618
619 fn daemon(id: &DaemonId, pid: Option<u32>, dir: Option<PathBuf>) -> Daemon {
620 Daemon {
621 id: id.clone(),
622 pid,
623 dir,
624 ..Daemon::default()
625 }
626 }
627
628 #[test]
629 fn clean_filters_intersect_and_preserve_running_daemons() {
630 let api = DaemonId::new("project-a", "api");
631 let worker = DaemonId::new("project-a", "worker");
632 let namespaces = vec!["project-a".to_string()];
633 let daemons = vec![api.clone()];
634
635 assert!(should_clean_daemon(
636 &api,
637 &daemon(&api, None, None),
638 &namespaces,
639 &daemons,
640 ));
641 assert!(!should_clean_daemon(
642 &worker,
643 &daemon(&worker, None, None),
644 &namespaces,
645 &daemons,
646 ));
647 assert!(!should_clean_daemon(
648 &api,
649 &daemon(&api, Some(42), None),
650 &namespaces,
651 &daemons,
652 ));
653 }
654
655 #[test]
656 fn prune_candidates_require_a_recorded_directory() {
657 let id = DaemonId::new("project-a", "api");
658 assert!(prune_candidate(&id, &daemon(&id, None, None), &[], &[]).is_none());
659 assert_eq!(
660 prune_candidate(
661 &id,
662 &daemon(&id, None, Some(PathBuf::from("missing"))),
663 &[],
664 &[]
665 ),
666 Some((id, PathBuf::from("missing")))
667 );
668 }
669}