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