use std::ffi::OsStr;
use std::io::IsTerminal;
use shep_client::{Client, ConnectError};
use shep_core::config::{DaemonConfig, DaemonConfigError, DaemonOverrides};
use shep_core::paths::ShepPaths;
use shep_core::protocol::DogSource;
use shep_core::values::UpDuration;
use shep_daemon::boot::{self, BootError, BootOptions, RunningDaemon, Shepherd, boot};
use shep_daemon::dogs::DogSpec;
#[cfg(unix)]
use shep_daemon::notify::NOTIFY_SOCKET_ENV;
use shep_daemon::tokio_runner::TokioRunner;
use tracing_subscriber::EnvFilter;
use crate::cli::DaemonArgs;
use crate::commands::{admin, muster};
use crate::exit::ExitCode;
use crate::output::Streams;
#[derive(Debug)]
pub enum DaemonRunError {
Config(DaemonConfigError),
Boot(BootError),
Run(BootError),
}
impl core::fmt::Display for DaemonRunError {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
match self {
Self::Config(err) => write!(f, "invalid daemon configuration: {err}"),
Self::Boot(err) => write!(f, "the daemon failed to boot: {err}"),
Self::Run(err) => write!(f, "the daemon failed while running: {err}"),
}
}
}
impl core::error::Error for DaemonRunError {
fn source(&self) -> Option<&(dyn core::error::Error + 'static)> {
match self {
Self::Config(err) => Some(err),
Self::Boot(err) | Self::Run(err) => Some(err),
}
}
}
impl From<DaemonConfigError> for DaemonRunError {
fn from(source: DaemonConfigError) -> Self {
Self::Config(source)
}
}
fn read_daemon_config_source(paths: &ShepPaths) -> Result<Option<String>, DaemonRunError> {
match std::fs::read_to_string(&paths.daemon_config) {
Ok(src) => Ok(Some(src)),
Err(source) if source.kind() == std::io::ErrorKind::NotFound => Ok(None),
Err(source) => Err(DaemonRunError::Boot(BootError::Io {
path: paths.daemon_config.clone(),
source,
})),
}
}
fn install_log_subscriber(config: &DaemonConfig) {
let builder = tracing_subscriber::fmt()
.with_env_filter(EnvFilter::new(config.daemon.log_level.as_str()))
.with_writer(std::io::stderr);
let installed = if config.daemon.log_json {
builder.json().try_init()
} else {
builder
.with_ansi(ansi_enabled(
std::io::stderr().is_terminal(),
std::env::var_os("NO_COLOR").as_deref(),
))
.try_init()
};
if let Err(err) = installed {
eprintln!("shep: the daemon's own logs are not being rendered: {err}");
}
}
fn ansi_enabled(stderr_is_terminal: bool, no_color: Option<&OsStr>) -> bool {
stderr_is_terminal && no_color.is_none_or(OsStr::is_empty)
}
pub async fn boot_supervisor(
paths: ShepPaths,
args: &DaemonArgs,
delete_flock_on_shutdown: bool,
) -> Result<RunningDaemon, DaemonRunError> {
let env = |key: &str| std::env::var(key).ok();
let file_source = read_daemon_config_source(&paths)?;
let overrides = daemon_overrides(args);
let config = DaemonConfig::load_layered(file_source.as_deref(), &env, &overrides)?;
install_log_subscriber(&config);
#[cfg(unix)]
let notify_socket = std::env::var_os(NOTIFY_SOCKET_ENV);
#[cfg(windows)]
let notify_socket: Option<std::ffi::OsString> = None;
let mut options = boot_options(&config, args, notify_socket.as_deref());
options.delete_flock_on_shutdown = delete_flock_on_shutdown;
boot(TokioRunner::new(), paths, options)
.await
.map_err(DaemonRunError::Boot)
}
pub async fn run_daemon(paths: ShepPaths, args: &DaemonArgs) -> Result<(), DaemonRunError> {
boot_supervisor(paths, args, false)
.await?
.run()
.await
.map_err(DaemonRunError::Run)
}
#[must_use]
pub fn daemon_overrides(args: &DaemonArgs) -> DaemonOverrides {
DaemonOverrides::new()
.log_json(args.log_json)
.log_level(args.log_level)
.socket(args.socket.clone())
.max_cron_sleep(args.max_cron_sleep)
}
#[must_use]
pub fn boot_options(
config: &DaemonConfig,
args: &DaemonArgs,
notify_socket: Option<&OsStr>,
) -> BootOptions {
BootOptions {
socket: config.daemon.socket.clone(),
ready_fd: None,
restore: !args.no_restore,
max_cron_sleep: config.daemon.max_cron_sleep.map(UpDuration::as_duration),
notify_socket: notify_socket
.filter(|_| args.foreground)
.map(OsStr::to_os_string),
dogs: config
.daemon
.enabled_dogs
.iter()
.map(|name| {
let source = match config.daemon.adopted_dogs.get(name) {
Some(path) => DogSource::Adopted {
path: path.display().to_string(),
},
None => DogSource::BuiltIn,
};
DogSpec {
name: name.clone(),
source,
}
})
.collect(),
delete_flock_on_shutdown: false,
}
}
#[must_use]
pub fn daemon_exit_code(err: &DaemonRunError) -> ExitCode {
match err {
DaemonRunError::Config(_) => ExitCode::InvalidConfig,
DaemonRunError::Boot(boot_err) => match boot_err {
BootError::AlreadyRunning { .. } => ExitCode::DaemonAlreadyRunning,
_ => ExitCode::Failure,
},
DaemonRunError::Run(_) => ExitCode::Failure,
}
}
#[derive(Debug, PartialEq, Eq)]
enum Arm {
#[expect(
dead_code,
reason = "phase 2 constructs this; phase 1 ships no handover for any version to support"
)]
Handover,
StopAndStart,
}
impl Arm {
fn for_daemon(_daemon_version: Option<&str>) -> Self {
Self::StopAndStart
}
}
fn version_from_refusal(err: &ConnectError) -> Option<&str> {
match err {
ConnectError::ProtocolMismatch { daemon_version, .. } => daemon_version.as_deref(),
_ => None,
}
}
pub async fn reload(
streams: &mut Streams<'_>,
paths: &ShepPaths,
guard: crate::VersionGuard,
) -> ExitCode {
reload_with_wait(streams, paths, guard, admin::KILL_TEARDOWN_WAIT).await
}
async fn reload_with_wait(
streams: &mut Streams<'_>,
paths: &ShepPaths,
guard: crate::VersionGuard,
wait: std::time::Duration,
) -> ExitCode {
let running_version = match Client::connect(&paths.socket).await {
Ok(client) => Some(client.daemon().daemon_version.clone()),
Err(err) => version_from_refusal(&err).map(str::to_owned),
};
match Arm::for_daemon(running_version.as_deref()) {
Arm::Handover => unreachable!("the handover arm lands in phase 2"),
Arm::StopAndStart => stop_and_start(streams, paths, guard, wait).await,
}
}
async fn stop_and_start(
streams: &mut Streams<'_>,
paths: &ShepPaths,
guard: crate::VersionGuard,
wait: std::time::Duration,
) -> ExitCode {
let pid = match boot::daemon_liveness(paths) {
Ok(Shepherd::Running(pid)) => pid,
Ok(Shepherd::Booting) => {
let message = "a shepherd is starting up and has not recorded its pid yet; try again";
return streams.fail(ExitCode::DaemonUnreachable, message);
}
Ok(Shepherd::Absent) => {
let message = format!(
"no shepherd is running, so there is nothing to reload (nothing holds the lock \
on `{}`). `shep muster` brings the flock up from the roll",
boot::pidfile(paths).display()
);
return streams.fail(ExitCode::DaemonUnreachable, &message);
}
Err(err) => return streams.fail(ExitCode::Failure, &err.to_string()),
};
if let Err((code, message)) = admin::signal_graceful_stop(pid) {
return streams.fail(code, &message);
}
if !admin::wait_for_socket_to_disappear(&paths.socket, wait).await {
let message = "the shepherd was signalled, but teardown is still in progress; \
nothing has been started in its place";
return streams.fail(ExitCode::DeadlineExceeded, message);
}
let client = match crate::connect_or_spawn_client(streams, paths, guard).await {
Ok(client) => client,
Err(code) => return code,
};
report_reload(&client, streams).await
}
async fn report_reload(client: &Client, streams: &mut Streams<'_>) -> ExitCode {
let shepherd = client.daemon();
let message = format!(
"the shepherd is now {} (pid {})",
shepherd.daemon_version, shepherd.pid
);
streams.aside("reload", &message);
muster::muster(client, streams).await
}
#[cfg(test)]
mod tests {
use super::*;
use crate::VersionGuard;
use crate::cli::Format;
use shep_core::config::LogLevel;
use shep_core::protocol::Response;
#[test]
fn every_daemon_flag_reaches_the_config() {
let args = DaemonArgs {
cmd: None,
no_restore: false,
foreground: false,
log_json: Some(true),
log_level: Some(LogLevel::Trace),
socket: Some(std::path::PathBuf::from("/tmp/flag.sock")),
max_cron_sleep: Some(UpDuration::from_millis(120_000)),
};
let cfg = DaemonConfig::load_layered(
Some(
"[daemon]\nlog_json = false\nlog_level = \"error\"\nsocket = \"/tmp/file.sock\"\n",
),
&|_| None,
&daemon_overrides(&args),
)
.unwrap();
assert!(cfg.daemon.log_json);
assert_eq!(cfg.daemon.log_level, LogLevel::Trace);
assert_eq!(
cfg.daemon.socket,
Some(std::path::PathBuf::from("/tmp/flag.sock"))
);
assert_eq!(
cfg.daemon.max_cron_sleep,
Some(UpDuration::from_millis(120_000))
);
}
#[test]
fn colour_needs_a_terminal_and_no_no_color() {
assert!(ansi_enabled(true, None));
assert!(!ansi_enabled(false, None), "a file never gets escape codes");
assert!(!ansi_enabled(true, Some(OsStr::new("1"))));
assert!(
ansi_enabled(true, Some(OsStr::new(""))),
"an empty NO_COLOR is an unset NO_COLOR"
);
assert!(
!ansi_enabled(false, Some(OsStr::new("1"))),
"the two reasons to suppress colour must not cancel out"
);
}
#[test]
fn boot_options_pass_ready_fd_none_and_the_configured_socket() {
let config =
DaemonConfig::load(Some("[daemon]\nsocket = \"/tmp/custom.sock\"\n"), &|_| None)
.unwrap();
let opts = boot_options(
&config,
&DaemonArgs {
cmd: None,
no_restore: false,
foreground: false,
log_json: None,
log_level: None,
socket: None,
max_cron_sleep: None,
},
None,
);
assert!(
opts.ready_fd.is_none(),
"readiness is a handshake in this phase"
);
assert_eq!(
opts.socket.as_deref(),
Some(std::path::Path::new("/tmp/custom.sock"))
);
assert!(opts.restore, "the default is to restore the muster roll");
}
#[test]
fn boot_options_carry_every_enabled_dog_with_the_source_the_file_names() {
let src = r#"
[daemon]
enabled_dogs = ["metrics", "otel"]
[daemon.adopted_dogs]
otel = "/usr/local/bin/shep-otel"
"#;
let config = DaemonConfig::load(Some(src), &|_| None).unwrap();
let opts = boot_options(
&config,
&DaemonArgs {
cmd: None,
no_restore: false,
foreground: false,
log_json: None,
log_level: None,
socket: None,
max_cron_sleep: None,
},
None,
);
assert_eq!(
opts.dogs,
vec![
DogSpec {
name: "metrics".into(),
source: DogSource::BuiltIn
},
DogSpec {
name: "otel".into(),
source: DogSource::Adopted {
path: "/usr/local/bin/shep-otel".into()
}
},
]
);
}
#[test]
fn boot_options_carry_the_configured_max_cron_sleep_and_invent_none() {
let configured =
DaemonConfig::load(Some("[daemon]\nmax_cron_sleep = \"5m\"\n"), &|_| None).unwrap();
assert_eq!(
boot_options(
&configured,
&DaemonArgs {
cmd: None,
no_restore: false,
foreground: false,
log_json: None,
log_level: None,
socket: None,
max_cron_sleep: None,
},
None
)
.max_cron_sleep,
Some(core::time::Duration::from_secs(300))
);
let unset = DaemonConfig::load(None, &|_| None).unwrap();
assert_eq!(
boot_options(
&unset,
&DaemonArgs {
cmd: None,
no_restore: false,
foreground: false,
log_json: None,
log_level: None,
socket: None,
max_cron_sleep: None,
},
None
)
.max_cron_sleep,
None,
"an unset knob must stay None: the daemon owns the default"
);
}
#[test]
fn no_restore_boots_without_the_muster_roll() {
let config = DaemonConfig::load(None, &|_| None).unwrap();
let opts = boot_options(
&config,
&DaemonArgs {
cmd: None,
no_restore: true,
foreground: false,
log_json: None,
log_level: None,
socket: None,
max_cron_sleep: None,
},
None,
);
assert!(!opts.restore);
}
#[test]
fn the_foreground_flag_reaches_the_boot_options() {
let config = DaemonConfig::load(None, &|_| None).unwrap();
let bare = boot_options(
&config,
&DaemonArgs {
cmd: None,
no_restore: false,
foreground: false,
log_json: None,
log_level: None,
socket: None,
max_cron_sleep: None,
},
None,
);
assert!(
bare.notify_socket.is_none(),
"an autostarted daemon reports to nobody"
);
let supervised = boot_options(
&config,
&DaemonArgs {
cmd: None,
no_restore: false,
foreground: true,
log_json: None,
log_level: None,
socket: None,
max_cron_sleep: None,
},
Some(OsStr::new("/run/systemd/notify")),
);
assert_eq!(
supervised.notify_socket.as_deref(),
Some(OsStr::new("/run/systemd/notify"))
);
let unflagged = boot_options(
&config,
&DaemonArgs {
cmd: None,
no_restore: false,
foreground: true,
log_json: None,
log_level: None,
socket: None,
max_cron_sleep: None,
},
None,
);
assert!(unflagged.notify_socket.is_none());
let inherited = boot_options(
&config,
&DaemonArgs {
cmd: None,
no_restore: false,
foreground: false,
log_json: None,
log_level: None,
socket: None,
max_cron_sleep: None,
},
Some(OsStr::new("/run/systemd/notify")),
);
assert!(inherited.notify_socket.is_none());
}
#[test]
fn foreground_and_no_restore_are_independent() {
let config = DaemonConfig::load(None, &|_| None).unwrap();
let opts = boot_options(
&config,
&DaemonArgs {
cmd: None,
no_restore: false,
foreground: true,
log_json: None,
log_level: None,
socket: None,
max_cron_sleep: None,
},
Some(OsStr::new("/run/systemd/notify")),
);
assert!(opts.restore, "a supervised daemon still musters its roll");
}
#[test]
fn already_running_gets_its_own_exit_code_and_everything_else_is_failure() {
use DaemonRunError::{Boot, Config, Run};
assert_eq!(
daemon_exit_code(&Boot(BootError::AlreadyRunning { pid: Some(7) })),
ExitCode::DaemonAlreadyRunning
);
assert_eq!(
daemon_exit_code(&Boot(BootError::AlreadyRunning { pid: None })),
ExitCode::DaemonAlreadyRunning
);
assert_eq!(
daemon_exit_code(&Boot(BootError::Io {
path: "/x".into(),
source: std::io::Error::other("x"),
})),
ExitCode::Failure
);
assert_eq!(
daemon_exit_code(&Run(BootError::Io {
path: "/x".into(),
source: std::io::Error::other("x"),
})),
ExitCode::Failure
);
assert_eq!(
daemon_exit_code(&Config(DaemonConfigError::Toml("expected `=`".into()))),
ExitCode::InvalidConfig
);
assert_eq!(
daemon_exit_code(&Config(DaemonConfigError::BadEnvValue(
"SHEP_LOG_JSON",
"maybe".into()
))),
ExitCode::InvalidConfig
);
}
#[test]
fn boot_and_run_report_different_phases_for_the_same_underlying_error() {
use DaemonRunError::{Boot, Run};
let io_err = || BootError::Io {
path: "/x".into(),
source: std::io::Error::other("x"),
};
let boot_msg = Boot(io_err()).to_string();
let run_msg = Run(io_err()).to_string();
assert_ne!(boot_msg, run_msg);
assert!(boot_msg.starts_with("the daemon failed to boot"));
assert!(
!run_msg.starts_with("the daemon failed to boot"),
"a run-phase failure must not still claim to be a boot failure: {run_msg:?}"
);
assert!(run_msg.starts_with("the daemon failed while running"));
}
#[test]
fn reload_picks_the_stop_arm_against_a_daemon_too_old_to_hand_over() {
assert_eq!(Arm::for_daemon(Some("0.1.8")), Arm::StopAndStart);
}
#[test]
fn reload_picks_the_stop_arm_against_a_daemon_of_this_very_version() {
assert_eq!(
Arm::for_daemon(Some(env!("CARGO_PKG_VERSION"))),
Arm::StopAndStart
);
}
#[test]
fn reload_picks_the_stop_arm_when_the_handshake_is_refused_without_a_version() {
let refusal = ConnectError::ProtocolMismatch {
client: shep_core::protocol::PROTOCOL_VERSION,
daemon_version: None,
message: "this daemon speaks protocol 1".to_string(),
};
assert_eq!(version_from_refusal(&refusal), None);
assert_eq!(
Arm::for_daemon(version_from_refusal(&refusal)),
Arm::StopAndStart
);
}
#[test]
fn a_refusal_that_names_a_version_yields_it_for_the_arm_choice() {
let refusal = ConnectError::ProtocolMismatch {
client: shep_core::protocol::PROTOCOL_VERSION,
daemon_version: Some("0.1.8".to_string()),
message: "this daemon speaks protocol 1".to_string(),
};
assert_eq!(version_from_refusal(&refusal), Some("0.1.8"));
}
#[test]
fn a_connect_failure_that_is_not_a_refusal_names_no_version() {
let err = ConnectError::HandshakeClosed;
assert_eq!(version_from_refusal(&err), None);
}
#[tokio::test]
async fn reload_reports_each_sheep_rather_than_announcing_the_flock_stopped() {
let dir = tempfile::tempdir().unwrap();
let addr = shep_client::testing::control_address(dir.path());
let (client, _envelopes) = shep_client::testing::fake_client_answering(&addr, |_req| {
Response::Mustered(vec![shep_client::testing::sample_info()])
})
.await;
let mut out = Vec::new();
let mut err = Vec::new();
let code = {
let mut streams = Streams {
out: &mut out,
err: &mut err,
style: crate::style::Presentation::BARE,
fmt: Format::Table,
};
report_reload(&client, &mut streams).await
};
assert_eq!(code, ExitCode::Success);
let text = format!(
"{}{}",
String::from_utf8(out).unwrap(),
String::from_utf8(err).unwrap()
);
assert!(text.contains("web"), "{text}");
assert!(!text.to_lowercase().contains("flock stopped"), "{text}");
}
#[tokio::test]
async fn reload_refuses_a_home_no_shepherd_owns() {
let dir = tempfile::tempdir().unwrap();
let paths = ShepPaths::resolve(&|_| None, dir.path());
std::fs::create_dir_all(&paths.run).unwrap();
std::fs::create_dir_all(&paths.pids).unwrap();
let mut out = Vec::new();
let mut err = Vec::new();
let code = {
let mut streams = Streams {
out: &mut out,
err: &mut err,
style: crate::style::Presentation::BARE,
fmt: Format::Table,
};
reload(&mut streams, &paths, VersionGuard::Exempt).await
};
assert_eq!(code, ExitCode::DaemonUnreachable);
let text = String::from_utf8(err).unwrap();
assert!(text.contains("no shepherd"), "{text}");
}
}