use shep_client::{Client, START_DEADLINE};
use shep_core::config::AppConfig;
use shep_core::paths::ShepPaths;
use shep_core::protocol::{Request, Response, SelectorSpec};
use crate::cli::{BleatsArgs, DaemonArgs};
use crate::commands::bleats::bleats_with_signal;
use crate::commands::daemon::{boot_supervisor, daemon_exit_code};
use crate::commands::empty::{Sample, sample, watch_until_empty};
use crate::exit::ExitCode;
use crate::output::Streams;
pub struct ForegroundOptions {
pub paths: ShepPaths,
pub apps: Vec<AppConfig>,
pub tidy_up: bool,
}
enum Ending {
Empty(Sample),
SupervisorExited { failed: bool },
}
fn refuse_if_skewed(streams: &mut Streams<'_>, client: &Client) -> Result<(), ExitCode> {
crate::refuse_version_skew(streams, client, crate::VersionGuard::Enforce)
}
pub async fn run(streams: &mut Streams<'_>, quiet: bool, options: ForegroundOptions) -> ExitCode {
let ForegroundOptions {
paths,
apps,
tidy_up,
} = options;
let daemon_args = DaemonArgs {
cmd: None,
no_restore: true,
foreground: true,
log_json: None,
log_level: None,
socket: None,
max_cron_sleep: None,
};
let daemon = match boot_supervisor(paths.clone(), &daemon_args, tidy_up).await {
Ok(daemon) => daemon,
Err(err) => {
let code = daemon_exit_code(&err);
return streams.fail(code, &err.to_string());
}
};
let mut supervisor = tokio::spawn(daemon.run());
let client = match Client::connect(&paths.socket).await {
Ok(client) => {
if let Err(code) = refuse_if_skewed(streams, &client) {
supervisor.abort();
return code;
}
client
}
Err(err) => {
supervisor.abort();
let code = ExitCode::from(&err);
return streams.fail(code, &err.to_string());
}
};
if let Err(err) = client
.request_with_deadline(Request::Start { apps }, Some(START_DEADLINE))
.await
{
let code = ExitCode::from(&err);
streams.fail(code, &err.to_string());
let _ = client.request(Request::KillDaemon).await;
let _ = supervisor.await;
return code;
}
let mut ending: Option<Ending> = None;
{
let watcher = watch_until_empty(|| async {
match client.request(Request::ListFlock).await {
Ok(Response::Flock(procs)) => sample(&procs),
_ => Sample::Busy,
}
});
tokio::pin!(watcher);
let interrupt = async {
tokio::select! {
reading = &mut watcher => {
ending = Some(Ending::Empty(reading));
}
result = &mut supervisor => {
let failed = !matches!(result, Ok(Ok(())));
ending = Some(Ending::SupervisorExited { failed });
}
}
};
bleats_with_signal(
&client,
streams,
quiet,
&BleatsArgs {
selector: "all".to_string(),
no_follow: false,
lines: 0,
err: false,
out: false,
},
interrupt,
)
.await;
}
let already_consumed = matches!(ending, Some(Ending::SupervisorExited { .. }));
if tidy_up && !already_consumed {
let _ = client
.request(Request::Stop {
selector: SelectorSpec::All,
})
.await;
let _ = client
.request(Request::Delete {
selector: SelectorSpec::All,
})
.await;
}
let _ = client.request(Request::KillDaemon).await;
let supervisor_failed = if already_consumed {
matches!(ending, Some(Ending::SupervisorExited { failed: true }))
} else {
!matches!(supervisor.await, Ok(Ok(())))
};
if supervisor_failed {
return ExitCode::Failure;
}
match ending {
Some(Ending::Empty(Sample::EmptyFailed)) => ExitCode::FlockEmpty,
Some(Ending::Empty(Sample::EmptyClean)) => ExitCode::Success,
Some(Ending::Empty(Sample::Busy)) | Some(Ending::SupervisorExited { .. }) | None => {
ExitCode::Success
}
}
}
#[cfg(test)]
mod tests {
use crate::cli::Format;
use crate::exit::ExitCode;
use crate::output::Streams;
use crate::style;
fn buffered_streams<'a>(out: &'a mut Vec<u8>, err: &'a mut Vec<u8>) -> Streams<'a> {
Streams {
out,
err,
style: style::Presentation::BARE,
fmt: Format::Table,
}
}
#[tokio::test]
async fn a_version_skewed_shepherd_is_refused() {
let dir = tempfile::tempdir().unwrap();
let addr = shep_client::testing::control_address(dir.path());
let ack = shep_core::protocol::HelloAck {
daemon_version: "0.1.8".to_string(),
protocol: shep_core::protocol::PROTOCOL_VERSION,
pid: 4242,
};
let (client, _fake) = shep_client::testing::fake_client_with_ack(&addr, ack).await;
let mut out = Vec::new();
let mut err = Vec::new();
let mut streams = buffered_streams(&mut out, &mut err);
let code = super::refuse_if_skewed(&mut streams, &client)
.expect_err("a differing crate version must be refused");
assert_eq!(code, ExitCode::VersionSkew);
}
#[tokio::test]
async fn a_matching_version_proceeds() {
let dir = tempfile::tempdir().unwrap();
let addr = shep_client::testing::control_address(dir.path());
let ack = shep_core::protocol::HelloAck {
daemon_version: env!("CARGO_PKG_VERSION").to_string(),
protocol: shep_core::protocol::PROTOCOL_VERSION,
pid: 4242,
};
let (client, _fake) = shep_client::testing::fake_client_with_ack(&addr, ack).await;
let mut out = Vec::new();
let mut err = Vec::new();
let mut streams = buffered_streams(&mut out, &mut err);
super::refuse_if_skewed(&mut streams, &client).expect("a matching version is not a skew");
}
}