use crate::Result;
use crate::daemon::Daemon;
use crate::daemon_id::{DaemonId, validate_namespace};
use crate::daemon_status::DaemonStatus;
use crate::ipc::client::IpcClient;
use crate::pitchfork_toml::{NamespaceEntry, PitchforkToml};
use indexmap::IndexMap;
use std::collections::HashSet;
pub fn completed_oneshots() -> HashSet<DaemonId> {
crate::state_file::StateFile::get()
.daemons
.iter()
.filter(|(_, d)| d.oneshot && d.status.is_completed())
.map(|(id, _)| id.clone())
.collect()
}
#[derive(Debug, Clone, Default)]
pub struct NamespaceFilter {
namespaces: Vec<String>,
}
impl NamespaceFilter {
pub fn new(mut namespaces: Vec<String>) -> Self {
namespaces.sort();
namespaces.dedup();
Self { namespaces }
}
pub fn from_flags(namespaces: &[String], project: bool) -> Result<Self> {
let mut all = Vec::with_capacity(namespaces.len() + 1);
for ns in namespaces {
validate_namespace(ns)?;
all.push(ns.clone());
}
if project {
all.push(PitchforkToml::namespace_for_dir(&crate::env::CWD)?);
}
Ok(Self::new(all))
}
pub fn is_empty(&self) -> bool {
self.namespaces.is_empty()
}
pub fn matches(&self, id: &DaemonId) -> bool {
self.namespaces.is_empty() || self.namespaces.iter().any(|ns| ns == id.namespace())
}
pub fn single(&self) -> Option<&str> {
match self.namespaces.as_slice() {
[ns] => Some(ns),
_ => None,
}
}
}
#[derive(Debug, Clone)]
pub struct DaemonListEntry {
pub id: DaemonId,
pub daemon: Daemon,
pub is_disabled: bool,
pub is_available: bool, }
pub async fn get_all_daemons(
client: &IpcClient,
filter: &NamespaceFilter,
) -> Result<Vec<DaemonListEntry>> {
let config = PitchforkToml::all_merged()?;
let state_file = crate::state_file::StateFile::read(&*crate::env::PITCHFORK_STATE_FILE)?;
let state_daemons: Vec<Daemon> = state_file.daemons.values().cloned().collect();
let disabled_daemons = client.get_disabled_daemons().await?;
let disabled_set: HashSet<DaemonId> = disabled_daemons.into_iter().collect();
build_daemon_list(
state_daemons,
disabled_set,
config,
PitchforkToml::read_global_namespaces(),
filter,
)
}
pub async fn get_all_daemons_direct(
supervisor: &crate::supervisor::Supervisor,
) -> Result<Vec<DaemonListEntry>> {
let config = PitchforkToml::all_merged()?;
let state_file = supervisor.state_file.lock().await;
let state_daemons: Vec<Daemon> = state_file.daemons.values().cloned().collect();
let disabled_set: HashSet<DaemonId> = state_file.disabled.clone().into_iter().collect();
drop(state_file);
build_daemon_list(
state_daemons,
disabled_set,
config,
PitchforkToml::read_global_namespaces(),
&NamespaceFilter::default(),
)
}
pub async fn get_daemon_direct(
supervisor: &crate::supervisor::Supervisor,
id: &DaemonId,
) -> Result<Option<DaemonListEntry>> {
let pitchfork_id = DaemonId::pitchfork();
if *id == pitchfork_id {
return Ok(None);
}
let state_file = supervisor.state_file.lock().await;
if let Some(daemon) = state_file.daemons.get(id).cloned() {
let is_disabled = state_file.disabled.contains(id);
drop(state_file);
return Ok(Some(DaemonListEntry {
id: id.clone(),
is_available: daemon.config_registered,
daemon,
is_disabled,
}));
}
let is_disabled = state_file.disabled.contains(id);
drop(state_file);
let config = PitchforkToml::all_merged()?;
if let Some(daemon_config) = config.daemons.get(id) {
return Ok(Some(DaemonListEntry {
id: id.clone(),
daemon: build_placeholder_daemon(id, daemon_config),
is_disabled,
is_available: true,
}));
}
let namespaces = PitchforkToml::read_global_namespaces();
for (_, entry) in namespaces {
match PitchforkToml::all_merged_from(&entry.dir) {
Ok(ns_config) => {
if let Some(daemon_config) = ns_config.daemons.get(id) {
return Ok(Some(DaemonListEntry {
id: id.clone(),
daemon: build_placeholder_daemon(id, daemon_config),
is_disabled,
is_available: true,
}));
}
}
Err(e) => {
log::warn!("Failed to load namespace from {}: {e}", entry.dir.display());
}
}
}
Ok(None)
}
pub fn build_placeholder_daemon(
id: &DaemonId,
daemon_config: &crate::pitchfork_toml::PitchforkTomlDaemon,
) -> Daemon {
Daemon {
id: id.clone(),
status: DaemonStatus::Stopped,
oneshot: daemon_config.is_oneshot(),
port: daemon_config.port.clone(),
depends: vec![],
env: None,
watch: vec![],
watch_mode: daemon_config.watch_mode,
watch_base_dir: None,
mise: daemon_config.mise,
user: daemon_config.user.clone(),
active_port: None,
slug: None,
proxy: None,
memory_limit: daemon_config.memory_limit,
cpu_limit: daemon_config.cpu_limit,
cron_schedule: daemon_config.cron.as_ref().map(|c| c.schedule.clone()),
cron_retrigger: daemon_config.cron.as_ref().map(|c| c.retrigger),
cron_immediate: daemon_config.cron.as_ref().map(|c| c.immediate),
..Daemon::default()
}
}
fn build_daemon_list(
state_daemons: Vec<Daemon>,
disabled_set: HashSet<DaemonId>,
config: PitchforkToml,
ns_registry: IndexMap<String, NamespaceEntry>,
filter: &NamespaceFilter,
) -> Result<Vec<DaemonListEntry>> {
let mut entries = Vec::new();
let mut seen_ids = HashSet::new();
let pitchfork_id = DaemonId::pitchfork();
for daemon in state_daemons {
if daemon.id == pitchfork_id || !filter.matches(&daemon.id) {
continue; }
seen_ids.insert(daemon.id.clone());
entries.push(DaemonListEntry {
id: daemon.id.clone(),
is_disabled: disabled_set.contains(&daemon.id),
is_available: daemon.config_registered,
daemon,
});
}
for (daemon_id, daemon_config) in &config.daemons {
if *daemon_id == pitchfork_id || seen_ids.contains(daemon_id) || !filter.matches(daemon_id)
{
continue;
}
let placeholder = build_placeholder_daemon(daemon_id, daemon_config);
entries.push(DaemonListEntry {
id: daemon_id.clone(),
daemon: placeholder,
is_disabled: disabled_set.contains(daemon_id),
is_available: true,
});
seen_ids.insert(daemon_id.clone());
}
for (ns_name, entry) in ns_registry {
match PitchforkToml::all_merged_from(&entry.dir) {
Ok(ns_config) => {
for (daemon_id, daemon_config) in &ns_config.daemons {
if *daemon_id == pitchfork_id
|| seen_ids.contains(daemon_id)
|| !filter.matches(daemon_id)
{
continue;
}
let placeholder = build_placeholder_daemon(daemon_id, daemon_config);
entries.push(DaemonListEntry {
id: daemon_id.clone(),
daemon: placeholder,
is_disabled: disabled_set.contains(daemon_id),
is_available: true,
});
seen_ids.insert(daemon_id.clone());
}
}
Err(e) => {
log::warn!(
"Failed to load namespace '{ns_name}' from {}: {e}",
entry.dir.display()
);
}
}
}
Ok(entries)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::pitchfork_toml::PitchforkTomlDaemon;
use std::path::PathBuf;
fn state_daemon(ns: &str, name: &str) -> Daemon {
Daemon {
id: DaemonId::new(ns, name),
..Daemon::default()
}
}
fn config_with(daemons: &[(&str, &str)]) -> PitchforkToml {
let mut pt = PitchforkToml::new(PathBuf::from("/tmp/pitchfork.toml"));
for (ns, name) in daemons {
pt.daemons
.insert(DaemonId::new(*ns, *name), PitchforkTomlDaemon::default());
}
pt
}
fn qualified_ids(entries: &[DaemonListEntry]) -> Vec<String> {
entries.iter().map(|e| e.id.qualified()).collect()
}
#[test]
fn test_empty_filter_matches_everything() {
let filter = NamespaceFilter::default();
assert!(filter.is_empty());
assert!(filter.matches(&DaemonId::new("frontend", "api")));
assert!(filter.matches(&DaemonId::new("global", "postgres")));
assert_eq!(filter.single(), None);
}
#[test]
fn test_filter_matches_only_listed_namespaces() {
let filter = NamespaceFilter::new(vec!["frontend".to_string()]);
assert!(!filter.is_empty());
assert!(filter.matches(&DaemonId::new("frontend", "api")));
assert!(!filter.matches(&DaemonId::new("backend", "api")));
assert!(!filter.matches(&DaemonId::new("global", "postgres")));
}
#[test]
fn test_filter_multiple_namespaces_union() {
let filter = NamespaceFilter::new(vec!["frontend".to_string(), "backend".to_string()]);
assert!(filter.matches(&DaemonId::new("frontend", "api")));
assert!(filter.matches(&DaemonId::new("backend", "api")));
assert!(!filter.matches(&DaemonId::new("global", "postgres")));
assert_eq!(filter.single(), None);
}
#[test]
fn test_filter_single() {
let filter = NamespaceFilter::new(vec!["frontend".to_string()]);
assert_eq!(filter.single(), Some("frontend"));
let filter = NamespaceFilter::new(vec!["frontend".to_string(), "frontend".to_string()]);
assert_eq!(filter.single(), Some("frontend"));
}
#[test]
fn test_from_flags_validates_namespaces() {
let filter = NamespaceFilter::from_flags(&["frontend".to_string()], false).unwrap();
assert_eq!(filter.single(), Some("frontend"));
assert!(NamespaceFilter::from_flags(&["my--ns".to_string()], false).is_err());
assert!(NamespaceFilter::from_flags(&["has space".to_string()], false).is_err());
assert!(NamespaceFilter::from_flags(&["a/b".to_string()], false).is_err());
assert!(NamespaceFilter::from_flags(&[String::new()], false).is_err());
}
#[test]
fn test_from_flags_dedups() {
let filter =
NamespaceFilter::from_flags(&["frontend".to_string(), "frontend".to_string()], false)
.unwrap();
assert_eq!(filter.single(), Some("frontend"));
}
#[test]
fn test_build_daemon_list_unfiltered_keeps_all_namespaces() {
let state = vec![
state_daemon("frontend", "api"),
state_daemon("backend", "api"),
];
let config = config_with(&[("frontend", "worker")]);
let entries = build_daemon_list(
state,
HashSet::new(),
config,
IndexMap::new(),
&NamespaceFilter::default(),
)
.unwrap();
let ids = qualified_ids(&entries);
assert!(ids.contains(&"frontend/api".to_string()));
assert!(ids.contains(&"backend/api".to_string()));
assert!(ids.contains(&"frontend/worker".to_string()));
}
#[test]
fn test_build_daemon_list_filters_state_and_config_daemons() {
let state = vec![
state_daemon("frontend", "api"),
state_daemon("backend", "api"),
];
let config = config_with(&[("frontend", "worker"), ("backend", "worker")]);
let entries = build_daemon_list(
state,
HashSet::new(),
config,
IndexMap::new(),
&NamespaceFilter::new(vec!["frontend".to_string()]),
)
.unwrap();
let ids = qualified_ids(&entries);
assert_eq!(ids, vec!["frontend/api", "frontend/worker"]);
}
#[test]
fn test_build_daemon_list_filter_union_of_namespaces() {
let state = vec![
state_daemon("frontend", "api"),
state_daemon("backend", "api"),
state_daemon("global", "postgres"),
];
let entries = build_daemon_list(
state,
HashSet::new(),
config_with(&[]),
IndexMap::new(),
&NamespaceFilter::new(vec!["frontend".to_string(), "global".to_string()]),
)
.unwrap();
let ids = qualified_ids(&entries);
assert!(ids.contains(&"frontend/api".to_string()));
assert!(ids.contains(&"global/postgres".to_string()));
assert!(!ids.contains(&"backend/api".to_string()));
}
#[test]
fn test_build_daemon_list_filter_preserves_disabled_flag() {
let state = vec![state_daemon("frontend", "api")];
let disabled: HashSet<DaemonId> = [DaemonId::new("frontend", "api")].into_iter().collect();
let entries = build_daemon_list(
state,
disabled,
config_with(&[]),
IndexMap::new(),
&NamespaceFilter::new(vec!["frontend".to_string()]),
)
.unwrap();
assert_eq!(entries.len(), 1);
assert!(entries[0].is_disabled);
}
}