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 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#[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 pub autostop: bool,
70 pub oneshot: Option<bool>,
74 pub cron_schedule: Option<String>,
75 pub cron_retrigger: Option<CronRetrigger>,
76 pub cron_immediate: Option<bool>,
77 pub last_exit_success: Option<bool>,
78 pub retry: Option<Retry>,
79 pub retry_count: Option<u32>,
80 pub ready_delay: Option<u64>,
81 pub ready_output: Option<ReadyOutput>,
82 pub ready_http: Option<ReadyHttp>,
83 pub ready_port: Option<ReadyPort>,
84 pub ready_cmd: Option<ReadyCmd>,
85 pub health_cmd: Option<HealthCmd>,
86 pub health_http: Option<HealthHttp>,
87 pub health_port: Option<HealthPort>,
88 pub port: Option<PortConfig>,
90 pub resolved_port: Option<Vec<u16>>,
94 pub active_port: Option<u16>,
96 pub slug: Option<String>,
98 pub proxy: Option<bool>,
100 pub depends: Option<Vec<DaemonId>>,
101 pub env: Option<IndexMap<String, String>>,
102 pub watch: Option<Vec<String>>,
103 pub watch_mode: Option<WatchMode>,
104 pub watch_base_dir: Option<PathBuf>,
105 pub mise: Option<bool>,
106 pub user: Option<String>,
108 pub memory_limit: Option<MemoryLimit>,
110 pub cpu_limit: Option<CpuLimit>,
112 pub stop_signal: Option<StopConfig>,
114 pub archive_hook: Option<String>,
116 pub log_format: Option<String>,
118 pub pty: Option<bool>,
120 pub config_registered: bool,
122 pub proxy_idle_timeout_ms: Option<Option<u64>>,
126}
127
128#[derive(Debug)]
140pub(crate) struct UpsertDaemonOptsBuilder {
141 pub opts: UpsertDaemonOpts,
142}
143
144impl UpsertDaemonOpts {
145 pub fn builder(id: DaemonId) -> UpsertDaemonOptsBuilder {
147 UpsertDaemonOptsBuilder {
148 opts: UpsertDaemonOpts {
149 id,
150 ..Default::default()
151 },
152 }
153 }
154
155 pub(crate) fn from_run_options(
159 opts: &RunOptions,
160 status: DaemonStatus,
161 ) -> UpsertDaemonOptsBuilder {
162 UpsertDaemonOpts::builder(opts.id.clone()).set(|o| {
163 o.status = status;
164 o.shell_pid = opts.shell_pid;
165 o.dir = Some(opts.dir.0.clone());
166 o.cmd = Some(opts.cmd.clone());
167 o.run = opts.run.clone();
168 o.autostop = opts.autostop;
169 o.oneshot = Some(opts.oneshot);
170 o.cron_schedule = opts.cron_schedule.clone();
171 o.cron_retrigger = opts.cron_retrigger;
172 o.cron_immediate = opts.cron_immediate;
173 o.retry = Some(opts.retry);
174 o.retry_count = Some(opts.retry_count);
175 o.ready_delay = opts.ready_delay;
176 o.ready_output = opts.ready_output.clone();
177 o.ready_http = opts.ready_http.clone();
178 o.ready_port = opts.ready_port.clone();
179 o.ready_cmd = opts.ready_cmd.clone();
180 o.health_cmd = opts.health_cmd.clone();
181 o.health_http = opts.health_http.clone();
182 o.health_port = opts.health_port.clone();
183 o.port = opts.port.clone();
184 o.depends = Some(opts.depends.clone());
185 o.env = opts.env.clone();
186 o.watch = Some(opts.watch.clone());
187 o.watch_mode = Some(opts.watch_mode);
188 o.watch_base_dir = opts.watch_base_dir.clone();
189 o.mise = opts.mise;
190 o.user = opts.user.clone();
191 o.memory_limit = opts.memory_limit;
192 o.cpu_limit = opts.cpu_limit;
193 o.stop_signal = opts.stop_signal;
194 o.pty = opts.pty;
195 o.archive_hook = opts.archive_hook.clone();
196 o.log_format = opts.log_format.clone();
197 o.proxy_idle_timeout_ms = Some(opts.proxy_idle_timeout_ms);
198 })
199 }
200}
201
202impl UpsertDaemonOptsBuilder {
203 pub fn set<F: FnOnce(&mut UpsertDaemonOpts)>(mut self, f: F) -> Self {
205 f(&mut self.opts);
206 self
207 }
208
209 pub fn build(self) -> UpsertDaemonOpts {
211 self.opts
212 }
213}
214
215impl Supervisor {
216 pub(crate) async fn restore_own_record(&self) {
226 if self
227 .shutting_down
228 .load(std::sync::atomic::Ordering::Acquire)
229 {
230 return;
231 }
232 let pitchfork_id = DaemonId::pitchfork();
233 let state = self.state_file.lock().await;
234 if self
238 .shutting_down
239 .load(std::sync::atomic::Ordering::Acquire)
240 {
241 return;
242 }
243 let Some(own) = state
244 .daemons
245 .get(&pitchfork_id)
246 .filter(|d| d.pid.is_some())
247 .cloned()
248 else {
249 return;
250 };
251 let own_pid = own.pid.unwrap_or_default();
252 let path = state.path.clone();
253 let result = tokio::task::spawn_blocking(move || {
256 crate::state_file::StateFile::restore_daemon_in_file(&path, &own)
257 })
258 .await;
259 match result {
260 Ok(Ok(DiskRecord::Present)) => {}
261 Ok(Ok(DiskRecord::Restored)) => {
262 if state.is_dirty() {
265 debug!("recorded this supervisor in the state file ahead of the next flush");
266 } else {
267 warn!(
268 "state file {} no longer recorded this supervisor (pid {own_pid}); restored it",
269 state.path.display()
270 );
271 }
272 state.forget_written_snapshot();
275 }
276 Ok(Ok(DiskRecord::Unparseable)) => {
277 warn!(
278 "state file {} cannot be parsed; rewriting it from the supervisor's state",
279 state.path.display()
280 );
281 state.force_next_write();
282 }
283 Ok(Err(e)) => warn!("failed to restore the supervisor record in the state file: {e}"),
284 Err(e) => warn!("failed to restore the supervisor record in the state file: {e}"),
285 }
286 }
287
288 pub(crate) async fn upsert_daemon(&self, opts: UpsertDaemonOpts) -> Result<Daemon> {
290 info!(
291 "upserting daemon: {} pid: {} status: {}",
292 opts.id,
293 opts.pid.unwrap_or(0),
294 opts.status
295 );
296 let mut state_file = self.state_file.lock().await;
297 let existing = state_file.daemons.get(&opts.id);
298 let daemon = Daemon {
299 id: opts.id.clone(),
300 title: opts.pid.and_then(|pid| {
306 PROCS.title(pid).or_else(|| {
307 existing
308 .filter(|d| d.pid == Some(pid))
309 .and_then(|d| d.title.clone())
310 })
311 }),
312 start_time: opts.pid.and_then(|pid| {
313 PROCS.start_time(pid).or_else(|| {
314 existing
315 .filter(|d| d.pid == Some(pid))
316 .and_then(|d| d.start_time)
317 })
318 }),
319 boot_time: opts.pid.map(|_| PROCS.boot_time()),
323 pid: opts.pid,
324 status: opts.status,
325 shell_pid: opts.shell_pid,
326 autostop: opts.autostop || existing.is_some_and(|d| d.autostop),
327 oneshot: opts
332 .oneshot
333 .unwrap_or_else(|| existing.is_some_and(|d| d.oneshot)),
334 dir: opts.dir.or(existing.and_then(|d| d.dir.clone())),
335 cmd: opts.cmd.or(existing.and_then(|d| d.cmd.clone())),
336 run: opts.run.or(existing.and_then(|d| d.run.clone())),
337 cron_schedule: opts
338 .cron_schedule
339 .or(existing.and_then(|d| d.cron_schedule.clone())),
340 cron_retrigger: opts
341 .cron_retrigger
342 .or(existing.and_then(|d| d.cron_retrigger)),
343 cron_immediate: opts
344 .cron_immediate
345 .or(existing.and_then(|d| d.cron_immediate)),
346 last_cron_triggered: existing.and_then(|d| d.last_cron_triggered),
347 last_exit_success: opts
348 .last_exit_success
349 .or(existing.and_then(|d| d.last_exit_success)),
350 retry: opts
351 .retry
352 .unwrap_or_else(|| existing.map(|d| d.retry).unwrap_or_default()),
353 retry_count: opts
354 .retry_count
355 .unwrap_or(existing.map(|d| d.retry_count).unwrap_or(0)),
356 ready_delay: opts.ready_delay.or(existing.and_then(|d| d.ready_delay)),
357 ready_output: opts
358 .ready_output
359 .or(existing.and_then(|d| d.ready_output.clone())),
360 ready_http: opts
361 .ready_http
362 .or(existing.and_then(|d| d.ready_http.clone())),
363 ready_port: opts
364 .ready_port
365 .or(existing.and_then(|d| d.ready_port.clone())),
366 ready_cmd: opts
367 .ready_cmd
368 .or(existing.and_then(|d| d.ready_cmd.clone())),
369 health_cmd: opts
370 .health_cmd
371 .or(existing.and_then(|d| d.health_cmd.clone())),
372 health_http: opts
373 .health_http
374 .or(existing.and_then(|d| d.health_http.clone())),
375 health_port: opts
376 .health_port
377 .or(existing.and_then(|d| d.health_port.clone())),
378 port: opts.port.or_else(|| existing.and_then(|d| d.port.clone())),
379 resolved_port: match opts.resolved_port {
380 Some(ports) => ports,
381 None => existing
382 .map(|d| d.resolved_port.clone())
383 .unwrap_or_default(),
384 },
385 depends: opts
386 .depends
387 .unwrap_or_else(|| existing.map(|d| d.depends.clone()).unwrap_or_default()),
388 env: opts.env.or(existing.and_then(|d| d.env.clone())),
389 watch: opts
390 .watch
391 .unwrap_or_else(|| existing.map(|d| d.watch.clone()).unwrap_or_default()),
392 watch_mode: opts
393 .watch_mode
394 .unwrap_or_else(|| existing.map(|d| d.watch_mode).unwrap_or_default()),
395 watch_base_dir: opts
396 .watch_base_dir
397 .or(existing.and_then(|d| d.watch_base_dir.clone())),
398 mise: opts.mise.or(existing.and_then(|d| d.mise)),
399 user: opts.user.or(existing.and_then(|d| d.user.clone())),
400 proxy: opts.proxy.or(existing.and_then(|d| d.proxy)),
401 active_port: opts.active_port,
407 slug: opts.slug.or(existing.and_then(|d| d.slug.clone())),
408 memory_limit: opts.memory_limit.or(existing.and_then(|d| d.memory_limit)),
409 cpu_limit: opts.cpu_limit.or(existing.and_then(|d| d.cpu_limit)),
410 stop_signal: opts.stop_signal.or(existing.and_then(|d| d.stop_signal)),
411 archive_hook: opts
412 .archive_hook
413 .or(existing.and_then(|d| d.archive_hook.clone())),
414 log_format: opts
415 .log_format
416 .or(existing.and_then(|d| d.log_format.clone())),
417 pty: opts.pty.or(existing.and_then(|d| d.pty)),
418 config_registered: opts.config_registered,
419 proxy_idle_timeout_ms: opts
420 .proxy_idle_timeout_ms
421 .unwrap_or_else(|| existing.and_then(|d| d.proxy_idle_timeout_ms)),
422 };
423 state_file.insert_daemon(&opts.id, daemon.clone());
424 Ok(daemon)
425 }
426
427 pub async fn enable(&self, id: &DaemonId) -> Result<bool> {
429 info!("enabling daemon: {id}");
430 let config = PitchforkToml::all_merged_all_namespaces()?;
431 let mut state_file = self.state_file.lock().await;
432 let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
433 if !exists {
434 return Err(miette::miette!("daemon '{}' not found", id));
435 }
436 let result = state_file.enable_daemon(id);
437 Ok(result)
438 }
439
440 pub async fn disable(&self, id: &DaemonId) -> Result<bool> {
442 info!("disabling daemon: {id}");
443 let config = PitchforkToml::all_merged_all_namespaces()?;
444 let mut state_file = self.state_file.lock().await;
445 let exists = state_file.daemons.contains_key(id) || config.daemons.contains_key(id);
446 if !exists {
447 return Err(miette::miette!("daemon '{}' not found", id));
448 }
449 let result = state_file.disable_daemon(id);
450 Ok(result)
451 }
452
453 pub(crate) async fn get_daemon(&self, id: &DaemonId) -> Option<Daemon> {
455 self.state_file.lock().await.daemons.get(id).cloned()
456 }
457
458 pub(crate) async fn active_daemons(&self) -> Vec<Daemon> {
460 let pitchfork_id = DaemonId::pitchfork();
461 self.state_file
462 .lock()
463 .await
464 .daemons
465 .values()
466 .filter(|d| d.pid.is_some() && d.id != pitchfork_id)
467 .cloned()
468 .collect()
469 }
470
471 pub(crate) async fn remove_daemon(&self, id: &DaemonId) -> Result<()> {
473 let mut state_file = self.state_file.lock().await;
474 state_file.remove_daemon(id);
475 Ok(())
476 }
477
478 pub(crate) async fn set_shell_dir(&self, shell_pid: u32, dir: PathBuf) -> Result<()> {
480 let mut state_file = self.state_file.lock().await;
481 state_file.set_shell_dir(shell_pid, dir);
482 Ok(())
483 }
484
485 pub(crate) async fn get_shell_dir(&self, shell_pid: u32) -> Option<PathBuf> {
487 self.state_file
488 .lock()
489 .await
490 .shell_dirs
491 .get(&shell_pid.to_string())
492 .cloned()
493 }
494
495 #[cfg(unix)]
497 pub(crate) async fn remove_shell_pid(&self, shell_pid: u32) -> Result<()> {
498 let mut state_file = self.state_file.lock().await;
499 state_file.remove_shell_dir(shell_pid);
500 Ok(())
501 }
502
503 pub(crate) async fn get_dirs_with_shell_pids(&self) -> HashMap<PathBuf, Vec<u32>> {
505 self.state_file.lock().await.shell_dirs.iter().fold(
506 HashMap::new(),
507 |mut acc, (pid, dir)| {
508 if let Ok(pid) = pid.parse() {
509 acc.entry(dir.clone()).or_default().push(pid);
510 }
511 acc
512 },
513 )
514 }
515
516 pub(crate) async fn get_notifications(&self) -> Vec<(log::LevelFilter, String)> {
518 self.pending_notifications.lock().await.drain(..).collect()
519 }
520
521 pub(crate) async fn clean(&self) -> Result<()> {
523 self.clean_filtered(&[], &[], false).await?;
524 Ok(())
525 }
526
527 pub(crate) async fn clean_filtered(
529 &self,
530 namespaces: &[String],
531 daemons: &[DaemonId],
532 prune: bool,
533 ) -> Result<u64> {
534 if prune {
535 let candidates: Vec<(DaemonId, PathBuf)> = {
536 let state_file = self.state_file.lock().await;
537 state_file
538 .daemons
539 .iter()
540 .filter_map(|(id, daemon)| prune_candidate(id, daemon, namespaces, daemons))
541 .collect()
542 };
543
544 let mut missing = HashMap::new();
545 for (id, dir) in candidates {
546 if !tokio::fs::try_exists(&dir)
547 .await
548 .map_err(|source| FileError::ReadError {
549 path: dir.clone(),
550 source,
551 })?
552 {
553 missing.insert(id, dir);
554 }
555 }
556
557 let mut removed = 0;
558 let mut state_file = self.state_file.lock().await;
559 state_file.retain_daemons(|id, daemon| {
560 let remove = should_clean_daemon(id, daemon, namespaces, daemons)
561 && missing
562 .get(id)
563 .is_some_and(|missing_dir| daemon.dir.as_ref() == Some(missing_dir));
564 removed += u64::from(remove);
565 !remove
566 });
567 return Ok(removed);
568 }
569
570 let mut removed = 0;
571 let mut state_file = self.state_file.lock().await;
572 state_file.retain_daemons(|id, daemon| {
573 let remove = should_clean_daemon(id, daemon, namespaces, daemons);
574 removed += u64::from(remove);
575 !remove
576 });
577 Ok(removed)
578 }
579
580 pub(crate) async fn get_active_directories(&self) -> Vec<PathBuf> {
584 self.state_file.lock().await.active_directories()
585 }
586
587 pub(crate) async fn get_liveness_sessions(&self) -> Vec<(u32, PathBuf, Option<String>)> {
591 self.state_file
592 .lock()
593 .await
594 .iter_project_sessions()
595 .into_iter()
596 .filter_map(|(pid_str, dir, session)| {
597 pid_str
598 .parse()
599 .ok()
600 .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
601 })
602 .collect()
603 }
604
605 pub(crate) async fn get_project_sessions_info(&self) -> Vec<crate::ipc::ProjectSessionInfo> {
609 let sessions: Vec<(u32, PathBuf, Option<String>)> = self
610 .state_file
611 .lock()
612 .await
613 .iter_project_sessions()
614 .into_iter()
615 .filter_map(|(pid_str, dir, session)| {
616 pid_str
617 .parse()
618 .ok()
619 .map(|pid| (pid, dir.clone(), session.liveness_title.clone()))
620 })
621 .collect();
622 let pids: Vec<u32> = sessions.iter().map(|(pid, _, _)| *pid).collect();
623 if !pids.is_empty() {
624 PROCS.refresh_pids(&pids);
625 }
626 sessions
627 .into_iter()
628 .map(
629 |(pid, directory, liveness_title)| crate::ipc::ProjectSessionInfo {
630 pid,
631 directory,
632 liveness_title,
633 alive: PROCS.is_running(pid),
634 current_title: PROCS.title(pid),
635 },
636 )
637 .collect()
638 }
639
640 pub(crate) async fn enter_project_session(
644 &self,
645 pid: u32,
646 dir: PathBuf,
647 ) -> Result<Option<crate::state_file::ProjectSession>> {
648 if pid == 0 {
649 return Err(miette::miette!("invalid host PID 0"));
650 }
651 if pid > i32::MAX as u32 {
652 return Err(miette::miette!("host PID {pid} exceeds i32::MAX"));
653 }
654 PROCS.refresh_pids(&[pid]);
655 let liveness_title = PROCS.title(pid);
656 #[cfg(unix)]
662 if !PROCS.is_running(pid) {
663 return Err(miette::miette!("host PID {pid} is not running"));
664 }
665 let mut state_file = self.state_file.lock().await;
666 let previous = state_file.set_project_session(
667 pid,
668 dir,
669 crate::state_file::ProjectSession { liveness_title },
670 );
671 Ok(previous)
672 }
673
674 pub(crate) async fn leave_project_session(
678 &self,
679 pid: u32,
680 dir: &std::path::Path,
681 ) -> Result<Option<PathBuf>> {
682 let mut state_file = self.state_file.lock().await;
683 if state_file.remove_project_session(pid, dir).is_some() {
684 Ok(Some(dir.to_path_buf()))
685 } else {
686 Ok(None)
687 }
688 }
689}
690
691#[cfg(test)]
692mod tests {
693 use super::*;
694
695 fn daemon(id: &DaemonId, pid: Option<u32>, dir: Option<PathBuf>) -> Daemon {
696 Daemon {
697 id: id.clone(),
698 pid,
699 dir,
700 ..Daemon::default()
701 }
702 }
703
704 #[test]
705 fn clean_filters_intersect_and_preserve_running_daemons() {
706 let api = DaemonId::new("project-a", "api");
707 let worker = DaemonId::new("project-a", "worker");
708 let namespaces = vec!["project-a".to_string()];
709 let daemons = vec![api.clone()];
710
711 assert!(should_clean_daemon(
712 &api,
713 &daemon(&api, None, None),
714 &namespaces,
715 &daemons,
716 ));
717 assert!(!should_clean_daemon(
718 &worker,
719 &daemon(&worker, None, None),
720 &namespaces,
721 &daemons,
722 ));
723 assert!(!should_clean_daemon(
724 &api,
725 &daemon(&api, Some(42), None),
726 &namespaces,
727 &daemons,
728 ));
729 }
730
731 #[test]
732 fn prune_candidates_require_a_recorded_directory() {
733 let id = DaemonId::new("project-a", "api");
734 assert!(prune_candidate(&id, &daemon(&id, None, None), &[], &[]).is_none());
735 assert_eq!(
736 prune_candidate(
737 &id,
738 &daemon(&id, None, Some(PathBuf::from("missing"))),
739 &[],
740 &[]
741 ),
742 Some((id, PathBuf::from("missing")))
743 );
744 }
745}