use std::cell::RefCell;
use std::collections::HashSet;
use std::time::{Duration, Instant};
use super::expose::{self, MeshExposer};
use super::runtime::{container_name, deployment_labels, ContainerRuntime, RunSpec};
use super::state::LocalState;
use super::{Desired, Healthcheck};
#[derive(Debug, Clone, PartialEq)]
pub enum Action {
Pull {
id: String,
},
Start {
id: String,
},
Health {
id: String,
},
StopOld {
id: String,
},
Stop {
id: String,
},
Remove {
id: String,
container_id: String,
},
Report {
id: String,
phase: String,
message: Option<String>,
},
}
impl Action {
fn deployment_id(&self) -> &str {
match self {
Action::Pull { id }
| Action::Start { id }
| Action::Health { id }
| Action::StopOld { id }
| Action::Stop { id }
| Action::Remove { id, .. }
| Action::Report { id, .. } => id,
}
}
}
pub fn plan(desired: &[Desired], local: &LocalState) -> Vec<Action> {
let mut actions = Vec::new();
for d in desired {
let record = local.deployments.get(&d.id);
if d.desired_state == "stopped" {
if let Some(rec) = record {
if rec.container_id.is_some() {
actions.push(Action::Stop { id: d.id.clone() });
}
}
continue;
}
let deploying = match record {
None => {
actions.push(Action::Pull { id: d.id.clone() });
actions.push(Action::Start { id: d.id.clone() });
actions.push(Action::Health { id: d.id.clone() });
true
}
Some(rec) if rec.version != d.version => {
actions.push(Action::Pull { id: d.id.clone() });
actions.push(Action::Start { id: d.id.clone() });
actions.push(Action::Health { id: d.id.clone() });
actions.push(Action::StopOld { id: d.id.clone() });
true
}
Some(rec) if rec.phase != "healthy" || rec.container_id.is_none() => {
actions.push(Action::Start { id: d.id.clone() });
actions.push(Action::Health { id: d.id.clone() });
true
}
Some(_) => false, };
if deploying {
if let Some(prev) = &d.previous {
if prev.id != d.id {
actions.push(Action::Stop {
id: prev.id.clone(),
});
}
}
}
}
actions
}
pub fn plan_failure(
id: &str,
failed_container_id: Option<&str>,
had_previous: bool,
) -> Vec<Action> {
let mut actions = Vec::new();
if let Some(cid) = failed_container_id {
actions.push(Action::Remove {
id: id.to_string(),
container_id: cid.to_string(),
});
}
actions.push(Action::Report {
id: id.to_string(),
phase: "failed".to_string(),
message: None,
});
if had_previous {
actions.push(Action::Report {
id: id.to_string(),
phase: "rolled_back".to_string(),
message: None,
});
}
actions
}
pub trait StatusReporter {
#[allow(clippy::too_many_arguments)]
fn report(
&self,
id: &str,
version: u64,
phase: &str,
message: &str,
attempt: u32,
container_id: Option<&str>,
endpoint: Option<&str>,
mesh_endpoint: Option<&str>,
);
}
fn wait_healthy(rt: &dyn ContainerRuntime, hc: &Option<Healthcheck>, container_id: &str) -> bool {
let Some(hc) = hc else { return true };
let deadline = Instant::now() + Duration::from_secs(hc.timeout_s.max(1));
loop {
if rt.health_check(container_id, hc) {
return true;
}
if Instant::now() >= deadline {
return false;
}
std::thread::sleep(Duration::from_secs(2));
}
}
pub fn drive(
rt: &dyn ContainerRuntime,
reporter: &dyn StatusReporter,
desired: &[Desired],
local: &mut LocalState,
) {
let actions = plan(desired, local);
let mut failed: std::collections::HashSet<String> = Default::default();
let attempt: u32 = 1;
for action in actions {
let id = action.deployment_id().to_string();
if failed.contains(&id) {
continue;
}
let Some(d) = desired.iter().find(|d| d.id == id) else {
if let Action::Stop { id } = &action {
let version = local.deployments.get(id).map(|r| r.version).unwrap_or(0);
do_stop(rt, reporter, local, id, version);
}
continue;
};
match action {
Action::Pull { id } => {
reporter.report(&id, d.version, "pulling", "", attempt, None, None, None);
if let Err(e) = rt.pull(&d.image) {
run_failure(reporter, local, &id, d.version, attempt, None, &e);
failed.insert(id);
}
}
Action::Start { id } => {
reporter.report(&id, d.version, "starting", "", attempt, None, None, None);
if !d.warm {
println!(" [DEPLOY] warm=false not implemented, keeping warm ({id})");
}
let name = container_name(&id, d.version);
let labels = deployment_labels(&id, d.version);
let mut env = d.env.clone();
if let Some(sa) = &d.service_account {
env.insert("ZAKURO_SERVICE_ACCOUNT".to_string(), sa.name.clone());
env.insert("ZAKURO_GRANTS".to_string(), sa.grants.join(","));
}
let spec = RunSpec {
name,
image: &d.image,
cmd: d.cmd.as_deref(),
ports: &d.ports,
env: &env,
labels,
};
match rt.run(&spec) {
Ok(container_id) => {
let rec = local.deployments.entry(id.clone()).or_default();
rec.previous_container_id = rec.container_id.take();
rec.container_id = Some(container_id);
rec.image = d.image.clone();
rec.version = d.version;
rec.price_per_second = d.price_per_second;
rec.service_account = d.service_account.clone();
rec.warm = d.warm;
rec.phase = "starting".to_string();
let _ = local.save();
}
Err(e) => {
run_failure(reporter, local, &id, d.version, attempt, None, &e);
failed.insert(id);
}
}
}
Action::Health { id } => {
let container_id = local
.deployments
.get(&id)
.and_then(|r| r.container_id.clone());
let Some(container_id) = container_id else {
failed.insert(id);
continue;
};
if wait_healthy(rt, &d.healthcheck, &container_id) {
let container_port = d
.healthcheck
.as_ref()
.map(|hc| hc.port)
.or_else(|| d.ports.first().map(|p| p.container));
let ip = rt.container_ip(&container_id).ok();
let endpoint = match (&ip, container_port) {
(Some(ip), Some(port)) => Some(format!("{ip}:{port}")),
_ => None,
};
if let Some(rec) = local.deployments.get_mut(&id) {
rec.phase = "healthy".to_string();
rec.ip = ip;
rec.port = container_port;
rec.endpoint = endpoint.clone();
}
let _ = local.save();
reporter.report(
&id,
d.version,
"healthy",
"",
attempt,
Some(&container_id),
endpoint.as_deref(),
None,
);
} else {
let logs = rt.logs_tail(&container_id, 20);
let _ = rt.stop(&container_id);
let _ = rt.rm(&container_id);
let previous = local
.deployments
.get_mut(&id)
.and_then(|rec| rec.previous_container_id.take());
run_failure(
reporter,
local,
&id,
d.version,
attempt,
Some(&container_id),
&logs,
);
if let Some(prev) = previous {
let _ = rt.restart(&prev);
restore_previous_container(reporter, local, &id, d.version, attempt, prev);
}
failed.insert(id);
}
}
Action::StopOld { id } => {
if let Some(rec) = local.deployments.get_mut(&id) {
if let Some(old) = rec.previous_container_id.take() {
let _ = rt.stop(&old);
let _ = rt.rm(&old);
}
}
let _ = local.save();
}
Action::Stop { id } => do_stop(rt, reporter, local, &id, d.version),
Action::Remove { container_id, .. } => {
let _ = rt.stop(&container_id);
let _ = rt.rm(&container_id);
}
Action::Report { id, phase, message } => {
reporter.report(
&id,
d.version,
&phase,
message.as_deref().unwrap_or(""),
attempt,
None,
None,
None,
);
}
}
}
}
pub fn drive_and_expose(
rt: &dyn ContainerRuntime,
reporter: &dyn StatusReporter,
exposer: &dyn MeshExposer,
mesh_ip: Option<&str>,
desired: &[Desired],
local: &mut LocalState,
) {
let tracked = ReportedIds {
inner: reporter,
ids: RefCell::default(),
};
drive(rt, &tracked, desired, local);
let mut changed = false;
for id in tracked.ids.into_inner() {
if let Some(rec) = local.deployments.get_mut(&id) {
changed |= rec.mesh_endpoint.take().is_some();
}
}
let synced = expose::sync(local, exposer, mesh_ip, reporter);
if synced || changed {
let _ = local.save();
}
}
struct ReportedIds<'a> {
inner: &'a dyn StatusReporter,
ids: RefCell<HashSet<String>>,
}
impl StatusReporter for ReportedIds<'_> {
fn report(
&self,
id: &str,
version: u64,
phase: &str,
message: &str,
attempt: u32,
container_id: Option<&str>,
endpoint: Option<&str>,
mesh_endpoint: Option<&str>,
) {
self.ids.borrow_mut().insert(id.to_string());
self.inner.report(
id,
version,
phase,
message,
attempt,
container_id,
endpoint,
mesh_endpoint,
);
}
}
fn do_stop(
rt: &dyn ContainerRuntime,
reporter: &dyn StatusReporter,
local: &mut LocalState,
id: &str,
version: u64,
) {
if let Some(rec) = local.deployments.get_mut(id) {
if let Some(cid) = rec.container_id.take() {
let _ = rt.stop(&cid);
let _ = rt.rm(&cid);
}
rec.phase = "stopped".to_string();
rec.endpoint = None;
rec.port = None;
}
let _ = local.save();
reporter.report(id, version, "stopped", "", 1, None, None, None);
}
fn run_failure(
reporter: &dyn StatusReporter,
local: &mut LocalState,
id: &str,
version: u64,
attempt: u32,
failed_container_id: Option<&str>,
message: &str,
) {
reporter.report(
id,
version,
"failed",
message,
attempt,
failed_container_id,
None,
None,
);
let rec = local.deployments.entry(id.to_string()).or_default();
if failed_container_id.is_some() || rec.container_id.is_none() {
rec.container_id = None;
rec.phase = "failed".to_string();
}
let _ = local.save();
}
fn restore_previous_container(
reporter: &dyn StatusReporter,
local: &mut LocalState,
id: &str,
version: u64,
attempt: u32,
previous_container_id: String,
) {
let endpoint = local.deployments.get(id).and_then(|r| r.endpoint.clone());
if let Some(rec) = local.deployments.get_mut(id) {
rec.container_id = Some(previous_container_id.clone());
rec.phase = "rolled_back".to_string();
}
let _ = local.save();
reporter.report(
id,
version,
"rolled_back",
"",
attempt,
Some(&previous_container_id),
endpoint.as_deref(),
None,
);
}
#[cfg(test)]
mod tests {
use super::*;
use crate::broker::deploy::state::DeploymentRecord;
use std::collections::HashMap;
fn desired(id: &str, version: u64, desired_state: &str) -> Desired {
Desired {
id: id.to_string(),
version,
name: format!("name-{id}"),
image: "python:3.12-slim".to_string(),
cmd: None,
ports: vec![super::super::Port {
container: 8000,
protocol: "http".to_string(),
}],
env: HashMap::new(),
healthcheck: None,
price_per_second: 0.0001,
desired_state: desired_state.to_string(),
service_account: None,
warm: true,
previous: None,
}
}
fn record(version: u64, phase: &str, container_id: Option<&str>) -> DeploymentRecord {
DeploymentRecord {
version,
container_id: container_id.map(|s| s.to_string()),
image: "python:3.12-slim".to_string(),
endpoint: None,
ip: None,
port: None,
price_per_second: 0.0001,
service_account: None,
warm: true,
phase: phase.to_string(),
previous_container_id: None,
mesh_endpoint: None,
mesh_port: None,
}
}
#[test]
fn plan_new_deployment_pulls_starts_and_checks_health() {
let d = desired("dep_1", 1, "running");
let local = LocalState::default();
let actions = plan(&[d], &local);
assert_eq!(
actions,
vec![
Action::Pull { id: "dep_1".into() },
Action::Start { id: "dep_1".into() },
Action::Health { id: "dep_1".into() },
]
);
}
#[test]
fn plan_version_bump_appends_stop_old_after_health() {
let d = desired("dep_1", 2, "running");
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(1, "healthy", Some("old-container")));
let actions = plan(&[d], &local);
assert_eq!(
actions,
vec![
Action::Pull { id: "dep_1".into() },
Action::Start { id: "dep_1".into() },
Action::Health { id: "dep_1".into() },
Action::StopOld { id: "dep_1".into() },
]
);
}
#[test]
fn plan_stopped_desired_state_stops_running_container() {
let d = desired("dep_1", 1, "stopped");
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(1, "healthy", Some("c1")));
let actions = plan(&[d], &local);
assert_eq!(actions, vec![Action::Stop { id: "dep_1".into() }]);
}
#[test]
fn plan_stopped_desired_state_with_nothing_running_is_a_noop() {
let d = desired("dep_1", 1, "stopped");
let local = LocalState::default();
assert_eq!(plan(&[d], &local), vec![]);
}
#[test]
fn plan_unchanged_and_healthy_is_a_noop() {
let d = desired("dep_1", 1, "running");
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(1, "healthy", Some("c1")));
assert_eq!(plan(&[d], &local), vec![]);
}
#[test]
fn plan_resumes_a_stuck_mid_flight_deployment() {
let d = desired("dep_1", 1, "running");
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(1, "starting", Some("c1")));
let actions = plan(&[d], &local);
assert_eq!(
actions,
vec![
Action::Start { id: "dep_1".into() },
Action::Health { id: "dep_1".into() },
]
);
}
#[test]
fn plan_previous_naming_another_deployment_stops_it_on_handoff() {
let mut d = desired("dep_2", 1, "running");
d.previous = Some(Box::new(super::super::PreviousDeployment {
id: "dep_1".into(),
image: "old-image".into(),
cmd: None,
ports: vec![],
env: HashMap::new(),
healthcheck: None,
}));
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(1, "healthy", Some("old-container")));
let actions = plan(&[d], &local);
assert_eq!(
actions,
vec![
Action::Pull { id: "dep_2".into() },
Action::Start { id: "dep_2".into() },
Action::Health { id: "dep_2".into() },
Action::Stop { id: "dep_1".into() },
]
);
}
#[test]
fn plan_previous_naming_the_same_deployment_does_not_double_stop() {
let mut d = desired("dep_1", 2, "running");
d.previous = Some(Box::new(super::super::PreviousDeployment {
id: "dep_1".into(),
image: "python:3.12-slim".into(),
cmd: None,
ports: vec![],
env: HashMap::new(),
healthcheck: None,
}));
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(1, "healthy", Some("old-container")));
let actions = plan(&[d], &local);
assert_eq!(
actions,
vec![
Action::Pull { id: "dep_1".into() },
Action::Start { id: "dep_1".into() },
Action::Health { id: "dep_1".into() },
Action::StopOld { id: "dep_1".into() },
]
);
}
#[test]
fn plan_failure_with_container_and_previous_removes_reports_failed_then_rolled_back() {
let actions = plan_failure("dep_1", Some("new-c"), true);
assert_eq!(
actions,
vec![
Action::Remove {
id: "dep_1".into(),
container_id: "new-c".into(),
},
Action::Report {
id: "dep_1".into(),
phase: "failed".into(),
message: None,
},
Action::Report {
id: "dep_1".into(),
phase: "rolled_back".into(),
message: None,
},
]
);
}
#[test]
fn plan_failure_without_container_or_previous_is_just_a_failed_report() {
let actions = plan_failure("dep_1", None, false);
assert_eq!(
actions,
vec![Action::Report {
id: "dep_1".into(),
phase: "failed".into(),
message: None,
}]
);
}
use std::sync::Mutex;
#[derive(Default)]
struct FakeRuntime {
next_container_id: Mutex<u32>,
run_should_fail: Mutex<bool>,
health_should_fail: Mutex<bool>,
stopped: Mutex<Vec<String>>,
removed: Mutex<Vec<String>>,
restarted: Mutex<Vec<String>>,
last_run_env: Mutex<HashMap<String, String>>,
}
impl ContainerRuntime for FakeRuntime {
fn available(&self) -> bool {
true
}
fn pull(&self, _image: &str) -> Result<(), String> {
Ok(())
}
fn run(&self, spec: &RunSpec) -> Result<String, String> {
*self.last_run_env.lock().unwrap() = spec.env.clone();
if *self.run_should_fail.lock().unwrap() {
return Err("run failed".to_string());
}
let mut n = self.next_container_id.lock().unwrap();
*n += 1;
Ok(format!("c{n}"))
}
fn stop(&self, container_id: &str) -> Result<(), String> {
self.stopped.lock().unwrap().push(container_id.to_string());
Ok(())
}
fn rm(&self, container_id: &str) -> Result<(), String> {
self.removed.lock().unwrap().push(container_id.to_string());
Ok(())
}
fn restart(&self, container_id: &str) -> Result<(), String> {
self.restarted
.lock()
.unwrap()
.push(container_id.to_string());
Ok(())
}
fn is_running(&self, _container_id: &str) -> bool {
true
}
fn container_ip(&self, _container_id: &str) -> Result<String, String> {
Ok("172.17.0.5".to_string())
}
fn health_check(&self, _container_id: &str, _hc: &Healthcheck) -> bool {
!*self.health_should_fail.lock().unwrap()
}
fn logs_tail(&self, _container_id: &str, _lines: u32) -> String {
"boom".to_string()
}
fn list_by_label(&self, _id: &str) -> Vec<String> {
vec![]
}
}
#[derive(Default)]
struct RecordingReporter {
calls: Mutex<Vec<(String, u64, String)>>, meshes: Mutex<Vec<Option<String>>>, }
impl StatusReporter for RecordingReporter {
fn report(
&self,
id: &str,
version: u64,
phase: &str,
_message: &str,
_attempt: u32,
_container_id: Option<&str>,
_endpoint: Option<&str>,
mesh_endpoint: Option<&str>,
) {
self.calls
.lock()
.unwrap()
.push((id.to_string(), version, phase.to_string()));
self.meshes
.lock()
.unwrap()
.push(mesh_endpoint.map(str::to_string));
}
}
#[test]
fn drive_new_deployment_goes_healthy_and_persists_endpoint() {
let rt = FakeRuntime::default();
let reporter = RecordingReporter::default();
let d = desired("dep_1", 1, "running");
let mut local = LocalState::default();
drive(&rt, &reporter, &[d], &mut local);
let rec = local.deployments.get("dep_1").expect("recorded");
assert_eq!(rec.phase, "healthy");
assert_eq!(rec.container_id.as_deref(), Some("c1"));
assert_eq!(rec.endpoint.as_deref(), Some("172.17.0.5:8000"));
assert_eq!(rec.ip.as_deref(), Some("172.17.0.5"));
let calls = reporter.calls.lock().unwrap();
let phases: Vec<&str> = calls.iter().map(|(_, _, p)| p.as_str()).collect();
assert_eq!(phases, vec!["pulling", "starting", "healthy"]);
}
#[test]
fn drive_passes_service_account_to_the_container_as_env_and_persists_it() {
let rt = FakeRuntime::default();
let reporter = RecordingReporter::default();
let mut d = desired("dep_1", 1, "running");
d.service_account = Some(super::super::ServiceAccount {
name: "billing-svc".to_string(),
grants: vec!["kv:read".to_string(), "queue:publish".to_string()],
});
let mut local = LocalState::default();
drive(&rt, &reporter, &[d], &mut local);
let env = rt.last_run_env.lock().unwrap();
assert_eq!(
env.get("ZAKURO_SERVICE_ACCOUNT").map(String::as_str),
Some("billing-svc")
);
assert_eq!(
env.get("ZAKURO_GRANTS").map(String::as_str),
Some("kv:read,queue:publish")
);
let rec = local.deployments.get("dep_1").unwrap();
assert_eq!(
rec.service_account.as_ref().map(|sa| sa.name.as_str()),
Some("billing-svc")
);
}
#[test]
fn drive_version_bump_stops_old_container_after_healthy() {
let rt = FakeRuntime::default();
let reporter = RecordingReporter::default();
let d = desired("dep_1", 2, "running");
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(1, "healthy", Some("old-container")));
drive(&rt, &reporter, &[d], &mut local);
assert!(rt
.stopped
.lock()
.unwrap()
.contains(&"old-container".to_string()));
assert!(rt
.removed
.lock()
.unwrap()
.contains(&"old-container".to_string()));
let rec = local.deployments.get("dep_1").unwrap();
assert_eq!(rec.phase, "healthy");
assert_eq!(rec.version, 2);
assert!(rec.previous_container_id.is_none());
}
#[test]
fn drive_start_failure_leaves_existing_container_untouched() {
let rt = FakeRuntime::default();
*rt.run_should_fail.lock().unwrap() = true;
let reporter = RecordingReporter::default();
let d = desired("dep_1", 2, "running");
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(1, "healthy", Some("old-container")));
drive(&rt, &reporter, &[d], &mut local);
let rec = local.deployments.get("dep_1").unwrap();
assert_eq!(rec.phase, "healthy");
assert_eq!(rec.version, 1);
assert_eq!(rec.container_id.as_deref(), Some("old-container"));
assert!(rt.restarted.lock().unwrap().is_empty());
assert!(rt.stopped.lock().unwrap().is_empty());
assert!(rt.removed.lock().unwrap().is_empty());
let calls = reporter.calls.lock().unwrap();
let phases: Vec<&str> = calls.iter().map(|(_, _, p)| p.as_str()).collect();
assert_eq!(phases, vec!["pulling", "starting", "failed"]);
}
#[test]
fn drive_health_failure_rolls_back_to_previous_container() {
let rt = FakeRuntime::default();
*rt.health_should_fail.lock().unwrap() = true;
let reporter = RecordingReporter::default();
let mut d = desired("dep_1", 2, "running");
d.healthcheck = Some(Healthcheck {
kind: "http".to_string(),
port: 8000,
path: None,
timeout_s: 1,
});
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(1, "healthy", Some("old-container")));
drive(&rt, &reporter, &[d], &mut local);
let rec = local.deployments.get("dep_1").unwrap();
assert_eq!(rec.phase, "rolled_back");
assert_eq!(rec.container_id.as_deref(), Some("old-container"));
assert!(rt
.restarted
.lock()
.unwrap()
.contains(&"old-container".to_string()));
assert!(rt.stopped.lock().unwrap().contains(&"c1".to_string()));
assert!(rt.removed.lock().unwrap().contains(&"c1".to_string()));
assert!(!rt
.removed
.lock()
.unwrap()
.contains(&"old-container".to_string()));
let calls = reporter.calls.lock().unwrap();
let phases: Vec<&str> = calls.iter().map(|(_, _, p)| p.as_str()).collect();
assert_eq!(phases, vec!["pulling", "starting", "failed", "rolled_back"]);
}
#[test]
fn drive_run_failure_with_no_previous_leaves_deployment_failed() {
let rt = FakeRuntime::default();
*rt.run_should_fail.lock().unwrap() = true;
let reporter = RecordingReporter::default();
let d = desired("dep_1", 1, "running");
let mut local = LocalState::default();
drive(&rt, &reporter, &[d], &mut local);
let rec = local.deployments.get("dep_1").unwrap();
assert_eq!(rec.phase, "failed");
assert!(rec.container_id.is_none());
let calls = reporter.calls.lock().unwrap();
let phases: Vec<&str> = calls.iter().map(|(_, _, p)| p.as_str()).collect();
assert_eq!(phases, vec!["pulling", "starting", "failed"]);
}
#[test]
fn drive_stopped_desired_state_stops_and_reports() {
let rt = FakeRuntime::default();
let reporter = RecordingReporter::default();
let d = desired("dep_1", 1, "stopped");
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(1, "healthy", Some("c1")));
drive(&rt, &reporter, &[d], &mut local);
assert!(rt.stopped.lock().unwrap().contains(&"c1".to_string()));
let rec = local.deployments.get("dep_1").unwrap();
assert_eq!(rec.phase, "stopped");
assert!(rec.container_id.is_none());
let calls = reporter.calls.lock().unwrap();
let phases: Vec<&str> = calls.iter().map(|(_, _, p)| p.as_str()).collect();
assert_eq!(phases, vec!["stopped"]);
}
fn exposed_record(ip: &str, container_id: &str) -> DeploymentRecord {
let mut rec = record(1, "healthy", Some(container_id));
rec.ip = Some(ip.to_string());
rec.port = Some(8000);
rec.endpoint = Some(format!("{ip}:8000"));
rec.mesh_endpoint = Some("10.13.13.22:8000".to_string());
rec.mesh_port = Some(8000);
rec
}
#[test]
fn drive_and_expose_reports_a_new_deployment_healthy_then_its_mesh_endpoint() {
let _home = expose::TempZakuroHome::new();
let rt = FakeRuntime::default();
let reporter = RecordingReporter::default();
let exposer = expose::FakeExposer::default();
let mut local = LocalState::default();
drive_and_expose(
&rt,
&reporter,
&exposer,
Some("10.13.13.22"),
&[desired("dep_1", 1, "running")],
&mut local,
);
let calls = reporter.calls.lock().unwrap();
let phases: Vec<&str> = calls.iter().map(|(_, _, p)| p.as_str()).collect();
assert_eq!(phases, vec!["pulling", "starting", "healthy", "healthy"]);
assert_eq!(
*reporter.meshes.lock().unwrap(),
vec![None, None, None, Some("10.13.13.22:8000".to_string())]
);
assert_eq!(
exposer.target_of("dep_1").as_deref(),
Some("172.17.0.5:8000")
);
assert_eq!(
local.deployments["dep_1"].mesh_endpoint.as_deref(),
Some("10.13.13.22:8000")
);
}
#[test]
fn drive_and_expose_re_reports_the_mesh_endpoint_after_a_version_bump() {
let _home = expose::TempZakuroHome::new();
let rt = FakeRuntime::default();
let reporter = RecordingReporter::default();
let exposer = expose::FakeExposer::default();
exposer.expose("dep_1", "10.13.13.22", "172.17.0.4:8000", 8000);
let mut local = LocalState::default();
local.deployments.insert(
"dep_1".into(),
exposed_record("172.17.0.4", "old-container"),
);
drive_and_expose(
&rt,
&reporter,
&exposer,
Some("10.13.13.22"),
&[desired("dep_1", 2, "running")],
&mut local,
);
let last = reporter.calls.lock().unwrap().last().cloned();
assert_eq!(last, Some(("dep_1".to_string(), 2, "healthy".to_string())));
let last_mesh = reporter.meshes.lock().unwrap().last().cloned().flatten();
assert_eq!(last_mesh.as_deref(), Some("10.13.13.22:8000"));
assert_eq!(
exposer.target_of("dep_1").as_deref(),
Some("172.17.0.5:8000")
);
}
#[test]
fn drive_and_expose_releases_a_stopped_deployment() {
let _home = expose::TempZakuroHome::new();
let rt = FakeRuntime::default();
let reporter = RecordingReporter::default();
let exposer = expose::FakeExposer::default();
exposer.expose("dep_1", "10.13.13.22", "172.17.0.4:8000", 8000);
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), exposed_record("172.17.0.4", "c1"));
drive_and_expose(
&rt,
&reporter,
&exposer,
Some("10.13.13.22"),
&[desired("dep_1", 1, "stopped")],
&mut local,
);
assert!(exposer.exposed_ids().is_empty());
assert_eq!(local.deployments["dep_1"].mesh_endpoint, None);
assert_eq!(*reporter.meshes.lock().unwrap(), vec![None]);
}
#[test]
fn drive_stop_reports_the_desired_version_not_the_local_one() {
let _home = expose::TempZakuroHome::new();
let rt = FakeRuntime::default();
let reporter = RecordingReporter::default();
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(1, "healthy", Some("c1")));
drive(
&rt,
&reporter,
&[desired("dep_1", 2, "stopped")],
&mut local,
);
assert_eq!(
*reporter.calls.lock().unwrap(),
vec![("dep_1".to_string(), 2, "stopped".to_string())]
);
}
#[test]
fn drive_handoff_stop_of_an_unlisted_deployment_reports_its_local_version() {
let _home = expose::TempZakuroHome::new();
let rt = FakeRuntime::default();
let reporter = RecordingReporter::default();
let mut d = desired("dep_2", 1, "running");
d.previous = Some(Box::new(super::super::PreviousDeployment {
id: "dep_1".into(),
image: "old-image".into(),
cmd: None,
ports: vec![],
env: HashMap::new(),
healthcheck: None,
}));
let mut local = LocalState::default();
local
.deployments
.insert("dep_1".into(), record(3, "healthy", Some("old-container")));
drive(&rt, &reporter, &[d], &mut local);
let calls = reporter.calls.lock().unwrap();
assert!(
calls.contains(&("dep_1".to_string(), 3, "stopped".to_string())),
"{calls:?}"
);
}
}