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,
handover: true,
}
}
#[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 {
#[cfg(unix)]
Handover,
StopAndStart,
}
#[cfg_attr(windows, allow(dead_code))]
const HANDOVER_SINCE: &str = "0.1.18";
#[cfg_attr(windows, allow(dead_code))]
fn version_parts(version: &str) -> Option<(u64, u64, u64)> {
let mut parts = version.split('.');
let major = parts.next()?.parse().ok()?;
let minor = parts.next()?.parse().ok()?;
let patch = parts.next()?;
let patch = patch
.split_once(['-', '+'])
.map_or(patch, |(number, _suffix)| number);
Some((major, minor, patch.parse().ok()?))
}
impl Arm {
fn for_daemon(daemon_version: Option<&str>) -> Self {
#[cfg(unix)]
{
let Some(running) = daemon_version.and_then(version_parts) else {
return Self::StopAndStart;
};
let floor = version_parts(HANDOVER_SINCE)
.expect("HANDOVER_SINCE is a literal three-number version");
let own = version_parts(env!("CARGO_PKG_VERSION"))
.expect("this crate's own version is a three-number version");
if running >= floor || running == own {
Self::Handover
} else {
Self::StopAndStart
}
}
#[cfg(windows)]
{
let _ = daemon_version;
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 connected = match Client::connect(&paths.socket).await {
Ok(client) => Ok(client),
Err(err) => Err(version_from_refusal(&err).map(str::to_owned)),
};
#[cfg(unix)]
let running_version = match &connected {
Ok(client) => Some(client.daemon().daemon_version.clone()),
Err(from_refusal) => from_refusal.clone(),
};
#[cfg(unix)]
if Arm::for_daemon(running_version.as_deref()) == Arm::Handover
&& let Ok(client) = &connected
{
match ask_fitness(client).await {
Fitness::Carryable => {
drop(connected);
return hand_over(streams, paths, guard, wait).await;
}
Fitness::Refused(reason) => {
streams.aside("reload", &reason);
}
}
}
drop(connected);
stop_and_start(streams, paths, guard, wait).await
}
#[cfg(unix)]
#[derive(Debug)]
enum Fitness {
Carryable,
Refused(String),
}
#[cfg(unix)]
async fn ask_fitness(client: &Client) -> Fitness {
match client
.request(shep_core::protocol::Request::HandoverFitness)
.await
{
Ok(shep_core::protocol::Response::HandoverFitness { refusal: None }) => Fitness::Carryable,
Ok(shep_core::protocol::Response::HandoverFitness {
refusal: Some(reason),
}) => Fitness::Refused(reason),
Ok(other) => Fitness::Refused(format!(
"this shepherd answered a handover question with {other:?}, so its flock is being \
stopped and started instead"
)),
Err(err) => Fitness::Refused(format!(
"this shepherd could not say whether its flock can be handed over ({err}), so it is \
being stopped and started instead"
)),
}
}
#[cfg(unix)]
async fn hand_over(
streams: &mut Streams<'_>,
paths: &ShepPaths,
guard: crate::VersionGuard,
wait: std::time::Duration,
) -> ExitCode {
let pid = match proven_shepherd(streams, paths) {
Ok(pid) => pid,
Err(code) => return code,
};
let Ok(witness) = Client::connect(&paths.socket).await else {
let message = "could not hold a connection across the handover signal; \
stopping and starting instead";
streams.aside("reload", message);
return stop_and_start(streams, paths, guard, wait).await;
};
if let Err((code, message)) = signal_handover(pid) {
return streams.fail(code, &message);
}
match await_successor(paths, &witness, wait).await {
Some(client) => report_reload(&client, streams, false).await,
None => {
let message = "the shepherd did not come back on this version after the handover \
signal; starting one instead";
streams.aside("reload", message);
let client = match crate::connect_or_spawn_client(streams, paths, guard).await {
Ok(client) => client,
Err(code) => return code,
};
report_reload(&client, streams, true).await
}
}
}
#[cfg(unix)]
fn signal_handover(pid: u32) -> Result<(), (ExitCode, String)> {
use nix::sys::signal::{self, Signal};
use nix::unistd::Pid;
let Ok(target) = i32::try_from(pid) else {
let message = format!("the recorded pid {pid} is not one this platform can signal");
return Err((ExitCode::Internal, message));
};
signal::kill(Pid::from_raw(target), Signal::SIGHUP).map_err(|errno| {
let message = format!("could not signal the shepherd at pid {pid}: {errno}");
(ExitCode::Failure, message)
})
}
#[cfg(unix)]
async fn await_successor(
paths: &ShepPaths,
witness: &Client,
wait: std::time::Duration,
) -> Option<Client> {
let deadline = tokio::time::Instant::now() + wait;
while witness
.request(shep_core::protocol::Request::ListFlock)
.await
.is_ok()
{
if tokio::time::Instant::now() >= deadline {
return None;
}
tokio::time::sleep(SUCCESSOR_POLL_INTERVAL).await;
}
loop {
if let Ok(client) = Client::connect(&paths.socket).await
&& client.daemon().daemon_version == env!("CARGO_PKG_VERSION")
&& client
.request(shep_core::protocol::Request::ListFlock)
.await
.is_ok()
{
return Some(client);
}
if tokio::time::Instant::now() >= deadline {
return None;
}
tokio::time::sleep(SUCCESSOR_POLL_INTERVAL).await;
}
}
#[cfg(unix)]
const SUCCESSOR_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(20);
fn proven_shepherd(streams: &mut Streams<'_>, paths: &ShepPaths) -> Result<u32, ExitCode> {
match boot::daemon_liveness(paths) {
Ok(Shepherd::Running(pid)) => Ok(pid),
Ok(Shepherd::Booting) => {
let message = "a shepherd is starting up and has not recorded its pid yet; try again";
Err(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()
);
Err(streams.fail(ExitCode::DaemonUnreachable, &message))
}
Err(err) => Err(streams.fail(ExitCode::Failure, &err.to_string())),
}
}
async fn stop_and_start(
streams: &mut Streams<'_>,
paths: &ShepPaths,
guard: crate::VersionGuard,
wait: std::time::Duration,
) -> ExitCode {
let pid = match proven_shepherd(streams, paths) {
Ok(pid) => pid,
Err(code) => return code,
};
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, true).await
}
async fn report_reload(client: &Client, streams: &mut Streams<'_>, restored: bool) -> ExitCode {
report_reload_waiting(client, streams, restored, DOG_SETTLE_WAIT).await
}
async fn report_reload_waiting(
client: &Client,
streams: &mut Streams<'_>,
restored: bool,
dog_wait: std::time::Duration,
) -> ExitCode {
let shepherd = client.daemon();
let message = format!(
"the shepherd is now {} (pid {})",
shepherd.daemon_version, shepherd.pid
);
streams.aside("reload", &message);
report_dog_staleness(client, streams, &shepherd.daemon_version, dog_wait).await;
if restored {
return muster::muster(client, streams).await;
}
crate::commands::query::flock(client, streams).await
}
const DOG_SETTLE_WAIT: std::time::Duration = std::time::Duration::from_secs(3);
const DOG_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(50);
async fn report_dog_staleness(
client: &Client,
streams: &mut Streams<'_>,
daemon_version: &str,
wait: std::time::Duration,
) {
let deadline = tokio::time::Instant::now() + wait;
loop {
let Ok(shep_core::protocol::Response::DogStaleness { stale, pending }) = client
.request(shep_core::protocol::Request::DogStaleness)
.await
else {
return;
};
let out_of_time = tokio::time::Instant::now() >= deadline;
if pending.is_empty() || out_of_time {
if !stale.is_empty() {
streams.aside("reload", &stale_dog_report(&stale, daemon_version));
}
if out_of_time && !pending.is_empty() {
streams.aside("reload", &unsettled_dog_report(&pending, wait));
}
return;
}
tokio::time::sleep(DOG_POLL_INTERVAL).await;
}
}
fn stale_dog_report(stale: &[String], daemon_version: &str) -> String {
match stale {
[only] => format!(
"the `{only}` dog cannot talk to this shepherd; restarting it from the binary on \
disk did not help, so rebuild or reinstall it against shep {daemon_version}, then \
restart it"
),
many => format!(
"these dogs cannot talk to this shepherd: {}; restarting them from the binaries on \
disk did not help, so rebuild or reinstall them against shep {daemon_version}, then \
restart them",
quoted_names(many)
),
}
}
fn unsettled_dog_report(pending: &[String], wait: std::time::Duration) -> String {
let budget = shep_daemon::dogs::DOG_SILENCE_BUDGET;
match pending {
[only] => format!(
"the `{only}` dog has not answered this shepherd after {wait:?}; a dog silent \
past {budget:?} is restarted once from the binary on disk and then reported \
stale, and `shep bleats {only}` shows why"
),
many => format!(
"these dogs have not answered this shepherd after {wait:?}: {}; a dog silent past \
{budget:?} is restarted once from the binary on disk and then reported stale, and \
`shep bleats <dog>` shows why for each",
quoted_names(many)
),
}
}
fn quoted_names(names: &[String]) -> String {
names
.iter()
.map(|name| format!("`{name}`"))
.collect::<Vec<_>>()
.join(", ")
}
#[cfg(test)]
mod tests {
use super::*;
use crate::VersionGuard;
use crate::cli::Format;
use shep_core::config::LogLevel;
use shep_core::protocol::{Request, 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_at_every_version_below_the_floor() {
assert_eq!(Arm::for_daemon(Some("0.1.9")), Arm::StopAndStart);
assert_eq!(Arm::for_daemon(Some("0.1.16")), Arm::StopAndStart);
assert_eq!(
Arm::for_daemon(Some("not a version")),
Arm::StopAndStart,
"a version this CLI cannot read is unknown, and unknown is the safe arm"
);
}
#[cfg(unix)]
#[test]
fn a_shepherd_of_this_binarys_own_version_answers_for_itself() {
assert_eq!(
Arm::for_daemon(Some(env!("CARGO_PKG_VERSION"))),
Arm::Handover
);
}
#[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, true).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}");
}
#[cfg(unix)]
#[test]
fn reload_picks_the_handover_against_a_daemon_new_enough_to_carry_its_flock() {
assert_eq!(Arm::for_daemon(Some("9.9.9")), Arm::Handover);
assert_eq!(Arm::for_daemon(Some(HANDOVER_SINCE)), Arm::Handover);
}
#[tokio::test]
async fn a_reload_against_an_older_daemon_never_sends_the_query() {
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 sent = shep_client::testing::fake_daemon_answering_with_ack(
&paths.socket,
ack_naming("0.1.8"),
|_| Response::Pong,
)
.await;
let mut out = Vec::new();
let mut err = Vec::new();
{
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;
}
let asked: Vec<Request> = std::iter::from_fn(|| sent.try_recv().ok())
.map(|envelope| envelope.body)
.collect();
assert!(
!asked
.iter()
.any(|req| matches!(req, Request::HandoverFitness)),
"a daemon that cannot parse the query must never be asked it: {asked:?}"
);
}
#[cfg(unix)]
#[tokio::test]
async fn a_reload_against_a_newer_daemon_asks_before_it_signals() {
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 sent = shep_client::testing::fake_daemon_answering_with_ack(
&paths.socket,
ack_naming("9.9.9"),
|_| Response::HandoverFitness { refusal: None },
)
.await;
let mut out = Vec::new();
let mut err = Vec::new();
{
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;
}
let asked: Vec<Request> = std::iter::from_fn(|| sent.try_recv().ok())
.map(|envelope| envelope.body)
.collect();
assert_eq!(
asked
.iter()
.filter(|req| matches!(req, Request::HandoverFitness))
.count(),
1,
"exactly one fitness query, and it is the first thing asked: {asked:?}"
);
}
#[cfg(unix)]
#[tokio::test]
async fn a_refused_flock_prints_the_reason_and_falls_back() {
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 _sent = shep_client::testing::fake_daemon_answering_with_ack(
&paths.socket,
ack_naming("9.9.9"),
|_| Response::HandoverFitness {
refusal: Some("sheep 'clustered' has more than one instance".to_string()),
},
)
.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,
};
reload(&mut streams, &paths, VersionGuard::Exempt).await
};
let text = format!(
"{}{}",
String::from_utf8(out).unwrap(),
String::from_utf8(err).unwrap()
);
assert!(text.contains("more than one instance"), "{text}");
assert_eq!(code, ExitCode::DaemonUnreachable, "{text}");
}
#[tokio::test]
async fn a_reload_waits_for_a_pending_dog_before_it_reports_staleness() {
let dir = tempfile::tempdir().unwrap();
let addr = shep_client::testing::control_address(dir.path());
let asks = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let counted = std::sync::Arc::clone(&asks);
let (client, _envelopes) = shep_client::testing::fake_client_answering(&addr, move |req| {
if !matches!(req, Request::DogStaleness) {
return Response::Mustered(vec![]);
}
let seen = counted.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if seen < 2 {
Response::DogStaleness {
stale: vec![],
pending: vec!["metrics".to_string()],
}
} else {
Response::DogStaleness {
stale: vec!["metrics".to_string()],
pending: vec![],
}
}
})
.await;
let mut out = Vec::new();
let mut err = Vec::new();
{
let mut streams = Streams {
out: &mut out,
err: &mut err,
style: crate::style::Presentation::BARE,
fmt: Format::Table,
};
report_reload_waiting(
&client,
&mut streams,
true,
std::time::Duration::from_secs(3),
)
.await;
}
let text = String::from_utf8(err).unwrap();
assert!(
text.contains("metrics") && text.contains("rebuild or reinstall"),
"the dog that could not come back must be named: {text}"
);
assert!(
asks.load(std::sync::atomic::Ordering::SeqCst) >= 3,
"an answer taken on the first ask is a claim about a dog that had not spoken"
);
}
#[tokio::test]
async fn a_reload_whose_dogs_all_answered_says_nothing_about_them() {
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| {
if matches!(req, Request::DogStaleness) {
Response::DogStaleness {
stale: vec![],
pending: vec![],
}
} else {
Response::Mustered(vec![shep_client::testing::sample_info()])
}
})
.await;
let mut out = Vec::new();
let mut err = Vec::new();
{
let mut streams = Streams {
out: &mut out,
err: &mut err,
style: crate::style::Presentation::BARE,
fmt: Format::Table,
};
report_reload_waiting(
&client,
&mut streams,
true,
std::time::Duration::from_secs(3),
)
.await;
}
let text = format!(
"{}{}",
String::from_utf8(out).unwrap(),
String::from_utf8(err).unwrap()
);
assert!(
!text.contains("dog"),
"a flock whose dogs all came back has nothing to say about them: {text}"
);
}
#[tokio::test]
async fn a_reload_stops_waiting_for_a_dog_that_never_answers() {
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| {
if matches!(req, Request::DogStaleness) {
Response::DogStaleness {
stale: vec![],
pending: vec!["metrics".to_string()],
}
} else {
Response::Mustered(vec![])
}
})
.await;
let mut out = Vec::new();
let mut err = Vec::new();
{
let mut streams = Streams {
out: &mut out,
err: &mut err,
style: crate::style::Presentation::BARE,
fmt: Format::Table,
};
tokio::time::timeout(
std::time::Duration::from_secs(5),
report_reload_waiting(
&client,
&mut streams,
true,
std::time::Duration::from_millis(150),
),
)
.await
.expect("a dog that never answers must not hold the verb open");
}
let text = String::from_utf8(err).unwrap();
assert!(
text.contains("metrics") && text.contains("shep bleats metrics"),
"an unanswered dog is reported as unanswered, not as healthy: {text}"
);
}
#[test]
fn the_stale_report_says_what_happened_and_never_reads_the_disk() {
let one = stale_dog_report(&["metrics".to_string()], "0.1.22");
assert_eq!(
one,
"the `metrics` dog cannot talk to this shepherd; restarting it from the binary on \
disk did not help, so rebuild or reinstall it against shep 0.1.22, then restart it"
);
let two = stale_dog_report(&["bark".to_string(), "metrics".to_string()], "0.1.22");
assert_eq!(
two,
"these dogs cannot talk to this shepherd: `bark`, `metrics`; restarting them from \
the binaries on disk did not help, so rebuild or reinstall them against shep \
0.1.22, then restart them"
);
}
#[test]
fn the_unsettled_report_says_what_to_check_and_never_claims_a_verdict() {
let one = unsettled_dog_report(&["metrics".to_string()], std::time::Duration::from_secs(3));
assert_eq!(
one,
"the `metrics` dog has not answered this shepherd after 3s; a dog silent past 5s \
is restarted once from the binary on disk and then reported stale, and `shep \
bleats metrics` shows why"
);
let two = unsettled_dog_report(
&["bark".to_string(), "metrics".to_string()],
std::time::Duration::from_secs(3),
);
assert_eq!(
two,
"these dogs have not answered this shepherd after 3s: `bark`, `metrics`; a dog \
silent past 5s is restarted once from the binary on disk and then reported \
stale, and `shep bleats <dog>` shows why for each"
);
}
fn ack_naming(version: &str) -> shep_core::protocol::HelloAck {
shep_core::protocol::HelloAck {
daemon_version: version.to_string(),
protocol: shep_core::protocol::PROTOCOL_VERSION,
pid: 4242,
}
}
}