use std::collections::HashMap;
use std::process::Command;
use std::time::Duration;
use super::{Healthcheck, Port};
pub struct RunSpec<'a> {
pub name: String,
pub image: &'a str,
pub cmd: Option<&'a str>,
pub ports: &'a [Port],
pub env: &'a HashMap<String, String>,
pub labels: Vec<(String, String)>,
}
pub fn deploy_network() -> String {
std::env::var("ZAKURO_DEPLOY_NETWORK")
.ok()
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.unwrap_or_else(|| "bridge".to_string())
}
pub trait ContainerRuntime {
fn available(&self) -> bool;
fn pull(&self, image: &str) -> Result<(), String>;
fn run(&self, spec: &RunSpec) -> Result<String, String>;
fn stop(&self, container_id: &str) -> Result<(), String>;
fn rm(&self, container_id: &str) -> Result<(), String>;
fn restart(&self, container_id: &str) -> Result<(), String>;
fn is_running(&self, container_id: &str) -> bool;
fn container_ip(&self, container_id: &str) -> Result<String, String>;
fn health_check(&self, container_id: &str, hc: &Healthcheck) -> bool;
fn logs_tail(&self, container_id: &str, lines: u32) -> String;
fn list_by_label(&self, id: &str) -> Vec<String>;
}
pub fn container_name(id: &str, version: u64) -> String {
format!("zk-dep-{id}-v{version}")
}
pub fn deployment_labels(id: &str, version: u64) -> Vec<(String, String)> {
vec![
("zakuro.deployment".to_string(), id.to_string()),
("zakuro.version".to_string(), version.to_string()),
]
}
pub fn health_http_url(host: &str, hc: &Healthcheck) -> String {
let path = hc.path.as_deref().unwrap_or("/");
let path = if path.starts_with('/') {
path.to_string()
} else {
format!("/{path}")
};
format!("http://{host}:{}{path}", hc.port)
}
pub struct DockerCli;
pub fn runtime_available() -> bool {
DockerCli.available()
}
fn docker(args: &[&str]) -> Result<String, String> {
let out = Command::new("docker")
.args(args)
.output()
.map_err(|e| format!("docker not available: {e}"))?;
if out.status.success() {
Ok(String::from_utf8_lossy(&out.stdout).trim().to_string())
} else {
let stderr = String::from_utf8_lossy(&out.stderr).trim().to_string();
Err(if stderr.is_empty() {
format!("docker {} failed", args.first().copied().unwrap_or(""))
} else {
stderr
})
}
}
const CONTAINER_IP_RETRIES: u32 = 10;
const CONTAINER_IP_RETRY_DELAY: Duration = Duration::from_millis(300);
impl ContainerRuntime for DockerCli {
fn available(&self) -> bool {
Command::new("docker")
.arg("info")
.output()
.map(|o| o.status.success())
.unwrap_or(false)
}
fn pull(&self, image: &str) -> Result<(), String> {
docker(&["pull", image]).map(|_| ())
}
fn run(&self, spec: &RunSpec) -> Result<String, String> {
let _ = docker(&["rm", "-f", &spec.name]);
let network = deploy_network();
let mut args: Vec<String> = vec![
"run".into(),
"-d".into(),
"--name".into(),
spec.name.clone(),
"--network".into(),
network,
];
for (k, v) in &spec.labels {
args.push("--label".into());
args.push(format!("{k}={v}"));
}
for (k, v) in spec.env {
args.push("-e".into());
args.push(format!("{k}={v}"));
}
args.push(spec.image.to_string());
if let Some(cmd) = spec.cmd {
args.extend(cmd.split_whitespace().map(|s| s.to_string()));
}
let arg_refs: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
docker(&arg_refs)
}
fn stop(&self, container_id: &str) -> Result<(), String> {
docker(&["stop", container_id]).map(|_| ())
}
fn rm(&self, container_id: &str) -> Result<(), String> {
docker(&["rm", "-f", container_id]).map(|_| ())
}
fn restart(&self, container_id: &str) -> Result<(), String> {
docker(&["start", container_id]).map(|_| ())
}
fn is_running(&self, container_id: &str) -> bool {
docker(&["inspect", "-f", "{{.State.Running}}", container_id])
.map(|out| out == "true")
.unwrap_or(false)
}
fn container_ip(&self, container_id: &str) -> Result<String, String> {
for attempt in 0..CONTAINER_IP_RETRIES {
let out = docker(&[
"inspect",
"-f",
"{{range .NetworkSettings.Networks}}{{.IPAddress}} {{end}}",
container_id,
])?;
if let Some(ip) = out.split_whitespace().find(|s| !s.is_empty()) {
return Ok(ip.to_string());
}
if attempt + 1 < CONTAINER_IP_RETRIES {
std::thread::sleep(CONTAINER_IP_RETRY_DELAY);
}
}
Err(format!(
"container {container_id} has no network IP after {CONTAINER_IP_RETRIES} attempts \
(host network mode, or the container exited)"
))
}
fn health_check(&self, container_id: &str, hc: &Healthcheck) -> bool {
let host = self
.container_ip(container_id)
.unwrap_or_else(|_| "127.0.0.1".to_string());
match hc.kind.as_str() {
"tcp" => std::net::TcpStream::connect((host.as_str(), hc.port)).is_ok(),
_ => ureq::get(health_http_url(&host, hc))
.config()
.timeout_global(Some(Duration::from_secs(3)))
.http_status_as_error(false)
.build()
.call()
.map(|r| (200..300).contains(&r.status().as_u16()))
.unwrap_or(false),
}
}
fn logs_tail(&self, container_id: &str, lines: u32) -> String {
docker(&["logs", "--tail", &lines.to_string(), container_id]).unwrap_or_default()
}
fn list_by_label(&self, id: &str) -> Vec<String> {
docker(&[
"ps",
"-a",
"-q",
"--filter",
&format!("label=zakuro.deployment={id}"),
])
.map(|out| out.lines().map(|l| l.trim().to_string()).collect())
.unwrap_or_default()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn container_name_is_unique_per_version() {
assert_eq!(container_name("dep_1", 3), "zk-dep-dep_1-v3");
assert_ne!(container_name("dep_1", 3), container_name("dep_1", 4));
}
#[test]
fn deployment_labels_carry_id_and_version() {
let labels = deployment_labels("dep_1", 3);
assert!(labels.contains(&("zakuro.deployment".to_string(), "dep_1".to_string())));
assert!(labels.contains(&("zakuro.version".to_string(), "3".to_string())));
}
fn hc(kind: &str, path: Option<&str>) -> Healthcheck {
Healthcheck {
kind: kind.to_string(),
port: 8000,
path: path.map(|s| s.to_string()),
timeout_s: 60,
}
}
#[test]
fn health_http_url_defaults_path_to_root() {
assert_eq!(
health_http_url("172.17.0.5", &hc("http", None)),
"http://172.17.0.5:8000/"
);
}
#[test]
fn health_http_url_tolerates_a_path_without_leading_slash() {
assert_eq!(
health_http_url("172.17.0.5", &hc("http", Some("healthz"))),
"http://172.17.0.5:8000/healthz"
);
}
#[test]
fn health_http_url_keeps_an_already_leading_slash() {
assert_eq!(
health_http_url("172.17.0.5", &hc("http", Some("/api/health"))),
"http://172.17.0.5:8000/api/health"
);
}
#[test]
fn health_http_url_uses_the_container_ip_not_loopback() {
let url = health_http_url("10.13.13.10", &hc("http", None));
assert!(url.starts_with("http://10.13.13.10:"));
assert!(!url.contains("127.0.0.1"));
}
#[test]
fn deploy_network_defaults_to_bridge_when_unset() {
let _lock = crate::credentials::HOME_ENV_LOCK.lock();
let prev = std::env::var_os("ZAKURO_DEPLOY_NETWORK");
std::env::remove_var("ZAKURO_DEPLOY_NETWORK");
assert_eq!(deploy_network(), "bridge");
if let Some(v) = prev {
std::env::set_var("ZAKURO_DEPLOY_NETWORK", v);
}
}
#[test]
fn deploy_network_honours_the_env_override() {
let _lock = crate::credentials::HOME_ENV_LOCK.lock();
let prev = std::env::var_os("ZAKURO_DEPLOY_NETWORK");
std::env::set_var("ZAKURO_DEPLOY_NETWORK", "hubnet");
assert_eq!(deploy_network(), "hubnet");
match prev {
Some(v) => std::env::set_var("ZAKURO_DEPLOY_NETWORK", v),
None => std::env::remove_var("ZAKURO_DEPLOY_NETWORK"),
}
}
#[test]
fn deploy_network_blank_value_falls_back_to_bridge() {
let _lock = crate::credentials::HOME_ENV_LOCK.lock();
let prev = std::env::var_os("ZAKURO_DEPLOY_NETWORK");
std::env::set_var("ZAKURO_DEPLOY_NETWORK", " ");
assert_eq!(deploy_network(), "bridge");
match prev {
Some(v) => std::env::set_var("ZAKURO_DEPLOY_NETWORK", v),
None => std::env::remove_var("ZAKURO_DEPLOY_NETWORK"),
}
}
}