1use crate::Result;
2use crate::daemon::Daemon;
3use crate::daemon_id::{DaemonId, validate_namespace};
4use crate::daemon_status::DaemonStatus;
5use crate::ipc::client::IpcClient;
6use crate::pitchfork_toml::{NamespaceEntry, PitchforkToml};
7use indexmap::IndexMap;
8use std::collections::HashSet;
9
10pub fn completed_oneshots() -> HashSet<DaemonId> {
18 crate::state_file::StateFile::get()
19 .daemons
20 .iter()
21 .filter(|(_, d)| d.oneshot && d.status.is_completed())
22 .map(|(id, _)| id.clone())
23 .collect()
24}
25
26#[derive(Debug, Clone, Default)]
31pub struct NamespaceFilter {
32 namespaces: Vec<String>,
33}
34
35impl NamespaceFilter {
36 pub fn new(mut namespaces: Vec<String>) -> Self {
38 namespaces.sort();
39 namespaces.dedup();
40 Self { namespaces }
41 }
42
43 pub fn from_flags(namespaces: &[String], project: bool) -> Result<Self> {
49 let mut all = Vec::with_capacity(namespaces.len() + 1);
50 for ns in namespaces {
51 validate_namespace(ns)?;
52 all.push(ns.clone());
53 }
54 if project {
55 all.push(PitchforkToml::namespace_for_dir(&crate::env::CWD)?);
56 }
57 Ok(Self::new(all))
58 }
59
60 pub fn is_empty(&self) -> bool {
62 self.namespaces.is_empty()
63 }
64
65 pub fn matches(&self, id: &DaemonId) -> bool {
67 self.namespaces.is_empty() || self.namespaces.iter().any(|ns| ns == id.namespace())
68 }
69
70 pub fn single(&self) -> Option<&str> {
72 match self.namespaces.as_slice() {
73 [ns] => Some(ns),
74 _ => None,
75 }
76 }
77}
78
79#[derive(Debug, Clone)]
81pub struct DaemonListEntry {
82 pub id: DaemonId,
83 pub daemon: Daemon,
84 pub is_disabled: bool,
85 pub is_available: bool, }
87
88pub async fn get_all_daemons(
105 client: &IpcClient,
106 filter: &NamespaceFilter,
107) -> Result<Vec<DaemonListEntry>> {
108 let config = PitchforkToml::all_merged()?;
109
110 let state_file = crate::state_file::StateFile::read(&*crate::env::PITCHFORK_STATE_FILE)?;
112 let state_daemons: Vec<Daemon> = state_file.daemons.values().cloned().collect();
113
114 let disabled_daemons = client.get_disabled_daemons().await?;
115 let disabled_set: HashSet<DaemonId> = disabled_daemons.into_iter().collect();
116
117 build_daemon_list(
118 state_daemons,
119 disabled_set,
120 config,
121 PitchforkToml::read_global_namespaces(),
122 filter,
123 )
124}
125
126pub async fn get_all_daemons_direct(
137 supervisor: &crate::supervisor::Supervisor,
138) -> Result<Vec<DaemonListEntry>> {
139 let config = PitchforkToml::all_merged()?;
140
141 let state_file = supervisor.state_file.lock().await;
143 let state_daemons: Vec<Daemon> = state_file.daemons.values().cloned().collect();
144 let disabled_set: HashSet<DaemonId> = state_file.disabled.clone().into_iter().collect();
145 drop(state_file); build_daemon_list(
148 state_daemons,
149 disabled_set,
150 config,
151 PitchforkToml::read_global_namespaces(),
152 &NamespaceFilter::default(),
153 )
154}
155
156pub async fn get_daemon_direct(
161 supervisor: &crate::supervisor::Supervisor,
162 id: &DaemonId,
163) -> Result<Option<DaemonListEntry>> {
164 let pitchfork_id = DaemonId::pitchfork();
165 if *id == pitchfork_id {
166 return Ok(None);
167 }
168
169 let state_file = supervisor.state_file.lock().await;
171 if let Some(daemon) = state_file.daemons.get(id).cloned() {
172 let is_disabled = state_file.disabled.contains(id);
173 drop(state_file);
174 return Ok(Some(DaemonListEntry {
175 id: id.clone(),
176 is_available: daemon.config_registered,
177 daemon,
178 is_disabled,
179 }));
180 }
181 let is_disabled = state_file.disabled.contains(id);
182 drop(state_file);
183
184 let config = PitchforkToml::all_merged()?;
186 if let Some(daemon_config) = config.daemons.get(id) {
187 return Ok(Some(DaemonListEntry {
188 id: id.clone(),
189 daemon: build_placeholder_daemon(id, daemon_config),
190 is_disabled,
191 is_available: true,
192 }));
193 }
194
195 let namespaces = PitchforkToml::read_global_namespaces();
197 for (_, entry) in namespaces {
198 match PitchforkToml::all_merged_from(&entry.dir) {
199 Ok(ns_config) => {
200 if let Some(daemon_config) = ns_config.daemons.get(id) {
201 return Ok(Some(DaemonListEntry {
202 id: id.clone(),
203 daemon: build_placeholder_daemon(id, daemon_config),
204 is_disabled,
205 is_available: true,
206 }));
207 }
208 }
209 Err(e) => {
210 log::warn!("Failed to load namespace from {}: {e}", entry.dir.display());
211 }
212 }
213 }
214
215 Ok(None)
216}
217
218pub fn build_placeholder_daemon(
220 id: &DaemonId,
221 daemon_config: &crate::pitchfork_toml::PitchforkTomlDaemon,
222) -> Daemon {
223 Daemon {
224 id: id.clone(),
225 status: DaemonStatus::Stopped,
226 oneshot: daemon_config.is_oneshot(),
227 port: daemon_config.port.clone(),
228 depends: vec![],
229 env: None,
230 watch: vec![],
231 watch_mode: daemon_config.watch_mode,
232 watch_base_dir: None,
233 mise: daemon_config.mise,
234 user: daemon_config.user.clone(),
235 active_port: None,
236 slug: None,
237 proxy: None,
238 memory_limit: daemon_config.memory_limit,
239 cpu_limit: daemon_config.cpu_limit,
240 cron_schedule: daemon_config.cron.as_ref().map(|c| c.schedule.clone()),
245 cron_retrigger: daemon_config.cron.as_ref().map(|c| c.retrigger),
246 cron_immediate: daemon_config.cron.as_ref().map(|c| c.immediate),
247 ..Daemon::default()
248 }
249}
250
251fn build_daemon_list(
257 state_daemons: Vec<Daemon>,
258 disabled_set: HashSet<DaemonId>,
259 config: PitchforkToml,
260 ns_registry: IndexMap<String, NamespaceEntry>,
261 filter: &NamespaceFilter,
262) -> Result<Vec<DaemonListEntry>> {
263 let mut entries = Vec::new();
264 let mut seen_ids = HashSet::new();
265
266 let pitchfork_id = DaemonId::pitchfork();
268
269 for daemon in state_daemons {
271 if daemon.id == pitchfork_id || !filter.matches(&daemon.id) {
272 continue; }
274
275 seen_ids.insert(daemon.id.clone());
280 entries.push(DaemonListEntry {
281 id: daemon.id.clone(),
282 is_disabled: disabled_set.contains(&daemon.id),
283 is_available: daemon.config_registered,
284 daemon,
285 });
286 }
287
288 for (daemon_id, daemon_config) in &config.daemons {
290 if *daemon_id == pitchfork_id || seen_ids.contains(daemon_id) || !filter.matches(daemon_id)
291 {
292 continue;
293 }
294
295 let placeholder = build_placeholder_daemon(daemon_id, daemon_config);
296
297 entries.push(DaemonListEntry {
298 id: daemon_id.clone(),
299 daemon: placeholder,
300 is_disabled: disabled_set.contains(daemon_id),
301 is_available: true,
302 });
303 seen_ids.insert(daemon_id.clone());
304 }
305
306 for (ns_name, entry) in ns_registry {
310 match PitchforkToml::all_merged_from(&entry.dir) {
311 Ok(ns_config) => {
312 for (daemon_id, daemon_config) in &ns_config.daemons {
313 if *daemon_id == pitchfork_id
314 || seen_ids.contains(daemon_id)
315 || !filter.matches(daemon_id)
316 {
317 continue;
318 }
319 let placeholder = build_placeholder_daemon(daemon_id, daemon_config);
320 entries.push(DaemonListEntry {
321 id: daemon_id.clone(),
322 daemon: placeholder,
323 is_disabled: disabled_set.contains(daemon_id),
324 is_available: true,
325 });
326 seen_ids.insert(daemon_id.clone());
327 }
328 }
329 Err(e) => {
330 log::warn!(
331 "Failed to load namespace '{ns_name}' from {}: {e}",
332 entry.dir.display()
333 );
334 }
335 }
336 }
337
338 Ok(entries)
339}
340
341#[cfg(test)]
342mod tests {
343 use super::*;
344 use crate::pitchfork_toml::PitchforkTomlDaemon;
345 use std::path::PathBuf;
346
347 fn state_daemon(ns: &str, name: &str) -> Daemon {
348 Daemon {
349 id: DaemonId::new(ns, name),
350 ..Daemon::default()
351 }
352 }
353
354 fn config_with(daemons: &[(&str, &str)]) -> PitchforkToml {
355 let mut pt = PitchforkToml::new(PathBuf::from("/tmp/pitchfork.toml"));
356 for (ns, name) in daemons {
357 pt.daemons
358 .insert(DaemonId::new(*ns, *name), PitchforkTomlDaemon::default());
359 }
360 pt
361 }
362
363 fn qualified_ids(entries: &[DaemonListEntry]) -> Vec<String> {
364 entries.iter().map(|e| e.id.qualified()).collect()
365 }
366
367 #[test]
368 fn test_empty_filter_matches_everything() {
369 let filter = NamespaceFilter::default();
370 assert!(filter.is_empty());
371 assert!(filter.matches(&DaemonId::new("frontend", "api")));
372 assert!(filter.matches(&DaemonId::new("global", "postgres")));
373 assert_eq!(filter.single(), None);
374 }
375
376 #[test]
377 fn test_filter_matches_only_listed_namespaces() {
378 let filter = NamespaceFilter::new(vec!["frontend".to_string()]);
379 assert!(!filter.is_empty());
380 assert!(filter.matches(&DaemonId::new("frontend", "api")));
381 assert!(!filter.matches(&DaemonId::new("backend", "api")));
382 assert!(!filter.matches(&DaemonId::new("global", "postgres")));
383 }
384
385 #[test]
386 fn test_filter_multiple_namespaces_union() {
387 let filter = NamespaceFilter::new(vec!["frontend".to_string(), "backend".to_string()]);
388 assert!(filter.matches(&DaemonId::new("frontend", "api")));
389 assert!(filter.matches(&DaemonId::new("backend", "api")));
390 assert!(!filter.matches(&DaemonId::new("global", "postgres")));
391 assert_eq!(filter.single(), None);
392 }
393
394 #[test]
395 fn test_filter_single() {
396 let filter = NamespaceFilter::new(vec!["frontend".to_string()]);
397 assert_eq!(filter.single(), Some("frontend"));
398
399 let filter = NamespaceFilter::new(vec!["frontend".to_string(), "frontend".to_string()]);
401 assert_eq!(filter.single(), Some("frontend"));
402 }
403
404 #[test]
405 fn test_from_flags_validates_namespaces() {
406 let filter = NamespaceFilter::from_flags(&["frontend".to_string()], false).unwrap();
408 assert_eq!(filter.single(), Some("frontend"));
409
410 assert!(NamespaceFilter::from_flags(&["my--ns".to_string()], false).is_err());
412 assert!(NamespaceFilter::from_flags(&["has space".to_string()], false).is_err());
413 assert!(NamespaceFilter::from_flags(&["a/b".to_string()], false).is_err());
414 assert!(NamespaceFilter::from_flags(&[String::new()], false).is_err());
415 }
416
417 #[test]
418 fn test_from_flags_dedups() {
419 let filter =
420 NamespaceFilter::from_flags(&["frontend".to_string(), "frontend".to_string()], false)
421 .unwrap();
422 assert_eq!(filter.single(), Some("frontend"));
423 }
424
425 #[test]
426 fn test_build_daemon_list_unfiltered_keeps_all_namespaces() {
427 let state = vec![
428 state_daemon("frontend", "api"),
429 state_daemon("backend", "api"),
430 ];
431 let config = config_with(&[("frontend", "worker")]);
432 let entries = build_daemon_list(
433 state,
434 HashSet::new(),
435 config,
436 IndexMap::new(),
437 &NamespaceFilter::default(),
438 )
439 .unwrap();
440 let ids = qualified_ids(&entries);
441 assert!(ids.contains(&"frontend/api".to_string()));
442 assert!(ids.contains(&"backend/api".to_string()));
443 assert!(ids.contains(&"frontend/worker".to_string()));
444 }
445
446 #[test]
447 fn test_build_daemon_list_filters_state_and_config_daemons() {
448 let state = vec![
449 state_daemon("frontend", "api"),
450 state_daemon("backend", "api"),
451 ];
452 let config = config_with(&[("frontend", "worker"), ("backend", "worker")]);
453 let entries = build_daemon_list(
454 state,
455 HashSet::new(),
456 config,
457 IndexMap::new(),
458 &NamespaceFilter::new(vec!["frontend".to_string()]),
459 )
460 .unwrap();
461 let ids = qualified_ids(&entries);
462 assert_eq!(ids, vec!["frontend/api", "frontend/worker"]);
463 }
464
465 #[test]
466 fn test_build_daemon_list_filter_union_of_namespaces() {
467 let state = vec![
468 state_daemon("frontend", "api"),
469 state_daemon("backend", "api"),
470 state_daemon("global", "postgres"),
471 ];
472 let entries = build_daemon_list(
473 state,
474 HashSet::new(),
475 config_with(&[]),
476 IndexMap::new(),
477 &NamespaceFilter::new(vec!["frontend".to_string(), "global".to_string()]),
478 )
479 .unwrap();
480 let ids = qualified_ids(&entries);
481 assert!(ids.contains(&"frontend/api".to_string()));
482 assert!(ids.contains(&"global/postgres".to_string()));
483 assert!(!ids.contains(&"backend/api".to_string()));
484 }
485
486 #[test]
487 fn test_build_daemon_list_filter_preserves_disabled_flag() {
488 let state = vec![state_daemon("frontend", "api")];
489 let disabled: HashSet<DaemonId> = [DaemonId::new("frontend", "api")].into_iter().collect();
490 let entries = build_daemon_list(
491 state,
492 disabled,
493 config_with(&[]),
494 IndexMap::new(),
495 &NamespaceFilter::new(vec!["frontend".to_string()]),
496 )
497 .unwrap();
498 assert_eq!(entries.len(), 1);
499 assert!(entries[0].is_disabled);
500 }
501}