#![allow(dead_code)]
extern crate serde;
extern crate toml;
use std::env;
mod async_exec;
mod autobroker;
mod broker;
mod command;
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 tui;
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, tui: 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, tui_mode: tui,
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 if tui {
if let Err(e) = broker::start_server(config) {
eprintln!("{} {}", "Broker error:".red(), 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. Use {} for interactive TUI.",
"Ctrl+C".yellow().bold(),
"--tui".cyan()
);
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!(
" {} Open the guided terminal (default)",
"zc".cyan()
);
println!(
" {} Sign in — opens the setup wizard",
"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!();
if full {
println!("{}", "Broker".bold());
println!(" broker [host] [port] Start a broker (-t TUI dashboard, -d daemon).");
println!(" brokers List mesh brokers by zc://node-<fp> id.");
println!(" attach zc://<node> Attach to a broker's live TUI.");
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 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(vpn::mesh_proxy())
.build(),
);
let is_local_target = broker_url.contains("localhost") || broker_url.contains("127.0.0.1");
let response = match agent.get(&url).call() {
Ok(r) => r,
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,
tui_mode: 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;
}
const W_NAME: usize = 18;
const W_STATUS: usize = 9; const W_NODE: usize = 16; const W_ADDR: usize = 24; const W_CPU: usize = 5;
const W_GPU: usize = 4;
const W_MEM: usize = 9; const W_PRICE: usize = 9; const W_WIN: usize = 11;
fn cell(body: &str, vis: usize, width: usize, right: bool) -> String {
let pad = " ".repeat(width.saturating_sub(vis));
if right {
format!("{pad}{body}")
} else {
format!("{body}{pad}")
}
}
let rule: String = "─".repeat(
W_NAME + W_STATUS + W_NODE + W_ADDR + W_CPU + W_GPU + W_MEM + W_PRICE + W_WIN * 3 + 10, );
println!();
println!(
" {} {} worker(s) at {}",
"Workers:".bold(),
total,
broker_url
);
println!(" {}", rule.dimmed());
println!(
" {} {} {} {} {} {} {} {} {} {} {}",
cell(&"NAME".bold().to_string(), 4, W_NAME, false),
cell(&"STATUS".bold().to_string(), 6, W_STATUS, false),
cell(&"NODE".bold().to_string(), 4, W_NODE, false),
cell(&"ADDRESS".bold().to_string(), 7, W_ADDR, false),
cell(&"CPUs".bold().to_string(), 4, W_CPU, true),
cell(&"GPUs".bold().to_string(), 4, W_GPU, true),
cell(&"MEM".bold().to_string(), 3, W_MEM, true),
cell(&"PRICE".bold().to_string(), 5, W_PRICE, true),
cell(&"5h".bold().to_string(), 2, W_WIN, true),
cell(&"1w".bold().to_string(), 2, W_WIN, true),
cell(&"1m".bold().to_string(), 2, W_WIN, true),
);
println!(" {}", rule.dimmed());
for w in workers {
let name = w.get("name").and_then(|v| v.as_str()).unwrap_or("?");
let status = w.get("status").and_then(|v| v.as_str()).unwrap_or("?");
let uri = w.get("uri").and_then(|v| v.as_str()).unwrap_or("?");
let owned_addr;
let addr: &str = if uri.starts_with("zc://") {
uri
} else {
owned_addr = format!("zc://{}", name);
&owned_addr
};
let node = w.get("node").and_then(|v| v.as_str()).unwrap_or("—");
let cpus = w
.get("cpus_available")
.and_then(|v| v.as_f64())
.unwrap_or(0.0);
let gpus = w
.get("gpus_available")
.and_then(|v| v.as_u64())
.unwrap_or(0);
let mem_gib = w
.get("memory_available_gib")
.and_then(|v| v.as_f64())
.unwrap_or(0.0);
let price = w.get("price_per_hour").and_then(|v| v.as_f64());
let r5h = w.get("requests_5h").and_then(|v| v.as_u64()).unwrap_or(0);
let r1w = w.get("requests_1w").and_then(|v| v.as_u64()).unwrap_or(0);
let r1m = w.get("requests_1m").and_then(|v| v.as_u64()).unwrap_or(0);
let q5h = w.get("quota_5h").and_then(|v| v.as_u64()).unwrap_or(0);
let q1w = w.get("quota_1w").and_then(|v| v.as_u64()).unwrap_or(0);
let q1m = w.get("quota_1m").and_then(|v| v.as_u64()).unwrap_or(0);
let (status_body, status_vis) = match status {
"healthy" => (status.green().to_string(), status.len()),
"unhealthy" => (status.red().to_string(), status.len()),
"busy" => (status.yellow().to_string(), status.len()),
_ => (status.to_string(), status.len()),
};
let fmt_window = |count: u64, quota: u64| -> (String, usize) {
if quota == 0 {
let s = format!("{} / ∞", count);
let vis = s.chars().count();
(s, vis)
} else {
let pct = count as f64 / quota as f64 * 100.0;
let s = format!("{} / {} ({:.0}%)", count, quota, pct);
let vis = s.chars().count();
let colored = if pct >= 90.0 {
s.red().to_string()
} else if pct >= 70.0 {
s.yellow().to_string()
} else {
s.green().to_string()
};
(colored, vis)
}
};
let cpu_s = format!("{:.1}", cpus);
let gpu_s = gpus.to_string();
let mem_s = format!("{:.1}GiB", mem_gib);
let price_s = match price {
Some(p) => format!("{:.3}", p),
None => "—".to_string(),
};
let (w5, v5) = fmt_window(r5h, q5h);
let (ww, vw) = fmt_window(r1w, q1w);
let (wm, vm) = fmt_window(r1m, q1m);
println!(
" {} {} {} {} {} {} {} {} {} {} {}",
cell(name, name.chars().count(), W_NAME, false),
cell(&status_body, status_vis, W_STATUS, false),
cell(node, node.chars().count(), W_NODE, false),
cell(addr, addr.chars().count(), W_ADDR, false),
cell(&cpu_s, cpu_s.chars().count(), W_CPU, true),
cell(&gpu_s, gpu_s.chars().count(), W_GPU, true),
cell(&mem_s, mem_s.chars().count(), W_MEM, true),
cell(&price_s, price_s.chars().count(), W_PRICE, true),
cell(&w5, v5, W_WIN, true),
cell(&ww, vw, W_WIN, true),
cell(&wm, vm, W_WIN, true),
);
}
println!(" {}", rule.dimmed());
println!();
}
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(vpn::mesh_proxy())
.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!();
}
fn show_price(broker_url: &str) {
use colored::Colorize;
use std::time::Duration;
let url = format!("{}/price", 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(vpn::mesh_proxy())
.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;
}
};
match data.get("price_per_hour").and_then(|v| v.as_f64()) {
Some(price) => println!("current price: {} credits/hour", price),
None => eprintln!("{} unexpected response: {}", "Error:".red(), body),
}
}
fn set_price(broker_url: &str, value: f64) {
use colored::Colorize;
use std::time::Duration;
let url = format!("{}/price/set", 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(vpn::mesh_proxy())
.build(),
);
let body = serde_json::json!({ "price_per_hour": value }).to_string();
let response = match agent
.post(&url)
.header("Content-Type", "application/json")
.send(body.as_bytes())
{
Ok(r) => r,
Err(e) => {
eprintln!("{} {}", "Error:".red(), e);
return;
}
};
let resp_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(&resp_body) {
Ok(v) => v,
Err(e) => {
eprintln!("{} {}", "Error parsing response:".red(), e);
return;
}
};
match data.get("price_per_hour").and_then(|v| v.as_f64()) {
Some(price) => println!("price set: {} credits/hour", price),
None => eprintln!("{} unexpected response: {}", "Error:".red(), resp_body),
}
}
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(vpn::mesh_proxy())
.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(vpn::mesh_proxy())
.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(vpn::mesh_proxy())
.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 {
use std::io::IsTerminal;
if std::io::stdout().is_terminal() {
let has_key = env::var("ZAKURO_API_KEY")
.map(|k| !k.trim().is_empty())
.unwrap_or(false);
if !has_key {
tui::onboard::run();
}
if let Err(e) = tui::run_shell() {
broker::cleanup_terminal();
eprintln!("shell error: {}", e);
}
} else {
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();
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 == "vpn").unwrap_or(false) {
vpn::run_cli(&args[2..]);
return;
}
if args.get(1).map(|a| a == "update").unwrap_or(false) {
update::run_cli(&args[2..]);
return;
}
if args.get(1).map(|a| a == "down").unwrap_or(false) {
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) {
use std::io::IsTerminal;
let force = args.iter().any(|a| a == "--force");
let has_key = env::var("ZAKURO_API_KEY")
.map(|k| !k.trim().is_empty())
.unwrap_or(false);
if std::io::stdout().is_terminal() && (force || !has_key) {
let connected = tui::onboard::run();
std::process::exit(if connected { 0 } else { 1 });
}
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" => {
let broker_url = match resolve_default_broker(None) {
Ok(u) => u,
Err(e) => {
eprintln!("Error: {}", e);
std::process::exit(1);
}
};
show_price(&broker_url);
}
"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, false); }
"shell" | "repl" => {
if let Err(e) = tui::run_shell() {
broker::cleanup_terminal();
eprintln!("shell error: {}", e);
}
}
"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, false); }
_ => {
let _ = exec::zk0(&arg1[..]);
}
},
"-t" | "--tui" => match &arg1[..] {
"broker" => {
start_broker(None, None, false, true); }
_ => {
help(false);
}
},
"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, 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" => {
match arg1.parse::<f64>() {
Ok(value) => {
let broker_url = match resolve_default_broker(None) {
Ok(u) => u,
Err(e) => {
eprintln!("Error: {}", e);
std::process::exit(1);
}
};
set_price(&broker_url, value);
}
Err(_) => {
eprintln!("Error: invalid price value '{}'", arg1);
std::process::exit(1);
}
}
}
"bench" => {
broker::bench::run_from_args(std::slice::from_ref(arg1));
}
"attach" => {
match &zakuro_auth {
Some(auth) => {
let url = match broker::uri::resolve(arg1) {
Ok(u) => u,
Err(e) => {
eprintln!("Error: {}", e);
std::process::exit(1);
}
};
let default_panic = std::panic::take_hook();
std::panic::set_hook(Box::new(move |info| {
broker::cleanup_terminal();
default_panic(info);
}));
println!("Attaching to broker at {}...", url);
let api_key = Some(auth.clone());
if let Err(e) = broker::run_remote_tui(url, api_key) {
broker::cleanup_terminal();
eprintln!("TUI error: {}", e);
}
}
None => require_auth(),
}
}
_ => 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, false); }
_ => help(false),
},
"-t" | "--tui" => {
match &arg1[..] {
"broker" => {
let port: u16 = arg2.parse().unwrap_or(9000);
if up::is_port_open(port) {
use colored::Colorize;
println!(" {} Port {} is already in use — attaching to existing broker...", "•".cyan(), port);
let url = format!("http://localhost:{}", port);
match &zakuro_auth {
Some(api_key) => {
let default_panic = std::panic::take_hook();
std::panic::set_hook(Box::new(move |info| {
broker::cleanup_terminal();
default_panic(info);
}));
if let Err(e) =
broker::run_remote_tui(url, Some(api_key.clone()))
{
broker::cleanup_terminal();
eprintln!("{} {}", "TUI error:".red(), e);
}
}
None => {
eprintln!(
" {} Set {} to attach to the broker.",
"Error:".red(),
"ZAKURO_API_KEY".yellow()
);
}
}
} else {
start_broker(None, Some(port), false, true); }
}
_ => help(false),
}
}
"broker" => {
let port: u16 = arg2.parse().unwrap_or(9000);
start_broker(Some(arg1), Some(port), false, false); }
"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"));
}
}