#[cfg(test)]
pub(crate) mod integration {
use std::net::TcpListener;
use std::sync::Arc;
use std::thread;
use std::time::Duration;
use crate::broker::{discovery::DiscoveryConfig, server::start_server, BrokerConfig};
static NEXT_PORT: std::sync::atomic::AtomicU16 = std::sync::atomic::AtomicU16::new(20000);
fn port_range(n: u16) -> u16 {
use std::sync::atomic::Ordering;
loop {
let base = NEXT_PORT.fetch_add(n, Ordering::Relaxed);
assert!(
base.checked_add(n).is_some_and(|end| end < 60000),
"port_range exhausted the test port range"
);
if (base..base + n).all(|p| TcpListener::bind(("127.0.0.1", p)).is_ok()) {
return base;
}
}
}
fn free_port() -> u16 {
port_range(1)
}
pub(crate) fn bind_loopback_http() -> (tiny_http::Server, u16) {
let mut last = None;
for _ in 0..10 {
let port = free_port();
match tiny_http::Server::http(("127.0.0.1", port)) {
Ok(server) => return (server, port),
Err(e) => last = Some(e),
}
}
panic!("tiny_http loopback bind failed after 10 retries: {last:?}");
}
fn retry<T, E>(
attempts: u32,
is_transient: impl Fn(&E) -> bool,
mut f: impl FnMut() -> Result<T, E>,
) -> Result<T, E> {
let mut attempt = 1u32;
loop {
match f() {
Ok(v) => return Ok(v),
Err(e) if is_transient(&e) && attempt < attempts => {
thread::sleep(Duration::from_millis(25 * attempt as u64));
attempt += 1;
}
Err(e) => return Err(e),
}
}
}
fn is_transient_ureq(e: &ureq::Error) -> bool {
matches!(
e,
ureq::Error::Timeout(_) | ureq::Error::Io(_) | ureq::Error::ConnectionFailed
)
}
#[test]
fn retry_returns_ok_after_transient_failures() {
use std::cell::Cell;
let calls = Cell::new(0u32);
let result: Result<&str, &str> = retry(
3,
|_e| true,
|| {
let n = calls.get() + 1;
calls.set(n);
if n < 3 {
Err("transient")
} else {
Ok("ok")
}
},
);
assert_eq!(result, Ok("ok"));
assert_eq!(calls.get(), 3, "should retry up to the first success");
}
#[test]
fn retry_does_not_retry_non_transient() {
use std::cell::Cell;
let calls = Cell::new(0u32);
let result: Result<&str, &str> = retry(
3,
|_e| false,
|| {
calls.set(calls.get() + 1);
Err("fatal")
},
);
assert_eq!(result, Err("fatal"));
assert_eq!(calls.get(), 1, "a non-transient error must not be retried");
}
#[test]
fn retry_gives_up_after_attempts() {
use std::cell::Cell;
let calls = Cell::new(0u32);
let result: Result<&str, &str> = retry(
3,
|_e| true,
|| {
calls.set(calls.get() + 1);
Err("always")
},
);
assert_eq!(result, Err("always"));
assert_eq!(calls.get(), 3, "should attempt exactly `attempts` times");
}
#[test]
fn free_port_hands_out_unique_ports() {
use std::collections::HashSet;
let mut handles = vec![];
for _ in 0..16 {
handles.push(thread::spawn(|| {
(0..20).map(|_| free_port()).collect::<Vec<u16>>()
}));
}
let mut all = HashSet::new();
for h in handles {
for p in h.join().unwrap() {
assert!(all.insert(p), "free_port handed out a duplicate: {p}");
}
}
assert_eq!(all.len(), 16 * 20, "every allocated port must be unique");
}
fn test_config(port: u16) -> BrokerConfig {
BrokerConfig {
host: "127.0.0.1".to_string(),
port,
verbose: false,
daemon: false,
tui_mode: false,
enable_discovery: true,
health_check_interval: 1,
worker_timeout: 5,
min_credits: 0.0001,
enable_p2p: false,
peer_key: None,
owner_user_id: None,
node_name: None,
worker_key: None,
api_url: None,
api_key: None,
wireguard_ip_override: None,
quic_port: None,
runtime_worker_threads: None,
discovery: DiscoveryConfig {
subnet: "10.13.13".to_string(),
worker_port: 3960,
extra_ports: vec![],
scan_port_range: None,
interval_secs: 60,
enable_scan: false,
enable_dns: false,
peers: vec![],
local_workers: vec![],
},
}
}
fn start_broker(port: u16) -> String {
let cfg = test_config(port);
thread::spawn(move || {
let _ = start_server(cfg);
});
let url = format!("http://127.0.0.1:{}", port);
for _ in 0..300 {
if ureq::get(&format!("{}/health", url))
.config()
.timeout_global(Some(Duration::from_millis(150)))
.build()
.call()
.is_ok()
{
return url;
}
thread::sleep(Duration::from_millis(100));
}
panic!("Broker on port {} did not start within 30s", port);
}
fn get(url: &str) -> ureq::http::Response<ureq::Body> {
retry(3, is_transient_ureq, || {
ureq::get(url)
.config()
.timeout_global(Some(Duration::from_secs(5)))
.build()
.call()
})
.unwrap_or_else(|e| panic!("GET {url} failed after retries: {e:?}"))
}
fn get_json(url: &str) -> serde_json::Value {
get(url).into_body().read_json().unwrap()
}
fn post_json(url: &str, body: serde_json::Value) -> ureq::http::Response<ureq::Body> {
retry(3, is_transient_ureq, || {
ureq::post(url)
.header("Content-Type", "application/json")
.config()
.timeout_global(Some(Duration::from_secs(5)))
.build()
.send_json(body.clone())
})
.unwrap_or_else(|e| panic!("POST {url} failed after retries: {e:?}"))
}
fn post_json_err(url: &str, body: serde_json::Value) -> ureq::http::Response<ureq::Body> {
retry(3, is_transient_ureq, || {
ureq::post(url)
.header("Content-Type", "application/json")
.config()
.timeout_global(Some(Duration::from_secs(5)))
.http_status_as_error(false)
.build()
.send_json(body.clone())
})
.unwrap_or_else(|e| panic!("POST(err) {url} failed after retries: {e:?}"))
}
fn post_json_with(
url: &str,
body: serde_json::Value,
headers: &[(&str, &str)],
) -> ureq::http::Response<ureq::Body> {
retry(3, is_transient_ureq, || {
let mut req = ureq::post(url)
.header("Content-Type", "application/json")
.config()
.timeout_global(Some(Duration::from_secs(5)))
.http_status_as_error(false)
.build();
for (k, v) in headers {
req = req.header(*k, *v);
}
req.send_json(body.clone())
})
.unwrap_or_else(|e| panic!("POST(with) {url} failed after retries: {e:?}"))
}
fn add_credits(base: &str, user: &str, amount: f64) {
let url = format!("{}/credits/{}/add", base, user);
let resp = retry(3, is_transient_ureq, || {
ureq::post(&url)
.header("Content-Type", "application/json")
.header("X-Api-Key", "test-master")
.config()
.timeout_global(Some(Duration::from_secs(5)))
.http_status_as_error(false)
.build()
.send_json(serde_json::json!({"amount": amount}))
})
.unwrap_or_else(|e| panic!("add_credits {url} failed after retries: {e:?}"));
assert!(
resp.status().as_u16() == 200,
"add_credits returned {}: {:?}",
resp.status(),
resp.into_body().read_to_string().ok()
);
}
fn get_balance(base: &str, user: &str) -> f64 {
let url = format!("{}/credits/{}", base, user);
let resp = retry(3, is_transient_ureq, || {
ureq::get(&url)
.header("Authorization", &format!("Bearer zk_{}_test", user))
.config()
.timeout_global(Some(Duration::from_secs(5)))
.build()
.call()
})
.unwrap_or_else(|e| panic!("get_balance {url} failed after retries: {e:?}"));
let json: serde_json::Value = resp.into_body().read_json().unwrap();
json["balance"].as_f64().unwrap()
}
fn register_worker(
base: &str,
name: &str,
cpus: f64,
memory_gib: u64,
price_per_hour: f64,
) -> String {
let body = serde_json::json!({
"name": name,
"uri": format!("http://127.0.0.1:{}/{}", free_port(), name),
"worker_type": "zakuro",
"resources": {
"cpus_total": cpus,
"cpus_available": cpus,
"memory_total": memory_gib * 1024 * 1024 * 1024,
"memory_available": memory_gib * 1024 * 1024 * 1024,
"gpus_total": 0,
"gpus_available": 0
},
"pricing": {
"price_per_hour": price_per_hour,
"min_charge": 0.001
}
});
let resp = post_json(&format!("{}/workers", base), body);
let json: serde_json::Value = resp.into_body().read_json().unwrap();
json["id"].as_str().unwrap().to_string()
}
#[test]
fn test_health_check_no_auth_required() {
let url = start_broker(free_port());
let body = get_json(&format!("{}/health", url));
assert_eq!(body["status"], "healthy");
}
#[test]
fn test_worker_list_empty_on_fresh_broker() {
let url = start_broker(free_port());
let body = get_json(&format!("{}/workers", url));
assert_eq!(body["total"], 0);
assert!(body["workers"].as_array().unwrap().is_empty());
}
#[test]
fn test_register_worker_appears_in_list() {
let url = start_broker(free_port());
let id = register_worker(&url, "compute-1", 4.0, 8, 0.001);
assert!(!id.is_empty());
let list = get_json(&format!("{}/workers", url));
assert_eq!(list["total"], 1);
let workers = list["workers"].as_array().unwrap();
assert_eq!(workers[0]["name"], "compute-1");
assert_eq!(workers[0]["id"], id);
}
#[test]
fn test_multiple_workers_registered() {
let url = start_broker(free_port());
register_worker(&url, "w1", 4.0, 8, 0.001);
register_worker(&url, "w2", 8.0, 16, 0.002);
let list = get_json(&format!("{}/workers", url));
assert_eq!(list["total"], 2);
}
#[test]
fn test_worker_heartbeat_accepted() {
let url = start_broker(free_port());
let id = register_worker(&url, "hb-worker", 4.0, 8, 0.001);
let resp = post_json(
&format!("{}/workers/heartbeat", url),
serde_json::json!({"worker_id": id, "active_requests": 3}),
);
assert_eq!(resp.status().as_u16(), 200);
let workers = get_json(&format!("{}/workers", url));
let w = &workers["workers"].as_array().unwrap()[0];
assert_eq!(w["active_requests"], 3);
}
#[test]
fn test_unregister_worker_removes_from_list() {
let url = start_broker(free_port());
let id = register_worker(&url, "temp-worker", 2.0, 4, 0.001);
assert_eq!(get_json(&format!("{}/workers", url))["total"], 1);
ureq::delete(&format!("{}/workers/{}", url, id))
.config()
.timeout_global(Some(Duration::from_secs(5)))
.build()
.call()
.ok();
assert_eq!(get_json(&format!("{}/workers", url))["total"], 0);
}
#[test]
fn test_add_credits_and_read_balance() {
let url = start_broker(free_port());
add_credits(&url, "alice", 100.0);
let balance = get_balance(&url, "alice");
assert!((balance - 100.0).abs() < 0.001);
}
#[test]
fn test_multiple_topups_accumulate() {
let url = start_broker(free_port());
add_credits(&url, "bob", 50.0);
add_credits(&url, "bob", 30.0);
add_credits(&url, "bob", 20.0);
let balance = get_balance(&url, "bob");
assert!((balance - 100.0).abs() < 0.001);
}
#[test]
fn test_user_balances_are_independent() {
let url = start_broker(free_port());
add_credits(&url, "user-a", 100.0);
add_credits(&url, "user-b", 200.0);
assert!((get_balance(&url, "user-a") - 100.0).abs() < 0.001);
assert!((get_balance(&url, "user-b") - 200.0).abs() < 0.001);
}
#[test]
fn test_price_estimate_with_workers() {
let url = start_broker(free_port());
register_worker(&url, "price-w", 4.0, 8, 0.001);
let resp = post_json_err(
&format!("{}/price", url),
serde_json::json!({
"cpus": 1.0,
"memory_bytes": 1073741824,
"estimated_duration_secs": 10.0
}),
);
assert_eq!(resp.status().as_u16(), 200);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert!(body["min_cost"].as_f64().unwrap() > 0.0);
assert!(body["matching_workers"].as_u64().unwrap() >= 1);
}
#[test]
fn test_price_estimate_no_workers_returns_error() {
let url = start_broker(free_port());
let resp = post_json_err(&format!("{}/price", url), serde_json::json!({"cpus": 1.0}));
assert_eq!(resp.status().as_u16(), 503);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["code"], "NO_CAPACITY");
}
#[test]
fn test_worker_price_affects_estimate() {
let url = start_broker(free_port());
register_worker(&url, "cheap-w", 4.0, 8, 0.001);
register_worker(&url, "pricey-w", 4.0, 8, 1.0);
let resp = post_json_err(
&format!("{}/price", url),
serde_json::json!({
"cpus": 1.0,
"memory_bytes": 1073741824,
"estimated_duration_secs": 3600.0
}),
);
assert_eq!(resp.status().as_u16(), 200);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
let min = body["min_cost"].as_f64().unwrap();
let max = body["max_cost"].as_f64().unwrap();
assert!(min < max, "min={} max={}", min, max);
assert_eq!(body["matching_workers"], 2);
}
#[test]
fn test_execute_local_mode_no_worker_returns_no_workers() {
let url = start_broker(free_port());
let resp = post_json_err(
&format!("{}/execute", url),
serde_json::json!({"fn": "test_fn", "args": []}),
);
assert_eq!(resp.status().as_u16(), 503);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["code"], "NO_WORKERS");
}
#[test]
fn test_execute_no_auth_returns_unauthorized_when_local_mode_disabled() {
let url = start_broker(free_port());
let resp = post_json_with(
&format!("{}/execute", url),
serde_json::json!({"fn": "greet"}),
&[("Authorization", "Bearer test-master")],
);
assert!(resp.status().as_u16() >= 400 || resp.status().as_u16() == 503);
}
#[test]
fn test_execute_insufficient_credits_rejected() {
let url = start_broker(free_port());
register_worker(&url, "gpu-w", 8.0, 32, 10.0);
let resp = post_json_with(
&format!("{}/execute", url),
serde_json::json!({
"fn": "train",
"cpus": 8.0,
"estimated_duration_secs": 3600.0
}),
&[("X-Zakuro-User", "broke-user")],
);
assert!(resp.status().as_u16() >= 400 || resp.status().as_u16() == 503);
}
#[test]
fn test_broker_state_local_mode_flag() {
use crate::broker::BrokerState;
let state = BrokerState::new();
assert!(!state.is_local_mode()); state.set_local_mode(true);
assert!(state.is_local_mode());
}
#[test]
fn test_is_local_worker_localhost() {
use crate::broker::BrokerState;
let state = BrokerState::new();
assert!(state.is_local_worker("http://127.0.0.1:3960"));
assert!(state.is_local_worker("http://localhost:3960"));
assert!(!state.is_local_worker("http://10.13.13.5:3960"));
}
#[test]
fn test_is_local_worker_all_local_in_local_mode() {
use crate::broker::BrokerState;
let state = BrokerState::new();
state.set_local_mode(true);
assert!(state.is_local_worker("http://10.13.13.5:3960"));
}
#[test]
fn test_rate_limit_per_second_blocks_excess_requests() {
use crate::broker::credits::CreditManager;
let mgr = CreditManager::new();
mgr.get_or_create("alice", 1000.0);
mgr.set_rate_limits("alice", Some(3), None, None); mgr.set_rate_limit("alice", 1000);
assert!(mgr.check_rate_limit("alice")); assert!(mgr.check_rate_limit("alice")); assert!(mgr.check_rate_limit("alice")); assert!(!mgr.check_rate_limit("alice")); }
#[test]
fn test_stale_worker_marked_unhealthy() {
use crate::broker::worker::WorkerRegistration;
use crate::broker::worker::{HardwareInfo, WorkerPricing, WorkerResources};
use crate::broker::worker::{WorkerRegistry, WorkerStatus};
let registry = WorkerRegistry::new();
let w = registry.register(WorkerRegistration {
name: "stale-w".to_string(),
uri: "http://127.0.0.1:3960".to_string(),
worker_type: "zakuro".to_string(),
resources: WorkerResources::default(),
pricing: WorkerPricing::default(),
tags: vec![],
max_timeout_secs: 0.0,
hardware: HardwareInfo::default(),
wireguard_ip: None,
is_docker: None,
source_node: None,
explicit_local: false,
provider_type: Default::default(),
served_models: vec![],
price_per_mtok: 0.0,
});
assert_eq!(registry.get(&w.id).unwrap().status, WorkerStatus::Healthy);
registry.mark_stale(-1);
assert_eq!(registry.get(&w.id).unwrap().status, WorkerStatus::Unhealthy);
}
#[test]
fn test_scenario_register_worker_and_verify_price() {
let url = start_broker(free_port());
register_worker(&url, "precision-w", 8.0, 16, 7.2);
let resp = post_json_err(
&format!("{}/price", url),
serde_json::json!({
"cpus": 2.0,
"memory_bytes": 2147483648u64, "estimated_duration_secs": 3600.0
}),
);
assert_eq!(resp.status().as_u16(), 200);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
let min = body["min_cost"].as_f64().unwrap();
assert!(
(min - 7.2).abs() < 0.01,
"expected min_cost ≈ 7.2, got {}",
min
);
}
#[test]
fn test_scenario_add_credits_then_execute_path() {
let url = start_broker(free_port());
add_credits(&url, "charlie", 50.0);
assert!((get_balance(&url, "charlie") - 50.0).abs() < 0.001);
let resp = post_json_err(
&format!("{}/execute", url),
serde_json::json!({"fn": "hello"}),
);
assert_eq!(resp.status().as_u16(), 503);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["code"], "NO_WORKERS");
assert!((get_balance(&url, "charlie") - 50.0).abs() < 0.001);
}
#[test]
fn test_budget_credits_derives_correct_timeout_and_reservation() {
use crate::broker::worker::WorkerPricing;
let pricing = WorkerPricing {
price_per_hour: 3600.0,
min_charge: 0.001,
};
let budget = 1.0_f64;
let price_per_sec = pricing.price_per_hour / 3600.0;
let max_secs = if price_per_sec > 0.0 {
budget / price_per_sec
} else {
86_400.0
};
let effective_timeout = max_secs;
assert!(
(effective_timeout - 1.0).abs() < 1e-9,
"expected effective_timeout = 1.0s, got {}",
effective_timeout
);
let reservation = pricing.estimate_cost(effective_timeout);
assert!(
(reservation - 1.0).abs() < 1e-9,
"expected reservation = 1.0 credit, got {}",
reservation
);
}
#[test]
fn test_execute_negative_budget_credits_rejected() {
let url = start_broker(free_port());
let resp = post_json_with(
&format!("{}/execute", url),
serde_json::json!({"fn": "test"}),
&[("X-Zakuro-Requirements", r#"{"budget_credits": -1.0}"#)],
);
assert_eq!(resp.status().as_u16(), 400);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["code"], "BAD_REQUEST");
}
fn test_config_p2p(port: u16) -> BrokerConfig {
BrokerConfig {
peer_key: Some("test-peer-key".to_string()),
enable_p2p: true,
owner_user_id: Some("test-owner".to_string()),
..test_config(port)
}
}
fn start_broker_p2p(port: u16) -> String {
let cfg = test_config_p2p(port);
thread::spawn(move || {
let _ = start_server(cfg);
});
let url = format!("http://127.0.0.1:{}", port);
for _ in 0..300 {
if ureq::get(&format!("{}/health", url))
.config()
.timeout_global(Some(Duration::from_millis(150)))
.build()
.call()
.is_ok()
{
return url;
}
thread::sleep(Duration::from_millis(100));
}
panic!("P2P Broker on port {} did not start within 30s", port);
}
fn start_broker_remote(port: u16) -> String {
let cfg = BrokerConfig {
enable_discovery: false,
..test_config(port)
};
thread::spawn(move || {
let _ = start_server(cfg);
});
let url = format!("http://127.0.0.1:{}", port);
for _ in 0..300 {
if ureq::get(&format!("{}/health", url))
.config()
.timeout_global(Some(Duration::from_millis(150)))
.build()
.call()
.is_ok()
{
return url;
}
thread::sleep(Duration::from_millis(100));
}
panic!("Remote broker on port {} did not start within 30s", port);
}
fn peer_post_json(
url: &str,
body: serde_json::Value,
peer_key: &str,
) -> ureq::http::Response<ureq::Body> {
retry(3, is_transient_ureq, || {
ureq::post(url)
.header("Content-Type", "application/json")
.header("X-Peer-Key", peer_key)
.config()
.timeout_global(Some(Duration::from_secs(5)))
.http_status_as_error(false)
.build()
.send_json(body.clone())
})
.unwrap_or_else(|e| panic!("peer POST {url} failed after retries: {e:?}"))
}
fn peer_get_json(url: &str, peer_key: &str) -> ureq::http::Response<ureq::Body> {
retry(3, is_transient_ureq, || {
ureq::get(url)
.header("X-Peer-Key", peer_key)
.config()
.timeout_global(Some(Duration::from_secs(5)))
.http_status_as_error(false)
.build()
.call()
})
.unwrap_or_else(|e| panic!("peer GET {url} failed after retries: {e:?}"))
}
fn start_mock_worker() -> (String, u16) {
let (raw, port) = bind_loopback_http();
let server = std::sync::Arc::new(raw);
for _ in 0..4 {
let server = server.clone();
thread::spawn(move || {
for request in server.incoming_requests() {
let response = tiny_http::Response::from_data(b"mock-result".to_vec())
.with_status_code(200)
.with_header(
tiny_http::Header::from_bytes(
"Content-Type",
"application/octet-stream",
)
.unwrap(),
)
.with_header(tiny_http::Header::from_bytes("Connection", "close").unwrap());
let _ = request.respond(response);
}
});
}
(format!("http://127.0.0.1:{}", port), port)
}
fn register_worker_at(base: &str, name: &str, uri: &str, price: f64) -> String {
let body = serde_json::json!({
"name": name,
"uri": uri,
"worker_type": "zakuro",
"resources": {
"cpus_total": 4.0,
"cpus_available": 4.0,
"memory_total": 8589934592u64,
"memory_available": 8589934592u64,
"gpus_total": 0,
"gpus_available": 0
},
"pricing": {
"price_per_hour": price,
"min_charge": 0.001
}
});
let resp = post_json(&format!("{}/workers", base), body);
let json: serde_json::Value = resp.into_body().read_json().unwrap();
json["id"].as_str().unwrap().to_string()
}
#[test]
fn test_api_key_format_zk_user_hex_resolves_user_id() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
assert_eq!(
ledger
.resolve_user_from_api_key("zk_9000000001_abc123")
.unwrap(),
"9000000001"
);
assert_eq!(
ledger
.resolve_user_from_api_key("zk_alice_deadbeef")
.unwrap(),
"alice"
);
}
#[test]
fn test_api_key_invalid_format_rejected() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
assert!(ledger.resolve_user_from_api_key("").is_err());
assert!(ledger.resolve_user_from_api_key("zk_").is_err());
assert!(ledger
.resolve_user_from_api_key("zk__nouserportion")
.is_err());
}
#[test]
fn test_ledger_local_reserve_commit_refunds_difference() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
ledger.local_add_credits("alice", 100.0);
assert!((ledger.load_balance_if_needed("alice") - 100.0).abs() < 1e-9);
let (rid, balance_before) = ledger.local_reserve("alice", 40.0, "req-1").unwrap();
assert!((balance_before - 100.0).abs() < 1e-9);
assert!((ledger.load_balance_if_needed("alice") - 60.0).abs() < 1e-9);
let balance_after = ledger.local_commit(&rid, 25.0).unwrap();
assert!(
(balance_after - 75.0).abs() < 1e-9,
"expected 75.0 (reserved 40, charged 25, refund 15), got {}",
balance_after
);
}
#[test]
fn test_ledger_local_reserve_insufficient_rejects() {
use crate::broker::ledger::{Ledger, LedgerError};
let ledger = Ledger::new(None, None);
ledger.local_add_credits("bob", 10.0);
match ledger.local_reserve("bob", 20.0, "req-x") {
Err(LedgerError::InsufficientCredits {
required,
available,
}) => {
assert!((required - 20.0).abs() < 1e-9);
assert!((available - 10.0).abs() < 1e-9);
}
other => panic!("expected InsufficientCredits, got {:?}", other),
}
assert!((ledger.load_balance_if_needed("bob") - 10.0).abs() < 1e-9);
}
#[test]
fn test_ledger_unknown_balance_is_not_insufficient_credits() {
use crate::broker::ledger::{Ledger, LedgerError};
let ledger = Ledger::new(
Some("http://127.0.0.1:1/unreachable".to_string()),
Some("zk_unknown_key".to_string()),
);
assert!(
ledger.try_get_balance("nobody").is_none(),
"an unreachable authority with no cache must report unknown, not Some(0.0)"
);
match ledger.local_reserve("nobody", 1.0, "req-unknown") {
Err(LedgerError::BalanceUnavailable { user_id }) => {
assert_eq!(user_id, "nobody");
}
other => panic!("expected BalanceUnavailable, got {:?}", other),
}
let known_zero = Ledger::new(None, None);
known_zero.local_add_credits("alice", 0.0);
match known_zero.local_reserve("alice", 5.0, "req-zero") {
Err(LedgerError::InsufficientCredits { available, .. }) => {
assert!(available.abs() < 1e-9, "expected 0.0, got {}", available);
}
other => panic!(
"expected InsufficientCredits for a known-zero balance, got {:?}",
other
),
}
}
#[test]
fn test_ledger_local_cancel_restores_full_balance() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
ledger.local_add_credits("carol", 50.0);
let (rid, _) = ledger.local_reserve("carol", 30.0, "req-c").unwrap();
assert!((ledger.load_balance_if_needed("carol") - 20.0).abs() < 1e-9);
ledger.local_cancel(&rid).unwrap();
assert!((ledger.load_balance_if_needed("carol") - 50.0).abs() < 1e-9);
}
#[test]
fn test_ledger_local_add_credits_accumulates() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
assert!((ledger.local_add_credits("dave", 50.0) - 50.0).abs() < 1e-9);
assert!((ledger.local_add_credits("dave", 30.0) - 80.0).abs() < 1e-9);
assert!((ledger.local_add_credits("dave", 20.0) - 100.0).abs() < 1e-9);
}
#[test]
fn test_ledger_reserve_rejects_non_finite_and_negative() {
use crate::broker::ledger::{Ledger, LedgerError};
let ledger = Ledger::new(None, None);
ledger.local_credits.insert("u".to_string(), 100.0);
for bad in [f64::NAN, f64::INFINITY, f64::NEG_INFINITY, -1.0] {
match ledger.reserve("u", bad, "r") {
Err(LedgerError::InvalidAmount(_)) => {}
other => panic!("expected InvalidAmount for {}, got {:?}", bad, other),
}
}
assert!((ledger.local_credits.get("u").map(|v| *v).unwrap_or(0.0) - 100.0).abs() < 1e-9);
}
#[test]
fn test_ledger_local_reserve_rejects_non_finite_and_negative() {
use crate::broker::ledger::{Ledger, LedgerError};
let ledger = Ledger::new(None, None);
ledger.local_add_credits("u", 100.0);
for bad in [f64::NAN, f64::INFINITY, -5.0] {
match ledger.local_reserve("u", bad, "r") {
Err(LedgerError::InvalidAmount(_)) => {}
other => panic!("expected InvalidAmount for {}, got {:?}", bad, other),
}
}
assert!((ledger.load_balance_if_needed("u") - 100.0).abs() < 1e-9);
}
#[test]
fn test_ledger_commit_caps_actual_to_reserved_range() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
ledger.local_credits.insert("u".to_string(), 100.0);
let rid = ledger.reserve("u", 40.0, "r").unwrap();
assert!((ledger.local_credits.get("u").map(|v| *v).unwrap_or(0.0) - 60.0).abs() < 1e-9);
let bal = ledger.commit(&rid, 999.0).unwrap();
assert!(
(bal - 60.0).abs() < 1e-9,
"charge capped at reserved 40, balance stays 60, got {}",
bal
);
let rid2 = ledger.reserve("u", 30.0, "r2").unwrap();
let bal2 = ledger.commit(&rid2, -10.0).unwrap();
assert!(
(bal2 - 60.0).abs() < 1e-9,
"negative charge clamped to 0, full refund, got {}",
bal2
);
}
#[test]
fn test_ledger_local_commit_caps_actual_to_reserved_range() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
ledger.local_add_credits("u", 100.0);
let (rid, _) = ledger.local_reserve("u", 40.0, "r").unwrap();
let bal = ledger.local_commit(&rid, 999.0).unwrap();
assert!((bal - 60.0).abs() < 1e-9, "expected 60, got {}", bal);
let (rid2, _) = ledger.local_reserve("u", 20.0, "r2").unwrap();
let bal2 = ledger.local_commit(&rid2, f64::NAN).unwrap();
assert!(
(bal2 - 40.0).abs() < 1e-9,
"NaN actual charges full reserved, got {}",
bal2
);
}
#[test]
fn test_ledger_reserve_loads_absent_local_credits_key() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
ledger.authoritative_balances.insert("u".to_string(), 50.0);
let rid = ledger
.reserve("u", 20.0, "r")
.expect("reserve should load from authoritative");
assert!((ledger.local_credits.get("u").map(|v| *v).unwrap_or(0.0) - 30.0).abs() < 1e-9);
let bal = ledger.commit(&rid, 20.0).unwrap();
assert!((bal - 30.0).abs() < 1e-9, "expected 30, got {}", bal);
}
#[test]
fn body_over_cap_returns_413() {
let port = free_port();
let base = start_broker(port);
let big = "A".repeat(6 * 1024 * 1024);
let body = serde_json::json!({
"name": big,
"uri": "http://127.0.0.1:3960",
"worker_type": "compute",
"pricing": { "price_per_hour": 1.0, "min_charge": 0.0 },
"resources": {
"cpus_available": 1.0, "cpus_total": 1.0,
"memory_available": 0, "memory_total": 0,
"gpus_available": 0, "gpus_total": 0,
"storage_gb": 0.0
}
});
let result = ureq::post(&format!("{}/workers", base))
.header("Content-Type", "application/json")
.config()
.timeout_global(Some(Duration::from_secs(5)))
.http_status_as_error(false)
.build()
.send_json(body);
match result {
Ok(resp) => assert_eq!(
resp.status().as_u16(),
413,
"oversized body must be rejected with 413, got {}",
resp.status()
),
Err(ureq::Error::Io(_)) | Err(ureq::Error::Timeout(_)) => {}
Err(other) => panic!("unexpected error sending oversized body: {other}"),
}
}
#[test]
fn test_p2p_authority_deterministic_with_two_brokers() {
use crate::broker::peer::{Authority, PeerManager};
let pm = PeerManager::new(
Some("10.0.0.1"),
&["10.0.0.2:3960".to_string()],
9000,
"secret".to_string(),
true,
);
let a1 = pm.determine_authority("user-1");
let a2 = pm.determine_authority("user-1");
match (&a1, &a2) {
(Authority::Local, Authority::Local) => {}
(Authority::Standalone, Authority::Standalone) => {}
(Authority::Peer(a), Authority::Peer(b)) => assert_eq!(a, b),
_ => panic!("authority not deterministic: {:?} vs {:?}", a1, a2),
}
let mut local_count = 0;
let mut fallback_count = 0;
for i in 0..100 {
match pm.determine_authority(&format!("user-{}", i)) {
Authority::Local => local_count += 1,
Authority::Standalone => fallback_count += 1,
_ => {}
}
}
assert!(
local_count >= 20,
"expected >=20 local, got {}",
local_count
);
assert!(
fallback_count >= 20,
"expected >=20 fallback, got {}",
fallback_count
);
}
#[test]
fn test_peer_health_accessible_without_key() {
let url = start_broker(free_port());
let resp = ureq::get(&format!("{}/peer/health", url))
.config()
.timeout_global(Some(Duration::from_secs(5)))
.build()
.call()
.unwrap();
assert_eq!(resp.status().as_u16(), 200);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["status"], "healthy");
}
#[test]
fn test_peer_credit_surface_fail_closed_without_key() {
let url = start_broker(free_port());
for (path, body) in [
(
"/peer/reserve",
serde_json::json!({"user_id": "u", "amount": 1.0, "request_id": "r"}),
),
(
"/peer/commit",
serde_json::json!({"reservation_id": "x", "actual_cost": 1.0}),
),
("/peer/cancel", serde_json::json!({"reservation_id": "x"})),
(
"/peer/earn",
serde_json::json!({"amount": 1.0, "duration_ms": 0.0, "worker_id": "w", "requesting_user": "u", "request_id": "e"}),
),
(
"/peer/tasks/offer",
serde_json::json!({"task_id": "t", "payload_b64": "", "max_price_per_hour": 1.0, "timeout_secs": 1.0}),
),
] {
let resp = ureq::post(&format!("{}{}", url, path))
.header("Content-Type", "application/json")
.config()
.timeout_global(Some(Duration::from_secs(5)))
.http_status_as_error(false)
.build()
.send_json(body)
.unwrap();
assert_eq!(
resp.status().as_u16(),
401,
"POST {} must be 401 when no peer key configured",
path
);
}
for path in ["/peer/balance?user_id=u", "/peer/workers"] {
let resp = ureq::get(&format!("{}{}", url, path))
.config()
.timeout_global(Some(Duration::from_secs(5)))
.http_status_as_error(false)
.build()
.call()
.unwrap();
assert_eq!(
resp.status().as_u16(),
401,
"GET {} must be 401 when no peer key configured",
path
);
}
}
#[test]
fn test_peer_credit_surface_rejects_any_key_when_none_configured() {
let url = start_broker(free_port()); let resp = peer_post_json(
&format!("{}/peer/reserve", url),
serde_json::json!({"user_id": "u", "amount": 1.0, "request_id": "r"}),
"attacker-guessed-key",
);
assert_eq!(resp.status().as_u16(), 401);
}
#[test]
fn test_peer_health_rejects_wrong_key() {
let url = start_broker_p2p(free_port());
let resp = peer_get_json(&format!("{}/peer/health", url), "wrong-key");
assert_eq!(resp.status().as_u16(), 401);
let resp = peer_get_json(&format!("{}/peer/health", url), "test-peer-key");
assert_eq!(resp.status().as_u16(), 200);
}
#[test]
fn test_peers_requires_key_when_configured() {
let url = start_broker_p2p(free_port()); let resp = ureq::get(&format!("{}/peers", url))
.config()
.timeout_global(Some(Duration::from_secs(5)))
.http_status_as_error(false)
.build()
.call()
.unwrap();
assert_eq!(
resp.status().as_u16(),
401,
"/peers must 401 without the peer key"
);
let ok = peer_get_json(&format!("{}/peers", url), "test-peer-key");
assert_eq!(ok.status().as_u16(), 200);
}
#[test]
fn test_two_phase_offer_reserves_without_executing() {
let url = start_broker_p2p(free_port());
register_worker_at(&url, "local-w", "http://127.0.0.1:3960", 1.0);
let offer = serde_json::json!({
"task_id": "tp-1", "payload_b64": "", "max_price_per_hour": 100.0,
"estimated_duration_secs": 1.0, "timeout_secs": 10.0,
"requester_user_id": "u", "source_broker": "n", "two_phase": true
});
let resp = peer_post_json(&format!("{}/peer/tasks/offer", url), offer, "test-peer-key");
assert_eq!(resp.status().as_u16(), 200);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert!(
body.get("estimated_cost").is_some(),
"expected TaskAccept, got {:?}",
body
);
assert!(
body.get("payload_b64").is_none(),
"must not have executed the task"
);
let c = peer_post_json(
&format!("{}/peer/tasks/cancel", url),
serde_json::json!({"task_id": "tp-1"}),
"test-peer-key",
);
assert_eq!(c.status().as_u16(), 200);
}
#[test]
fn test_two_phase_commit_unknown_rejects() {
let url = start_broker_p2p(free_port());
let resp = peer_post_json(
&format!("{}/peer/tasks/commit", url),
serde_json::json!({"task_id": "nope"}),
"test-peer-key",
);
assert_eq!(resp.status().as_u16(), 404);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["task_id"], "nope");
}
#[test]
fn test_two_phase_endpoints_require_peer_key() {
let url = start_broker_p2p(free_port());
let resp = ureq::post(&format!("{}/peer/tasks/commit", url))
.header("Content-Type", "application/json")
.config()
.timeout_global(Some(Duration::from_secs(5)))
.http_status_as_error(false)
.build()
.send_json(serde_json::json!({"task_id": "x"}))
.unwrap();
assert_eq!(resp.status().as_u16(), 401);
}
#[test]
fn test_peer_reserve_commit_cycle_via_http() {
let url = start_broker_p2p(free_port());
let pk = "test-peer-key";
let resp = peer_post_json(
&format!("{}/peer/earn", url),
serde_json::json!({
"amount": 100.0,
"duration_ms": 1000.0,
"worker_id": "test-w",
"requesting_user": "someone",
"request_id": "earn-1"
}),
pk,
);
assert_eq!(resp.status().as_u16(), 200);
let resp = peer_get_json(&format!("{}/peer/balance?user_id=test-owner", url), pk);
assert_eq!(resp.status().as_u16(), 200);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert!((body["balance"].as_f64().unwrap() - 100.0).abs() < 1e-4);
let resp = peer_post_json(
&format!("{}/peer/reserve", url),
serde_json::json!({"user_id": "test-owner", "amount": 40.0, "request_id": "res-1"}),
pk,
);
assert_eq!(resp.status().as_u16(), 200);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
let reservation_id = body["reservation_id"].as_str().unwrap().to_string();
assert!((body["balance_before"].as_f64().unwrap() - 100.0).abs() < 1e-4);
let resp = peer_get_json(&format!("{}/peer/balance?user_id=test-owner", url), pk);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert!((body["balance"].as_f64().unwrap() - 60.0).abs() < 1e-4);
let resp = peer_post_json(
&format!("{}/peer/commit", url),
serde_json::json!({"reservation_id": reservation_id, "actual_cost": 25.0}),
pk,
);
assert_eq!(resp.status().as_u16(), 200);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert!((body["balance_after"].as_f64().unwrap() - 75.0).abs() < 1e-4);
}
#[test]
fn test_peer_cancel_refunds_via_http() {
let url = start_broker_p2p(free_port());
let pk = "test-peer-key";
peer_post_json(
&format!("{}/peer/earn", url),
serde_json::json!({"amount": 50.0, "duration_ms": 0.0, "worker_id": "w", "requesting_user": "u", "request_id": "e1"}),
pk,
);
let resp = peer_post_json(
&format!("{}/peer/reserve", url),
serde_json::json!({"user_id": "test-owner", "amount": 30.0, "request_id": "r2"}),
pk,
);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
let rid = body["reservation_id"].as_str().unwrap().to_string();
let resp = peer_post_json(
&format!("{}/peer/cancel", url),
serde_json::json!({"reservation_id": rid}),
pk,
);
assert_eq!(resp.status().as_u16(), 200);
let resp = peer_get_json(&format!("{}/peer/balance?user_id=test-owner", url), pk);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert!((body["balance"].as_f64().unwrap() - 50.0).abs() < 1e-4);
}
#[test]
fn test_peer_reserve_insufficient_returns_402() {
let url = start_broker_p2p(free_port());
let pk = "test-peer-key";
let resp = peer_post_json(
&format!("{}/peer/reserve", url),
serde_json::json!({"user_id": "test-owner", "amount": 100.0, "request_id": "r3"}),
pk,
);
assert_eq!(resp.status().as_u16(), 402);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["code"], "INSUFFICIENT_CREDITS");
}
#[test]
fn test_transaction_buffer_push_and_drain() {
use crate::broker::flush::{BufferedTransaction, TransactionBuffer};
let buf = TransactionBuffer::new();
assert_eq!(buf.pending_count(), 0);
assert_eq!(buf.total_flushed(), 0);
buf.push_transaction(BufferedTransaction {
request_id: "tx-1".to_string(),
user_id: "alice".to_string(),
tx_type: "commit".to_string(),
amount: 0.05,
balance_after: 99.95,
worker_id: "w-1".to_string(),
duration_ms: 150.0,
source_node: Some("node-1".to_string()),
worker_name: Some("worker-1".to_string()),
worker_uri: Some("http://10.13.13.5:3960".to_string()),
price_per_hour: 3.6,
});
buf.push_transaction(BufferedTransaction {
request_id: "tx-2".to_string(),
user_id: "alice".to_string(),
tx_type: "commit".to_string(),
amount: 0.10,
balance_after: 99.85,
worker_id: "w-1".to_string(),
duration_ms: 200.0,
source_node: None,
worker_name: None,
worker_uri: None,
price_per_hour: 0.0,
});
assert_eq!(buf.pending_count(), 2);
}
#[test]
fn test_transaction_buffer_balance_snapshots() {
use crate::broker::flush::TransactionBuffer;
let buf = TransactionBuffer::new();
buf.snapshot_balance("alice", 100.0);
buf.snapshot_balance("alice", 95.0);
buf.snapshot_balance("bob", 200.0);
assert_eq!(buf.pending_count(), 0);
}
#[test]
fn test_credits_endpoint_returns_balance_status() {
let url = start_broker(free_port());
add_credits(&url, "status-user", 50.0);
let resp = ureq::get(&format!("{}/credits/status-user", url))
.header("Authorization", "Bearer zk_status-user_test")
.config()
.timeout_global(Some(Duration::from_secs(5)))
.build()
.call()
.unwrap();
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["balance_status"], "authoritative");
assert!(body["last_prefetched"].is_null());
}
#[test]
fn test_local_mode_mock_worker_execute_no_credits_charged() {
let (worker_url, _) = start_mock_worker();
let broker_url = start_broker(free_port());
add_credits(&broker_url, "local-user", 100.0);
let balance_before = get_balance(&broker_url, "local-user");
assert!((balance_before - 100.0).abs() < 0.001);
register_worker_at(&broker_url, "mock-w", &worker_url, 3.6);
let resp = post_json_with(
&format!("{}/execute", broker_url),
serde_json::json!({"fn": "test_fn", "args": []}),
&[("X-Zakuro-User", "local-user")],
);
assert_eq!(resp.status().as_u16(), 200);
let balance_after = get_balance(&broker_url, "local-user");
assert!(
(balance_after - 100.0).abs() < 0.001,
"balance should be unchanged in local mode, got {}",
balance_after
);
}
#[test]
fn test_async_execute_route_round_trips_worker_response() {
let (worker_url, _) = start_mock_worker();
let broker_url = start_broker(free_port());
register_worker_at(&broker_url, "async-mock-w", &worker_url, 3.6);
let resp = ureq::post(&format!("{}/execute", broker_url))
.header("Content-Type", "application/json")
.header("X-Zakuro-User", "async-test-user")
.config()
.timeout_global(Some(std::time::Duration::from_secs(10)))
.http_status_as_error(false)
.build()
.send_json(serde_json::json!({"fn": "test_fn", "args": []}))
.unwrap();
assert_eq!(
resp.status().as_u16(),
200,
"async /execute route must return 200"
);
let body = resp.into_body().read_to_string().unwrap();
assert_eq!(
body, "mock-result",
"async /execute route must return the worker's body verbatim"
);
}
#[test]
fn test_remote_mode_requires_bearer_auth() {
let url = start_broker_remote(free_port());
let resp = post_json_err(
&format!("{}/execute", url),
serde_json::json!({"fn": "test"}),
);
assert_eq!(resp.status().as_u16(), 401);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["code"], "UNAUTHORIZED");
}
#[test]
fn test_remote_mode_enforces_credits_for_non_local_worker() {
let url_free = start_broker_remote(free_port());
register_worker_at(&url_free, "remote-w", "http://10.99.99.1:3960", 3.6);
let status_free = ureq::post(&format!("{}/execute", url_free))
.header("Authorization", "Bearer zk_nocreds_abc123def")
.config()
.timeout_global(Some(Duration::from_secs(5)))
.http_status_as_error(false)
.build()
.send_json(serde_json::json!({"fn": "test"}))
.map(|r| r.status().as_u16())
.unwrap_or(0); assert_ne!(
status_free, 402,
"no API credentials → billing disabled → must never get 402"
);
let port_billed = free_port();
let cfg = BrokerConfig {
enable_discovery: false,
api_url: Some("http://test-authority".to_string()),
api_key: Some("test-key".to_string()),
..test_config(port_billed)
};
thread::spawn(move || {
let _ = start_server(cfg);
});
let url_billed = format!("http://127.0.0.1:{}", port_billed);
for _ in 0..300 {
if ureq::get(&format!("{}/health", url_billed))
.config()
.timeout_global(Some(Duration::from_millis(150)))
.build()
.call()
.is_ok()
{
break;
}
thread::sleep(Duration::from_millis(100));
}
register_worker_at(&url_billed, "remote-w", "http://10.99.99.1:3960", 3.6);
let resp = post_json_with(
&format!("{}/execute", url_billed),
serde_json::json!({"fn": "test"}),
&[("Authorization", "Bearer zk_nocreds_abc123def")],
);
assert_eq!(
resp.status().as_u16(),
402,
"billing enabled → must get 402 for zero credits"
);
}
#[test]
fn test_p2p_dispatch_rejects_insufficient_credits_before_dispatch() {
let port = free_port();
let cfg = BrokerConfig {
enable_discovery: false,
enable_p2p: true,
peer_key: Some("test-peer-key".to_string()),
owner_user_id: Some("test-owner".to_string()),
api_url: Some("http://test-authority".to_string()),
api_key: Some("test-key".to_string()),
..test_config(port)
};
thread::spawn(move || {
let _ = start_server(cfg);
});
let url = format!("http://127.0.0.1:{}", port);
for _ in 0..300 {
if ureq::get(&format!("{}/health", url))
.config()
.timeout_global(Some(Duration::from_millis(150)))
.build()
.call()
.is_ok()
{
break;
}
thread::sleep(Duration::from_millis(100));
}
let resp = post_json_with(
&format!("{}/execute", url),
serde_json::json!({"fn": "test"}),
&[
("Authorization", "Bearer zk_nocreds_abc123def"),
(
"X-Zakuro-Requirements",
r#"{"remote_only": true, "estimated_duration_secs": 60.0, "budget_credits": 1.0}"#,
),
],
);
assert_eq!(
resp.status().as_u16(),
402,
"P2P dispatch with zero credits must return 402, not dispatch the task"
);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["code"], "INSUFFICIENT_CREDITS");
}
#[test]
fn test_is_local_worker_comprehensive() {
use crate::broker::BrokerState;
let state = BrokerState::new();
assert!(state.is_local_worker("http://127.0.0.1:3960"));
assert!(state.is_local_worker("http://127.0.0.1:8089"));
assert!(state.is_local_worker("http://localhost:3960"));
assert!(state.is_local_worker("http://localhost:9000/execute"));
assert!(!state.is_local_worker("http://10.13.13.5:3960"));
assert!(!state.is_local_worker("http://100.64.0.1:3960"));
assert!(!state.is_local_worker("http://192.168.1.100:3960"));
state.set_local_mode(true);
assert!(state.is_local_worker("http://10.13.13.5:3960"));
assert!(state.is_local_worker("http://100.64.0.1:3960"));
state.set_local_mode(false);
assert!(!state.is_local_worker("http://10.13.13.5:3960"));
}
#[test]
fn test_peer_workers_exports_own_workers_never_peer_synced_ones() {
let url = start_broker_p2p(free_port());
let pk = "test-peer-key";
register_worker_at(&url, "local-w", "http://127.0.0.1:3960", 1.0);
register_worker_at(&url, "advertised-w", "http://10.13.13.5:3960", 2.0);
post_json(
&format!("{}/workers", url),
serde_json::json!({
"name": "peer-w",
"uri": "http://10.13.13.9:3960",
"worker_type": "zakuro",
"source_node": "zc://node-deadbeef"
}),
);
let body = get_json(&format!("{}/workers", url));
assert_eq!(body["total"], 3);
let resp = peer_get_json(&format!("{}/peer/workers", url), pk);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["total"], 2);
let names: Vec<&str> = body["workers"]
.as_array()
.unwrap()
.iter()
.map(|w| w["name"].as_str().unwrap())
.collect();
assert!(names.contains(&"local-w") && names.contains(&"advertised-w"));
assert!(!names.contains(&"peer-w"));
}
#[test]
fn test_wal_proof_of_execution_full_lifecycle() {
use crate::broker::wal::{Wal, WalEntry, WalStatus};
let path = format!("/tmp/zc_test_wal_proof_{}.jsonl", free_port());
let _ = std::fs::remove_file(&path);
let wal = Wal::open(&path).unwrap();
wal.append(&WalEntry {
request_id: "proof-1".to_string(),
user_id: "alice".to_string(),
reservation_id: "res-proof-1".to_string(),
estimated_cost: 0.05,
actual_cost: None,
worker_id: "worker-a".to_string(),
duration_ms: None,
timestamp: chrono::Utc::now(),
status: WalStatus::Reserved,
})
.unwrap();
wal.flush_buffer().unwrap();
let uncommitted = wal.read_uncommitted().unwrap();
assert!(uncommitted
.iter()
.any(|e| e.request_id == "proof-1" && e.status == WalStatus::Reserved));
wal.update_status("proof-1", WalStatus::Executed, Some(0.03), Some(250.0))
.unwrap();
wal.flush_buffer().unwrap();
let uncommitted = wal.read_uncommitted().unwrap();
assert!(uncommitted
.iter()
.any(|e| e.request_id == "proof-1" && e.status == WalStatus::Executed));
wal.update_status("proof-1", WalStatus::Committed, Some(0.03), Some(250.0))
.unwrap();
wal.flush_buffer().unwrap();
let uncommitted = wal.read_uncommitted().unwrap();
assert!(
!uncommitted.iter().any(|e| e.request_id == "proof-1"),
"committed entry must not appear in uncommitted"
);
let _ = std::fs::remove_file(&path);
}
#[test]
fn test_wal_recovery_cancels_stale_reservations() {
use crate::broker::ledger::Ledger;
use crate::broker::recovery;
use crate::broker::wal::{Wal, WalEntry, WalStatus};
let path = format!("/tmp/zc_test_wal_cancel_{}.jsonl", free_port());
let _ = std::fs::remove_file(&path);
let ledger = Ledger::new(None, None);
let wal = Wal::open(&path).unwrap();
wal.append(&WalEntry {
request_id: "stale-req".to_string(),
user_id: "alice".to_string(),
reservation_id: "res-stale".to_string(),
estimated_cost: 30.0,
actual_cost: None,
worker_id: "w1".to_string(),
duration_ms: None,
timestamp: chrono::Utc::now(),
status: WalStatus::Reserved,
})
.unwrap();
wal.flush_buffer().unwrap();
let tx_buffer = crate::broker::flush::TransactionBuffer::new();
recovery::replay_wal(&wal, &ledger, &tx_buffer, None);
let uncommitted = wal.read_uncommitted().unwrap();
assert!(
uncommitted.is_empty(),
"all entries should be resolved after replay"
);
let _ = std::fs::remove_file(&path);
}
#[test]
fn test_wal_recovery_commits_executed_entries() {
use crate::broker::ledger::Ledger;
use crate::broker::recovery;
use crate::broker::wal::{Wal, WalEntry, WalStatus};
let path = format!("/tmp/zc_test_wal_commit_{}.jsonl", free_port());
let _ = std::fs::remove_file(&path);
let ledger = Ledger::new(None, None);
let wal = Wal::open(&path).unwrap();
wal.append(&WalEntry {
request_id: "exec-req".to_string(),
user_id: "bob".to_string(),
reservation_id: "res-exec".to_string(),
estimated_cost: 0.03,
actual_cost: Some(0.03),
worker_id: "w2".to_string(),
duration_ms: Some(150.0),
timestamp: chrono::Utc::now(),
status: WalStatus::Executed,
})
.unwrap();
wal.flush_buffer().unwrap();
let tx_buffer = crate::broker::flush::TransactionBuffer::new();
recovery::replay_wal(&wal, &ledger, &tx_buffer, None);
let uncommitted = wal.read_uncommitted().unwrap();
assert!(
uncommitted.is_empty(),
"executed entry should be committed after replay"
);
let _ = std::fs::remove_file(&path);
}
#[test]
fn test_wal_recovery_does_not_double_credit_committed_earn() {
use crate::broker::flush::TransactionBuffer;
use crate::broker::ledger::Ledger;
use crate::broker::recovery;
use crate::broker::wal::{Wal, WalEntry, WalStatus};
let path = format!("/tmp/zc_test_wal_earn_committed_{}.jsonl", free_port());
let _ = std::fs::remove_file(&path);
let ledger = Ledger::new(None, None);
let wal = Wal::open(&path).unwrap();
wal.append(&WalEntry {
request_id: "earn-req-committed".to_string(),
user_id: "owner1".to_string(),
reservation_id: String::new(),
estimated_cost: 0.05,
actual_cost: Some(0.05),
worker_id: "worker-1".to_string(),
duration_ms: Some(1000.0),
timestamp: chrono::Utc::now(),
status: WalStatus::Earned,
})
.unwrap();
ledger.local_add_credits("owner1", 0.05);
wal.update_status("earn-req-committed", WalStatus::Committed, Some(0.05), None)
.unwrap();
wal.flush_buffer().unwrap();
let tx_buffer = TransactionBuffer::new();
recovery::replay_wal(&wal, &ledger, &tx_buffer, None);
let balance = ledger
.authoritative_balances
.get("owner1")
.map(|v| *v)
.unwrap_or(0.0);
assert!(
(balance - 0.05).abs() < 1e-9,
"earn must be applied exactly once, got balance {}",
balance
);
assert_eq!(
tx_buffer.pending_count(),
0,
"a committed earn must not be re-queued for delivery"
);
let _ = std::fs::remove_file(&path);
}
#[test]
fn test_wal_recovery_replays_unflushed_earn() {
use crate::broker::flush::TransactionBuffer;
use crate::broker::ledger::Ledger;
use crate::broker::recovery;
use crate::broker::wal::{Wal, WalEntry, WalStatus};
let path = format!("/tmp/zc_test_wal_earn_uncommitted_{}.jsonl", free_port());
let _ = std::fs::remove_file(&path);
let ledger = Ledger::new(None, None);
let wal = Wal::open(&path).unwrap();
wal.append(&WalEntry {
request_id: "earn-req-lost".to_string(),
user_id: "owner2".to_string(),
reservation_id: String::new(),
estimated_cost: 0.07,
actual_cost: Some(0.07),
worker_id: "worker-2".to_string(),
duration_ms: Some(2000.0),
timestamp: chrono::Utc::now(),
status: WalStatus::Earned,
})
.unwrap();
wal.flush_buffer().unwrap();
let tx_buffer = TransactionBuffer::new();
recovery::replay_wal(&wal, &ledger, &tx_buffer, None);
assert!(
ledger.authoritative_balances.get("owner2").is_none(),
"replay must not mutate in-memory balance (dashboard is authoritative)"
);
assert_eq!(
tx_buffer.pending_count(),
1,
"unflushed earn must be re-queued for dashboard delivery after replay"
);
let uncommitted = wal.read_uncommitted().unwrap();
assert!(
uncommitted.iter().any(|e| e.request_id == "earn-req-lost"),
"replayed-but-not-yet-flushed earn must remain uncommitted"
);
wal.update_status(
"earn-req-lost",
WalStatus::Committed,
Some(0.07),
Some(2000.0),
)
.unwrap();
wal.flush_buffer().unwrap();
let uncommitted = wal.read_uncommitted().unwrap();
assert!(
!uncommitted.iter().any(|e| e.request_id == "earn-req-lost"),
"confirmed earn must no longer be replay-eligible"
);
let _ = std::fs::remove_file(&path);
}
#[test]
fn test_flush_per_item_results_marks_only_durable_committed() {
use crate::broker::flush::{BufferedTransaction, TransactionBuffer};
use crate::broker::wal::{Wal, WalEntry, WalStatus};
let path = format!("/tmp/zc_test_flush_per_item_{}.jsonl", free_port());
let _ = std::fs::remove_file(&path);
let wal = Wal::open(&path).unwrap();
let mk = |rid: &str| WalEntry {
request_id: rid.to_string(),
user_id: "owner3".to_string(),
reservation_id: String::new(),
estimated_cost: 1.0,
actual_cost: Some(1.0),
worker_id: "w".to_string(),
duration_ms: Some(10.0),
timestamp: chrono::Utc::now(),
status: WalStatus::Earned,
};
for rid in ["earn-a", "earn-b", "earn-c"] {
wal.append(&mk(rid)).unwrap();
}
wal.flush_buffer().unwrap();
let buf = TransactionBuffer::new();
for rid in ["earn-a", "earn-b", "earn-c"] {
buf.push_transaction(BufferedTransaction {
request_id: rid.to_string(),
user_id: "owner3".to_string(),
tx_type: "credit".to_string(),
amount: 1.0,
balance_after: 0.0,
worker_id: "w".to_string(),
duration_ms: 10.0,
source_node: None,
worker_name: Some("w".to_string()),
worker_uri: None,
price_per_hour: 0.0,
});
}
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
let handle = std::thread::spawn(move || {
use std::io::{Read, Write};
let (mut stream, _) = listener.accept().unwrap();
let mut buf = [0u8; 8192];
let _ = stream.read(&mut buf);
let body = r#"{"results":[
{"request_id":"earn-a","status":"applied","error":null},
{"request_id":"earn-b","status":"duplicate","error":null},
{"request_id":"earn-c","status":"failed","error":"boom"}],
"applied":1,"duplicates":1,"failed":1,"success":false,
"inserted":2,"errors":[{"job_name":"x","error":"boom"}]}"#;
let resp = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
body.len(),
body
);
let _ = stream.write_all(resp.as_bytes());
let _ = stream.flush();
});
let api_url = format!("http://{}", addr);
buf.flush_to_api(&api_url, "test-key", &wal);
handle.join().unwrap();
wal.flush_buffer().unwrap();
let uncommitted: std::collections::HashSet<String> = wal
.read_uncommitted()
.unwrap()
.into_iter()
.map(|e| e.request_id)
.collect();
assert!(
!uncommitted.contains("earn-a"),
"applied earn must be Committed"
);
assert!(
!uncommitted.contains("earn-b"),
"duplicate earn must be Committed"
);
assert!(
uncommitted.contains("earn-c"),
"failed earn must remain Earned/replay-eligible"
);
assert_eq!(
buf.pending_count(),
1,
"failed item must be re-buffered for retry"
);
let _ = std::fs::remove_file(&path);
}
#[test]
fn test_prefetched_balance_status_transitions() {
use crate::broker::credits::{BalanceStatus, CreditManager};
let mgr = CreditManager::new();
mgr.get_or_create("peer-user", 100.0);
assert_eq!(
mgr.get("peer-user").unwrap().balance_status,
BalanceStatus::Authoritative
);
mgr.set_prefetched_balance("peer-user", 95.0);
let uc = mgr.get("peer-user").unwrap();
assert_eq!(uc.balance_status, BalanceStatus::Prefetched);
assert_eq!(uc.balance, 95.0);
assert!(uc.last_prefetched.is_some());
mgr.set_reconciling("peer-user");
assert_eq!(
mgr.get("peer-user").unwrap().balance_status,
BalanceStatus::Reconciling
);
mgr.set_authoritative("peer-user");
let uc = mgr.get("peer-user").unwrap();
assert_eq!(uc.balance_status, BalanceStatus::Authoritative);
assert!(uc.last_prefetched.is_none());
}
#[test]
fn test_reconciliation_needed_after_staleness() {
use crate::broker::credits::CreditManager;
let mgr = CreditManager::new();
mgr.get_or_create("stale-user", 100.0);
assert!(!mgr.needs_reconciliation("stale-user", 0));
mgr.set_prefetched_balance("stale-user", 90.0);
assert!(mgr.needs_reconciliation("stale-user", 0));
assert!(!mgr.needs_reconciliation("stale-user", 999999));
assert!(!mgr.needs_reconciliation("ghost-user", 0));
}
#[test]
fn test_retry_marks_unreachable_workers_unhealthy() {
let url = start_broker(free_port());
register_worker_at(&url, "unreachable-1", "http://127.0.0.1:19991", 0.001);
register_worker_at(&url, "unreachable-2", "http://127.0.0.1:19992", 0.001);
let workers = get_json(&format!("{}/workers", url));
assert_eq!(workers["total"], 2);
for w in workers["workers"].as_array().unwrap() {
assert_eq!(w["status"], "healthy");
}
let resp = post_json_err(
&format!("{}/execute", url),
serde_json::json!({"fn": "unreachable_test"}),
);
assert_eq!(resp.status().as_u16(), 503);
let workers = get_json(&format!("{}/workers", url));
let unhealthy_count = workers["workers"]
.as_array()
.unwrap()
.iter()
.filter(|w| w["status"] == "unhealthy")
.count();
assert_eq!(
unhealthy_count, 2,
"both unreachable workers should be marked unhealthy after retry"
);
}
#[test]
fn test_retry_exhausts_all_workers_returns_503() {
let url = start_broker(free_port());
register_worker_at(&url, "dead-1", "http://127.0.0.1:19993", 0.001);
register_worker_at(&url, "dead-2", "http://127.0.0.1:19994", 0.001);
register_worker_at(&url, "dead-3", "http://127.0.0.1:19995", 0.001);
let resp = post_json_err(
&format!("{}/execute", url),
serde_json::json!({"fn": "fail_test"}),
);
assert_eq!(resp.status().as_u16(), 503);
let body: serde_json::Value = resp.into_body().read_json().unwrap();
assert_eq!(body["code"], "NO_WORKERS");
let error_msg = body["error"].as_str().unwrap();
assert!(
error_msg.contains("unreachable") || error_msg.contains("retries"),
"error should mention retry exhaustion: {}",
error_msg
);
}
#[test]
fn test_retry_failure_then_new_healthy_worker_succeeds() {
let url = start_broker(free_port());
register_worker_at(&url, "bad-w", "http://127.0.0.1:19996", 0.001);
let resp = post_json_err(
&format!("{}/execute", url),
serde_json::json!({"fn": "fail"}),
);
assert_eq!(resp.status().as_u16(), 503);
let (mock_url, _) = start_mock_worker();
register_worker_at(&url, "good-w", &mock_url, 0.001);
let resp = post_json_err(
&format!("{}/execute", url),
serde_json::json!({"fn": "succeed"}),
);
assert_eq!(
resp.status().as_u16(),
200,
"should succeed with healthy worker"
);
}
#[test]
fn test_multi_window_rate_limits_enforced_together() {
use crate::broker::credits::CreditManager;
let mgr = CreditManager::new();
mgr.get_or_create("limited-user", 1000.0);
mgr.set_rate_limits("limited-user", Some(5), Some(10), None);
mgr.set_rate_limit("limited-user", 1000);
for _ in 0..5 {
assert!(mgr.check_rate_limit("limited-user"));
}
assert!(!mgr.check_rate_limit("limited-user"));
}
#[test]
fn test_ledger_local_fallback_full_cycle() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
ledger.local_credits.insert("user1".to_string(), 100.0);
assert_eq!(ledger.get_balance("user1"), 100.0);
let rid = ledger.reserve("user1", 40.0, "req-1").unwrap();
assert_eq!(
ledger.local_credits.get("user1").map(|v| *v).unwrap_or(0.0),
60.0
);
let _balance = ledger.commit(&rid, 25.0).unwrap();
assert_eq!(
ledger.local_credits.get("user1").map(|v| *v).unwrap_or(0.0),
75.0
);
assert!(ledger.cancel("non-existent").is_ok());
}
#[test]
fn test_api_mode_skips_pg_pool_creation() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(
Some("https://my.zakuro-ai.com".to_string()),
Some("zk_test_abc123".to_string()),
);
assert!(ledger.is_api_mode());
}
#[test]
fn test_api_mode_credit_cycle_uses_local_only() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(
Some("https://my.zakuro-ai.com".to_string()),
Some("zk_test_abc123".to_string()),
);
assert!(ledger.is_api_mode());
ledger.local_credits.insert("api-user".to_string(), 100.0);
let rid = ledger.reserve("api-user", 40.0, "req-1").unwrap();
assert_eq!(
ledger
.local_credits
.get("api-user")
.map(|v| *v)
.unwrap_or(0.0),
60.0
);
let _ = ledger.commit(&rid, 25.0).unwrap();
assert_eq!(
ledger
.local_credits
.get("api-user")
.map(|v| *v)
.unwrap_or(0.0),
75.0
);
assert!(ledger.cancel("non-existent").is_ok());
}
#[test]
fn test_api_mode_wal_recovery_cancel_uses_local_fallback() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(
Some("https://my.zakuro-ai.com".to_string()),
Some("zk_test_abc123".to_string()),
);
assert!(ledger.cancel_from_wal("user-x", 50.0).is_ok());
let local_balance = ledger
.local_credits
.get("user-x")
.map(|v| *v)
.unwrap_or(0.0);
assert_eq!(local_balance, 50.0);
}
#[test]
fn test_api_mode_wal_recovery_commit_uses_local_fallback() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(
Some("https://my.zakuro-ai.com".to_string()),
Some("zk_test_abc123".to_string()),
);
let balance = ledger.commit_from_wal("user-y", 0.10, 0.06).unwrap();
assert!((balance - 0.04).abs() < 1e-9);
}
#[test]
fn test_non_api_mode_attempts_pg_connection() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
assert!(!ledger.is_api_mode());
}
fn fibonacci(n: u64) -> u64 {
if n <= 1 {
return n;
}
let (mut a, mut b) = (0u64, 1u64);
for _ in 2..=n {
let c = a + b;
a = b;
b = c;
}
b
}
fn start_fibonacci_worker(bind_ip: &str) -> (String, u16) {
let addr = format!("{}:0", bind_ip);
let server = tiny_http::Server::http(&addr)
.unwrap_or_else(|e| panic!("bind {}:0 failed: {}", bind_ip, e));
let port = server.server_addr().to_ip().unwrap().port();
thread::spawn(move || {
for mut req in server.incoming_requests() {
let mut buf = Vec::new();
let _ = req.as_reader().read_to_end(&mut buf);
let n = serde_json::from_slice::<serde_json::Value>(&buf)
.ok()
.and_then(|v| v["n"].as_u64())
.unwrap_or(4);
let body = serde_json::json!({"result": fibonacci(n), "n": n});
let _ = req.respond(
tiny_http::Response::from_data(serde_json::to_vec(&body).unwrap())
.with_status_code(200)
.with_header(
tiny_http::Header::from_bytes("Content-Type", "application/json")
.unwrap(),
),
);
}
});
(format!("http://{}:{}", bind_ip, port), port)
}
fn start_mesh_broker(
port: u16,
own_ip: &str,
owner: &str,
peer_key: &str,
peers: Vec<String>,
api_url: Option<String>,
api_key: Option<String>,
) -> String {
start_mesh_broker_with_quic(port, own_ip, owner, peer_key, peers, api_url, api_key, None)
}
#[allow(clippy::too_many_arguments)]
fn start_mesh_broker_with_quic(
port: u16,
own_ip: &str,
owner: &str,
peer_key: &str,
peers: Vec<String>,
api_url: Option<String>,
api_key: Option<String>,
quic_port: Option<u16>,
) -> String {
let cfg = BrokerConfig {
host: "0.0.0.0".to_string(),
port,
health_check_interval: 2,
worker_timeout: 30,
min_credits: 0.0001,
daemon: true,
verbose: false,
tui_mode: false,
enable_discovery: false,
enable_p2p: true,
peer_key: Some(peer_key.to_string()),
owner_user_id: Some(owner.to_string()),
node_name: Some(format!("mesh-{}", own_ip)),
worker_key: None,
api_url,
api_key,
wireguard_ip_override: Some(own_ip.to_string()),
quic_port,
runtime_worker_threads: None,
discovery: DiscoveryConfig {
subnet: "10.13.13".to_string(),
worker_port: 3960,
extra_ports: vec![],
scan_port_range: None,
interval_secs: 60,
enable_scan: false,
enable_dns: false,
peers,
local_workers: vec![],
},
};
thread::spawn(move || {
let _ = start_server(cfg);
});
let url = format!("http://127.0.0.1:{}", port);
for _ in 0..300 {
if ureq::get(&format!("{}/health", url))
.config()
.timeout_global(Some(Duration::from_millis(200)))
.build()
.call()
.is_ok()
{
return url;
}
thread::sleep(Duration::from_millis(100));
}
panic!("Mesh broker on port {} did not start within 30s", port);
}
#[test]
fn test_p2p_mesh_fibonacci_500() {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Instant;
let peer_key = "fib-mesh-key";
let users = ["alice", "bob", "carol", "dave"];
let total_req = 500usize;
let threads = 10usize;
let per_thread = total_req / threads;
let (w1_uri, w1_port) = start_fibonacci_worker("127.0.0.2");
let (w2_uri, w2_port) = start_fibonacci_worker("127.0.0.3");
let b1_port = free_port();
let b2_port = free_port();
let b1_url = start_mesh_broker(
b1_port,
"127.0.0.2",
"node01-owner",
peer_key,
vec![format!("http://127.0.0.3:{}", b2_port)],
None,
None,
);
let b2_url = start_mesh_broker(
b2_port,
"127.0.0.3",
"node02-owner",
peer_key,
vec![format!("http://127.0.0.2:{}", b1_port)],
None,
None,
);
for base in [&b1_url, &b2_url] {
register_worker_at(base, "fib-node01", &w1_uri, 0.001);
register_worker_at(base, "fib-node02", &w2_uri, 0.001);
}
struct Row {
user: String,
broker: &'static str,
latency_ms: f64,
ok: bool,
fib_ok: bool,
}
let rows = Arc::new(std::sync::Mutex::new(Vec::<Row>::with_capacity(total_req)));
let idx = Arc::new(AtomicUsize::new(0));
let clock = Instant::now();
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_global(Some(Duration::from_secs(30)))
.http_status_as_error(false)
.build(),
);
let mut handles = Vec::new();
for _ in 0..threads {
let rows = rows.clone();
let idx = idx.clone();
let agent = agent.clone();
let b1 = b1_url.clone();
let b2 = b2_url.clone();
handles.push(thread::spawn(move || {
for _ in 0..per_thread {
let i = idx.fetch_add(1, Ordering::SeqCst);
let user = users[i % users.len()];
let broker = if i.is_multiple_of(2) {
(&b1, "Node01")
} else {
(&b2, "Node02")
};
let key = format!("zk_{}_0001", user);
let t0 = Instant::now();
let resp = agent
.post(&format!("{}/execute", broker.0))
.header("Authorization", &format!("Bearer {}", key))
.header(
"X-Zakuro-Requirements",
r#"{"strategy":"round_robin","estimated_duration_secs":0.01}"#,
)
.send_json(serde_json::json!({"n": 4}));
let ms = t0.elapsed().as_secs_f64() * 1000.0;
let (ok, fib_ok) = match resp {
Ok(r) if r.status().as_u16() == 200 => {
let body: serde_json::Value =
r.into_body().read_json().unwrap_or_default();
(true, body["result"].as_u64() == Some(3))
}
Ok(r) => {
let code = r.status().as_u16();
let msg = r.into_body().read_to_string().unwrap_or_default();
eprintln!(" [WARN] HTTP {} for {}: {}", code, user, msg);
(false, false)
}
Err(e) => {
eprintln!(" [WARN] transport error: {}", e);
(false, false)
}
};
rows.lock().unwrap().push(Row {
user: user.to_string(),
broker: broker.1,
latency_ms: ms,
ok,
fib_ok,
});
}
}));
}
for h in handles {
h.join().unwrap();
}
let wall = clock.elapsed();
let rows = rows.lock().unwrap();
let n = rows.len();
let ok = rows.iter().filter(|r| r.ok).count();
let fib_ok = rows.iter().filter(|r| r.fib_ok).count();
let rps = n as f64 / wall.as_secs_f64();
let mut lats: Vec<f64> = rows.iter().map(|r| r.latency_ms).collect();
lats.sort_by(|a, b| a.partial_cmp(b).unwrap());
let pct = |p: usize| lats[lats.len() * p / 100];
let node_stats = |name: &str| {
let rs: Vec<&Row> = rows.iter().filter(|r| r.broker == name).collect();
let cnt = rs.len();
let avg_ms = if cnt > 0 {
rs.iter().map(|r| r.latency_ms).sum::<f64>() / cnt as f64
} else {
0.0
};
(cnt, avg_ms)
};
let (n1_cnt, n1_lat) = node_stats("Node01");
let (n2_cnt, n2_lat) = node_stats("Node02");
let u_stats: Vec<(&str, usize)> = users
.iter()
.map(|u| {
let cnt = rows.iter().filter(|r| r.user == *u).count();
(*u, cnt)
})
.collect();
println!("\n{}", "═".repeat(72));
println!(" P2P Mesh Fibonacci(4) — 500 Requests (no billing)");
println!("{}", "═".repeat(72));
println!("\n Cluster Topology");
println!(" {}", "─".repeat(68));
println!(
" Node01 broker=0.0.0.0:{} own_ip=127.0.0.2 worker=127.0.0.2:{}",
b1_port, w1_port
);
println!(
" Node02 broker=0.0.0.0:{} own_ip=127.0.0.3 worker=127.0.0.3:{}",
b2_port, w2_port
);
println!(" Mode: broker-to-broker (no API, billing disabled)");
println!("\n Execution Summary");
println!(" {}", "─".repeat(68));
println!(" Total requests: {}", n);
println!(
" Successful: {} ({:.1}%)",
ok,
ok as f64 / n as f64 * 100.0
);
println!(" fibonacci(4) = 3: {}/{} correct", fib_ok, n);
println!(" Wall-clock time: {:.2}s", wall.as_secs_f64());
println!(" Throughput: {:.1} req/s", rps);
println!("\n Latency (ms)");
println!(" {}", "─".repeat(68));
println!(" Min: {:.2}", lats[0]);
println!(" Avg: {:.2}", lats.iter().sum::<f64>() / n as f64);
println!(" p50: {:.2}", pct(50));
println!(" p90: {:.2}", pct(90));
println!(" p99: {:.2}", pct(99));
println!(" Max: {:.2}", lats[n - 1]);
println!("\n Per-Node Breakdown");
println!(" {}", "─".repeat(68));
println!(" Node01: {} reqs avg_lat={:.2}ms", n1_cnt, n1_lat);
println!(" Node02: {} reqs avg_lat={:.2}ms", n2_cnt, n2_lat);
println!("\n Per-User Breakdown");
println!(" {}", "─".repeat(68));
for (u, cnt) in &u_stats {
println!(" {:8} {} reqs", u, cnt);
}
println!("\n {}", "━".repeat(68));
println!(" API VERIFICATION (GET /stats, /workers)");
println!(" {}", "━".repeat(68));
for (label, base) in [("Node01", &b1_url), ("Node02", &b2_url)] {
println!("\n ┌─ {} ({})", label, base);
let workers_raw = match agent.get(&format!("{}/workers", base)).call() {
Ok(r) => r.into_body().read_to_string().unwrap_or_default(),
Err(e) => format!("ERR: {}", e),
};
let workers_resp: serde_json::Value =
serde_json::from_str(&workers_raw).unwrap_or_default();
let workers_arr = workers_resp["workers"]
.as_array()
.or_else(|| workers_resp.as_array());
println!(
" │ /workers ({} registered)",
workers_arr.map(|a| a.len()).unwrap_or(0)
);
if let Some(ws) = workers_arr {
for w in ws {
println!(
" │ {} {} status={}",
w["name"].as_str().unwrap_or("?"),
w["uri"].as_str().unwrap_or("?"),
w["status"].as_str().unwrap_or("?")
);
}
}
let admin_key = "zk_admin_test-master";
let stats_raw = agent
.get(&format!("{}/stats", base))
.header("Authorization", &format!("Bearer {}", admin_key))
.call()
.map(|r| r.into_body().read_to_string().unwrap_or_default())
.unwrap_or_default();
let stats_resp: serde_json::Value =
serde_json::from_str(&stats_raw).unwrap_or_default();
let m = &stats_resp["metrics"];
println!(" │ /stats");
println!(" │ total_requests: {}", m["total_requests"]);
println!(" │ successful: {}", m["successful_requests"]);
println!(" │ failed: {}", m["failed_requests"]);
println!(
" │ avg_latency_ms: {:.2}",
m["avg_latency_ms"].as_f64().unwrap_or(0.0)
);
println!(" │ active_workers: {}", m["active_workers"]);
println!(" └─");
}
println!("\n{}\n", "═".repeat(72));
assert_eq!(ok, n, "all requests must succeed");
assert_eq!(fib_ok, n, "all fibonacci results must equal 3");
assert!(rps > 5.0, "throughput must exceed 5 req/s, got {:.1}", rps);
assert!(n1_cnt > 0, "Node01 must handle requests");
assert!(n2_cnt > 0, "Node02 must handle requests");
for (label, base) in [("Node01", &b1_url), ("Node02", &b2_url)] {
let raw = agent
.get(&format!("{}/stats", base))
.header("Authorization", "Bearer zk_admin_test-master")
.call()
.map(|r| r.into_body().read_to_string().unwrap_or_default())
.unwrap_or_default();
let stats: serde_json::Value = serde_json::from_str(&raw).unwrap_or_default();
let api_total = stats["metrics"]["total_requests"].as_u64().unwrap_or(0);
let api_ok = stats["metrics"]["successful_requests"]
.as_u64()
.unwrap_or(0);
assert!(api_total > 0, "{} /stats must report requests", label);
assert_eq!(
api_total, api_ok,
"{} all requests must be successful",
label
);
}
}
#[test]
fn test_p2p_mesh_fibonacci_500_api() {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Instant;
let api_url = match std::env::var("ZAKURO_API_URL") {
Ok(v) if !v.is_empty() => v,
_ => {
println!("SKIP: ZAKURO_API_URL not set");
return;
}
};
let api_key = match std::env::var("ZAKURO_API_KEY") {
Ok(v) if !v.is_empty() => v,
_ => {
println!("SKIP: ZAKURO_API_KEY not set");
return;
}
};
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_global(Some(Duration::from_secs(10)))
.http_status_as_error(false)
.build(),
);
let me_raw = match agent
.get(&format!(
"{}/api/auth/me/api-key",
api_url.trim_end_matches('/')
))
.header("Authorization", &format!("Bearer {}", api_key))
.call()
{
Ok(resp) => resp.into_body().read_to_string().unwrap_or_default(),
Err(e) => {
println!("SKIP: /api/auth/me/api-key failed ({})", e);
return;
}
};
let me: serde_json::Value = serde_json::from_str(&me_raw).unwrap_or_default();
let zakuro_uid = match me["zakuro_user_id"].as_str() {
Some(uid) => uid,
None => {
println!("SKIP: no zakuro_user_id in response");
return;
}
};
let balance_before = me["credits_balance"].as_f64().unwrap_or(0.0);
let username = me["username"].as_str().unwrap_or("?");
println!(
" Dashboard user: {} (zakuro_user_id={})",
username, zakuro_uid
);
println!(" Balance before: {:.6} credits", balance_before);
let peer_key = "fib-mesh-api-key";
let master = "test-master";
let total_req = 500usize;
let threads = 10usize;
let per_thread = total_req / threads;
let user_key = format!("zk_{}_0001", zakuro_uid);
let (w1_uri, _w1_port) = start_fibonacci_worker("127.0.0.2");
let (w2_uri, _w2_port) = start_fibonacci_worker("127.0.0.3");
let b1_port = free_port();
let b2_port = free_port();
let b1_url = start_mesh_broker(
b1_port,
"127.0.0.2",
"placeholder",
peer_key,
vec![format!("http://127.0.0.3:{}", b2_port)],
Some(api_url.clone()),
Some(api_key.clone()),
);
let b2_url = start_mesh_broker(
b2_port,
"127.0.0.3",
"placeholder",
peer_key,
vec![format!("http://127.0.0.2:{}", b1_port)],
Some(api_url.clone()),
Some(api_key.clone()),
);
register_worker_at(&b1_url, "fib-node01", &w1_uri, 0.001);
register_worker_at(&b2_url, "fib-node02", &w2_uri, 0.001);
struct Row {
broker: &'static str,
worker: String,
latency_ms: f64,
cost: f64,
ok: bool,
fib_ok: bool,
}
let rows = Arc::new(std::sync::Mutex::new(Vec::<Row>::with_capacity(total_req)));
let idx = Arc::new(AtomicUsize::new(0));
let clock = Instant::now();
let mut handles = Vec::new();
for _ in 0..threads {
let rows = rows.clone();
let idx = idx.clone();
let agent = agent.clone();
let b1 = b1_url.clone();
let b2 = b2_url.clone();
let key = user_key.clone();
handles.push(thread::spawn(move || {
for _ in 0..per_thread {
let i = idx.fetch_add(1, Ordering::SeqCst);
let broker = if i.is_multiple_of(2) {
(&b1, "Node01")
} else {
(&b2, "Node02")
};
let t0 = Instant::now();
let resp = agent
.post(&format!("{}/execute", broker.0))
.header("Authorization", &format!("Bearer {}", key))
.header(
"X-Zakuro-Requirements",
r#"{"strategy":"round_robin","estimated_duration_secs":0.01}"#,
)
.send_json(serde_json::json!({"n": 4}));
let ms = t0.elapsed().as_secs_f64() * 1000.0;
let (ok, cost, fib_ok, worker) = match resp {
Ok(r) if r.status().as_u16() == 200 => {
let c = r
.headers()
.get("X-Zakuro-Cost")
.and_then(|v| v.to_str().ok())
.and_then(|v| v.parse::<f64>().ok())
.unwrap_or(0.0);
let w = r
.headers()
.get("X-Zakuro-Worker")
.and_then(|v| v.to_str().ok())
.unwrap_or("?")
.to_string();
let body: serde_json::Value =
r.into_body().read_json().unwrap_or_default();
(true, c, body["result"].as_u64() == Some(3), w)
}
Ok(r) => {
let code = r.status().as_u16();
let msg = r.into_body().read_to_string().unwrap_or_default();
eprintln!(" [WARN] HTTP {} : {}", code, msg);
(false, 0.0, false, "?".to_string())
}
Err(e) => {
eprintln!(" [WARN] transport error: {}", e);
(false, 0.0, false, "?".to_string())
}
};
rows.lock().unwrap().push(Row {
broker: broker.1,
worker,
latency_ms: ms,
cost,
ok,
fib_ok,
});
}
}));
}
for h in handles {
h.join().unwrap();
}
let wall = clock.elapsed();
let rows = rows.lock().unwrap();
let n = rows.len();
let ok = rows.iter().filter(|r| r.ok).count();
let fib_ok = rows.iter().filter(|r| r.fib_ok).count();
let rps = n as f64 / wall.as_secs_f64();
let charged = rows.iter().filter(|r| r.cost > 0.0).count();
let _free = rows.iter().filter(|r| r.cost == 0.0 && r.ok).count();
let tot_cost: f64 = rows.iter().map(|r| r.cost).sum();
let n1_to_n2 = rows
.iter()
.filter(|r| r.broker == "Node01" && r.worker == "fib-node02")
.count();
let n2_to_n1 = rows
.iter()
.filter(|r| r.broker == "Node02" && r.worker == "fib-node01")
.count();
let n1_local = rows
.iter()
.filter(|r| r.broker == "Node01" && r.worker == "fib-node01")
.count();
let n2_local = rows
.iter()
.filter(|r| r.broker == "Node02" && r.worker == "fib-node02")
.count();
let mut lats: Vec<f64> = rows.iter().map(|r| r.latency_ms).collect();
lats.sort_by(|a, b| a.partial_cmp(b).unwrap());
let pct = |p: usize| lats[lats.len() * p / 100];
let node_stats = |name: &str| {
let rs: Vec<&Row> = rows.iter().filter(|r| r.broker == name).collect();
let cnt = rs.len();
let c: f64 = rs.iter().map(|r| r.cost).sum();
let ch = rs.iter().filter(|r| r.cost > 0.0).count();
let fr = rs.iter().filter(|r| r.cost == 0.0 && r.ok).count();
let avg_ms = if cnt > 0 {
rs.iter().map(|r| r.latency_ms).sum::<f64>() / cnt as f64
} else {
0.0
};
(cnt, ch, fr, c, avg_ms)
};
let (n1_cnt, _n1_ch, _n1_fr, n1_cost, n1_lat) = node_stats("Node01");
let (n2_cnt, _n2_ch, _n2_fr, n2_cost, n2_lat) = node_stats("Node02");
println!("\n{}", "═".repeat(72));
println!(" P2P Publish-Lock Fibonacci(4) — 500 Requests (API mode)");
println!("{}", "═".repeat(72));
println!("\n Cluster Topology (publish-lock)");
println!(" {}", "─".repeat(68));
println!(
" zk0node01 broker=0.0.0.0:{} ip=127.0.0.2 local_worker=fib-node01",
b1_port
);
println!(
" zk0node02 broker=0.0.0.0:{} ip=127.0.0.3 local_worker=fib-node02",
b2_port
);
println!(" Mode: API (verified owner via dashboard handshake)");
println!(" Owner: {} (zakuro_user_id={})", username, zakuro_uid);
println!(" Model: publish-lock — each broker executes only on its own workers");
println!("\n Execution Routing");
println!(" {}", "─".repeat(68));
println!(
" zk0node01 → fib-node01 (local): {} requests",
n1_local
);
println!(
" zk0node01 → fib-node02 (peer): {} requests (published to zk0node02)",
n1_to_n2
);
println!(
" zk0node02 → fib-node02 (local): {} requests",
n2_local
);
println!(
" zk0node02 → fib-node01 (peer): {} requests (published to zk0node01)",
n2_to_n1
);
println!("\n Execution Summary");
println!(" {}", "─".repeat(68));
println!(" Total requests: {}", n);
println!(
" Successful: {} ({:.1}%)",
ok,
ok as f64 / n as f64 * 100.0
);
println!(" fibonacci(4) = 3: {}/{} correct", fib_ok, n);
println!(" Wall-clock time: {:.2}s", wall.as_secs_f64());
println!(" Throughput: {:.1} req/s", rps);
println!("\n Latency (ms)");
println!(" {}", "─".repeat(68));
println!(" Min: {:.2}", lats[0]);
println!(" Avg: {:.2}", lats.iter().sum::<f64>() / n as f64);
println!(" p50: {:.2}", pct(50));
println!(" p90: {:.2}", pct(90));
println!(" p99: {:.2}", pct(99));
println!(" Max: {:.2}", lats[n - 1]);
println!("\n Credits (self-execution = free, verified owner)");
println!(" {}", "─".repeat(68));
println!(" All executions free (same verified owner on all nodes)");
println!(" Total cost: {:.6} credits", tot_cost);
println!("\n Per-Node Breakdown");
println!(" {}", "─".repeat(68));
println!(
" zk0node01: {} reqs cost={:.6} avg_lat={:.2}ms",
n1_cnt, n1_cost, n1_lat
);
println!(
" zk0node02: {} reqs cost={:.6} avg_lat={:.2}ms",
n2_cnt, n2_cost, n2_lat
);
println!("\n {}", "━".repeat(68));
println!(" BROKER API (GET /stats, /workers, /peer/identity)");
println!(" {}", "━".repeat(68));
let admin_key = format!("zk_admin_{}", master);
for (label, base) in [("zk0node01", &b1_url), ("zk0node02", &b2_url)] {
println!("\n ┌─ {} ({})", label, base);
let workers_raw = match agent.get(&format!("{}/workers", base)).call() {
Ok(r) => r.into_body().read_to_string().unwrap_or_default(),
Err(e) => format!("ERR: {}", e),
};
let workers_resp: serde_json::Value =
serde_json::from_str(&workers_raw).unwrap_or_default();
let workers_arr = workers_resp["workers"]
.as_array()
.or_else(|| workers_resp.as_array());
println!(
" │ /workers ({} registered)",
workers_arr.map(|a| a.len()).unwrap_or(0)
);
let stats_raw = agent
.get(&format!("{}/stats", base))
.header("Authorization", &format!("Bearer {}", admin_key))
.call()
.map(|r| r.into_body().read_to_string().unwrap_or_default())
.unwrap_or_default();
let stats_resp: serde_json::Value =
serde_json::from_str(&stats_raw).unwrap_or_default();
let m = &stats_resp["metrics"];
println!(
" │ /stats total={} ok={} spent={:.6}",
m["total_requests"],
m["successful_requests"],
m["total_credits_spent"].as_f64().unwrap_or(0.0)
);
let id_raw = agent
.get(&format!("{}/peer/identity", base))
.header("X-Peer-Key", peer_key)
.call()
.map(|r| r.into_body().read_to_string().unwrap_or_default())
.unwrap_or_default();
let id: serde_json::Value = serde_json::from_str(&id_raw).unwrap_or_default();
println!(
" │ /peer/identity owner={} verified={} workers={}",
id["owner_user_id"].as_str().unwrap_or("?"),
id["verified"].as_bool().unwrap_or(false),
id["workers"].as_array().map(|a| a.len()).unwrap_or(0)
);
println!(" └─");
}
println!("\n {}", "━".repeat(68));
println!(" DASHBOARD VERIFICATION ({})", api_url);
println!(" {}", "━".repeat(68));
println!(" Waiting for flush cycles (5s)...");
thread::sleep(Duration::from_secs(5));
let me_after = agent
.get(&format!(
"{}/api/auth/me/api-key",
api_url.trim_end_matches('/')
))
.header("Authorization", &format!("Bearer {}", api_key))
.call()
.expect("failed to call /api/auth/me/api-key after flush")
.into_body()
.read_to_string()
.unwrap();
let me_after: serde_json::Value = serde_json::from_str(&me_after).unwrap();
let balance_after = me_after["credits_balance"].as_f64().unwrap_or(0.0);
let deducted = balance_before - balance_after;
println!(" Balance before: {:>12.6} credits", balance_before);
println!(" Balance after: {:>12.6} credits", balance_after);
println!(" Deducted: {:>12.6} credits", deducted);
println!(
" Owner verified: {} (handshake confirmed at startup)",
zakuro_uid
);
if (balance_after - balance_before).abs() < 0.0001 {
println!(
" Status: CORRECT — no credits moved (self-execution, verified owner)"
);
} else {
println!(" Status: UNEXPECTED — balance changed for self-execution");
}
println!("\n{}\n", "═".repeat(72));
assert_eq!(ok, n, "all requests must succeed");
assert_eq!(fib_ok, n, "all fibonacci results must equal 3");
assert!(rps > 5.0, "throughput must exceed 5 req/s, got {:.1}", rps);
assert!(n1_cnt > 0, "zk0node01 must handle requests");
assert!(n2_cnt > 0, "zk0node02 must handle requests");
assert_eq!(charged, 0, "self-execution must not charge credits");
assert!(tot_cost == 0.0, "total cost must be 0 for self-execution");
assert!(
(balance_after - balance_before).abs() < 0.0001,
"balance must not change for self-execution (before={:.6}, after={:.6})",
balance_before,
balance_after
);
}
#[test]
fn test_quic_throughput_benchmark() {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Instant;
let api_url = match std::env::var("ZAKURO_API_URL") {
Ok(v) if !v.is_empty() => v,
_ => {
println!("SKIP: ZAKURO_API_URL not set");
return;
}
};
let api_key = match std::env::var("ZAKURO_API_KEY") {
Ok(v) if !v.is_empty() => v,
_ => {
println!("SKIP: ZAKURO_API_KEY not set");
return;
}
};
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_global(Some(Duration::from_secs(30)))
.http_status_as_error(false)
.build(),
);
let me_raw = match agent
.get(&format!(
"{}/api/auth/me/api-key",
api_url.trim_end_matches('/')
))
.header("Authorization", &format!("Bearer {}", api_key))
.call()
{
Ok(resp) => resp.into_body().read_to_string().unwrap_or_default(),
Err(e) => {
println!("SKIP: /api/auth/me/api-key failed ({})", e);
return;
}
};
let me: serde_json::Value = serde_json::from_str(&me_raw).unwrap_or_default();
let zakuro_uid = match me["zakuro_user_id"].as_str() {
Some(uid) => uid,
None => {
println!("SKIP: no zakuro_user_id in response");
return;
}
};
let username = me["username"].as_str().unwrap_or("?");
let peer_key = "quic-bench-key";
let total_req = 500usize;
let threads = 10usize;
let per_thread = total_req / threads;
let user_key = format!("zk_{}_0001", zakuro_uid);
let (_w2_uri, w2_port) = start_fibonacci_worker("127.0.0.3");
let w2_uri = format!("http://127.0.0.3:{}", w2_port);
let b1_port = free_port();
let q1_port = free_port();
let b2_port = free_port();
let q2_port = free_port();
let b1_url = start_mesh_broker_with_quic(
b1_port,
"127.0.0.2",
"placeholder",
peer_key,
vec![format!("http://127.0.0.3:{}", b2_port)],
Some(api_url.clone()),
Some(api_key.clone()),
Some(q1_port),
);
let b2_url = start_mesh_broker_with_quic(
b2_port,
"127.0.0.3",
"placeholder",
peer_key,
vec![format!("http://127.0.0.2:{}", b1_port)],
Some(api_url.clone()),
Some(api_key.clone()),
Some(q2_port),
);
register_worker_at(&b2_url, "fib-node02", &w2_uri, 0.001);
thread::sleep(Duration::from_millis(1500));
#[derive(Debug)]
struct Row {
latency_ms: f64,
ok: bool,
fib_ok: bool,
worker: String,
transport: String,
}
let run_phase = |broker_url: &str, n: usize| -> Vec<Row> {
let rows = Arc::new(std::sync::Mutex::new(Vec::<Row>::with_capacity(n)));
let idx = Arc::new(AtomicUsize::new(0));
let mut handles = Vec::new();
for _ in 0..threads {
let rows = rows.clone();
let idx = idx.clone();
let agent = agent.clone();
let url = broker_url.to_string();
let key = user_key.clone();
handles.push(thread::spawn(move || {
for _ in 0..per_thread {
let _ = idx.fetch_add(1, Ordering::SeqCst);
let t0 = Instant::now();
let resp = agent
.post(&format!("{}/execute", url))
.header("Authorization", &format!("Bearer {}", key))
.header(
"X-Zakuro-Requirements",
r#"{"strategy":"round_robin","estimated_duration_secs":0.01}"#,
)
.send_json(serde_json::json!({"n": 4}));
let ms = t0.elapsed().as_secs_f64() * 1000.0;
let (ok, fib_ok, worker, transport) = match resp {
Ok(r) if r.status().as_u16() == 200 => {
let w = r
.headers()
.get("X-Zakuro-Worker")
.and_then(|v| v.to_str().ok())
.unwrap_or("?")
.to_string();
let t = r
.headers()
.get("X-Zakuro-Transport")
.and_then(|v| v.to_str().ok())
.unwrap_or("?")
.to_string();
let body: serde_json::Value =
r.into_body().read_json().unwrap_or_default();
(true, body["result"].as_u64() == Some(3), w, t)
}
Ok(r) => {
let code = r.status().as_u16();
let msg = r.into_body().read_to_string().unwrap_or_default();
eprintln!(" [WARN] HTTP {} : {}", code, msg);
(false, false, "?".to_string(), "?".to_string())
}
Err(e) => {
eprintln!(" [WARN] transport error: {}", e);
(false, false, "?".to_string(), "?".to_string())
}
};
rows.lock().unwrap().push(Row {
latency_ms: ms,
ok,
fib_ok,
worker,
transport,
});
}
}));
}
for h in handles {
h.join().unwrap();
}
Arc::try_unwrap(rows).unwrap().into_inner().unwrap()
};
let clock1 = Instant::now();
let quic_rows = run_phase(&b1_url, total_req);
let wall1 = clock1.elapsed();
let clock2 = Instant::now();
let local_rows = run_phase(&b2_url, total_req);
let wall2 = clock2.elapsed();
let stats = |rows: &[Row]| {
let n = rows.len();
let ok = rows.iter().filter(|r| r.ok).count();
let fib = rows.iter().filter(|r| r.fib_ok).count();
let mut lats: Vec<f64> = rows.iter().map(|r| r.latency_ms).collect();
lats.sort_by(|a, b| a.partial_cmp(b).unwrap());
let avg = lats.iter().sum::<f64>() / n.max(1) as f64;
let p50 = lats[n * 50 / 100];
let p90 = lats[n * 90 / 100];
let p99 = lats[n * 99 / 100];
(n, ok, fib, avg, p50, p90, p99, lats[0], lats[n - 1])
};
let (qn, qok, qfib, qavg, qp50, qp90, qp99, qmin, qmax) = stats(&quic_rows);
let (ln, lok, lfib, lavg, lp50, lp90, lp99, lmin, lmax) = stats(&local_rows);
let qrps = qn as f64 / wall1.as_secs_f64();
let lrps = ln as f64 / wall2.as_secs_f64();
let quic_count = quic_rows.iter().filter(|r| r.transport == "quic").count();
let local_count = local_rows.iter().filter(|r| r.transport == "local").count();
println!("\n{}", "═".repeat(72));
println!(
" QUIC Throughput Benchmark — {} requests × 2 phases",
total_req
);
println!("{}", "═".repeat(72));
println!("\n Cluster Topology");
println!(" {}", "─".repeat(68));
println!(
" Gateway (Node01) http=0.0.0.0:{} quic=0.0.0.0:{} workers=NONE",
b1_port, q1_port
);
println!(
" Executor (Node02) http=0.0.0.0:{} quic=0.0.0.0:{} worker=127.0.0.3:{}",
b2_port, q2_port, w2_port
);
println!(
" Owner: {} (zakuro_user_id={})",
username, zakuro_uid
);
println!(" Protocol: QUIC (quinn, multiplexed UDP streams)");
println!(
"\n Phase 1: QUIC Relay ({} requests → Gateway → QUIC → Executor)",
total_req
);
println!(" {}", "─".repeat(68));
println!(
" Transport: {} QUIC / {} other",
quic_count,
qn - quic_count
);
println!(
" Success: {}/{} ({:.1}%) fib(4)=3: {}/{}",
qok,
qn,
qok as f64 / qn as f64 * 100.0,
qfib,
qn
);
println!(
" Wall clock: {:.2}s Throughput: {:.1} req/s",
wall1.as_secs_f64(),
qrps
);
println!(
" Latency: min={:.2} avg={:.2} p50={:.2} p90={:.2} p99={:.2} max={:.2}",
qmin, qavg, qp50, qp90, qp99, qmax
);
println!(
"\n Phase 2: Local Baseline ({} requests → Executor directly)",
total_req
);
println!(" {}", "─".repeat(68));
println!(
" Transport: {} local / {} other",
local_count,
ln - local_count
);
println!(
" Success: {}/{} ({:.1}%) fib(4)=3: {}/{}",
lok,
ln,
lok as f64 / ln as f64 * 100.0,
lfib,
ln
);
println!(
" Wall clock: {:.2}s Throughput: {:.1} req/s",
wall2.as_secs_f64(),
lrps
);
println!(
" Latency: min={:.2} avg={:.2} p50={:.2} p90={:.2} p99={:.2} max={:.2}",
lmin, lavg, lp50, lp90, lp99, lmax
);
println!("\n Comparison");
println!(" {}", "─".repeat(68));
let overhead = if lavg > 0.0 {
(qavg - lavg) / lavg * 100.0
} else {
0.0
};
println!(
" QUIC relay overhead: {:.2}ms avg ({:+.1}% vs local)",
qavg - lavg,
overhead
);
println!(
" Throughput ratio: {:.2}x (local/QUIC)",
lrps / qrps.max(0.01)
);
println!("\n{}\n", "═".repeat(72));
assert_eq!(qok, qn, "all QUIC-relayed requests must succeed");
assert_eq!(qfib, qn, "all QUIC fibonacci results must be correct");
assert_eq!(lok, ln, "all local requests must succeed");
assert_eq!(lfib, ln, "all local fibonacci results must be correct");
assert!(
qrps > 5.0,
"QUIC throughput must exceed 5 req/s, got {:.1}",
qrps
);
assert!(
quic_count > 0,
"at least some requests must use QUIC transport (got {})",
quic_count
);
}
#[test]
fn test_standalone_billing_cycle_reserve_commit() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
ledger
.local_credits
.insert("user_billing_01".to_string(), 50.0);
let res = ledger.reserve("user_billing_01", 10.0, "req-standalone-01");
assert!(
res.is_ok(),
"reserve must succeed after balance is in local_credits"
);
let balance_mid = ledger
.local_credits
.get("user_billing_01")
.map(|v| *v)
.unwrap_or(0.0);
assert_eq!(balance_mid, 40.0, "10 credits held in reservation");
let balance_after = ledger.commit(&res.unwrap(), 7.5).unwrap();
assert!(
(balance_after - 42.5).abs() < 1e-9,
"balance_after should be 42.5, got {}",
balance_after
);
}
#[test]
fn test_standalone_billing_cancel_full_refund() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
ledger
.local_credits
.insert("user_billing_02".to_string(), 30.0);
let res_id = ledger
.reserve("user_billing_02", 20.0, "req-cancel-01")
.unwrap();
assert_eq!(
ledger
.local_credits
.get("user_billing_02")
.map(|v| *v)
.unwrap_or(0.0),
10.0
);
ledger.cancel(&res_id).unwrap();
assert_eq!(
ledger
.local_credits
.get("user_billing_02")
.map(|v| *v)
.unwrap_or(0.0),
30.0,
"full refund on cancel"
);
}
#[test]
fn test_standalone_billing_insufficient_credits() {
use crate::broker::ledger::{Ledger, LedgerError};
let ledger = Ledger::new(None, None);
ledger.local_credits.insert("user_broke".to_string(), 2.5);
let result = ledger.reserve("user_broke", 5.0, "req-broke-01");
match result {
Err(LedgerError::InsufficientCredits {
required,
available,
}) => {
assert_eq!(required, 5.0);
assert_eq!(available, 2.5);
}
_ => panic!("expected InsufficientCredits, got {:?}", result.map(|_| ())),
}
assert_eq!(
ledger
.local_credits
.get("user_broke")
.map(|v| *v)
.unwrap_or(0.0),
2.5
);
}
#[test]
fn test_reserve_uses_authoritative_balance_as_fallback() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
ledger.local_add_credits("user_authbal", 80.0);
let res = ledger.reserve("user_authbal", 20.0, "req-authbal-01");
assert!(
res.is_ok(),
"reserve should work via authoritative_balances fallback"
);
}
#[test]
fn test_cost_calculation_price_consistency() {
use crate::broker::worker::WorkerPricing;
let pricing = WorkerPricing {
price_per_hour: 3.6,
min_charge: 0.001,
};
let cost_10s = pricing.estimate_cost(10.0);
assert!(
(cost_10s - 0.01).abs() < 1e-9,
"10s at 3.6/hr should cost exactly 0.01, got {}",
cost_10s
);
let cost_tiny = pricing.estimate_cost(0.1);
assert_eq!(cost_tiny, pricing.min_charge, "min_charge must apply");
let initial = 100.0_f64;
let reserved = pricing.estimate_cost(300.0); let actual = pricing.estimate_cost(47.0); let balance_after = initial - actual;
let simulated = initial - reserved + (reserved - actual);
assert!(
(simulated - balance_after).abs() < 1e-10,
"reserve+refund must equal direct deduction"
);
}
#[test]
fn test_p2p_authority_is_deterministic() {
use crate::broker::peer::{Authority, PeerManager};
let peers = vec!["100.64.0.2:9000".to_string()];
let pm = PeerManager::new(
Some("100.64.0.1"),
&peers,
9000,
"test-key".to_string(),
true,
);
for _ in 0..100 {
let a1 = pm.determine_authority("9000000001");
let a2 = pm.determine_authority("9000000001");
let label = |a: &Authority| match a {
Authority::Local => "local",
Authority::Peer(_) => "peer",
Authority::Standalone => "standalone",
};
assert_eq!(
label(&a1),
label(&a2),
"authority changed between calls for same user"
);
}
}
#[test]
fn test_p2p_authority_peer_url_stable() {
use crate::broker::peer::{Authority, PeerManager};
let peers = vec![
"100.64.0.3:9000".to_string(),
"100.64.0.2:9000".to_string(),
"100.64.0.4:9000".to_string(),
];
let pm1 = PeerManager::new(Some("100.64.0.1"), &peers, 9000, "k".to_string(), true);
let pm2 = PeerManager::new(Some("100.64.0.1"), &peers, 9000, "k".to_string(), true);
for uid in &["uid_alpha", "uid_beta", "uid_gamma", "uid_delta"] {
let a1 = pm1.determine_authority(uid);
let a2 = pm2.determine_authority(uid);
if let (Authority::Peer(u1), Authority::Peer(u2)) = (&a1, &a2) {
assert_eq!(
u1, u2,
"user {} peer URL must be stable: {} vs {}",
uid, u1, u2
);
}
let type1 = match &a1 {
Authority::Local => 0,
Authority::Peer(_) => 1,
Authority::Standalone => 2,
};
let type2 = match &a2 {
Authority::Local => 0,
Authority::Peer(_) => 1,
Authority::Standalone => 2,
};
assert_eq!(type1, type2, "user {} authority type must be stable", uid);
}
}
#[test]
fn test_p2p_authority_distribution_two_brokers() {
use crate::broker::peer::{Authority, PeerManager};
let pm = PeerManager::new(
Some("100.64.0.1"),
&["100.64.0.2:9000".to_string()],
9000,
"k".to_string(),
true,
);
let mut local = 0u32;
let mut remote = 0u32;
for i in 9000000001u64..9000000101 {
match pm.determine_authority(&i.to_string()) {
Authority::Local => local += 1,
_ => remote += 1,
}
}
assert!(local >= 25, "expected ~50 local, got {}", local);
assert!(remote >= 25, "expected ~50 remote, got {}", remote);
assert_eq!(local + remote, 100);
}
#[test]
fn test_p2p_local_reserve_commit_balance_consistency() {
use crate::broker::ledger::Ledger;
let ledger = Ledger::new(None, None);
ledger.local_add_credits("user_p2p_01", 200.0);
let (res_id, bal_before) = ledger
.local_reserve("user_p2p_01", 50.0, "p2p-req-01")
.unwrap();
assert_eq!(bal_before, 200.0);
let bal_after = ledger.local_commit(&res_id, 35.0).unwrap();
assert!(
(bal_after - 165.0).abs() < 1e-9,
"balance after commit: expected 165.0, got {}",
bal_after
);
assert!(
((200.0 - bal_after) - 35.0).abs() < 1e-9,
"only actual_cost should be debited"
);
}
#[test]
fn test_two_broker_balance_sheet_consistency() {
use crate::broker::ledger::Ledger;
use crate::broker::worker::WorkerPricing;
let ledger_a = Ledger::new(None, None);
ledger_a.local_add_credits("user_X", 100.0);
let price = WorkerPricing {
price_per_hour: 3.6,
min_charge: 0.001,
};
let timeout_reservation = price.estimate_cost(300.0);
let (res_id, bal_before) = ledger_a
.local_reserve("user_X", timeout_reservation, "cross-001")
.unwrap();
assert_eq!(bal_before, 100.0);
let actual_cost = price.estimate_cost(47.3);
let bal_after = ledger_a.local_commit(&res_id, actual_cost).unwrap();
let expected = 100.0 - actual_cost;
assert!(
(bal_after - expected).abs() < 1e-9,
"balance mismatch: expected {:.6}, got {:.6}",
expected,
bal_after
);
let total_debited = 100.0 - bal_after;
assert!(
(total_debited - actual_cost).abs() < 1e-9,
"dashboard must see actual_cost ({:.6}), not reservation ({:.6})",
actual_cost,
timeout_reservation
);
}
#[test]
fn test_concurrent_reservations_no_overdraft() {
use crate::broker::ledger::Ledger;
use std::sync::Arc;
let ledger = Arc::new(Ledger::new(None, None));
ledger
.local_credits
.insert("user_concurrent".to_string(), 5.0);
let handles: Vec<_> = (0..20)
.map(|i| {
let l = Arc::clone(&ledger);
thread::spawn(move || {
let _ = l.reserve("user_concurrent", 1.0, &format!("req-c{}", i));
})
})
.collect();
for h in handles {
h.join().unwrap();
}
let bal = ledger
.local_credits
.get("user_concurrent")
.map(|v| *v)
.unwrap_or(0.0);
assert!(bal >= 0.0, "balance went negative: {}", bal);
assert!(bal <= 5.0, "balance exceeded initial: {}", bal);
}
#[test]
fn test_round_robin_even_distribution() {
use crate::broker::credits::CreditManager;
use crate::broker::router::{ResourceRequirements, Router, RoutingStrategy};
use crate::broker::worker::{
WorkerPricing, WorkerRegistration, WorkerRegistry, WorkerResources,
};
let registry = WorkerRegistry::new();
let credits = CreditManager::new();
for i in 0..2 {
let reg = WorkerRegistration {
name: format!("rr-worker-{}", i),
uri: format!("http://127.0.0.1:{}", 13960 + i),
worker_type: "zakuro".to_string(),
resources: WorkerResources {
cpus_available: 4.0,
cpus_total: 4.0,
memory_available: 8 * 1024 * 1024 * 1024,
memory_total: 8 * 1024 * 1024 * 1024,
gpus_available: 0,
gpus_total: 0,
},
pricing: WorkerPricing {
price_per_hour: 3.6,
min_charge: 0.001,
},
tags: vec![],
max_timeout_secs: 0.0,
hardware: Default::default(),
wireguard_ip: None,
is_docker: None,
source_node: None,
explicit_local: false,
provider_type: Default::default(),
served_models: vec![],
price_per_mtok: 0.0,
};
registry.register(reg);
}
credits.get_or_create("user_rr", 1000.0);
let router = Router::new();
let req = ResourceRequirements {
strategy: RoutingStrategy::RoundRobin,
..Default::default()
};
let mut counts: std::collections::HashMap<String, u32> = std::collections::HashMap::new();
for _ in 0..20 {
if let Ok(d) = router.select_worker(®istry, &credits, "user_rr", 1000.0, &req) {
*counts.entry(d.worker.name.clone()).or_insert(0) += 1;
}
}
assert_eq!(counts.len(), 2, "both workers should receive requests");
for (name, count) in &counts {
assert!(
*count >= 5,
"worker {} only selected {} times in 20",
name,
count
);
}
}
#[test]
fn test_wal_recovery_cancels_reserved_entries() {
use crate::broker::ledger::Ledger;
use crate::broker::recovery;
use crate::broker::wal::{Wal, WalEntry, WalStatus};
let tmp = format!("/tmp/zc2-wal-test-{}.jsonl", uuid::Uuid::new_v4());
let wal = Wal::open(&tmp).expect("WAL open failed");
let ledger = Ledger::new(None, None);
ledger
.local_credits
.insert("user_wal_crash".to_string(), 0.0);
let entry = WalEntry {
request_id: "wal-crash-req-001".to_string(),
user_id: "user_wal_crash".to_string(),
reservation_id: "wal-crash-res-001".to_string(),
estimated_cost: 7.5,
actual_cost: None,
worker_id: "worker-x".to_string(),
duration_ms: None,
timestamp: chrono::Utc::now(),
status: WalStatus::Reserved,
};
wal.append(&entry).expect("append failed");
wal.flush_buffer().expect("flush failed");
let tx_buffer = crate::broker::flush::TransactionBuffer::new();
recovery::replay_wal(&wal, &ledger, &tx_buffer, None);
let balance = ledger
.local_credits
.get("user_wal_crash")
.map(|v| *v)
.unwrap_or(0.0);
assert_eq!(
balance, 7.5,
"crashed reservation must be refunded on recovery"
);
let _ = std::fs::remove_file(&tmp);
}
#[test]
fn test_two_broker_round_robin_credit_accounting() {
let port_a = free_port();
let port_b = free_port();
let cfg_a = {
let mut c = test_config(port_a);
c.enable_p2p = true;
c.peer_key = Some("shared-peer-key".to_string());
c.discovery.peers = vec![format!("127.0.0.1:{}", port_b)];
c
};
let cfg_b = {
let mut c = test_config(port_b);
c.enable_p2p = true;
c.peer_key = Some("shared-peer-key".to_string());
c.discovery.peers = vec![format!("127.0.0.1:{}", port_a)];
c
};
thread::spawn(move || {
let _ = crate::broker::server::start_server(cfg_a);
});
thread::spawn(move || {
let _ = crate::broker::server::start_server(cfg_b);
});
let url_a = format!("http://127.0.0.1:{}", port_a);
let url_b = format!("http://127.0.0.1:{}", port_b);
for (url, port) in [(&url_a, port_a), (&url_b, port_b)] {
let mut up = false;
for _ in 0..300 {
if ureq::get(&format!("{}/health", url))
.config()
.timeout_global(Some(Duration::from_millis(150)))
.build()
.call()
.is_ok()
{
up = true;
break;
}
thread::sleep(Duration::from_millis(100));
}
assert!(up, "Broker on port {} did not start within 30s", port);
}
add_credits(&url_a, "user_2broker", 100.0);
register_worker(&url_a, "node-a-worker", 4.0, 8, 3.6);
register_worker(&url_b, "node-b-worker", 4.0, 8, 7.2);
thread::sleep(Duration::from_millis(500));
let health_a = ureq::get(&format!("{}/peer/health", url_a))
.header("X-Peer-Key", "shared-peer-key")
.config()
.timeout_global(Some(Duration::from_secs(3)))
.build()
.call();
assert!(
health_a
.map(|r| r.status().as_u16() == 200)
.unwrap_or(false),
"broker_a /peer/health must respond 200"
);
let health_b = ureq::get(&format!("{}/peer/health", url_b))
.header("X-Peer-Key", "shared-peer-key")
.config()
.timeout_global(Some(Duration::from_secs(3)))
.build()
.call();
assert!(
health_b
.map(|r| r.status().as_u16() == 200)
.unwrap_or(false),
"broker_b /peer/health must respond 200"
);
let workers_a_via_peer = ureq::get(&format!("{}/peer/workers", url_a))
.header("X-Peer-Key", "shared-peer-key")
.config()
.timeout_global(Some(Duration::from_secs(3)))
.build()
.call()
.unwrap()
.into_body()
.read_json::<serde_json::Value>()
.unwrap();
assert!(
workers_a_via_peer["workers"]
.as_array()
.map(|a| !a.is_empty())
.unwrap_or(false),
"broker_a must expose its workers via /peer/workers"
);
let workers_b_via_peer = ureq::get(&format!("{}/peer/workers", url_b))
.header("X-Peer-Key", "shared-peer-key")
.config()
.timeout_global(Some(Duration::from_secs(3)))
.build()
.call()
.unwrap()
.into_body()
.read_json::<serde_json::Value>()
.unwrap();
assert!(
workers_b_via_peer["workers"]
.as_array()
.map(|a| !a.is_empty())
.unwrap_or(false),
"broker_b must expose its workers via /peer/workers"
);
}
#[cfg(unix)]
mod agent_it {
use super::{free_port, port_range};
use crate::agent::config::AgentConfig;
use crate::agent::core::tests::{hub_json, row};
use crate::agent::files::tests::tmp;
use crate::agent::hub::tests::mock_hub;
use crate::agent::run::start_background;
use crate::agent::testbin::self_exec_argv;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
const TEST_MAX_WORKERS: u32 = 3;
fn port_block() -> u16 {
port_range(TEST_MAX_WORKERS as u16 + 20)
}
fn agent_cfg(workers: u32, sharing: bool) -> (AgentConfig, PathBuf) {
let state = tmp("agent-it");
std::fs::create_dir_all(state.join("agent")).unwrap();
std::fs::write(
state.join("agent").join("state.json"),
format!(r#"{{"sharing":{sharing},"workers":{workers},"known_worker_ids":[]}}"#),
)
.unwrap();
let mut cfg = AgentConfig::for_dirs(state.clone());
cfg.port = 0;
cfg.broker_port = free_port();
cfg.base_worker_port = port_block();
cfg.max_workers = TEST_MAX_WORKERS;
cfg.broker_argv = self_exec_argv("broker_main");
cfg.worker_template = self_exec_argv("fake_worker_main");
cfg.worker_template_overridden = true;
cfg.child_env = vec![
("ZAKURO_HOME".into(), state.display().to_string()),
(
"ZAKURO_WAL_PATH".into(),
state.join("wal.jsonl").display().to_string(),
),
("ZAKURO_SCAN_INTERVAL".into(), "1".into()),
("ZAKURO_P2P".into(), "false".into()),
];
cfg.timings.reconcile = Duration::from_millis(300);
cfg.timings.hub_poll = Duration::from_millis(500);
cfg.timings.kill_grace = Duration::from_secs(2);
cfg.timings.broker_stop_wait = Duration::from_secs(5);
cfg.timings.shutdown_wait = Duration::from_secs(3);
(cfg, state)
}
fn sign_in(state: &Path, hub: &str) {
std::fs::write(
state.join("credentials"),
format!("api_key=zk_1_test\napi_url={hub}\n"),
)
.unwrap();
}
fn hanging_hub() -> String {
let listener = std::net::TcpListener::bind(("127.0.0.1", 0)).unwrap();
let port = listener.local_addr().unwrap().port();
std::thread::spawn(move || {
for stream in listener.incoming().flatten() {
std::thread::sleep(Duration::from_secs(120));
drop(stream);
}
});
format!("http://127.0.0.1:{port}")
}
struct Agent {
port: u16,
token: String,
stop: Option<std::sync::mpsc::Sender<()>>,
thread: Option<std::thread::JoinHandle<()>>,
}
impl Drop for Agent {
fn drop(&mut self) {
drop(self.stop.take());
if let Some(t) = self.thread.take() {
let _ = t.join();
}
}
}
fn start(cfg: AgentConfig) -> Agent {
let (port, token, stop, thread) = start_background(cfg).expect("agent starts");
Agent {
port,
token,
stop: Some(stop),
thread: Some(thread),
}
}
fn http() -> ureq::Agent {
ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.http_status_as_error(false)
.timeout_global(Some(Duration::from_secs(10)))
.build(),
)
}
fn summary(a: &Agent) -> serde_json::Value {
http()
.get(&format!("http://127.0.0.1:{}/v1/summary", a.port))
.header("Authorization", format!("Bearer {}", a.token).as_str())
.call()
.unwrap()
.into_body()
.read_json()
.unwrap()
}
fn send(
a: &Agent,
method: &str,
path: &str,
body: serde_json::Value,
) -> (u16, serde_json::Value) {
let url = format!("http://127.0.0.1:{}{path}", a.port);
let auth = format!("Bearer {}", a.token);
let req = if method == "PUT" {
http().put(&url)
} else {
http().post(&url)
};
let resp = req
.header("Authorization", auth.as_str())
.send_json(body)
.unwrap();
let status = resp.status().as_u16();
(status, resp.into_body().read_json().unwrap_or_default())
}
fn wait_until(what: &str, timeout: Duration, mut done: impl FnMut() -> bool) {
let end = Instant::now() + timeout;
while Instant::now() < end {
if done() {
return;
}
std::thread::sleep(Duration::from_millis(200));
}
panic!("timed out waiting for: {what}");
}
fn broker_workers(broker_port: u16) -> Vec<serde_json::Value> {
let body: serde_json::Value = match http()
.get(&format!("http://127.0.0.1:{broker_port}/workers"))
.call()
{
Ok(r) => r.into_body().read_json().unwrap_or_default(),
Err(_) => return vec![],
};
body["workers"].as_array().cloned().unwrap_or_default()
}
fn broker_worker_names(broker_port: u16) -> Vec<String> {
let mut names: Vec<String> = broker_workers(broker_port)
.iter()
.filter_map(|w| w["name"].as_str().map(String::from))
.collect();
names.sort();
names
}
fn worker_status(broker_port: u16, name: &str) -> Option<String> {
broker_workers(broker_port)
.into_iter()
.find(|w| w["name"] == name)
.and_then(|w| w["status"].as_str().map(String::from))
}
fn worker_id(broker_port: u16, name: &str) -> Option<String> {
broker_workers(broker_port)
.into_iter()
.find(|w| w["name"] == name)
.and_then(|w| w["id"].as_str().map(String::from))
}
fn fake_worker_count(worker_port: u16, endpoint: &str) -> u64 {
http()
.get(&format!("http://127.0.0.1:{worker_port}/{endpoint}"))
.call()
.unwrap()
.into_body()
.read_json::<serde_json::Value>()
.unwrap()["count"]
.as_u64()
.unwrap()
}
fn fake_worker_execute_count(worker_port: u16) -> u64 {
fake_worker_count(worker_port, "execute_count")
}
fn fake_worker_info_count(worker_port: u16) -> u64 {
fake_worker_count(worker_port, "info_count")
}
fn broker_execute(broker_port: u16) {
let resp = http()
.post(&format!("http://127.0.0.1:{broker_port}/execute"))
.send_json(serde_json::json!({"fn": "noop", "args": []}))
.unwrap();
assert_eq!(
resp.status().as_u16(),
200,
"the broker must accept a keyless local-mode /execute"
);
}
#[test]
fn agent_serves_a_summary_merged_from_the_hub_and_its_broker() {
let (mut cfg, state) = agent_cfg(1, true);
let pk = cfg.node_pubkey();
let (hub, _) = mock_hub(move |_, url, _| {
if url.starts_with("/api/provide/summary") {
(200, hub_json(&pk, &row(11, Some(18.0))))
} else {
(404, "{}".into())
}
});
sign_in(&state, &hub);
cfg.api_url = hub.clone();
let a = start(cfg);
wait_until(
"the broker and one worker running",
Duration::from_secs(40),
|| {
let s = summary(&a);
s["this_mac"]["broker"] == "running" && s["this_mac"]["workers"]["running"] == 1
},
);
let s = summary(&a);
assert_eq!(s["v"], 1);
assert_eq!(s["env"], "custom");
assert_eq!(s["account"]["username"], "jean");
assert_eq!(s["devices"][0]["this_mac"], true);
assert_eq!(s["prices"]["this_mac"]["state"], "set");
assert_eq!(s["links"]["hub"], "https://stg.hub.zakuro-ai.com/provide");
let children: crate::agent::files::Children =
crate::agent::files::load_json(&state.join("agent").join("children.json")).unwrap();
assert!(children.broker.is_some() && children.workers.contains_key("0"));
}
#[test]
fn scaling_two_to_one_drains_then_stops_the_highest_worker_and_pause_stops_all() {
let (mut cfg, state) = agent_cfg(2, true);
let pk = cfg.node_pubkey();
let (hub, _) = mock_hub(move |_, _, _| (200, hub_json(&pk, "")));
sign_in(&state, &hub);
cfg.api_url = hub;
let (broker_port, fp8) = (cfg.broker_port, cfg.node_fp8());
let (w0, w1) = (format!("{fp8}-w0"), format!("{fp8}-w1"));
let a = start(cfg);
wait_until(
"the broker lists w0 and w1",
Duration::from_secs(40),
|| broker_worker_names(broker_port) == vec![w0.clone(), w1.clone()],
);
let (status, s) = send(&a, "PUT", "/v1/workers", serde_json::json!({"count": 1}));
assert_eq!(
(status, s["this_mac"]["workers"]["desired"].as_u64()),
(200, Some(1))
);
wait_until(
"w1 drained, stopped and deleted from the broker",
Duration::from_secs(30),
|| {
broker_worker_names(broker_port) == vec![w0.clone()]
&& summary(&a)["this_mac"]["workers"]["running"] == 1
},
);
let (status, _) = send(&a, "PUT", "/v1/sharing", serde_json::json!({"on": false}));
assert_eq!(status, 200);
wait_until("everything stopped", Duration::from_secs(30), || {
let s = summary(&a);
s["this_mac"]["broker"] == "stopped" && s["this_mac"]["workers"]["running"] == 0
});
}
#[test]
fn draining_a_worker_survives_discovery_rescans_and_stops_new_execute_calls() {
let local_mode = crate::broker::discovery::detect_discovery_mode(
&crate::broker::discovery::DiscoveryConfig::default().subnet,
) == crate::broker::discovery::DiscoveryMode::Local;
let (mut cfg, state) = agent_cfg(2, true);
let pk = cfg.node_pubkey();
let (hub, _) = mock_hub(move |_, _, _| (200, hub_json(&pk, "")));
sign_in(&state, &hub);
cfg.api_url = hub;
let (broker_port, base_worker_port, fp8) =
(cfg.broker_port, cfg.base_worker_port, cfg.node_fp8());
let (w0, w1) = (format!("{fp8}-w0"), format!("{fp8}-w1"));
let (w0_port, w1_port) = (base_worker_port, base_worker_port + 1);
let _a = start(cfg);
wait_until(
"the broker lists w0 and w1",
Duration::from_secs(40),
|| broker_worker_names(broker_port) == vec![w0.clone(), w1.clone()],
);
let id1 = worker_id(broker_port, &w1).expect("w1 registered");
let resp = http()
.post(&format!(
"http://127.0.0.1:{broker_port}/workers/{id1}/drain"
))
.send_empty()
.unwrap();
assert_eq!(resp.status().as_u16(), 200);
assert_eq!(worker_status(broker_port, &w1).as_deref(), Some("draining"));
let execute_at_drain = fake_worker_execute_count(w1_port);
let info_at_drain = fake_worker_info_count(w1_port);
wait_until(
"at least 3 more discovery scans reach w1",
Duration::from_secs(15),
|| fake_worker_info_count(w1_port) >= info_at_drain + 3,
);
assert_eq!(
worker_status(broker_port, &w1).as_deref(),
Some("draining"),
"discovery must not lift a draining worker back to healthy"
);
if !local_mode {
eprintln!(
"skipping the /execute positive control in \
draining_a_worker_survives_discovery_rescans_and_stops_new_execute_calls: \
this host is on the WireGuard mesh, so the test broker's /execute is not \
keyless local-mode and would 401"
);
return;
}
let w0_before = fake_worker_execute_count(w0_port);
const N: u64 = 3;
for _ in 0..N {
broker_execute(broker_port);
}
assert_eq!(
fake_worker_execute_count(w0_port),
w0_before + N,
"dispatch must still reach the healthy worker"
);
assert_eq!(
fake_worker_execute_count(w1_port),
execute_at_drain,
"a draining worker must receive no /execute calls"
);
}
#[test]
fn shutdown_signals_a_term_ignoring_worker_promptly_even_with_a_hanging_hub() {
let (mut cfg, state) = agent_cfg(1, true);
let hub = hanging_hub();
sign_in(&state, &hub);
cfg.api_url = hub;
cfg.worker_template = vec![
"sh".into(),
"-c".into(),
"trap '' TERM; exec sleep 30".into(),
];
cfg.worker_template_overridden = true;
cfg.timings.hub_timeout = Duration::from_secs(20);
cfg.timings.hub_poll = Duration::from_millis(300);
cfg.timings.shutdown_wait = Duration::from_secs(2);
let shutdown_wait = cfg.timings.shutdown_wait;
let a = start(cfg);
wait_until("the worker is running", Duration::from_secs(15), || {
summary(&a)["this_mac"]["workers"]["running"] == 1
});
let children: crate::agent::files::Children =
crate::agent::files::load_json(&state.join("agent").join("children.json")).unwrap();
let pid = children.workers.get("0").expect("worker 0 recorded").pid;
assert!(
crate::agent::supervisor::is_alive(pid),
"the worker must be running before shutdown"
);
let began = Instant::now();
drop(a);
let elapsed = began.elapsed();
let in_flight_tick = 2 * crate::agent::supervisor::BROKER_CHECK_TIMEOUT;
assert!(
elapsed < shutdown_wait + in_flight_tick + Duration::from_secs(3),
"shutdown took {elapsed:?}, expected to finish within shutdown_wait \
({shutdown_wait:?}) plus an in-flight tick's broker checks \
({in_flight_tick:?}) plus a small margin, even with a hanging hub \
and a TERM-ignoring worker"
);
assert!(
!crate::agent::supervisor::is_alive(pid),
"the TERM-ignoring worker group must be gone after shutdown"
);
let children_after: crate::agent::files::Children =
crate::agent::files::load_json(&state.join("agent").join("children.json")).unwrap();
assert!(
children_after.broker.is_none() && children_after.workers.is_empty(),
"no reconcile tick may respawn children after shutdown: {children_after:?}"
);
}
#[test]
fn device_price_fans_out_over_http_and_seeds_a_new_row() {
let (mut cfg, state) = agent_cfg(1, false); let pk = cfg.node_pubkey();
let rows = Arc::new(Mutex::new(format!("{},{}", row(11, None), row(12, None))));
let r = rows.clone();
let (hub, seen) = mock_hub(move |method, url, _| match method {
"GET" if url.starts_with("/api/provide/summary") => {
(200, hub_json(&pk, &r.lock().unwrap()))
}
"PUT" => (200, "{}".into()),
_ => (404, "{}".into()),
});
sign_in(&state, &hub);
cfg.api_url = hub;
let a = start(cfg);
wait_until("the first hub poll", Duration::from_secs(10), || {
summary(&a)["account"]["username"] == "jean"
});
let (status, _) = send(
&a,
"PUT",
"/v1/price",
serde_json::json!({"scope": "device", "price_per_hour": 18}),
);
assert_eq!(status, 200);
let put_urls = || -> Vec<String> {
seen.lock()
.unwrap()
.iter()
.filter(|e| e.0 == "PUT")
.map(|e| e.1.clone())
.collect()
};
let mut fanned_out = put_urls();
fanned_out.sort();
assert_eq!(
fanned_out,
vec!["/api/workers/11/price", "/api/workers/12/price"]
);
*rows.lock().unwrap() = format!(
"{},{},{}",
row(11, Some(18.0)),
row(12, Some(18.0)),
row(13, None)
);
let (status, _) = send(&a, "POST", "/v1/refresh", serde_json::json!({}));
assert_eq!(status, 200);
wait_until("row 13 seeded at 18", Duration::from_secs(10), || {
put_urls().last().map(String::as_str) == Some("/api/workers/13/price")
});
}
}
}