pub mod expose;
pub mod reconcile;
pub mod runtime;
pub mod serve;
pub mod state;
use std::sync::Arc;
use std::time::Duration;
use serde::{Deserialize, Serialize};
#[allow(unused_imports)]
pub use reconcile::{plan, plan_failure, Action};
#[allow(unused_imports)]
pub use runtime::{runtime_available, ContainerRuntime, DockerCli, RunSpec};
#[allow(unused_imports)]
pub use state::DeploymentRecord;
pub use state::LocalState;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Port {
pub container: u16,
#[serde(default = "default_protocol")]
pub protocol: String,
}
fn default_protocol() -> String {
"http".to_string()
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Healthcheck {
#[serde(rename = "type")]
pub kind: String, pub port: u16,
#[serde(default)]
pub path: Option<String>,
#[serde(default = "default_timeout_s")]
pub timeout_s: u64,
}
fn default_timeout_s() -> u64 {
60
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ServiceAccount {
pub name: String,
#[serde(default)]
pub grants: Vec<String>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct PreviousDeployment {
pub id: String,
pub image: String,
#[serde(default)]
pub cmd: Option<String>,
#[serde(default)]
pub ports: Vec<Port>,
#[serde(default)]
pub env: std::collections::HashMap<String, String>,
#[serde(default)]
pub healthcheck: Option<Healthcheck>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Desired {
pub id: String,
pub version: u64,
pub name: String,
pub image: String,
#[serde(default)]
pub cmd: Option<String>,
#[serde(default)]
pub ports: Vec<Port>,
#[serde(default)]
pub env: std::collections::HashMap<String, String>,
#[serde(default)]
pub healthcheck: Option<Healthcheck>,
pub price_per_second: f64,
pub desired_state: String, #[serde(default)]
pub service_account: Option<ServiceAccount>,
#[serde(default = "default_warm")]
pub warm: bool,
#[serde(default)]
pub previous: Option<Box<PreviousDeployment>>,
}
pub(crate) fn default_warm() -> bool {
true
}
#[derive(Debug, Deserialize)]
struct DesiredResponse {
deployments: Vec<Desired>,
}
pub fn fetch_desired(
api_url: &str,
api_key: &str,
node_pubkey: &str,
) -> Result<Vec<Desired>, String> {
let endpoint = format!(
"{}/api/broker/deployments/desired",
api_url.trim_end_matches('/')
);
let resp = ureq::get(&endpoint)
.config()
.timeout_global(Some(Duration::from_secs(10)))
.http_status_as_error(false)
.build()
.header("X-Broker-Api-Key", api_key)
.header("X-Node-Pubkey", node_pubkey)
.call()
.map_err(|e| e.to_string())?;
if resp.status().as_u16() != 200 {
return Err(format!("fetch_desired HTTP {}", resp.status().as_u16()));
}
let text = resp
.into_body()
.read_to_string()
.map_err(|e| e.to_string())?;
let parsed: DesiredResponse = serde_json::from_str(&text).map_err(|e| e.to_string())?;
Ok(parsed.deployments)
}
#[allow(clippy::too_many_arguments)]
fn status_payload(
node_pubkey: &str,
version: u64,
phase: &str,
message: &str,
attempt: u32,
container_id: Option<&str>,
endpoint: Option<&str>,
mesh_endpoint: Option<&str>,
) -> serde_json::Value {
serde_json::json!({
"node_pubkey": node_pubkey,
"version": version,
"phase": phase,
"message": message,
"attempt": attempt,
"container_id": container_id,
"endpoint": endpoint,
"mesh_endpoint": mesh_endpoint,
})
}
#[allow(clippy::too_many_arguments)]
pub fn report_status(
api_url: &str,
api_key: &str,
node_pubkey: &str,
id: &str,
version: u64,
phase: &str,
message: &str,
attempt: u32,
container_id: Option<&str>,
endpoint: Option<&str>,
mesh_endpoint: Option<&str>,
) {
let url = format!(
"{}/api/broker/deployments/{}/status",
api_url.trim_end_matches('/'),
id
);
let payload = status_payload(
node_pubkey,
version,
phase,
message,
attempt,
container_id,
endpoint,
mesh_endpoint,
);
let body = serde_json::to_string(&payload).unwrap_or_default();
let result = ureq::post(&url)
.config()
.timeout_global(Some(Duration::from_secs(10)))
.http_status_as_error(false)
.build()
.header("X-Broker-Api-Key", api_key)
.header("X-Node-Pubkey", node_pubkey)
.header("Content-Type", "application/json")
.send(body.as_str());
if let Err(e) = result {
eprintln!(" [DEPLOY] failed to report status for {id}: {e}");
}
}
pub struct HttpReporter<'a> {
pub api_url: &'a str,
pub api_key: &'a str,
pub node_pubkey: &'a str,
}
impl reconcile::StatusReporter for HttpReporter<'_> {
fn report(
&self,
id: &str,
version: u64,
phase: &str,
message: &str,
attempt: u32,
container_id: Option<&str>,
endpoint: Option<&str>,
mesh_endpoint: Option<&str>,
) {
report_status(
self.api_url,
self.api_key,
self.node_pubkey,
id,
version,
phase,
message,
attempt,
container_id,
endpoint,
mesh_endpoint,
);
}
}
pub fn tick(state: &Arc<crate::broker::BrokerState>) {
let (Some(api_url), Some(api_key)) = (&state.config.api_url, &state.config.api_key) else {
return;
};
let node_pubkey = state.node_key.public_b64();
let desired = match fetch_desired(api_url, api_key, &node_pubkey) {
Ok(d) => d,
Err(e) => {
eprintln!(" [DEPLOY] failed to fetch desired deployments: {e}");
return;
}
};
if desired.is_empty() {
return;
}
let rt = DockerCli;
if !rt.available() {
for d in &desired {
if d.desired_state == "running" {
report_status(
api_url,
api_key,
&node_pubkey,
&d.id,
d.version,
"failed",
"no container runtime on this host",
1,
None,
None,
None,
);
}
}
return;
}
let mut local = LocalState::load();
let reporter = HttpReporter {
api_url,
api_key,
node_pubkey: &node_pubkey,
};
let mesh_ip = expose::exposure_ip(
crate::broker::discovery::get_mesh_ip(),
std::env::var("ZAKURO_MESH_EXPOSE").ok().as_deref(),
);
if mesh_ip.is_some() {
expose::refresh_allowlist(api_url, api_key);
}
reconcile::drive_and_expose(
&rt,
&reporter,
expose::forwarder(),
mesh_ip.as_deref(),
&desired,
&mut local,
);
}
pub fn print_deployments() {
let local = LocalState::load();
if local.deployments.is_empty() {
println!("No local deployment state (nothing reconciled on this node yet).");
} else {
println!(
"{:<24} {:>7} {:<12} {:<22} {:<22} {:<6} {:<20} CONTAINER",
"DEPLOYMENT", "VERSION", "PHASE", "ENDPOINT", "MESH", "WARM", "SERVICE_ACCOUNT"
);
let mut ids: Vec<&String> = local.deployments.keys().collect();
ids.sort();
for id in ids {
let rec = &local.deployments[id];
let service_account = rec
.service_account
.as_ref()
.map(|sa| format!("{}[{}]", sa.name, sa.grants.join(",")))
.unwrap_or_else(|| "-".to_string());
println!(
"{:<24} {:>7} {:<12} {:<22} {:<22} {:<6} {:<20} {}",
id,
rec.version,
rec.phase,
rec.endpoint.as_deref().unwrap_or("-"),
rec.mesh_endpoint.as_deref().unwrap_or("-"),
rec.warm,
service_account,
rec.container_id.as_deref().unwrap_or("-"),
);
}
}
println!();
let rt = DockerCli;
if !rt.available() {
println!("(docker not available on this host)");
return;
}
match std::process::Command::new("docker")
.args([
"ps",
"-a",
"--filter",
"label=zakuro.deployment",
"--format",
"table {{.Names}}\t{{.Status}}\t{{.Ports}}",
])
.output()
{
Ok(out) if out.status.success() => {
print!("{}", String::from_utf8_lossy(&out.stdout));
}
Ok(out) => {
eprintln!(
"docker ps failed: {}",
String::from_utf8_lossy(&out.stderr).trim()
);
}
Err(e) => eprintln!("docker ps failed: {e}"),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn status_payload_carries_the_mesh_endpoint() {
let payload = status_payload(
"pk",
3,
"healthy",
"",
1,
Some("c1"),
Some("172.17.0.3:8888"),
Some("10.13.13.22:8888"),
);
assert_eq!(payload["endpoint"], "172.17.0.3:8888");
assert_eq!(payload["mesh_endpoint"], "10.13.13.22:8888");
}
#[test]
fn status_payload_sends_an_explicit_null_when_not_exposed() {
let payload = status_payload("pk", 3, "stopped", "", 1, None, None, None);
assert!(payload.get("mesh_endpoint").is_some_and(|v| v.is_null()));
}
}