use std::net::TcpStream;
use std::path::PathBuf;
use std::process::{Child, Command};
use std::sync::{
atomic::{AtomicBool, Ordering},
Arc,
};
use std::time::{Duration, Instant};
use colored::Colorize;
const BASE_WORKER_PORT: u16 = 3960;
const DEFAULT_BROKER_PORT: u16 = 9000;
const WORKER_STARTUP_TIMEOUT_SECS: u64 = 20;
struct WorkerHandle {
port: u16,
name: String,
child: Child,
}
impl Drop for WorkerHandle {
fn drop(&mut self) {
let _ = self.child.kill();
}
}
pub fn parse_args(args: &[String]) -> (usize, u16, u16, bool) {
let mut workers = 1usize;
let mut base_port = BASE_WORKER_PORT;
let mut broker_port = DEFAULT_BROKER_PORT;
let mut daemon = false;
let mut i = 0;
while i < args.len() {
match args[i].as_str() {
"--workers" | "-w" | "-n" => {
if let Some(v) = args.get(i + 1) {
workers = v.parse().unwrap_or(1);
i += 2;
continue;
}
}
"--port" | "-p" => {
if let Some(v) = args.get(i + 1) {
base_port = v.parse().unwrap_or(BASE_WORKER_PORT);
i += 2;
continue;
}
}
"--broker-port" | "-b" => {
if let Some(v) = args.get(i + 1) {
broker_port = v.parse().unwrap_or(DEFAULT_BROKER_PORT);
i += 2;
continue;
}
}
"-d" | "--daemon" | "--background" => {
daemon = true;
}
s if s.starts_with("--workers=") => {
workers = s.trim_start_matches("--workers=").parse().unwrap_or(1);
}
_ => {}
}
i += 1;
}
(workers, base_port, broker_port, daemon)
}
pub fn is_port_open(port: u16) -> bool {
TcpStream::connect_timeout(
&format!("127.0.0.1:{}", port).parse().unwrap(),
Duration::from_millis(150),
)
.is_ok()
}
fn find_zakuro_dir() -> Option<PathBuf> {
if let Ok(dir) = std::env::var("ZAKURO_WORKER_DIR") {
let p = PathBuf::from(dir);
if p.join("zakuro/worker/server.py").exists() {
return Some(p);
}
}
let candidates = [
std::env::current_dir()
.ok()
.and_then(|d| d.parent().map(|p| p.join("zak-zakuro"))),
Some(PathBuf::from("/opt/code/ZAK/zak-zakuro")),
std::env::var("HOME")
.ok()
.map(|h| PathBuf::from(h).join("zak-zakuro")),
];
candidates
.into_iter()
.flatten()
.find(|candidate| candidate.join("zakuro/worker/server.py").exists())
}
fn spawn_worker(port: u16, index: usize, zakuro_dir: &PathBuf) -> Result<WorkerHandle, String> {
let name = format!("local-worker-{}", index);
let child = Command::new("uv")
.args([
"run",
"python",
"-m",
"zakuro.worker.server",
"--port",
&port.to_string(),
"--worker-name",
&name,
])
.current_dir(zakuro_dir)
.env("ZAKURO_WORKER_TYPE", "zakuro")
.env("UV_NO_PROGRESS", "1")
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.map_err(|e| format!("Failed to spawn worker on port {}: {}", port, e))?;
Ok(WorkerHandle { port, name, child })
}
fn wait_for_worker(port: u16, timeout_secs: u64) -> bool {
let deadline = Instant::now() + Duration::from_secs(timeout_secs);
let url = format!("http://127.0.0.1:{}/health", port);
while Instant::now() < deadline {
if let Ok(resp) = ureq::get(&url)
.config()
.timeout_global(Some(Duration::from_millis(500)))
.build()
.call()
{
if resp.status().as_u16() == 200 {
return true;
}
}
std::thread::sleep(Duration::from_millis(300));
}
false
}
pub fn start_up(workers: usize, base_port: u16, broker_port: u16, daemon: bool) {
println!();
if daemon {
println!(
" {} {} worker(s) + broker on :{} {}\n",
"zc up".bold().cyan(),
workers,
broker_port,
"(background)".dimmed()
);
} else {
println!(
" {} {} worker(s) + broker on :{}\n",
"zc up".bold().cyan(),
workers,
broker_port
);
}
let zakuro_dir = match find_zakuro_dir() {
Some(d) => d,
None => {
eprintln!(
" {} Cannot find zak-zakuro directory.\n Set ZAKURO_WORKER_DIR or place it at ../zak-zakuro",
"Error:".red()
);
return;
}
};
println!(" {} {}", "Worker src:".bold(), zakuro_dir.display());
let broker_already_up = is_port_open(broker_port);
if broker_already_up {
println!(
" {} Broker already running on :{} (reusing)",
"Broker:".bold(),
broker_port
);
} else {
let scan_end = base_port + (workers as u16).max(1) + 19;
let scan_range = format!("{}-{}", base_port, scan_end);
broker::apply_user_broker_defaults();
if daemon {
let exe = std::env::current_exe().unwrap_or_else(|_| PathBuf::from("zc"));
let broker_child = Command::new(&exe)
.args(["broker", &broker_port.to_string()])
.env("ZAKURO_SCAN_RANGE", &scan_range)
.env("ZAKURO_SCAN_INTERVAL", "3")
.env(
"ZAKURO_P2P",
if broker::p2p_default() {
"true"
} else {
"false"
},
)
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn();
match broker_child {
Ok(_child) => { }
Err(e) => {
eprintln!(
" {} Failed to spawn broker subprocess: {}",
"Error:".red(),
e
);
return;
}
}
} else {
let discovery = broker::DiscoveryConfig {
worker_port: base_port,
extra_ports: vec![],
scan_port_range: Some((base_port, scan_end)),
enable_scan: true,
enable_dns: false,
interval_secs: 5,
..Default::default()
};
let config = broker::BrokerConfig {
host: "0.0.0.0".to_string(),
port: broker_port,
daemon: true,
verbose: false,
tui_mode: false,
enable_discovery: true,
discovery,
enable_p2p: broker::p2p_default(),
..Default::default()
};
std::thread::Builder::new()
.name("broker".into())
.spawn(move || {
if let Err(e) = broker::start_server(config) {
eprintln!(" [BROKER] Error: {}", e);
}
})
.expect("Failed to spawn broker thread");
}
print!(" {} Starting...", "Broker:".bold());
let start = Instant::now();
loop {
if is_port_open(broker_port) {
println!(
"\r {} http://127.0.0.1:{}{}",
"Broker:".bold(),
broker_port,
" ".repeat(20)
);
break;
}
if start.elapsed().as_secs() > 5 {
println!();
eprintln!(" {} Broker didn't start in time", "Error:".red());
return;
}
std::thread::sleep(Duration::from_millis(100));
}
}
println!();
let mut handles: Vec<WorkerHandle> = Vec::with_capacity(workers);
let mut ready_from_existing: usize = 0;
for i in 0..workers {
let port = base_port + i as u16;
if is_port_open(port) {
print!(
" {} Port {} already in use — checking ...",
format!("Worker {}:", i).bold(),
port
);
if wait_for_worker(port, 3) {
println!(" {}", "ready".green());
ready_from_existing += 1;
} else {
println!(" {}", "no response (not a worker?)".yellow());
}
continue;
}
match spawn_worker(port, i, &zakuro_dir) {
Ok(handle) => {
println!(
" {} Spawned on :{} (PID {})",
format!("Worker {}:", i).bold(),
port,
handle.child.id()
);
handles.push(handle);
}
Err(e) => {
eprintln!(" {} Failed to start worker {}: {}", "Error:".red(), i, e);
}
}
}
println!();
let mut ready = ready_from_existing;
for handle in &handles {
print!(
" Waiting for {} on :{} ...",
handle.name.cyan(),
handle.port
);
if wait_for_worker(handle.port, WORKER_STARTUP_TIMEOUT_SECS) {
println!(" {}", "ready".green());
ready += 1;
} else {
println!(" {}", "timeout".red());
}
}
if ready == 0 && (workers > 0) {
eprintln!(
"\n {} No workers became healthy. Aborting.",
"Error:".red()
);
return;
}
std::thread::sleep(Duration::from_secs(2));
println!();
println!(
" {}",
"─────────────────────────────────────────────".dimmed()
);
println!(" {} {}/{} ready", "Workers:".bold(), ready, workers);
println!(
" {} http://127.0.0.1:{} (route requests here)",
"Broker:".bold(),
broker_port
);
println!();
println!(
" Benchmark: {}",
format!("zc bench http://127.0.0.1:{}", broker_port).cyan()
);
if daemon {
println!(" Stop with: {}", "zc down".cyan());
println!(" Monitor: {}", "zc workers".cyan());
println!(
" {}",
"─────────────────────────────────────────────".dimmed()
);
println!();
for handle in handles {
std::mem::forget(handle);
}
std::process::exit(0);
} else {
println!(
" Press {} to stop all workers + broker.",
"Ctrl+C".yellow().bold()
);
println!(
" {}",
"─────────────────────────────────────────────".dimmed()
);
println!();
let running = Arc::new(AtomicBool::new(true));
let r = running.clone();
#[cfg(unix)]
unsafe {
libc_sigint(r);
}
#[cfg(not(unix))]
{
let _ = r;
}
while running.load(Ordering::SeqCst) {
for handle in &mut handles {
if let Ok(Some(status)) = handle.child.try_wait() {
eprintln!(
"\n {} Worker {} (:{}) exited unexpectedly: {}",
"Warning:".yellow(),
handle.name,
handle.port,
status
);
}
}
std::thread::sleep(Duration::from_secs(2));
}
println!("\n Shutting down {} worker(s)...", handles.len());
drop(handles);
println!(" Done.");
}
}
#[cfg(unix)]
unsafe fn libc_sigint(flag: Arc<AtomicBool>) {
use std::sync::Mutex;
static HANDLER: std::sync::OnceLock<Mutex<Option<Arc<AtomicBool>>>> =
std::sync::OnceLock::new();
HANDLER.get_or_init(|| Mutex::new(None));
if let Some(m) = HANDLER.get() {
*m.lock().unwrap() = Some(flag);
}
extern "C" fn handler(_: libc::c_int) {
use std::sync::OnceLock;
static HANDLER: OnceLock<Mutex<Option<Arc<AtomicBool>>>> = OnceLock::new();
println!("\n Received Ctrl+C");
std::process::exit(0);
}
libc::signal(libc::SIGINT, handler as *const () as libc::sighandler_t);
}
pub fn stop_up(broker_port: u16, worker_port_start: u16, worker_port_end: u16) {
println!();
println!(" {} Stopping cluster...", "zc down".bold().cyan());
println!();
let mut killed = 0usize;
let open_ports: Vec<u16> = {
let (tx, rx) = std::sync::mpsc::channel();
let mut handles = Vec::new();
for port in worker_port_start..=worker_port_end {
let tx = tx.clone();
handles.push(std::thread::spawn(move || {
let addr: std::net::SocketAddr = format!("127.0.0.1:{}", port).parse().unwrap();
if TcpStream::connect_timeout(&addr, Duration::from_millis(50)).is_ok() {
let _ = tx.send(port);
}
}));
}
drop(tx);
let mut open: Vec<u16> = rx.iter().collect();
for h in handles {
let _ = h.join();
}
open.sort_unstable();
open
};
for port in open_ports {
if port == broker_port {
continue;
} let url = format!("http://127.0.0.1:{}/health", port);
let is_worker = ureq::get(&url)
.config()
.timeout_global(Some(Duration::from_millis(300)))
.build()
.call()
.map(|r| r.status().as_u16() == 200)
.unwrap_or(false);
if is_worker {
print!(" Stopping worker on :{} ...", port);
if kill_port(port) {
println!(" {}", "done".green());
killed += 1;
} else {
println!(" {}", "failed (not found)".yellow());
}
}
}
if is_port_open(broker_port) {
print!(" Stopping broker on :{} ...", broker_port);
if kill_port(broker_port) {
println!(" {}", "done".green());
killed += 1;
} else {
println!(" {}", "failed (not found)".yellow());
}
}
println!();
if killed > 0 {
println!(" {} {} process(es) stopped.", "✓".green(), killed);
} else {
println!(
" {} Nothing was running on the expected ports.",
"Info:".bold()
);
}
println!();
}
fn kill_port(port: u16) -> bool {
let status = std::process::Command::new("fuser")
.args(["-k", &format!("{}/tcp", port)])
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.status();
if status.map(|s| s.success()).unwrap_or(false) {
return true;
}
let lsof = std::process::Command::new("lsof")
.args(["-ti", &format!("tcp:{}", port)])
.output();
if let Ok(out) = lsof {
let pids = String::from_utf8_lossy(&out.stdout);
for pid_str in pids.split_whitespace() {
if let Ok(pid) = pid_str.parse::<u32>() {
let _ = std::process::Command::new("kill")
.arg(pid.to_string())
.status();
return true;
}
}
}
false
}
use crate::broker;