use std::collections::BTreeSet;
use std::time::Duration;
use shep_client::shep_core::protocol::ProcessInfo;
use shep_client::shep_core::status::ProcStatus;
use tokio::time::{Instant, sleep};
use crate::daemon::Daemon;
use crate::error::Error;
use crate::state::Verify;
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct Generation {
pids: BTreeSet<u32>,
}
impl Generation {
pub async fn of<D: Daemon>(daemon: &D, sheep: &str) -> Result<Self, Error> {
let flock = daemon.describe(sheep).await?;
Ok(Self {
pids: flock.iter().filter_map(|info| info.pid).collect(),
})
}
#[must_use]
pub fn instances(&self) -> u32 {
u32::try_from(self.pids.len()).unwrap_or(u32::MAX)
}
pub(crate) fn of_infos(infos: &[&ProcessInfo]) -> Self {
Self {
pids: infos.iter().filter_map(|info| info.pid).collect(),
}
}
pub(crate) fn is_new(&self, info: &ProcessInfo) -> bool {
info.pid.is_some_and(|pid| !self.pids.contains(&pid))
}
pub(crate) fn holds(&self, info: &ProcessInfo) -> bool {
info.pid.is_some_and(|pid| self.pids.contains(&pid))
}
fn has_turned_over(&self, flock: &[ProcessInfo], accept: fn(&ProcessInfo) -> bool) -> bool {
!flock.is_empty() && flock.iter().all(|info| self.is_new(info) && accept(info))
}
}
pub(crate) const DWELL: Duration = Duration::from_secs(10);
pub(crate) const POLL: Duration = Duration::from_millis(100);
pub async fn wait<D: Daemon>(
daemon: &D,
sheep: &str,
mode: Verify,
before: &Generation,
budget: Duration,
) -> Result<bool, Error> {
let accept = match mode {
Verify::Probed => is_online,
Verify::Alive => is_alive,
};
let Some(settled) = turnover(daemon, sheep, before, accept, budget).await? else {
return Ok(false);
};
if mode == Verify::Probed {
return Ok(true);
}
sleep(DWELL).await;
let flock = daemon.describe(sheep).await?;
Ok(!flock.is_empty()
&& flock
.iter()
.all(|info| settled.holds(info) && is_alive(info)))
}
async fn turnover<D: Daemon>(
daemon: &D,
sheep: &str,
before: &Generation,
accept: fn(&ProcessInfo) -> bool,
budget: Duration,
) -> Result<Option<Generation>, Error> {
let deadline = Instant::now() + budget;
loop {
let flock = daemon.describe(sheep).await?;
if before.has_turned_over(&flock, accept) {
return Ok(Some(Generation {
pids: flock.iter().filter_map(|info| info.pid).collect(),
}));
}
if Instant::now() >= deadline {
return Ok(None);
}
sleep(POLL).await;
}
}
fn is_online(info: &ProcessInfo) -> bool {
info.status == ProcStatus::Online
}
pub(crate) fn is_alive(info: &ProcessInfo) -> bool {
matches!(info.status, ProcStatus::Starting | ProcStatus::Online)
}
#[cfg(test)]
mod tests {
use std::cell::Cell;
use shep_client::shep_core::config::AppConfig;
use super::*;
fn instance(id: u32, status: ProcStatus, pid: u32) -> ProcessInfo {
ProcessInfo::builder(id, "web", status)
.pid(Some(pid))
.build()
}
fn serving(pid: u32) -> Generation {
Generation {
pids: [pid].into_iter().collect(),
}
}
struct Listings {
sequence: Vec<Vec<ProcessInfo>>,
next: Cell<usize>,
}
impl Listings {
fn new(sequence: Vec<Vec<ProcessInfo>>) -> Self {
Self {
sequence,
next: Cell::new(0),
}
}
}
impl Daemon for Listings {
async fn dog_config(&self, _name: &str) -> Result<String, Error> {
unimplemented!()
}
async fn list_flock(&self) -> Result<Vec<ProcessInfo>, Error> {
unimplemented!()
}
async fn describe(&self, _sheep: &str) -> Result<Vec<ProcessInfo>, Error> {
let Some(last) = self.sequence.len().checked_sub(1) else {
return Ok(Vec::new());
};
let index = self.next.get().min(last);
self.next.set((index + 1).min(last));
Ok(self.sequence[index].clone())
}
async fn start(&self, _apps: Vec<AppConfig>) -> Result<(), Error> {
unimplemented!()
}
async fn delete(&self, _id: u32) -> Result<(), Error> {
unimplemented!()
}
async fn reload(&self, _sheep: &str) -> Result<(), Error> {
unimplemented!()
}
async fn restart(&self, _sheep: &str) -> Result<(), Error> {
unimplemented!()
}
async fn save_roll(&self) -> Result<std::path::PathBuf, Error> {
unimplemented!()
}
async fn set_smit(&self, _sheep: &str, _text: &str) -> Result<(), Error> {
unimplemented!()
}
}
#[tokio::test(start_paused = true)]
async fn the_old_instance_still_online_is_not_a_new_release() {
let daemon = Listings::new(vec![vec![instance(1, ProcStatus::Online, 12835)]]);
let ok = wait(
&daemon,
"web",
Verify::Probed,
&serving(12835),
Duration::from_millis(50),
)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test(start_paused = true)]
async fn a_reload_that_has_not_happened_yet_is_not_success() {
let daemon = Listings::new(vec![
vec![instance(1, ProcStatus::Online, 12835)],
vec![instance(1, ProcStatus::Online, 12835)],
]);
let ok = wait(
&daemon,
"web",
Verify::Probed,
&serving(12835),
Duration::from_millis(50),
)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test(start_paused = true)]
async fn a_new_instance_reaching_online_is_success() {
let daemon = Listings::new(vec![
vec![instance(1, ProcStatus::Online, 12835)],
vec![
instance(1, ProcStatus::Stopping, 12835),
instance(2, ProcStatus::Starting, 13002),
],
vec![instance(2, ProcStatus::Online, 13002)],
]);
assert!(
wait(
&daemon,
"web",
Verify::Probed,
&serving(12835),
Duration::from_secs(5)
)
.await
.unwrap()
);
}
#[tokio::test(start_paused = true)]
async fn a_draining_old_instance_is_not_a_finished_turnover() {
let daemon = Listings::new(vec![vec![
instance(1, ProcStatus::Stopping, 12835),
instance(2, ProcStatus::Online, 13002),
]]);
let ok = wait(
&daemon,
"web",
Verify::Probed,
&serving(12835),
Duration::from_millis(50),
)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test(start_paused = true)]
async fn a_new_instance_that_never_leaves_starting_is_not_success() {
let daemon = Listings::new(vec![vec![instance(2, ProcStatus::Starting, 13002)]]);
let ok = wait(
&daemon,
"web",
Verify::Probed,
&serving(12835),
Duration::from_millis(50),
)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test(start_paused = true)]
async fn alive_accepts_a_new_process_that_is_still_running() {
let daemon = Listings::new(vec![vec![instance(2, ProcStatus::Starting, 13002)]]);
assert!(
wait(
&daemon,
"web",
Verify::Alive,
&serving(12835),
Duration::from_millis(50)
)
.await
.unwrap()
);
}
#[tokio::test(start_paused = true)]
async fn alive_rejects_the_old_process_still_running() {
let daemon = Listings::new(vec![vec![instance(1, ProcStatus::Online, 12835)]]);
let ok = wait(
&daemon,
"web",
Verify::Alive,
&serving(12835),
Duration::from_millis(50),
)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test(start_paused = true)]
async fn alive_rejects_an_instance_waiting_to_be_restarted() {
let daemon = Listings::new(vec![vec![instance(2, ProcStatus::WaitingRestart, 13002)]]);
let ok = wait(
&daemon,
"web",
Verify::Alive,
&serving(12835),
Duration::from_millis(50),
)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test(start_paused = true)]
async fn alive_waits_out_a_reload_that_is_still_in_flight() {
let daemon = Listings::new(vec![
vec![instance(1, ProcStatus::Online, 12835)],
vec![instance(1, ProcStatus::Online, 12835)],
vec![instance(1, ProcStatus::Online, 12835)],
vec![instance(2, ProcStatus::Starting, 13002)],
]);
assert!(
wait(
&daemon,
"web",
Verify::Alive,
&serving(12835),
Duration::from_secs(120)
)
.await
.unwrap()
);
}
#[tokio::test(start_paused = true)]
async fn alive_rejects_a_process_that_dies_during_the_dwell() {
let daemon = Listings::new(vec![
vec![instance(2, ProcStatus::Starting, 13002)],
vec![instance(2, ProcStatus::Starting, 13456)],
]);
let ok = wait(
&daemon,
"web",
Verify::Alive,
&serving(12835),
Duration::from_secs(120),
)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test(start_paused = true)]
async fn alive_rejects_a_flock_that_empties_during_the_dwell() {
let daemon = Listings::new(vec![
vec![instance(2, ProcStatus::Starting, 13002)],
Vec::new(),
]);
let ok = wait(
&daemon,
"web",
Verify::Alive,
&serving(12835),
Duration::from_secs(120),
)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test(start_paused = true)]
async fn an_instance_with_no_pid_is_not_a_new_generation() {
let daemon = Listings::new(vec![vec![
ProcessInfo::builder(2, "web", ProcStatus::Online).build(),
]]);
let ok = wait(
&daemon,
"web",
Verify::Probed,
&serving(12835),
Duration::from_millis(50),
)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test(start_paused = true)]
async fn an_empty_describe_is_failure_in_both_modes() {
let empty = Listings::new(Vec::new());
assert!(
!wait(
&empty,
"web",
Verify::Probed,
&serving(12835),
Duration::from_millis(50)
)
.await
.unwrap()
);
let empty = Listings::new(Vec::new());
assert!(
!wait(
&empty,
"web",
Verify::Alive,
&serving(12835),
Duration::from_millis(50)
)
.await
.unwrap()
);
}
#[tokio::test(start_paused = true)]
async fn the_fixture_repeats_its_last_listing_past_exhaustion() {
let daemon = Listings::new(vec![
vec![instance(1, ProcStatus::Online, 12835)],
vec![instance(1, ProcStatus::Online, 12835)],
]);
let ok = wait(
&daemon,
"web",
Verify::Probed,
&serving(12835),
Duration::from_millis(350),
)
.await
.unwrap();
assert!(!ok);
}
#[tokio::test(start_paused = true)]
async fn a_generation_is_the_pids_that_are_serving() {
let daemon = Listings::new(vec![vec![
instance(1, ProcStatus::Online, 12835),
instance(2, ProcStatus::Online, 12836),
ProcessInfo::builder(3, "web", ProcStatus::WaitingRestart).build(),
]]);
let generation = Generation::of(&daemon, "web").await.unwrap();
assert_eq!(
generation,
Generation {
pids: [12835, 12836].into_iter().collect()
}
);
}
}