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