#[allow(warnings)]
extern crate serde;
extern crate toml;
use std::env;
mod async_exec;
mod broker;
mod common;
mod exec;
mod envs;
mod manager;
mod up;
#[cfg(test)]
mod integration_tests;
fn update() {
let _ = exec::tty("curl https://get.zakuro-ai.com/zc | sh");
}
fn start_broker(host: Option<&str>, port: Option<u16>, daemon: bool, tui: bool) {
use broker::BrokerConfig;
use colored::Colorize;
let config = BrokerConfig {
host: host.unwrap_or("0.0.0.0").to_string(),
port: port.unwrap_or(9000),
daemon,
verbose: !daemon, tui_mode: tui,
..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 = std::env::var("ZAKURO_API_URL")
.unwrap_or_else(|_| "https://my.zakuro-ai.com".to_string());
let agent = ureq::AgentBuilder::new()
.timeout_connect(Duration::from_secs(10))
.timeout_read(Duration::from_secs(10))
.build();
let endpoint = format!("{}/api/auth/me/api-key", api_url.trim_end_matches('/'));
let response = match agent
.get(&endpoint)
.set("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() != 200 {
eprintln!("{} HTTP {}", "API error:".red(), response.status());
if let Ok(body) = response.into_string() {
eprintln!(" {}", body);
}
return;
}
let body = match response.into_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 [OPTIONS] COMMAND", "Usage".bold());
println!();
println!("A self-sufficient runtime for zakuro");
println!();
println!("{}", "Options:".bold());
println!(" -d, --daemon Run broker in background (daemon mode).");
println!(" -t, --tui Run broker with interactive terminal dashboard.");
println!(" -v, --version Get the version of the current command line.");
println!(" -h, --help Print this help.");
println!(" --full Show all commands including legacy/broken ones.");
println!();
println!("{}", "Broker Commands:".bold());
println!(" broker Start broker with live transaction log.");
println!(" broker <port> Start broker on specified port.");
println!(" broker <host> <port> Start broker on specified host and port.");
println!(" -t broker Start broker with interactive TUI dashboard.");
println!(" -d broker Start broker in daemon mode (background).");
println!();
println!("{}", "Cluster Commands:".bold());
println!(" up Start 1 worker + broker locally (foreground).");
println!(" up --workers 4 Start 4 workers + broker (ports 3960-3963).");
println!(" up --workers 4 -d Start 4 workers in the background (daemon).");
println!(" up --workers 4 --port 4000 Start 4 workers on ports 4000-4003.");
println!(" up --broker-port 9001 Start broker on a custom port.");
println!(" down Stop all workers + broker started by 'zc up'.");
println!(" workers List workers on the local broker (ZAKURO_BROKER or zc://localhost).");
println!(" workers zc://node-name List workers on a named broker.");
println!(" workers zc://ip:port List workers on a broker at a specific address.");
println!();
println!("{}", "Benchmark Commands:".bold());
println!(" bench Benchmark broker (default: best_price strategy).");
println!(" bench -s best_latency Benchmark with fastest worker routing.");
println!(" bench --compare Compare all routing strategies.");
println!(" bench -c 50 -n 10000 Benchmark with 50 concurrent, 10k requests.");
println!(" bench zc://node-name Benchmark a specific named broker.");
println!();
println!("{}", "Identity:".bold());
println!(" me Show authenticated user info and credits.");
println!(" credits Alias for 'me'.");
println!();
println!("{}", "Monitoring:".bold());
println!(" attach zc://node-name Attach to a named broker's TUI.");
println!(" attach zc://localhost Attach to the local broker TUI.");
println!();
println!("{}", "Diagnostics:".bold());
println!(" info Show system info, clusters, and network status.");
println!(" show-network-auth Print Tailscale auth key from dashboard API (requires ZAKURO_API_KEY).");
println!();
println!(" The broker automatically discovers workers on the Tailscale network");
println!();
println!("{}", "Container Commands:".bold());
println!(" update Update the command line.");
println!(" pull Pull updated images.");
println!(" images List zakuro images built on the machine.");
println!(" ps List current running zakuro containers.");
println!(" kill Remove current running zakuro containers.");
println!(" restart Restart the containers with updated images.");
println!();
if full {
println!("{}", "Legacy Commands (broken - require zk0 container):".bold().red());
println!(" connect [BROKEN] Enter zk0 in interactive mode.");
println!(" rmi [BROKEN] Remove zakuro images (no confirmation, destructive).");
println!(" context <path> [BROKEN] Set new zakuro context.");
println!();
}
println!("To get more help, check out our guides at {}", "https://docs.zakuro.ai/".cyan().underline());
}
fn show_workers(broker_url: &str) {
use colored::Colorize;
use std::time::Duration;
let url = format!("{}/workers", broker_url.trim_end_matches('/'));
let agent = ureq::AgentBuilder::new()
.timeout_connect(Duration::from_secs(3))
.timeout_read(Duration::from_secs(5))
.build();
let response = match agent.get(&url).call() {
Ok(r) => r,
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_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);
println!();
println!(" {} {} worker(s) at {}", "Workers:".bold(), total, broker_url);
println!(
" {}",
"──────────────────────────────────────────────────────────────────────────────────────────────".dimmed()
);
println!(
" {:<20} {:<10} {:<26} {:>5} {:>5} {:>12} {:>12} {:>12}",
"NAME".bold(), "STATUS".bold(), "URI".bold(), "CPUs".bold(), "GPUs".bold(),
"5h".bold(), "1w".bold(), "1m".bold()
);
println!(
" {}",
"──────────────────────────────────────────────────────────────────────────────────────────────".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 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 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_colored = match status {
"healthy" => status.green(),
"unhealthy" => status.red(),
"busy" => status.yellow(),
_ => status.normal(),
};
let fmt_window = |count: u64, quota: u64| -> String {
if quota == 0 {
format!("{} / ∞", count)
} else {
let pct = count as f64 / quota as f64 * 100.0;
let s = format!("{} / {} ({:.0}%)", count, quota, pct);
if pct >= 90.0 {
s.red().to_string()
} else if pct >= 70.0 {
s.yellow().to_string()
} else {
s.green().to_string()
}
}
};
println!(
" {:<20} {:<19} {:<26} {:>5.1} {:>5} {:>21} {:>21} {:>21}",
name, status_colored, uri, cpus, gpus,
fmt_window(r5h, q5h),
fmt_window(r1w, q1w),
fmt_window(r1m, q1m),
);
}
println!(
" {}",
"──────────────────────────────────────────────────────────────────────────────────────────────".dimmed()
);
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");
}
fn main() {
envs::update();
let args: Vec<String> = env::args().collect();
let full_help = args.iter().any(|a| a == "--full");
if args.len() == 1 || 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;
}
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 == "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;
}
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 raw = env::var("ZAKURO_BROKER")
.unwrap_or_else(|_| "zc://localhost".to_string());
let broker_url = match broker::uri::resolve(&raw) {
Ok(u) => u,
Err(e) => { eprintln!("Error: {}", e); std::process::exit(1); }
};
show_workers(&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(),
"update" => update(),
"rmi" => manager::rmi(),
"--version" => common::version(),
"-v" => common::version(),
"vars" => {
common::context(None);
},
"broker" => {
start_broker(None, None, false, 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(),
}
},
"show-network-auth" => {
match broker::get_tailscale_auth_key() {
Ok(key) => println!("{}", key),
Err(e) => {
use colored::Colorize;
eprintln!("{} {}", "Error:".red(), e);
std::process::exit(1);
}
}
},
_ => {
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" => {
common::context(Some(&arg1[..]));
},
"build" => {
common::build(Some(args));
}
"broker" => {
let port: u16 = arg1.parse().unwrap_or(9000);
start_broker(None, Some(port), false, false); }
"workers" => {
let broker_url = match broker::uri::resolve(&arg1) {
Ok(u) => u,
Err(e) => { eprintln!("Error: {}", e); std::process::exit(1); }
};
show_workers(&broker_url);
}
"bench" => {
broker::bench::run_from_args(&[arg1.clone()]);
}
"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); }
_ => 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);
}
}
}
}
}