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 },
}
pub async fn run(streams: &mut Streams<'_>, quiet: bool, options: ForegroundOptions) -> ExitCode {
let ForegroundOptions {
paths,
apps,
tidy_up,
} = options;
let daemon_args = DaemonArgs {
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) => 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
}
}
}