use std::net::TcpStream;
use std::time::Duration;
use colored::Colorize;
#[derive(Debug, Clone)]
pub struct ClusterInfo {
pub name: &'static str,
pub uri: String,
pub available: bool,
pub version: Option<String>,
pub nodes: Option<u32>,
}
#[derive(Debug, Clone)]
pub struct NetworkInfo {
pub local_ip: Option<String>,
pub wireguard_ip: Option<String>,
pub hostname: String,
pub mode: NetworkMode,
}
#[derive(Debug, Clone, Copy)]
pub enum NetworkMode {
WireGuard,
Local,
Unknown,
}
impl NetworkMode {
pub fn as_str(&self) -> &'static str {
match self {
NetworkMode::WireGuard => "WireGuard P2P",
NetworkMode::Local => "Local",
NetworkMode::Unknown => "Unknown",
}
}
}
pub fn detect_clusters() -> Vec<ClusterInfo> {
vec![detect_ray(), detect_dask(), detect_spark()]
}
fn detect_ray() -> ClusterInfo {
let default_uri = "ray://127.0.0.1:10001";
let mut available = false;
let mut version = None;
let mut nodes = None;
if let Ok(resp) = ureq::get("http://127.0.0.1:8265/api/cluster_status")
.config()
.timeout_global(Some(Duration::from_secs(2)))
.build()
.call()
{
if let Ok(body) = resp.into_body().read_to_string() {
if let Ok(json) = serde_json::from_str::<serde_json::Value>(&body) {
available = true;
if let Some(active) =
json["data"]["clusterStatus"]["autoscalerReport"]["activeNodes"].as_object()
{
let n: u64 = active.values().filter_map(|v| v.as_u64()).sum();
nodes = Some(n as u32);
}
}
}
}
if available {
version = ureq::get("http://127.0.0.1:8265/api/version")
.config()
.timeout_global(Some(Duration::from_secs(2)))
.build()
.call()
.ok()
.and_then(|r| r.into_body().read_to_string().ok())
.and_then(|b| serde_json::from_str::<serde_json::Value>(&b).ok())
.and_then(|j| j["ray_version"].as_str().map(|s| s.to_string()))
.or_else(|| Some("2.x".to_string()));
}
ClusterInfo {
name: "Ray",
uri: default_uri.to_string(),
available,
version,
nodes,
}
}
fn detect_dask() -> ClusterInfo {
let default_uri = "dask://127.0.0.1:8786";
let mut available = false;
let mut version = None;
let mut nodes = None;
if let Ok(resp) = ureq::get("http://127.0.0.1:8787/json/identity.json")
.config()
.timeout_global(Some(Duration::from_secs(2)))
.build()
.call()
{
if let Ok(body) = resp.into_body().read_to_string() {
if let Ok(json) = serde_json::from_str::<serde_json::Value>(&body) {
if json["type"].as_str() == Some("Scheduler") {
available = true;
nodes = json["n_workers"].as_u64().map(|n| n as u32);
version = json["version"]
.as_str()
.map(|s| s.to_string())
.or_else(|| Some("distributed".to_string()));
}
}
}
}
ClusterInfo {
name: "Dask",
uri: default_uri.to_string(),
available,
version,
nodes,
}
}
fn detect_spark() -> ClusterInfo {
let default_uri = "spark://127.0.0.1:7077";
let available = check_port("127.0.0.1", 7077);
let mut version = None;
let mut nodes = None;
if available {
for ui_port in [8080u16, 8082] {
if let Ok(resp) = ureq::get(&format!("http://127.0.0.1:{}/json/", ui_port))
.config()
.timeout_global(Some(Duration::from_secs(2)))
.build()
.call()
{
if let Ok(body) = resp.into_body().read_to_string() {
if let Ok(json) = serde_json::from_str::<serde_json::Value>(&body) {
nodes = json["aliveworkers"].as_u64().map(|n| n as u32);
version = Some("3.x".to_string());
break;
}
}
}
}
}
ClusterInfo {
name: "Spark",
uri: default_uri.to_string(),
available,
version,
nodes,
}
}
fn check_port(host: &str, port: u16) -> bool {
let addr: std::net::SocketAddr = match format!("{}:{}", host, port).parse() {
Ok(a) => a,
Err(_) => return false,
};
TcpStream::connect_timeout(&addr, Duration::from_secs(1)).is_ok()
}
pub fn get_network_info() -> NetworkInfo {
let hostname = hostname::get()
.map(|h| h.to_string_lossy().to_string())
.unwrap_or_else(|_| "unknown".to_string());
let local_ip = get_local_ip();
let wireguard_ip = get_wireguard_ip();
let mode = if wireguard_ip.is_some() {
NetworkMode::WireGuard
} else if local_ip.is_some() {
NetworkMode::Local
} else {
NetworkMode::Unknown
};
NetworkInfo {
local_ip,
wireguard_ip,
hostname,
mode,
}
}
fn get_local_ip() -> Option<String> {
if let Ok(output) = std::process::Command::new("hostname").arg("-I").output() {
if output.status.success() {
let ips = String::from_utf8_lossy(&output.stdout);
return ips.split_whitespace().next().map(|s| s.to_string());
}
}
if let Ok(socket) = std::net::UdpSocket::bind("0.0.0.0:0") {
if socket.connect("8.8.8.8:80").is_ok() {
if let Ok(addr) = socket.local_addr() {
return Some(addr.ip().to_string());
}
}
}
None
}
fn get_wireguard_ip() -> Option<String> {
for iface in ["zakuro0", "wg0"] {
let Ok(output) = std::process::Command::new("ip")
.args(["addr", "show", iface])
.output()
else {
continue;
};
if !output.status.success() {
continue;
}
let output_str = String::from_utf8_lossy(&output.stdout);
for line in output_str.lines() {
if line.contains("inet ") && !line.contains("inet6") {
if let Some(inet_part) = line.split("inet ").nth(1) {
if let Some(ip) = inet_part.split('/').next() {
return Some(ip.trim().to_string());
}
}
}
}
}
None
}
pub fn print_info() {
println!();
println!(
" {}",
"╔═══════════════════════════════════════════╗".cyan()
);
println!(
" {} {} {}",
"║".cyan(),
"Zakuro System Info".bold().white(),
"║".cyan()
);
println!(
" {}",
"╚═══════════════════════════════════════════╝".cyan()
);
println!();
let network = get_network_info();
println!(" {}", "Network".bold());
println!(" {}", "─".repeat(50));
println!(" Hostname: {}", network.hostname);
println!(
" Mode: {}",
match network.mode {
NetworkMode::WireGuard => network.mode.as_str().green(),
NetworkMode::Local => network.mode.as_str().yellow(),
NetworkMode::Unknown => network.mode.as_str().red(),
}
);
if let Some(ip) = &network.local_ip {
println!(" Local IP: {}", ip);
}
if let Some(ip) = &network.wireguard_ip {
println!(" WireGuard IP: {}", ip.cyan());
}
println!();
let clusters = detect_clusters();
println!(" {}", "Compute Clusters".bold());
println!(" {}", "─".repeat(50));
for cluster in &clusters {
let status = if cluster.available {
"●".green()
} else {
"○".red()
};
let version_str = cluster.version.as_deref().unwrap_or("-");
let nodes_str = cluster
.nodes
.map(|n| format!("{} {}", n, if n == 1 { "node" } else { "nodes" }))
.unwrap_or("-".to_string());
println!(
" {} {:8} {:20} {} {}",
status,
cluster.name,
cluster.uri.dimmed(),
version_str,
if cluster.available {
nodes_str
} else {
"not running".dimmed().to_string()
}
);
}
println!();
println!(" {}", "Configuration".bold());
println!(" {}", "─".repeat(50));
let zakuro_auth = std::env::var("ZAKURO_API_KEY").ok();
let api_url = crate::credentials::default_api_url();
println!(" API URL: {}", api_url.cyan());
println!(
" ZAKURO_API_KEY: {}",
if zakuro_auth.is_some() {
"✓ set".green()
} else {
"✗ not set".red()
}
);
println!();
println!(" {}", "Services".bold());
println!(" {}", "─".repeat(50));
let services = [
("Production Broker", "hub.zakuro-ai.com", 443_u16, true),
("Local Broker", "127.0.0.1", 9000, false),
("Local Worker", "127.0.0.1", 3960, false),
];
for (name, host, port, is_https) in services {
let available = if is_https {
true
} else {
check_port(host, port)
};
let status = if available {
"●".green()
} else {
"○".dimmed()
};
let port_display = if is_https {
"https".to_string()
} else {
format!(":{}", port)
};
println!(
" {} {:20} {}{}",
status,
name,
host.dimmed(),
port_display.dimmed()
);
}
println!();
let mesh_ports: Vec<u16> = (9001..=9009)
.filter(|&p| check_port("127.0.0.1", p))
.collect();
if !mesh_ports.is_empty() {
println!(" {}", "Mesh Nodes".bold());
println!(" {}", "─".repeat(50));
for port in mesh_ports {
let url = format!("http://127.0.0.1:{}/health", port);
match ureq::get(&url)
.config()
.timeout_global(Some(Duration::from_millis(500)))
.build()
.call()
{
Ok(resp) => {
if let Ok(body) = resp.into_body().read_to_string() {
if let Ok(v) = serde_json::from_str::<serde_json::Value>(&body) {
let node_name = v["node_name"].as_str().unwrap_or("unknown");
let ts_ip = v["wireguard_ip"].as_str();
let ts_connected = v["wireguard_connected"].as_bool().unwrap_or(false);
let ts_status = if ts_connected {
format!("WireGuard {}", ts_ip.unwrap_or(""))
.green()
.to_string()
} else {
"no WireGuard".red().to_string()
};
println!(
" {} {:20} :{} {}",
"●".green(),
node_name,
port,
ts_status
);
continue;
}
}
println!(" {} :{}", "●".green(), port);
}
Err(_) => {
println!(" {} :{}", "○".dimmed(), port);
}
}
}
println!();
}
}
#[cfg(test)]
mod panic_fix_tests {
use super::check_port;
#[test]
fn check_port_bad_host_returns_false_not_panic() {
assert!(!check_port("not-an-ip-host", 9000));
assert!(!check_port("", 9000));
}
#[test]
fn check_port_unreachable_ip_returns_false() {
assert!(!check_port("127.0.0.1", 1));
}
}