use crate::daemon_id::DaemonId;
use crate::error::{ConfigParseError, DependencyError, FileError, find_similar_daemon};
use crate::settings::SettingsPartial;
use crate::settings::settings;
use crate::state_file::StateFile;
use crate::{Result, env};
use indexmap::IndexMap;
use miette::Context;
use once_cell::sync::Lazy;
use schemars::JsonSchema;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Mutex as StdMutex;
use std::time::SystemTime;
pub use crate::config_types::{
CpuLimit, CronRetrigger, Dir, MemoryLimit, OnOutputHook, PitchforkTomlAuto, PitchforkTomlCron,
PitchforkTomlHooks, PortBump, PortConfig, ReadyCmd, ReadyHttp, ReadyOutput, ReadyPort, Retry,
StopConfig, StopSignal, WatchMode,
};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, JsonSchema)]
pub struct SlugEntryRaw {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub dir: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub namespace: Option<String>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub daemon: Option<String>,
}
#[derive(Debug, Clone)]
pub struct SlugEntry {
pub dir: Option<PathBuf>,
pub namespace: Option<String>,
pub daemon: Option<String>,
}
impl SlugEntry {
pub fn resolve_dir(&self) -> Option<PathBuf> {
self.dir.clone().or_else(|| {
self.namespace.as_ref().and_then(|ns| {
let namespaces = PitchforkToml::read_global_namespaces();
namespaces.get(ns).map(|entry| entry.dir.clone())
})
})
}
pub fn resolve_namespace(&self) -> Option<String> {
self.namespace.clone().or_else(|| {
self.resolve_dir()
.and_then(|dir| PitchforkToml::namespace_for_dir(&dir).ok())
})
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, JsonSchema)]
pub struct GroupEntryRaw {
#[schemars(with = "Vec<DaemonId>")]
pub daemons: Vec<String>,
}
#[derive(Debug, Clone)]
pub struct GroupEntry {
pub daemons: Vec<DaemonId>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, JsonSchema)]
pub struct NamespaceEntryRaw {
pub dir: String,
}
#[derive(Debug, Clone)]
pub struct NamespaceEntry {
pub dir: PathBuf,
}
#[derive(Debug, Default, serde::Serialize, serde::Deserialize)]
struct PitchforkTomlRaw {
#[serde(skip_serializing_if = "Option::is_none", default)]
pub namespace: Option<String>,
#[serde(default)]
pub daemons: IndexMap<String, PitchforkTomlDaemonRaw>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub env: Option<IndexMap<String, String>>,
#[serde(default)]
pub settings: Option<SettingsPartial>,
#[serde(skip_serializing_if = "IndexMap::is_empty", default)]
pub slugs: IndexMap<String, SlugEntryRaw>,
#[serde(skip_serializing_if = "IndexMap::is_empty", default)]
pub groups: IndexMap<String, GroupEntryRaw>,
#[serde(skip_serializing_if = "IndexMap::is_empty", default)]
pub namespaces: IndexMap<String, NamespaceEntryRaw>,
}
#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize, schemars::JsonSchema)]
pub struct PitchforkTomlDaemonLogs {
#[serde(skip_serializing_if = "Option::is_none", default)]
pub log_format: Option<String>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub time_retention: Option<String>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub line_retention: Option<i64>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub archive_hook: Option<String>,
}
#[derive(Debug, serde::Serialize, serde::Deserialize)]
struct PitchforkTomlDaemonRaw {
pub run: String,
#[serde(skip_serializing_if = "Vec::is_empty", default)]
pub auto: Vec<PitchforkTomlAuto>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub cron: Option<PitchforkTomlCron>,
#[serde(default)]
pub retry: Retry,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub ready_delay: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub ready_output: Option<ReadyOutput>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub ready_http: Option<ReadyHttp>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub ready_port: Option<ReadyPort>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub ready_cmd: Option<ReadyCmd>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub port: Option<PortConfig>,
#[serde(skip_serializing_if = "Vec::is_empty", default)]
pub expected_port: Vec<u16>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub auto_bump_port: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub port_bump_attempts: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub boot_start: Option<bool>,
#[serde(skip_serializing_if = "Vec::is_empty", default)]
pub depends: Vec<String>,
#[serde(skip_serializing_if = "Vec::is_empty", default)]
pub watch: Vec<String>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub watch_mode: Option<WatchMode>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub dir: Option<String>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub env: Option<IndexMap<String, String>>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub hooks: Option<PitchforkTomlHooks>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub mise: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub user: Option<String>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub memory_limit: Option<MemoryLimit>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub cpu_limit: Option<CpuLimit>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub stop_signal: Option<StopConfig>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub pty: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub time_retention: Option<String>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub line_retention: Option<i64>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub archive_hook: Option<String>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub logs: Option<PitchforkTomlDaemonLogs>,
}
#[derive(Debug, Clone, Default, JsonSchema)]
#[schemars(title = "Pitchfork Configuration")]
pub struct PitchforkToml {
#[serde(default)]
pub daemons: IndexMap<DaemonId, PitchforkTomlDaemon>,
#[serde(skip_serializing_if = "Option::is_none", default)]
pub env: Option<IndexMap<String, String>>,
pub namespace: Option<String>,
#[serde(default)]
pub(crate) settings: SettingsPartial,
#[schemars(default, with = "IndexMap<String, SlugEntryRaw>")]
pub slugs: IndexMap<String, SlugEntry>,
#[schemars(default, with = "IndexMap<String, GroupEntryRaw>")]
pub groups: IndexMap<String, GroupEntry>,
#[schemars(default, with = "IndexMap<String, NamespaceEntryRaw>")]
pub namespaces: IndexMap<String, NamespaceEntry>,
#[schemars(skip)]
pub path: Option<PathBuf>,
}
pub(crate) fn is_global_config(path: &Path) -> bool {
path == *env::PITCHFORK_GLOBAL_CONFIG_USER || path == *env::PITCHFORK_GLOBAL_CONFIG_SYSTEM
}
fn is_local_config(path: &Path) -> bool {
path.file_name()
.map(|n| n == "pitchfork.local.toml")
.unwrap_or(false)
}
pub(crate) fn is_dot_config_pitchfork(path: &Path) -> bool {
path.ends_with(".config/pitchfork.toml") || path.ends_with(".config/pitchfork.local.toml")
}
fn sibling_base_config(path: &Path) -> Option<PathBuf> {
if !is_local_config(path) {
return None;
}
path.parent().map(|p| p.join("pitchfork.toml"))
}
fn parse_namespace_override_from_content(path: &Path, content: &str) -> Result<Option<String>> {
use toml::Value;
let doc: Value = toml::from_str(content)
.map_err(|e| ConfigParseError::from_toml_error(path, content.to_string(), e))?;
let Some(value) = doc.get("namespace") else {
return Ok(None);
};
match value {
Value::String(s) => Ok(Some(s.clone())),
_ => Err(ConfigParseError::InvalidNamespace {
path: path.to_path_buf(),
namespace: value.to_string(),
reason: "top-level 'namespace' must be a string".to_string(),
}
.into()),
}
}
fn read_namespace_override_from_file(path: &Path) -> Result<Option<String>> {
if !path.exists() {
return Ok(None);
}
let content = std::fs::read_to_string(path).map_err(|e| FileError::ReadError {
path: path.to_path_buf(),
source: e,
})?;
parse_namespace_override_from_content(path, &content)
}
fn validate_namespace(path: &Path, namespace: &str) -> Result<String> {
if let Err(e) = DaemonId::try_new(namespace, "probe") {
return Err(ConfigParseError::InvalidNamespace {
path: path.to_path_buf(),
namespace: namespace.to_string(),
reason: e.to_string(),
}
.into());
}
Ok(namespace.to_string())
}
fn derive_namespace_from_dir(path: &Path) -> Result<String> {
let dir_for_namespace = if is_dot_config_pitchfork(path) {
path.parent().and_then(|p| p.parent())
} else {
path.parent()
};
let raw_namespace = dir_for_namespace
.and_then(|p| p.file_name())
.and_then(|n| n.to_str())
.ok_or_else(|| miette::miette!("cannot derive namespace from path '{}'", path.display()))?
.to_string();
validate_namespace(path, &raw_namespace).map_err(|e| {
ConfigParseError::InvalidNamespace {
path: path.to_path_buf(),
namespace: raw_namespace,
reason: format!(
"{e}. Set a valid top-level namespace, e.g. namespace = \"my-project\""
),
}
.into()
})
}
fn namespace_from_path_with_override(path: &Path, explicit: Option<&str>) -> Result<String> {
if is_global_config(path) {
if let Some(ns) = explicit
&& ns != "global"
{
return Err(ConfigParseError::InvalidNamespace {
path: path.to_path_buf(),
namespace: ns.to_string(),
reason: "global config files must use namespace 'global'".to_string(),
}
.into());
}
return Ok("global".to_string());
}
if let Some(ns) = explicit {
return validate_namespace(path, ns);
}
derive_namespace_from_dir(path)
}
fn namespace_from_file(path: &Path) -> Result<String> {
let explicit = read_namespace_override_from_file(path)?;
let base_explicit = sibling_base_config(path)
.filter(|p| p.exists())
.map(|p| read_namespace_override_from_file(&p))
.transpose()?
.flatten();
if let (Some(local_ns), Some(base_ns)) = (explicit.as_deref(), base_explicit.as_deref())
&& local_ns != base_ns
{
return Err(ConfigParseError::InvalidNamespace {
path: path.to_path_buf(),
namespace: local_ns.to_string(),
reason: format!(
"namespace '{local_ns}' does not match sibling pitchfork.toml namespace '{base_ns}'"
),
}
.into());
}
let effective_explicit = explicit.as_deref().or(base_explicit.as_deref());
namespace_from_path_with_override(path, effective_explicit)
}
pub fn namespace_from_path(path: &Path) -> Result<String> {
namespace_from_file(path)
}
struct ConfigCacheEntry {
config: PitchforkToml,
source_meta: Vec<(PathBuf, Option<(SystemTime, u64)>)>,
}
static CONFIG_CACHE: Lazy<StdMutex<HashMap<PathBuf, ConfigCacheEntry>>> =
Lazy::new(|| StdMutex::new(HashMap::new()));
fn meta_matches(paths: &[PathBuf], snapshot: &[(PathBuf, Option<(SystemTime, u64)>)]) -> bool {
if paths.len() != snapshot.len() {
return false;
}
paths
.iter()
.zip(snapshot.iter())
.all(|(p, (snap_p, snap_meta))| p == snap_p && current_meta(p) == *snap_meta)
}
fn current_meta(path: &Path) -> Option<(SystemTime, u64)> {
let md = std::fs::metadata(path).ok()?;
Some((md.modified().ok()?, md.len()))
}
fn snapshot_meta(paths: &[PathBuf]) -> Vec<(PathBuf, Option<(SystemTime, u64)>)> {
paths.iter().map(|p| (p.clone(), current_meta(p))).collect()
}
pub fn invalidate_config_cache() {
if let Ok(mut cache) = CONFIG_CACHE.lock() {
cache.clear();
}
}
impl PitchforkToml {
pub fn resolve_daemon_id(&self, user_id: &str) -> Result<Vec<DaemonId>> {
if user_id.contains('/') {
return match DaemonId::parse(user_id) {
Ok(id) => Ok(vec![id]),
Err(e) => Err(e), };
}
let global_slugs = Self::read_global_slugs();
if let Some(entry) = global_slugs.get(user_id) {
let daemon_name = entry.daemon.as_deref().unwrap_or(user_id);
if let Some(dir) = entry.resolve_dir()
&& let Ok(project_config) = Self::all_merged_from(&dir)
{
let matches: Vec<DaemonId> = project_config
.daemons
.keys()
.filter(|id| id.name() == daemon_name)
.cloned()
.collect();
match matches.as_slice() {
[] => {}
[id] => return Ok(vec![id.clone()]),
_ => {
let mut candidates: Vec<String> =
matches.iter().map(|id| id.qualified()).collect();
candidates.sort();
return Err(miette::miette!(
"slug '{}' maps to daemon '{}' which matches multiple daemons: {}",
user_id,
daemon_name,
candidates.join(", ")
));
}
}
}
}
let matches: Vec<DaemonId> = self
.daemons
.keys()
.filter(|id| id.name() == user_id)
.cloned()
.collect();
if matches.is_empty() {
let state_matches = Self::find_in_state_file(user_id);
match state_matches.as_slice() {
[] => {}
[id] => return Ok(vec![id.clone()]),
_ => {
let mut candidates: Vec<String> =
state_matches.iter().map(|id| id.qualified()).collect();
candidates.sort();
return Err(miette::miette!(
"daemon '{}' is ambiguous; matches: {}. Use a qualified daemon ID (namespace/name)",
user_id,
candidates.join(", ")
));
}
}
let _ = DaemonId::try_new("global", user_id)?;
}
Ok(matches)
}
fn find_in_state_file(short_name: &str) -> Vec<DaemonId> {
match StateFile::read(&*env::PITCHFORK_STATE_FILE) {
Ok(state) => state
.daemons
.keys()
.filter(|id| id.name() == short_name)
.cloned()
.collect(),
Err(e) => {
warn!("cannot read state file: {e}");
Vec::new()
}
}
}
#[allow(dead_code)]
pub fn resolve_daemon_id_prefer_local(
&self,
user_id: &str,
current_dir: &Path,
) -> Result<DaemonId> {
if user_id.contains('/') {
return DaemonId::parse(user_id);
}
let current_namespace = Self::namespace_for_dir(current_dir)?;
self.resolve_daemon_id_with_namespace(user_id, ¤t_namespace)
}
fn resolve_daemon_id_with_namespace(
&self,
user_id: &str,
current_namespace: &str,
) -> Result<DaemonId> {
let global_slugs = Self::read_global_slugs();
if let Some(entry) = global_slugs.get(user_id) {
let daemon_name = entry.daemon.as_deref().unwrap_or(user_id);
if let Some(dir) = entry.resolve_dir()
&& let Ok(project_config) = Self::all_merged_from(&dir)
{
let matches: Vec<DaemonId> = project_config
.daemons
.keys()
.filter(|id| id.name() == daemon_name)
.cloned()
.collect();
match matches.as_slice() {
[] => {}
[id] => return Ok(id.clone()),
_ => {
let mut candidates: Vec<String> =
matches.iter().map(|id| id.qualified()).collect();
candidates.sort();
return Err(miette::miette!(
"slug '{}' maps to daemon '{}' which matches multiple daemons: {}",
user_id,
daemon_name,
candidates.join(", ")
));
}
}
}
}
let preferred_id = DaemonId::try_new(current_namespace, user_id)?;
if self.daemons.contains_key(&preferred_id) {
return Ok(preferred_id);
}
let matches = self.resolve_daemon_id(user_id)?;
if matches.len() > 1 {
let mut candidates: Vec<String> = matches.iter().map(|id| id.qualified()).collect();
candidates.sort();
return Err(miette::miette!(
"daemon '{}' is ambiguous; matches: {}. Use a qualified daemon ID (namespace/name)",
user_id,
candidates.join(", ")
));
}
if let Some(id) = matches.into_iter().next() {
return Ok(id);
}
let global_id = DaemonId::try_new("global", user_id)?;
if self.daemons.contains_key(&global_id) {
return Ok(global_id);
}
let suggestion = find_similar_daemon(user_id, self.daemons.keys().map(|id| id.name()));
Err(DependencyError::DaemonNotFound {
name: user_id.to_string(),
suggestion,
}
.into())
}
pub fn namespace_for_dir(dir: &Path) -> Result<String> {
Ok(Self::list_paths_from(dir)
.iter()
.rfind(|p| p.exists()) .map(|p| namespace_from_path(p))
.transpose()?
.unwrap_or_else(|| "global".to_string()))
}
pub fn resolve_id(user_id: &str) -> Result<DaemonId> {
if user_id.contains('/') {
return DaemonId::parse(user_id);
}
let config = Self::all_merged()?;
let ns = Self::namespace_for_dir(&env::CWD)?;
config.resolve_daemon_id_with_namespace(user_id, &ns)
}
pub fn resolve_id_allow_adhoc(user_id: &str) -> Result<DaemonId> {
if user_id.contains('/') {
return DaemonId::parse(user_id);
}
let config = Self::all_merged()?;
let ns = Self::namespace_for_dir(&env::CWD)?;
let preferred_id = DaemonId::try_new(&ns, user_id)?;
if config.daemons.contains_key(&preferred_id) {
return Ok(preferred_id);
}
let matches = config.resolve_daemon_id(user_id)?;
if matches.len() > 1 {
let mut candidates: Vec<String> = matches.iter().map(|id| id.qualified()).collect();
candidates.sort();
return Err(miette::miette!(
"daemon '{}' is ambiguous; matches: {}. Use a qualified daemon ID (namespace/name)",
user_id,
candidates.join(", ")
));
}
if let Some(id) = matches.into_iter().next() {
return Ok(id);
}
DaemonId::try_new("global", user_id)
}
pub fn resolve_ids<S: AsRef<str>>(user_ids: &[S]) -> Result<Vec<DaemonId>> {
if user_ids.iter().all(|s| s.as_ref().contains('/')) {
return user_ids
.iter()
.map(|s| DaemonId::parse(s.as_ref()))
.collect();
}
let config = Self::all_merged()?;
let ns = Self::namespace_for_dir(&env::CWD)?;
user_ids
.iter()
.map(|s| {
let id = s.as_ref();
if id.contains('/') {
DaemonId::parse(id)
} else {
config.resolve_daemon_id_with_namespace(id, &ns)
}
})
.collect()
}
pub fn resolve_ids_and_group<S: AsRef<str>>(
user_ids: &[S],
group_name: Option<&str>,
) -> Result<Vec<DaemonId>> {
let config = Self::all_merged()?;
let ns = Self::namespace_for_dir(&env::CWD)?;
let mut ids = Vec::new();
let mut seen = std::collections::HashSet::new();
for id in user_ids {
let id_str = id.as_ref();
let daemon_id = if id_str.contains('/') {
DaemonId::parse(id_str)?
} else {
config.resolve_daemon_id_with_namespace(id_str, &ns)?
};
if seen.insert(daemon_id.clone()) {
ids.push(daemon_id);
}
}
if let Some(name) = group_name {
match config.groups.get(name) {
Some(group) => {
let missing: Vec<String> = group
.daemons
.iter()
.filter(|id| !config.daemons.contains_key(*id))
.map(|id| id.qualified())
.collect();
if !missing.is_empty() {
return Err(miette::miette!(
"group '{}' references undefined daemon{}: {}",
name,
if missing.len() > 1 { "s" } else { "" },
missing.join(", ")
));
}
for daemon_id in &group.daemons {
if seen.insert(daemon_id.clone()) {
ids.push(daemon_id.clone());
}
}
}
None => {
let suggestion =
find_similar_daemon(name, config.groups.keys().map(|s| s.as_str()));
return Err(miette::miette!(
"group '{}' not found in configuration{}",
name,
suggestion.map(|s| format!(", {s}")).unwrap_or_default()
));
}
}
}
Ok(ids)
}
pub fn list_paths() -> Vec<PathBuf> {
Self::list_paths_from(&env::CWD)
}
pub fn list_paths_from(cwd: &Path) -> Vec<PathBuf> {
let mut paths = Vec::new();
paths.push(env::PITCHFORK_GLOBAL_CONFIG_SYSTEM.clone());
paths.push(env::PITCHFORK_GLOBAL_CONFIG_USER.clone());
let mut project_paths = xx::file::find_up_all(
cwd,
&[
"pitchfork.local.toml",
"pitchfork.toml",
".config/pitchfork.local.toml",
".config/pitchfork.toml",
],
);
project_paths.reverse();
paths.extend(project_paths);
paths
}
pub fn all_merged() -> Result<PitchforkToml> {
Self::all_merged_from(&env::CWD)
}
pub fn all_merged_all_namespaces() -> Result<Self> {
let mut pt = Self::all_merged_from(&env::CWD)?;
let namespaces = Self::read_global_namespaces();
for (ns_name, entry) in namespaces {
match Self::all_merged_from(&entry.dir) {
Ok(ns_config) => {
for (daemon_id, daemon_config) in ns_config.daemons {
if !pt.daemons.contains_key(&daemon_id) {
pt.daemons.insert(daemon_id, daemon_config);
}
}
pt.settings.merge_from(&ns_config.settings);
}
Err(e) => {
log::warn!(
"Failed to load namespace '{ns_name}' from {}: {e}",
entry.dir.display()
);
}
}
}
Ok(pt)
}
pub fn all_merged_from(cwd: &Path) -> Result<PitchforkToml> {
let paths = Self::list_paths_from(cwd);
let cache_key = cwd.canonicalize().unwrap_or_else(|_| cwd.to_path_buf());
{
let cache = CONFIG_CACHE.lock().unwrap_or_else(|e| e.into_inner());
if let Some(entry) = cache.get(&cache_key)
&& meta_matches(&paths, &entry.source_meta)
{
return Ok(entry.config.clone());
}
}
let snapshot = snapshot_meta(&paths);
let pt = Self::all_merged_from_uncached(&paths)?;
let mut cache = CONFIG_CACHE.lock().unwrap_or_else(|e| e.into_inner());
cache.insert(
cache_key,
ConfigCacheEntry {
config: pt.clone(),
source_meta: snapshot,
},
);
Ok(pt)
}
fn all_merged_from_uncached(paths: &[PathBuf]) -> Result<PitchforkToml> {
use std::collections::HashMap as StdHashMap;
let mut ns_to_origin: StdHashMap<String, (PathBuf, PathBuf)> = StdHashMap::new();
let mut pt = Self::default();
for p in paths {
match Self::read(p) {
Ok(pt2) => {
if p.exists() && !is_global_config(p) {
let ns = namespace_from_path(p)?;
let origin_dir = if is_dot_config_pitchfork(p) {
p.parent().and_then(|d| d.parent())
} else {
p.parent()
}
.map(|dir| dir.canonicalize().unwrap_or_else(|_| dir.to_path_buf()))
.unwrap_or_else(|| p.clone());
if let Some((other_path, other_dir)) = ns_to_origin.get(ns.as_str())
&& *other_dir != origin_dir
{
return Err(crate::error::ConfigParseError::NamespaceCollision {
path_a: other_path.clone(),
path_b: p.clone(),
ns,
}
.into());
}
ns_to_origin.insert(ns, (p.clone(), origin_dir));
}
pt.merge(pt2)
}
Err(e) => return Err(e.wrap_err(format!("error reading {}", p.display()))),
}
}
Ok(pt)
}
}
impl PitchforkToml {
pub fn new(path: PathBuf) -> Self {
Self {
daemons: Default::default(),
env: None,
namespace: None,
settings: SettingsPartial::default(),
slugs: IndexMap::new(),
groups: IndexMap::new(),
namespaces: IndexMap::new(),
path: Some(path),
}
}
pub fn parse_str(content: &str, path: &Path) -> Result<Self> {
let raw_config: PitchforkTomlRaw = toml::from_str(content)
.map_err(|e| ConfigParseError::from_toml_error(path, content.to_string(), e))?;
let namespace = {
let base_explicit = sibling_base_config(path)
.filter(|p| p.exists())
.map(|p| read_namespace_override_from_file(&p))
.transpose()?
.flatten();
if is_local_config(path)
&& let (Some(local_ns), Some(base_ns)) =
(raw_config.namespace.as_deref(), base_explicit.as_deref())
&& local_ns != base_ns
{
return Err(ConfigParseError::InvalidNamespace {
path: path.to_path_buf(),
namespace: local_ns.to_string(),
reason: format!(
"namespace '{local_ns}' does not match sibling pitchfork.toml namespace '{base_ns}'"
),
}
.into());
}
let explicit = raw_config.namespace.as_deref().or(base_explicit.as_deref());
namespace_from_path_with_override(path, explicit)?
};
let mut pt = Self::new(path.to_path_buf());
pt.namespace = raw_config.namespace.clone();
for (short_name, raw_daemon) in raw_config.daemons {
let id = match DaemonId::try_new(&namespace, &short_name) {
Ok(id) => id,
Err(e) => {
return Err(ConfigParseError::InvalidDaemonName {
name: short_name,
path: path.to_path_buf(),
reason: e.to_string(),
}
.into());
}
};
let mut depends = Vec::new();
for dep in raw_daemon.depends {
let dep_id = if dep.contains('/') {
match DaemonId::parse(&dep) {
Ok(id) => id,
Err(e) => {
return Err(ConfigParseError::InvalidDependency {
daemon: short_name.clone(),
dependency: dep,
path: path.to_path_buf(),
reason: e.to_string(),
}
.into());
}
}
} else {
match DaemonId::try_new(&namespace, &dep) {
Ok(id) => id,
Err(e) => {
return Err(ConfigParseError::InvalidDependency {
daemon: short_name.clone(),
dependency: dep,
path: path.to_path_buf(),
reason: e.to_string(),
}
.into());
}
}
};
depends.push(dep_id);
}
let has_deprecated = !raw_daemon.expected_port.is_empty()
|| raw_daemon.auto_bump_port.is_some()
|| raw_daemon.port_bump_attempts.is_some();
let port = if let Some(port) = raw_daemon.port {
if has_deprecated {
warn!(
"daemon {short_name}: both `port` and deprecated expected_port/auto_bump_port/port_bump_attempts are set; ignoring deprecated fields"
);
}
Some(port)
} else if has_deprecated {
warn!(
"daemon {short_name}: expected_port/auto_bump_port/port_bump_attempts are deprecated, use [daemons.{short_name}.port] instead"
);
let bump = if raw_daemon.auto_bump_port.unwrap_or(false) {
PortBump(
raw_daemon
.port_bump_attempts
.unwrap_or_else(|| settings().default_port_bump_attempts()),
)
} else {
PortBump(0)
};
Some(PortConfig {
expect: raw_daemon.expected_port,
bump,
})
} else {
None
};
let daemon = PitchforkTomlDaemon {
run: raw_daemon.run,
auto: raw_daemon.auto,
cron: raw_daemon.cron,
retry: raw_daemon.retry,
ready_delay: raw_daemon.ready_delay,
ready_output: raw_daemon.ready_output,
ready_http: raw_daemon.ready_http,
ready_port: raw_daemon.ready_port,
ready_cmd: raw_daemon.ready_cmd,
port,
boot_start: raw_daemon.boot_start,
depends,
watch: raw_daemon.watch,
watch_mode: raw_daemon.watch_mode.unwrap_or_default(),
dir: raw_daemon.dir,
env: raw_daemon.env,
hooks: raw_daemon.hooks,
mise: raw_daemon.mise,
user: raw_daemon.user,
memory_limit: raw_daemon.memory_limit,
cpu_limit: raw_daemon.cpu_limit,
stop_signal: raw_daemon.stop_signal,
pty: raw_daemon.pty,
time_retention: raw_daemon.time_retention,
line_retention: raw_daemon.line_retention,
archive_hook: raw_daemon.archive_hook,
logs: raw_daemon.logs,
path: Some(path.to_path_buf()),
};
pt.daemons.insert(id, daemon);
}
if let Some(settings) = raw_config.settings {
pt.settings = settings;
}
pt.env = raw_config.env;
for (slug, entry) in raw_config.slugs {
pt.slugs.insert(
slug,
SlugEntry {
dir: entry.dir.map(env::expand_tilde),
namespace: entry.namespace,
daemon: entry.daemon,
},
);
}
for (name, entry) in raw_config.namespaces {
pt.namespaces.insert(
name,
NamespaceEntry {
dir: env::expand_tilde(entry.dir),
},
);
}
for (group_name, raw_group) in raw_config.groups {
let mut daemons = Vec::new();
for daemon_name in &raw_group.daemons {
let id = if daemon_name.contains('/') {
DaemonId::parse(daemon_name).map_err(|e| {
ConfigParseError::InvalidDependency {
daemon: group_name.clone(),
dependency: daemon_name.clone(),
path: path.to_path_buf(),
reason: e.to_string(),
}
})?
} else {
DaemonId::try_new(&namespace, daemon_name).map_err(|e| {
ConfigParseError::InvalidDaemonName {
name: daemon_name.clone(),
path: path.to_path_buf(),
reason: e.to_string(),
}
})?
};
daemons.push(id);
}
pt.groups.insert(group_name, GroupEntry { daemons });
}
Ok(pt)
}
pub fn read<P: AsRef<Path>>(path: P) -> Result<Self> {
let path = path.as_ref();
if !path.exists() {
return Ok(Self::new(path.to_path_buf()));
}
let _lock = xx::fslock::get(path, false)
.wrap_err_with(|| format!("failed to acquire lock on {}", path.display()))?;
let raw = std::fs::read_to_string(path).map_err(|e| FileError::ReadError {
path: path.to_path_buf(),
source: e,
})?;
Self::parse_str(&raw, path)
}
pub fn write(&self) -> Result<()> {
if let Some(path) = &self.path {
let _lock = xx::fslock::get(path, false)
.wrap_err_with(|| format!("failed to acquire lock on {}", path.display()))?;
self.write_unlocked()
} else {
Err(FileError::NoPath.into())
}
}
fn write_unlocked(&self) -> Result<()> {
if let Some(path) = &self.path {
let config_namespace = if path.exists() {
namespace_from_path(path)?
} else {
namespace_from_path_with_override(path, self.namespace.as_deref())?
};
let mut raw = PitchforkTomlRaw {
namespace: self.namespace.clone(),
env: self.env.clone(),
settings: (!self.settings.is_empty()).then(|| self.settings.clone()),
..PitchforkTomlRaw::default()
};
for (id, daemon) in &self.daemons {
if id.namespace() != config_namespace {
return Err(miette::miette!(
"cannot write daemon '{}' to {}: daemon belongs to namespace '{}' but file namespace is '{}'",
id,
path.display(),
id.namespace(),
config_namespace
));
}
let port = daemon.port.as_ref();
let raw_daemon = PitchforkTomlDaemonRaw {
run: daemon.run.clone(),
auto: daemon.auto.clone(),
cron: daemon.cron.clone(),
retry: daemon.retry,
ready_delay: daemon.ready_delay,
ready_output: daemon.ready_output.clone(),
ready_http: daemon.ready_http.clone(),
ready_port: daemon.ready_port.clone(),
ready_cmd: daemon.ready_cmd.clone(),
port: port.cloned(),
expected_port: port.map(|p| p.expect.clone()).unwrap_or_default(),
auto_bump_port: port.filter(|p| p.auto_bump()).map(|_| true),
port_bump_attempts: port
.filter(|p| p.auto_bump())
.map(|p| p.max_bump_attempts()),
boot_start: daemon.boot_start,
depends: daemon
.depends
.iter()
.map(|d| {
if d.namespace() == config_namespace {
d.name().to_string()
} else {
d.qualified()
}
})
.collect(),
watch: daemon.watch.clone(),
watch_mode: match daemon.watch_mode {
WatchMode::Native => None,
mode => Some(mode),
},
dir: daemon.dir.clone(),
env: daemon.env.clone(),
hooks: daemon.hooks.clone(),
mise: daemon.mise,
user: daemon.user.clone(),
memory_limit: daemon.memory_limit,
cpu_limit: daemon.cpu_limit,
stop_signal: daemon.stop_signal,
pty: daemon.pty,
time_retention: daemon.time_retention.clone(),
line_retention: daemon.line_retention,
archive_hook: daemon.archive_hook.clone(),
logs: daemon.logs.clone(),
};
raw.daemons.insert(id.name().to_string(), raw_daemon);
}
for (slug, entry) in &self.slugs {
raw.slugs.insert(
slug.clone(),
SlugEntryRaw {
dir: entry.dir.as_ref().map(|d| d.to_string_lossy().to_string()),
namespace: entry.namespace.clone(),
daemon: entry.daemon.clone(),
},
);
}
for (name, group) in &self.groups {
let raw_daemons: Vec<String> = group
.daemons
.iter()
.map(|id| {
if id.namespace() == config_namespace {
id.name().to_string()
} else {
id.qualified()
}
})
.collect();
raw.groups.insert(
name.clone(),
GroupEntryRaw {
daemons: raw_daemons,
},
);
}
for (name, entry) in &self.namespaces {
raw.namespaces.insert(
name.clone(),
NamespaceEntryRaw {
dir: entry.dir.to_string_lossy().to_string(),
},
);
}
let raw_str = toml::to_string(&raw).map_err(|e| FileError::SerializeError {
path: path.clone(),
source: e,
})?;
xx::file::write(path, &raw_str).map_err(|e| FileError::WriteError {
path: path.clone(),
details: Some(e.to_string()),
})?;
invalidate_config_cache();
Ok(())
} else {
Err(FileError::NoPath.into())
}
}
pub fn merge(&mut self, pt: Self) {
for (id, d) in pt.daemons {
self.daemons.insert(id, d);
}
if let Some(env) = pt.env {
let merged = self.env.get_or_insert_with(IndexMap::new);
for (k, v) in env {
merged.insert(k, v);
}
}
for (slug, entry) in pt.slugs {
self.slugs.insert(slug, entry);
}
for (name, group) in pt.groups {
self.groups.insert(name, group);
}
for (name, entry) in pt.namespaces {
self.namespaces.insert(name, entry);
}
self.settings.merge_from(&pt.settings);
}
pub fn read_global_slugs() -> IndexMap<String, SlugEntry> {
match Self::read(&*env::PITCHFORK_GLOBAL_CONFIG_USER) {
Ok(pt) => pt.slugs,
Err(_) => IndexMap::new(),
}
}
pub fn find_slug_for_daemon_in_registry(
daemon_id: &DaemonId,
global_slugs: &IndexMap<String, SlugEntry>,
) -> Option<String> {
global_slugs
.iter()
.find(|(slug, entry)| {
let daemon_name = entry.daemon.as_deref().unwrap_or(slug);
if daemon_id.name() != daemon_name {
return false;
}
match entry.resolve_namespace() {
Some(namespace) => daemon_id.namespace() == namespace,
None => false,
}
})
.map(|(slug, _)| slug.clone())
}
#[allow(dead_code)]
pub fn is_slug_registered(slug: &str) -> bool {
Self::read_global_slugs().contains_key(slug)
}
pub fn add_slug_with_namespace(
slug: &str,
namespace: Option<&str>,
daemon: Option<&str>,
) -> Result<()> {
let global_path = &*env::PITCHFORK_GLOBAL_CONFIG_USER;
if let Some(parent) = global_path.parent() {
std::fs::create_dir_all(parent).map_err(|e| {
miette::miette!(
"Failed to create config directory {}: {e}",
parent.display()
)
})?;
}
let _lock = xx::fslock::get(global_path, false)
.wrap_err_with(|| format!("failed to acquire lock on {}", global_path.display()))?;
let mut pt = if global_path.exists() {
let raw = std::fs::read_to_string(global_path).map_err(|e| FileError::ReadError {
path: global_path.to_path_buf(),
source: e,
})?;
Self::parse_str(&raw, global_path)?
} else {
Self::new(global_path.to_path_buf())
};
if let Some(ns) = namespace
&& !pt.namespaces.contains_key(ns)
{
let dir = pt
.slugs
.get(slug)
.and_then(|e| e.resolve_dir())
.or_else(|| namespace.and_then(|_| env::CWD.as_path().canonicalize().ok()));
if let Some(ref d) = dir {
pt.namespaces
.insert(ns.to_string(), NamespaceEntry { dir: d.clone() });
}
}
pt.slugs.insert(
slug.to_string(),
SlugEntry {
dir: None,
namespace: namespace.map(str::to_string),
daemon: daemon.map(str::to_string),
},
);
pt.write_unlocked()?;
crate::proxy::hosts::sync_hosts_from_settings();
Ok(())
}
pub fn remove_slug(slug: &str) -> Result<bool> {
let global_path = &*env::PITCHFORK_GLOBAL_CONFIG_USER;
if !global_path.exists() {
return Ok(false);
}
let _lock = xx::fslock::get(global_path, false)
.wrap_err_with(|| format!("failed to acquire lock on {}", global_path.display()))?;
let raw = std::fs::read_to_string(global_path).map_err(|e| FileError::ReadError {
path: global_path.to_path_buf(),
source: e,
})?;
let mut pt = Self::parse_str(&raw, global_path)?;
let removed = pt.slugs.shift_remove(slug).is_some();
if removed {
pt.write_unlocked()?;
crate::proxy::hosts::sync_hosts_from_settings();
}
Ok(removed)
}
pub fn read_global_namespaces() -> IndexMap<String, NamespaceEntry> {
match Self::read(&*env::PITCHFORK_GLOBAL_CONFIG_USER) {
Ok(pt) => pt.namespaces,
Err(_) => IndexMap::new(),
}
}
pub fn register_namespace(name: &str, dir: &str) -> crate::Result<()> {
let global_path = &*crate::env::PITCHFORK_GLOBAL_CONFIG_USER;
if let Some(parent) = global_path.parent() {
std::fs::create_dir_all(parent).map_err(|e| {
miette::miette!(
"Failed to create config directory {}: {e}",
parent.display()
)
})?;
}
let _lock = xx::fslock::get(global_path, false)
.wrap_err_with(|| format!("failed to acquire lock on {}", global_path.display()))?;
let mut pt = if global_path.exists() {
let raw = std::fs::read_to_string(global_path).map_err(|e| {
crate::error::FileError::ReadError {
path: global_path.to_path_buf(),
source: e,
}
})?;
Self::parse_str(&raw, global_path)?
} else {
Self::new(global_path.to_path_buf())
};
pt.namespaces.insert(
name.to_string(),
NamespaceEntry {
dir: env::expand_tilde(dir),
},
);
pt.write_unlocked()?;
Ok(())
}
pub fn remove_namespace(name: &str) -> crate::Result<bool> {
let global_path = &*crate::env::PITCHFORK_GLOBAL_CONFIG_USER;
if !global_path.exists() {
return Ok(false);
}
let _lock = xx::fslock::get(global_path, false)
.wrap_err_with(|| format!("failed to acquire lock on {}", global_path.display()))?;
let raw = std::fs::read_to_string(global_path).map_err(|e| {
crate::error::FileError::ReadError {
path: global_path.to_path_buf(),
source: e,
}
})?;
let mut pt = Self::parse_str(&raw, global_path)?;
let removed = pt.namespaces.shift_remove(name).is_some();
if removed {
pt.write_unlocked()?;
}
Ok(removed)
}
}
#[derive(Debug, Clone, JsonSchema, Default)]
pub struct PitchforkTomlDaemon {
#[schemars(example = example_run_command())]
pub run: String,
#[schemars(default)]
pub auto: Vec<PitchforkTomlAuto>,
pub cron: Option<PitchforkTomlCron>,
#[schemars(default)]
pub retry: Retry,
pub ready_delay: Option<u64>,
pub ready_output: Option<ReadyOutput>,
pub ready_http: Option<ReadyHttp>,
pub ready_port: Option<ReadyPort>,
pub ready_cmd: Option<ReadyCmd>,
pub port: Option<PortConfig>,
pub boot_start: Option<bool>,
#[schemars(default)]
pub depends: Vec<DaemonId>,
#[schemars(default)]
pub watch: Vec<String>,
#[schemars(default)]
pub watch_mode: WatchMode,
pub dir: Option<String>,
pub env: Option<IndexMap<String, String>>,
pub hooks: Option<PitchforkTomlHooks>,
pub mise: Option<bool>,
pub user: Option<String>,
pub memory_limit: Option<MemoryLimit>,
pub cpu_limit: Option<CpuLimit>,
pub stop_signal: Option<StopConfig>,
pub pty: Option<bool>,
pub time_retention: Option<String>,
pub line_retention: Option<i64>,
pub archive_hook: Option<String>,
pub logs: Option<PitchforkTomlDaemonLogs>,
#[schemars(skip)]
pub path: Option<PathBuf>,
}
impl PitchforkTomlDaemon {
pub fn effective_user(&self) -> Option<String> {
let daemon_user = self
.user
.as_deref()
.map(str::trim)
.filter(|u| !u.is_empty());
daemon_user.map(str::to_owned).or_else(|| {
let s = crate::settings::settings();
let su = s.supervisor.user.trim();
(!su.is_empty()).then(|| su.to_owned())
})
}
pub fn to_run_options(
&self,
id: &crate::daemon_id::DaemonId,
cmd: Vec<String>,
) -> crate::daemon::RunOptions {
use crate::daemon::RunOptions;
let effective_user = self.effective_user();
let dir = crate::ipc::batch::resolve_daemon_dir(
self.dir.as_deref(),
self.path.as_deref(),
effective_user.as_deref(),
);
let slug = crate::pitchfork_toml::PitchforkToml::read_global_slugs()
.into_iter()
.find(|(slug, entry)| {
let daemon_name = entry.daemon.as_deref().unwrap_or(slug);
if daemon_name != id.name() {
return false;
}
match entry.resolve_namespace() {
Some(namespace) => namespace == id.namespace(),
None => false,
}
})
.map(|(slug, _)| slug);
RunOptions {
id: id.clone(),
cmd,
run: Some(self.run.clone()),
force: false,
shell_pid: None,
dir: Dir(dir),
autostop: self.auto.contains(&PitchforkTomlAuto::Stop),
cron_schedule: self.cron.as_ref().map(|c| c.schedule.clone()),
cron_retrigger: self.cron.as_ref().map(|c| c.retrigger),
cron_immediate: self.cron.as_ref().map(|c| c.immediate),
retry: self.retry,
retry_count: 0,
ready_delay: self.ready_delay,
ready_output: self.ready_output.clone(),
ready_http: self.ready_http.clone(),
ready_port: self.ready_port.clone(),
ready_cmd: self.ready_cmd.clone(),
port: self.port.clone(),
wait_ready: false,
depends: self.depends.clone(),
env: self.env.clone(),
watch: self.watch.clone(),
watch_mode: self.watch_mode,
watch_base_dir: Some(crate::ipc::batch::resolve_config_base_dir(
self.path.as_deref(),
)),
mise: self.mise,
slug,
proxy: None,
user: self.user.clone(),
memory_limit: self.memory_limit,
cpu_limit: self.cpu_limit,
stop_signal: self.stop_signal,
archive_hook: self
.logs
.as_ref()
.and_then(|l| l.archive_hook.clone())
.or_else(|| self.archive_hook.clone()),
log_format: self.logs.as_ref().and_then(|l| l.log_format.clone()),
on_output_hook: self.hooks.as_ref().and_then(|h| h.on_output.clone()),
pty: self.pty,
}
}
}
fn example_run_command() -> &'static str {
"exec node server.js"
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::Path;
#[test]
fn test_daemon_user_parses_and_flows_to_run_options() {
let pt = PitchforkToml::parse_str(
r#"
[daemons.api]
run = "node server.js"
user = "postgres"
"#,
Path::new("/tmp/my-project/pitchfork.toml"),
)
.unwrap();
let id = DaemonId::new("my-project", "api");
let daemon = pt.daemons.get(&id).unwrap();
assert_eq!(daemon.user.as_deref(), Some("postgres"));
let opts = daemon.to_run_options(&id, vec!["node".to_string(), "server.js".to_string()]);
assert_eq!(opts.user.as_deref(), Some("postgres"));
}
#[test]
fn test_daemon_user_write_roundtrip() {
let temp = tempfile::tempdir().unwrap();
let path = temp.path().join("pitchfork.toml");
let mut pt = PitchforkToml::new(path.clone());
pt.namespace = Some("test-project".to_string());
pt.daemons.insert(
DaemonId::new("test-project", "api"),
PitchforkTomlDaemon {
run: "node server.js".to_string(),
user: Some("postgres".to_string()),
..PitchforkTomlDaemon::default()
},
);
pt.write().unwrap();
let raw = std::fs::read_to_string(&path).unwrap();
assert!(raw.contains("user = \"postgres\""));
let parsed = PitchforkToml::read(&path).unwrap();
let daemon = parsed
.daemons
.get(&DaemonId::new("test-project", "api"))
.unwrap();
assert_eq!(daemon.user.as_deref(), Some("postgres"));
}
#[test]
fn test_registry_dirs_expand_tilde() {
let pt = PitchforkToml::parse_str(
r#"
[slugs.api]
dir = "~/projects/api"
[namespaces.web]
dir = "~/projects/web"
"#,
Path::new("/tmp/config.toml"),
)
.unwrap();
assert_eq!(
pt.slugs["api"].dir,
Some(crate::env::HOME_DIR.join("projects/api"))
);
assert_eq!(
pt.namespaces["web"].dir,
crate::env::HOME_DIR.join("projects/web")
);
}
#[test]
fn test_settings_write_roundtrip() {
let temp = tempfile::tempdir().unwrap();
let path = temp.path().join("pitchfork.toml");
let mut pt = PitchforkToml::new(path.clone());
pt.namespace = Some("test-project".to_string());
pt.settings.web.auto_start = Some(true);
pt.settings.general.log_level = Some("debug".to_string());
pt.write().unwrap();
let raw = std::fs::read_to_string(&path).unwrap();
assert!(
raw.contains("[settings.web]"),
"settings.web section should be written, got:\n{raw}"
);
assert!(raw.contains("auto_start = true"));
assert!(raw.contains("log_level = \"debug\""));
let parsed = PitchforkToml::read(&path).unwrap();
assert_eq!(parsed.settings.web.auto_start, Some(true));
assert_eq!(parsed.settings.general.log_level.as_deref(), Some("debug"));
}
#[test]
fn test_settings_preserved_on_unrelated_write() {
let temp = tempfile::tempdir().unwrap();
let path = temp.path().join("pitchfork.toml");
std::fs::write(&path, "[settings.web]\nauto_start = true\n").unwrap();
let mut pt = PitchforkToml::read(&path).unwrap();
pt.slugs.insert(
"api".to_string(),
SlugEntry {
dir: None,
namespace: Some("myproject".to_string()),
daemon: None,
},
);
pt.namespaces.insert(
"myproject".to_string(),
NamespaceEntry {
dir: PathBuf::from("/tmp/myproject"),
},
);
pt.write().unwrap();
let raw = std::fs::read_to_string(&path).unwrap();
assert!(
raw.contains("[settings.web]"),
"existing settings must be preserved, got:\n{raw}"
);
assert!(raw.contains("auto_start = true"));
assert!(raw.contains("[slugs.api]"));
let parsed = PitchforkToml::read(&path).unwrap();
assert_eq!(parsed.settings.web.auto_start, Some(true));
assert!(parsed.slugs.contains_key("api"));
}
#[test]
fn test_config_cache_hit_and_invalidation() {
let temp = tempfile::tempdir().unwrap();
let dir = temp.path();
let config_path = dir.join("pitchfork.toml");
std::fs::write(&config_path, "[daemons.api]\nrun = \"echo v1\"\n").unwrap();
super::invalidate_config_cache();
let pt1 = PitchforkToml::all_merged_from(dir).unwrap();
let daemon_id = DaemonId::new(namespace_from_path(&config_path).unwrap(), "api");
assert_eq!(pt1.daemons[&daemon_id].run, "echo v1");
let pt2 = PitchforkToml::all_merged_from(dir).unwrap();
assert_eq!(pt2.daemons[&daemon_id].run, "echo v1");
std::thread::sleep(std::time::Duration::from_millis(50));
std::fs::write(&config_path, "[daemons.api]\nrun = \"echo v2\"\n").unwrap();
let pt3 = PitchforkToml::all_merged_from(dir).unwrap();
assert_eq!(pt3.daemons[&daemon_id].run, "echo v2");
super::invalidate_config_cache();
let pt4 = PitchforkToml::all_merged_from(dir).unwrap();
assert_eq!(pt4.daemons[&daemon_id].run, "echo v2");
super::invalidate_config_cache();
}
#[test]
fn test_config_cache_invalidation_on_write() {
let temp = tempfile::tempdir().unwrap();
let dir = temp.path();
let config_path = dir.join("pitchfork.toml");
std::fs::write(&config_path, "[daemons.api]\nrun = \"echo v1\"\n").unwrap();
super::invalidate_config_cache();
let pt1 = PitchforkToml::all_merged_from(dir).unwrap();
let daemon_id = DaemonId::new(namespace_from_path(&config_path).unwrap(), "api");
assert_eq!(pt1.daemons[&daemon_id].run, "echo v1");
let mut pt = PitchforkToml::read(&config_path).unwrap();
pt.daemons.get_mut(&daemon_id).unwrap().run = "echo v3".to_string();
let _ = pt.write();
let pt2 = PitchforkToml::all_merged_from(dir).unwrap();
assert_eq!(pt2.daemons[&daemon_id].run, "echo v3");
super::invalidate_config_cache();
}
#[test]
fn test_config_cache_size_invalidation() {
let temp = tempfile::tempdir().unwrap();
let dir = temp.path();
let config_path = dir.join("pitchfork.toml");
std::fs::write(&config_path, "[daemons.api]\nrun = \"echo v1\"\n").unwrap();
super::invalidate_config_cache();
let pt1 = PitchforkToml::all_merged_from(dir).unwrap();
let daemon_id = DaemonId::new(namespace_from_path(&config_path).unwrap(), "api");
assert_eq!(pt1.daemons[&daemon_id].run, "echo v1");
let original_mtime = std::fs::metadata(&config_path).unwrap().modified().unwrap();
std::fs::write(&config_path, "[daemons.api]\nrun = \"echo different\"\n").unwrap();
let file = std::fs::OpenOptions::new()
.write(true)
.open(&config_path)
.unwrap();
let times = std::fs::FileTimes::new().set_modified(original_mtime);
file.set_times(times).unwrap();
let pt2 = PitchforkToml::all_merged_from(dir).unwrap();
assert_eq!(
pt2.daemons[&daemon_id].run, "echo different",
"cache should invalidate on size change even with identical mtime"
);
super::invalidate_config_cache();
}
}