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
10#[derive(Debug, Clone, Default)]
15pub struct NamespaceFilter {
16 namespaces: Vec<String>,
17}
18
19impl NamespaceFilter {
20 pub fn new(mut namespaces: Vec<String>) -> Self {
22 namespaces.sort();
23 namespaces.dedup();
24 Self { namespaces }
25 }
26
27 pub fn from_flags(namespaces: &[String], project: bool) -> Result<Self> {
33 let mut all = Vec::with_capacity(namespaces.len() + 1);
34 for ns in namespaces {
35 validate_namespace(ns)?;
36 all.push(ns.clone());
37 }
38 if project {
39 all.push(PitchforkToml::namespace_for_dir(&crate::env::CWD)?);
40 }
41 Ok(Self::new(all))
42 }
43
44 pub fn is_empty(&self) -> bool {
46 self.namespaces.is_empty()
47 }
48
49 pub fn matches(&self, id: &DaemonId) -> bool {
51 self.namespaces.is_empty() || self.namespaces.iter().any(|ns| ns == id.namespace())
52 }
53
54 pub fn single(&self) -> Option<&str> {
56 match self.namespaces.as_slice() {
57 [ns] => Some(ns),
58 _ => None,
59 }
60 }
61}
62
63#[derive(Debug, Clone)]
65pub struct DaemonListEntry {
66 pub id: DaemonId,
67 pub daemon: Daemon,
68 pub is_disabled: bool,
69 pub is_available: bool, }
71
72pub async fn get_all_daemons(
89 client: &IpcClient,
90 filter: &NamespaceFilter,
91) -> Result<Vec<DaemonListEntry>> {
92 let config = PitchforkToml::all_merged()?;
93
94 let state_file = crate::state_file::StateFile::read(&*crate::env::PITCHFORK_STATE_FILE)?;
96 let state_daemons: Vec<Daemon> = state_file.daemons.values().cloned().collect();
97
98 let disabled_daemons = client.get_disabled_daemons().await?;
99 let disabled_set: HashSet<DaemonId> = disabled_daemons.into_iter().collect();
100
101 build_daemon_list(
102 state_daemons,
103 disabled_set,
104 config,
105 PitchforkToml::read_global_namespaces(),
106 filter,
107 )
108}
109
110pub async fn get_all_daemons_direct(
121 supervisor: &crate::supervisor::Supervisor,
122) -> Result<Vec<DaemonListEntry>> {
123 let config = PitchforkToml::all_merged()?;
124
125 let state_file = supervisor.state_file.lock().await;
127 let state_daemons: Vec<Daemon> = state_file.daemons.values().cloned().collect();
128 let disabled_set: HashSet<DaemonId> = state_file.disabled.clone().into_iter().collect();
129 drop(state_file); build_daemon_list(
132 state_daemons,
133 disabled_set,
134 config,
135 PitchforkToml::read_global_namespaces(),
136 &NamespaceFilter::default(),
137 )
138}
139
140pub async fn get_daemon_direct(
145 supervisor: &crate::supervisor::Supervisor,
146 id: &DaemonId,
147) -> Result<Option<DaemonListEntry>> {
148 let pitchfork_id = DaemonId::pitchfork();
149 if *id == pitchfork_id {
150 return Ok(None);
151 }
152
153 let state_file = supervisor.state_file.lock().await;
155 if let Some(daemon) = state_file.daemons.get(id).cloned() {
156 let is_disabled = state_file.disabled.contains(id);
157 drop(state_file);
158 return Ok(Some(DaemonListEntry {
159 id: id.clone(),
160 is_available: daemon.config_registered,
161 daemon,
162 is_disabled,
163 }));
164 }
165 let is_disabled = state_file.disabled.contains(id);
166 drop(state_file);
167
168 let config = PitchforkToml::all_merged()?;
170 if let Some(daemon_config) = config.daemons.get(id) {
171 return Ok(Some(DaemonListEntry {
172 id: id.clone(),
173 daemon: build_placeholder_daemon(id, daemon_config),
174 is_disabled,
175 is_available: true,
176 }));
177 }
178
179 let namespaces = PitchforkToml::read_global_namespaces();
181 for (_, entry) in namespaces {
182 match PitchforkToml::all_merged_from(&entry.dir) {
183 Ok(ns_config) => {
184 if let Some(daemon_config) = ns_config.daemons.get(id) {
185 return Ok(Some(DaemonListEntry {
186 id: id.clone(),
187 daemon: build_placeholder_daemon(id, daemon_config),
188 is_disabled,
189 is_available: true,
190 }));
191 }
192 }
193 Err(e) => {
194 log::warn!("Failed to load namespace from {}: {e}", entry.dir.display());
195 }
196 }
197 }
198
199 Ok(None)
200}
201
202pub fn build_placeholder_daemon(
204 id: &DaemonId,
205 daemon_config: &crate::pitchfork_toml::PitchforkTomlDaemon,
206) -> Daemon {
207 Daemon {
208 id: id.clone(),
209 status: DaemonStatus::Stopped,
210 port: daemon_config.port.clone(),
211 depends: vec![],
212 env: None,
213 watch: vec![],
214 watch_mode: daemon_config.watch_mode,
215 watch_base_dir: None,
216 mise: daemon_config.mise,
217 user: daemon_config.user.clone(),
218 active_port: None,
219 slug: None,
220 proxy: None,
221 memory_limit: daemon_config.memory_limit,
222 cpu_limit: daemon_config.cpu_limit,
223 ..Daemon::default()
224 }
225}
226
227fn build_daemon_list(
233 state_daemons: Vec<Daemon>,
234 disabled_set: HashSet<DaemonId>,
235 config: PitchforkToml,
236 ns_registry: IndexMap<String, NamespaceEntry>,
237 filter: &NamespaceFilter,
238) -> Result<Vec<DaemonListEntry>> {
239 let mut entries = Vec::new();
240 let mut seen_ids = HashSet::new();
241
242 let pitchfork_id = DaemonId::pitchfork();
244
245 for daemon in state_daemons {
247 if daemon.id == pitchfork_id || !filter.matches(&daemon.id) {
248 continue; }
250
251 seen_ids.insert(daemon.id.clone());
256 entries.push(DaemonListEntry {
257 id: daemon.id.clone(),
258 is_disabled: disabled_set.contains(&daemon.id),
259 is_available: daemon.config_registered,
260 daemon,
261 });
262 }
263
264 for (daemon_id, daemon_config) in &config.daemons {
266 if *daemon_id == pitchfork_id || seen_ids.contains(daemon_id) || !filter.matches(daemon_id)
267 {
268 continue;
269 }
270
271 let placeholder = build_placeholder_daemon(daemon_id, daemon_config);
272
273 entries.push(DaemonListEntry {
274 id: daemon_id.clone(),
275 daemon: placeholder,
276 is_disabled: disabled_set.contains(daemon_id),
277 is_available: true,
278 });
279 seen_ids.insert(daemon_id.clone());
280 }
281
282 for (ns_name, entry) in ns_registry {
286 match PitchforkToml::all_merged_from(&entry.dir) {
287 Ok(ns_config) => {
288 for (daemon_id, daemon_config) in &ns_config.daemons {
289 if *daemon_id == pitchfork_id
290 || seen_ids.contains(daemon_id)
291 || !filter.matches(daemon_id)
292 {
293 continue;
294 }
295 let placeholder = build_placeholder_daemon(daemon_id, daemon_config);
296 entries.push(DaemonListEntry {
297 id: daemon_id.clone(),
298 daemon: placeholder,
299 is_disabled: disabled_set.contains(daemon_id),
300 is_available: true,
301 });
302 seen_ids.insert(daemon_id.clone());
303 }
304 }
305 Err(e) => {
306 log::warn!(
307 "Failed to load namespace '{ns_name}' from {}: {e}",
308 entry.dir.display()
309 );
310 }
311 }
312 }
313
314 Ok(entries)
315}
316
317#[cfg(test)]
318mod tests {
319 use super::*;
320 use crate::pitchfork_toml::PitchforkTomlDaemon;
321 use std::path::PathBuf;
322
323 fn state_daemon(ns: &str, name: &str) -> Daemon {
324 Daemon {
325 id: DaemonId::new(ns, name),
326 ..Daemon::default()
327 }
328 }
329
330 fn config_with(daemons: &[(&str, &str)]) -> PitchforkToml {
331 let mut pt = PitchforkToml::new(PathBuf::from("/tmp/pitchfork.toml"));
332 for (ns, name) in daemons {
333 pt.daemons
334 .insert(DaemonId::new(*ns, *name), PitchforkTomlDaemon::default());
335 }
336 pt
337 }
338
339 fn qualified_ids(entries: &[DaemonListEntry]) -> Vec<String> {
340 entries.iter().map(|e| e.id.qualified()).collect()
341 }
342
343 #[test]
344 fn test_empty_filter_matches_everything() {
345 let filter = NamespaceFilter::default();
346 assert!(filter.is_empty());
347 assert!(filter.matches(&DaemonId::new("frontend", "api")));
348 assert!(filter.matches(&DaemonId::new("global", "postgres")));
349 assert_eq!(filter.single(), None);
350 }
351
352 #[test]
353 fn test_filter_matches_only_listed_namespaces() {
354 let filter = NamespaceFilter::new(vec!["frontend".to_string()]);
355 assert!(!filter.is_empty());
356 assert!(filter.matches(&DaemonId::new("frontend", "api")));
357 assert!(!filter.matches(&DaemonId::new("backend", "api")));
358 assert!(!filter.matches(&DaemonId::new("global", "postgres")));
359 }
360
361 #[test]
362 fn test_filter_multiple_namespaces_union() {
363 let filter = NamespaceFilter::new(vec!["frontend".to_string(), "backend".to_string()]);
364 assert!(filter.matches(&DaemonId::new("frontend", "api")));
365 assert!(filter.matches(&DaemonId::new("backend", "api")));
366 assert!(!filter.matches(&DaemonId::new("global", "postgres")));
367 assert_eq!(filter.single(), None);
368 }
369
370 #[test]
371 fn test_filter_single() {
372 let filter = NamespaceFilter::new(vec!["frontend".to_string()]);
373 assert_eq!(filter.single(), Some("frontend"));
374
375 let filter = NamespaceFilter::new(vec!["frontend".to_string(), "frontend".to_string()]);
377 assert_eq!(filter.single(), Some("frontend"));
378 }
379
380 #[test]
381 fn test_from_flags_validates_namespaces() {
382 let filter = NamespaceFilter::from_flags(&["frontend".to_string()], false).unwrap();
384 assert_eq!(filter.single(), Some("frontend"));
385
386 assert!(NamespaceFilter::from_flags(&["my--ns".to_string()], false).is_err());
388 assert!(NamespaceFilter::from_flags(&["has space".to_string()], false).is_err());
389 assert!(NamespaceFilter::from_flags(&["a/b".to_string()], false).is_err());
390 assert!(NamespaceFilter::from_flags(&[String::new()], false).is_err());
391 }
392
393 #[test]
394 fn test_from_flags_dedups() {
395 let filter =
396 NamespaceFilter::from_flags(&["frontend".to_string(), "frontend".to_string()], false)
397 .unwrap();
398 assert_eq!(filter.single(), Some("frontend"));
399 }
400
401 #[test]
402 fn test_build_daemon_list_unfiltered_keeps_all_namespaces() {
403 let state = vec![
404 state_daemon("frontend", "api"),
405 state_daemon("backend", "api"),
406 ];
407 let config = config_with(&[("frontend", "worker")]);
408 let entries = build_daemon_list(
409 state,
410 HashSet::new(),
411 config,
412 IndexMap::new(),
413 &NamespaceFilter::default(),
414 )
415 .unwrap();
416 let ids = qualified_ids(&entries);
417 assert!(ids.contains(&"frontend/api".to_string()));
418 assert!(ids.contains(&"backend/api".to_string()));
419 assert!(ids.contains(&"frontend/worker".to_string()));
420 }
421
422 #[test]
423 fn test_build_daemon_list_filters_state_and_config_daemons() {
424 let state = vec![
425 state_daemon("frontend", "api"),
426 state_daemon("backend", "api"),
427 ];
428 let config = config_with(&[("frontend", "worker"), ("backend", "worker")]);
429 let entries = build_daemon_list(
430 state,
431 HashSet::new(),
432 config,
433 IndexMap::new(),
434 &NamespaceFilter::new(vec!["frontend".to_string()]),
435 )
436 .unwrap();
437 let ids = qualified_ids(&entries);
438 assert_eq!(ids, vec!["frontend/api", "frontend/worker"]);
439 }
440
441 #[test]
442 fn test_build_daemon_list_filter_union_of_namespaces() {
443 let state = vec![
444 state_daemon("frontend", "api"),
445 state_daemon("backend", "api"),
446 state_daemon("global", "postgres"),
447 ];
448 let entries = build_daemon_list(
449 state,
450 HashSet::new(),
451 config_with(&[]),
452 IndexMap::new(),
453 &NamespaceFilter::new(vec!["frontend".to_string(), "global".to_string()]),
454 )
455 .unwrap();
456 let ids = qualified_ids(&entries);
457 assert!(ids.contains(&"frontend/api".to_string()));
458 assert!(ids.contains(&"global/postgres".to_string()));
459 assert!(!ids.contains(&"backend/api".to_string()));
460 }
461
462 #[test]
463 fn test_build_daemon_list_filter_preserves_disabled_flag() {
464 let state = vec![state_daemon("frontend", "api")];
465 let disabled: HashSet<DaemonId> = [DaemonId::new("frontend", "api")].into_iter().collect();
466 let entries = build_daemon_list(
467 state,
468 disabled,
469 config_with(&[]),
470 IndexMap::new(),
471 &NamespaceFilter::new(vec!["frontend".to_string()]),
472 )
473 .unwrap();
474 assert_eq!(entries.len(), 1);
475 assert!(entries[0].is_disabled);
476 }
477}