#![allow(dead_code)]
extern crate serde;
extern crate toml;
use std::env;
#[cfg(unix)]
mod agent;
#[cfg(not(unix))]
mod agent {
pub fn run_cli(_: &[String], _: bool) -> i32 {
eprintln!("zc agent is not supported on this platform yet");
2
}
pub mod handoff {
pub fn try_up(_: &[String]) -> Option<i32> {
None
}
pub fn try_down() -> Option<i32> {
None
}
pub fn try_price(_: crate::PriceCmd) -> Option<i32> {
None
}
pub fn price_direct(_: crate::PriceCmd) -> i32 {
eprintln!("zc price needs the hub client, which is not built on this platform yet");
2
}
}
}
mod async_exec;
mod autobroker;
mod broker;
mod common;
mod credentials;
mod dataset;
mod enroll;
mod envs;
mod exec;
mod infer;
mod init;
mod manager;
mod mesh;
mod mesh_dir;
mod model_uri;
mod serve;
mod serve_general;
mod up;
mod update;
mod vpn;
mod wire;
#[cfg(test)]
mod integration_tests;
#[cfg(windows)]
const HOOKS_BIN: &str = "zc-hooks.exe";
#[cfg(not(windows))]
const HOOKS_BIN: &str = "zc-hooks";
fn start_broker(host: Option<&str>, port: Option<u16>, daemon: bool) {
use broker::BrokerConfig;
use colored::Colorize;
broker::apply_user_broker_defaults();
let config = BrokerConfig {
host: host.unwrap_or("0.0.0.0").to_string(),
port: port.unwrap_or(9000),
daemon,
verbose: !daemon, enable_p2p: broker::p2p_default(),
..Default::default()
};
if daemon {
println!("Broker starting on {}:{}...", config.host, config.port);
if let Err(e) = broker::start_server(config) {
eprintln!("Broker error: {}", e);
}
} else {
println!();
println!(
" {}",
"╔═══════════════════════════════════════════╗".cyan()
);
println!(
" {} {} {}",
"║".cyan(),
"Zakuro Compute Broker".bold().white(),
"║".cyan()
);
println!(
" {}",
"╚═══════════════════════════════════════════╝".cyan()
);
println!();
println!(
" {} http://{}:{}",
"Listening:".green(),
config.host,
config.port
);
println!();
println!(" {}", "API Endpoints:".bold());
println!(
" {}",
"─────────────────────────────────────────────".dimmed()
);
println!(" {} /health Health check", "GET ".green());
println!(" {} /workers List workers", "GET ".green());
println!(" {} /workers Register worker", "POST".blue());
println!(
" {} /workers/heartbeat Worker heartbeat",
"POST".blue()
);
println!(
" {} /workers/:id Unregister worker",
"DEL ".red()
);
println!(" {} /execute Execute request", "POST".blue());
println!(" {} /price Estimate price", "POST".blue());
println!(" {} /credits/:user Get balance", "GET ".green());
println!(" {} /credits/:user/add Add credits", "POST".blue());
println!(
" {}",
"─────────────────────────────────────────────".dimmed()
);
println!();
println!(" Press {} to stop.", "Ctrl+C".yellow().bold());
println!();
if let Err(e) = broker::start_server(config) {
eprintln!("{} {}", "Broker error:".red(), e);
}
}
}
fn launch() {
manager::pull();
manager::kill();
manager::restart();
}
fn setup() {
common::download_auth();
common::download_conf();
launch();
}
fn show_user_info(api_key: &str) {
use colored::Colorize;
use std::time::Duration;
let api_url = crate::credentials::default_api_url();
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_connect(Some(Duration::from_secs(10)))
.timeout_global(Some(Duration::from_secs(10)))
.build(),
);
let endpoint = format!("{}/api/auth/me/api-key", api_url.trim_end_matches('/'));
let response = match agent
.get(&endpoint)
.header("Authorization", &format!("Bearer {}", api_key))
.call()
{
Ok(r) => r,
Err(e) => {
eprintln!("{} {}", "Failed to connect to API:".red(), e);
eprintln!(" Endpoint: {}", endpoint);
eprintln!();
eprintln!(" Make sure:");
eprintln!(" • ZAKURO_API_KEY is set to your API key");
eprintln!(" • You have internet connectivity");
eprintln!(" • API is reachable: {}", api_url);
return;
}
};
if response.status().as_u16() != 200 {
eprintln!("{} HTTP {}", "API error:".red(), response.status());
if let Ok(body) = response.into_body().read_to_string() {
eprintln!(" {}", body);
}
return;
}
let body = match response.into_body().read_to_string() {
Ok(s) => s,
Err(e) => {
eprintln!("{} {}", "Failed to read response:".red(), e);
return;
}
};
let data: serde_json::Value = match serde_json::from_str(&body) {
Ok(v) => v,
Err(e) => {
eprintln!("{} {}", "Failed to parse response:".red(), e);
return;
}
};
println!();
if let Some(user_id) = data.get("zakuro_user_id").and_then(|v| v.as_str()) {
println!(" {} {}", "User ID:".bold(), user_id.cyan());
}
if let Some(username) = data.get("username").and_then(|v| v.as_str()) {
println!(" {} {}", "Username:".bold(), username);
}
if let Some(email) = data.get("email").and_then(|v| v.as_str()) {
println!(" {} {}", "Email:".bold(), email);
}
if let Some(balance) = data.get("credits_balance").and_then(|v| v.as_f64()) {
println!(
" {} {}",
"Balance:".bold(),
if balance >= 0.0 {
format!("{:.4} credits", balance).green()
} else {
format!("{:.4} credits", balance).red()
}
);
}
if let Some(disabled) = data.get("disabled").and_then(|v| v.as_bool()) {
println!(
" {} {}",
"Status:".bold(),
if disabled {
"Disabled".red()
} else {
"Active".green()
}
);
}
println!(" {} {}", "API:".bold(), api_url.green());
println!();
}
fn help(full: bool) {
use colored::Colorize;
println!();
println!(
"{} {}",
"zc".bold().cyan(),
"— the zakuro compute client".dimmed()
);
println!("{}: zc [COMMAND]", "Usage".bold());
println!();
println!("{}", "Start here".bold());
println!(" {} Sign in with your browser", "zc login".cyan());
println!(" {} Your account & credits", "zc me".cyan());
println!();
println!("{}", "Connect to the mesh".bold());
println!(
" {} Join the zakuro WireGuard mesh",
"zc connect".cyan()
);
println!(" {} Leave the mesh", "zc disconnect".cyan());
println!(
" {} Show your connection & peer",
"zc status".cyan()
);
println!();
println!("{}", "Use compute".bold());
println!(" {} List available workers", "zc ls".cyan());
println!(
" {} Machine-readable listing (price + capabilities)",
"zc ls --json".cyan()
);
println!(" {} Benchmark the mesh", "zc bench".cyan());
println!(
" {} zc://<uuid> -m \"...\" One-shot model inference",
"zc infer".cyan()
);
println!(
" {} zc://<uuid> Interactive chat with a model",
"zc chat".cyan()
);
println!();
println!("{}", "Share your machine".bold());
println!(
" {} Offer this machine to the mesh",
"zc share".cyan()
);
println!(
" {} [N] Offer N workers (default 1)",
"zc share --workers".cyan()
);
println!(" {} Stop sharing", "zc unshare".cyan());
println!(
" {} Run the menu-bar/widget agent at login",
"zc agent install".cyan()
);
println!(
" {} What the agent is doing (add --json)",
"zc agent status".cyan()
);
println!(
" {} What this Mac charges (also: inherit, default <c/h>)",
"zc price".cyan()
);
println!();
if full {
println!("{}", "Broker".bold());
println!(" broker [host] [port] Start a broker (-d daemon).");
println!(" brokers List mesh brokers by zc://node-<fp> id.");
println!();
println!("{}", "Local Docker mesh".bold());
println!(" mesh up [N] Launch an N-node Docker compute mesh + broker.");
println!(" mesh status | mesh down Inspect / tear down the local mesh.");
println!();
println!("{}", "Headless enrollment".bold());
println!(" token create Mint a one-time join token for a headless node.");
println!(" join <token> Enroll this node using a join token.");
println!();
{
println!("{}", "Zakuro Drive Hooks".bold());
println!(
" hooks <add|get|list|update|remove|logs> Manage drive file-event hooks."
);
println!();
}
println!("{}", "Maintenance".bold());
println!(" update [--from-source] Self-update zc (binary, or build from source).");
println!(" pull | images | ps | kill | restart Manage local zakuro containers.");
println!(" info System info, clusters, and network status.");
println!();
println!("{}", "Aliases".dimmed());
println!(" {} login=init · connect/disconnect/status=vpn · share=up · unshare=down · ls=workers",
"(old names still work)".dimmed());
println!();
} else {
println!(
" {} Show advanced commands (broker, mesh, tokens, maintenance)",
"zc --full".dimmed()
);
println!();
}
println!(
"Guides: {}",
"https://docs.zakuro-ai.com/".cyan().underline()
);
}
fn is_loopback_url(url: &str) -> bool {
let Ok(uri) = url.parse::<http::Uri>() else {
return false;
};
let Some(host) = uri.host() else {
return false;
};
let host = host.trim_start_matches('[').trim_end_matches(']');
host.eq_ignore_ascii_case("localhost")
|| host
.parse::<std::net::IpAddr>()
.is_ok_and(|ip| ip.is_loopback())
}
fn mesh_proxy_for(url: &str) -> Option<ureq::Proxy> {
if is_loopback_url(url) {
None
} else {
vpn::mesh_proxy()
}
}
fn broker_unreachable(err: &ureq::Error) -> bool {
matches!(
err,
ureq::Error::ConnectionFailed
| ureq::Error::HostNotFound
| ureq::Error::Io(_)
| ureq::Error::Timeout(_)
)
}
fn render_workers_table(workers: &[serde_json::Value], total: u64, broker_url: &str) -> String {
use colored::Colorize;
struct Cell {
body: String,
vis: usize,
}
impl Cell {
fn plain(body: impl Into<String>) -> Self {
let body = body.into();
let vis = body.chars().count();
Cell { body, vis }
}
fn painted(body: String, vis: usize) -> Self {
Cell { body, vis }
}
fn pad(&self, width: usize, right: bool) -> String {
let fill = " ".repeat(width.saturating_sub(self.vis));
if right {
format!("{fill}{}", self.body)
} else {
format!("{}{fill}", self.body)
}
}
}
fn num(v: f64, decimals: usize) -> String {
let s = format!("{v:.decimals$}");
match s.contains('.') {
true => s.trim_end_matches('0').trim_end_matches('.').to_string(),
false => s,
}
}
const HEADS: [&str; 9] = [
"WORKER", "STATUS", "CPUs", "GPUs", "MEM", "PRICE", "5h", "1w", "1m",
];
const RIGHT: [bool; 9] = [false, false, true, true, true, true, true, true, true];
const GAP: &str = " ";
const INDENT: &str = " ";
let f64_of = |w: &serde_json::Value, k: &str| w.get(k).and_then(|v| v.as_f64()).unwrap_or(0.0);
let u64_of = |w: &serde_json::Value, k: &str| w.get(k).and_then(|v| v.as_u64()).unwrap_or(0);
let str_of = |w: &serde_json::Value, k: &str| {
w.get(k).and_then(|v| v.as_str()).unwrap_or("?").to_string()
};
let window = |count: u64, quota: u64| -> Cell {
if quota == 0 {
return Cell::plain(format!("{count} / ∞"));
}
let pct = count as f64 / quota as f64 * 100.0;
let s = format!("{count} / {quota} ({pct:.0}%)");
let vis = s.chars().count();
let painted = if pct >= 90.0 {
s.red()
} else if pct >= 70.0 {
s.yellow()
} else {
s.green()
};
Cell::painted(painted.to_string(), vis)
};
let mut groups: std::collections::BTreeMap<String, Vec<&serde_json::Value>> =
std::collections::BTreeMap::new();
for w in workers {
groups.entry(str_of(w, "node")).or_default().push(w);
}
for owned in groups.values_mut() {
owned.sort_by_key(|w| str_of(w, "name"));
}
let row_of = |w: &serde_json::Value| -> Vec<Cell> {
let status = str_of(w, "status");
let status_cell = match status.as_str() {
"healthy" => Cell::painted(status.green().to_string(), status.chars().count()),
"unhealthy" => Cell::painted(status.red().to_string(), status.chars().count()),
"busy" => Cell::painted(status.yellow().to_string(), status.chars().count()),
_ => Cell::plain(status),
};
let price = match w.get("price_per_hour").and_then(|v| v.as_f64()) {
Some(p) => num(p, 3),
None => "—".to_string(),
};
vec![
Cell::plain(str_of(w, "name")),
status_cell,
Cell::plain(num(f64_of(w, "cpus_available"), 1)),
Cell::plain(u64_of(w, "gpus_available").to_string()),
Cell::plain(format!("{}GiB", num(f64_of(w, "memory_available_gib"), 1))),
Cell::plain(price),
window(u64_of(w, "requests_5h"), u64_of(w, "quota_5h")),
window(u64_of(w, "requests_1w"), u64_of(w, "quota_1w")),
window(u64_of(w, "requests_1m"), u64_of(w, "quota_1m")),
]
};
let rows: Vec<(String, Vec<Cell>)> = groups
.iter()
.flat_map(|(broker, owned)| owned.iter().map(|w| (broker.clone(), row_of(w))))
.collect();
let widths: Vec<usize> = HEADS
.iter()
.enumerate()
.map(|(i, h)| {
rows.iter()
.map(|(_, r)| r[i].vis)
.chain(std::iter::once(h.chars().count()))
.max()
.unwrap_or(0)
})
.collect();
let line = |cells: &[Cell]| -> String {
let padded: Vec<String> = cells
.iter()
.enumerate()
.map(|(i, c)| c.pad(widths[i], RIGHT[i]))
.collect();
format!("{INDENT}{}", padded.join(GAP))
};
let table_width = INDENT.len() + widths.iter().sum::<usize>() + GAP.len() * (widths.len() - 1);
let rule = format!(" {}", "─".repeat(table_width - 2).dimmed());
let header: Vec<Cell> = HEADS
.iter()
.map(|h| Cell::painted(h.bold().to_string(), h.chars().count()))
.collect();
let mut out = String::new();
out.push('\n');
out.push_str(&format!(
" {} {} broker(s) · {} worker(s) at {}\n",
"Mesh:".bold(),
groups.len(),
total,
broker_url
));
out.push_str(&format!("{rule}\n"));
out.push_str(&format!("{}\n", line(&header)));
out.push_str(&format!("{rule}\n"));
for (broker, owned) in &groups {
let cpus: f64 = owned.iter().map(|w| f64_of(w, "cpus_available")).sum();
let gpus: u64 = owned.iter().map(|w| u64_of(w, "gpus_available")).sum();
let mem: f64 = owned
.iter()
.map(|w| f64_of(w, "memory_available_gib"))
.sum();
let plural = if owned.len() == 1 {
"worker"
} else {
"workers"
};
let gpu_part = match gpus {
0 => String::new(),
n => format!(" · {n} GPUs"),
};
out.push_str(&format!(
" {} {} {plural} · {} CPUs{gpu_part} · {}GiB\n",
broker.cyan(),
owned.len(),
num(cpus, 1),
num(mem, 1)
));
for w in owned {
out.push_str(&format!("{}\n", line(&row_of(w))));
}
}
out.push_str(&format!("{rule}\n\n"));
out
}
fn show_workers(broker_url: &str, node_filter: Option<&str>, json: bool) {
use colored::Colorize;
use std::time::Duration;
let url = match node_filter {
Some(n) => format!(
"{}/workers?node={}",
broker_url.trim_end_matches('/'),
n.strip_prefix("zc://").unwrap_or(n)
),
None => format!("{}/workers", broker_url.trim_end_matches('/')),
};
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_connect(Some(Duration::from_secs(3)))
.timeout_global(Some(Duration::from_secs(5)))
.proxy(mesh_proxy_for(broker_url))
.build(),
);
let is_local_target = is_loopback_url(broker_url);
let response = match agent.get(&url).call() {
Ok(r) => r,
Err(e) if !broker_unreachable(&e) => {
eprintln!(
"{} broker at {} answered: {}",
"Error:".red(),
broker_url,
e
);
return;
}
Err(e) if !is_local_target => {
eprintln!("{} broker unreachable: {}", "Error:".red(), e);
return;
}
Err(_) => {
let broker_port: u16 = broker_url
.rsplit(':')
.next()
.and_then(|p| p.trim_end_matches('/').parse().ok())
.unwrap_or(9000);
eprintln!(" Broker not running — starting on port {}...", broker_port);
let config = broker::BrokerConfig {
host: "0.0.0.0".to_string(),
port: broker_port,
daemon: true,
verbose: false,
..Default::default()
};
async_exec::spawn_detached(move || {
let _ = broker::start_server(config);
});
let health_url = format!("{}/health", broker_url.trim_end_matches('/'));
let ready = (0..50).any(|_| {
std::thread::sleep(std::time::Duration::from_millis(100));
agent.get(&health_url).call().is_ok()
});
if !ready {
eprintln!("{} Broker failed to start within 5s", "Error:".red());
return;
}
match agent.get(&url).call() {
Ok(r) => r,
Err(e) => {
eprintln!("{} {}", "Error:".red(), e);
return;
}
}
}
};
let body = match response.into_body().read_to_string() {
Ok(s) => s,
Err(e) => {
eprintln!("{} {}", "Error reading response:".red(), e);
return;
}
};
let data: serde_json::Value = match serde_json::from_str(&body) {
Ok(v) => v,
Err(e) => {
eprintln!("{} {}", "Error parsing response:".red(), e);
return;
}
};
let workers = match data.get("workers").and_then(|v| v.as_array()) {
Some(w) => w,
None => {
eprintln!("No workers field in response");
return;
}
};
let total = data
.get("total")
.and_then(|v| v.as_u64())
.unwrap_or(workers.len() as u64);
if json {
println!("{}", serde_json::to_string(&data).unwrap_or_default());
return;
}
print!("{}", render_workers_table(workers, total, broker_url));
}
fn render_brokers_output(data: &serde_json::Value) -> String {
use colored::Colorize;
let mut out = String::new();
let self_id = data.get("self").and_then(|v| v.as_str()).unwrap_or("?");
out.push_str(&format!(" {} {}\n", "Self:".bold(), self_id));
let brokers = data
.get("brokers")
.and_then(|v| v.as_array())
.cloned()
.unwrap_or_default();
if brokers.is_empty() {
out.push_str(" No other brokers known.\n");
return out;
}
out.push_str(&format!(" {} peer broker(s):\n", brokers.len()));
for b in &brokers {
let id = b.get("id").and_then(|v| v.as_str()).unwrap_or("?");
let reachable = b
.get("reachable")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let dot = if reachable {
"●".green().to_string()
} else {
"●".red().to_string()
};
out.push_str(&format!(" {} {}\n", dot, id));
}
out
}
fn show_brokers_cmd() {
use colored::Colorize;
match mesh_dir::directory() {
Ok(list) => {
let offline = list
.iter()
.filter(|b| !b.revoked && b.endpoint.is_none())
.count();
let live: Vec<mesh_dir::MeshBroker> = list
.into_iter()
.filter(|b| !b.revoked && b.endpoint.is_some())
.collect();
let probes = mesh_dir::probe_all(live, std::time::Duration::from_secs(4));
let self_fp = (9000..=9010)
.filter_map(broker::uri::probe_local)
.find_map(|h| h.node_id)
.map(|id| broker::node_identity::strip_node_arg(&id).to_string());
println!();
print!("{}", mesh_dir::render(&probes, self_fp.as_deref(), offline));
println!();
}
Err(e) => {
eprintln!(
" {} mesh directory unavailable ({}); showing the local broker's view",
"note:".yellow(),
e
);
match resolve_default_broker(None) {
Ok(url) => show_brokers(&url),
Err(e) => {
eprintln!("Error: {}", e);
std::process::exit(1);
}
}
}
}
}
fn show_brokers(broker_url: &str) {
use colored::Colorize;
use std::time::Duration;
let url = format!("{}/brokers", broker_url.trim_end_matches('/'));
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_connect(Some(Duration::from_secs(3)))
.timeout_global(Some(Duration::from_secs(5)))
.proxy(mesh_proxy_for(broker_url))
.build(),
);
let response = match agent.get(&url).call() {
Ok(r) => r,
Err(e) => {
eprintln!("{} {}", "Error:".red(), e);
return;
}
};
let body = match response.into_body().read_to_string() {
Ok(s) => s,
Err(e) => {
eprintln!("{} {}", "Error reading response:".red(), e);
return;
}
};
let data: serde_json::Value = match serde_json::from_str(&body) {
Ok(v) => v,
Err(e) => {
eprintln!("{} {}", "Error parsing response:".red(), e);
return;
}
};
println!();
print!("{}", render_brokers_output(&data));
println!();
}
pub(crate) fn resolve_default_broker(explicit_cli: Option<&str>) -> Result<String, String> {
let env_broker = env::var("ZAKURO_BROKER").ok();
if autobroker::should_auto_spawn(explicit_cli, env_broker.as_deref()) {
autobroker::ensure_local_broker(true)
} else {
let raw = explicit_cli
.map(|s| s.to_string())
.or(env_broker)
.unwrap_or_else(|| "zc://localhost".to_string());
broker::uri::resolve(&raw)
}
}
fn run_serve(args: &[String]) {
let parsed = match serve::parse_serve_args(args) {
Ok(p) => p,
Err(e) => {
eprintln!("Error: {}", e);
eprintln!("Usage: zc serve zc://<owner>/<name> [more…] [--price <zkcr_per_mtok>] [--port <base_port>] [--api-url <url>] [--advertise <host>] [--general]");
std::process::exit(2);
}
};
let broker_url = match resolve_default_broker(None) {
Ok(u) => u,
Err(e) => {
eprintln!("Error resolving broker: {}", e);
std::process::exit(1);
}
};
let api_url = parsed
.api_url
.clone()
.unwrap_or_else(credentials::default_api_url);
let auth_bearer = env::var("ZAKURO_API_KEY").ok();
let advertise = parsed.advertise.clone().or_else(|| {
env::var("ZAKURO_ADVERTISE_ADDR")
.ok()
.filter(|v| !v.is_empty())
});
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_connect(Some(std::time::Duration::from_secs(10)))
.timeout_global(Some(std::time::Duration::from_secs(600)))
.proxy(mesh_proxy_for(&broker_url))
.build(),
);
let mut children = Vec::new();
for (i, address) in parsed.model_addresses.iter().enumerate() {
let model_uuid = resolve_model_address_or_exit(address);
let model_uuid = &model_uuid;
let port = parsed.base_port + i as u16;
println!("Serving {} on port {}...", model_uuid, port);
match serve::serve_one_model(
&agent,
&api_url,
&broker_url,
model_uuid,
port,
parsed.price_per_mtok,
auth_bearer.as_deref(),
advertise.as_deref(),
) {
Ok((child, worker_name)) => {
println!(
" registered worker '{}' with broker {}",
worker_name, broker_url
);
children.push(child);
}
Err(e) => {
eprintln!("Error serving {}: {}", model_uuid, e);
}
}
}
if parsed.general {
let general_uri = serve::worker_uri(advertise.as_deref(), parsed.base_port);
let registration = serve::worker_registration_general(
"zc-serve-general",
&general_uri,
parsed.price_per_mtok,
);
match serve::register_worker(&agent, &broker_url, ®istration) {
Err(e) => {
eprintln!("Error registering general provider: {e}");
std::process::exit(1);
}
Ok(worker_id) => {
println!(
"registered general-provider worker with broker {}",
broker_url
);
let hb_agent = agent.clone();
let hb_broker = broker_url.clone();
let _ = std::thread::Builder::new()
.name("zc-serve-general-heartbeat".into())
.spawn(move || {
let never = std::sync::atomic::AtomicBool::new(false);
serve::run_heartbeat_loop(
&never,
serve::HEARTBEAT_INTERVAL,
std::thread::sleep,
|| {
if let Err(e) =
serve::send_heartbeat(&hb_agent, &hb_broker, &worker_id)
{
eprintln!(" [SERVE] general heartbeat failed: {e}");
}
},
);
});
let hooks = serve_general::production_hooks(
agent.clone(),
api_url.clone(),
auth_bearer.clone(),
);
let loader = serve_general::GeneralLoader::new(
hooks,
parsed.base_port + 1,
serve_general::GeneralLoader::max_loaded_from_env(),
);
if let Err(e) = serve_general::run_general_server(loader, parsed.base_port) {
eprintln!("Error running general provider server: {e}");
std::process::exit(1);
}
}
}
}
if children.is_empty() && !parsed.general {
std::process::exit(1);
}
for mut child in children {
let _ = child.wait();
}
}
fn resolve_infer_broker(explicit: Option<&str>) -> Result<String, String> {
match explicit {
Some(url) if url.starts_with("zc://") => broker::uri::resolve(url),
Some(url) => Ok(url.to_string()),
None => resolve_default_broker(None),
}
}
fn bearer_header() -> Option<String> {
credentials::load_into_env();
env::var("ZAKURO_API_KEY")
.ok()
.filter(|k| !k.trim().is_empty())
.map(|k| format!("Bearer {}", k.trim()))
}
fn post_infer(
agent: &ureq::Agent,
broker_url: &str,
body: &serde_json::Value,
) -> Result<serde_json::Value, String> {
let url = format!("{}/infer", broker_url.trim_end_matches('/'));
let mut req = agent.post(&url).header("Content-Type", "application/json");
if let Some(auth) = bearer_header() {
req = req.header("Authorization", &auth);
}
let response = req
.send_json(body)
.map_err(|e| format!("broker request failed: {e}"))?;
let status = response.status();
let text = response
.into_body()
.read_to_string()
.map_err(|e| format!("reading broker response: {e}"))?;
if !status.is_success() {
return Err(text);
}
serde_json::from_str::<serde_json::Value>(&text)
.map_err(|e| format!("parsing broker response: {e} (body: {text})"))
}
fn run_infer(args: &[String]) {
let parsed = match infer::parse_infer_args(args) {
Ok(p) => p,
Err(e) => {
eprintln!("Error: {}", e);
eprintln!("Usage: zc infer zc://<owner>/<name> -m \"<prompt>\" [--max-tokens N] [--system \"<sys>\"] [--broker <url>]");
std::process::exit(2);
}
};
let broker_url = match resolve_infer_broker(parsed.broker.as_deref()) {
Ok(u) => u,
Err(e) => {
eprintln!("Error resolving broker: {}", e);
std::process::exit(1);
}
};
let model_uuid = resolve_model_address_or_exit(&parsed.model);
let body = infer::build_infer_body(
&model_uuid,
&parsed.prompt,
parsed.system.as_deref(),
parsed.max_tokens,
);
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_connect(Some(std::time::Duration::from_secs(10)))
.timeout_global(Some(std::time::Duration::from_secs(300)))
.http_status_as_error(false)
.proxy(mesh_proxy_for(&broker_url))
.build(),
);
match post_infer(&agent, &broker_url, &body) {
Ok(resp) => {
println!("{}", infer::format_infer_output(&resp));
}
Err(e) => {
eprintln!("Error: {}", infer::format_infer_error(&e));
std::process::exit(1);
}
}
}
fn resolve_model_address_or_exit(addr: &crate::model_uri::ModelAddress) -> String {
credentials::load_into_env();
let api_url = std::env::var("ZAKURO_API_URL")
.ok()
.filter(|u| !u.trim().is_empty())
.unwrap_or_else(credentials::default_api_url);
let api_key = std::env::var("ZAKURO_API_KEY").ok();
match crate::model_uri::resolve_address(addr, &api_url, api_key.as_deref()) {
Ok(u) => u,
Err(e) => {
eprintln!("Error: {}", e);
std::process::exit(1);
}
}
}
fn run_chat(args: &[String]) {
use std::io::{self, BufRead, Write};
let parsed = match infer::parse_chat_args(args) {
Ok(p) => p,
Err(e) => {
eprintln!("Error: {}", e);
eprintln!("Usage: zc chat zc://<owner>/<name> [--system \"<sys>\"] [--broker <url>]");
std::process::exit(2);
}
};
let broker_url = match resolve_infer_broker(parsed.broker.as_deref()) {
Ok(u) => u,
Err(e) => {
eprintln!("Error resolving broker: {}", e);
std::process::exit(1);
}
};
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_connect(Some(std::time::Duration::from_secs(10)))
.timeout_global(Some(std::time::Duration::from_secs(300)))
.http_status_as_error(false)
.proxy(mesh_proxy_for(&broker_url))
.build(),
);
let model_uuid = resolve_model_address_or_exit(&parsed.model);
let mut messages: Vec<serde_json::Value> = Vec::new();
if let Some(sys) = &parsed.system {
messages.push(serde_json::json!({"role": "system", "content": sys}));
}
println!(
"zc chat — model zc://{} (broker {})",
model_uuid, broker_url
);
println!("Type your message and press Enter. /quit or Ctrl-D to exit.\n");
let stdin = io::stdin();
loop {
print!("> ");
let _ = io::stdout().flush();
let mut line = String::new();
let bytes_read = match stdin.lock().read_line(&mut line) {
Ok(n) => n,
Err(e) => {
eprintln!("Error reading input: {}", e);
break;
}
};
if bytes_read == 0 {
break;
}
let line = line.trim_end_matches('\n').trim_end_matches('\r');
if line.trim() == "/quit" {
break;
}
if line.trim().is_empty() {
continue;
}
messages.push(serde_json::json!({"role": "user", "content": line}));
let body = infer::build_chat_body(&model_uuid, &messages, None);
match post_infer(&agent, &broker_url, &body) {
Ok(resp) => {
println!("{}", infer::format_infer_output(&resp));
if let Some(content) = resp.get("content").and_then(|v| v.as_str()) {
messages.push(serde_json::json!({"role": "assistant", "content": content}));
}
}
Err(e) => {
eprintln!("Error: {}", infer::format_infer_error(&e));
messages.pop();
}
}
println!();
}
}
fn require_auth() {
eprintln!("This command requires authentication.\n");
eprintln!("Set your token with:");
eprintln!(" export ZAKURO_API_KEY=\"your-token-here\"\n");
eprintln!("Don't have a token? Request access at https://zakuro-ai.com");
}
fn normalize_aliases(mut args: Vec<String>) -> Vec<String> {
let Some(verb) = args.get(1).cloned() else {
return args;
};
let rest = || args[2..].to_vec();
let rebuilt: Option<Vec<String>> = match verb.as_str() {
"login" => Some([vec![args[0].clone(), "init".into()], rest()].concat()),
"share" => Some([vec![args[0].clone(), "up".into()], rest()].concat()),
"unshare" => Some([vec![args[0].clone(), "down".into()], rest()].concat()),
"ls" => Some([vec![args[0].clone(), "workers".into()], rest()].concat()),
"connect" => Some(
[
vec![args[0].clone(), "vpn".into(), "connect".into()],
rest(),
]
.concat(),
),
"disconnect" => Some(
[
vec![args[0].clone(), "vpn".into(), "disconnect".into()],
rest(),
]
.concat(),
),
"status" => Some([vec![args[0].clone(), "vpn".into(), "status".into()], rest()].concat()),
_ => None,
};
if let Some(r) = rebuilt {
args = r;
}
args
}
fn extract_json_flag(mut args: Vec<String>) -> (Vec<String>, bool) {
let present = args.iter().any(|a| a == "--json");
if present {
args.retain(|a| a != "--json");
}
(args, present)
}
fn main() {
credentials::load_into_env();
let _ = rustls::crypto::ring::default_provider().install_default();
envs::update();
let raw_args: Vec<String> = normalize_aliases(env::args().collect());
let (args, json_flag) = extract_json_flag(raw_args);
let full_help = args.iter().any(|a| a == "--full");
if args
.get(1)
.map(|a| a == "-h" || a == "--help")
.unwrap_or(false)
{
help(full_help);
return;
}
if args.get(1).map(|a| a == "--full").unwrap_or(false) {
help(true);
return;
}
if args.len() == 1 {
help(full_help);
return;
}
let zakuro_auth = env::var("ZAKURO_API_KEY").ok();
let _broker_url = env::var("ZAKURO_BROKER").unwrap_or_else(|_| "zc://localhost".to_string());
if args.get(1).map(|a| a == "up").unwrap_or(false) {
let up_args: Vec<String> = args[2..].to_vec();
if let Some(code) = agent::handoff::try_up(&up_args) {
std::process::exit(code);
}
let (workers, base_port, broker_port, daemon) = up::parse_args(&up_args);
up::start_up(workers, base_port, broker_port, daemon);
return;
}
if args.get(1).map(|a| a == "mesh").unwrap_or(false) {
match args.get(2).map(|s| s.as_str()).unwrap_or("status") {
"up" => {
let n = args
.get(3)
.and_then(|s| s.parse::<usize>().ok())
.unwrap_or(4);
mesh::up(n);
}
"down" => mesh::down(),
"status" | "ls" | "ps" => mesh::status(),
other => eprintln!("usage: zc mesh [up <N>|down|status] (got '{}')", other),
}
return;
}
if args.get(1).map(|a| a == "agent").unwrap_or(false) {
std::process::exit(agent::run_cli(&args[2..], json_flag));
}
if args.get(1).map(|a| a == "deployments").unwrap_or(false) {
broker::deploy::print_deployments();
return;
}
if args.get(1).map(|a| a == "vpn").unwrap_or(false) {
vpn::run_cli(&args[2..]);
return;
}
if args.get(1).map(|a| a == "update").unwrap_or(false) {
std::process::exit(update::run_cli(&args[2..]));
}
if args.get(1).map(|a| a == "down").unwrap_or(false) {
if let Some(code) = agent::handoff::try_down() {
std::process::exit(code);
}
let down_args: Vec<String> = args[2..].to_vec();
let (_, base_port, broker_port, _) = up::parse_args(&down_args);
up::stop_up(broker_port, base_port, base_port + 39);
return;
}
if args.get(1).map(|a| a == "hooks").unwrap_or(false) {
let sibling = std::env::current_exe()
.ok()
.and_then(|p| p.parent().map(|d| d.join(HOOKS_BIN)));
let sidecar = sibling
.clone()
.filter(|p| p.is_file())
.unwrap_or_else(|| std::path::PathBuf::from(HOOKS_BIN));
match std::process::Command::new(&sidecar)
.args(&args[2..])
.status()
{
Ok(st) => std::process::exit(st.code().unwrap_or(1)),
Err(e) => {
eprintln!("`zc hooks` needs the `{HOOKS_BIN}` helper, which was not found.");
match &sibling {
Some(p) => eprintln!(" tried: {} then $PATH", p.display()),
None => eprintln!(" tried: $PATH"),
}
eprintln!(" reason: {e}");
eprintln!();
eprintln!("Install it alongside zc, or build it from the repo:");
eprintln!(" cargo build --release -p zc-hooks (requires zakuro-drive access)");
std::process::exit(2);
}
}
}
if args.get(1).map(|a| a == "init").unwrap_or(false) {
std::process::exit(init::run());
}
if args.get(1).map(|a| a == "token").unwrap_or(false)
&& args.get(2).map(|a| a == "create").unwrap_or(false)
{
std::process::exit(enroll::run_token_create(&args[3..]));
}
if args.get(1).map(|a| a == "join").unwrap_or(false) {
std::process::exit(enroll::run_join(&args[2..]));
}
if args.get(1).map(|a| a == "serve").unwrap_or(false) {
run_serve(&args[2..]);
return;
}
if args.get(1).map(|a| a == "infer").unwrap_or(false) {
run_infer(&args[2..]);
return;
}
if args.get(1).map(|a| a == "chat").unwrap_or(false) {
run_chat(&args[2..]);
return;
}
if args.get(1).map(|a| a == "dataset").unwrap_or(false) {
std::process::exit(dataset::run(&args[2..]));
}
match args.len() {
1 => unreachable!(),
2 => {
let arg0 = &args[1];
match &arg0[..] {
"-h" => help(false),
"--help" => help(false),
"ps" => manager::ps(),
"dist" => match common::dist() {
Ok(res) => {
println!("{}", res);
}
Err(why) => {
eprintln!("{}", why);
}
},
"images" => manager::images(),
"launch" => launch(),
"download_conf" => common::download_conf(),
"download_auth" => common::download_auth(),
"setup" => setup(),
"pull" => manager::pull(),
"workers" => {
let broker_url = match resolve_default_broker(None) {
Ok(u) => u,
Err(e) => {
eprintln!("Error: {}", e);
std::process::exit(1);
}
};
show_workers(&broker_url, None, json_flag);
}
"brokers" => show_brokers_cmd(),
"price" => std::process::exit(run_price(&args[2..])),
"logs" => common::logs(true),
"nodes" => manager::nodes(),
"restart" => manager::restart(),
"servers" => manager::server_list(),
"add_worker" => manager::add_worker(),
"kill" => manager::kill(),
"connect" => manager::connect(),
"rmi" => manager::rmi(),
"--version" => common::version(),
"-v" => common::version(),
"vars" => {
if let Err(e) = common::context(None) {
eprintln!("Error reading context: {}", e);
}
}
"broker" => {
start_broker(None, None, false); }
"up" => {
up::start_up(1, 3960, 9000, false);
}
"bench" => {
broker::bench::run_from_args(&[]);
}
"info" => {
broker::info::print_info();
}
"me" | "whoami" | "credits" => match &zakuro_auth {
Some(auth) => show_user_info(auth),
None => require_auth(),
},
_ => help(false),
}
}
3 => {
let arg0 = &args[1];
let arg1 = &args[2];
match &arg0[..] {
"--docker" => match &arg1[..] {
"rm" => manager::remove_container(),
_ => {
let _ = exec::zk0(&arg1[..]);
}
},
"-d" | "--daemon" => match &arg1[..] {
"rm" => manager::remove_container(),
"broker" => {
start_broker(None, None, true); }
_ => {
let _ = exec::zk0(&arg1[..]);
}
},
"push" => {
manager::push(Some(&arg1[..]));
}
"context" => {
if let Err(e) = common::context(Some(&arg1[..])) {
eprintln!("Error setting context: {}", e);
}
}
"build" => {
common::build(Some(args));
}
"broker" => {
let port: u16 = arg1.parse().unwrap_or(9000);
start_broker(None, Some(port), false); }
"workers" => {
if let Some(node) = arg1
.strip_prefix("zc://")
.filter(|rest| rest.starts_with("node-"))
{
let broker_url = match resolve_default_broker(None) {
Ok(u) => u,
Err(e) => {
eprintln!("Error: {}", e);
std::process::exit(1);
}
};
show_workers(&broker_url, Some(node), json_flag);
} else {
let broker_url = match broker::uri::resolve(arg1) {
Ok(u) => u,
Err(e) => {
eprintln!("Error: {}", e);
std::process::exit(1);
}
};
show_workers(&broker_url, None, json_flag);
}
}
"price" => std::process::exit(run_price(&args[2..])),
"bench" => {
broker::bench::run_from_args(std::slice::from_ref(arg1));
}
_ => help(false),
}
}
4 => {
let arg0 = &args[1];
let arg1 = &args[2];
let arg2 = &args[3];
match &arg0[..] {
"-d" | "--daemon" => match &arg1[..] {
"broker" => {
let port: u16 = arg2.parse().unwrap_or(9000);
start_broker(None, Some(port), true); }
_ => help(false),
},
"broker" => {
let port: u16 = arg2.parse().unwrap_or(9000);
start_broker(Some(arg1), Some(port), false); }
"price" => std::process::exit(run_price(&args[2..])),
"bench" => {
broker::bench::run_from_args(&args[2..]);
}
_ => help(false),
}
}
_ => {
let arg0 = &args[1];
match &arg0[..] {
"build" => {
common::build(Some(args));
}
"bench" => {
let bench_args: Vec<String> = args[2..].to_vec();
broker::bench::run_from_args(&bench_args);
}
"up" => {
let up_args: Vec<String> = args[2..].to_vec();
let (workers, base_port, broker_port, daemon) = up::parse_args(&up_args);
up::start_up(workers, base_port, broker_port, daemon);
}
_ => {
help(false);
}
}
}
}
}
#[cfg(test)]
mod alias_tests {
use super::normalize_aliases;
fn norm(cmd: &[&str]) -> Vec<String> {
normalize_aliases(cmd.iter().map(|s| s.to_string()).collect())
}
#[test]
fn maps_primary_verbs_to_canonical() {
assert_eq!(norm(&["zc", "login"]), vec!["zc", "init"]);
assert_eq!(
norm(&["zc", "share", "--workers", "4"]),
vec!["zc", "up", "--workers", "4"]
);
assert_eq!(norm(&["zc", "unshare"]), vec!["zc", "down"]);
assert_eq!(norm(&["zc", "ls"]), vec!["zc", "workers"]);
}
#[test]
fn mesh_verbs_expand_to_vpn_subcommand() {
assert_eq!(
norm(&["zc", "connect", "--docker"]),
vec!["zc", "vpn", "connect", "--docker"]
);
assert_eq!(norm(&["zc", "disconnect"]), vec!["zc", "vpn", "disconnect"]);
assert_eq!(norm(&["zc", "status"]), vec!["zc", "vpn", "status"]);
}
#[test]
fn unknown_and_empty_pass_through() {
assert_eq!(norm(&["zc", "me"]), vec!["zc", "me"]);
assert_eq!(norm(&["zc"]), vec!["zc"]);
assert_eq!(
norm(&["zc", "vpn", "connect"]),
vec!["zc", "vpn", "connect"]
);
}
}
#[cfg(test)]
mod brokers_cli_tests {
use super::render_brokers_output;
use serde_json::json;
fn sample_brokers_json_with_internal_ips() -> serde_json::Value {
json!({
"self": "zc://node-aaaaaaaaaaaaaaaa",
"brokers": [
{"id": "zc://node-deadbeefcafebabe", "reachable": true, "peer_url": "http://10.13.13.5:9000"},
{"id": "zc://node-0123456789abcdef", "reachable": false, "peer_url": "http://127.0.0.1:9001"}
]
})
}
#[test]
fn brokers_output_is_ip_free() {
let body = render_brokers_output(&sample_brokers_json_with_internal_ips());
assert!(!body.contains("10.13.13."));
assert!(!body.contains("127.0.0.1"));
assert!(body.contains("zc://node-"));
}
}
#[cfg(test)]
mod json_listing_tests {
use super::extract_json_flag;
use serde_json::json;
fn strip(cmd: &[&str]) -> (Vec<String>, bool) {
extract_json_flag(cmd.iter().map(|s| s.to_string()).collect())
}
#[test]
fn detects_and_strips_trailing_json_flag() {
let (args, present) = strip(&["zc", "workers", "--json"]);
assert!(present);
assert_eq!(args, vec!["zc", "workers"]);
}
#[test]
fn detects_and_strips_json_flag_before_positional_args() {
let (args, present) = strip(&["zc", "workers", "--json", "zc://node-abc123"]);
assert!(present);
assert_eq!(args, vec!["zc", "workers", "zc://node-abc123"]);
}
#[test]
fn absent_json_flag_leaves_args_untouched() {
let (args, present) = strip(&["zc", "workers"]);
assert!(!present);
assert_eq!(args, vec!["zc", "workers"]);
}
#[test]
fn worker_listing_json_carries_price_and_capability_fields() {
let body = json!({
"total": 1,
"workers": [{
"id": "w1",
"name": "worker-1",
"uri": "zc://worker-node-abc-1",
"worker_type": "cpu",
"status": "healthy",
"cpus_available": 8.0,
"cpus_total": 8.0,
"memory_available_gib": 32.0,
"memory_total_gib": 32.0,
"gpus_available": 0,
"gpus_total": 0,
"price_per_hour": 3.6,
"min_charge": 0.0,
"active_requests": 0,
"avg_latency_ms": 0.0,
"max_timeout_secs": 60.0,
"requests_5h": 0, "requests_1w": 0, "requests_1m": 0,
"quota_5h": 0, "quota_1w": 0, "quota_1m": 0,
"node": "zc://node-abc123",
}],
});
let workers = body.get("workers").and_then(|v| v.as_array()).unwrap();
let w = &workers[0];
assert!(w.get("node").and_then(|v| v.as_str()).is_some());
assert!(w.get("price_per_hour").and_then(|v| v.as_f64()).is_some());
assert!(w.get("cpus_available").and_then(|v| v.as_f64()).is_some());
assert!(w
.get("memory_available_gib")
.and_then(|v| v.as_f64())
.is_some());
assert!(!body.to_string().contains("127.0.0.1"));
}
}
fn run_price(args: &[String]) -> i32 {
match parse_price_args(args) {
Ok(cmd) => {
agent::handoff::try_price(cmd).unwrap_or_else(|| agent::handoff::price_direct(cmd))
}
Err(e) => {
eprintln!("Error: {e}");
eprintln!("{PRICE_USAGE}");
2
}
}
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub enum PriceCmd {
Show,
Device(f64),
Inherit,
Default(f64),
}
impl PriceCmd {
pub fn value(self) -> Option<f64> {
match self {
PriceCmd::Device(v) | PriceCmd::Default(v) => Some(v),
PriceCmd::Show | PriceCmd::Inherit => None,
}
}
}
pub const PRICE_USAGE: &str = "Usage: zc price What this Mac charges
zc price <credits/hour> Set what this Mac charges
zc price inherit Follow the account default again
zc price default <c/h> Set the account default for every device";
fn parse_price_args(args: &[String]) -> Result<PriceCmd, String> {
let words: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
match words[..] {
[] => Ok(PriceCmd::Show),
["inherit"] => Ok(PriceCmd::Inherit),
["default"] => Err("`zc price default` needs a value in credits/hour".to_string()),
["default", v] => parse_price_arg(v).map(PriceCmd::Default),
[v] => parse_price_arg(v).map(PriceCmd::Device),
_ => Err(format!(
"`zc price {}` is not a price command",
words.join(" ")
)),
}
}
fn parse_price_arg(s: &str) -> Result<f64, String> {
match s.parse::<f64>() {
Ok(v) if v.is_finite() => Ok(v),
Ok(_) => Err(format!(
"invalid price value '{s}': the price must be a finite number of credits/hour"
)),
Err(_) => Err(format!("invalid price value '{s}'")),
}
}
#[cfg(test)]
mod price_command_tests {
use super::{parse_price_args, PriceCmd};
fn parse(words: &[&str]) -> Result<PriceCmd, String> {
let owned: Vec<String> = words.iter().map(|w| w.to_string()).collect();
parse_price_args(&owned)
}
#[test]
fn each_form_maps_to_one_command() {
assert_eq!(parse(&[]), Ok(PriceCmd::Show));
assert_eq!(parse(&["18"]), Ok(PriceCmd::Device(18.0)));
assert_eq!(parse(&["3.6"]), Ok(PriceCmd::Device(3.6)));
assert_eq!(parse(&["inherit"]), Ok(PriceCmd::Inherit));
assert_eq!(parse(&["default", "42"]), Ok(PriceCmd::Default(42.0)));
}
#[test]
fn a_form_that_would_clear_or_mangle_a_price_is_refused() {
for bad in [
vec!["nan"],
vec!["inf"],
vec!["-inf"],
vec!["abc"],
vec![""],
vec!["default"],
vec!["default", "nan"],
vec!["default", "abc"],
vec!["inherit", "18"],
vec!["18", "42"],
vec!["default", "42", "18"],
] {
assert!(parse(&bad).is_err(), "{bad:?} must be refused");
}
assert!(parse(&["default"]).unwrap_err().contains("credits/hour"));
}
}
#[cfg(test)]
mod price_arg_tests {
use super::parse_price_arg;
#[test]
fn a_price_argument_must_be_a_finite_number() {
assert_eq!(parse_price_arg("18"), Ok(18.0));
assert_eq!(parse_price_arg("3.6"), Ok(3.6));
for bad in ["nan", "NaN", "inf", "-inf", "infinity", "abc", ""] {
assert!(parse_price_arg(bad).is_err(), "{bad:?} must be refused");
}
let e = parse_price_arg("nan").unwrap_err();
assert!(e.contains("finite"), "{e}");
}
}
#[cfg(test)]
mod broker_target_tests {
use super::{broker_unreachable, is_loopback_url};
#[test]
fn only_this_machine_counts_as_a_loopback_broker() {
for (url, want) in [
("http://localhost:9001", true),
("http://LOCALHOST:9000", true),
("http://127.0.0.1:9001/workers", true),
("http://127.0.0.53:9000", true),
("http://[::1]:9000/health", true),
("http://10.13.13.5:9000", false),
("https://stg.api.zakuro-ai.com", false),
("http://localhost.example.com:9000", false),
("http://my-127.0.0.1.example:9000", false),
("not a url", false),
] {
assert_eq!(is_loopback_url(url), want, "{url}");
}
}
#[test]
fn only_a_broker_that_never_answered_is_unreachable() {
use std::io::{Error, ErrorKind};
assert!(broker_unreachable(&ureq::Error::ConnectionFailed));
assert!(broker_unreachable(&ureq::Error::HostNotFound));
assert!(broker_unreachable(&ureq::Error::Io(Error::from(
ErrorKind::ConnectionRefused
))));
assert!(broker_unreachable(&ureq::Error::Timeout(
ureq::Timeout::Connect
)));
assert!(!broker_unreachable(&ureq::Error::StatusCode(500)));
assert!(!broker_unreachable(&ureq::Error::StatusCode(404)));
}
}
#[cfg(test)]
mod workers_table_tests {
use super::render_workers_table;
fn visible(table: &str) -> Vec<String> {
let mut out = Vec::new();
for line in table.lines() {
let mut s = String::new();
let mut chars = line.chars();
while let Some(c) = chars.next() {
if c == '\u{1b}' {
for e in chars.by_ref() {
if e == 'm' {
break;
}
}
} else {
s.push(c);
}
}
out.push(s);
}
out
}
fn grid(table: &str) -> Vec<String> {
visible(table)
.into_iter()
.filter(|l| l.contains('─') || l.contains("WORKER") || l.starts_with(" "))
.collect()
}
fn with_node(mut w: serde_json::Value, node: &str) -> serde_json::Value {
w["node"] = serde_json::json!(node);
w
}
fn worker(name: &str, uri: &str, cpus: f64, mem: f64, price: f64) -> serde_json::Value {
serde_json::json!({
"name": name, "status": "healthy", "uri": uri,
"node": "zc://node-dd715b7f6d94d44e",
"cpus_available": cpus, "gpus_available": 0,
"memory_available_gib": mem, "price_per_hour": price,
})
}
fn staging_rows() -> Vec<serde_json::Value> {
vec![
worker(
"dd715b7f-w0",
"zc://worker-dd715b7f6d94d44e-3960",
10.0,
3.4,
20.0,
),
worker(
"worker-zc-worker-cloudcompute-2",
"zc://worker-af185d3fcc82649f-3960",
8.0,
15.5,
3.6,
),
]
}
#[test]
fn every_row_lines_up_with_the_header() {
let rows = staging_rows();
let table = render_workers_table(&rows, rows.len() as u64, "http://localhost:9000");
let grid = grid(&table);
assert!(
grid.len() >= 5,
"two rules, a header and two rows: {grid:?}"
);
let widths: Vec<usize> = grid.iter().map(|l| l.chars().count()).collect();
assert!(
widths.windows(2).all(|w| w[0] == w[1]),
"every line of the grid is the same width: {widths:?}\n{table}"
);
let header = grid
.iter()
.find(|l| l.contains("WORKER"))
.expect("a header");
assert!(
!header.contains("BROKER"),
"no per-row broker column: {header}"
);
assert!(
widths[0] <= 130,
"the table stays within a wide terminal: {} columns\n{table}",
widths[0]
);
let name_at = header.find("WORKER").expect("a WORKER column");
let price_end = header.find("PRICE").expect("a PRICE column") + "PRICE".len();
let mut worker_lines = grid
.iter()
.filter(|l| l.starts_with(" ") && !l.contains("WORKER"));
for w in rows.iter() {
let line = worker_lines.next().expect("a row per worker");
let name = w["name"].as_str().unwrap();
assert_eq!(
&line[name_at..name_at + name.len()],
name,
"NAME starts in its column: {line}"
);
assert!(
line[..price_end].ends_with(&["20", "3.6"][usize::from(name.len() > 11)]),
"PRICE ends in its column: {line}"
);
}
}
#[test]
fn workers_are_listed_under_the_broker_that_owns_them() {
let i9 = "zc://node-e18d8c0ed6ec7b5a";
let x399 = "zc://node-6188100f5daf463e";
let rows = vec![
with_node(
worker(
"worker-zc-worker-i9-1",
"zc://worker-e18d8c0ed6ec7b5a-3960",
20.0,
31.0,
3.6,
),
i9,
),
with_node(
worker(
"worker-zc-worker-x399-1",
"zc://worker-6188100f5daf463e-3960",
32.0,
62.6,
3.6,
),
x399,
),
with_node(
worker(
"worker-zc-worker-i9-2",
"zc://worker-d2f66b53e3fe399f-3960",
20.0,
31.0,
3.6,
),
i9,
),
];
let table = render_workers_table(&rows, rows.len() as u64, "http://localhost:9000");
let lines = visible(&table);
assert!(
!table.contains("zc://worker-"),
"no worker address is offered as a target:\n{table}"
);
for id in [i9, x399] {
assert_eq!(
lines.iter().filter(|l| l.contains(id)).count(),
1,
"{id} names its group once:\n{table}"
);
}
let at = |needle: &str| {
lines
.iter()
.position(|l| l.contains(needle))
.unwrap_or_else(|| panic!("{needle} is listed:\n{table}"))
};
assert_eq!(
at(i9) + 1,
at("worker-zc-worker-i9-1"),
"the i9 broker heads its own workers:\n{table}"
);
assert_eq!(
at("worker-zc-worker-i9-2"),
at("worker-zc-worker-i9-1") + 1,
"and they stay together:\n{table}"
);
assert_eq!(
at(x399) + 1,
at("worker-zc-worker-x399-1"),
"the x399 broker heads its own worker:\n{table}"
);
let head = lines
.iter()
.find(|l| l.contains("broker"))
.expect("a summary line");
assert!(
head.contains('2') && head.contains('3'),
"2 brokers, 3 workers: {head}"
);
}
#[test]
fn numbers_carry_no_digits_they_do_not_need() {
let rows = staging_rows();
let table = render_workers_table(&rows, rows.len() as u64, "http://localhost:9000");
for noise in ["20.000", "3.600", "10.0 ", "8.0 "] {
assert!(
!table.contains(noise),
"{noise:?} is padding, not information:\n{table}"
);
}
let priced = vec![worker("w0", "zc://worker-a-3960", 2.0, 1.0, 3.625)];
assert!(
render_workers_table(&priced, 1, "http://localhost:9000").contains("3.625"),
"a real fraction survives"
);
}
}