use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use shep_client::shep_core::config::AppConfig;
use crate::daemon::Daemon;
use crate::error::Error;
use crate::paths::{self, Tree};
use crate::roll;
use crate::state::State;
#[derive(Debug)]
pub enum Restored {
Returned {
sheep: String,
to: PathBuf,
},
NeverMoved {
sheep: String,
at: PathBuf,
},
LeftRunning {
sheep: String,
from: PathBuf,
},
Failed {
sheep: String,
why: String,
},
Reset {
sheep: String,
why: String,
},
PartlyDeleted {
sheep: String,
why: String,
},
Lost {
sheep: String,
why: String,
},
RollUnreadable {
sheep: Vec<String>,
why: String,
},
}
pub async fn all<D: Daemon>(daemon: &D, shep_home: &Path) -> Vec<Restored> {
let Ok(names) = paths::targets(shep_home) else {
return Vec::new();
};
let registered = roll::registered(daemon).await;
let mut results = Vec::new();
let mut blocked_by_roll = Vec::new();
for sheep in names {
let tree = Tree::for_sheep(shep_home, &sheep);
let state = match State::read(&tree.state_file()) {
Ok(state) => state,
Err(err) => {
results.push(Restored::Failed {
sheep,
why: err.to_string(),
});
continue;
}
};
let never_moved = state.deployed.is_none();
let (Some(cwd), Some(script)) = (state.origin_cwd, state.origin_script) else {
results.push(Restored::LeftRunning {
sheep,
from: tree.current(),
});
continue;
};
let Ok(registered) = ®istered else {
blocked_by_roll.push(sheep);
continue;
};
if never_moved && !registered_against(registered, &sheep, &tree) {
results.push(Restored::NeverMoved { sheep, at: cwd });
continue;
}
results.push(
match put_back(daemon, &sheep, registered, &cwd, &script).await {
PutBack::Done => Restored::Returned { sheep, to: cwd },
PutBack::Untouched(err) => Restored::Failed {
sheep,
why: err.to_string(),
},
PutBack::Reset(err) => Restored::Reset {
sheep,
why: err.to_string(),
},
PutBack::PartlyDeleted(err) => Restored::PartlyDeleted {
sheep,
why: err.to_string(),
},
PutBack::Deleted(err) => Restored::Lost {
sheep,
why: err.to_string(),
},
},
);
}
if let Err(err) = ®istered
&& !blocked_by_roll.is_empty()
{
results.push(Restored::RollUnreadable {
sheep: blocked_by_roll,
why: err.to_string(),
});
}
results
}
enum PutBack {
Done,
Untouched(Error),
Reset(Error),
PartlyDeleted(Error),
Deleted(Error),
}
async fn put_back<D: Daemon>(
daemon: &D,
sheep: &str,
registered: &BTreeMap<String, AppConfig>,
cwd: &Path,
script: &str,
) -> PutBack {
let Some(current) = registered.get(sheep).cloned() else {
return PutBack::Untouched(Error::Config(format!(
"{sheep} is no longer registered, so there is nothing to put back"
)));
};
let mut restored = current.clone();
restored.cwd = Some(cwd.display().to_string());
restored.script = script.to_owned();
let live = match daemon.describe(sheep).await {
Ok(live) => live,
Err(err) => return PutBack::Untouched(err),
};
let mut any_delete_landed = false;
for info in &live {
if let Err(err) = daemon.delete(info.id).await {
return if any_delete_landed {
PutBack::PartlyDeleted(err)
} else {
PutBack::Untouched(err)
};
}
any_delete_landed = true;
}
match daemon.start(vec![restored]).await {
Ok(()) => PutBack::Done,
Err(err) => {
if daemon.start(vec![current]).await.is_ok() {
PutBack::Reset(err)
} else {
PutBack::Deleted(err)
}
}
}
}
#[must_use]
pub fn report(results: &[Restored]) -> String {
results
.iter()
.map(|result| match result {
Restored::Returned { sheep, to } => {
format!("{sheep} restored to {}\n", to.display())
}
Restored::NeverMoved { sheep, at } => {
format!(
"{sheep} was never moved - its cutover did not land - and is still running \
from {}\n",
at.display()
)
}
Restored::LeftRunning { sheep, from } => {
format!("{sheep} still running from {}\n", from.display())
}
Restored::Failed { sheep, why } => {
format!("{sheep} could not be restored and was left as it is: {why}\n")
}
Restored::Reset { sheep, why } => format!(
"{sheep} could not be restored ({why}), so its previous configuration was put \
back instead - doing that stopped it and started it again, so it is running \
with the same config as before but nothing mid-flight survived\n"
),
Restored::PartlyDeleted { sheep, why } => format!(
"{sheep}: only SOME of its instances could be stopped before a delete failed \
({why}), so it is neither fully running nor fully removed. Run `shep describe \
{sheep}` to see what is actually still there, and `shep flock` for whether it \
is still registered, before assuming either.\n"
),
Restored::Lost { sheep, why } => format!(
"{sheep} IS NO LONGER REGISTERED: restoring it failed ({why}) and so did \
putting its previous configuration back, so it is stopped and gone from the \
flock. It will not come back on its own after a restart. Re-register it from \
its own Flockfile.\n"
),
Restored::RollUnreadable { sheep, why } => format!(
"{}: none of these could be checked against the muster roll, so none of them \
could be restored ({why})\n",
sheep.join(", ")
),
})
.collect()
}
fn registered_against(registered: &BTreeMap<String, AppConfig>, sheep: &str, tree: &Tree) -> bool {
registered
.get(sheep)
.and_then(|app| app.cwd.as_deref())
.is_some_and(|cwd| Path::new(cwd) == tree.current())
}
#[cfg(test)]
mod tests {
#[tokio::test]
async fn a_record_that_cannot_be_read_is_reported_rather_than_skipped() {
let home = tempfile::tempdir().expect("tempdir");
let tree = Tree::for_sheep(home.path(), "garbled");
std::fs::create_dir_all(tree.root()).expect("create the tree");
std::fs::write(tree.state_file(), "this is not toml").expect("write deploy.toml");
let results = all(&Recording::new(&[], Refuse::Never), home.path()).await;
assert_eq!(results.len(), 1, "the target must be reported, not skipped");
assert!(
matches!(&results[0], Restored::Failed { sheep, .. } if sheep == "garbled"),
"an unreadable record must be a reported failure, got: {:?}",
results[0]
);
}
use std::cell::{Cell, RefCell};
use std::fs;
use shep_client::RequestError;
use shep_client::shep_core::protocol::{ProcessInfo, RpcError, RpcErrorCode};
use shep_client::shep_core::status::ProcStatus;
use super::*;
use crate::state::{Verify, Watch};
fn write_target_with_origin(home: &Path, sheep: &str, origin_cwd: &str, origin_script: &str) {
let tree = Tree::for_sheep(home, sheep);
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: Some("a1b2c3d".to_owned()),
failed: None,
verify: Verify::default(),
watch: Watch::default(),
origin_cwd: Some(PathBuf::from(origin_cwd)),
origin_script: Some(origin_script.to_owned()),
checkout: PathBuf::from(origin_cwd),
};
state.write(&tree.state_file()).expect("write state");
}
fn write_target_with_origin_absent(home: &Path, sheep: &str) {
let tree = Tree::for_sheep(home, sheep);
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: Some("a1b2c3d".to_owned()),
failed: None,
verify: Verify::default(),
watch: Watch::default(),
origin_cwd: None,
origin_script: None,
checkout: PathBuf::from("/srv/deploy-tree"),
};
state.write(&tree.state_file()).expect("write state");
}
enum Refuse {
Never,
FirstOnly,
Always,
}
struct Recording {
sheep: Vec<&'static str>,
refuse: Refuse,
instances: u32,
describe_fails: bool,
delete_fails_at: Option<u32>,
unreadable_roll: bool,
registered_cwd: Option<String>,
calls: RefCell<Vec<&'static str>>,
starts: RefCell<Vec<AppConfig>>,
attempts: Cell<usize>,
delete_attempts: Cell<u32>,
}
impl Recording {
fn new(sheep: &[&'static str], refuse: Refuse) -> Self {
Self {
sheep: sheep.to_vec(),
refuse,
instances: 1,
describe_fails: false,
delete_fails_at: None,
unreadable_roll: false,
registered_cwd: None,
calls: RefCell::new(Vec::new()),
starts: RefCell::new(Vec::new()),
attempts: Cell::new(0),
delete_attempts: Cell::new(0),
}
}
fn with_registered(sheep: &[&'static str]) -> Self {
Self::new(sheep, Refuse::Never)
}
fn registered_at(mut self, cwd: &std::path::Path) -> Self {
self.registered_cwd = Some(cwd.display().to_string());
self
}
fn refusing_first_start_only(sheep: &[&'static str]) -> Self {
Self::new(sheep, Refuse::FirstOnly)
}
fn refusing_every_start(sheep: &[&'static str]) -> Self {
Self::new(sheep, Refuse::Always)
}
fn with_unreadable_roll() -> Self {
let mut this = Self::new(&[], Refuse::Never);
this.unreadable_roll = true;
this
}
fn refusing_describe(sheep: &[&'static str]) -> Self {
let mut this = Self::new(sheep, Refuse::Never);
this.describe_fails = true;
this
}
fn refusing_first_delete_only(sheep: &[&'static str]) -> Self {
let mut this = Self::new(sheep, Refuse::Never);
this.delete_fails_at = Some(1);
this
}
fn refusing_delete_partway(sheep: &[&'static str]) -> Self {
let mut this = Self::new(sheep, Refuse::Never);
this.instances = 2;
this.delete_fails_at = Some(2);
this
}
fn calls(&self) -> Vec<&'static str> {
self.calls.borrow().clone()
}
fn started(&self) -> Vec<AppConfig> {
self.starts.borrow().clone()
}
}
impl Daemon for Recording {
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> {
if self.describe_fails {
return Err(Error::Request(RequestError::Rpc(RpcError {
code: RpcErrorCode::Internal,
message: "describe refused".to_owned(),
daemon_version: None,
})));
}
Ok((0..self.instances)
.map(|offset| {
ProcessInfo::builder(offset + 1, sheep, ProcStatus::Online)
.pid(Some(1000 + offset))
.build()
})
.collect())
}
async fn start(&self, apps: Vec<AppConfig>) -> Result<(), Error> {
let attempt = self.attempts.get();
self.attempts.set(attempt + 1);
let refused = match self.refuse {
Refuse::Never => false,
Refuse::FirstOnly => attempt == 0,
Refuse::Always => true,
};
self.starts.borrow_mut().extend(apps);
if refused {
return Err(Error::Request(RequestError::Rpc(RpcError {
code: RpcErrorCode::Internal,
message: "refused".to_owned(),
daemon_version: None,
})));
}
self.calls.borrow_mut().push("start");
Ok(())
}
async fn delete(&self, _id: u32) -> Result<(), Error> {
let attempt = self.delete_attempts.get() + 1;
self.delete_attempts.set(attempt);
if self.delete_fails_at == Some(attempt) {
return Err(Error::Request(RequestError::Rpc(RpcError {
code: RpcErrorCode::Internal,
message: "delete refused".to_owned(),
daemon_version: None,
})));
}
self.calls.borrow_mut().push("delete");
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> {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.keep().join("flock.json");
if self.unreadable_roll {
fs::write(
&path,
"{\"apps\":[{\"app\":{\"name\":\"w\",\"a_field_from_the_future\":1}}]}",
)
.expect("write roll");
return Ok(path);
}
let apps: Vec<String> = self
.sheep
.iter()
.map(|name| {
let cwd = self
.registered_cwd
.clone()
.unwrap_or_else(|| "/srv/deploy-tree/current".to_owned());
format!(
"{{\"app\":{{\"name\":{name:?},\"script\":\"the-shepherds-own-script\",\
\"cwd\":{cwd:?}}}}}"
)
})
.collect();
fs::write(&path, format!("{{\"apps\":[{}]}}", apps.join(","))).expect("write roll");
Ok(path)
}
async fn set_smit(&self, _sheep: &str, _text: &str) -> Result<(), Error> {
unimplemented!()
}
}
#[tokio::test]
async fn a_pre_existing_sheep_goes_back_to_its_own_checkout() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "bpm", "/srv/reactmap", "bun .");
let daemon = Recording::with_registered(&["bpm"]);
let results = all(&daemon, home.path()).await;
assert!(matches!(results[0], Restored::Returned { .. }));
let started = daemon.started();
assert_eq!(started[0].cwd.as_deref(), Some("/srv/reactmap"));
assert_eq!(started[0].script, "bun .");
}
#[tokio::test]
async fn the_old_registration_is_removed_before_the_new_one_is_started() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "bpm", "/srv/reactmap", "bun .");
let daemon = Recording::with_registered(&["bpm"]);
all(&daemon, home.path()).await;
assert_eq!(
daemon.calls(),
vec!["delete", "start"],
"deleting after starting would leave two registrations"
);
}
#[tokio::test]
async fn a_bootstrapped_sheep_is_left_running_and_named_in_the_report() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin_absent(home.path(), "ctm");
let daemon = Recording::with_registered(&["ctm"]);
let results = all(&daemon, home.path()).await;
assert!(daemon.calls().is_empty(), "nothing is stopped or started");
let text = report(&results);
assert!(text.contains("ctm still running from"), "{text}");
assert!(text.contains("deploy/ctm/current"), "{text}");
}
#[tokio::test]
async fn a_cutover_that_never_landed_is_not_restored() {
let home = tempfile::tempdir().expect("tempdir");
write_target_never_cut_over(home.path(), "ctm", "/srv/ctm", "./run.sh");
let daemon = Recording::with_registered(&["ctm"]);
let results = all(&daemon, home.path()).await;
assert!(
daemon.calls().is_empty(),
"a sheep that never moved must not be stopped or started: {:?}",
daemon.calls()
);
let text = report(&results);
assert!(text.contains("ctm was never moved"), "{text}");
assert!(text.contains("/srv/ctm"), "{text}");
}
#[tokio::test]
async fn a_cutover_that_registered_is_restored_even_without_a_deployed_sha() {
let home = tempfile::tempdir().expect("tempdir");
write_target_never_cut_over(home.path(), "ctm", "/srv/ctm", "./run.sh");
let tree = Tree::for_sheep(home.path(), "ctm");
let daemon = Recording::with_registered(&["ctm"]).registered_at(&tree.current());
let results = all(&daemon, home.path()).await;
let text = report(&results);
assert!(
!text.contains("never moved"),
"a sheep the cutover registered against the tree was moved: {text}"
);
assert!(
!daemon.calls().is_empty(),
"it must actually be put back, not merely reported"
);
}
fn write_target_never_cut_over(home: &Path, sheep: &str, origin_cwd: &str, script: &str) {
let tree = Tree::for_sheep(home, sheep);
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: None,
failed: None,
verify: Verify::default(),
watch: Watch::default(),
origin_cwd: Some(PathBuf::from(origin_cwd)),
origin_script: Some(script.to_owned()),
checkout: PathBuf::from(origin_cwd),
};
state.write(&tree.state_file()).expect("write state");
}
#[tokio::test]
async fn a_failure_is_reported_and_the_rest_still_run() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "aaa", "/srv/a", "./a");
write_target_with_origin(home.path(), "zzz", "/srv/z", "./z");
let daemon = Recording::refusing_first_start_only(&["aaa", "zzz"]);
let results = all(&daemon, home.path()).await;
assert_eq!(results.len(), 2);
assert!(
matches!(results[0], Restored::Reset { .. }),
"{:?}",
results[0]
);
assert!(matches!(results[1], Restored::Returned { .. }));
assert!(report(&results).contains("aaa"), "the failure is named");
}
#[tokio::test]
async fn a_refused_restore_puts_the_shepherds_own_config_back() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "bpm", "/srv/reactmap", "bun .");
let daemon = Recording::refusing_first_start_only(&["bpm"]);
let results = all(&daemon, home.path()).await;
assert!(
matches!(results[0], Restored::Reset { .. }),
"{:?}",
results[0]
);
assert_eq!(daemon.started().len(), 2, "the restore, then the fallback");
assert_eq!(
daemon.started()[1].cwd.as_deref(),
Some("/srv/deploy-tree/current"),
"the fallback re-registers what the shepherd had, not the restore"
);
}
#[tokio::test]
async fn a_rescued_restore_is_worded_as_a_reset_not_left_alone() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "bpm", "/srv/reactmap", "bun .");
let daemon = Recording::refusing_first_start_only(&["bpm"]);
let results = all(&daemon, home.path()).await;
let text = report(&results);
assert!(!text.contains("left as it is"), "{text}");
assert!(text.contains("stopped it and started it again"), "{text}");
}
#[tokio::test]
async fn a_sheep_left_deleted_says_so_in_those_words() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "bpm", "/srv/reactmap", "bun .");
let daemon = Recording::refusing_every_start(&["bpm"]);
let results = all(&daemon, home.path()).await;
assert!(
matches!(results[0], Restored::Lost { .. }),
"{:?}",
results[0]
);
let text = report(&results);
assert!(text.contains("NO LONGER REGISTERED"), "{text}");
assert!(text.contains("will not come back"), "{text}");
}
#[tokio::test]
async fn the_deploy_tree_is_left_on_disk() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "bpm", "/srv/reactmap", "bun .");
all(&Recording::with_registered(&["bpm"]), home.path()).await;
assert!(home.path().join("deploy/bpm/deploy.toml").is_file());
}
#[tokio::test]
async fn a_roll_read_failure_is_reported_once_dog_wide_with_the_real_cause() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "aaa", "/srv/a", "./a");
write_target_with_origin(home.path(), "zzz", "/srv/z", "./z");
let daemon = Recording::with_unreadable_roll();
let results = all(&daemon, home.path()).await;
assert_eq!(
results.len(),
1,
"one row for the whole roll failure, not one per sheep: {results:?}"
);
assert!(matches!(results[0], Restored::RollUnreadable { .. }));
let text = report(&results);
assert!(text.contains("aaa"), "{text}");
assert!(text.contains("zzz"), "{text}");
assert!(text.contains("newer"), "{text}");
assert!(
daemon.calls().is_empty(),
"nothing was deleted or started without a readable roll"
);
}
#[tokio::test]
async fn a_describe_that_is_refused_leaves_the_sheep_untouched() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "bpm", "/srv/reactmap", "bun .");
let daemon = Recording::refusing_describe(&["bpm"]);
let results = all(&daemon, home.path()).await;
assert!(
matches!(results[0], Restored::Failed { .. }),
"{:?}",
results[0]
);
assert!(daemon.calls().is_empty(), "nothing was deleted or started");
}
#[tokio::test]
async fn a_delete_that_fails_partway_through_is_reported_as_partly_deleted() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "bpm", "/srv/reactmap", "bun .");
let daemon = Recording::refusing_delete_partway(&["bpm"]);
let results = all(&daemon, home.path()).await;
assert!(
matches!(results[0], Restored::PartlyDeleted { .. }),
"{:?}",
results[0]
);
assert_eq!(
daemon.calls(),
vec!["delete"],
"the first delete lands before the second refuses, and nothing after it runs"
);
assert!(
daemon.started().is_empty(),
"start is never reached once a delete fails"
);
}
#[tokio::test]
async fn a_partly_deleted_sheep_is_worded_as_neither_running_nor_removed() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "bpm", "/srv/reactmap", "bun .");
let daemon = Recording::refusing_delete_partway(&["bpm"]);
let results = all(&daemon, home.path()).await;
let text = report(&results);
assert!(!text.contains("left as it is"), "{text}");
assert!(!text.contains("gone from the flock"), "{text}");
assert!(!text.contains("never touched"), "{text}");
assert!(
text.contains("neither fully running nor fully removed"),
"{text}"
);
assert!(text.contains("shep describe bpm"), "{text}");
assert!(text.contains("shep flock"), "{text}");
}
#[tokio::test]
async fn a_delete_refused_with_nothing_landed_yet_stays_failed() {
let home = tempfile::tempdir().expect("tempdir");
write_target_with_origin(home.path(), "bpm", "/srv/reactmap", "bun .");
let daemon = Recording::refusing_first_delete_only(&["bpm"]);
let results = all(&daemon, home.path()).await;
assert!(
matches!(results[0], Restored::Failed { .. }),
"{:?}",
results[0]
);
let text = report(&results);
assert!(text.contains("left as it is"), "{text}");
assert!(
!text.contains("neither fully running nor fully removed"),
"{text}"
);
}
#[test]
fn a_report_with_a_failure_still_names_every_other_row() {
let results = vec![
Restored::Failed {
sheep: "aaa".to_owned(),
why: "refused".to_owned(),
},
Restored::Returned {
sheep: "zzz".to_owned(),
to: PathBuf::from("/srv/z"),
},
];
let text = report(&results);
assert!(text.contains("aaa"), "{text}");
assert!(text.contains("zzz"), "{text}");
}
}