use crate::Result;
use crate::cli::logs::{ReadyCheckType, create_ready_check_job, stream_startup_logs};
use crate::daemon::RunOptions;
use crate::daemon_id::DaemonId;
use crate::deps::{compute_reverse_stop_order, resolve_dependencies};
use crate::ipc::client::IpcClient;
use crate::pitchfork_toml::{
HealthCmd, HealthHttp, HealthPort, PitchforkToml, PitchforkTomlDaemon, ReadyCmd, ReadyHttp,
ReadyOutput, ReadyPort, project_dir_for_config,
};
use chrono::{DateTime, Local};
use indexmap::IndexMap;
use std::collections::{HashMap, HashSet};
use std::io::IsTerminal;
use std::path::{Path, PathBuf};
use std::sync::Arc;
#[derive(Debug, Clone)]
pub struct RunResult {
pub started: bool,
pub exit_code: Option<i32>,
pub start_time: DateTime<Local>,
pub resolved_ports: Vec<u16>,
pub error_message: Option<String>,
}
pub struct StartResult {
pub started: Vec<(DaemonId, DateTime<Local>, Vec<u16>)>,
pub any_failed: bool,
pub pending_job_updates: Vec<PendingJobUpdate>,
}
#[derive(Debug)]
pub struct StopResult {
pub any_failed: bool,
}
pub struct PendingJobUpdate {
pub job: Option<Arc<clx::progress::ProgressJob>>,
pub id: DaemonId,
pub run_result: std::result::Result<RunResult, miette::Report>,
}
pub struct SpawnTaskResult {
pub id: DaemonId,
pub job: Option<Arc<clx::progress::ProgressJob>>,
pub run_result: std::result::Result<RunResult, miette::Report>,
}
#[derive(Debug, Clone, Default)]
pub struct StartOptions {
pub force: bool,
pub shell_pid: Option<u32>,
pub delay: Option<u64>,
pub output: Option<String>,
pub http: Option<String>,
pub port: Option<u16>,
pub cmd: Option<String>,
pub health_cmd: Option<String>,
pub health_http: Option<String>,
pub health_port: Option<u16>,
pub expected_port: Option<Vec<u16>>,
pub auto_bump_port: Option<crate::config_types::PortBump>,
pub retry: Option<crate::config_types::Retry>,
pub quiet: bool,
}
pub async fn build_run_options(
id: &DaemonId,
daemon_config: &PitchforkTomlDaemon,
overrides: Option<&StartOptions>,
) -> std::result::Result<RunOptions, String> {
let cmd = shell_words::split(&daemon_config.run)
.map_err(|e| format!("Failed to parse command: {e}"))?;
let mut run_opts = daemon_config.to_run_options(id, cmd);
run_opts.wait_ready = true;
if let Some(opts) = overrides {
run_opts.shell_pid = opts.shell_pid;
run_opts.force = opts.force;
run_opts.ready_delay = opts.delay.or(run_opts.ready_delay);
run_opts.ready_output =
merge_ready_output_override(run_opts.ready_output, opts.output.clone());
run_opts.ready_http = merge_ready_http_override(run_opts.ready_http, opts.http.clone());
run_opts.ready_port = if let Some(port) = opts.port {
Some(ReadyPort {
port: Some(port),
template: None,
timeout: run_opts.ready_port.as_ref().and_then(|p| p.timeout),
})
} else {
run_opts.ready_port.clone()
};
run_opts.ready_cmd = merge_ready_cmd_override(run_opts.ready_cmd, opts.cmd.clone());
run_opts.health_cmd =
merge_health_cmd_override(run_opts.health_cmd, opts.health_cmd.clone());
run_opts.health_http =
merge_health_http_override(run_opts.health_http, opts.health_http.clone());
run_opts.health_port = merge_health_port_override(run_opts.health_port, opts.health_port);
if let Some(ref expected) = opts.expected_port {
run_opts.port.get_or_insert_with(Default::default).expect = expected.clone();
}
if let Some(bump) = opts.auto_bump_port {
run_opts.port.get_or_insert_with(Default::default).bump = bump;
}
}
if run_opts.mise.is_none() || should_inject_default_ready_delay(&run_opts) {
let project_dir = resolve_config_base_dir(daemon_config.path.as_deref());
let project_settings = tokio::task::spawn_blocking(move || {
crate::settings::Settings::load_from_dir(&project_dir)
})
.await
.map_err(|e| format!("Failed to load project settings: {e}"))?;
if run_opts.mise.is_none() {
run_opts.mise = Some(project_settings.general.mise);
}
if should_inject_default_ready_delay(&run_opts) {
run_opts.ready_delay = Some(project_settings.general_ready_delay_secs()?);
}
}
Ok(run_opts)
}
fn should_inject_default_ready_delay(opts: &RunOptions) -> bool {
opts.ready_delay.is_none()
&& opts.ready_output.is_none()
&& opts.ready_http.is_none()
&& opts.ready_port.is_none()
&& opts.ready_cmd.is_none()
}
pub(crate) fn render_daemon_config(
id: &DaemonId,
daemon_config: &mut PitchforkTomlDaemon,
pt: &PitchforkToml,
) -> Result<()> {
let resolved_daemons = crate::state_file::StateFile::read(&*crate::env::PITCHFORK_STATE_FILE)
.map(|sf| {
sf.daemons
.iter()
.filter_map(|(id, d)| {
if d.resolved_port.is_empty() {
None
} else {
Some((id.clone(), d.resolved_port.clone()))
}
})
.collect::<HashMap<DaemonId, Vec<u16>>>()
})
.unwrap_or_default();
let mut ctx =
crate::template::TemplateContext::new(id, daemon_config, &resolved_daemons, &pt.daemons);
crate::template::render_daemon_templates(daemon_config, &mut ctx, pt.env.as_ref())
.map_err(|e| miette::miette!("Template render error for daemon {id}: {e}"))
}
fn merge_ready_http_override(
configured: Option<ReadyHttp>,
override_url: Option<String>,
) -> Option<ReadyHttp> {
match (configured, override_url) {
(Some(mut ready_http), Some(url)) => {
ready_http.url = url;
Some(ready_http)
}
(None, Some(url)) => Some(ReadyHttp::new(url)),
(ready_http, None) => ready_http,
}
}
fn merge_ready_cmd_override(
configured: Option<ReadyCmd>,
override_cmd: Option<String>,
) -> Option<ReadyCmd> {
match (configured, override_cmd) {
(Some(mut ready_cmd), Some(cmd)) => {
ready_cmd.run = cmd;
Some(ready_cmd)
}
(None, Some(cmd)) => Some(ReadyCmd::new(cmd)),
(ready_cmd, None) => ready_cmd,
}
}
fn merge_health_cmd_override(
configured: Option<HealthCmd>,
override_cmd: Option<String>,
) -> Option<HealthCmd> {
match (configured, override_cmd) {
(Some(mut health_cmd), Some(cmd)) => {
health_cmd.run = cmd;
Some(health_cmd)
}
(None, Some(cmd)) => Some(HealthCmd::new(cmd)),
(health_cmd, None) => health_cmd,
}
}
fn merge_health_http_override(
configured: Option<HealthHttp>,
override_url: Option<String>,
) -> Option<HealthHttp> {
match (configured, override_url) {
(Some(mut health_http), Some(url)) => {
health_http.url = url;
Some(health_http)
}
(None, Some(url)) => Some(HealthHttp::new(url)),
(health_http, None) => health_http,
}
}
fn merge_health_port_override(
configured: Option<HealthPort>,
override_port: Option<u16>,
) -> Option<HealthPort> {
match (configured, override_port) {
(Some(mut health_port), Some(port)) => {
health_port.port = Some(port);
health_port.template = None;
Some(health_port)
}
(None, Some(port)) => Some(HealthPort::new(port)),
(health_port, None) => health_port,
}
}
fn merge_ready_output_override(
configured: Option<ReadyOutput>,
override_pattern: Option<String>,
) -> Option<ReadyOutput> {
match (configured, override_pattern) {
(Some(mut ready_output), Some(pattern)) => {
ready_output.pattern = pattern;
Some(ready_output)
}
(None, Some(pattern)) => Some(ReadyOutput::new(pattern)),
(ready_output, None) => ready_output,
}
}
fn ready_check_type(opts: &RunOptions) -> ReadyCheckType {
if let Some(ref output) = opts.ready_output {
ReadyCheckType::Output(output.pattern.clone())
} else if let Some(ref http) = opts.ready_http {
ReadyCheckType::Http(http.url.clone())
} else if let Some(port) = opts.ready_port.as_ref().and_then(|p| p.as_port()) {
ReadyCheckType::Port(port)
} else if let Some(ref cmd) = opts.ready_cmd {
ReadyCheckType::Cmd(cmd.run.clone())
} else if let Some(secs) = opts.ready_delay {
ReadyCheckType::Delay(secs)
} else {
ReadyCheckType::Default
}
}
pub fn update_job_with_result(
job: Option<&clx::progress::ProgressJob>,
id: &DaemonId,
result: &std::result::Result<RunResult, miette::Report>,
) {
use clx::progress::ProgressStatus;
let id_label = {
let is_tty = std::io::stderr().is_terminal();
let colors_enabled = is_tty && console::colors_enabled_stderr();
crate::cli::logs::colored_id_label(&id.qualified(), colors_enabled)
};
let show_ts = crate::settings::settings().general.startup_log_timestamps;
let prefix = if show_ts {
let now = chrono::Local::now();
format!("{}", crate::ui::style::edim(now.format("%H:%M:%S")))
} else {
"{{spinner()}}".to_string()
};
if let Some(job) = job {
match result {
Ok(run_result) if run_result.started => {
let body = if run_result.resolved_ports.is_empty() {
format!("{prefix} {id_label} started")
} else {
let port_str = run_result
.resolved_ports
.iter()
.map(ToString::to_string)
.collect::<Vec<_>>()
.join(", ");
let port_label = if run_result.resolved_ports.len() == 1 {
"port"
} else {
"ports"
};
format!(
"{prefix} {id_label} started on {port_label} {}",
crate::ui::style::ncyan(&port_str)
)
};
job.set_body(body);
job.set_status(ProgressStatus::Done);
}
Ok(run_result) => {
if run_result.exit_code.is_none() && !run_result.started {
job.remove();
return;
}
let exit_info = run_result
.exit_code
.map(|c| format!(" (exit code {c})"))
.unwrap_or_default();
let error_detail = run_result
.error_message
.as_ref()
.map(|msg| format!(": {msg}"))
.unwrap_or_default();
job.set_body(format!(
"{prefix} {id_label} failed{exit_info}{error_detail}"
));
job.set_status(ProgressStatus::Failed);
}
Err(e) => {
job.set_body(format!("{prefix} {id_label} failed: {e}"));
job.set_status(ProgressStatus::Failed);
}
}
} else if let Ok(run_result) = result {
if !run_result.started && run_result.exit_code.is_some() {
if let Some(ref msg) = run_result.error_message {
error!("{msg}");
}
if let Ok(lines) = crate::cli::logs::collect_startup_logs(id, run_result.start_time) {
crate::cli::logs::print_error_logs_block(&lines);
}
}
} else if let Err(e) = result {
error!("Failed to start daemon {id}: {e}");
}
}
impl IpcClient {
pub fn get_all_configured_daemons() -> Result<Vec<DaemonId>> {
Ok(PitchforkToml::all_merged()?
.daemons
.keys()
.cloned()
.collect())
}
pub fn get_local_configured_daemons() -> Result<Vec<DaemonId>> {
Self::get_configured_daemons_filtered(|id| id.namespace() != "global")
}
pub fn get_global_configured_daemons() -> Result<Vec<DaemonId>> {
Self::get_configured_daemons_filtered(|id| id.namespace() == "global")
}
fn get_configured_daemons_filtered<F>(predicate: F) -> Result<Vec<DaemonId>>
where
F: Fn(&DaemonId) -> bool,
{
Ok(PitchforkToml::all_merged()?
.daemons
.into_keys()
.filter(predicate)
.collect())
}
pub async fn get_running_daemons(&self) -> Result<Vec<DaemonId>> {
Ok(self
.active_daemons()
.await?
.iter()
.filter(|d| d.status.is_running() || d.status.is_waiting())
.map(|d| d.id.clone())
.collect())
}
pub async fn get_running_configured_daemons(&self, global: bool) -> Result<Vec<DaemonId>> {
let configured: HashSet<DaemonId> = if global {
Self::get_global_configured_daemons()?
} else {
Self::get_local_configured_daemons()?
}
.into_iter()
.collect();
Ok(self
.get_running_daemons()
.await?
.into_iter()
.filter(|id| configured.contains(id))
.collect())
}
pub async fn start_daemons(
self: &Arc<Self>,
ids: &[DaemonId],
opts: StartOptions,
) -> Result<StartResult> {
let pt = PitchforkToml::all_merged_all_namespaces()?;
let disabled_daemons = self.get_disabled_daemons().await?;
let all_daemons = self.active_daemons().await?;
let adhoc_daemons: HashMap<DaemonId, crate::daemon::Daemon> = all_daemons
.into_iter()
.filter(|d| !pt.daemons.contains_key(&d.id))
.map(|d| (d.id.clone(), d))
.collect();
let requested_ids: Vec<DaemonId> = ids
.iter()
.filter(|id| {
if disabled_daemons.contains(id) {
warn!("Daemon {id} is disabled");
false
} else {
true
}
})
.cloned()
.collect();
if requested_ids.is_empty() {
return Ok(StartResult {
started: vec![],
any_failed: false,
pending_job_updates: vec![],
});
}
let (config_ids, adhoc_ids): (Vec<DaemonId>, Vec<DaemonId>) = requested_ids
.into_iter()
.partition(|id| pt.daemons.contains_key(id));
let active_daemons = self.active_daemons().await?;
let running_daemons: HashSet<DaemonId> = active_daemons
.iter()
.filter(|d| d.status.is_running() || d.status.is_waiting())
.map(|d| d.id.clone())
.collect();
let running_ports_map: HashMap<DaemonId, Vec<u16>> = active_daemons
.into_iter()
.filter(|d| {
(d.status.is_running() || d.status.is_waiting()) && !d.resolved_port.is_empty()
})
.map(|d| (d.id, d.resolved_port))
.collect();
let explicitly_requested: HashSet<DaemonId> = ids.iter().cloned().collect();
let mut any_failed = false;
let mut successful_daemons: Vec<(DaemonId, DateTime<Local>, Vec<u16>)> = Vec::new();
let mut resolved_ports_map: std::collections::HashMap<DaemonId, Vec<u16>> =
std::collections::HashMap::new();
let mut pending_job_updates: Vec<PendingJobUpdate> = Vec::new();
if !config_ids.is_empty() {
let dep_order = resolve_dependencies(&config_ids, &pt.daemons)?;
for (level_idx, level) in dep_order.levels.iter().enumerate() {
let is_last_level = level_idx == dep_order.levels.len() - 1;
let mut successful_this_level: Vec<(DaemonId, Vec<u16>)> = Vec::new();
let to_start: Vec<DaemonId> = level
.iter()
.filter(|&id| {
if disabled_daemons.contains(id) {
warn!("Skipping disabled daemon {id} (dependency)");
return false;
}
if running_daemons.contains(id) {
if opts.force && explicitly_requested.contains(id) {
debug!("Force restarting explicitly requested daemon: {id}");
true } else {
if explicitly_requested.contains(id) {
info!("Daemon {id} is already running, use --force to restart");
} else {
debug!("Skipping already running daemon {id}");
}
false
}
} else {
true
}
})
.cloned()
.collect();
for id in level {
if let Some(ports) = running_ports_map.get(id) {
resolved_ports_map.insert(id.clone(), ports.clone());
}
}
if to_start.is_empty() {
continue;
}
let mut tasks = Vec::new();
for id in to_start {
if let Some(daemon_config) = pt.daemons.get(&id) {
let mut rendered_config = daemon_config.clone();
let mut template_ctx = crate::template::TemplateContext::new(
&id,
daemon_config,
&resolved_ports_map,
&pt.daemons,
);
match crate::template::render_daemon_templates(
&mut rendered_config,
&mut template_ctx,
pt.env.as_ref(),
) {
Ok(()) => {}
Err(e) => {
error!("Template render error for daemon {id}: {e}");
any_failed = true;
continue;
}
}
let is_explicit = explicitly_requested.contains(&id);
let task = Self::spawn_start_task(id, &rendered_config, is_explicit, &opts);
tasks.push(task);
}
}
for task in tasks {
match task.await {
Ok(result) => {
let SpawnTaskResult {
id,
job,
run_result,
..
} = result;
match &run_result {
Ok(rr) if rr.started => {
successful_this_level
.push((id.clone(), rr.resolved_ports.clone()));
successful_daemons.push((
id.clone(),
rr.start_time,
rr.resolved_ports.clone(),
));
}
Ok(rr) => {
if rr.exit_code.is_some() {
any_failed = true;
error!("Daemon {} failed to start", id);
}
}
Err(_) => {
any_failed = true;
}
}
pending_job_updates.push(PendingJobUpdate {
job: job.clone(),
id: id.clone(),
run_result,
});
}
Err(e) => {
error!("Task panicked: {e}");
any_failed = true;
}
}
}
for (id, ports) in successful_this_level {
resolved_ports_map.insert(id, ports);
}
if any_failed {
if !is_last_level {
error!("Dependency failed, aborting remaining starts");
}
break;
}
}
}
if !any_failed && !adhoc_ids.is_empty() {
let mut tasks = Vec::new();
for id in adhoc_ids {
if running_daemons.contains(&id) {
if opts.force && explicitly_requested.contains(&id) {
debug!("Force restarting ad-hoc daemon: {id}");
} else {
if explicitly_requested.contains(&id) {
info!("Ad-hoc daemon {id} is already running, use --force to restart");
}
continue;
}
}
if let Some(adhoc_daemon) = adhoc_daemons.get(&id) {
if let Some(ref cmd) = adhoc_daemon.cmd {
let is_explicit = explicitly_requested.contains(&id);
let task = Self::spawn_adhoc_start_task(
id,
cmd.clone(),
adhoc_daemon.dir.clone().unwrap_or_default(),
adhoc_daemon.env.clone(),
adhoc_daemon.ready_http.clone(),
adhoc_daemon.health_cmd.clone(),
adhoc_daemon.health_http.clone(),
adhoc_daemon.health_port.clone(),
is_explicit,
&opts,
);
tasks.push(task);
} else {
warn!("Ad-hoc daemon {id} has no saved command, cannot restart");
}
} else {
warn!("Daemon {id} not found in config or state");
}
}
for task in tasks {
match task.await {
Ok(result) => {
let SpawnTaskResult {
id,
job,
run_result,
..
} = result;
match &run_result {
Ok(rr) if rr.started => {
successful_daemons.push((
id.clone(),
rr.start_time,
rr.resolved_ports.clone(),
));
}
Ok(rr) => {
if rr.exit_code.is_some() {
any_failed = true;
}
}
Err(_) => {
any_failed = true;
}
}
pending_job_updates.push(PendingJobUpdate {
job: job.clone(),
id: id.clone(),
run_result,
});
}
Err(e) => {
error!("Task panicked: {e}");
any_failed = true;
}
}
}
}
Ok(StartResult {
started: successful_daemons,
any_failed,
pending_job_updates,
})
}
async fn connect_dedicated() -> Result<Self> {
Self::connect(false).await
}
fn spawn_start_task(
id: DaemonId,
daemon_config: &PitchforkTomlDaemon,
is_explicitly_requested: bool,
opts: &StartOptions,
) -> tokio::task::JoinHandle<SpawnTaskResult> {
let mut start_opts = opts.clone();
start_opts.force = opts.force && is_explicitly_requested;
let daemon_config = daemon_config.clone();
let quiet = opts.quiet;
tokio::spawn(async move {
let run_opts = match build_run_options(&id, &daemon_config, Some(&start_opts)).await {
Ok(opts) => opts,
Err(e) => {
return SpawnTaskResult {
id,
job: None,
run_result: Err(miette::miette!("Failed to parse command: {e}")),
};
}
};
let check_type = ready_check_type(&run_opts);
let job = if !quiet {
let job = create_ready_check_job(&id, &check_type);
Some(job)
} else {
None
};
let (log_stop_tx, log_handle) = if let Some(ref job) = job {
let (tx, handle) = stream_startup_logs(&id, job.clone());
(Some(tx), Some(handle))
} else {
(None, None)
};
let result = match Self::connect_dedicated().await {
Ok(ipc) => ipc.run(run_opts).await,
Err(e) => Err(e),
};
if let Some(tx) = &log_stop_tx {
let _ = tx.send(true);
}
if let Some(handle) = log_handle {
let _ = handle.await;
}
SpawnTaskResult {
id,
job,
run_result: result,
}
})
}
#[allow(clippy::too_many_arguments)]
fn spawn_adhoc_start_task(
id: DaemonId,
cmd: Vec<String>,
dir: PathBuf,
env: Option<IndexMap<String, String>>,
ready_http: Option<ReadyHttp>,
health_cmd: Option<HealthCmd>,
health_http: Option<HealthHttp>,
health_port: Option<HealthPort>,
is_explicitly_requested: bool,
opts: &StartOptions,
) -> tokio::task::JoinHandle<SpawnTaskResult> {
let force = opts.force && is_explicitly_requested;
let delay = opts.delay;
let output = opts.output.clone();
let http = merge_ready_http_override(ready_http, opts.http.clone());
let port = opts.port;
let ready_cmd = opts.cmd.clone().map(ReadyCmd::new);
let health_cmd = merge_health_cmd_override(health_cmd, opts.health_cmd.clone());
let health_http = merge_health_http_override(health_http, opts.health_http.clone());
let health_port = merge_health_port_override(health_port, opts.health_port);
let expected_port = opts.expected_port.clone();
let auto_bump_port = opts.auto_bump_port;
let retry = opts.retry.unwrap_or_default();
let shell_pid = opts.shell_pid;
let quiet = opts.quiet;
tokio::spawn(async move {
let mut run_opts = RunOptions {
id: id.clone(),
cmd,
force,
shell_pid,
dir: crate::config_types::Dir(dir),
retry,
ready_delay: delay,
ready_output: output.map(ReadyOutput::new),
ready_http: http,
ready_port: port.map(ReadyPort::new),
ready_cmd,
health_cmd,
health_http,
health_port,
port: crate::config_types::PortConfig::from_parts(
expected_port.unwrap_or_default(),
auto_bump_port.unwrap_or_default(),
),
wait_ready: true,
env,
watch: vec![],
watch_base_dir: None,
mise: None,
slug: None,
proxy: None,
..RunOptions::default()
};
if should_inject_default_ready_delay(&run_opts) {
run_opts.ready_delay = Some(
match crate::settings::settings().general_ready_delay_secs() {
Ok(delay) => delay,
Err(e) => {
return SpawnTaskResult {
id,
job: None,
run_result: Err(miette::miette!("{e}")),
};
}
},
);
}
let check_type = ready_check_type(&run_opts);
let job = if !quiet {
let job = create_ready_check_job(&id, &check_type);
Some(job)
} else {
None
};
let (log_stop_tx, log_handle) = if let Some(ref job) = job {
let (tx, handle) = stream_startup_logs(&id, job.clone());
(Some(tx), Some(handle))
} else {
(None, None)
};
let result = match Self::connect_dedicated().await {
Ok(ipc) => ipc.run(run_opts).await,
Err(e) => Err(e),
};
if let Some(tx) = &log_stop_tx {
let _ = tx.send(true);
}
if let Some(handle) = log_handle {
let _ = handle.await;
}
SpawnTaskResult {
id,
job,
run_result: result,
}
})
}
fn spawn_stop_task(id: DaemonId) -> tokio::task::JoinHandle<(DaemonId, Result<()>)> {
tokio::spawn(async move {
let result = match Self::connect_dedicated().await {
Ok(ipc) => ipc.stop(id.clone()).await.map(|_| ()),
Err(e) => Err(e),
};
(id, result)
})
}
pub async fn start_daemon(
&self,
id: &DaemonId,
overrides: Option<&StartOptions>,
) -> Result<RunResult> {
let pt = PitchforkToml::all_merged_all_namespaces()?;
let mut daemon_config = pt
.daemons
.get(id)
.cloned()
.ok_or_else(|| miette::miette!("Daemon config not found for {id}"))?;
render_daemon_config(id, &mut daemon_config, &pt)?;
let run_opts = build_run_options(id, &daemon_config, overrides)
.await
.map_err(|e| miette::miette!("{e}"))?;
self.run(run_opts).await
}
pub async fn restart_daemon(
&self,
id: &DaemonId,
overrides: Option<&StartOptions>,
) -> Result<RunResult> {
let _ = self.stop(id.clone()).await;
self.start_daemon(id, overrides).await
}
pub async fn stop_daemons(self: &Arc<Self>, ids: &[DaemonId]) -> Result<StopResult> {
let running_daemons: HashSet<DaemonId> = self
.active_daemons()
.await?
.iter()
.filter(|d| d.status.is_running() || d.status.is_waiting())
.map(|d| d.id.clone())
.collect();
let requested_ids: Vec<DaemonId> = ids
.iter()
.filter(|id| {
if !running_daemons.contains(*id) {
warn!("Daemon {id} is not running");
false
} else {
true
}
})
.cloned()
.collect();
if requested_ids.is_empty() {
info!("No running daemons to stop");
return Ok(StopResult { any_failed: false });
}
let mut any_failed = false;
let stop_levels = compute_reverse_stop_order(&requested_ids);
for level in stop_levels {
let to_stop: Vec<DaemonId> = level
.into_iter()
.filter(|id| running_daemons.contains(id))
.collect();
if to_stop.is_empty() {
continue;
}
let mut tasks = Vec::new();
for id in to_stop {
let task = Self::spawn_stop_task(id);
tasks.push(task);
}
for task in tasks {
match task.await {
Ok((id, Ok(()))) => {
debug!("Successfully stopped daemon {id}");
}
Ok((id, Err(e))) => {
error!("Failed to stop daemon {id}: {e}");
any_failed = true;
}
Err(e) => {
error!("Stop task panicked: {e}");
any_failed = true;
}
}
}
}
Ok(StopResult { any_failed })
}
pub async fn run_adhoc(
&self,
id: DaemonId,
cmd: Vec<String>,
dir: PathBuf,
opts: StartOptions,
) -> Result<RunResult> {
let mut run_opts = RunOptions {
id,
cmd,
shell_pid: opts.shell_pid,
force: opts.force,
dir: crate::config_types::Dir(dir),
retry: opts.retry.unwrap_or_default(),
ready_delay: opts.delay,
ready_output: opts.output.map(ReadyOutput::new),
ready_http: merge_ready_http_override(None, opts.http),
ready_port: opts.port.map(ReadyPort::new),
ready_cmd: opts.cmd.clone().map(ReadyCmd::new),
health_cmd: opts.health_cmd.clone().map(HealthCmd::new),
health_http: opts.health_http.clone().map(HealthHttp::new),
health_port: opts.health_port.map(HealthPort::new),
port: crate::config_types::PortConfig::from_parts(
opts.expected_port.unwrap_or_default(),
opts.auto_bump_port.unwrap_or_default(),
),
wait_ready: true,
mise: None,
slug: None,
proxy: None,
..RunOptions::default()
};
if should_inject_default_ready_delay(&run_opts) {
run_opts.ready_delay = Some(
crate::settings::settings()
.general_ready_delay_secs()
.map_err(|e| miette::miette!("{e}"))?,
);
}
self.run(run_opts).await
}
}
pub fn resolve_config_base_dir(config_path: Option<&Path>) -> PathBuf {
config_path
.and_then(project_dir_for_config)
.unwrap_or_else(|| crate::env::CWD.to_path_buf())
}
pub fn resolve_daemon_dir(
dir: Option<&str>,
config_path: Option<&Path>,
user: Option<&str>,
) -> PathBuf {
let base_dir = resolve_config_base_dir(config_path);
match dir {
Some(d) => base_dir.join(crate::env::expand_tilde_for_user(d, user)),
None => base_dir,
}
}
#[cfg(test)]
mod tests {
use std::time::Duration;
use crate::env;
use super::*;
#[test]
fn http_override_preserves_configured_status_codes() {
let configured = Some(ReadyHttp {
url: "http://localhost:3000/original".to_string(),
status: vec![401],
timeout: None,
});
let ready_http =
merge_ready_http_override(configured, Some("http://localhost:3000/health".to_string()))
.unwrap();
assert_eq!(ready_http.url, "http://localhost:3000/health");
assert_eq!(ready_http.status, vec![401]);
assert!(ready_http.accepts_status(401));
assert!(!ready_http.accepts_status(200));
}
#[tokio::test]
async fn build_run_options_preserves_ready_http_status_for_cli_http_override() {
let id = DaemonId::try_new("project", "api").unwrap();
let daemon_config = PitchforkTomlDaemon {
run: "echo ready".to_string(),
ready_http: Some(ReadyHttp {
url: "http://localhost:3000/original".to_string(),
status: vec![401],
timeout: None,
}),
..PitchforkTomlDaemon::default()
};
let opts = StartOptions {
http: Some("http://localhost:3000/health".to_string()),
..StartOptions::default()
};
let run_opts = build_run_options(&id, &daemon_config, Some(&opts))
.await
.unwrap();
let ready_http = run_opts.ready_http.unwrap();
assert_eq!(ready_http.url, "http://localhost:3000/health");
assert_eq!(ready_http.status, vec![401]);
}
#[tokio::test]
async fn build_run_options_resolves_global_mise_for_supervisor() {
let id = DaemonId::try_new("project", "api").unwrap();
let daemon_config = PitchforkTomlDaemon {
run: "echo ready".to_string(),
..PitchforkTomlDaemon::default()
};
let run_opts = build_run_options(&id, &daemon_config, None).await.unwrap();
assert_eq!(
run_opts.mise,
Some(crate::settings::settings().general.mise)
);
}
#[tokio::test]
async fn build_run_options_resolves_mise_from_daemon_project() {
run_project_mise_test_in_sanitized_child("project").await;
}
#[tokio::test]
async fn build_run_options_resolves_mise_from_daemon_dot_config() {
run_project_mise_test_in_sanitized_child("dot-config").await;
}
async fn run_project_mise_test_in_sanitized_child(mode: &str) {
const CHILD_SENTINEL: &str = "project-mise-sanitized-child-ran";
let output = tokio::process::Command::new(std::env::current_exe().unwrap())
.args([
"--exact",
"ipc::batch::tests::build_run_options_resolves_mise_in_sanitized_child",
"--nocapture",
])
.env("PITCHFORK_TEST_PROJECT_MISE_MODE", mode)
.env_remove("PITCHFORK_MISE")
.output()
.await
.unwrap();
assert!(
output.status.success(),
"sanitized child failed:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
assert!(
String::from_utf8_lossy(&output.stderr).contains(CHILD_SENTINEL),
"sanitized child test did not run:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
}
#[tokio::test]
async fn build_run_options_resolves_mise_in_sanitized_child() {
let Ok(mode) = std::env::var("PITCHFORK_TEST_PROJECT_MISE_MODE") else {
return;
};
eprintln!("project-mise-sanitized-child-ran");
let project = tempfile::tempdir().unwrap();
let config_dir = match mode.as_str() {
"project" => project.path().to_path_buf(),
"dot-config" => {
let config_dir = project.path().join(".config");
tokio::fs::create_dir(&config_dir).await.unwrap();
config_dir
}
_ => panic!("unknown project mise test mode: {mode}"),
};
let config_path = config_dir.join("pitchfork.toml");
let project_mise = !crate::settings::settings().general.mise;
tokio::fs::write(
&config_path,
format!("[settings.general]\nmise = {project_mise}\n"),
)
.await
.unwrap();
let id = DaemonId::try_new("other-project", "api").unwrap();
let daemon_config = PitchforkTomlDaemon {
run: "echo ready".to_string(),
path: Some(config_path),
..PitchforkTomlDaemon::default()
};
let run_opts = build_run_options(&id, &daemon_config, None).await.unwrap();
assert_eq!(run_opts.mise, Some(project_mise));
}
#[tokio::test]
async fn build_run_options_preserves_daemon_mise_override() {
let id = DaemonId::try_new("project", "api").unwrap();
let daemon_config = PitchforkTomlDaemon {
run: "echo ready".to_string(),
mise: Some(!crate::settings::settings().general.mise),
..PitchforkTomlDaemon::default()
};
let run_opts = build_run_options(&id, &daemon_config, None).await.unwrap();
assert_eq!(run_opts.mise, daemon_config.mise);
}
#[tokio::test]
async fn build_run_options_falls_back_to_global_ready_delay() {
let id = DaemonId::try_new("project", "api").unwrap();
let daemon_config = PitchforkTomlDaemon {
run: "echo ready".to_string(),
..PitchforkTomlDaemon::default()
};
let run_opts = build_run_options(&id, &daemon_config, None).await.unwrap();
assert_eq!(
run_opts.ready_delay,
Some(crate::settings::settings().general_ready_delay().as_secs())
);
}
#[tokio::test]
async fn build_run_options_preserves_daemon_ready_delay_override() {
let id = DaemonId::try_new("project", "api").unwrap();
let daemon_config = PitchforkTomlDaemon {
run: "echo ready".to_string(),
ready_delay: Some(5),
..PitchforkTomlDaemon::default()
};
let run_opts = build_run_options(&id, &daemon_config, None).await.unwrap();
assert_eq!(run_opts.ready_delay, Some(5));
}
#[tokio::test]
async fn build_run_options_skips_default_delay_for_configured_ready_checks() {
let id = DaemonId::try_new("project", "api").unwrap();
let configured_checks = [
PitchforkTomlDaemon {
run: "echo ready".to_string(),
ready_output: Some(ReadyOutput::new("ready")),
..PitchforkTomlDaemon::default()
},
PitchforkTomlDaemon {
run: "echo ready".to_string(),
ready_http: Some(ReadyHttp::new("http://localhost/health")),
..PitchforkTomlDaemon::default()
},
PitchforkTomlDaemon {
run: "echo ready".to_string(),
ready_port: Some(ReadyPort::new(3000)),
..PitchforkTomlDaemon::default()
},
PitchforkTomlDaemon {
run: "echo ready".to_string(),
ready_cmd: Some(ReadyCmd::new("true")),
..PitchforkTomlDaemon::default()
},
];
for daemon_config in configured_checks {
let run_opts = build_run_options(&id, &daemon_config, None).await.unwrap();
assert_eq!(run_opts.ready_delay, None);
}
}
#[tokio::test]
async fn build_run_options_skips_default_delay_for_cli_ready_check_overrides() {
let id = DaemonId::try_new("project", "api").unwrap();
let daemon_config = PitchforkTomlDaemon {
run: "echo ready".to_string(),
..PitchforkTomlDaemon::default()
};
let cli_overrides = [
StartOptions {
output: Some("ready".to_string()),
..StartOptions::default()
},
StartOptions {
http: Some("http://localhost/health".to_string()),
..StartOptions::default()
},
StartOptions {
port: Some(3000),
..StartOptions::default()
},
StartOptions {
cmd: Some("true".to_string()),
..StartOptions::default()
},
];
for overrides in &cli_overrides {
let run_opts = build_run_options(&id, &daemon_config, Some(overrides))
.await
.unwrap();
assert_eq!(run_opts.ready_delay, None);
}
}
#[tokio::test]
async fn build_run_options_resolves_ready_delay_from_daemon_project() {
const CHILD_SENTINEL: &str = "project-ready-delay-sanitized-child-ran";
let output = tokio::process::Command::new(std::env::current_exe().unwrap())
.args([
"--exact",
"ipc::batch::tests::build_run_options_resolves_ready_delay_in_sanitized_child",
"--nocapture",
])
.env(
"PITCHFORK_TEST_PROJECT_READY_DELAY_MODE",
"ready-delay-project",
)
.env_remove("PITCHFORK_READY_DELAY")
.output()
.await
.unwrap();
assert!(
output.status.success(),
"sanitized child failed:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
assert!(
String::from_utf8_lossy(&output.stderr).contains(CHILD_SENTINEL),
"sanitized child test did not run:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
}
#[tokio::test]
async fn build_run_options_resolves_ready_delay_in_sanitized_child() {
let Ok(_mode) = std::env::var("PITCHFORK_TEST_PROJECT_READY_DELAY_MODE") else {
return;
};
eprintln!("project-ready-delay-sanitized-child-ran");
let project = tempfile::tempdir().unwrap();
let config_path = project.path().join("pitchfork.toml");
tokio::fs::write(&config_path, "[settings.general]\nready_delay = \"7s\"\n")
.await
.unwrap();
let id = DaemonId::try_new("other-project", "api").unwrap();
let daemon_config = PitchforkTomlDaemon {
run: "echo ready".to_string(),
path: Some(config_path),
..PitchforkTomlDaemon::default()
};
let run_opts = build_run_options(&id, &daemon_config, None).await.unwrap();
assert_eq!(run_opts.ready_delay, Some(7));
}
#[tokio::test]
async fn build_run_options_rejects_subsecond_ready_delay() {
const CHILD_SENTINEL: &str = "subsecond-ready-delay-child-ran";
let output = tokio::process::Command::new(std::env::current_exe().unwrap())
.args([
"--exact",
"ipc::batch::tests::build_run_options_rejects_subsecond_ready_delay_in_child",
"--nocapture",
])
.env(
"PITCHFORK_TEST_PROJECT_READY_DELAY_MODE",
"subsecond-ready-delay-project",
)
.env_remove("PITCHFORK_READY_DELAY")
.output()
.await
.unwrap();
assert!(
output.status.success(),
"sanitized child failed:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
assert!(
String::from_utf8_lossy(&output.stderr).contains(CHILD_SENTINEL),
"sanitized child test did not run:\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
}
#[tokio::test]
async fn build_run_options_rejects_subsecond_ready_delay_in_child() {
let Ok(mode) = std::env::var("PITCHFORK_TEST_PROJECT_READY_DELAY_MODE") else {
return;
};
if mode != "subsecond-ready-delay-project" {
return;
}
eprintln!("subsecond-ready-delay-child-ran");
let project = tempfile::tempdir().unwrap();
let config_path = project.path().join("pitchfork.toml");
tokio::fs::write(
&config_path,
"[settings.general]\nready_delay = \"500ms\"\n",
)
.await
.unwrap();
let id = DaemonId::try_new("other-project", "api").unwrap();
let daemon_config = PitchforkTomlDaemon {
run: "echo ready".to_string(),
path: Some(config_path),
..PitchforkTomlDaemon::default()
};
let err = build_run_options(&id, &daemon_config, None)
.await
.unwrap_err();
assert!(
err.contains("whole number of seconds"),
"unexpected error: {err}"
);
}
#[tokio::test]
async fn build_run_options_preserves_ready_port_timeout_for_cli_port_override() {
let id = DaemonId::try_new("project", "api").unwrap();
let daemon_config = PitchforkTomlDaemon {
run: "echo ready".to_string(),
ready_port: Some(ReadyPort {
port: Some(3000),
template: None,
timeout: Some(Duration::from_secs(45)),
}),
..PitchforkTomlDaemon::default()
};
let opts = StartOptions {
port: Some(4000),
..StartOptions::default()
};
let run_opts = build_run_options(&id, &daemon_config, Some(&opts))
.await
.unwrap();
let ready_port = run_opts.ready_port.unwrap();
assert_eq!(ready_port.port, Some(4000));
assert_eq!(ready_port.timeout, Some(Duration::from_secs(45)));
}
#[tokio::test]
async fn build_run_options_preserves_ready_output_timeout_for_cli_output_override() {
let id = DaemonId::try_new("project", "api").unwrap();
let daemon_config = PitchforkTomlDaemon {
run: "echo ready".to_string(),
ready_output: Some(ReadyOutput {
pattern: "READY".to_string(),
timeout: Some(Duration::from_secs(45)),
}),
..PitchforkTomlDaemon::default()
};
let opts = StartOptions {
output: Some("DONE".to_string()),
..StartOptions::default()
};
let run_opts = build_run_options(&id, &daemon_config, Some(&opts))
.await
.unwrap();
let ready_output = run_opts.ready_output.unwrap();
assert_eq!(ready_output.pattern, "DONE");
assert_eq!(ready_output.timeout, Some(Duration::from_secs(45)));
}
#[test]
fn test_resolve_daemon_dir_none() {
let result = resolve_daemon_dir(
None,
Some(Path::new("/projects/myapp/pitchfork.toml")),
None,
);
assert_eq!(result, PathBuf::from("/projects/myapp"));
}
#[test]
fn test_resolve_daemon_dir_relative() {
let result = resolve_daemon_dir(
Some("frontend"),
Some(Path::new("/projects/myapp/pitchfork.toml")),
None,
);
assert_eq!(result, PathBuf::from("/projects/myapp/frontend"));
}
#[test]
fn test_resolve_daemon_dir_absolute() {
let result = resolve_daemon_dir(
Some("/opt/myapp"),
Some(Path::new("/projects/myapp/pitchfork.toml")),
None,
);
assert_eq!(result, PathBuf::from("/opt/myapp"));
}
#[test]
fn test_resolve_daemon_dir_tilde() {
let result = resolve_daemon_dir(
Some("~/projects/myapp"),
Some(Path::new("/projects/other/pitchfork.toml")),
None,
);
assert_eq!(result, crate::env::HOME_DIR.join("projects/myapp"));
}
#[test]
fn test_resolve_daemon_dir_tilde_with_user() {
let result = resolve_daemon_dir(
Some("~/data"),
Some(Path::new("/projects/other/pitchfork.toml")),
Some("nonexistent_user_xyz"),
);
assert_eq!(result, crate::env::HOME_DIR.join("data"));
}
#[test]
fn test_resolve_daemon_dir_tilde_with_current_user() {
let current_user = std::env::var("USER").unwrap_or_else(|_| "root".to_string());
let expected = crate::env::home_dir_for_effective_user(Some(¤t_user)).join("data");
let result = resolve_daemon_dir(
Some("~/data"),
Some(Path::new("/projects/other/pitchfork.toml")),
Some(¤t_user),
);
assert_eq!(result, expected);
}
#[test]
fn test_resolve_daemon_dir_no_config_path() {
let result = resolve_daemon_dir(None, None, None);
assert_eq!(result, crate::env::CWD.to_path_buf());
}
#[test]
fn test_resolve_daemon_dir_relative_no_config_path() {
let result = resolve_daemon_dir(Some("subdir"), None, None);
assert_eq!(result, crate::env::CWD.join("subdir"));
}
#[test]
fn test_resolve_daemon_dir_nested_relative() {
let result = resolve_daemon_dir(
Some("services/api"),
Some(Path::new("/projects/myapp/pitchfork.toml")),
None,
);
assert_eq!(result, PathBuf::from("/projects/myapp/services/api"));
}
#[test]
fn test_resolve_daemon_dir_dot_config_none() {
let result = resolve_daemon_dir(
None,
Some(Path::new("/projects/myapp/.config/pitchfork.toml")),
None,
);
assert_eq!(
result,
PathBuf::from("/projects/myapp"),
".config/pitchfork.toml should resolve to project dir"
);
}
#[test]
fn test_resolve_daemon_dir_dot_config_local_none() {
let result = resolve_daemon_dir(
None,
Some(Path::new("/projects/myapp/.config/pitchfork.local.toml")),
None,
);
assert_eq!(
result,
PathBuf::from("/projects/myapp"),
".config/pitchfork.local.toml should resolve to project dir"
);
}
#[test]
fn test_resolve_daemon_dir_dot_config_relative() {
let result = resolve_daemon_dir(
Some("frontend"),
Some(Path::new("/projects/myapp/.config/pitchfork.toml")),
None,
);
assert_eq!(
result,
PathBuf::from("/projects/myapp/frontend"),
"Relative dir should resolve from project dir"
);
}
#[test]
fn test_resolve_daemon_dir_dot_config_local_relative() {
let result = resolve_daemon_dir(
Some("frontend"),
Some(Path::new("/projects/myapp/.config/pitchfork.local.toml")),
None,
);
assert_eq!(
result,
PathBuf::from("/projects/myapp/frontend"),
"Relative dir should resolve from project dir"
);
}
#[test]
fn test_resolve_daemon_dir_dot_config_absolute() {
let result = resolve_daemon_dir(
Some("/opt/service"),
Some(Path::new("/projects/myapp/.config/pitchfork.toml")),
None,
);
assert_eq!(
result,
PathBuf::from("/opt/service"),
"Absolute dir should override project dir"
);
}
#[test]
fn test_resolve_daemon_dir_dot_config_local_absolute() {
let result = resolve_daemon_dir(
Some("/opt/service"),
Some(Path::new("/projects/myapp/.config/pitchfork.local.toml")),
None,
);
assert_eq!(
result,
PathBuf::from("/opt/service"),
"Absolute dir should override project dir"
);
}
#[test]
fn test_resolve_daemon_dir_global_config_normal() {
let global_path = env::PITCHFORK_GLOBAL_CONFIG_USER.as_path();
let result = resolve_daemon_dir(None, Some(global_path), None);
assert_eq!(
result,
global_path.parent().unwrap(),
"Global config should use parent directory"
);
}
}