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 ..Daemon::default()
241 }
242}
243
244fn build_daemon_list(
250 state_daemons: Vec<Daemon>,
251 disabled_set: HashSet<DaemonId>,
252 config: PitchforkToml,
253 ns_registry: IndexMap<String, NamespaceEntry>,
254 filter: &NamespaceFilter,
255) -> Result<Vec<DaemonListEntry>> {
256 let mut entries = Vec::new();
257 let mut seen_ids = HashSet::new();
258
259 let pitchfork_id = DaemonId::pitchfork();
261
262 for daemon in state_daemons {
264 if daemon.id == pitchfork_id || !filter.matches(&daemon.id) {
265 continue; }
267
268 seen_ids.insert(daemon.id.clone());
273 entries.push(DaemonListEntry {
274 id: daemon.id.clone(),
275 is_disabled: disabled_set.contains(&daemon.id),
276 is_available: daemon.config_registered,
277 daemon,
278 });
279 }
280
281 for (daemon_id, daemon_config) in &config.daemons {
283 if *daemon_id == pitchfork_id || seen_ids.contains(daemon_id) || !filter.matches(daemon_id)
284 {
285 continue;
286 }
287
288 let placeholder = build_placeholder_daemon(daemon_id, daemon_config);
289
290 entries.push(DaemonListEntry {
291 id: daemon_id.clone(),
292 daemon: placeholder,
293 is_disabled: disabled_set.contains(daemon_id),
294 is_available: true,
295 });
296 seen_ids.insert(daemon_id.clone());
297 }
298
299 for (ns_name, entry) in ns_registry {
303 match PitchforkToml::all_merged_from(&entry.dir) {
304 Ok(ns_config) => {
305 for (daemon_id, daemon_config) in &ns_config.daemons {
306 if *daemon_id == pitchfork_id
307 || seen_ids.contains(daemon_id)
308 || !filter.matches(daemon_id)
309 {
310 continue;
311 }
312 let placeholder = build_placeholder_daemon(daemon_id, daemon_config);
313 entries.push(DaemonListEntry {
314 id: daemon_id.clone(),
315 daemon: placeholder,
316 is_disabled: disabled_set.contains(daemon_id),
317 is_available: true,
318 });
319 seen_ids.insert(daemon_id.clone());
320 }
321 }
322 Err(e) => {
323 log::warn!(
324 "Failed to load namespace '{ns_name}' from {}: {e}",
325 entry.dir.display()
326 );
327 }
328 }
329 }
330
331 Ok(entries)
332}
333
334#[cfg(test)]
335mod tests {
336 use super::*;
337 use crate::pitchfork_toml::PitchforkTomlDaemon;
338 use std::path::PathBuf;
339
340 fn state_daemon(ns: &str, name: &str) -> Daemon {
341 Daemon {
342 id: DaemonId::new(ns, name),
343 ..Daemon::default()
344 }
345 }
346
347 fn config_with(daemons: &[(&str, &str)]) -> PitchforkToml {
348 let mut pt = PitchforkToml::new(PathBuf::from("/tmp/pitchfork.toml"));
349 for (ns, name) in daemons {
350 pt.daemons
351 .insert(DaemonId::new(*ns, *name), PitchforkTomlDaemon::default());
352 }
353 pt
354 }
355
356 fn qualified_ids(entries: &[DaemonListEntry]) -> Vec<String> {
357 entries.iter().map(|e| e.id.qualified()).collect()
358 }
359
360 #[test]
361 fn test_empty_filter_matches_everything() {
362 let filter = NamespaceFilter::default();
363 assert!(filter.is_empty());
364 assert!(filter.matches(&DaemonId::new("frontend", "api")));
365 assert!(filter.matches(&DaemonId::new("global", "postgres")));
366 assert_eq!(filter.single(), None);
367 }
368
369 #[test]
370 fn test_filter_matches_only_listed_namespaces() {
371 let filter = NamespaceFilter::new(vec!["frontend".to_string()]);
372 assert!(!filter.is_empty());
373 assert!(filter.matches(&DaemonId::new("frontend", "api")));
374 assert!(!filter.matches(&DaemonId::new("backend", "api")));
375 assert!(!filter.matches(&DaemonId::new("global", "postgres")));
376 }
377
378 #[test]
379 fn test_filter_multiple_namespaces_union() {
380 let filter = NamespaceFilter::new(vec!["frontend".to_string(), "backend".to_string()]);
381 assert!(filter.matches(&DaemonId::new("frontend", "api")));
382 assert!(filter.matches(&DaemonId::new("backend", "api")));
383 assert!(!filter.matches(&DaemonId::new("global", "postgres")));
384 assert_eq!(filter.single(), None);
385 }
386
387 #[test]
388 fn test_filter_single() {
389 let filter = NamespaceFilter::new(vec!["frontend".to_string()]);
390 assert_eq!(filter.single(), Some("frontend"));
391
392 let filter = NamespaceFilter::new(vec!["frontend".to_string(), "frontend".to_string()]);
394 assert_eq!(filter.single(), Some("frontend"));
395 }
396
397 #[test]
398 fn test_from_flags_validates_namespaces() {
399 let filter = NamespaceFilter::from_flags(&["frontend".to_string()], false).unwrap();
401 assert_eq!(filter.single(), Some("frontend"));
402
403 assert!(NamespaceFilter::from_flags(&["my--ns".to_string()], false).is_err());
405 assert!(NamespaceFilter::from_flags(&["has space".to_string()], false).is_err());
406 assert!(NamespaceFilter::from_flags(&["a/b".to_string()], false).is_err());
407 assert!(NamespaceFilter::from_flags(&[String::new()], false).is_err());
408 }
409
410 #[test]
411 fn test_from_flags_dedups() {
412 let filter =
413 NamespaceFilter::from_flags(&["frontend".to_string(), "frontend".to_string()], false)
414 .unwrap();
415 assert_eq!(filter.single(), Some("frontend"));
416 }
417
418 #[test]
419 fn test_build_daemon_list_unfiltered_keeps_all_namespaces() {
420 let state = vec![
421 state_daemon("frontend", "api"),
422 state_daemon("backend", "api"),
423 ];
424 let config = config_with(&[("frontend", "worker")]);
425 let entries = build_daemon_list(
426 state,
427 HashSet::new(),
428 config,
429 IndexMap::new(),
430 &NamespaceFilter::default(),
431 )
432 .unwrap();
433 let ids = qualified_ids(&entries);
434 assert!(ids.contains(&"frontend/api".to_string()));
435 assert!(ids.contains(&"backend/api".to_string()));
436 assert!(ids.contains(&"frontend/worker".to_string()));
437 }
438
439 #[test]
440 fn test_build_daemon_list_filters_state_and_config_daemons() {
441 let state = vec![
442 state_daemon("frontend", "api"),
443 state_daemon("backend", "api"),
444 ];
445 let config = config_with(&[("frontend", "worker"), ("backend", "worker")]);
446 let entries = build_daemon_list(
447 state,
448 HashSet::new(),
449 config,
450 IndexMap::new(),
451 &NamespaceFilter::new(vec!["frontend".to_string()]),
452 )
453 .unwrap();
454 let ids = qualified_ids(&entries);
455 assert_eq!(ids, vec!["frontend/api", "frontend/worker"]);
456 }
457
458 #[test]
459 fn test_build_daemon_list_filter_union_of_namespaces() {
460 let state = vec![
461 state_daemon("frontend", "api"),
462 state_daemon("backend", "api"),
463 state_daemon("global", "postgres"),
464 ];
465 let entries = build_daemon_list(
466 state,
467 HashSet::new(),
468 config_with(&[]),
469 IndexMap::new(),
470 &NamespaceFilter::new(vec!["frontend".to_string(), "global".to_string()]),
471 )
472 .unwrap();
473 let ids = qualified_ids(&entries);
474 assert!(ids.contains(&"frontend/api".to_string()));
475 assert!(ids.contains(&"global/postgres".to_string()));
476 assert!(!ids.contains(&"backend/api".to_string()));
477 }
478
479 #[test]
480 fn test_build_daemon_list_filter_preserves_disabled_flag() {
481 let state = vec![state_daemon("frontend", "api")];
482 let disabled: HashSet<DaemonId> = [DaemonId::new("frontend", "api")].into_iter().collect();
483 let entries = build_daemon_list(
484 state,
485 disabled,
486 config_with(&[]),
487 IndexMap::new(),
488 &NamespaceFilter::new(vec!["frontend".to_string()]),
489 )
490 .unwrap();
491 assert_eq!(entries.len(), 1);
492 assert!(entries[0].is_disabled);
493 }
494}