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 {
let count_holds = self.instances() == 0
|| u32::try_from(flock.len()).is_ok_and(|listed| listed == self.instances());
!flock.is_empty()
&& count_holds
&& 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 deadline = Instant::now() + DWELL;
loop {
let flock = describe_within(daemon, sheep, DWELL).await?;
let settled_and_ready = u32::try_from(flock.len()).is_ok_and(|n| n == settled.instances())
&& flock
.iter()
.all(|info| settled.holds(info) && is_online(info));
if settled_and_ready {
return Ok(true);
}
if flock.is_empty()
|| !flock
.iter()
.all(|info| settled.holds(info) && is_alive(info))
{
return Ok(false);
}
if Instant::now() >= deadline {
return Ok(false);
}
sleep(POLL).await;
}
}
async fn describe_within<D: Daemon>(
daemon: &D,
sheep: &str,
budget: Duration,
) -> Result<Vec<ProcessInfo>, Error> {
let deadline = Instant::now() + budget;
loop {
match daemon.describe(sheep).await {
Ok(flock) => return Ok(flock),
Err(err) if !err.is_retryable() => return Err(err),
Err(err) => {
let left = deadline.saturating_duration_since(Instant::now());
if left.is_zero() {
return Err(err);
}
sleep(POLL.min(left)).await;
}
}
}
}
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 = describe_within(
daemon,
sheep,
deadline.saturating_duration_since(Instant::now()),
)
.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::RequestError;
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),
}
}
}
struct Blips {
calls: Cell<u32>,
after: Vec<ProcessInfo>,
}
impl Daemon for Blips {
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 n = self.calls.get();
self.calls.set(n + 1);
if n == 0 {
return Err(Error::Request(RequestError::Timeout {
after: Duration::from_secs(1),
}));
}
Ok(self.after.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]
async fn a_transient_describe_failure_does_not_fail_verification() {
let before = Generation {
pids: [111].into_iter().collect(),
};
let daemon = Blips {
calls: Cell::new(0),
after: vec![instance(0, ProcStatus::Online, 222)],
};
let seen = turnover(&daemon, "web", &before, is_online, Duration::from_secs(10))
.await
.expect("a retryable blip must not end verification");
assert!(
seen.is_some(),
"the turnover after the blip must still be seen"
);
assert!(daemon.calls.get() >= 2, "it must have asked again");
}
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)],
vec![instance(2, ProcStatus::Online, 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)],
vec![instance(2, ProcStatus::Online, 13002)],
]);
assert!(
wait(
&daemon,
"web",
Verify::Alive,
&serving(12835),
Duration::from_secs(120)
)
.await
.unwrap()
);
}
#[tokio::test(start_paused = true)]
async fn a_turnover_that_lost_a_replica_is_not_a_turnover() {
let daemon = Listings::new(vec![vec![instance(1, ProcStatus::Online, 200)]]);
let before = Generation {
pids: [100, 101].into_iter().collect(),
};
let ok = wait(
&daemon,
"web",
Verify::Probed,
&before,
Duration::from_millis(50),
)
.await
.unwrap();
assert!(!ok, "one replacement for two instances is not a turnover");
}
#[tokio::test(start_paused = true)]
async fn alive_rejects_a_flock_that_came_back_smaller() {
let daemon = Listings::new(vec![
vec![
instance(1, ProcStatus::Online, 200),
instance(2, ProcStatus::Online, 201),
],
vec![instance(1, ProcStatus::Online, 200)],
]);
let before = Generation {
pids: [100, 101].into_iter().collect(),
};
let ok = wait(
&daemon,
"web",
Verify::Alive,
&before,
Duration::from_secs(120),
)
.await
.unwrap();
assert!(!ok, "half a flock is not a deployed release");
}
#[tokio::test(start_paused = true)]
async fn alive_rejects_a_replacement_shep_gave_up_on() {
let daemon = Listings::new(vec![vec![instance(2, ProcStatus::Starting, 13002)]]);
let ok = wait(
&daemon,
"web",
Verify::Alive,
&serving(12835),
Duration::from_secs(120),
)
.await
.unwrap();
assert!(
!ok,
"a replacement left at `starting` never came up, whatever verify mode asked"
);
}
#[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_one_replica_of_two_restarting_during_the_dwell() {
let daemon = Listings::new(vec![
vec![
instance(1, ProcStatus::Starting, 200),
instance(2, ProcStatus::Starting, 201),
],
vec![
instance(1, ProcStatus::Starting, 200),
instance(2, ProcStatus::Starting, 202),
],
]);
let before = Generation {
pids: [100, 101].into_iter().collect(),
};
let ok = wait(
&daemon,
"web",
Verify::Alive,
&before,
Duration::from_secs(120),
)
.await
.unwrap();
assert!(
!ok,
"every instance must still be the one that came up, not merely one of them"
);
}
#[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()
}
);
}
}