use std::path::{Path, PathBuf};
use std::time::Duration;
use shep_client::shep_core::config::AppConfig;
use shep_client::shep_core::protocol::ProcessInfo;
use tokio::time::{Instant, sleep};
use crate::build;
use crate::config::DogConfig;
use crate::daemon::Daemon;
use crate::deploy::RELOAD_DEADLINE_SLACK;
use crate::error::Error;
use crate::flockfile;
use crate::git;
use crate::lock;
use crate::paths::Tree;
use crate::roll;
use crate::shared;
use crate::state::{State, Verify, Watch};
use crate::swap;
use crate::verify::{DWELL, Generation, POLL, is_alive};
#[derive(Debug, Clone)]
pub struct Prepared {
pub tree: Tree,
pub state: State,
pub sha: String,
pub app: AppConfig,
pub previous_config: AppConfig,
}
pub async fn prepare<D: Daemon>(
daemon: &D,
shep_home: &Path,
sheep: &str,
config: &DogConfig,
) -> Result<Prepared, Error> {
let tree = Tree::for_sheep(shep_home, sheep);
if tree.state_file().is_file() {
let state = State::read(&tree.state_file()).map_err(|source| {
Error::Config(format!(
"{sheep} has a deploy tree at {} and its record at {} cannot be read. \
Refusing rather than guessing. That record is the only thing that says \
whether {sheep} was ever cut over, and the two answers need opposite \
handling. Do NOT remove the tree until you know which it is: if {sheep} WAS \
cut over, its working directory is inside it and removing it takes a running \
service's cwd with it. Restore or repair that file first - `shep describe \
{sheep}` shows whether the sheep is running from inside this tree. Reading it \
failed with: {source}",
tree.root().display(),
tree.state_file().display()
))
})?;
return Err(Error::Config(if state.deployed.is_some() {
format!(
"{sheep} is already a deploy target: its tree is at {}. Deploy it with \
`shep deploy {sheep}`, or change how it is watched with \
`shep deploy {sheep} --watch auto|manual`.",
tree.releases().display()
)
} else {
format!(
"{sheep} has a deploy tree at {} but was never cut over to it: its record names \
no deployed release, so nothing has ever been served from that tree. An \
abandoned first cutover leaves exactly this. Do NOT run `shep deploy {sheep}` \
against it - that would reload the sheep at its own checkout and report success \
for a release nothing served. Remove {} and run `shep-deploy setup {sheep}` \
again once the cause of the first failure is fixed.",
tree.releases().display(),
tree.root().display()
)
}));
}
let registered = roll::registered(daemon).await?;
let previous_config = registered.get(sheep).cloned().ok_or_else(|| {
Error::Config(format!(
"the shepherd has no sheep named {sheep:?} registered, so there is nothing to take \
over: `shep-deploy survey` lists every registered sheep and where it stands"
))
})?;
let checkout = PathBuf::from(previous_config.cwd.as_deref().ok_or_else(|| {
Error::Config(format!(
"shep records no working directory for {sheep}, so there is no checkout to deploy \
from"
))
})?);
let remote = git::remote_url(&checkout)?;
let branch = git::current_branch(&checkout)?;
let state = State {
remote,
branch,
deployed: None,
failed: None,
verify: Verify::default(),
watch: Watch::Manual,
origin_cwd: Some(checkout.clone()),
origin_script: Some(previous_config.script.clone()),
checkout,
};
let _deploying = lock::hold(&tree)?;
std::fs::create_dir_all(tree.releases()).map_err(|source| Error::Io {
path: tree.releases(),
source,
})?;
git::init_bare(&tree.git())?;
git::fetch(&tree.git(), &state.remote, config.git_timeout)?;
let sha = git::remote_head(&tree.git(), &state.branch)?;
let release = tree.release(&sha);
crate::deploy::checkout_release(&tree, &sha)?;
shared::link_cache(&release, &tree.cache_target())?;
let shared_paths = shared::to_link(&state.checkout)?;
shared::link_into(&release, &state.checkout, &shared_paths)?;
let app = flockfile::app_config(&release, sheep, &shared_paths)?;
let spec = flockfile::build_spec(&release, &shared_paths)?;
build::run(
sheep,
&release,
&spec,
app.user.as_deref(),
&config.passthrough,
&tree.cache_target(),
config.build_timeout,
)
.await?;
swap::point_at(&tree.current(), &release)?;
state.write(&tree.state_file())?;
Ok(Prepared {
tree,
state,
sha,
app,
previous_config,
})
}
pub async fn cut_over<D: Daemon>(daemon: &D, prepared: Prepared) -> Result<String, Error> {
let Prepared {
tree,
mut state,
sha,
app,
previous_config,
} = prepared;
let sheep = tree.sheep().to_owned();
let before = daemon.describe(&sheep).await?;
let rows: Vec<&ProcessInfo> = before.iter().collect();
let previous: Vec<u32> = rows.iter().map(|info| info.id).collect();
let generation = Generation::of_infos(&rows);
let mut registering = app;
registering.cwd = Some(tree.current().display().to_string());
match attempt(daemon, &sheep, registering, &generation).await {
CutOver::Done => {
let mut stranded = Vec::new();
for id in previous {
if daemon.delete(id).await.is_err() {
stranded.push(id);
}
}
state.deployed = Some(sha.clone());
if stranded.is_empty() {
state.watch = Watch::Auto;
}
state.write(&tree.state_file())?;
if stranded.is_empty() {
Ok(sha)
} else {
Err(Error::Stranded {
sheep,
sha,
ids: stranded,
})
}
}
CutOver::NotStarted(source) => Err(source),
CutOver::NotVerified(why) => {
let undone = undo_start(daemon, &sheep, &previous, previous_config).await;
Err(Error::CutOver {
sheep,
why,
removed: undone.removed,
repaired: undone.recorded,
tree: tree.root().to_owned(),
source: None,
})
}
CutOver::Failed(source) => {
let undone = undo_start(daemon, &sheep, &previous, previous_config).await;
Err(Error::CutOver {
sheep,
why: SHEPHERD_QUIET.to_owned(),
removed: undone.removed,
repaired: undone.recorded,
tree: tree.root().to_owned(),
source: Some(Box::new(source)),
})
}
}
}
enum CutOver {
Done,
NotStarted(Error),
NotVerified(String),
Failed(Error),
}
const PORT_COLLISION: &str = "The first cutover is the one deploy that runs two instances at \
once, so an app that does not bind with SO_REUSEPORT cannot take its own port while the \
original still holds it. Every deploy after the first replaces the instance rather than \
joining it, and does not meet this. If the app cannot set SO_REUSEPORT, `shep stop` it, \
remove the tree named above, and run setup again: with the port free the newcomer binds, \
and the cutover is the one deploy allowed to be down for a moment anyway.";
const SHEPHERD_QUIET: &str = "The shepherd stopped answering while the new instance was being \
watched, so nothing was established about it either way.";
fn cutover_budget(app: &AppConfig) -> Duration {
app.listen_timeout.as_duration() + RELOAD_DEADLINE_SLACK
}
async fn attempt<D: Daemon>(
daemon: &D,
sheep: &str,
app: AppConfig,
before: &Generation,
) -> CutOver {
let patience = cutover_budget(&app);
if let Err(source) = daemon.start(vec![app]).await {
return CutOver::NotStarted(source);
}
let started_at = Instant::now();
let deadline = started_at + patience;
let arrived = loop {
let flock = match daemon.describe(sheep).await {
Ok(flock) => flock,
Err(source) => return CutOver::Failed(source),
};
let newcomers: Vec<&ProcessInfo> =
flock.iter().filter(|info| before.is_new(info)).collect();
if !newcomers.is_empty() && newcomers.iter().all(|info| is_alive(info)) {
break Generation::of_infos(&newcomers);
}
if newcomers.iter().any(|info| !is_alive(info)) {
return CutOver::NotVerified(format!(
"The new instance failed before it finished starting. {PORT_COLLISION}"
));
}
if Instant::now() >= deadline {
return CutOver::NotVerified(format!(
"No new instance appeared within {}s, although the shepherd accepted the start.",
started_at.elapsed().as_secs()
));
}
sleep(POLL).await;
};
sleep(DWELL).await;
let flock = match daemon.describe(sheep).await {
Ok(flock) => flock,
Err(source) => return CutOver::Failed(source),
};
let survivors: Vec<&ProcessInfo> = flock
.iter()
.filter(|info| arrived.holds(info) && is_alive(info))
.collect();
if survivors.len() != arrived.instances() as usize {
return CutOver::NotVerified(format!(
"The new instance did not stay up for {}s after starting. {PORT_COLLISION}",
DWELL.as_secs()
));
}
if survivors.iter().any(|info| info.restarts > 0) {
return CutOver::NotVerified(format!(
"The instance this cutover adopted had already restarted {}s later, so what is \
running is not the process the start spawned. {PORT_COLLISION}",
DWELL.as_secs()
));
}
CutOver::Done
}
async fn undo_start<D: Daemon>(
daemon: &D,
sheep: &str,
previous: &[u32],
original: AppConfig,
) -> Undone {
let Ok(flock) = daemon.describe(sheep).await else {
return Undone {
removed: false,
recorded: false,
};
};
let mut removed = drain(daemon, &flock, previous).await;
if daemon.start(vec![original]).await.is_err() {
return Undone {
removed,
recorded: false,
};
}
let Ok(flock) = daemon.describe(sheep).await else {
return Undone {
removed,
recorded: false,
};
};
removed &= drain(daemon, &flock, previous).await;
Undone {
removed,
recorded: true,
}
}
struct Undone {
removed: bool,
recorded: bool,
}
async fn drain<D: Daemon>(daemon: &D, flock: &[ProcessInfo], previous: &[u32]) -> bool {
let mut all = true;
for info in flock.iter().filter(|info| !previous.contains(&info.id)) {
if daemon.delete(info.id).await.is_err() {
all = false;
}
}
all
}
#[cfg(test)]
mod tests {
use crate::fixtures;
fn test_config() -> crate::config::DogConfig {
crate::config::DogConfig {
interval: std::time::Duration::from_secs(30),
retention: 5,
git_timeout: std::time::Duration::from_secs(60),
build_timeout: std::time::Duration::from_secs(60),
passthrough: Vec::new(),
}
}
use core::error::Error as _;
use std::cell::{Cell, RefCell};
use std::time::Duration;
use shep_client::RequestError;
use shep_client::shep_core::protocol::{ProcessInfo, RpcError, RpcErrorCode};
use shep_client::shep_core::status::ProcStatus;
use tokio::time::Instant;
use super::*;
struct RollOf<'a>(&'a [(&'a str, &'a Path)]);
impl Daemon for RollOf<'_> {
async fn dog_config(&self, _name: &str) -> Result<String, Error> {
unimplemented!()
}
async fn list_flock(
&self,
) -> Result<Vec<shep_client::shep_core::protocol::ProcessInfo>, Error> {
unimplemented!()
}
async fn describe(
&self,
_sheep: &str,
) -> Result<Vec<shep_client::shep_core::protocol::ProcessInfo>, Error> {
unimplemented!()
}
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<PathBuf, Error> {
let dir = tempfile::tempdir().expect("tempdir");
let apps: Vec<String> = self
.0
.iter()
.map(|(name, cwd)| {
format!(
"{{\"app\":{{\"name\":{name:?},\"script\":\"./run.sh\",\"cwd\":{:?}}}}}",
cwd.to_str().expect("utf-8 cwd")
)
})
.collect();
let path = dir.keep().join("flock.json");
std::fs::write(&path, format!("{{\"apps\":[{}]}}", apps.join(",")))
.expect("write roll");
Ok(path)
}
async fn set_smit(&self, _sheep: &str, _text: &str) -> Result<(), Error> {
unimplemented!()
}
}
fn checkout_with_commit() -> tempfile::TempDir {
let dir = tempfile::tempdir().expect("tempdir");
fixtures::run_git(dir.path(), &["init", "-q", "-b", "main"]);
fixtures::run_git(dir.path(), &["config", "user.email", "test@example.com"]);
fixtures::run_git(dir.path(), &["config", "user.name", "test"]);
fixtures::run_git(
dir.path(),
&[
"remote",
"add",
"origin",
dir.path().to_str().expect("utf-8 tempdir path"),
],
);
std::fs::write(
dir.path().join("Flockfile.toml"),
"[[app]]\nname = \"bpm\"\nscript = \"./run.sh\"\n",
)
.expect("write Flockfile");
std::fs::write(dir.path().join("run.sh"), "#!/bin/sh\necho hi\n").expect("write run.sh");
fixtures::run_git(dir.path(), &["add", "."]);
fixtures::run_git(dir.path(), &["commit", "-q", "-m", "initial"]);
dir
}
fn write_target(home: &Path, sheep: &str, watch: Watch, sha: Option<&str>) {
let tree = Tree::for_sheep(home, sheep);
std::fs::create_dir_all(tree.state_file().parent().expect("has a parent"))
.expect("create target dir");
let state = State {
remote: "https://example.com/x".to_owned(),
branch: "main".to_owned(),
deployed: sha.map(str::to_owned),
failed: None,
verify: Verify::default(),
watch,
origin_cwd: None,
origin_script: None,
checkout: PathBuf::from("/srv/x"),
};
state.write(&tree.state_file()).expect("write state");
}
#[tokio::test]
async fn opting_in_twice_is_refused_by_name() {
let home = tempfile::tempdir().expect("tempdir");
let checkout = checkout_with_commit();
write_target(home.path(), "bpm", Watch::Auto, Some("old"));
let entries = [("bpm", checkout.path())];
let daemon = RollOf(&entries);
let err = prepare(&daemon, home.path(), "bpm", &test_config())
.await
.expect_err("refuses");
let shown = err.to_string();
assert!(shown.contains("bpm"), "{shown}");
assert!(shown.contains("already"), "{shown}");
}
#[tokio::test]
async fn the_pre_adoption_cwd_and_script_are_recorded() {
let home = tempfile::tempdir().expect("tempdir");
let checkout = checkout_with_commit();
let entries = [("bpm", checkout.path())];
let daemon = RollOf(&entries);
let prepared = prepare(&daemon, home.path(), "bpm", &test_config())
.await
.expect("prepares");
assert_eq!(prepared.state.origin_cwd.as_deref(), Some(checkout.path()));
assert_eq!(prepared.state.origin_script.as_deref(), Some("./run.sh"));
assert_eq!(prepared.state.checkout, checkout.path());
}
#[tokio::test]
async fn the_branch_comes_from_the_checkouts_own_head() {
let home = tempfile::tempdir().expect("tempdir");
let checkout = checkout_with_commit();
fixtures::run_git(checkout.path(), &["checkout", "-q", "-b", "stable"]);
let entries = [("bpm", checkout.path())];
let daemon = RollOf(&entries);
let prepared = prepare(&daemon, home.path(), "bpm", &test_config())
.await
.expect("prepares");
assert_eq!(prepared.state.branch, "stable");
}
#[tokio::test]
async fn a_prepare_that_died_before_writing_its_record_runs_again() {
let home = tempfile::tempdir().expect("tempdir");
let checkout = checkout_with_commit();
let entries = [("bpm", checkout.path())];
let daemon = RollOf(&entries);
let first = prepare(&daemon, home.path(), "bpm", &test_config())
.await
.expect("prepares");
std::fs::remove_file(first.tree.state_file()).expect("the record a kill never wrote");
let again = prepare(&daemon, home.path(), "bpm", &test_config())
.await
.expect("a tree with no record must be resumable, as the doc says");
assert_eq!(again.sha, first.sha);
assert!(
again.tree.state_file().is_file(),
"the second run must leave the record the first one never wrote"
);
}
#[tokio::test]
async fn current_ends_up_on_a_release_carrying_the_shared_files() {
let home = tempfile::tempdir().expect("tempdir");
let checkout = checkout_with_commit();
std::fs::write(checkout.path().join(".gitignore"), "local.json\n").expect("write");
std::fs::write(checkout.path().join("local.json"), "{}").expect("write");
let entries = [("bpm", checkout.path())];
let daemon = RollOf(&entries);
let prepared = prepare(&daemon, home.path(), "bpm", &test_config())
.await
.expect("prepares");
let current = swap::resolve(&prepared.tree.current())
.expect("reads")
.expect("current is set");
assert_eq!(current, prepared.tree.release(&prepared.sha));
assert_eq!(
std::fs::read_to_string(current.join("local.json")).expect("reads through the link"),
"{}"
);
}
#[tokio::test]
async fn a_detached_checkout_is_refused() {
let home = tempfile::tempdir().expect("tempdir");
let checkout = checkout_with_commit();
let head = fixtures::head_of(checkout.path());
fixtures::run_git(checkout.path(), &["checkout", "-q", &head]);
let entries = [("bpm", checkout.path())];
let daemon = RollOf(&entries);
let err = prepare(&daemon, home.path(), "bpm", &test_config())
.await
.expect_err("refuses");
assert!(err.to_string().contains("detached"), "{err}");
}
const ORIGINAL_PID: u32 = 100;
const NEWCOMER_ID: u32 = 99;
const NEWCOMER_PID: u32 = 200;
const REPAIR_ID: u32 = 500;
const REPAIR_PID: u32 = 700;
#[derive(Debug, Clone, Copy)]
enum Shape {
Absent,
Up {
status: ProcStatus,
pid: u32,
restarts: u32,
},
}
fn shape_at(script: &[(Duration, Shape)], elapsed: Duration) -> Shape {
script
.iter()
.rev()
.find(|(from, _)| *from <= elapsed)
.map_or(Shape::Absent, |(_, shape)| *shape)
}
fn originals_stay_online() -> Vec<(Duration, Shape)> {
vec![(
Duration::ZERO,
Shape::Up {
status: ProcStatus::Online,
pid: ORIGINAL_PID,
restarts: 0,
},
)]
}
fn originals_respawn_mid_cutover() -> Vec<(Duration, Shape)> {
vec![
(
Duration::ZERO,
Shape::Up {
status: ProcStatus::Online,
pid: ORIGINAL_PID,
restarts: 0,
},
),
(
Duration::from_secs(1),
Shape::Up {
status: ProcStatus::Online,
pid: ORIGINAL_PID + 50,
restarts: 1,
},
),
]
}
struct CutOverDouble {
originals: Vec<u32>,
original_script: Vec<(Duration, Shape)>,
script: Vec<(Duration, Shape)>,
refuses: Option<usize>,
refuses_deletes_from: Option<usize>,
deletes_seen: Cell<usize>,
mute_after: Option<Duration>,
starts: RefCell<Vec<AppConfig>>,
deletes: RefCell<Vec<u32>>,
attempts: Cell<usize>,
accepted_at: Cell<Option<Instant>>,
repairs: Cell<u32>,
}
impl CutOverDouble {
fn new(originals: &[u32], script: Vec<(Duration, Shape)>, refuses: Option<usize>) -> Self {
Self {
originals: originals.to_vec(),
original_script: originals_stay_online(),
script,
refuses,
refuses_deletes_from: None,
deletes_seen: Cell::new(0),
mute_after: None,
starts: RefCell::new(Vec::new()),
deletes: RefCell::new(Vec::new()),
attempts: Cell::new(0),
accepted_at: Cell::new(None),
repairs: Cell::new(0),
}
}
fn going_quiet_after(mut self, after: Duration) -> Self {
self.mute_after = Some(after);
self
}
fn refusing_deletes(self) -> Self {
self.refusing_deletes_from(0)
}
fn refusing_deletes_from(mut self, from: usize) -> Self {
self.refuses_deletes_from = Some(from);
self
}
fn while_the_originals(mut self, script: Vec<(Duration, Shape)>) -> Self {
self.original_script = script;
self
}
fn started(&self) -> Vec<AppConfig> {
self.starts.borrow().clone()
}
fn deleted(&self) -> Vec<u32> {
self.deletes.borrow().clone()
}
fn is_deleted(&self, id: u32) -> bool {
self.deletes.borrow().contains(&id)
}
fn row(&self, id: u32, status: ProcStatus, pid: u32, restarts: u32) -> ProcessInfo {
let _ = self;
ProcessInfo::builder(id, "bpm", status)
.pid(Some(pid))
.restarts(restarts)
.build()
}
}
impl Daemon for CutOverDouble {
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 mut flock = Vec::new();
let elapsed = self
.accepted_at
.get()
.map_or(Duration::ZERO, |accepted_at| Instant::now() - accepted_at);
if self.accepted_at.get().is_some()
&& self.mute_after.is_some_and(|after| elapsed >= after)
{
return Err(Error::Request(RequestError::Rpc(RpcError {
code: RpcErrorCode::Internal,
message: "the shepherd is not answering".to_owned(),
})));
}
if let Shape::Up {
status,
pid,
restarts,
} = shape_at(&self.original_script, elapsed)
{
for (offset, id) in self.originals.iter().enumerate() {
let offset = u32::try_from(offset).expect("a handful of instances");
if !self.is_deleted(*id) {
flock.push(self.row(*id, status, pid + offset, restarts));
}
}
}
if self.accepted_at.get().is_some()
&& let Shape::Up {
status,
pid,
restarts,
} = shape_at(&self.script, elapsed)
{
for offset in 0..self.originals.len() {
let offset = u32::try_from(offset).expect("a handful of instances");
let id = NEWCOMER_ID - offset;
if !self.is_deleted(id) {
flock.push(self.row(id, status, pid + offset, restarts));
}
}
}
for offset in 0..self.repairs.get() {
let id = REPAIR_ID + offset;
if !self.is_deleted(id) {
flock.push(self.row(id, ProcStatus::Online, REPAIR_PID + offset, 0));
}
}
Ok(flock)
}
async fn start(&self, apps: Vec<AppConfig>) -> Result<(), Error> {
let attempt = self.attempts.get();
self.attempts.set(attempt + 1);
if self.refuses == Some(attempt) {
return Err(Error::Request(RequestError::Rpc(RpcError {
code: RpcErrorCode::Internal,
message: "bpm cannot be started".to_owned(),
})));
}
self.starts.borrow_mut().extend(apps);
if attempt == 0 {
self.accepted_at.set(Some(Instant::now()));
} else {
self.repairs.set(self.repairs.get() + 1);
}
Ok(())
}
async fn delete(&self, id: u32) -> Result<(), Error> {
let seen = self.deletes_seen.get();
self.deletes_seen.set(seen + 1);
if self.refuses_deletes_from.is_some_and(|from| seen >= from) {
return Err(Error::Request(RequestError::Rpc(RpcError {
code: RpcErrorCode::Internal,
message: format!("instance {id} cannot be deleted"),
})));
}
self.deletes.borrow_mut().push(id);
Ok(())
}
async fn reload(&self, _sheep: &str) -> Result<(), Error> {
unimplemented!()
}
async fn restart(&self, _sheep: &str) -> Result<(), Error> {
unimplemented!()
}
async fn save_roll(&self) -> Result<PathBuf, Error> {
unimplemented!()
}
async fn set_smit(&self, _sheep: &str, _text: &str) -> Result<(), Error> {
unimplemented!()
}
}
struct Dirs {
_home: tempfile::TempDir,
_checkout: tempfile::TempDir,
}
async fn cutover_fixture_of(
originals: &[u32],
script: Vec<(Duration, Shape)>,
refuses: Option<usize>,
) -> (CutOverDouble, Prepared, Dirs) {
let home = tempfile::tempdir().expect("tempdir");
let checkout = checkout_with_commit();
let entries = [("bpm", checkout.path())];
let roll = RollOf(&entries);
let prepared = prepare(&roll, home.path(), "bpm", &test_config())
.await
.expect("prepares");
(
CutOverDouble::new(originals, script, refuses),
prepared,
Dirs {
_home: home,
_checkout: checkout,
},
)
}
fn comes_up() -> Vec<(Duration, Shape)> {
vec![
(Duration::ZERO, Shape::Absent),
(
Duration::from_millis(250),
Shape::Up {
status: ProcStatus::Starting,
pid: NEWCOMER_PID,
restarts: 0,
},
),
(
Duration::from_millis(450),
Shape::Up {
status: ProcStatus::Online,
pid: NEWCOMER_PID,
restarts: 0,
},
),
]
}
async fn cutover_fixture() -> (CutOverDouble, Prepared, Dirs) {
cutover_fixture_of(&[7], comes_up(), None).await
}
async fn cutover_fixture_with_instances(originals: &[u32]) -> (CutOverDouble, Prepared, Dirs) {
cutover_fixture_of(originals, comes_up(), None).await
}
async fn cutover_fixture_online_then_gone() -> (CutOverDouble, Prepared, Dirs) {
cutover_fixture_of(
&[7],
vec![
(
Duration::ZERO,
Shape::Up {
status: ProcStatus::Online,
pid: NEWCOMER_PID,
restarts: 0,
},
),
(Duration::from_secs(5), Shape::Absent),
],
None,
)
.await
}
async fn cutover_fixture_crash_looping() -> (CutOverDouble, Prepared, Dirs) {
cutover_fixture_of(
&[7],
vec![
(Duration::ZERO, Shape::Absent),
(
Duration::from_millis(250),
Shape::Up {
status: ProcStatus::Online,
pid: NEWCOMER_PID,
restarts: 0,
},
),
(
Duration::from_secs(5),
Shape::Up {
status: ProcStatus::Online,
pid: NEWCOMER_PID + 50,
restarts: 1,
},
),
],
None,
)
.await
}
fn dies_during_dwell() -> Vec<(Duration, Shape)> {
vec![
(Duration::ZERO, Shape::Absent),
(
Duration::from_millis(250),
Shape::Up {
status: ProcStatus::Online,
pid: NEWCOMER_PID,
restarts: 0,
},
),
(
Duration::from_secs(5),
Shape::Up {
status: ProcStatus::Errored,
pid: NEWCOMER_PID,
restarts: 0,
},
),
]
}
async fn cutover_fixture_dies_during_dwell() -> (CutOverDouble, Prepared, Dirs) {
cutover_fixture_of(&[7], dies_during_dwell(), None).await
}
async fn cutover_fixture_dies_and_refuses_repair() -> (CutOverDouble, Prepared, Dirs) {
cutover_fixture_of(&[7], dies_during_dwell(), Some(1)).await
}
async fn cutover_fixture_never_appears() -> (CutOverDouble, Prepared, Dirs) {
cutover_fixture_of(&[7], vec![(Duration::ZERO, Shape::Absent)], None).await
}
async fn cutover_fixture_comes_up_but_refuses_deletes() -> (CutOverDouble, Prepared, Dirs) {
let (daemon, prepared, dirs) = cutover_fixture_of(&[7, 8], comes_up(), None).await;
(daemon.refusing_deletes(), prepared, dirs)
}
async fn cutover_fixture_shepherd_goes_quiet() -> (CutOverDouble, Prepared, Dirs) {
let (daemon, prepared, dirs) = cutover_fixture_of(&[7], comes_up(), None).await;
(daemon.going_quiet_after(Duration::ZERO), prepared, dirs)
}
async fn cutover_fixture_dies_and_refuses_the_second_delete() -> (CutOverDouble, Prepared, Dirs)
{
let (daemon, prepared, dirs) = cutover_fixture_of(&[7], dies_during_dwell(), None).await;
(daemon.refusing_deletes_from(1), prepared, dirs)
}
async fn cutover_fixture_dies_and_refuses_deletes() -> (CutOverDouble, Prepared, Dirs) {
let (daemon, prepared, dirs) = cutover_fixture_of(&[7], dies_during_dwell(), None).await;
(daemon.refusing_deletes(), prepared, dirs)
}
async fn cutover_fixture_original_respawns() -> (CutOverDouble, Prepared, Dirs) {
let (daemon, prepared, dirs) =
cutover_fixture_of(&[7], vec![(Duration::ZERO, Shape::Absent)], None).await;
(
daemon.while_the_originals(originals_respawn_mid_cutover()),
prepared,
dirs,
)
}
async fn cutover_fixture_refusing_start() -> (CutOverDouble, Prepared, Dirs) {
cutover_fixture_of(&[7], comes_up(), Some(0)).await
}
#[tokio::test(start_paused = true)]
async fn the_new_registration_names_current_explicitly() {
let (daemon, prepared, _dirs) = cutover_fixture().await;
let current = prepared.tree.current();
cut_over(&daemon, prepared).await.expect("cuts over");
let started = daemon.started();
assert_eq!(started.len(), 1, "exactly one Start");
assert_eq!(
started[0].cwd.as_deref(),
current.to_str(),
"cwd must be the current symlink itself, not a release and not None"
);
}
#[tokio::test(start_paused = true)]
async fn only_the_old_instances_are_deleted_and_by_id() {
let (daemon, prepared, _dirs) = cutover_fixture().await;
cut_over(&daemon, prepared).await.expect("cuts over");
assert_eq!(daemon.deleted(), vec![7], "the pre-existing instance's id");
assert!(
!daemon.deleted().contains(&NEWCOMER_ID),
"99 is the newcomer"
);
}
#[tokio::test(start_paused = true)]
async fn every_replaced_instance_is_deleted() {
let (daemon, prepared, _dirs) = cutover_fixture_with_instances(&[7, 8]).await;
cut_over(&daemon, prepared).await.expect("cuts over");
assert_eq!(daemon.deleted(), vec![7, 8]);
}
#[tokio::test(start_paused = true)]
async fn a_newcomer_that_dies_during_the_dwell_is_deleted_and_the_old_one_kept() {
let (daemon, prepared, _dirs) = cutover_fixture_dies_during_dwell().await;
let err = cut_over(&daemon, prepared).await.expect_err("gives up");
assert!(
daemon.deleted().contains(&NEWCOMER_ID),
"the newcomer is removed"
);
assert!(!daemon.deleted().contains(&7), "the old instance is kept");
let shown = err.to_string();
assert!(shown.contains("SO_REUSEPORT"), "{shown}");
}
#[tokio::test(start_paused = true)]
async fn a_newcomer_that_went_online_without_serving_is_still_rejected() {
let (daemon, prepared, _dirs) = cutover_fixture_online_then_gone().await;
let err = cut_over(&daemon, prepared)
.await
.expect_err("the dwell catches it");
assert!(matches!(err, Error::CutOver { .. }), "{err}");
assert!(err.to_string().contains("did not stay up"), "{err}");
}
#[tokio::test(start_paused = true)]
async fn a_newcomer_that_crash_loops_through_the_dwell_is_rejected() {
let (daemon, prepared, _dirs) = cutover_fixture_crash_looping().await;
let err = cut_over(&daemon, prepared)
.await
.expect_err("the pids moved");
assert!(matches!(err, Error::CutOver { .. }), "{err}");
assert!(err.to_string().contains("did not stay up"), "{err}");
}
#[tokio::test(start_paused = true)]
async fn an_abandoned_cutover_puts_the_original_config_back_in_the_roll() {
let (daemon, prepared, _dirs) = cutover_fixture_dies_during_dwell().await;
let original = prepared.previous_config.clone();
cut_over(&daemon, prepared).await.expect_err("gives up");
let last = daemon.started().last().cloned().expect("a repair Start");
assert_eq!(
last.cwd, original.cwd,
"the roll is re-recorded at the old cwd"
);
assert_eq!(
daemon.deleted().len(),
2,
"the newcomer, and the instance the repair Start spawned"
);
}
#[tokio::test(start_paused = true)]
async fn a_failed_roll_repair_is_named_not_glossed() {
let (daemon, prepared, _dirs) = cutover_fixture_dies_and_refuses_repair().await;
let err = cut_over(&daemon, prepared).await.expect_err("gives up");
let shown = err.to_string();
assert!(shown.contains("muster"), "{shown}");
assert!(
shown.contains("reboot") || shown.contains("restart"),
"{shown}"
);
}
#[tokio::test(start_paused = true)]
async fn a_refused_start_deletes_nothing() {
let (daemon, prepared, _dirs) = cutover_fixture_refusing_start().await;
let err = cut_over(&daemon, prepared).await.expect_err("refused");
assert!(matches!(err, Error::Request(_)), "{err}");
assert!(daemon.deleted().is_empty());
}
#[tokio::test]
async fn a_record_that_cannot_be_read_refuses_rather_than_guessing() {
let home = tempfile::tempdir().expect("tempdir");
let checkout = checkout_with_commit();
write_target(home.path(), "bpm", Watch::Auto, Some("old"));
let tree = Tree::for_sheep(home.path(), "bpm");
std::fs::write(tree.state_file(), "this is not toml = = =").expect("corrupt the record");
let entries = [("bpm", checkout.path())];
let daemon = RollOf(&entries);
let err = prepare(&daemon, home.path(), "bpm", &test_config())
.await
.expect_err("refuses");
let shown = err.to_string();
assert!(shown.contains("cannot be read"), "{shown}");
assert!(
shown.contains("Do NOT remove the tree"),
"the irreversible instruction is withheld: {shown}"
);
assert!(
!shown.contains("was never cut over"),
"it must not claim the cutover was abandoned: {shown}"
);
}
#[tokio::test]
async fn an_abandoned_target_is_not_pointed_at_shep_deploy() {
let home = tempfile::tempdir().expect("tempdir");
let checkout = checkout_with_commit();
write_target(home.path(), "bpm", Watch::Auto, None);
let entries = [("bpm", checkout.path())];
let daemon = RollOf(&entries);
let err = prepare(&daemon, home.path(), "bpm", &test_config())
.await
.expect_err("refuses");
let shown = err.to_string();
assert!(shown.contains("never cut over"), "{shown}");
assert!(shown.contains("Do NOT run `shep deploy bpm`"), "{shown}");
assert!(
shown.contains("shep-deploy setup bpm"),
"it says how to try again: {shown}"
);
}
#[tokio::test(start_paused = true)]
async fn a_refused_delete_names_every_instance_it_left_behind() {
let (daemon, prepared, _dirs) = cutover_fixture_comes_up_but_refuses_deletes().await;
let path = prepared.tree.state_file();
let sha = prepared.sha.clone();
let err = cut_over(&daemon, prepared).await.expect_err("names them");
let shown = err.to_string();
assert!(shown.contains("shep delete 7"), "{shown}");
assert!(
shown.contains("shep delete 8"),
"not just the first one: {shown}"
);
assert_eq!(
State::read(&path).expect("reads").deployed,
Some(sha),
"the release is serving, so the record names it"
);
}
#[tokio::test(start_paused = true)]
async fn a_shepherd_that_goes_quiet_still_reports_the_repair_it_could_not_make() {
let (daemon, prepared, _dirs) = cutover_fixture_shepherd_goes_quiet().await;
let err = cut_over(&daemon, prepared).await.expect_err("gives up");
assert!(
matches!(
err,
Error::CutOver {
repaired: false,
..
}
),
"{err}"
);
assert!(err.source().is_some(), "the shepherd's own error is kept");
let shown = err.to_string();
assert!(shown.contains("stopped answering"), "{shown}");
assert!(
shown.contains("muster"),
"the roll paragraph fires: {shown}"
);
}
#[tokio::test(start_paused = true)]
async fn an_abandoned_cutover_says_the_target_is_not_deployable_yet() {
let (daemon, prepared, _dirs) = cutover_fixture_dies_during_dwell().await;
let tree = prepared.tree.root().to_owned();
let err = cut_over(&daemon, prepared).await.expect_err("gives up");
let shown = err.to_string();
assert!(shown.contains("is NOT a deploy target"), "{shown}");
assert!(shown.contains("Do NOT run `shep deploy bpm`"), "{shown}");
assert!(shown.contains("shep-deploy setup bpm"), "{shown}");
assert!(
shown.contains(&format!("remove {}", tree.display())),
"{shown}"
);
assert!(!shown.contains("$SHEP_HOME"), "{shown}");
}
#[tokio::test(start_paused = true)]
async fn a_cutover_that_spawned_nothing_does_not_blame_the_port() {
let (daemon, prepared, _dirs) = cutover_fixture_never_appears().await;
let err = cut_over(&daemon, prepared).await.expect_err("gives up");
let shown = err.to_string();
assert!(!shown.contains("SO_REUSEPORT"), "{shown}");
assert!(shown.contains("No new instance appeared"), "{shown}");
}
#[tokio::test(start_paused = true)]
async fn a_cutover_that_added_nothing_does_not_claim_to_have_removed_it() {
let (daemon, prepared, _dirs) = cutover_fixture_never_appears().await;
let err = cut_over(&daemon, prepared).await.expect_err("gives up");
let shown = err.to_string();
assert!(shown.contains("No new instance appeared"), "{shown}");
assert!(
!shown.contains("instance it added has been removed"),
"it cannot have removed what it never added: {shown}"
);
assert!(
shown.contains("Nothing this cutover started is left registered"),
"{shown}"
);
}
#[tokio::test(start_paused = true)]
async fn a_repair_instance_left_behind_is_not_described_by_its_cwd() {
let (daemon, prepared, _dirs) = cutover_fixture_dies_and_refuses_the_second_delete().await;
let err = cut_over(&daemon, prepared).await.expect_err("gives up");
assert!(
matches!(
err,
Error::CutOver {
removed: false,
repaired: true,
..
}
),
"{err}"
);
let shown = err.to_string();
assert!(
!shown.contains("cwd under the deploy tree"),
"the hint would point at the wrong instance: {shown}"
);
assert!(shown.contains("shep delete <id>"), "{shown}");
}
#[tokio::test(start_paused = true)]
async fn a_newcomer_that_could_not_be_removed_is_named_not_claimed_gone() {
let (daemon, prepared, _dirs) = cutover_fixture_dies_and_refuses_deletes().await;
let err = cut_over(&daemon, prepared).await.expect_err("gives up");
let shown = err.to_string();
assert!(shown.contains("could NOT be removed"), "{shown}");
assert!(shown.contains("shep describe bpm"), "{shown}");
}
#[tokio::test(start_paused = true)]
async fn an_original_that_respawns_is_not_mistaken_for_the_newcomer() {
let (daemon, prepared, _dirs) = cutover_fixture_original_respawns().await;
let err = cut_over(&daemon, prepared).await.expect_err("gives up");
assert!(matches!(err, Error::CutOver { .. }), "{err}");
assert!(err.to_string().contains("already restarted"), "{err}");
assert!(
!daemon.deleted().contains(&7),
"the healthy original is kept: {:?}",
daemon.deleted()
);
}
#[tokio::test(start_paused = true)]
async fn the_record_advances_only_after_the_newcomer_is_online() {
let (daemon, prepared, _dirs) = cutover_fixture_never_appears().await;
let path = prepared.tree.state_file();
let err = cut_over(&daemon, prepared).await.expect_err("gives up");
assert!(matches!(err, Error::CutOver { .. }), "{err}");
assert_eq!(State::read(&path).expect("reads").deployed, None);
}
#[tokio::test(start_paused = true)]
async fn a_prepared_target_is_not_watched_until_the_cutover_lands() {
let (daemon, prepared, _dirs) = cutover_fixture().await;
let path = prepared.tree.state_file();
assert_eq!(
State::read(&path).expect("reads").watch,
Watch::Manual,
"prepared, and nothing has served from the tree yet"
);
cut_over(&daemon, prepared).await.expect("cuts over");
assert_eq!(
State::read(&path).expect("reads").watch,
Watch::Auto,
"the cutover landed, so the poll loop may have it"
);
}
#[tokio::test(start_paused = true)]
async fn an_abandoned_cutover_leaves_the_target_unwatched() {
let (daemon, prepared, _dirs) = cutover_fixture_never_appears().await;
let path = prepared.tree.state_file();
cut_over(&daemon, prepared).await.expect_err("gives up");
assert_eq!(State::read(&path).expect("reads").watch, Watch::Manual);
}
#[tokio::test(start_paused = true)]
async fn a_stranded_cutover_does_not_start_the_loop_off_on_its_own() {
let (daemon, prepared, _dirs) = cutover_fixture_comes_up_but_refuses_deletes().await;
let path = prepared.tree.state_file();
let err = cut_over(&daemon, prepared).await.expect_err("is stranded");
assert!(matches!(err, Error::Stranded { .. }), "{err}");
let state = State::read(&path).expect("reads");
assert!(state.deployed.is_some(), "the release is serving");
assert_eq!(state.watch, Watch::Manual);
}
#[tokio::test(start_paused = true)]
async fn the_record_and_the_answer_name_the_release_that_was_cut_over() {
let (daemon, prepared, _dirs) = cutover_fixture().await;
let path = prepared.tree.state_file();
let expected = prepared.sha.clone();
let sha = cut_over(&daemon, prepared).await.expect("cuts over");
assert_eq!(sha, expected, "the sha it answers with");
assert_eq!(State::read(&path).expect("reads").deployed, Some(expected));
}
}