use std::net::TcpStream;
use std::path::{Path, PathBuf};
use std::process::{Child, Command};
use std::sync::{
atomic::{AtomicBool, Ordering},
Arc,
};
use std::time::{Duration, Instant};
use colored::Colorize;
pub(crate) 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()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ZakuroDirMissing {
pub invalid_setting: Option<String>,
}
impl ZakuroDirMissing {
pub fn detail(&self) -> Option<String> {
self.invalid_setting
.as_ref()
.map(|p| format!("ZAKURO_WORKER_DIR={p} has no zakuro/worker/server.py"))
}
}
fn is_valid_zakuro_dir(p: &Path) -> bool {
p.join("zakuro/worker/server.py").exists()
}
pub fn resolve_zakuro_dir(
env_worker_dir: Option<String>,
file_worker_dir: Option<String>,
zc_dir: Option<&Path>,
legacy_candidates: &[PathBuf],
) -> Result<PathBuf, ZakuroDirMissing> {
let mut invalid_setting = None;
for v in [env_worker_dir, file_worker_dir].into_iter().flatten() {
if v.is_empty() {
continue;
}
let p = PathBuf::from(&v);
if is_valid_zakuro_dir(&p) {
return Ok(p);
}
invalid_setting.get_or_insert(v);
}
if let Some(zc) = zc_dir {
let p = zc.join("zak-zakuro");
if is_valid_zakuro_dir(&p) {
return Ok(p);
}
}
legacy_candidates
.iter()
.find(|c| is_valid_zakuro_dir(c))
.cloned()
.ok_or(ZakuroDirMissing { invalid_setting })
}
pub fn legacy_zakuro_dir_candidates_from_env() -> Vec<PathBuf> {
let mut out = Vec::new();
if let Some(parent) = std::env::current_dir()
.ok()
.and_then(|d| d.parent().map(PathBuf::from))
{
out.push(parent.join("zak-zakuro"));
}
out.push(PathBuf::from("/opt/code/ZAK/zak-zakuro"));
if let Some(h) = std::env::var_os("HOME") {
out.push(PathBuf::from(h).join("zak-zakuro"));
}
out
}
pub(crate) fn resolve_zakuro_dir_for(
env_worker_dir: Option<String>,
zc_dir: Option<&Path>,
legacy_candidates: &[PathBuf],
) -> Result<PathBuf, ZakuroDirMissing> {
let file_worker_dir = zc_dir
.and_then(|d| std::fs::read_to_string(d.join("env")).ok())
.and_then(|text| {
crate::credentials::parse(&text)
.get("ZAKURO_WORKER_DIR")
.cloned()
});
resolve_zakuro_dir(env_worker_dir, file_worker_dir, zc_dir, legacy_candidates)
}
pub(crate) fn find_zakuro_dir() -> Option<PathBuf> {
let zc_dir = crate::credentials::dir();
resolve_zakuro_dir_for(
std::env::var("ZAKURO_WORKER_DIR").ok(),
zc_dir.as_deref(),
&legacy_zakuro_dir_candidates_from_env(),
)
.ok()
}
pub(crate) fn worker_name(node_fingerprint: &str, index: usize) -> String {
let fp8: String = node_fingerprint.chars().take(8).collect();
format!("{fp8}-w{index}")
}
pub(crate) fn worker_argv(port: u16, name: &str) -> Vec<String> {
[
"uv",
"run",
"--extra",
"worker",
"python",
"-m",
"zakuro.worker.server",
"--port",
]
.iter()
.map(|s| s.to_string())
.chain([
port.to_string(),
"--worker-name".to_string(),
name.to_string(),
])
.collect()
}
fn spawn_worker(port: u16, name: String, zakuro_dir: &PathBuf) -> Result<WorkerHandle, String> {
let argv = worker_argv(port, &name);
let child = Command::new(&argv[0])
.args(&argv[1..])
.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 Clone it into ~/.zakuro/zak-zakuro, or set ZAKURO_WORKER_DIR",
"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 node_fp = broker::node_identity::NodeKey::load_or_create().fingerprint();
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, worker_name(&node_fp, 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;
#[cfg(test)]
mod tests {
use super::{resolve_zakuro_dir, worker_argv, worker_name, ZakuroDirMissing};
use crate::agent::files::tests::tmp;
use crate::broker::node_identity::NodeKey;
use std::path::PathBuf;
fn valid_dir(tag: &str) -> PathBuf {
let d = tmp(&format!("up-{tag}"));
std::fs::create_dir_all(d.join("zakuro/worker")).unwrap();
std::fs::write(d.join("zakuro/worker/server.py"), b"").unwrap();
d
}
#[test]
fn env_var_wins_when_valid() {
let d = valid_dir("env");
assert_eq!(
resolve_zakuro_dir(Some(d.display().to_string()), None, None, &[]),
Ok(d)
);
}
#[test]
fn an_invalid_env_setting_falls_through_to_the_file_setting() {
let d = valid_dir("file");
assert_eq!(
resolve_zakuro_dir(
Some("/nonexistent/bogus".to_string()),
Some(d.display().to_string()),
None,
&[]
),
Ok(d)
);
}
#[test]
fn both_settings_invalid_falls_through_to_the_zc_dir() {
let zc = valid_dir("zc-dir-parent");
let zak = zc.join("zak-zakuro");
std::fs::create_dir_all(zak.join("zakuro/worker")).unwrap();
std::fs::write(zak.join("zakuro/worker/server.py"), b"").unwrap();
assert_eq!(
resolve_zakuro_dir(
Some("/nonexistent/a".to_string()),
Some("/nonexistent/b".to_string()),
Some(&zc),
&[]
),
Ok(zak)
);
}
#[test]
fn legacy_candidates_are_tried_in_the_given_order() {
let zc = valid_dir("zc-dir-no-zak");
let first = tmp("legacy-first"); std::fs::create_dir_all(&first).unwrap();
let second_parent = valid_dir("legacy-second");
let second = second_parent.join("zak-zakuro");
std::fs::create_dir_all(second.join("zakuro/worker")).unwrap();
std::fs::write(second.join("zakuro/worker/server.py"), b"").unwrap();
let third_parent = valid_dir("legacy-third");
let third = third_parent.join("zak-zakuro");
std::fs::create_dir_all(third.join("zakuro/worker")).unwrap();
std::fs::write(third.join("zakuro/worker/server.py"), b"").unwrap();
assert_eq!(
resolve_zakuro_dir(
None,
None,
Some(&zc),
&[first.clone(), second.clone(), third.clone()]
),
Ok(second),
"the first candidate is invalid, so the second (not the third) wins"
);
}
#[test]
fn nothing_found_reports_the_first_invalid_setting() {
assert_eq!(
resolve_zakuro_dir(
Some("/nonexistent/env-val".to_string()),
Some("/nonexistent/file-val".to_string()),
None,
&[]
),
Err(ZakuroDirMissing {
invalid_setting: Some("/nonexistent/env-val".to_string())
})
);
}
#[test]
fn nothing_found_and_no_setting_given_has_no_detail() {
let missing = resolve_zakuro_dir(None, None, None, &[]).unwrap_err();
assert_eq!(missing.invalid_setting, None);
assert_eq!(missing.detail(), None);
}
#[test]
fn an_invalid_setting_produces_the_expected_detail_text() {
let missing =
resolve_zakuro_dir(Some("/bad/path".to_string()), None, None, &[]).unwrap_err();
assert_eq!(
missing.detail().as_deref(),
Some("ZAKURO_WORKER_DIR=/bad/path has no zakuro/worker/server.py")
);
}
#[test]
fn legacy_zakuro_dir_candidates_from_env_orders_cwd_dev_path_then_home() {
let candidates = super::legacy_zakuro_dir_candidates_from_env();
assert!(candidates.contains(&PathBuf::from("/opt/code/ZAK/zak-zakuro")));
let dev_path_index = candidates
.iter()
.position(|p| p == &PathBuf::from("/opt/code/ZAK/zak-zakuro"))
.unwrap();
if let Some(home) = std::env::var_os("HOME") {
let home_path = PathBuf::from(home).join("zak-zakuro");
if let Some(home_index) = candidates.iter().position(|p| p == &home_path) {
assert!(
dev_path_index < home_index,
"the dev path comes before $HOME/zak-zakuro"
);
}
}
}
#[test]
fn worker_name_is_fp8_dash_index() {
assert_eq!(worker_name("a1b2c3d4e5f60718", 0), "a1b2c3d4-w0");
assert_eq!(worker_name("a1b2c3d4e5f60718", 12), "a1b2c3d4-w12");
}
#[test]
fn two_macs_never_share_a_worker_name() {
let (a, b) = (NodeKey::generate(), NodeKey::generate());
assert_ne!(
worker_name(&a.fingerprint(), 0),
worker_name(&b.fingerprint(), 0)
);
}
#[test]
fn worker_argv_is_the_python_entrypoint() {
assert_eq!(
worker_argv(3961, "a1b2c3d4-w1"),
vec![
"uv",
"run",
"--extra",
"worker",
"python",
"-m",
"zakuro.worker.server",
"--port",
"3961",
"--worker-name",
"a1b2c3d4-w1",
]
);
}
}