use std::io::Read;
use std::net::SocketAddr;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::Arc;
use std::thread;
use std::time::Instant;
use crate::async_exec;
use chrono::Local;
use colored::Colorize;
use serde::{Deserialize, Serialize};
use tiny_http::{Header, Method};
use super::http_adapter::{BrokerRequest, BrokerResponse};
use super::{
discovery::{detect_discovery_mode, Discovery, DiscoveryMode},
flush::BufferedTransaction,
ledger, node_identity,
peer::{self, Authority},
recovery,
router::{ResourceRequirements, RoutingDecision},
stats::{StatsResponse, TaskOfferRecord, TransactionRecord, TransactionStatus, WorkerStats},
task_board,
wal::{WalEntry, WalStatus},
worker::{self, ModelResolveError, Worker, WorkerHeartbeat, WorkerRegistration, WorkerStatus},
BrokerConfig, BrokerState,
};
use crate::model_uri::parse_model_uri;
static REQUEST_COUNTER: AtomicU64 = AtomicU64::new(0);
const PENDING_OFFER_TTL_SECS: u64 = 10;
fn header_or(name: &str, value: &str) -> (String, String) {
if Header::from_bytes(name.as_bytes(), value.as_bytes()).is_ok() {
(name.to_string(), value.to_string())
} else {
("X-Zakuro-Invalid".to_string(), "1".to_string())
}
}
fn capped_charge(computed: f64, reserved: f64) -> f64 {
if computed.is_finite() {
computed.clamp(0.0, reserved.max(0.0))
} else {
0.0
}
}
fn format_credits(amount: f64) -> String {
if amount >= 1.0 {
format!("{:.4}", amount).green().to_string()
} else if amount >= 0.01 {
format!("{:.4}", amount).yellow().to_string()
} else {
format!("{:.6}", amount).cyan().to_string()
}
}
#[allow(clippy::too_many_arguments)]
fn log_transaction(
state: &Arc<BrokerState>,
verbose: bool,
tx_num: u64,
user: &str,
action: &str,
cost: f64,
balance: f64,
worker: Option<&str>,
duration_ms: f64,
status: &str,
) {
log_transaction_with_worker_info(
state,
verbose,
tx_num,
user,
action,
cost,
balance,
worker,
duration_ms,
status,
None,
None,
None,
None,
None,
);
}
#[allow(clippy::too_many_arguments)]
fn log_transaction_with_worker_info(
state: &Arc<BrokerState>,
verbose: bool,
tx_num: u64,
user: &str,
action: &str,
cost: f64,
balance: f64,
worker: Option<&str>,
duration_ms: f64,
status: &str,
price_per_hour: Option<f64>,
owner_id: Option<&str>,
worker_pid: Option<&str>,
worker_ip: Option<&str>,
request_id: Option<&str>,
) {
let now = Local::now();
let tx_status = match status {
"OK" => TransactionStatus::Ok,
"FAIL" => TransactionStatus::Fail,
"PENDING" => TransactionStatus::Pending,
_ => TransactionStatus::Pending,
};
let record = TransactionRecord {
tx_num,
timestamp: now,
user_id: user.to_string(),
action: action.to_string(),
cost,
balance,
worker: worker.map(|s| s.to_string()),
duration_ms,
status: tx_status,
price_per_hour,
owner_id: owner_id.map(|s| s.to_string()),
worker_pid: worker_pid.map(|s| s.to_string()),
worker_ip: worker_ip.map(|s| s.to_string()),
request_id: request_id.map(|s| s.to_string()),
};
state.stats.record_transaction(record);
if !verbose {
return;
}
let time_str = now.format("%H:%M:%S").to_string();
let status_icon = match status {
"OK" => "✓".green(),
"FAIL" => "✗".red(),
"PENDING" => "⋯".yellow(),
_ => "•".white(),
};
let worker_str = worker
.map(|w| format!(" → {}", w.cyan()))
.unwrap_or_default();
println!(
" {} {} #{:<4} {:>12} {:>8} cost:{} bal:{} {:>6.1}ms{}",
time_str.dimmed(),
status_icon,
tx_num,
user.blue(),
action.bold(),
format_credits(cost),
format_credits(balance),
duration_ms,
worker_str,
);
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkerListResponse {
pub workers: Vec<WorkerInfo>,
pub total: usize,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkerInfo {
pub id: String,
pub name: String,
pub uri: String,
pub worker_type: String,
pub status: String,
pub cpus_available: f64,
pub cpus_total: f64,
pub memory_available_gib: f64,
pub memory_total_gib: f64,
pub gpus_available: u32,
pub gpus_total: u32,
pub price_per_hour: f64,
pub min_charge: f64,
pub active_requests: u32,
pub avg_latency_ms: f64,
pub max_timeout_secs: f64,
pub gpu_model: Option<String>,
pub gpu_vram_gb: Option<u32>,
pub cpu_model: Option<String>,
pub storage_gb: Option<u32>,
pub requests_5h: u64,
pub requests_1w: u64,
pub requests_1m: u64,
pub quota_5h: u64,
pub quota_1w: u64,
pub quota_1m: u64,
pub wireguard_ip: Option<String>,
pub is_docker: Option<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub node: Option<String>,
#[serde(default)]
pub provider_type: worker::ProviderType,
#[serde(default)]
pub served_models: Vec<String>,
#[serde(default)]
pub price_per_mtok: f64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub label: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PriceEstimateResponse {
pub min_cost: f64,
pub max_cost: f64,
pub matching_workers: usize,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CreditBalanceResponse {
pub user_id: String,
pub balance: f64,
pub total_spent: f64,
pub rate_limit: u32,
pub rate_limit_per_second: Option<u32>,
pub rate_limit_per_day: Option<u32>,
pub rate_limit_per_month: Option<u32>,
}
#[derive(Debug, Serialize)]
struct ErrorResponse {
error: String,
code: String,
}
#[derive(Debug, Deserialize)]
struct AddCreditsRequest {
amount: f64,
#[serde(default)]
description: Option<String>,
}
fn json_response<T: Serialize>(data: &T, status: u16) -> BrokerResponse {
let body = serde_json::to_vec(data).unwrap_or_default();
BrokerResponse::json_bytes(body, status)
}
fn error_response(error: &str, code: &str, status: u16) -> BrokerResponse {
json_response(
&ErrorResponse {
error: error.to_string(),
code: code.to_string(),
},
status,
)
}
fn quic_offer(
qt: &std::sync::Arc<super::quic::QuicTransport>,
peer_http_url: &str,
quic_port: u16,
offer: &task_board::TaskOffer,
) -> Result<task_board::TaskResult, String> {
let url_trimmed = peer_http_url
.trim_start_matches("http://")
.trim_start_matches("https://");
let host = url_trimmed.split(':').next().unwrap_or("127.0.0.1");
let connect_host = if host.starts_with("127.") {
"127.0.0.1"
} else {
host
};
let addr: SocketAddr = format!("{}:{}", connect_host, quic_port)
.parse()
.map_err(|e| format!("bad quic addr: {}", e))?;
qt.offer_task(addr, offer)
}
use super::http_adapter::MAX_BODY_BYTES;
fn body_or_413(req: &BrokerRequest) -> Result<&[u8], BrokerResponse> {
if req.body.len() > MAX_BODY_BYTES {
Err(error_response(
"Request body too large",
"PAYLOAD_TOO_LARGE",
413,
))
} else {
Ok(&req.body)
}
}
fn authenticate_bearer(
request: &BrokerRequest,
state: &BrokerState,
) -> Result<String, BrokerResponse> {
match request.header("Authorization") {
Some(auth) => match auth.strip_prefix("Bearer ") {
Some(token) => state
.ledger
.resolve_user_from_api_key(token)
.map_err(|e| error_response(&e.to_string(), "UNAUTHORIZED", 401)),
None => Err(error_response(
"Invalid Authorization header",
"UNAUTHORIZED",
401,
)),
},
None => Err(error_response(
"Authorization required. Set Authorization: Bearer <key>",
"UNAUTHORIZED",
401,
)),
}
}
fn verify_worker_key(request: &BrokerRequest, state: &BrokerState) -> bool {
match &state.config.worker_key {
None => true, Some(expected) => request
.header("X-Worker-Key")
.map(|k| k == *expected)
.unwrap_or(false),
}
}
fn node_handle_for(key: &super::node_identity::NodeKey, _label: Option<&str>) -> String {
key.node_uri()
}
fn worker_to_info(
w: &worker::Worker,
registry: &worker::WorkerRegistry,
node_key: &super::node_identity::NodeKey,
client_facing: bool,
) -> WorkerInfo {
use chrono::Duration as CDuration;
let reqs_5h = registry.requests_in_window(&w.id, CDuration::hours(5));
let reqs_1w = registry.requests_in_window(&w.id, CDuration::weeks(1));
let reqs_1m = registry.requests_in_window(&w.id, CDuration::days(30));
let quotas = registry.get_quotas(&w.id);
let bare = |n: &str| n.strip_prefix("node-").unwrap_or(n).to_string();
let node = Some(match &w.source_node {
Some(n) => format!("zc://node-{}", bare(n)),
None => node_handle_for(node_key, Some(&super::node_name_or_default())),
});
let (uri, wireguard_ip) = if client_facing {
(w.zc_uri(), None)
} else {
(w.uri.clone(), w.wireguard_ip.clone())
};
WorkerInfo {
id: w.id.clone(),
name: w.name.clone(),
uri,
worker_type: w.worker_type.clone(),
status: format!("{:?}", w.status).to_lowercase(),
cpus_available: w.resources.cpus_available,
cpus_total: w.resources.cpus_total,
memory_available_gib: w.resources.memory_available as f64 / (1024.0 * 1024.0 * 1024.0),
memory_total_gib: w.resources.memory_total as f64 / (1024.0 * 1024.0 * 1024.0),
gpus_available: w.resources.gpus_available,
gpus_total: w.resources.gpus_total,
price_per_hour: w.pricing.price_per_hour,
min_charge: w.pricing.min_charge,
active_requests: w.active_requests,
avg_latency_ms: w.avg_latency_ms,
max_timeout_secs: w.max_timeout_secs,
gpu_model: w.hardware.gpu_model.clone(),
gpu_vram_gb: w.hardware.gpu_vram_gb,
cpu_model: w.hardware.cpu_model.clone(),
storage_gb: w.hardware.storage_gb,
requests_5h: reqs_5h,
requests_1w: reqs_1w,
requests_1m: reqs_1m,
quota_5h: quotas.per_5h,
quota_1w: quotas.per_week,
quota_1m: quotas.per_month,
wireguard_ip,
is_docker: w.is_docker,
node,
provider_type: w.provider_type,
served_models: w.served_models.clone(),
price_per_mtok: w.price_per_mtok,
label: Some(w.name.clone()),
}
}
pub(crate) fn mesh_endpoint_for(mesh_ip: Option<String>, port: u16) -> Option<String> {
mesh_ip.map(|ip| format!("{ip}:{port}"))
}
const NO_MESH_ENDPOINT_LOG_INTERVAL_SECS: u64 = 30 * 60;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum MeshEndpointLogAction {
None,
NoEndpoint,
Recovered,
}
pub(crate) struct MeshEndpointLogGate {
had_endpoint: AtomicBool,
last_no_endpoint_log_secs: AtomicU64,
}
impl MeshEndpointLogGate {
pub(crate) const fn new() -> Self {
Self {
had_endpoint: AtomicBool::new(true),
last_no_endpoint_log_secs: AtomicU64::new(0),
}
}
pub(crate) fn on_tick(&self, has_endpoint: bool, now_secs: u64) -> MeshEndpointLogAction {
let had_before = self.had_endpoint.swap(has_endpoint, Ordering::SeqCst);
if has_endpoint {
return if had_before {
MeshEndpointLogAction::None
} else {
MeshEndpointLogAction::Recovered
};
}
if had_before {
self.last_no_endpoint_log_secs
.store(now_secs, Ordering::SeqCst);
return MeshEndpointLogAction::NoEndpoint;
}
let last = self.last_no_endpoint_log_secs.load(Ordering::SeqCst);
if now_secs.saturating_sub(last) >= NO_MESH_ENDPOINT_LOG_INTERVAL_SECS {
self.last_no_endpoint_log_secs
.store(now_secs, Ordering::SeqCst);
MeshEndpointLogAction::NoEndpoint
} else {
MeshEndpointLogAction::None
}
}
}
static MESH_ENDPOINT_LOG_GATE: MeshEndpointLogGate = MeshEndpointLogGate::new();
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct RosterPeer {
pub fp: String,
pub url: String,
}
pub(crate) fn roster_peer_entries(
cache: &super::roster_cache::RosterCache,
self_pubkey: Option<&str>,
self_endpoint: Option<&str>,
) -> Vec<RosterPeer> {
let self_fp = self_pubkey.and_then(super::node_identity::fingerprint_of_pubkey_b64);
cache
.entries()
.iter()
.filter(|(_, revoked)| !revoked)
.filter(|(pk, _)| Some(pk.as_str()) != self_pubkey)
.filter_map(|(pk, _)| {
let fp = super::node_identity::fingerprint_of_pubkey_b64(pk)?;
if self_fp.as_deref() == Some(fp.as_str()) {
return None;
}
let addr = cache.endpoint_for_fingerprint(&fp)?;
Some((fp, addr))
})
.filter(|(_, addr)| Some(addr.as_str()) != self_endpoint)
.filter(|(_, addr)| !addr.trim().is_empty())
.map(|(fp, addr)| {
let url = if addr.contains("://") {
addr
} else {
format!("http://{addr}")
};
RosterPeer { fp, url }
})
.collect()
}
pub fn handle(state: Arc<BrokerState>, req: &BrokerRequest, verbose: bool) -> BrokerResponse {
let path = if req.query.is_empty() {
req.path.clone()
} else {
format!("{}?{}", req.path, req.query)
};
let method = req.method.clone();
let response = match (method, path.as_str()) {
(_, p) if p.starts_with("/serve/") => super::deploy::serve::handle_serve(&state, req),
(Method::Get, "/health") => {
let live_ts_ip =
super::discovery::get_mesh_ip().or_else(|| state.own_wireguard_ip.clone());
let ts_connected = live_ts_ip
.as_ref()
.map(|ip| ip.starts_with("10.13.13."))
.unwrap_or(false);
json_response(
&serde_json::json!({
"status": "healthy",
"service": "zakuro-broker",
"wireguard_connected": ts_connected,
"node_name": state.config.node_name,
"node_id": state.node_key.node_uri(),
"mesh_tunnel": if ts_connected { "up" } else { "down" },
}),
200,
)
}
(Method::Get, path) if path == "/stats" || path.starts_with("/stats?") => {
match authenticate_bearer(req, &state) {
Err(resp) => resp,
Ok(caller) => {
let is_admin = caller == "admin";
let mut user_filter: Option<String> = path.find('?').and_then(|qpos| {
let query = &path[qpos + 1..];
query
.split('&')
.find(|p| p.starts_with("user="))
.map(|p| p[5..].to_string())
});
if !is_admin {
user_filter = Some(caller.clone());
}
let workers = state.workers.list();
let active = workers
.iter()
.filter(|w| w.status == WorkerStatus::Healthy)
.count();
let worker_stats: Vec<WorkerStats> = workers
.iter()
.map(|w| WorkerStats {
id: w.id.clone(),
name: w.name.clone(),
uri: w.zc_uri(),
status: format!("{:?}", w.status).to_lowercase(),
cpus_available: w.resources.cpus_available,
memory_available_gib: w.resources.memory_available as f64
/ (1024.0 * 1024.0 * 1024.0),
gpus_available: w.resources.gpus_available,
price_per_hour: w.pricing.price_per_hour,
active_requests: w.active_requests,
avg_latency_ms: w.avg_latency_ms,
})
.collect();
let wireguard_connected = state.own_wireguard_ip.is_some();
let metrics = state.stats.metrics(
active,
workers.len(),
state.is_local_mode(),
state.is_billing_enabled(),
None,
);
let mut transactions = state.stats.recent_transactions(50);
if let Some(ref uid) = user_filter {
transactions.retain(|tx| tx.user_id == *uid);
}
let stats_resp = StatsResponse {
host: state.config.host.clone(),
port: state.config.port,
transactions,
task_offers: state.stats.recent_task_offers(30),
workers: worker_stats,
metrics,
rps_history: state.stats.rps_history(),
wireguard_ip: None,
wireguard_connected,
};
json_response(&stats_resp, 200)
}
}
}
(Method::Get, path) if path == "/workers" || path.starts_with("/workers?") => {
let node_filter = req.query.split('&').find_map(|kv| kv.strip_prefix("node="));
let reserved = pending_reservations(&state);
let infos: Vec<WorkerInfo> = state
.workers
.list()
.iter()
.map(|w| {
let mut info = worker_to_info(w, &state.workers, &state.node_key, true);
info.active_requests += reserved.get(&w.id).copied().unwrap_or(0);
info
})
.filter(|wi| {
node_filter.is_none_or(|nf| {
wi.node.as_deref().is_some_and(|n| {
node_identity::fp_matches_node_arg(node_identity::strip_node_arg(n), nf)
})
})
})
.collect();
json_response(
&WorkerListResponse {
total: infos.len(),
workers: infos,
},
200,
)
}
(Method::Post, "/workers") => {
if !verify_worker_key(req, &state) {
error_response(
"Worker key required. Set X-Worker-Key header",
"UNAUTHORIZED",
401,
)
} else {
match body_or_413(req) {
Err(resp) => resp,
Ok(body) => match serde_json::from_slice::<WorkerRegistration>(body) {
Ok(registration) => {
if registration.pricing.price_per_hour < 0.0 {
error_response("price_per_hour must be >= 0", "BAD_REQUEST", 400)
} else {
println!(
" [BROKER] Registering worker: {} at {}",
registration.name, registration.uri
);
let mut worker = state.workers.register(registration);
if worker.source_node.is_none() {
if let Some(updated) = state
.workers
.set_node_fp(&worker.id, &state.node_key.fingerprint())
{
worker = updated;
}
}
if let Some(ref owner_id) = state.config.owner_user_id {
let node_name = state.config.node_name.as_deref();
let node_pubkey = state.node_key.public_b64();
if let (Some(ref api_url), Some(ref api_key)) =
(&state.config.api_url, &state.config.api_key)
{
match super::ledger::Ledger::sync_workers_via_api(
owner_id,
std::slice::from_ref(&worker),
api_url,
api_key,
node_name,
state.own_wireguard_ip.as_deref(),
Some(node_pubkey.as_str()),
) {
Ok(outcome) => {
state.workers.apply_hub_prices(&outcome.prices);
println!(" [WORKER_SYNC] Worker {} synced to dashboard via API", worker.name)
}
Err(e) => {
eprintln!(
" [WORKER_SYNC] API sync failed for {}: {}",
worker.name, e
);
}
}
}
}
json_response(&worker, 201)
} }
Err(e) => {
error_response(&format!("Invalid request: {}", e), "BAD_REQUEST", 400)
}
},
}
}
}
(Method::Post, "/workers/heartbeat") => {
if !verify_worker_key(req, &state) {
error_response(
"Worker key required. Set X-Worker-Key header",
"UNAUTHORIZED",
401,
)
} else {
match body_or_413(req) {
Err(resp) => resp,
Ok(body) => match serde_json::from_slice::<WorkerHeartbeat>(body) {
Ok(heartbeat) => match state.workers.heartbeat(heartbeat) {
Some(worker) => json_response(
&serde_json::json!({"status": "ok", "worker_id": worker.id}),
200,
),
None => error_response("Worker not found", "NOT_FOUND", 404),
},
Err(e) => {
error_response(&format!("Invalid request: {}", e), "BAD_REQUEST", 400)
}
},
}
}
}
(Method::Post, "/price") => match body_or_413(req) {
Err(resp) => resp,
Ok(body) => match serde_json::from_slice::<ResourceRequirements>(body) {
Ok(requirements) => {
match state.router.estimate_cost(&state.workers, &requirements) {
Some((min_cost, max_cost)) => {
let matching =
state.router.list_matching(&state.workers, &requirements);
json_response(
&PriceEstimateResponse {
min_cost,
max_cost,
matching_workers: matching.len(),
},
200,
)
}
None => error_response("No workers available", "NO_CAPACITY", 503),
}
}
Err(e) => error_response(&format!("Invalid request: {}", e), "BAD_REQUEST", 400),
},
},
(Method::Get, "/peers") => {
match peer_handshake_auth(&state, req) {
Some(resp) => resp,
None => {
state.peer_manager.ensure_fingerprints_cached();
let mut peer_urls = state.peer_manager.peer_urls();
peer_urls.sort();
let peers: Vec<String> = peer_urls
.into_iter()
.filter_map(|url| state.peer_manager.peer_fingerprint(&url))
.collect();
json_response(
&serde_json::json!({
"node": state.node_key.node_uri(),
"p2p_enabled": state.peer_manager.is_enabled(),
"peers": peers,
"total": peers.len(),
}),
200,
)
}
}
}
(Method::Get, "/brokers") => handle_brokers(&state),
(Method::Get, "/price") => handle_get_price(&state),
(Method::Post, "/price/set") => handle_set_price(&state, req),
(Method::Get, "/peer/health") => handle_peer_health(&state, req),
(Method::Get, "/peer/workers") => handle_peer_workers(&state, req),
(Method::Get, "/peer/peers") => handle_peer_peers(&state, req),
(Method::Get, "/peer/advert") => handle_peer_advert(&state, req),
(Method::Post, "/peer/reserve") => handle_peer_reserve(state.clone(), req),
(Method::Post, "/peer/commit") => handle_peer_commit(state.clone(), req),
(Method::Post, "/peer/cancel") => handle_peer_cancel(state.clone(), req),
(Method::Get, p) if p.starts_with("/peer/balance") => handle_peer_balance(&state, req),
(Method::Post, "/peer/earn") => handle_peer_earn(state.clone(), req),
(Method::Get, "/peer/identity") => handle_peer_identity(&state, req),
(Method::Post, "/peer/tasks/offer") => handle_peer_task_offer(state.clone(), req),
(Method::Post, "/peer/tasks/commit") => handle_peer_task_commit(state.clone(), req),
(Method::Post, "/peer/tasks/cancel") => handle_peer_task_cancel(state.clone(), req),
(Method::Post, "/infer") => handle_infer(&state, req),
(Method::Post, "/execute") => handle_execute(state, req, verbose),
(Method::Get, "/instances") => match authenticate_bearer(req, &state) {
Err(resp) => resp,
Ok(caller) if caller != "admin" => {
error_response("Forbidden: admin only", "FORBIDDEN", 403)
}
Ok(_) => {
let instances: Vec<serde_json::Value> = state
.instance_registry
.iter()
.map(|entry| {
serde_json::json!({
"instance_id": entry.key().clone(),
"worker_id": entry.value().clone(),
})
})
.collect();
json_response(
&serde_json::json!({
"instances": instances,
"total": instances.len(),
}),
200,
)
}
},
(Method::Delete, path) if path.starts_with("/instances/") => {
match authenticate_bearer(req, &state) {
Err(resp) => resp,
Ok(caller) if caller != "admin" => {
error_response("Forbidden: admin only", "FORBIDDEN", 403)
}
Ok(_) => {
let instance_id = path.trim_start_matches("/instances/");
match state.instance_registry.remove(instance_id) {
Some(_) => json_response(
&serde_json::json!({"status": "ok", "instance_id": instance_id}),
200,
),
None => error_response("Instance not found", "NOT_FOUND", 404),
}
}
}
}
(Method::Get, path) if path.starts_with("/credits/") && !path.contains("/add") => {
match authenticate_bearer(req, &state) {
Err(resp) => resp,
Ok(caller) => {
let user_id = path.trim_start_matches("/credits/");
if caller != "admin" && caller != user_id {
error_response("Forbidden: can only view own balance", "FORBIDDEN", 403)
} else {
let info = state.ledger.get_user_info(user_id);
let credits = state.credits.get(info.user_id.as_str());
let live_balance =
credits.as_ref().map(|c| c.balance).unwrap_or(info.balance);
json_response(
&serde_json::json!({
"user_id": info.user_id,
"balance": live_balance,
"balance_status": credits.as_ref().map(|c| c.balance_status.as_str()).unwrap_or("authoritative"),
"last_prefetched": credits.as_ref().and_then(|c| c.last_prefetched).map(|t| t.to_rfc3339()),
}),
200,
)
}
}
}
}
(Method::Post, path) if path.starts_with("/credits/") && path.ends_with("/add") => {
let user_id = path
.trim_start_matches("/credits/")
.trim_end_matches("/add");
if let Some(api_key) = req.header("X-Api-Key") {
match body_or_413(req) {
Err(resp) => resp,
Ok(body) => match serde_json::from_slice::<AddCreditsRequest>(body) {
Ok(req) => {
if state.ledger.is_api_mode() {
error_response("Adding credits directly via broker is not supported in API mode. Use the dashboard API.", "NOT_SUPPORTED", 501)
} else {
let master_key =
std::env::var("ZAKURO_MASTER_KEY").unwrap_or_default();
let authorized = master_key.is_empty() || api_key == master_key;
if !authorized {
error_response("Invalid master key", "UNAUTHORIZED", 401)
} else {
let new_balance = {
let mut entry = state
.ledger
.local_credits
.entry(user_id.to_string())
.or_insert(0.0);
*entry += req.amount;
*entry
};
state
.ledger
.authoritative_balances
.entry(user_id.to_string())
.and_modify(|b| *b += req.amount)
.or_insert(new_balance);
state.ledger.publish_transaction(
&uuid::Uuid::new_v4().to_string(),
user_id,
"credit",
req.amount,
new_balance,
"",
0.0,
state.config.node_name.as_deref(),
);
json_response(
&serde_json::json!({
"status": "ok",
"new_balance": new_balance,
"ledger": "local"
}),
200,
)
}
}
}
Err(e) => {
error_response(&format!("Invalid request: {}", e), "BAD_REQUEST", 400)
}
},
}
} else {
error_response(
"API key required. Set X-Api-Key header with ZAKURO_MASTER_KEY",
"UNAUTHORIZED",
401,
)
}
}
(Method::Get, "/ledger/status") => match authenticate_bearer(req, &state) {
Err(resp) => resp,
Ok(caller) if caller != "admin" => {
error_response("Forbidden: admin only", "FORBIDDEN", 403)
}
Ok(_) => json_response(
&serde_json::json!({
"api_mode": state.ledger.is_api_mode(),
}),
200,
),
},
(Method::Post, path) if path.starts_with("/workers/") && path.ends_with("/drain") => {
let worker_id = path
.trim_start_matches("/workers/")
.trim_end_matches("/drain");
handle_drain_worker(&state, req, worker_id)
}
(Method::Delete, path) if path.starts_with("/workers/") => {
if !verify_worker_key(req, &state) {
error_response(
"Worker key required. Set X-Worker-Key header",
"UNAUTHORIZED",
401,
)
} else {
let worker_id = path.trim_start_matches("/workers/");
match state.workers.remove(worker_id) {
Some(worker) => {
println!(" [BROKER] Worker unregistered: {}", worker_id);
json_response(
&serde_json::json!({"status": "ok", "worker_id": worker.id}),
200,
)
}
None => error_response("Worker not found", "NOT_FOUND", 404),
}
}
}
(Method::Get, "/quotas") => {
match authenticate_bearer(req, &state) {
Err(resp) => resp,
Ok(caller) => {
let user_credits = state.credits.get(&caller);
let now = chrono::Utc::now();
match user_credits {
Some(uc) => {
let sec_reset = 1.0_f64
- now
.signed_duration_since(uc.second_window_start)
.num_milliseconds() as f64
/ 1000.0;
let day_reset = 86400.0
- now.signed_duration_since(uc.day_window_start).num_seconds()
as f64;
let month_reset = 2_592_000.0
- now
.signed_duration_since(uc.month_window_start)
.num_seconds() as f64;
json_response(
&serde_json::json!({
"user_id": caller,
"rate_limit_per_second": uc.rate_limit_per_second,
"rate_limit_per_day": uc.rate_limit_per_day,
"rate_limit_per_month": uc.rate_limit_per_month,
"requests_this_second": uc.requests_this_second,
"requests_this_day": uc.requests_this_day,
"requests_this_month": uc.requests_this_month,
"seconds_until_second_reset": sec_reset.max(0.0),
"seconds_until_day_reset": day_reset.max(0.0),
"seconds_until_month_reset": month_reset.max(0.0),
}),
200,
)
}
None => json_response(
&serde_json::json!({
"user_id": caller,
"rate_limit_per_second": serde_json::Value::Null,
"rate_limit_per_day": serde_json::Value::Null,
"rate_limit_per_month": serde_json::Value::Null,
"requests_this_second": 0,
"requests_this_day": 0,
"requests_this_month": 0,
"seconds_until_second_reset": 0,
"seconds_until_day_reset": 0,
"seconds_until_month_reset": 0,
}),
200,
),
}
}
}
}
(Method::Get, "/me") => {
if let Some(auth) = req.header("Authorization") {
match auth.strip_prefix("Bearer ") {
Some(token) => match state.ledger.resolve_user_from_api_key(token) {
Ok(uid) => {
let balance = state.ledger.load_balance_if_needed(&uid);
json_response(
&serde_json::json!({
"user_id": uid,
"balance": balance,
"local_mode": state.is_local_mode(),
}),
200,
)
}
Err(e) => error_response(&e.to_string(), "UNAUTHORIZED", 401),
},
None => error_response("Invalid Authorization header", "UNAUTHORIZED", 401),
}
} else {
error_response(
"API key required. Set Authorization: Bearer <key>",
"UNAUTHORIZED",
401,
)
}
}
_ => error_response("Not found", "NOT_FOUND", 404),
};
response
}
fn peer_credit_auth(state: &Arc<BrokerState>, request: &BrokerRequest) -> Option<BrokerResponse> {
if let Some(reject) = node_sig_gate(state, request) {
return Some(reject);
}
if peer::verify_peer_key(
request.header("X-Peer-Key").as_deref(),
state.peer_manager.peer_key(),
) {
None
} else {
Some(error_response(
"Peer authentication required (set ZAKURO_PEER_KEY and send X-Peer-Key)",
"UNAUTHORIZED",
401,
))
}
}
fn check_node_sig(
roster: &dashmap::DashMap<String, bool>,
guard: &node_identity::ReplayGuard,
request: &BrokerRequest,
now: u64,
) -> Result<String, String> {
let is_rostered = |id: &str| roster.get(id).map(|r| !*r.value()).unwrap_or(false);
let header = |name: &str| request.header(name);
let method = format!("{:?}", request.method).to_uppercase(); node_identity::verify_request(
&is_rostered,
guard,
&method,
&request.path,
&request.body,
&header,
now,
)
}
fn voucher_required() -> bool {
std::env::var("ZAKURO_REQUIRE_VOUCHER")
.map(|v| v == "1" || v == "true")
.unwrap_or(false)
}
fn redeem_offer_voucher(state: &BrokerState, offer: &task_board::TaskOffer, actual_cost: f64) {
if !voucher_required() || offer.voucher_signed_json.is_empty() {
return;
}
let nonce = match serde_json::from_str::<serde_json::Value>(&offer.voucher_signed_json) {
Ok(v) => v
.get("task_nonce")
.and_then(|n| n.as_str())
.map(|s| s.to_string()),
Err(_) => None,
};
if let (Some(nonce), Some(api_url), Some(api_key)) =
(nonce, &state.config.api_url, &state.config.api_key)
{
match super::node_sync::redeem_voucher(api_url, api_key, &nonce, actual_cost) {
Ok(true) => {}
Ok(false) => {
eprintln!(" [VOUCHER] nonce already spent (double-spend attempt?): {nonce}")
}
Err(e) => eprintln!(" [VOUCHER] redeem failed: {e}"),
}
}
}
fn verify_offer_voucher(
state: &Arc<BrokerState>,
offer: &task_board::TaskOffer,
) -> Result<(), String> {
if !voucher_required() {
return Ok(());
}
let pubkey = state
.dash_voucher_pubkey
.read()
.ok()
.and_then(|g| g.clone())
.ok_or("no dashboard voucher pubkey cached yet")?;
if offer.voucher_signed_json.is_empty() {
return Err("voucher required but none supplied".into());
}
let price_estimate =
offer.max_price_per_hour * (offer.estimated_duration_secs / 3600.0).max(0.0);
super::voucher::verify_voucher_now(
&pubkey,
&offer.voucher_signed_json,
&offer.voucher_sig,
price_estimate,
)
.map(|_| ())
}
fn node_sig_gate(state: &Arc<BrokerState>, request: &BrokerRequest) -> Option<BrokerResponse> {
let enforce = std::env::var("ZAKURO_ENFORCE_NODE_SIG")
.map(|v| v == "1" || v == "true")
.unwrap_or(false);
match check_node_sig(
&state.node_roster,
&state.node_sig_guard,
request,
node_identity::now_secs(),
) {
Ok(_) => None,
Err(e) => {
if enforce {
Some(error_response(
&format!("node signature check failed: {e}"),
"UNAUTHORIZED",
401,
))
} else {
if request.header("X-Node-Id").is_some() {
eprintln!(
" [NODE_SIG] permissive: {} {} — {}",
format!("{:?}", request.method).to_uppercase(),
request.path,
e
);
}
None
}
}
}
}
fn peer_handshake_auth(
state: &Arc<BrokerState>,
request: &BrokerRequest,
) -> Option<BrokerResponse> {
let key = state.peer_manager.peer_key();
if key.is_empty() {
if request.is_loopback() {
None
} else {
Some(error_response(
"Peer authentication required (set ZAKURO_PEER_KEY)",
"UNAUTHORIZED",
401,
))
}
} else if peer::verify_peer_key(request.header("X-Peer-Key").as_deref(), key) {
None
} else {
Some(error_response("Invalid peer key", "UNAUTHORIZED", 401))
}
}
#[derive(serde::Serialize)]
struct BrokerEntry {
id: String,
reachable: bool,
}
#[derive(serde::Serialize)]
struct BrokersResponse {
#[serde(rename = "self")]
self_id: String,
brokers: Vec<BrokerEntry>,
}
fn handle_brokers(state: &Arc<BrokerState>) -> BrokerResponse {
state.peer_manager.ensure_fingerprints_cached();
let mut peer_urls = state.peer_manager.peer_urls();
peer_urls.sort();
let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
let brokers: Vec<BrokerEntry> = peer_urls
.into_iter()
.filter_map(|url| {
let fp_id = state.peer_manager.peer_fingerprint(&url)?;
let reachable = state
.peer_manager
.get_client(&url)
.map(|c| c.is_reachable())
.unwrap_or(false);
Some((fp_id, reachable))
})
.fold(Vec::<BrokerEntry>::new(), |mut acc, (fp_id, reachable)| {
if seen.insert(fp_id.clone()) {
acc.push(BrokerEntry {
id: fp_id,
reachable,
});
} else if reachable {
if let Some(e) = acc.iter_mut().find(|e| e.id == fp_id) {
e.reachable = true;
}
}
acc
});
json_response(
&BrokersResponse {
self_id: state.node_key.node_uri(),
brokers,
},
200,
)
}
#[derive(Debug, Deserialize)]
struct SetPriceRequest {
price_per_hour: f64,
}
fn handle_get_price(state: &Arc<BrokerState>) -> BrokerResponse {
json_response(
&serde_json::json!({ "price_per_hour": state.price.get() }),
200,
)
}
fn handle_set_price(state: &Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if !request.is_loopback() {
return error_response("Setting price is local-only (loopback)", "FORBIDDEN", 403);
}
match body_or_413(request) {
Err(resp) => resp,
Ok(body) => match serde_json::from_slice::<SetPriceRequest>(body) {
Ok(req) => {
state.price.set(req.price_per_hour);
json_response(
&serde_json::json!({ "price_per_hour": state.price.get() }),
200,
)
}
Err(e) => error_response(&format!("Invalid request body: {e}"), "BAD_REQUEST", 400),
},
}
}
fn handle_drain_worker(
state: &Arc<BrokerState>,
request: &BrokerRequest,
worker_id: &str,
) -> BrokerResponse {
if !request.is_loopback() {
return error_response("Draining is local-only (loopback)", "FORBIDDEN", 403);
}
match state.workers.drain(worker_id) {
Some(w) => json_response(
&serde_json::json!({
"status": "draining",
"worker_id": w.id,
"active_requests": w.active_requests,
}),
200,
),
None => error_response("Worker not found", "NOT_FOUND", 404),
}
}
#[derive(Debug, Deserialize)]
struct InferRequest {
model: String,
#[serde(flatten)]
rest: serde_json::Value,
}
#[derive(Debug, Deserialize)]
struct OpenAiChatResponse {
choices: Vec<OpenAiChoice>,
#[serde(default)]
usage: OpenAiUsage,
}
#[derive(Debug, Deserialize)]
struct OpenAiChoice {
message: OpenAiMessage,
}
#[derive(Debug, Deserialize)]
struct OpenAiMessage {
#[serde(default)]
content: String,
}
#[derive(Debug, Default, Deserialize)]
struct OpenAiUsage {
#[serde(default)]
prompt_tokens: u64,
#[serde(default)]
completion_tokens: u64,
#[serde(default)]
total_tokens: u64,
}
#[derive(Debug)]
enum ForwardError {
Transport(String),
BadResponse(String),
}
impl std::fmt::Display for ForwardError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Transport(_) => write!(f, "provider unreachable"),
Self::BadResponse(_) => write!(f, "provider returned an invalid response"),
}
}
}
fn forward_chat_live(
worker_uri: &str,
body: &serde_json::Value,
) -> Result<OpenAiChatResponse, ForwardError> {
let endpoint = format!("{}/v1/chat/completions", worker_uri.trim_end_matches('/'));
let payload = serde_json::to_string(body)
.map_err(|e| ForwardError::BadResponse(format!("failed to serialize request: {e}")))?;
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_connect(Some(std::time::Duration::from_secs(10)))
.timeout_recv_response(Some(std::time::Duration::from_secs(120)))
.build(),
);
let response = agent
.post(&endpoint)
.config()
.http_status_as_error(false)
.build()
.header("Content-Type", "application/json")
.send(payload.as_str())
.map_err(|e| ForwardError::Transport(e.to_string()))?;
let status = response.status().as_u16();
let text = response
.into_body()
.read_to_string()
.map_err(|e| ForwardError::Transport(e.to_string()))?;
if status != 200 {
return Err(ForwardError::BadResponse(format!(
"provider returned status {status}"
)));
}
serde_json::from_str::<OpenAiChatResponse>(&text)
.map_err(|e| ForwardError::BadResponse(format!("failed to parse provider response: {e}")))
}
fn infer_authenticate(
state: &Arc<BrokerState>,
request: &BrokerRequest,
) -> Result<Option<String>, BrokerResponse> {
let bearer = request
.header("Authorization")
.and_then(|a| a.strip_prefix("Bearer ").map(|t| t.trim().to_string()))
.filter(|t| !t.is_empty());
if let Some(token) = bearer {
state
.ledger
.resolve_user_from_api_key(&token)
.map(Some)
.map_err(|e| error_response(&e.to_string(), "UNAUTHORIZED", 401))
} else if state.is_local_mode() {
Ok(None)
} else {
Err(error_response(
"API key required. Set Authorization: Bearer <key>",
"UNAUTHORIZED",
401,
))
}
}
const INFER_DEFAULT_MAX_TOKENS: u64 = 4096;
fn infer_cost_ceiling(price_per_mtok: f64, body_bytes: usize, max_tokens: Option<u64>) -> f64 {
let prompt_est = (body_bytes as f64 / 3.0).ceil();
let gen_budget = max_tokens.unwrap_or(INFER_DEFAULT_MAX_TOKENS) as f64;
(prompt_est + gen_budget) * price_per_mtok / 1_000_000.0
}
#[allow(clippy::too_many_arguments)]
fn settle_own_provider_payout(
state: &Arc<BrokerState>,
requester_uid: &str,
worker: &Worker,
actual_cost: f64,
commission_c: f64,
voucher_nonce: &str,
request_id: &str,
duration_ms: f64,
) {
let (payout, commission) = split_commission(actual_cost, commission_c);
if commission > 0.0 && state.is_billing_enabled() {
let platform_uid =
std::env::var("ZAKURO_PLATFORM_UID").unwrap_or_else(|_| "platform-house".into());
state
.ledger
.local_credits
.entry(platform_uid.clone())
.and_modify(|b| *b += commission)
.or_insert(commission);
state
.ledger
.authoritative_balances
.entry(platform_uid)
.and_modify(|b| *b += commission)
.or_insert(commission);
}
match state.config.owner_user_id.as_deref() {
Some(owner) if payout > 0.0 => {
let earn_request_id = format!("earn-{request_id}");
let dashboard_configured =
state.config.api_url.is_some() && state.config.api_key.is_some();
if dashboard_configured {
let _ = state.wal.append(&WalEntry {
request_id: earn_request_id.clone(),
user_id: owner.to_string(),
reservation_id: String::new(),
estimated_cost: payout,
actual_cost: Some(payout),
worker_id: worker.id.clone(),
duration_ms: Some(duration_ms),
timestamp: chrono::Utc::now(),
status: WalStatus::Earned,
});
}
let balance_after = state.ledger.local_add_credits(owner, payout);
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: earn_request_id,
user_id: owner.to_string(),
tx_type: "credit".to_string(),
amount: payout,
balance_after,
worker_id: worker.id.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(worker.name.clone()),
worker_uri: Some(worker.uri.clone()),
price_per_hour: worker.pricing.price_per_hour,
});
}
Some(_) => {}
None => {
eprintln!(
" [INFER] no owner configured on this broker: refunding provider payout \
{payout:.6} for request {request_id} to the requester"
);
if payout > 0.0 && state.is_billing_enabled() {
state
.ledger
.local_credits
.entry(requester_uid.to_string())
.and_modify(|b| *b += payout)
.or_insert(payout);
state
.ledger
.authoritative_balances
.entry(requester_uid.to_string())
.and_modify(|b| *b += payout)
.or_insert(payout);
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: format!("{request_id}-earn-refund"),
user_id: requester_uid.to_string(),
tx_type: "refund".to_string(),
amount: payout,
balance_after: 0.0,
worker_id: worker.id.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(worker.name.clone()),
worker_uri: Some(worker.uri.clone()),
price_per_hour: worker.pricing.price_per_hour,
});
}
}
}
if !voucher_nonce.is_empty() {
if let (Some(api_url), Some(api_key)) = (&state.config.api_url, &state.config.api_key) {
match super::node_sync::redeem_voucher(api_url, api_key, voucher_nonce, actual_cost) {
Ok(true) => {}
Ok(false) => {
eprintln!(" [VOUCHER] infer nonce already spent: {voucher_nonce}")
}
Err(e) => eprintln!(" [VOUCHER] infer redeem failed: {e}"),
}
}
}
}
struct InferBilling {
res_id: String,
ceiling: f64,
request_id: String,
uid: String,
commission_c: f64,
voucher_nonce: String,
}
fn infer_actual_cost(price_per_mtok: f64, total_tokens: u64) -> f64 {
total_tokens as f64 * price_per_mtok / 1_000_000.0
}
fn handle_infer(state: &Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
handle_infer_with(state, request, forward_chat_live)
}
fn handle_infer_with(
state: &Arc<BrokerState>,
request: &BrokerRequest,
forward: impl FnOnce(&str, &serde_json::Value) -> Result<OpenAiChatResponse, ForwardError>,
) -> BrokerResponse {
let user_id = match infer_authenticate(state, request) {
Ok(uid) => uid,
Err(resp) => return resp,
};
let body = match body_or_413(request) {
Ok(b) => b,
Err(resp) => return resp,
};
let body_len = body.len();
let parsed: InferRequest = match serde_json::from_slice(body) {
Ok(v) => v,
Err(_) => return error_response("invalid model uri", "BAD_REQUEST", 400),
};
let model_uuid = match parse_model_uri(&parsed.model) {
Some(uuid) => uuid,
None => return error_response("invalid model uri", "BAD_REQUEST", 400),
};
let worker = match state.workers.resolve_model(&model_uuid) {
Ok(w) => w,
Err(ModelResolveError::NoProvider { model_uuid }) => {
return error_response(
&format!("no provider serving zc://{model_uuid}"),
"NOT_FOUND",
404,
);
}
};
let mut forward_body = parsed.rest;
if let serde_json::Value::Object(ref mut map) = forward_body {
let model_field = match worker.provider_type {
crate::broker::worker::ProviderType::General => model_uuid.clone(),
crate::broker::worker::ProviderType::Specialized => worker.name.clone(),
};
map.insert("model".to_string(), serde_json::Value::String(model_field));
}
let max_tokens = forward_body.get("max_tokens").and_then(|v| v.as_u64());
let billable = state.is_billing_enabled() && worker.price_per_mtok > 0.0;
let reservation = match (&user_id, billable) {
(Some(uid), true) => {
let ceiling = infer_cost_ceiling(worker.price_per_mtok, body_len, max_tokens);
let (commission_c, voucher_nonce) = if voucher_required() {
match mint_and_verify_remote_voucher(state, ceiling, None) {
Ok(v) => v,
Err(resp) => return resp,
}
} else {
(0.0, String::new())
};
let request_id = uuid::Uuid::new_v4().to_string();
match state.ledger.local_reserve(uid, ceiling, &request_id) {
Ok((res_id, _balance)) => Some(InferBilling {
res_id,
ceiling,
request_id,
uid: uid.clone(),
commission_c,
voucher_nonce,
}),
Err(ledger::LedgerError::InsufficientCredits {
required,
available,
}) => {
return error_response(
&format!("Insufficient credits: need {required:.6}, have {available:.6}"),
"INSUFFICIENT_CREDITS",
402,
);
}
Err(e) => return error_response(&e.to_string(), "LEDGER_ERROR", 503),
}
}
_ => None,
};
let started = Instant::now();
match forward(&worker.uri, &forward_body) {
Ok(resp) => {
let content = resp
.choices
.first()
.map(|c| c.message.content.clone())
.unwrap_or_default();
let charged = match &reservation {
Some(bill) => {
let InferBilling {
res_id,
ceiling,
request_id,
uid,
commission_c,
voucher_nonce,
} = bill;
let actual = capped_charge(
infer_actual_cost(worker.price_per_mtok, resp.usage.total_tokens),
*ceiling,
);
let balance_after = match state.ledger.local_commit(res_id, actual) {
Ok(b) => b,
Err(e) => {
eprintln!(" [INFER] commit error for {request_id}: {e}");
state.ledger.load_balance_if_needed(uid)
}
};
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: request_id.clone(),
user_id: uid.clone(),
tx_type: "infer".to_string(),
amount: actual,
balance_after,
worker_id: worker.id.clone(),
duration_ms: started.elapsed().as_millis() as f64,
source_node: state.config.node_name.clone(),
worker_name: Some(worker.name.clone()),
worker_uri: Some(worker.uri.clone()),
price_per_hour: worker.pricing.price_per_hour,
});
let own_fp = state.node_key.fingerprint();
let exec_fp = executor_fp(&worker);
if exec_fp == Some(own_fp.as_str()) {
settle_own_provider_payout(
state,
uid,
&worker,
actual,
*commission_c,
voucher_nonce,
request_id,
started.elapsed().as_millis() as f64,
);
actual
} else {
match exec_fp.and_then(|fp| state.peer_manager.get_url_for_fingerprint(fp))
{
Some(peer_url) => {
settle_executor_payout(
state,
uid,
actual,
*commission_c,
&peer_url,
&worker.id,
request_id,
started.elapsed().as_millis() as f64,
);
if !voucher_nonce.is_empty() {
if let (Some(api_url), Some(api_key)) =
(&state.config.api_url, &state.config.api_key)
{
match super::node_sync::redeem_voucher(
api_url,
api_key,
voucher_nonce,
actual,
) {
Ok(true) => {}
Ok(false) => eprintln!(
" [VOUCHER] infer nonce already spent: {voucher_nonce}"
),
Err(e) => {
eprintln!(" [VOUCHER] infer redeem failed: {e}")
}
}
}
}
actual
}
None => {
eprintln!(
" [INFER] provider's broker unresolved for worker {} (request {request_id}): refunding {actual:.6} in full",
worker.name
);
state
.ledger
.local_credits
.entry(uid.to_string())
.and_modify(|b| *b += actual)
.or_insert(actual);
state
.ledger
.authoritative_balances
.entry(uid.to_string())
.and_modify(|b| *b += actual)
.or_insert(actual);
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: format!("{request_id}-unresolved-refund"),
user_id: uid.to_string(),
tx_type: "refund".to_string(),
amount: actual,
balance_after: 0.0,
worker_id: worker.id.clone(),
duration_ms: started.elapsed().as_millis() as f64,
source_node: state.config.node_name.clone(),
worker_name: Some(worker.name.clone()),
worker_uri: Some(worker.uri.clone()),
price_per_hour: worker.pricing.price_per_hour,
});
0.0
}
}
}
}
None => 0.0,
};
json_response(
&serde_json::json!({
"content": content,
"usage": {
"prompt_tokens": resp.usage.prompt_tokens,
"completion_tokens": resp.usage.completion_tokens,
"total_tokens": resp.usage.total_tokens,
},
"provider": worker.name,
"charged_zkcr": charged,
}),
200,
)
}
Err(e) => {
if let Some(bill) = &reservation {
if let Err(cancel_err) = state.ledger.local_cancel(&bill.res_id) {
eprintln!(
" [INFER] cancel error for {}: {cancel_err}",
bill.request_id
);
}
}
error_response(&format!("provider request failed: {e}"), "BAD_GATEWAY", 502)
}
}
}
#[cfg(test)]
mod infer_tests {
use super::*;
use crate::broker::worker::{
HardwareInfo, ProviderType, WorkerPricing, WorkerRegistration, WorkerResources,
};
use crate::broker::BrokerState;
use std::sync::Arc;
use tiny_http::Method;
const UUID: &str = "ae5f3db4-437a-40d1-93ec-c8258315d69a";
fn state_with_provider(uuid: &str, name: &str) -> Arc<BrokerState> {
let state = BrokerState::new();
state.workers.register(WorkerRegistration {
name: name.to_string(),
uri: "http://10.0.0.5:9000".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: ProviderType::Specialized,
served_models: vec![uuid.to_string()],
price_per_mtok: 1.0,
});
state.set_local_mode(true);
Arc::new(state)
}
fn infer_req(body: &str) -> BrokerRequest {
infer_req_with_headers(body, vec![])
}
fn infer_req_with_headers(body: &str, headers: Vec<(String, String)>) -> BrokerRequest {
BrokerRequest {
method: Method::Post,
path: "/infer".to_string(),
query: String::new(),
headers,
body: body.as_bytes().to_vec(),
peer_addr: None,
}
}
fn canned_ok(
_worker_uri: &str,
_body: &serde_json::Value,
) -> Result<OpenAiChatResponse, ForwardError> {
Ok(OpenAiChatResponse {
choices: vec![OpenAiChoice {
message: OpenAiMessage {
content: "hello there".to_string(),
},
}],
usage: OpenAiUsage {
prompt_tokens: 10,
completion_tokens: 5,
total_tokens: 15,
},
})
}
fn canned_err(
_worker_uri: &str,
_body: &serde_json::Value,
) -> Result<OpenAiChatResponse, ForwardError> {
Err(ForwardError::Transport(
"connection refused to 10.0.0.5:9000 secret-token".to_string(),
))
}
#[test]
fn served_uuid_forwards_and_maps_response() {
let state = state_with_provider(UUID, "gpu-worker-1");
let req = infer_req(&format!(
r#"{{"model":"zc://{UUID}","messages":[{{"role":"user","content":"hi"}}]}}"#
));
let resp = handle_infer_with(&state, &req, canned_ok);
assert_eq!(resp.status, 200);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert_eq!(v["content"], "hello there");
assert_eq!(v["usage"]["prompt_tokens"], 10);
assert_eq!(v["usage"]["completion_tokens"], 5);
assert_eq!(v["usage"]["total_tokens"], 15);
assert_eq!(v["provider"], "gpu-worker-1");
assert_eq!(v["charged_zkcr"], 0.0);
}
#[test]
fn general_provider_receives_the_model_uuid_not_the_worker_name() {
let state = BrokerState::new();
state.set_local_mode(true);
state.workers.register(WorkerRegistration {
name: "general-1".to_string(),
uri: "http://10.0.0.6:8600".to_string(),
worker_type: "llm".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: ProviderType::General,
served_models: vec![crate::broker::worker::MODEL_WILDCARD.to_string()],
price_per_mtok: 0.0,
});
let state = Arc::new(state);
let req = infer_req(&format!(r#"{{"model":"zc://{UUID}","messages":[]}}"#));
let seen = std::sync::Mutex::new(String::new());
let resp = handle_infer_with(&state, &req, |_, body| {
*seen.lock().unwrap() = body["model"].as_str().unwrap_or("").to_string();
canned_ok("", body)
});
assert_eq!(resp.status, 200);
assert_eq!(
*seen.lock().unwrap(),
UUID,
"a general provider must be told WHICH model to load"
);
}
#[test]
fn unserved_uuid_returns_404() {
let state = BrokerState::new();
state.set_local_mode(true);
let state = Arc::new(state);
let req = infer_req(&format!(r#"{{"model":"zc://{UUID}","messages":[]}}"#));
let resp = handle_infer_with(&state, &req, canned_ok);
assert_eq!(resp.status, 404);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert_eq!(v["error"], format!("no provider serving zc://{UUID}"));
}
#[test]
fn invalid_model_uri_returns_400() {
let state = BrokerState::new();
state.set_local_mode(true);
let state = Arc::new(state);
let req = infer_req(r#"{"model":"zc://node-abc","messages":[]}"#);
let resp = handle_infer_with(&state, &req, canned_ok);
assert_eq!(resp.status, 400);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert_eq!(v["error"], "invalid model uri");
}
#[test]
fn provider_forward_error_returns_502_without_leaking_internals() {
let state = state_with_provider(UUID, "gpu-worker-1");
let req = infer_req(&format!(
r#"{{"model":"{UUID}","messages":[{{"role":"user","content":"hi"}}]}}"#
));
let resp = handle_infer_with(&state, &req, canned_err);
assert_eq!(resp.status, 502);
let body = String::from_utf8(resp.body.clone()).unwrap();
assert!(
!body.contains("10.0.0.5") && !body.contains("secret-token"),
"502 body must not leak transport internals: {body}"
);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert!(v["error"]
.as_str()
.unwrap()
.starts_with("provider request failed"));
}
fn billing_state_with_priced_provider(
uuid: &str,
price_per_mtok: f64,
uid: &str,
balance: f64,
) -> Arc<BrokerState> {
let mut state = BrokerState::new();
state.config.api_url = Some("http://dashboard.test".to_string());
state.config.api_key = Some("svc-key".to_string());
state.config.owner_user_id = Some("provider-owner".to_string());
let w = state.workers.register(WorkerRegistration {
name: "priced-worker".to_string(),
uri: "http://10.0.0.5:9000".to_string(),
worker_type: "llm".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: ProviderType::Specialized,
served_models: vec![uuid.to_string()],
price_per_mtok,
});
let fp = state.node_key.fingerprint();
state.workers.set_node_fp(&w.id, &fp);
state.set_local_mode(true); state
.ledger
.authoritative_balances
.insert(uid.to_string(), balance);
Arc::new(state)
}
fn bearer(uid: &str) -> Vec<(String, String)> {
vec![(
"Authorization".to_string(),
format!("Bearer zk_{uid}_deadbeef"),
)]
}
#[test]
fn cost_ceiling_and_actual_cost_math() {
assert!((infer_cost_ceiling(1000.0, 300, Some(200)) - 0.3).abs() < 1e-9);
let c = infer_cost_ceiling(1000.0, 300, None);
assert!((c - (100.0 + INFER_DEFAULT_MAX_TOKENS as f64) * 1000.0 / 1e6).abs() < 1e-9);
assert!((infer_actual_cost(1000.0, 15) - 0.015).abs() < 1e-12);
assert_eq!(infer_cost_ceiling(0.0, 300, Some(200)), 0.0);
assert_eq!(infer_actual_cost(0.0, 15), 0.0);
}
#[test]
fn billed_infer_charges_actual_tokens_and_reports_it() {
let state = billing_state_with_priced_provider(UUID, 1000.0, "u1", 10.0);
let body = format!(
r#"{{"model":"zc://{UUID}","messages":[{{"role":"user","content":"hi"}}],"max_tokens":100}}"#
);
let req = infer_req_with_headers(&body, bearer("u1"));
let resp = handle_infer_with(&state, &req, canned_ok);
assert_eq!(resp.status, 200);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert!((v["charged_zkcr"].as_f64().unwrap() - 0.015).abs() < 1e-9);
let bal = *state.ledger.authoritative_balances.get("u1").unwrap();
assert!((bal - (10.0 - 0.015)).abs() < 1e-9, "balance {bal}");
let owner = *state
.ledger
.authoritative_balances
.get("provider-owner")
.expect("owner credited");
assert!((owner - 0.015).abs() < 1e-9, "owner payout {owner}");
}
#[test]
fn no_owner_refunds_payout_to_requester() {
let state = billing_state_with_priced_provider(UUID, 1000.0, "u7", 10.0);
let mut raw = BrokerState::new();
raw.config.api_url = Some("http://dashboard.test".to_string());
raw.config.api_key = Some("svc-key".to_string());
raw.config.owner_user_id = None;
let w = raw.workers.register(WorkerRegistration {
name: "priced-worker".to_string(),
uri: "http://10.0.0.5:9000".to_string(),
worker_type: "llm".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: ProviderType::Specialized,
served_models: vec![UUID.to_string()],
price_per_mtok: 1000.0,
});
let fp = raw.node_key.fingerprint();
raw.workers.set_node_fp(&w.id, &fp);
raw.set_local_mode(true);
raw.ledger
.authoritative_balances
.insert("u7".to_string(), 10.0);
let state2 = Arc::new(raw);
drop(state);
let body = format!(
r#"{{"model":"zc://{UUID}","messages":[{{"role":"user","content":"hi"}}],"max_tokens":100}}"#
);
let req = infer_req_with_headers(&body, bearer("u7"));
let resp = handle_infer_with(&state2, &req, canned_ok);
assert_eq!(resp.status, 200);
let bal = *state2.ledger.authoritative_balances.get("u7").unwrap();
assert!(
(bal - 10.0).abs() < 1e-9,
"requester must net 0, balance {bal}"
);
}
#[test]
fn unstamped_provider_settlement_refunds_in_full() {
let state = billing_state_with_priced_provider(UUID, 1000.0, "u6", 10.0);
let wid = state.workers.list()[0].id.clone();
state.workers.set_node_fp(&wid, "");
let body = format!(
r#"{{"model":"zc://{UUID}","messages":[{{"role":"user","content":"hi"}}],"max_tokens":100}}"#
);
let req = infer_req_with_headers(&body, bearer("u6"));
let resp = handle_infer_with(&state, &req, canned_ok);
assert_eq!(resp.status, 200);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert_eq!(
v["charged_zkcr"], 0.0,
"refunded charge must be reported as 0"
);
let bal = *state.ledger.authoritative_balances.get("u6").unwrap();
assert!(
(bal - 10.0).abs() < 1e-9,
"requester made whole, balance {bal}"
);
assert!(
state
.ledger
.authoritative_balances
.get("provider-owner")
.map(|b| *b == 0.0)
.unwrap_or(true),
"nobody may be paid when the charge was refunded"
);
}
#[test]
fn insufficient_credits_rejects_before_forwarding() {
let state = billing_state_with_priced_provider(UUID, 1000.0, "u2", 0.0001);
let body = format!(
r#"{{"model":"zc://{UUID}","messages":[{{"role":"user","content":"hi"}}],"max_tokens":100}}"#
);
let req = infer_req_with_headers(&body, bearer("u2"));
let forwarded = std::sync::atomic::AtomicBool::new(false);
let resp = handle_infer_with(&state, &req, |_, _| {
forwarded.store(true, std::sync::atomic::Ordering::SeqCst);
canned_ok("", &serde_json::Value::Null)
});
assert_eq!(resp.status, 402);
assert!(
!forwarded.load(std::sync::atomic::Ordering::SeqCst),
"provider compute must not be spent for a caller who cannot pay"
);
let bal = *state.ledger.authoritative_balances.get("u2").unwrap();
assert!((bal - 0.0001).abs() < 1e-12);
}
#[test]
fn failed_forward_refunds_the_full_reservation() {
let state = billing_state_with_priced_provider(UUID, 1000.0, "u3", 5.0);
let body = format!(
r#"{{"model":"zc://{UUID}","messages":[{{"role":"user","content":"hi"}}],"max_tokens":100}}"#
);
let req = infer_req_with_headers(&body, bearer("u3"));
let resp = handle_infer_with(&state, &req, canned_err);
assert_eq!(resp.status, 502);
let bal = *state.ledger.authoritative_balances.get("u3").unwrap();
assert!(
(bal - 5.0).abs() < 1e-9,
"failed call must cost nothing, balance {bal}"
);
}
#[test]
fn unpriced_provider_stays_free_for_authed_caller() {
let state = billing_state_with_priced_provider(UUID, 0.0, "u4", 5.0);
let body = format!(r#"{{"model":"zc://{UUID}","messages":[]}}"#);
let req = infer_req_with_headers(&body, bearer("u4"));
let resp = handle_infer_with(&state, &req, canned_ok);
assert_eq!(resp.status, 200);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert_eq!(v["charged_zkcr"], 0.0);
let bal = *state.ledger.authoritative_balances.get("u4").unwrap();
assert!((bal - 5.0).abs() < 1e-12);
}
#[test]
fn keyless_local_caller_stays_free_even_on_priced_provider() {
let state = billing_state_with_priced_provider(UUID, 1000.0, "unused", 5.0);
let body = format!(r#"{{"model":"zc://{UUID}","messages":[]}}"#);
let req = infer_req(&body);
let resp = handle_infer_with(&state, &req, canned_ok);
assert_eq!(resp.status, 200);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert_eq!(v["charged_zkcr"], 0.0);
}
#[test]
fn over_reporting_provider_cannot_charge_past_the_ceiling() {
let state = billing_state_with_priced_provider(UUID, 1000.0, "u5", 10.0);
let body = format!(
r#"{{"model":"zc://{UUID}","messages":[{{"role":"user","content":"hi"}}],"max_tokens":10}}"#
);
let req = infer_req_with_headers(&body, bearer("u5"));
let resp = handle_infer_with(&state, &req, |_, _| {
Ok(OpenAiChatResponse {
choices: vec![OpenAiChoice {
message: OpenAiMessage {
content: "x".to_string(),
},
}],
usage: OpenAiUsage {
prompt_tokens: 1_000_000,
completion_tokens: 1_000_000,
total_tokens: 2_000_000,
},
})
});
assert_eq!(resp.status, 200);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
let charged = v["charged_zkcr"].as_f64().unwrap();
let ceiling = infer_cost_ceiling(1000.0, req.body.len(), Some(10));
assert!(
charged <= ceiling + 1e-12,
"charged {charged} must not exceed reserved ceiling {ceiling}"
);
}
fn canned_panics(
_worker_uri: &str,
_body: &serde_json::Value,
) -> Result<OpenAiChatResponse, ForwardError> {
panic!("forward must not run when the auth gate rejects the request");
}
#[test]
fn remote_mode_without_bearer_returns_401_before_forward() {
let state = state_with_provider(UUID, "gpu-worker-1");
state.set_local_mode(false); let req = infer_req(&format!(
r#"{{"model":"zc://{UUID}","messages":[{{"role":"user","content":"hi"}}]}}"#
));
let resp = handle_infer_with(&state, &req, canned_panics);
assert_eq!(resp.status, 401);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert_eq!(
v["error"],
"API key required. Set Authorization: Bearer <key>"
);
}
#[test]
fn remote_mode_with_valid_bearer_proceeds() {
let state = state_with_provider(UUID, "gpu-worker-1");
state.set_local_mode(false); let req = infer_req_with_headers(
&format!(r#"{{"model":"zc://{UUID}","messages":[{{"role":"user","content":"hi"}}]}}"#),
vec![(
"Authorization".to_string(),
"Bearer zk_user123_abcdef".to_string(),
)],
);
let resp = handle_infer_with(&state, &req, canned_ok);
assert_eq!(resp.status, 200);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert_eq!(v["content"], "hello there");
assert_eq!(v["charged_zkcr"], 0.0);
}
#[test]
fn local_mode_without_bearer_is_allowed() {
let state = state_with_provider(UUID, "gpu-worker-1"); let req = infer_req(&format!(
r#"{{"model":"zc://{UUID}","messages":[{{"role":"user","content":"hi"}}]}}"#
));
let resp = handle_infer_with(&state, &req, canned_ok);
assert_eq!(resp.status, 200);
let v: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert_eq!(v["content"], "hello there");
}
}
fn handle_peer_health(state: &Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Some(resp) = peer_handshake_auth(state, request) {
return resp;
}
let body = serde_json::to_vec(&peer::PeerHealthResponse {
status: "healthy".to_string(),
broker_id: state.node_key.node_uri(),
})
.unwrap_or_default();
let sig = state.node_key.sign(&body);
BrokerResponse::json_bytes(body, 200)
.with_header("X-Node-Id", &state.node_key.public_b64())
.with_header("X-Node-Sig", &sig)
}
fn handle_peer_workers(state: &Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Some(resp) = peer_credit_auth(state, request) {
return resp;
}
let worker_infos: Vec<WorkerInfo> = state
.workers
.list()
.iter()
.filter(|w| w.source_node.is_none())
.map(|w| worker_to_info(w, &state.workers, &state.node_key, false))
.collect();
json_response(
&WorkerListResponse {
total: worker_infos.len(),
workers: worker_infos,
},
200,
)
}
fn handle_peer_peers(state: &Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Err(e) = check_node_sig(
&state.node_roster,
&state.node_sig_guard,
request,
node_identity::now_secs(),
) {
return error_response(
&format!("node signature check failed: {e}"),
"UNAUTHORIZED",
401,
);
}
let entries: Vec<serde_json::Value> = state
.peer_manager
.gossip_entries()
.into_iter()
.map(|p| serde_json::json!({"fp": p.fp, "url": p.url, "epoch": p.last_seen}))
.collect();
let body = serde_json::to_vec(&serde_json::json!({ "peers": entries })).unwrap_or_default();
let sig = state.node_key.sign(&body);
BrokerResponse::json_bytes(body, 200)
.with_header("X-Node-Id", &state.node_key.public_b64())
.with_header("X-Node-Sig", &sig)
}
fn handle_peer_advert(state: &Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Err(e) = check_node_sig(
&state.node_roster,
&state.node_sig_guard,
request,
node_identity::now_secs(),
) {
return error_response(
&format!("node signature check failed: {e}"),
"UNAUTHORIZED",
401,
);
}
let advert = crate::broker::discovery::Advert {
fp: state.node_key.fingerprint(),
price_per_hour: state.price.get(),
resources: state.workers.aggregate_available(),
epoch: node_identity::now_secs(),
};
let body = serde_json::to_vec(&advert).unwrap_or_default();
let sig = state.node_key.sign(&body);
BrokerResponse::json_bytes(body, 200)
.with_header("X-Node-Id", &state.node_key.public_b64())
.with_header("X-Node-Sig", &sig)
}
fn handle_peer_reserve(state: Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Some(resp) = peer_credit_auth(&state, request) {
return resp;
}
let body = match body_or_413(request) {
Ok(b) => b,
Err(resp) => return resp,
};
let req: peer::PeerReserveRequest = match serde_json::from_slice(body) {
Ok(r) => r,
Err(e) => return error_response(&format!("Invalid request: {}", e), "BAD_REQUEST", 400),
};
match state
.ledger
.local_reserve(&req.user_id, req.amount, &req.request_id)
{
Ok((reservation_id, balance_before)) => json_response(
&peer::PeerReserveResponse {
reservation_id,
balance_before,
},
200,
),
Err(ledger::LedgerError::InsufficientCredits {
required,
available,
}) => json_response(
&peer::PeerErrorResponse {
error: format!(
"Insufficient credits: need {:.4}, have {:.4}",
required, available
),
code: "INSUFFICIENT_CREDITS".to_string(),
},
402,
),
Err(e @ ledger::LedgerError::BalanceUnavailable { .. }) => json_response(
&peer::PeerErrorResponse {
error: e.to_string(),
code: "BALANCE_UNAVAILABLE".to_string(),
},
503,
),
Err(e) => error_response(&e.to_string(), "LEDGER_ERROR", 500),
}
}
fn handle_peer_commit(state: Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Some(resp) = peer_credit_auth(&state, request) {
return resp;
}
let body = match body_or_413(request) {
Ok(b) => b,
Err(resp) => return resp,
};
let req: peer::PeerCommitRequest = match serde_json::from_slice(body) {
Ok(r) => r,
Err(e) => return error_response(&format!("Invalid request: {}", e), "BAD_REQUEST", 400),
};
match state
.ledger
.local_commit(&req.reservation_id, req.actual_cost)
{
Ok(balance_after) => json_response(&peer::PeerCommitResponse { balance_after }, 200),
Err(e) => error_response(&e.to_string(), "LEDGER_ERROR", 500),
}
}
fn handle_peer_cancel(state: Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Some(resp) = peer_credit_auth(&state, request) {
return resp;
}
let body = match body_or_413(request) {
Ok(b) => b,
Err(resp) => return resp,
};
let req: peer::PeerCancelRequest = match serde_json::from_slice(body) {
Ok(r) => r,
Err(e) => return error_response(&format!("Invalid request: {}", e), "BAD_REQUEST", 400),
};
if let Err(e) = state.ledger.local_cancel(&req.reservation_id) {
eprintln!(
" [BILLING] local_cancel failed for {}: {}",
req.reservation_id, e
);
}
json_response(&serde_json::json!({"status": "ok"}), 200)
}
fn handle_peer_balance(state: &Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Some(resp) = peer_credit_auth(state, request) {
return resp;
}
let user_id = request
.query
.split('&')
.find(|p| p.starts_with("user_id="))
.map(|p| p[8..].to_string())
.unwrap_or_default();
if user_id.is_empty() {
return error_response("Missing user_id parameter", "BAD_REQUEST", 400);
}
let balance = state.ledger.load_balance_if_needed(&user_id);
json_response(&peer::PeerBalanceResponse { user_id, balance }, 200)
}
fn handle_peer_earn(state: Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Some(resp) = peer_credit_auth(&state, request) {
return resp;
}
let body = match body_or_413(request) {
Ok(b) => b,
Err(resp) => return resp,
};
let req: peer::PeerEarnRequest = match serde_json::from_slice(body) {
Ok(r) => r,
Err(e) => return error_response(&format!("Invalid request: {}", e), "BAD_REQUEST", 400),
};
let owner_user_id = match state.config.owner_user_id.as_deref() {
Some(uid) => uid.to_string(),
None => {
return error_response("No owner configured on this broker", "NOT_CONFIGURED", 503);
}
};
let earn_request_id = format!("earn-{}", req.request_id);
let dashboard_configured = state.config.api_url.is_some() && state.config.api_key.is_some();
if dashboard_configured {
let _ = state.wal.append(&WalEntry {
request_id: earn_request_id.clone(),
user_id: owner_user_id.clone(),
reservation_id: String::new(),
estimated_cost: req.amount,
actual_cost: Some(req.amount),
worker_id: req.worker_id.clone(),
duration_ms: Some(req.duration_ms),
timestamp: chrono::Utc::now(),
status: WalStatus::Earned,
});
}
let balance_after = state.ledger.local_add_credits(&owner_user_id, req.amount);
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: earn_request_id,
user_id: owner_user_id.clone(),
tx_type: "credit".to_string(),
amount: req.amount,
balance_after,
worker_id: req.worker_id.clone(),
duration_ms: req.duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(req.worker_id.clone()),
worker_uri: None,
price_per_hour: 0.0,
});
state
.tx_buffer
.snapshot_balance(&owner_user_id, balance_after);
eprintln!(
" [EARN] Worker {} earned {:.6} credits from user {} → owner {} (balance now {:.6})",
req.worker_id, req.amount, req.requesting_user, owner_user_id, balance_after
);
json_response(
&serde_json::json!({"status": "ok", "balance_after": balance_after}),
200,
)
}
fn handle_peer_identity(state: &Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Some(resp) = peer_handshake_auth(state, request) {
return resp;
}
let workers: Vec<task_board::WorkerSummary> = state
.workers
.list()
.iter()
.filter(|w| state.is_local_worker(&w.uri))
.map(|w| task_board::WorkerSummary {
name: w.name.clone(),
price_per_hour: w.pricing.price_per_hour,
status: format!("{:?}", w.status),
cpus: w.resources.cpus_total,
memory_bytes: w.resources.memory_total,
gpus: w.resources.gpus_total,
})
.collect();
let identity = task_board::PeerIdentity {
owner_user_id: state.config.owner_user_id.clone().unwrap_or_default(),
node_name: state.config.node_name.clone().unwrap_or_default(),
verified: state.is_billing_enabled(),
workers,
quic_port: state.quic_port.load(Ordering::Relaxed),
supports_two_phase: true,
};
json_response(&identity, 200)
}
fn select_local_worker_for_offer(
state: &BrokerState,
offer: &task_board::TaskOffer,
) -> Option<Worker> {
let local_workers: Vec<Worker> = state
.workers
.healthy()
.into_iter()
.filter(|w| state.is_local_worker(&w.uri))
.filter(|w| w.pricing.price_per_hour <= offer.max_price_per_hour)
.filter(|w| {
offer
.worker_type
.as_ref()
.is_none_or(|wt| w.worker_type.as_str() == wt.as_str())
})
.collect();
if local_workers.is_empty() {
return None;
}
let idx = (REQUEST_COUNTER.fetch_add(1, Ordering::SeqCst) as usize) % local_workers.len();
Some(local_workers[idx].clone())
}
pub fn reserve_offer(
state: &BrokerState,
offer: &task_board::TaskOffer,
) -> Result<task_board::TaskAccept, task_board::TaskReject> {
let worker =
select_local_worker_for_offer(state, offer).ok_or_else(|| task_board::TaskReject {
task_id: offer.task_id.clone(),
reason: "No matching local worker".to_string(),
})?;
if !state.workers.try_reserve_quota(&worker.id) {
return Err(task_board::TaskReject {
task_id: offer.task_id.clone(),
reason: "Worker quota exhausted".to_string(),
});
}
let estimated_cost = worker.pricing.estimate_cost(offer.estimated_duration_secs);
state.pending_offers.insert(
offer.task_id.clone(),
super::PendingOffer {
offer: offer.clone(),
worker_id: worker.id.clone(),
created: Instant::now(),
},
);
Ok(task_board::TaskAccept {
task_id: offer.task_id.clone(),
estimated_cost,
worker_name: worker.name.clone(),
executor_owner: state.config.owner_user_id.clone().unwrap_or_default(),
price_per_hour: worker.pricing.price_per_hour,
})
}
pub fn pending_reservations(state: &BrokerState) -> std::collections::HashMap<String, u32> {
let mut counts = std::collections::HashMap::new();
for entry in state.pending_offers.iter() {
*counts.entry(entry.value().worker_id.clone()).or_insert(0) += 1;
}
counts
}
pub fn cancel_offer(state: &BrokerState, task_id: &str) {
if let Some((_, pending)) = state.pending_offers.remove(task_id) {
state.workers.cancel_quota_reservation(&pending.worker_id);
}
}
pub fn take_committed(state: &BrokerState, task_id: &str) -> Option<super::PendingOffer> {
state.pending_offers.remove(task_id).map(|(_, p)| p)
}
pub fn sweep_expired_offers(state: &BrokerState, ttl: std::time::Duration) {
let stale: Vec<String> = state
.pending_offers
.iter()
.filter(|e| e.value().created.elapsed() > ttl)
.map(|e| e.key().clone())
.collect();
for task_id in stale {
cancel_offer(state, &task_id);
}
}
pub fn execute_offer_locally(
state: &BrokerState,
offer: &task_board::TaskOffer,
) -> Result<task_board::TaskResult, task_board::TaskReject> {
let worker =
select_local_worker_for_offer(state, offer).ok_or_else(|| task_board::TaskReject {
task_id: offer.task_id.clone(),
reason: "No matching local worker".to_string(),
})?;
run_offer_on(state, offer, &worker)
}
fn run_offer_on(
state: &BrokerState,
offer: &task_board::TaskOffer,
worker: &Worker,
) -> Result<task_board::TaskResult, task_board::TaskReject> {
let _active_guard = state.workers.active_guard(&worker.id);
let payload = base64::Engine::decode(
&base64::engine::general_purpose::STANDARD,
&offer.payload_b64,
)
.map_err(|e| task_board::TaskReject {
task_id: offer.task_id.clone(),
reason: format!("Invalid payload encoding: {}", e),
})?;
let worker_uri = format!("{}/execute", worker.uri.trim_end_matches('/'));
let start = Instant::now();
let timeout = if offer.timeout_secs > 0.0 {
offer.timeout_secs
} else {
300.0
};
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_global(Some(std::time::Duration::from_secs_f64(timeout + 5.0)))
.build(),
);
let forward_result = agent
.post(&worker_uri)
.header("Content-Type", "application/octet-stream")
.header("X-Zakuro-Request-Id", &offer.task_id)
.header("X-Zakuro-Timeout-Secs", &format!("{:.1}", timeout))
.send(&payload);
let duration_ms = start.elapsed().as_secs_f64() * 1000.0;
match forward_result {
Ok(response) => {
let actual_cost = worker.pricing.estimate_cost(duration_ms / 1000.0);
let worker_pid = response
.headers()
.get("X-Zakuro-Pid")
.and_then(|v| v.to_str().ok())
.map(|s| s.to_string());
let mut response_body = Vec::new();
let _ = response
.into_body()
.into_reader()
.read_to_end(&mut response_body);
let result_b64 =
base64::Engine::encode(&base64::engine::general_purpose::STANDARD, &response_body);
let executor_owner = state.config.owner_user_id.clone().unwrap_or_default();
redeem_offer_voucher(state, offer, actual_cost);
let mut result = task_board::TaskResult {
task_id: offer.task_id.clone(),
payload_b64: result_b64,
duration_ms,
actual_cost,
worker_name: worker.name.clone(),
worker_uri: worker.zc_uri(),
price_per_hour: worker.pricing.price_per_hour,
executor_owner,
worker_pid,
executor_node_uri: state.node_key.node_uri(),
executor_node_id: None,
executor_sig: None,
};
result.executor_node_id = Some(state.node_key.public_b64());
result.executor_sig = Some(state.node_key.sign(&task_board::receipt_bytes(&result)));
Ok(result)
}
Err(e) => Err(task_board::TaskReject {
task_id: offer.task_id.clone(),
reason: format!("Worker execution failed: {}", e),
}),
}
}
fn handle_peer_task_offer(state: Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Some(resp) = peer_credit_auth(&state, request) {
return resp;
}
let body = match body_or_413(request) {
Ok(b) => b,
Err(resp) => return resp,
};
let offer: task_board::TaskOffer = match serde_json::from_slice(body) {
Ok(o) => o,
Err(e) => return error_response(&format!("Invalid task offer: {}", e), "BAD_REQUEST", 400),
};
if let Err(reason) = verify_offer_voucher(&state, &offer) {
return error_response(
&format!("voucher check failed: {reason}"),
"VOUCHER_INVALID",
402,
);
}
if offer.two_phase {
match reserve_offer(&state, &offer) {
Ok(accept) => json_response(&accept, 200),
Err(reject) => json_response(&reject, 404),
}
} else {
match execute_offer_locally(&state, &offer) {
Ok(result) => json_response(&result, 200),
Err(reject) => json_response(&reject, 404),
}
}
}
fn handle_peer_task_commit(state: Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Some(resp) = peer_credit_auth(&state, request) {
return resp;
}
let body = match body_or_413(request) {
Ok(b) => b,
Err(resp) => return resp,
};
let commit: task_board::TaskCommit = match serde_json::from_slice(body) {
Ok(c) => c,
Err(e) => return error_response(&format!("Invalid commit: {}", e), "BAD_REQUEST", 400),
};
match take_committed(&state, &commit.task_id) {
Some(pending) => match state.workers.get(&pending.worker_id) {
Some(worker) => match run_offer_on(&state, &pending.offer, &worker) {
Ok(result) => json_response(&result, 200),
Err(reject) => json_response(&reject, 404),
},
None => json_response(
&task_board::TaskReject {
task_id: commit.task_id.clone(),
reason: "Reserved worker no longer available".to_string(),
},
404,
),
},
None => json_response(
&task_board::TaskReject {
task_id: commit.task_id.clone(),
reason: "Unknown or expired reservation".to_string(),
},
404,
),
}
}
fn handle_peer_task_cancel(state: Arc<BrokerState>, request: &BrokerRequest) -> BrokerResponse {
if let Some(resp) = peer_credit_auth(&state, request) {
return resp;
}
let body = match body_or_413(request) {
Ok(b) => b,
Err(resp) => return resp,
};
let cancel: task_board::TaskCancel = match serde_json::from_slice(body) {
Ok(c) => c,
Err(e) => return error_response(&format!("Invalid cancel: {}", e), "BAD_REQUEST", 400),
};
cancel_offer(&state, &cancel.task_id);
json_response(
&serde_json::json!({ "status": "cancelled", "task_id": cancel.task_id }),
200,
)
}
fn forward_to_worker(
agent: &ureq::Agent,
worker_uri: &str,
body: &[u8],
request_id: &str,
effective_timeout: f64,
) -> Result<ureq::http::Response<ureq::Body>, ureq::Error> {
if effective_timeout > 0.0 {
agent
.post(worker_uri)
.config()
.timeout_global(Some(std::time::Duration::from_secs_f64(
effective_timeout + 5.0,
)))
.build()
.header("Content-Type", "application/octet-stream")
.header("X-Zakuro-Request-Id", request_id)
.header(
"X-Zakuro-Timeout-Secs",
&format!("{:.1}", effective_timeout),
)
.send(body)
} else {
agent
.post(worker_uri)
.header("Content-Type", "application/octet-stream")
.header("X-Zakuro-Request-Id", request_id)
.send(body)
}
}
struct WorkerReply {
status: u16,
headers: Vec<(String, String)>,
body: Vec<u8>,
}
async fn forward_to_worker_async(
client: &reqwest::Client,
worker_uri: &str,
body: Vec<u8>,
request_id: &str,
effective_timeout: f64,
) -> Result<WorkerReply, String> {
let mut req = client
.post(worker_uri)
.header("Content-Type", "application/octet-stream")
.header("X-Zakuro-Request-Id", request_id)
.body(body);
if effective_timeout > 0.0 {
req = req
.header("X-Zakuro-Timeout-Secs", format!("{:.1}", effective_timeout))
.timeout(std::time::Duration::from_secs_f64(effective_timeout + 5.0));
}
let resp = req.send().await.map_err(|e| e.to_string())?;
let status = resp.status().as_u16();
let headers = resp
.headers()
.iter()
.map(|(k, v)| {
(
k.as_str().to_lowercase(),
v.to_str().unwrap_or("").to_string(),
)
})
.collect();
let body = resp.bytes().await.map_err(|e| e.to_string())?.to_vec();
Ok(WorkerReply {
status,
headers,
body,
})
}
fn into_worker_reply(resp: ureq::http::Response<ureq::Body>) -> WorkerReply {
use std::io::Read as _;
let status = resp.status().as_u16();
let headers = resp
.headers()
.iter()
.map(|(k, v)| {
(
k.as_str().to_lowercase(),
v.to_str().unwrap_or("").to_string(),
)
})
.collect();
let mut body = Vec::new();
let _ = resp.into_body().into_reader().read_to_end(&mut body);
WorkerReply {
status,
headers,
body,
}
}
fn reply_header<'a>(reply: &'a WorkerReply, name: &str) -> Option<&'a str> {
let needle = name.to_lowercase();
reply
.headers
.iter()
.find(|(k, _)| *k == needle)
.map(|(_, v)| v.as_str())
}
fn execute_finish(
state: &Arc<BrokerState>,
verbose: bool,
p: PreparedForward,
forward: Result<WorkerReply, String>,
) -> BrokerResponse {
let worker_name = p.worker.name.clone();
let duration_ms = p.start.elapsed().as_secs_f64() * 1000.0;
match forward {
Ok(reply) => {
let actual_cost = if p.is_local {
0.0
} else {
let actual_duration_secs = duration_ms / 1000.0;
p.worker.pricing.estimate_cost(actual_duration_secs)
};
if !p.is_local {
let _ = state.wal.update_status(
&p.request_id,
WalStatus::Executed,
Some(actual_cost),
Some(duration_ms),
);
}
let credits_remaining = commit_credits(
state,
&p.authority,
p.is_local,
&p.reservation_id,
actual_cost,
p.balance_before,
&p.request_id,
&p.user_id,
&p.worker,
duration_ms,
p.commission_c,
&p.voucher_task_nonce,
);
if !p.is_local {
let _ = state.wal.update_status(
&p.request_id,
WalStatus::Committed,
Some(actual_cost),
Some(duration_ms),
);
}
state
.workers
.record_request(&p.worker.id, duration_ms, true);
state.active_requests.remove(&p.request_id);
let worker_pid = reply_header(&reply, "x-zakuro-pid").map(|s| s.to_string());
let worker_ip = reply_header(&reply, "x-zakuro-ip")
.map(|s| s.to_string())
.or_else(|| {
let u = p
.worker
.uri
.strip_prefix("http://")
.unwrap_or(&p.worker.uri);
u.split(':').next().map(|s| s.to_string())
});
let response_body = reply.body;
if p.instance_action.as_deref() == Some("create_instance") {
if let Some(ref iid) = p.instance_id {
state
.instance_registry
.insert(iid.clone(), p.worker.id.clone());
}
}
let action_label = match p.instance_action.as_deref() {
Some("create_instance") => "CREATE_INST",
Some("call_method") => "CALL_METHOD",
_ => "EXECUTE",
};
log_transaction_with_worker_info(
state,
verbose,
p.tx_num,
&p.user_id,
action_label,
actual_cost,
credits_remaining,
Some(&worker_name),
duration_ms,
"OK",
Some(p.worker.pricing.price_per_hour),
state.config.owner_user_id.as_deref(),
worker_pid.as_deref(),
worker_ip.as_deref(),
Some(&p.request_id),
);
let headers = [
header_or("Content-Type", "application/octet-stream"),
header_or("X-Zakuro-Request-Id", &p.request_id),
header_or("X-Zakuro-Worker", worker_name.as_str()),
header_or("X-Zakuro-Cost", &format!("{:.6}", actual_cost)),
header_or(
"X-Zakuro-Credits-Remaining",
&format!("{:.6}", credits_remaining),
),
header_or("X-Zakuro-Duration-Ms", &format!("{:.2}", duration_ms)),
header_or("X-Zakuro-Transport", "local"),
];
let mut resp = BrokerResponse {
status: 200,
headers: Vec::new(),
body: response_body,
};
resp.headers.extend(headers);
resp
}
Err(e) => {
let is_timeout = e.contains("timed out") || e.contains("Timeout");
let timeout_cost = if is_timeout && !p.is_local {
let elapsed_secs = duration_ms / 1000.0;
capped_charge(
p.worker.pricing.estimate_cost(elapsed_secs),
p.reservation_amount,
)
} else {
0.0
};
if is_timeout {
if !p.is_local && timeout_cost > 0.0 {
match &p.authority {
Authority::Local => {
if let Err(e) =
state.ledger.local_commit(&p.reservation_id, timeout_cost)
{
eprintln!(
" [BILLING] local_commit failed for {}: {}",
p.request_id, e
);
}
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: p.request_id.clone(),
user_id: p.user_id.clone(),
tx_type: "commit".to_string(),
amount: timeout_cost,
balance_after: p.balance_before - timeout_cost,
worker_id: p.worker.id.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(p.worker.name.clone()),
worker_uri: Some(p.worker.uri.clone()),
price_per_hour: p.worker.pricing.price_per_hour,
});
}
Authority::Peer(url) => {
if let Some(client) = state.peer_manager.get_client(url) {
let _ = client.commit(&p.reservation_id, timeout_cost);
} else {
let _ = state.ledger.commit(&p.reservation_id, timeout_cost);
}
state.ledger.publish_transaction(
&p.request_id,
&p.user_id,
"commit",
timeout_cost,
p.balance_before - timeout_cost,
&p.worker.id,
duration_ms,
state.config.node_name.as_deref(),
);
}
Authority::Standalone => {
let _ = state.ledger.commit(&p.reservation_id, timeout_cost);
state.ledger.publish_transaction(
&p.request_id,
&p.user_id,
"commit",
timeout_cost,
p.balance_before - timeout_cost,
&p.worker.id,
duration_ms,
state.config.node_name.as_deref(),
);
}
}
let _ = state.wal.update_status(
&p.request_id,
WalStatus::Committed,
Some(timeout_cost),
Some(duration_ms),
);
}
state.active_requests.remove(&p.request_id);
log_transaction(
state,
verbose,
p.tx_num,
&p.user_id,
"EXECUTE",
timeout_cost,
p.balance_before - timeout_cost,
Some(&worker_name),
duration_ms,
"FAIL",
);
state.workers.decrement_active(&p.worker.id);
return error_response(
&format!(
"Request timed out after {:.1}s (charged {:.6} credits)",
duration_ms / 1000.0,
timeout_cost
),
"TIMEOUT",
504,
);
}
state.workers.mark_unhealthy(&p.worker.id);
state.workers.cancel_quota_reservation(&p.worker.id);
if !p.is_local {
cancel_credits(
state,
&p.authority,
&p.reservation_id,
&p.request_id,
&p.user_id,
&p.worker.id,
p.balance_before,
duration_ms,
);
}
state.active_requests.remove(&p.request_id);
state.workers.decrement_active(&p.worker.id);
const MAX_RETRIES: usize = 2;
let mut tried_ids: Vec<String> = vec![p.worker.id.clone()];
for attempt in 1..=MAX_RETRIES {
let next_worker = {
let healthy = state.workers.healthy();
healthy.into_iter().find(|w| !tried_ids.contains(&w.id))
};
let next = match next_worker {
Some(w) => w,
None => {
log_transaction(
state,
verbose,
p.tx_num,
&p.user_id,
"EXECUTE",
0.0,
p.balance_before,
Some(&worker_name),
duration_ms,
"FAIL",
);
return error_response(
&format!(
"Worker unreachable and no fallback available (tried {} worker(s))",
tried_ids.len()
),
"NO_WORKERS",
503,
);
}
};
let retry_uri = format!("{}/execute", next.uri.trim_end_matches('/'));
let retry_name = next.name.clone();
tried_ids.push(next.id.clone());
if !state.workers.try_reserve_quota(&next.id) {
continue;
}
eprintln!(
" [RETRY] attempt {}/{} → {} ({})",
attempt, MAX_RETRIES, retry_name, retry_uri
);
let mut retry_guard = state.workers.active_guard(&next.id);
state
.active_requests
.insert(p.request_id.clone(), next.id.clone());
let retry_start = Instant::now();
let retry_result = forward_to_worker(
&state.http_client,
&retry_uri,
&p.body,
&p.request_id,
p.effective_timeout,
)
.map(into_worker_reply)
.map_err(|e| e.to_string());
let retry_duration_ms = retry_start.elapsed().as_secs_f64() * 1000.0;
match retry_result {
Ok(retry_reply) => {
let actual_cost = if p.is_local {
0.0
} else {
next.pricing.estimate_cost(retry_duration_ms / 1000.0)
};
let retry_reservation_id = if !p.is_local {
let res_id = format!("{}-retry{}", p.request_id, attempt);
match &p.authority {
Authority::Local => {
let _ = state.ledger.local_reserve(
&p.user_id,
actual_cost,
&res_id,
);
let _ = state.ledger.local_commit(&res_id, actual_cost);
}
Authority::Peer(url) => {
if let Some(client) = state.peer_manager.get_client(url) {
if client.reserve(&p.user_id, actual_cost, &res_id).is_ok()
{
let _ = client.commit(&res_id, actual_cost);
}
}
}
Authority::Standalone => {
if state
.ledger
.reserve(&p.user_id, actual_cost, &res_id)
.is_ok()
{
let _ = state.ledger.commit(&res_id, actual_cost);
}
}
}
res_id
} else {
String::new()
};
if !p.is_local {
let _ = state.wal.update_status(
&p.request_id,
WalStatus::Executed,
Some(actual_cost),
Some(retry_duration_ms),
);
let _ = state.wal.update_status(
&p.request_id,
WalStatus::Committed,
Some(actual_cost),
Some(retry_duration_ms),
);
}
state
.workers
.record_request(&next.id, retry_duration_ms, true);
retry_guard.disarm();
state.active_requests.remove(&p.request_id);
let credits_remaining = if p.is_local {
p.balance_before
} else {
p.balance_before - actual_cost
};
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: format!("{}-retry{}", p.request_id, attempt),
user_id: p.user_id.clone(),
tx_type: "commit".to_string(),
amount: actual_cost,
balance_after: credits_remaining,
worker_id: next.id.clone(),
duration_ms: retry_duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(next.name.clone()),
worker_uri: Some(next.uri.clone()),
price_per_hour: next.pricing.price_per_hour,
});
let action_label = match p.instance_action.as_deref() {
Some("create_instance") => "CREATE_INST",
Some("call_method") => "CALL_METHOD",
_ => "EXECUTE",
};
log_transaction(
state,
verbose,
p.tx_num,
&p.user_id,
action_label,
actual_cost,
credits_remaining,
Some(&retry_name),
retry_duration_ms,
"OK",
);
let response_body = retry_reply.body;
let _ = retry_reservation_id; let headers = [
header_or("Content-Type", "application/octet-stream"),
header_or("X-Zakuro-Request-Id", &p.request_id),
header_or("X-Zakuro-Cost", &format!("{:.6}", actual_cost)),
header_or(
"X-Zakuro-Credits-Remaining",
&format!("{:.6}", credits_remaining),
),
header_or("X-Zakuro-Duration-Ms", &format!("{:.2}", retry_duration_ms)),
header_or("X-Zakuro-Retries", &format!("{}", attempt)),
];
let mut resp = BrokerResponse {
status: 200,
headers: Vec::new(),
body: response_body,
};
resp.headers.extend(headers);
return resp;
}
Err(_retry_err) => {
state.workers.mark_unhealthy(&next.id);
state.workers.cancel_quota_reservation(&next.id);
state.active_requests.remove(&p.request_id);
}
}
}
log_transaction(
state,
verbose,
p.tx_num,
&p.user_id,
"EXECUTE",
0.0,
p.balance_before,
Some(&worker_name),
duration_ms,
"FAIL",
);
error_response(
&format!(
"All workers unreachable after {} retries (tried: {})",
MAX_RETRIES,
tried_ids.join(", ")
),
"NO_WORKERS",
503,
)
}
}
}
#[allow(clippy::too_many_arguments)]
fn reserve_credits(
state: &Arc<BrokerState>,
authority: &peer::Authority,
user_id: &str,
amount: f64,
request_id: &str,
verbose: bool,
tx_num: u64,
balance_before: f64,
) -> Result<String, BrokerResponse> {
let fail = |msg: &str, code: &str, status: u16| -> Result<String, _> {
log_transaction(
state,
verbose,
tx_num,
user_id,
"EXECUTE",
0.0,
balance_before,
None,
0.0,
"FAIL",
);
Err(error_response(msg, code, status))
};
match authority {
peer::Authority::Local => match state.ledger.local_reserve(user_id, amount, request_id) {
Ok((id, _)) => Ok(id),
Err(ledger::LedgerError::InsufficientCredits {
required,
available,
}) => fail(
&format!(
"Insufficient credits: need {:.4}, have {:.4}",
required, available
),
"INSUFFICIENT_CREDITS",
402,
),
Err(e @ ledger::LedgerError::BalanceUnavailable { .. }) => {
fail(&e.to_string(), "BALANCE_UNAVAILABLE", 503)
}
Err(e) => fail(&e.to_string(), "LEDGER_ERROR", 503),
},
peer::Authority::Peer(url) => {
let peer_result = state
.peer_manager
.get_client(url)
.and_then(|c| c.reserve(user_id, amount, request_id).ok())
.map(|r| r.reservation_id);
if let Some(id) = peer_result {
return Ok(id);
}
eprintln!(
" [P2P] Peer reserve failed ({}), falling back to local",
url
);
reserve_via_ledger(
state,
user_id,
amount,
request_id,
verbose,
tx_num,
balance_before,
)
}
peer::Authority::Standalone => reserve_via_ledger(
state,
user_id,
amount,
request_id,
verbose,
tx_num,
balance_before,
),
}
}
fn reserve_via_ledger(
state: &Arc<BrokerState>,
user_id: &str,
amount: f64,
request_id: &str,
verbose: bool,
tx_num: u64,
balance_before: f64,
) -> Result<String, BrokerResponse> {
let fail = |msg: &str, code: &str, status: u16| -> Result<String, _> {
log_transaction(
state,
verbose,
tx_num,
user_id,
"EXECUTE",
0.0,
balance_before,
None,
0.0,
"FAIL",
);
Err(error_response(msg, code, status))
};
match state.ledger.reserve(user_id, amount, request_id) {
Ok(id) => Ok(id),
Err(ledger::LedgerError::InsufficientCredits {
required,
available,
}) => fail(
&format!(
"Insufficient credits: need {:.4}, have {:.4}",
required, available
),
"INSUFFICIENT_CREDITS",
402,
),
Err(e @ ledger::LedgerError::BalanceUnavailable { .. }) => {
fail(&e.to_string(), "BALANCE_UNAVAILABLE", 503)
}
Err(e) => fail(&e.to_string(), "LEDGER_ERROR", 503),
}
}
#[allow(clippy::too_many_arguments)]
fn cancel_credits(
state: &Arc<BrokerState>,
authority: &peer::Authority,
reservation_id: &str,
request_id: &str,
user_id: &str,
worker_id: &str,
balance_before: f64,
duration_ms: f64,
) {
match authority {
peer::Authority::Local => {
if let Err(e) = state.ledger.local_cancel(reservation_id) {
eprintln!(" [BILLING] local_cancel failed for {}: {}", request_id, e);
}
}
peer::Authority::Peer(url) => {
if let Some(client) = state.peer_manager.get_client(url) {
if let Err(e) = client.cancel(reservation_id) {
eprintln!(" [BILLING] peer cancel failed for {}: {}", request_id, e);
}
} else if let Err(e) = state.ledger.cancel(reservation_id) {
eprintln!(" [BILLING] cancel failed for {}: {}", request_id, e);
}
}
peer::Authority::Standalone => {
if let Err(e) = state.ledger.cancel(reservation_id) {
eprintln!(" [BILLING] cancel failed for {}: {}", request_id, e);
}
}
}
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: request_id.to_string(),
user_id: user_id.to_string(),
tx_type: "cancel".to_string(),
amount: 0.0,
balance_after: balance_before,
worker_id: worker_id.to_string(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: None,
worker_uri: None,
price_per_hour: 0.0,
});
let _ = state.wal.update_status(
request_id,
super::wal::WalStatus::Failed,
None,
Some(duration_ms),
);
}
#[allow(clippy::too_many_arguments)]
fn settle_executor_payout(
state: &Arc<BrokerState>,
user_id: &str,
cost: f64,
commission_c: f64,
executor_peer_url: &str,
executor_worker_name: &str,
request_id: &str,
duration_ms: f64,
) {
let (payout, commission) = split_commission(cost, commission_c);
if commission > 0.0 && state.is_billing_enabled() {
let platform_uid =
std::env::var("ZAKURO_PLATFORM_UID").unwrap_or_else(|_| "platform-house".into());
state
.ledger
.local_credits
.entry(platform_uid.clone())
.and_modify(|b| *b += commission)
.or_insert(commission);
state
.ledger
.authoritative_balances
.entry(platform_uid)
.and_modify(|b| *b += commission)
.or_insert(commission);
}
let settle_url = state
.peer_manager
.get_url_for_endpoint(executor_peer_url)
.unwrap_or_else(|| executor_peer_url.to_string());
let earn_ok = match state.peer_manager.get_client(&settle_url) {
Some(client) => match client.earn(
payout,
duration_ms,
executor_worker_name,
user_id,
request_id,
) {
Ok(_) => true,
Err(e) => {
eprintln!(
" [EARN] Failed to notify peer {}: {}",
executor_peer_url, e
);
false
}
},
None => {
eprintln!(
" [EARN] No client for peer {} (resolved to {}, request {}): executor unreachable, refunding payout to requester",
executor_peer_url, settle_url, request_id
);
false
}
};
let refund = earn_refund_amount(payout, earn_ok);
if refund > 0.0 && state.is_billing_enabled() {
state
.ledger
.local_credits
.entry(user_id.to_string())
.and_modify(|b| *b += refund)
.or_insert(refund);
state
.ledger
.authoritative_balances
.entry(user_id.to_string())
.and_modify(|b| *b += refund)
.or_insert(refund);
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: format!("{request_id}-earn-refund"),
user_id: user_id.to_string(),
tx_type: "refund".to_string(),
amount: refund,
balance_after: 0.0,
worker_id: executor_worker_name.to_string(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(executor_worker_name.to_string()),
worker_uri: None,
price_per_hour: 0.0,
});
}
}
fn executor_fp(worker: &Worker) -> Option<&str> {
if worker.node_fp.is_empty() {
None
} else {
Some(worker.node_fp.as_str())
}
}
fn executor_is_self(state: &BrokerState, worker: &Worker) -> bool {
executor_fp(worker).is_some_and(|fp| fp == state.node_key.fingerprint())
}
#[allow(clippy::too_many_arguments)]
fn commit_credits(
state: &Arc<BrokerState>,
authority: &peer::Authority,
is_local: bool,
reservation_id: &str,
actual_cost: f64,
balance_before: f64,
request_id: &str,
user_id: &str,
worker: &worker::Worker,
duration_ms: f64,
commission_c: f64,
voucher_task_nonce: &str,
) -> f64 {
if is_local {
state.ledger.publish_transaction(
request_id,
user_id,
"commit",
0.0,
balance_before,
&worker.id,
duration_ms,
state.config.node_name.as_deref(),
);
if state.is_billing_enabled() {
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: request_id.to_string(),
user_id: user_id.to_string(),
tx_type: "commit".to_string(),
amount: 0.0,
balance_after: balance_before,
worker_id: worker.id.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(worker.name.clone()),
worker_uri: Some(worker.uri.clone()),
price_per_hour: worker.pricing.price_per_hour,
});
}
return balance_before;
}
match authority {
peer::Authority::Local => {
let balance = match state.ledger.local_commit(reservation_id, actual_cost) {
Ok(b) => b,
Err(e) => {
eprintln!(" [P2P] Local commit error for {}: {}", request_id, e);
balance_before - actual_cost
}
};
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: request_id.to_string(),
user_id: user_id.to_string(),
tx_type: "commit".to_string(),
amount: actual_cost,
balance_after: balance,
worker_id: worker.id.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(worker.name.clone()),
worker_uri: Some(worker.uri.clone()),
price_per_hour: worker.pricing.price_per_hour,
});
state.tx_buffer.snapshot_balance(user_id, balance);
if executor_is_self(state, worker) {
settle_self_payout(
state,
user_id,
actual_cost,
commission_c,
worker,
request_id,
duration_ms,
);
redeem_task_nonce(state, voucher_task_nonce, actual_cost);
return balance;
}
let exec_fp = executor_fp(worker);
match exec_fp.and_then(|fp| state.peer_manager.get_url_for_fingerprint(fp)) {
Some(peer_url) => {
settle_executor_payout(
state,
user_id,
actual_cost,
commission_c,
&peer_url,
&worker.id,
request_id,
duration_ms,
);
if !voucher_task_nonce.is_empty() {
if let (Some(api_url), Some(api_key)) =
(&state.config.api_url, &state.config.api_key)
{
match super::node_sync::redeem_voucher(
api_url,
api_key,
voucher_task_nonce,
actual_cost,
) {
Ok(true) => {}
Ok(false) => eprintln!(
" [VOUCHER] nonce already spent (double-spend attempt?): {voucher_task_nonce}"
),
Err(e) => eprintln!(" [VOUCHER] redeem failed: {e}"),
}
}
}
}
None => {
eprintln!(
" [EARN] UNRESOLVED executor peer for node={} \
(request {request_id}): no executor payout/redeem — refunding \
requester {actual_cost:.6} in full to conserve",
exec_fp
.map(|f| format!("zc://node-{f}"))
.unwrap_or_else(|| "<unstamped>".into())
);
if state.is_billing_enabled() {
state
.ledger
.local_credits
.entry(user_id.to_string())
.and_modify(|b| *b += actual_cost)
.or_insert(actual_cost);
state
.ledger
.authoritative_balances
.entry(user_id.to_string())
.and_modify(|b| *b += actual_cost)
.or_insert(actual_cost);
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: format!("{request_id}-unresolved-refund"),
user_id: user_id.to_string(),
tx_type: "refund".to_string(),
amount: actual_cost,
balance_after: balance + actual_cost,
worker_id: worker.id.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(worker.name.clone()),
worker_uri: Some(worker.uri.clone()),
price_per_hour: worker.pricing.price_per_hour,
});
}
}
}
balance
}
peer::Authority::Peer(url) => {
let peer_balance = if let Some(client) = state.peer_manager.get_client(url) {
match client.commit(reservation_id, actual_cost) {
Ok(resp) => {
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: request_id.to_string(),
user_id: user_id.to_string(),
tx_type: "commit".to_string(),
amount: actual_cost,
balance_after: resp.balance_after,
worker_id: worker.id.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(worker.name.clone()),
worker_uri: Some(worker.uri.clone()),
price_per_hour: worker.pricing.price_per_hour,
});
Some(resp.balance_after)
}
Err(e) => {
eprintln!(
" [P2P] Peer commit failed ({}), falling back to local: {}",
url, e
);
None
}
}
} else {
None
};
match peer_balance {
Some(b) => b,
None => {
state.ledger.publish_transaction(
request_id,
user_id,
"commit",
actual_cost,
balance_before - actual_cost,
&worker.id,
duration_ms,
state.config.node_name.as_deref(),
);
match state.ledger.commit(reservation_id, actual_cost) {
Ok(b) => b,
Err(e) => {
eprintln!(" [LEDGER] Commit error for {}: {}", request_id, e);
balance_before - actual_cost
}
}
}
}
}
peer::Authority::Standalone => {
let balance = match state.ledger.commit(reservation_id, actual_cost) {
Ok(b) => b,
Err(e) => {
eprintln!(
" [LEDGER] Commit error for {}: {} (WAL has record)",
request_id, e
);
balance_before - actual_cost
}
};
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: request_id.to_string(),
user_id: user_id.to_string(),
tx_type: "commit".to_string(),
amount: actual_cost,
balance_after: balance,
worker_id: worker.id.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(worker.name.clone()),
worker_uri: Some(worker.uri.clone()),
price_per_hour: worker.pricing.price_per_hour,
});
if executor_is_self(state, worker) {
settle_self_payout(
state,
user_id,
actual_cost,
commission_c,
worker,
request_id,
duration_ms,
);
redeem_task_nonce(state, voucher_task_nonce, actual_cost);
}
balance
}
}
}
fn redeem_task_nonce(state: &Arc<BrokerState>, voucher_task_nonce: &str, actual_cost: f64) {
if voucher_task_nonce.is_empty() {
return;
}
if let (Some(api_url), Some(api_key)) = (&state.config.api_url, &state.config.api_key) {
match super::node_sync::redeem_voucher(api_url, api_key, voucher_task_nonce, actual_cost) {
Ok(true) => {}
Ok(false) => eprintln!(
" [VOUCHER] nonce already spent (double-spend attempt?): {voucher_task_nonce}"
),
Err(e) => eprintln!(" [VOUCHER] redeem failed: {e}"),
}
}
}
fn settle_self_payout(
state: &Arc<BrokerState>,
user_id: &str,
cost: f64,
commission_c: f64,
worker: &worker::Worker,
request_id: &str,
duration_ms: f64,
) {
if !state.is_billing_enabled() || cost <= 0.0 {
return;
}
let (payout, commission) = split_commission(cost, commission_c);
if commission > 0.0 {
let platform_uid =
std::env::var("ZAKURO_PLATFORM_UID").unwrap_or_else(|_| "platform-house".into());
state
.ledger
.local_credits
.entry(platform_uid.clone())
.and_modify(|b| *b += commission)
.or_insert(commission);
state
.ledger
.authoritative_balances
.entry(platform_uid)
.and_modify(|b| *b += commission)
.or_insert(commission);
}
let Some(owner) = state.config.owner_user_id.clone() else {
eprintln!(
" [EARN] self-executed job {request_id} but no owner configured on this broker; \
refunding payout {payout:.6} to requester {user_id}"
);
state
.ledger
.local_credits
.entry(user_id.to_string())
.and_modify(|b| *b += payout)
.or_insert(payout);
state
.ledger
.authoritative_balances
.entry(user_id.to_string())
.and_modify(|b| *b += payout)
.or_insert(payout);
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: format!("{request_id}-earn-refund"),
user_id: user_id.to_string(),
tx_type: "refund".to_string(),
amount: payout,
balance_after: 0.0,
worker_id: worker.id.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(worker.name.clone()),
worker_uri: None,
price_per_hour: worker.pricing.price_per_hour,
});
return;
};
let earn_request_id = format!("earn-{request_id}");
let _ = state.wal.append(&WalEntry {
request_id: earn_request_id.clone(),
user_id: owner.clone(),
reservation_id: String::new(),
estimated_cost: payout,
actual_cost: Some(payout),
worker_id: worker.id.clone(),
duration_ms: Some(duration_ms),
timestamp: chrono::Utc::now(),
status: WalStatus::Earned,
});
let balance_after = state.ledger.local_add_credits(&owner, payout);
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: earn_request_id,
user_id: owner.clone(),
tx_type: "credit".to_string(),
amount: payout,
balance_after,
worker_id: worker.id.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(worker.name.clone()),
worker_uri: None,
price_per_hour: worker.pricing.price_per_hour,
});
state.tx_buffer.snapshot_balance(&owner, balance_after);
println!(
" [EARN] Worker {} earned {:.6} credits from user {} → owner {} (self-executed; commission {:.6}; balance now {:.6})",
worker.name, payout, user_id, owner, commission, balance_after
);
}
fn peer_dispatch_deadline(timeout_secs: f64) -> std::time::Duration {
if timeout_secs > 0.0 {
std::time::Duration::from_secs_f64(timeout_secs + 12.0)
} else {
std::time::Duration::from_secs(30)
}
}
fn broadcast_targets(peer_urls: &[String], abandoned_preferred_peer: Option<&str>) -> Vec<String> {
match abandoned_preferred_peer {
Some(abandoned) => peer_urls
.iter()
.filter(|url| url.as_str() != abandoned)
.cloned()
.collect(),
None => peer_urls.to_vec(),
}
}
fn marketplace_select_enabled() -> bool {
std::env::var("ZAKURO_MARKETPLACE_SELECT")
.map(|v| v != "0")
.unwrap_or(true)
}
fn preselect_executor_peer(state: &Arc<BrokerState>, req: &ResourceRequirements) -> Option<String> {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0);
let peers: Vec<(String, Option<super::discovery::Advert>, f64)> = state
.peer_manager
.peer_urls()
.into_iter()
.map(|url| {
let advert = state.peer_manager.peer_advert(&url);
let latency_ms = state.peer_manager.peer_latency_ms(&url).unwrap_or(f64::MAX);
(url, advert, latency_ms)
})
.collect();
const MAX_ADVERT_AGE_SECS: u64 = 120;
super::selection::select_executor_broker(&peers, req, MAX_ADVERT_AGE_SECS, now).map(|b| b.url)
}
#[allow(clippy::too_many_arguments)]
fn split_commission(gross: f64, c: f64) -> (f64, f64) {
let commission = gross * c;
let payout = gross - commission;
(payout, commission)
}
fn earn_refund_amount(payout: f64, earn_ok: bool) -> f64 {
if earn_ok {
0.0
} else {
payout
}
}
fn offer_commission_c(offer: &task_board::TaskOffer, dash_pubkey: Option<&str>) -> f64 {
if offer.voucher_signed_json.is_empty() || offer.voucher_sig.is_empty() {
return 0.0;
}
let Some(pubkey) = dash_pubkey else {
return 0.0;
};
match super::voucher::verify_voucher_now(
pubkey,
&offer.voucher_signed_json,
&offer.voucher_sig,
0.0,
) {
Ok(v) if (0.0..1.0).contains(&v.commission_c) => v.commission_c,
_ => 0.0,
}
}
#[allow(clippy::too_many_arguments)]
fn dispatch_to_peers(
state: &Arc<BrokerState>,
offer: task_board::TaskOffer,
start: Instant,
user_id: &str,
_is_self_owned: bool,
request_id: &str,
tx_num: u64,
verbose: bool,
balance_before: f64,
preferred_peer: Option<String>,
verified_commission_c: f64,
) -> BrokerResponse {
let peer_urls: Vec<String> = state.peer_manager.peer_urls();
if state.is_billing_enabled() {
let unbounded = offer.max_price_per_hour >= f64::MAX;
let max_cost = if unbounded {
0.0
} else {
offer.max_price_per_hour / 3600.0 * offer.estimated_duration_secs
};
let insufficient = if unbounded {
balance_before <= 0.0
} else {
balance_before < max_cost
};
if insufficient {
log_transaction(
state,
verbose,
tx_num,
user_id,
"EXECUTE",
0.0,
balance_before,
None,
0.0,
"FAIL",
);
let detail = if unbounded {
format!(
"Insufficient credits: a positive balance is required, have {:.4}",
balance_before
)
} else {
format!(
"Insufficient credits: need {:.4}, have {:.4}",
max_cost, balance_before
)
};
return error_response(&detail, "INSUFFICIENT_CREDITS", 402);
}
}
let quic_transport = state.quic.get().cloned();
let mut winning_peer: Option<String> = None;
let mut winning_result: Option<task_board::TaskResult> = None;
let mut winning_transport: Option<String> = None;
if let Some(ref qt) = quic_transport {
if let Some((url, res)) = qt.broadcast_offer_to_subscribers(&offer) {
winning_peer = Some(url);
winning_result = Some(res);
winning_transport = Some("subscribed".to_string());
}
}
let mut abandoned_preferred_peer: Option<String> = None;
if winning_peer.is_none() {
if let Some(ref peer_url) = preferred_peer {
if state.peer_manager.get_client(peer_url).is_some() {
let (tx, rx) = std::sync::mpsc::channel::<Result<task_board::TaskResult, String>>();
let state_c = state.clone();
let offer_c = offer.clone();
let peer_url_c = peer_url.clone();
std::thread::spawn(move || {
let result = match state_c.peer_manager.get_client(&peer_url_c) {
Some(client) => client.offer_task(&offer_c),
None => Err("peer unreachable".to_string()),
};
let _ = tx.send(result);
});
let deadline = peer_dispatch_deadline(offer.timeout_secs);
if let Ok(Ok(result)) = rx.recv_timeout(deadline) {
winning_peer = Some(peer_url.clone());
winning_result = Some(result);
winning_transport = Some("http".to_string());
} else {
abandoned_preferred_peer = Some(peer_url.clone());
}
}
}
}
let broadcast_peer_urls = broadcast_targets(&peer_urls, abandoned_preferred_peer.as_deref());
if winning_peer.is_none() && !broadcast_peer_urls.is_empty() {
let (tx, rx) =
std::sync::mpsc::channel::<Option<(String, task_board::TaskResult, String)>>();
for peer_url in &broadcast_peer_urls {
let tx = tx.clone();
let state_c = state.clone();
let offer_c = offer.clone();
let quic_c = quic_transport.clone();
let peer_url = peer_url.clone();
std::thread::spawn(move || {
let peer_quic_port = state_c.peer_manager.get_quic_port(&peer_url);
let (offer_result, transport_used) =
if let (Some(ref qt), p) = (&quic_c, peer_quic_port) {
if p > 0 {
(quic_offer(qt, &peer_url, p, &offer_c), "quic".to_string())
} else if let Some(client) = state_c.peer_manager.get_client(&peer_url) {
(client.offer_task(&offer_c), "http".to_string())
} else {
(Err("peer unreachable".into()), String::new())
}
} else if let Some(client) = state_c.peer_manager.get_client(&peer_url) {
(client.offer_task(&offer_c), "http".to_string())
} else {
(Err("peer unreachable".into()), String::new())
};
let _ = tx.send(offer_result.ok().map(|r| (peer_url, r, transport_used)));
});
}
drop(tx);
let deadline = peer_dispatch_deadline(offer.timeout_secs);
let dispatch_start = Instant::now();
for _ in 0..broadcast_peer_urls.len() {
let remaining = deadline.saturating_sub(dispatch_start.elapsed());
match rx.recv_timeout(remaining) {
Ok(Some((url, res, trans))) => {
winning_peer = Some(url);
winning_result = Some(res);
winning_transport = Some(trans);
break;
}
Ok(None) => continue, Err(_) => break, }
}
}
if let (Some(peer_url), Some(result), Some(transport_used)) =
(winning_peer, winning_result, winning_transport)
{
let duration_ms = start.elapsed().as_secs_f64() * 1000.0;
let cost = result.actual_cost;
let balance_after = balance_before - cost;
log_transaction_with_worker_info(
state,
verbose,
tx_num,
user_id,
"EXECUTE",
cost,
balance_after,
Some(&result.worker_name),
duration_ms,
"OK",
Some(result.price_per_hour),
Some(&result.executor_owner),
result.worker_pid.as_deref(),
None,
Some(request_id),
);
if state.is_billing_enabled() && cost > 0.0 {
state.credits.set_balance(user_id, balance_after);
if balance_after < 0.0 {
eprintln!(
" [BILLING] OVERSPEND user={} cost={:.4} balance_before={:.4} balance_after={:.4} request_id={}",
user_id, cost, balance_before, balance_after, request_id
);
}
state
.ledger
.local_credits
.entry(user_id.to_string())
.and_modify(|b| *b -= cost);
state
.ledger
.authoritative_balances
.entry(user_id.to_string())
.and_modify(|b| *b -= cost);
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: request_id.to_string(),
user_id: user_id.to_string(),
tx_type: "commit".to_string(),
amount: cost,
balance_after,
worker_id: result.worker_name.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(result.worker_name.clone()),
worker_uri: Some(result.worker_uri.clone()),
price_per_hour: result.price_per_hour,
});
}
if cost > 0.0 {
settle_executor_payout(
state,
user_id,
cost,
verified_commission_c,
&peer_url,
&result.worker_name,
request_id,
duration_ms,
);
}
let response_body = base64::Engine::decode(
&base64::engine::general_purpose::STANDARD,
&result.payload_b64,
)
.unwrap_or_default();
let headers = [
header_or("Content-Type", "application/octet-stream"),
header_or("X-Zakuro-Request-Id", request_id),
header_or("X-Zakuro-Worker", result.worker_name.as_str()),
header_or("X-Zakuro-Cost", &format!("{:.6}", cost)),
header_or(
"X-Zakuro-Credits-Remaining",
&format!("{:.6}", balance_after),
),
header_or("X-Zakuro-Duration-Ms", &format!("{:.2}", duration_ms)),
header_or("X-Zakuro-Transport", transport_used.as_str()),
];
let mut resp = BrokerResponse {
status: 200,
headers: Vec::new(),
body: response_body,
};
resp.headers.extend(headers);
return resp;
}
log_transaction(
state,
verbose,
tx_num,
user_id,
"EXECUTE",
0.0,
balance_before,
None,
0.0,
"FAIL",
);
error_response("No worker available (local or peer)", "NO_WORKERS", 503)
}
struct PreparedForward {
worker_uri: String,
body: Vec<u8>,
request_id: String,
tx_num: u64,
effective_timeout: f64,
user_id: String,
authority: Authority,
reservation_id: String,
balance_before: f64,
worker: Worker,
is_local: bool,
reservation_amount: f64,
instance_action: Option<String>,
instance_id: Option<String>,
start: Instant,
commission_c: f64,
voucher_task_nonce: String,
}
fn mint_and_verify_remote_voucher(
state: &Arc<BrokerState>,
budget: f64,
mint_key: Option<&str>,
) -> Result<(f64, String), BrokerResponse> {
let (_json, _sig, v) = mint_and_verify_voucher(state, budget, mint_key)?;
Ok((v.commission_c, v.task_nonce))
}
fn mint_and_verify_voucher(
state: &Arc<BrokerState>,
budget: f64,
mint_key: Option<&str>,
) -> Result<(String, String, super::voucher::Voucher), BrokerResponse> {
let (api_url, own_key) = match (&state.config.api_url, &state.config.api_key) {
(Some(u), Some(k)) => (u, k),
_ => {
return Err(error_response(
"Voucher required for remote dispatch but broker is not configured for \
dashboard vouchers (missing api_url/api_key)",
"VOUCHER_UNAVAILABLE",
503,
));
}
};
let api_key = mint_key.filter(|k| !k.trim().is_empty()).unwrap_or(own_key);
let (signed_json, sig) =
super::node_sync::obtain_voucher(api_url, api_key, budget.max(0.000_001)).map_err(|e| {
eprintln!(" [VOUCHER] issue failed: {e}");
error_response(
&format!("Voucher required for remote dispatch but mint failed: {e}"),
"VOUCHER_MINT_FAILED",
503,
)
})?;
let dash_pubkey = state
.dash_voucher_pubkey
.read()
.ok()
.and_then(|g| g.clone())
.ok_or_else(|| {
error_response(
"Voucher required for remote dispatch but no dashboard voucher pubkey cached yet",
"VOUCHER_UNAVAILABLE",
503,
)
})?;
let v =
super::voucher::verify_voucher_now(&dash_pubkey, &signed_json, &sig, 0.0).map_err(|e| {
eprintln!(" [VOUCHER] minted voucher failed verification: {e}");
error_response(
&format!("Voucher verification failed: {e}"),
"VOUCHER_VERIFY_FAILED",
503,
)
})?;
if !voucher_commission_valid(v.commission_c) {
eprintln!(
" [VOUCHER] minted voucher has invalid/zero commission_c={}, rejecting (fail-closed)",
v.commission_c
);
return Err(error_response(
"Voucher minted but commission_c invalid; refusing to dispatch uncommissioned",
"VOUCHER_COMMISSION_INVALID",
503,
));
}
Ok((signed_json, sig, v))
}
fn voucher_commission_valid(c: f64) -> bool {
c.is_finite() && c > 0.0 && c < 1.0
}
#[allow(clippy::large_enum_variant)]
enum Prepared {
Forward(PreparedForward),
Done(BrokerResponse),
}
fn execute_prepare(state: &Arc<BrokerState>, request: &BrokerRequest, verbose: bool) -> Prepared {
let start = Instant::now();
let request_id = uuid::Uuid::new_v4().to_string();
let tx_num = REQUEST_COUNTER.fetch_add(1, Ordering::SeqCst) + 1;
let instance_action = request.header("X-Zakuro-Instance-Action");
let instance_id = request.header("X-Zakuro-Instance-Id");
let bearer = request
.header("Authorization")
.and_then(|a| a.strip_prefix("Bearer ").map(|t| t.trim().to_string()))
.filter(|t| !t.is_empty());
let requester_key: Option<String> = bearer.clone();
let user_id = if let Some(token) = bearer {
match state.ledger.resolve_user_from_api_key(&token) {
Ok(uid) => uid,
Err(e) => return Prepared::Done(error_response(&e.to_string(), "UNAUTHORIZED", 401)),
}
} else if state.is_local_mode() {
request
.header("X-Zakuro-User")
.unwrap_or_else(|| "anonymous".to_string())
} else {
return Prepared::Done(error_response(
"API key required. Set Authorization: Bearer <key>",
"UNAUTHORIZED",
401,
));
};
let requirements: ResourceRequirements = request
.header("X-Zakuro-Requirements")
.and_then(|s| serde_json::from_str(&s).ok())
.unwrap_or_default();
if let Some(b) = requirements.budget_credits {
if b <= 0.0 {
return Prepared::Done(error_response(
"budget_credits must be positive",
"BAD_REQUEST",
400,
));
}
}
let authority = state.peer_manager.determine_authority(&user_id);
let balance_before = match &authority {
Authority::Local => {
state.ledger.load_balance_if_needed(&user_id)
}
Authority::Peer(url) => {
if let Some(client) = state.peer_manager.get_client(url) {
client
.get_balance(&user_id)
.unwrap_or_else(|_| state.ledger.get_balance(&user_id))
} else {
state.ledger.get_balance(&user_id)
}
}
Authority::Standalone => state.ledger.load_balance_if_needed(&user_id),
};
let _ = state.credits.get_or_create(&user_id, balance_before);
if matches!(authority, Authority::Peer(_)) {
state
.credits
.set_prefetched_balance(&user_id, balance_before);
} else {
state.credits.set_balance(&user_id, balance_before);
}
let pinned_worker = if instance_action.as_deref() == Some("call_method") {
if let Some(ref iid) = instance_id {
if let Some(worker_id) = state.instance_registry.get(iid) {
state.workers.get(worker_id.value())
} else {
log_transaction(
state,
verbose,
tx_num,
&user_id,
"CALL_METHOD",
0.0,
balance_before,
None,
0.0,
"FAIL",
);
return Prepared::Done(error_response(
&format!("Instance not found in broker registry: {}", iid),
"INSTANCE_NOT_FOUND",
404,
));
}
} else {
None
}
} else {
None
};
let routing = if let Some(worker) = pinned_worker {
if worker.status != WorkerStatus::Healthy {
log_transaction(
state,
verbose,
tx_num,
&user_id,
"CALL_METHOD",
0.0,
balance_before,
Some(&worker.name),
0.0,
"FAIL",
);
if let Some(ref iid) = instance_id {
state.instance_registry.remove(iid.as_str());
}
return Prepared::Done(error_response(
"Pinned worker is unhealthy, instance binding removed",
"WORKER_UNHEALTHY",
503,
));
}
RoutingDecision {
estimated_cost: worker
.pricing
.estimate_cost(requirements.estimated_duration_secs),
reason: format!("Instance affinity: pinned to {}", worker.name),
alternatives_count: 0,
worker,
}
} else if state.is_local_mode() || !state.is_billing_enabled() {
match state
.router
.select_worker_no_checks(&state.workers, &requirements)
{
Ok(r) => r,
Err(e) => {
log_transaction(
state,
verbose,
tx_num,
&user_id,
"EXECUTE",
0.0,
balance_before,
None,
0.0,
"FAIL",
);
let status = if e.code == "QUOTA_EXCEEDED" { 429 } else { 503 };
return Prepared::Done(error_response(&e.message, &e.code, status));
}
}
} else if state.peer_manager.is_enabled() {
let is_self_owned = state.config.owner_user_id.as_deref() == Some(user_id.as_str());
let target_node = requirements.target_node.as_deref();
let force_remote = requirements.remote_only
|| target_node
.is_some_and(|tn| !node_identity::node_filter_matches(&state.node_key, tn));
let try_local: Option<RoutingDecision> = if force_remote {
None
} else {
let live_ts_ip =
super::discovery::get_mesh_ip().or_else(|| state.own_wireguard_ip.clone());
state
.router
.select_local_worker(&state.workers, live_ts_ip.as_deref(), &requirements)
.ok()
};
if let Some(local_decision) = try_local {
local_decision
} else {
{
let body = match body_or_413(request) {
Ok(b) => b,
Err(resp) => return Prepared::Done(resp),
};
let (voucher_signed_json, voucher_sig, offer_commission_c_verified) =
if voucher_required() {
let budget = requirements.budget_credits.unwrap_or(0.0);
if budget <= 0.0 {
return Prepared::Done(error_response(
"Voucher required for remote dispatch but no positive \
budget_credits supplied",
"VOUCHER_UNAVAILABLE",
503,
));
}
match mint_and_verify_voucher(state, budget, requester_key.as_deref()) {
Ok((j, s, v)) => (j, s, v.commission_c),
Err(resp) => return Prepared::Done(resp),
}
} else {
(String::new(), String::new(), 0.0)
};
let offer = task_board::TaskOffer {
task_id: request_id.clone(),
payload_b64: base64::Engine::encode(
&base64::engine::general_purpose::STANDARD,
body,
),
max_price_per_hour: if let Some(budget) = requirements.budget_credits {
let hours = requirements.estimated_duration_secs / 3600.0;
if hours > 0.0 {
budget / hours
} else {
f64::MAX
}
} else {
f64::MAX
},
estimated_duration_secs: requirements.estimated_duration_secs,
timeout_secs: requirements.timeout_secs,
cpus: requirements.cpus,
memory_bytes: requirements.memory_bytes,
gpus: requirements.gpus,
worker_type: requirements.worker_type.clone(),
tags: requirements.tags.clone(),
requester_user_id: user_id.clone(),
source_broker: state.config.node_name.clone().unwrap_or_default(),
peer_key: state.peer_manager.peer_key().to_string(),
two_phase: false,
voucher_signed_json,
voucher_sig,
target_worker: String::new(),
};
state.stats.record_task_offer(TaskOfferRecord {
task_id: request_id.clone(),
timestamp: Local::now(),
requester_user_id: user_id.clone(),
max_price_per_hour: offer.max_price_per_hour,
source_broker: offer.source_broker.clone(),
});
let preferred_peer = if marketplace_select_enabled() {
preselect_executor_peer(state, &requirements)
} else {
None
};
return Prepared::Done(dispatch_to_peers(
state,
offer,
start,
&user_id,
is_self_owned,
&request_id,
tx_num,
verbose,
balance_before,
preferred_peer,
offer_commission_c_verified,
));
}
}
} else {
match state.router.select_worker(
&state.workers,
&state.credits,
&user_id,
balance_before,
&requirements,
) {
Ok(r) => r,
Err(ref e) if e.code == "INSUFFICIENT_CREDITS" || e.code == "TIMEOUT_INCOMPATIBLE" => {
if requirements.remote_only {
log_transaction(
state,
verbose,
tx_num,
&user_id,
"EXECUTE",
0.0,
balance_before,
None,
0.0,
"FAIL",
);
let status = if e.code == "INSUFFICIENT_CREDITS" {
402
} else {
503
};
return Prepared::Done(error_response(&e.message, &e.code, status));
}
let owner_here = state.config.owner_user_id.as_deref() == Some(user_id.as_str());
let healthy = if owner_here {
state.workers.healthy()
} else {
Vec::new()
};
match healthy.into_iter().find(|w| state.is_local_worker(&w.uri)) {
Some(local_worker) => RoutingDecision {
worker: local_worker,
estimated_cost: 0.0,
reason: "Local fallback - free execution".to_string(),
alternatives_count: 0,
},
None => {
log_transaction(
state,
verbose,
tx_num,
&user_id,
"EXECUTE",
0.0,
balance_before,
None,
0.0,
"FAIL",
);
let status = if e.code == "INSUFFICIENT_CREDITS" {
402
} else {
503
};
return Prepared::Done(error_response(&e.message, &e.code, status));
}
}
}
Err(e) => {
let status = match e.code.as_str() {
"INSUFFICIENT_CREDITS" => 402,
"RATE_LIMITED" | "QUOTA_EXCEEDED" => 429,
"NO_WORKERS" | "NO_CAPACITY" | "TIMEOUT_INCOMPATIBLE" => 503,
_ => 400,
};
log_transaction(
state,
verbose,
tx_num,
&user_id,
"EXECUTE",
0.0,
balance_before,
None,
0.0,
"FAIL",
);
return Prepared::Done(error_response(&e.message, &e.code, status));
}
}
};
let is_self_owned = state.config.owner_user_id.as_deref() == Some(&user_id);
let is_local = !state.is_billing_enabled() || is_self_owned;
let self_executed =
executor_is_self(state, &routing.worker) || state.is_local_worker(&routing.worker.uri);
let effective_timeout = if let Some(budget) = requirements.budget_credits {
let price_per_sec = routing.worker.pricing.price_per_hour / 3600.0;
let max_secs = if price_per_sec > 0.0 {
budget / price_per_sec
} else {
86_400.0
};
if requirements.timeout_secs > 0.0 {
max_secs.min(requirements.timeout_secs)
} else {
max_secs
}
} else {
requirements.timeout_secs
};
let reservation_amount = if !is_local && effective_timeout > 0.0 {
let timeout_cost = routing.worker.pricing.estimate_cost(effective_timeout);
timeout_cost.min(balance_before)
} else {
routing.estimated_cost
};
let (commission_c, voucher_task_nonce): (f64, String) = if !is_local && voucher_required() {
if self_executed || matches!(authority, Authority::Local) {
match mint_and_verify_remote_voucher(
state,
reservation_amount.max(routing.estimated_cost),
requester_key.as_deref(),
) {
Ok(v) => v,
Err(resp) => return Prepared::Done(resp),
}
} else {
return Prepared::Done(error_response(
"Voucher required for remote dispatch but this broker is not the \
settlement authority for the requester; refusing to run uncommissioned",
"VOUCHER_NON_LOCAL_AUTHORITY",
503,
));
}
} else {
(0.0, String::new())
};
let reservation_id = if is_local {
format!("local-{}", request_id)
} else {
match reserve_credits(
state,
&authority,
&user_id,
reservation_amount,
&request_id,
verbose,
tx_num,
balance_before,
) {
Ok(id) => id,
Err(resp) => return Prepared::Done(resp),
}
};
if !is_local {
let _ = state.wal.append(&WalEntry {
request_id: request_id.clone(),
user_id: user_id.clone(),
reservation_id: reservation_id.clone(),
estimated_cost: reservation_amount,
actual_cost: None,
worker_id: routing.worker.id.clone(),
duration_ms: None,
timestamp: chrono::Utc::now(),
status: WalStatus::Reserved,
});
}
let body = match body_or_413(request) {
Ok(b) => b.to_vec(),
Err(resp) => return Prepared::Done(resp),
};
state
.active_requests
.insert(request_id.clone(), routing.worker.id.clone());
let worker_uri = format!("{}/execute", routing.worker.uri.trim_end_matches('/'));
Prepared::Forward(PreparedForward {
worker_uri,
body,
request_id,
tx_num,
effective_timeout,
user_id,
authority,
reservation_id,
balance_before,
worker: routing.worker,
is_local,
reservation_amount,
instance_action,
instance_id,
start,
commission_c,
voucher_task_nonce,
})
}
fn handle_execute(
state: Arc<BrokerState>,
request: &BrokerRequest,
verbose: bool,
) -> BrokerResponse {
let p = match execute_prepare(&state, request, verbose) {
Prepared::Done(r) => return r,
Prepared::Forward(p) => p,
};
state.workers.increment_active(&p.worker.id);
let forward = forward_to_worker(
&state.http_client,
&p.worker_uri,
&p.body,
&p.request_id,
p.effective_timeout,
)
.map(into_worker_reply)
.map_err(|e| e.to_string());
execute_finish(&state, verbose, p, forward)
}
fn build_broker_runtime(worker_threads: Option<usize>) -> std::io::Result<tokio::runtime::Runtime> {
if matches!(worker_threads, Some(0)) {
return tokio::runtime::Builder::new_current_thread()
.enable_all()
.max_blocking_threads(64)
.build();
}
let mut builder = tokio::runtime::Builder::new_multi_thread();
builder.enable_all().max_blocking_threads(64);
if let Some(n) = worker_threads {
builder.worker_threads(n.max(1));
}
builder.build()
}
async fn shutdown_signal(running: Arc<AtomicBool>) {
let ctrl_c = async {
let _ = tokio::signal::ctrl_c().await;
};
#[cfg(unix)]
let terminate = async {
match tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) {
Ok(mut sig) => {
sig.recv().await;
}
Err(e) => {
eprintln!(" [SHUTDOWN] failed to install SIGTERM handler: {}", e);
std::future::pending::<()>().await;
}
}
};
#[cfg(not(unix))]
let terminate = std::future::pending::<()>();
tokio::select! {
_ = ctrl_c => {}
_ = terminate => {}
}
eprintln!(" [SHUTDOWN] signal received, starting graceful shutdown...");
running.store(false, Ordering::Relaxed);
}
fn serve_axum(
state: Arc<BrokerState>,
host: String,
port: u16,
verbose: bool,
running: Arc<AtomicBool>,
) -> std::io::Result<()> {
use axum::extract::ConnectInfo;
use axum::response::IntoResponse;
let rt = build_broker_runtime(state.config.runtime_worker_threads)?;
rt.block_on(async move {
let inline_safe = !state.ledger.is_api_mode() && !state.peer_manager.is_enabled();
let inline_forced = std::env::var("ZC_INLINE_EXECUTE").ok();
let inline_execute = match inline_forced.as_deref() {
Some(v) => v == "1" || v.eq_ignore_ascii_case("true"),
None => inline_safe,
};
if inline_execute {
eprintln!(
"[PERF] /execute inline fast path ENABLED ({})",
if inline_forced.is_some() {
"forced via ZC_INLINE_EXECUTE"
} else {
"auto: standalone ledger, no P2P"
}
);
}
let app = axum::Router::new()
.route(
"/execute",
axum::routing::post({
let state = state.clone();
move |ConnectInfo(peer): ConnectInfo<SocketAddr>,
req: axum::extract::Request| {
let state = state.clone();
async move {
let breq = BrokerRequest::from_axum(req, Some(peer)).await;
let prepared =
if inline_execute {
execute_prepare(&state, &breq, verbose)
} else {
let st = state.clone();
match tokio::task::spawn_blocking(move || {
execute_prepare(&st, &breq, verbose)
})
.await
{
Ok(p) => p,
Err(_) => return BrokerResponse::json_bytes(
br#"{"error":"handler panicked","code":"INTERNAL"}"#
.to_vec(),
500,
)
.into_response(),
}
};
let p = match prepared {
Prepared::Done(r) => return r.into_response(),
Prepared::Forward(p) => p,
};
state.workers.increment_active(&p.worker.id);
let forward = forward_to_worker_async(
&state.async_http,
&p.worker_uri,
p.body.clone(),
&p.request_id,
p.effective_timeout,
)
.await;
let worker_id = p.worker.id.clone();
if inline_execute {
match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
execute_finish(&state, verbose, p, forward)
})) {
Ok(r) => r.into_response(),
Err(_) => {
state.workers.decrement_active(&worker_id);
BrokerResponse::json_bytes(
br#"{"error":"handler panicked","code":"INTERNAL"}"#
.to_vec(),
500,
)
.into_response()
}
}
} else {
let st = state.clone();
match tokio::task::spawn_blocking(move || {
execute_finish(&st, verbose, p, forward)
})
.await
{
Ok(r) => r.into_response(),
Err(_) => {
state.workers.decrement_active(&worker_id);
BrokerResponse::json_bytes(
br#"{"error":"handler panicked","code":"INTERNAL"}"#
.to_vec(),
500,
)
.into_response()
}
}
}
}
}
}),
)
.fallback(
move |ConnectInfo(peer): ConnectInfo<SocketAddr>, req: axum::extract::Request| {
let state = state.clone();
async move {
let breq = BrokerRequest::from_axum(req, Some(peer)).await;
match tokio::task::spawn_blocking(move || handle(state, &breq, verbose))
.await
{
Ok(resp) => resp.into_response(),
Err(_join_err) => {
BrokerResponse::json_bytes(
br#"{"error":"handler panicked","code":"INTERNAL"}"#.to_vec(),
500,
)
.into_response()
}
}
}
},
);
let listener = tokio::net::TcpListener::bind((host.as_str(), port)).await?;
axum::serve(
listener,
app.into_make_service_with_connect_info::<SocketAddr>(),
)
.with_graceful_shutdown(shutdown_signal(running))
.await
})
}
pub fn start_server(config: BrokerConfig) -> std::io::Result<()> {
let addr = format!("{}:{}", config.host, config.port);
let bind_host = config.host.clone();
let bind_port = config.port;
let verbose = config.verbose;
let daemon = config.daemon;
let mut broker_state = BrokerState::with_config(config.clone());
if broker_state.is_billing_enabled() {
match broker_state.verify_owner_with_dashboard() {
Some(uid) => {
eprintln!(
" [HANDSHAKE] Verified owner: {} (dashboard confirmed)",
uid
);
}
None => {
eprintln!(" [HANDSHAKE] WARNING: Could not verify owner with dashboard — billing continues with configured credentials");
}
}
}
let state = Arc::new(broker_state);
for (pk, revoked) in super::roster_cache::RosterCache::load().entries() {
state.node_roster.insert(pk, revoked);
}
if state.peer_manager.is_enabled() {
state.peer_manager.load_persisted();
}
let quic_disabled = std::env::var("ZAKURO_DISABLE_QUIC")
.map(|v| v == "1" || v == "true")
.unwrap_or(false);
if state.peer_manager.is_enabled() && !quic_disabled {
let quic_port = config.quic_port.unwrap_or(config.port + 1);
let quic_addr = SocketAddr::from(([0, 0, 0, 0], quic_port));
match super::quic::QuicTransport::new(quic_addr, state.clone()) {
Ok(qt) => {
let actual_port = qt.local_addr().map(|a| a.port()).unwrap_or(quic_port);
state.quic_port.store(actual_port, Ordering::Relaxed);
let local = qt.local_addr().map(|a| a.to_string()).unwrap_or_default();
let _ = state.quic.set(qt);
eprintln!(" [QUIC] Listening on {} (peer task transport)", local);
let our_url = format!(
"http://{}:{}",
state.own_wireguard_ip.as_deref().unwrap_or("127.0.0.1"),
config.port
);
super::quic::run_subscription_client(
state.clone(),
our_url,
std::sync::Arc::new(execute_offer_locally),
);
}
Err(e) => {
eprintln!(" [QUIC] Failed to start: {} (falling back to HTTP)", e);
}
}
let state_for_discovery = state.clone();
async_exec::spawn_detached(move || {
thread::sleep(std::time::Duration::from_millis(500));
state_for_discovery.peer_manager.discover_quic_ports();
});
}
recovery::replay_wal(
&state.wal,
&state.ledger,
&state.tx_buffer,
state.config.node_name.as_deref(),
);
let running = Arc::new(AtomicBool::new(true));
let state_wal = state.clone();
let running_wal = running.clone();
async_exec::spawn_detached(move || {
state_wal.wal.run_flush_worker(&running_wal);
});
let state_clone = state.clone();
let running_cleanup = running.clone();
let cleanup_handle = std::thread::Builder::new()
.name("zc-cleanup".into())
.spawn(move || {
let mut compaction_counter: u64 = 0;
while running_cleanup.load(Ordering::Relaxed) {
thread::sleep(std::time::Duration::from_secs(
state_clone.config.health_check_interval,
));
state_clone
.workers
.mark_stale(state_clone.config.worker_timeout as i64);
let removed = state_clone
.workers
.remove_stale(state_clone.config.worker_timeout as i64 * 2);
if !removed.is_empty() {
eprintln!(" [HEALTH] Removed {} stale worker(s)", removed.len());
}
sweep_expired_offers(
&state_clone,
std::time::Duration::from_secs(PENDING_OFFER_TTL_SECS),
);
state_clone.stats.tick_rps();
if compaction_counter.is_multiple_of(3) {
super::deploy::tick(&state_clone);
}
if compaction_counter.is_multiple_of(12) {
if let Some(ref owner_id) = state_clone.config.owner_user_id {
let local_ip = super::discovery::get_effective_node_ip();
let workers: Vec<_> = state_clone
.workers
.list()
.into_iter()
.filter(|w| match w.wireguard_ip.as_deref() {
Some("127.0.0.1") | Some("::1") | Some("localhost") | None => true,
Some(ip) => local_ip.as_deref() == Some(ip),
})
.collect();
let node_name = state_clone.config.node_name.as_deref();
let node_pubkey = state_clone.node_key.public_b64();
if let (Some(ref api_url), Some(ref api_key)) =
(&state_clone.config.api_url, &state_clone.config.api_key)
{
match super::ledger::Ledger::sync_workers_via_api(
owner_id,
&workers,
api_url,
api_key,
node_name,
local_ip.as_deref(),
Some(node_pubkey.as_str()),
) {
Ok(outcome) => {
let repriced =
state_clone.workers.apply_hub_prices(&outcome.prices);
println!(
" [WORKER_SYNC] Synced {} worker(s) to {} via API ({} repriced by the hub)",
workers.len(),
api_url,
repriced
);
}
Err(e) => {
eprintln!(" [WORKER_SYNC] API sync failed: {}", e);
}
}
}
}
}
if compaction_counter.is_multiple_of(12) {
if let (Some(ref api_url), Some(ref api_key)) =
(&state_clone.config.api_url, &state_clone.config.api_key)
{
let pubkey = state_clone.node_key.public_b64();
let endpoint =
mesh_endpoint_for(crate::vpn::native::mesh_ip(), state_clone.config.port);
let self_endpoint = endpoint.clone();
match MESH_ENDPOINT_LOG_GATE
.on_tick(endpoint.is_some(), super::node_identity::now_secs())
{
MeshEndpointLogAction::NoEndpoint => eprintln!(
" [NODE_SYNC] no zakuro0 address - registering without a mesh endpoint"
),
MeshEndpointLogAction::Recovered => println!(
" [NODE_SYNC] mesh endpoint present again - registering as reachable"
),
MeshEndpointLogAction::None => {}
}
let deploy_runtime = super::deploy::runtime_available();
if let Err(e) = super::node_sync::register_node(
api_url,
api_key,
&pubkey,
endpoint.as_deref(),
Some(deploy_runtime),
) {
eprintln!(" [NODE_SYNC] register failed: {e}");
}
match super::node_sync::fetch_roster_full(api_url, api_key) {
Ok((pairs, endpoints)) => {
let fresh: std::collections::HashSet<String> =
pairs.iter().map(|(k, _)| k.clone()).collect();
for (pk, revoked) in pairs.clone() {
state_clone.node_roster.insert(pk, revoked);
}
state_clone.node_roster.retain(|k, _| fresh.contains(k));
let mut cache =
super::roster_cache::RosterCache::from_entries(pairs);
cache.set_endpoints(endpoints);
cache.persist();
let roster_peers = roster_peer_entries(
&cache,
Some(pubkey.as_str()),
self_endpoint.as_deref(),
);
state_clone.peer_manager.ensure_fingerprints_cached();
let mut keep: std::collections::HashSet<String> =
std::collections::HashSet::new();
for rp in &roster_peers {
if state_clone
.peer_manager
.configured_url_for_fingerprint(&rp.fp)
.is_some()
{
continue;
}
state_clone
.peer_manager
.register_roster_peer(rp.url.clone());
keep.insert(rp.url.clone());
}
let evicted =
state_clone.peer_manager.retain_roster_peers(&keep);
if !evicted.is_empty() {
println!(
" [NODE_SYNC] Evicted {} roster-derived peer(s) no longer on the roster",
evicted.len()
);
}
}
Err(e) => eprintln!(" [NODE_SYNC] roster fetch failed: {e}"),
}
match super::node_sync::fetch_voucher_pubkey(api_url) {
Ok(pk) => {
if let Ok(mut w) = state_clone.dash_voucher_pubkey.write() {
*w = Some(pk);
}
}
Err(e) => eprintln!(" [NODE_SYNC] voucher pubkey fetch failed: {e}"),
}
}
}
let unhealthy_ids: Vec<String> = state_clone
.workers
.list()
.iter()
.filter(|w| w.status == WorkerStatus::Unhealthy)
.map(|w| w.id.clone())
.collect();
if !unhealthy_ids.is_empty() {
state_clone
.instance_registry
.retain(|_, worker_id| !unhealthy_ids.contains(worker_id));
}
if let (Some(ref api_url), Some(ref api_key)) =
(&state_clone.config.api_url, &state_clone.config.api_key)
{
state_clone
.tx_buffer
.flush_to_api(api_url, api_key, &state_clone.wal);
}
if let Err(e) = state_clone.wal.flush_buffer() {
eprintln!(" [WAL] flush_buffer failed: {}", e);
}
if compaction_counter.is_multiple_of(12) && state_clone.peer_manager.is_enabled() {
state_clone.peer_manager.health_check_all();
}
if compaction_counter.is_multiple_of(6)
&& state_clone.peer_manager.is_enabled()
&& super::discovery::get_mesh_ip().is_some()
{
let roster = if state_clone.node_roster.is_empty() {
super::roster_cache::RosterCache::load()
} else {
super::roster_cache::RosterCache::from_entries(
state_clone
.node_roster
.iter()
.map(|e| (e.key().clone(), *e.value()))
.collect(),
)
};
super::discovery::run_discovery_round(&state_clone, &roster);
}
if (compaction_counter <= 1 || compaction_counter.is_multiple_of(12))
&& state_clone.peer_manager.is_enabled()
&& state_clone.quic.get().is_some()
{
state_clone.peer_manager.discover_quic_ports();
}
if compaction_counter.is_multiple_of(6) && state_clone.peer_manager.is_enabled() {
let user_ids = state_clone.credits.get_all_user_ids();
for uid in user_ids {
if state_clone.credits.needs_reconciliation(&uid, 30) {
state_clone.credits.set_reconciling(&uid);
let authority = state_clone.peer_manager.determine_authority(&uid);
if let Authority::Peer(ref url) = authority {
if let Some(client) = state_clone.peer_manager.get_client(url) {
match client.get_balance(&uid) {
Ok(balance) => {
state_clone
.credits
.set_prefetched_balance(&uid, balance);
}
Err(_) => {
}
}
}
}
}
}
}
compaction_counter += 1;
if compaction_counter.is_multiple_of(60) {
if let Err(e) = state_clone.wal.compact() {
eprintln!(" [WAL] Compaction failed: {}", e);
}
}
}
if let (Some(ref api_url), Some(ref api_key)) =
(&state_clone.config.api_url, &state_clone.config.api_key)
{
eprintln!(" [FLUSH] Shutdown: flushing remaining transactions...");
state_clone
.tx_buffer
.flush_to_api(api_url, api_key, &state_clone.wal);
}
if let Err(e) = state_clone.wal.flush_buffer() {
eprintln!(" [WAL] flush_buffer failed: {}", e);
}
})
.expect("failed to spawn zc-cleanup thread");
if config.enable_discovery {
let state_for_discovery = state.clone();
let discovery_config = config.discovery.clone();
let discovery_verbose = verbose;
let mode = detect_discovery_mode(&discovery_config.subnet);
let is_local = matches!(mode, DiscoveryMode::Local);
state.set_local_mode(is_local);
if verbose && !daemon {
match &mode {
DiscoveryMode::WireGuard { subnet } => {
println!(
" {} WireGuard network detected (subnet: {}.0/24)",
"[DISCOVERY]".cyan(),
subnet
);
}
DiscoveryMode::Local => {
println!(
" {} Local mode - free execution, scanning localhost:3960-3962",
"[DISCOVERY]".yellow()
);
}
}
}
async_exec::spawn_detached(move || {
let discovery = Discovery::new(discovery_config, state_for_discovery);
discovery.run(discovery_verbose);
});
}
if !daemon {
if state.ledger.is_api_mode() {
println!(" {} API mode - using dashboard API", "[LEDGER]".cyan());
} else {
println!(
" {} Standalone mode - using local in-memory operations",
"[LEDGER]".yellow()
);
}
if state.is_billing_enabled() {
println!(
" {} Billing enabled (dashboard API authority)",
"[BILLING]".cyan()
);
} else {
println!(
" {} Billing disabled — all executions are free (no API credentials)",
"[BILLING]".yellow()
);
}
if state.peer_manager.is_enabled() {
println!(
" {} P2P credit operations enabled ({} peers)",
"[P2P]".cyan(),
state.peer_manager.peer_count(),
);
}
println!(" {} Listening on {}", "[BROKER]".green().bold(), addr);
if verbose {
println!();
println!(" {}", "Live Transactions:".bold().underline());
println!(" {}", "─".repeat(80));
println!(
" {} {:>4} {:>12} {:>8} {:>10} {:>10} {:>8} {}",
"TIME".dimmed(),
"#".dimmed(),
"USER".dimmed(),
"ACTION".dimmed(),
"COST".dimmed(),
"BALANCE".dimmed(),
"LATENCY".dimmed(),
"WORKER".dimmed(),
);
println!(" {}", "─".repeat(80));
}
}
let serve_result = serve_axum(state, bind_host, bind_port, verbose, running);
let _ = cleanup_handle.join();
serve_result
}
#[cfg(test)]
mod dispatch_deadline_tests {
use super::peer_dispatch_deadline;
use std::time::Duration;
#[test]
fn unset_timeout_caps_at_30s_not_the_peer_clients_310s_default() {
assert_eq!(peer_dispatch_deadline(0.0), Duration::from_secs(30));
}
#[test]
fn explicit_timeout_is_honored_with_slack() {
assert_eq!(peer_dispatch_deadline(5.0), Duration::from_secs_f64(17.0));
}
#[test]
fn preselected_offer_wait_is_bounded_by_dispatch_deadline_not_offer_tasks_310s() {
use crate::broker::peer::PeerClient;
use crate::broker::task_board;
use std::io::Read;
use std::net::TcpListener;
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
std::thread::spawn(move || {
if let Ok((mut stream, _)) = listener.accept() {
let mut buf = [0u8; 1024];
let _ = stream.read(&mut buf);
std::thread::sleep(Duration::from_secs(600));
}
});
let client = PeerClient::new(format!("http://{addr}"), "test-key".to_string());
let offer: task_board::TaskOffer = serde_json::from_str(
r#"{"task_id":"t-bound","payload_b64":"","max_price_per_hour":100.0,"estimated_duration_secs":1.0,"timeout_secs":0.0,"requester_user_id":"u","source_broker":"n"}"#,
)
.unwrap();
let (tx, rx) = std::sync::mpsc::channel::<Result<task_board::TaskResult, String>>();
let offer_c = offer.clone();
std::thread::spawn(move || {
let _ = tx.send(client.offer_task(&offer_c));
});
let deadline = peer_dispatch_deadline(offer.timeout_secs);
assert_eq!(
deadline,
Duration::from_secs(30),
"sanity: unset timeout_secs -> 30s bound"
);
let start = std::time::Instant::now();
let outcome = rx.recv_timeout(deadline);
let elapsed = start.elapsed();
assert!(
outcome.is_err(),
"stub peer never responds, so the channel recv must itself time out"
);
assert!(
elapsed < Duration::from_secs(60),
"pre-selected offer wait was not bounded by peer_dispatch_deadline: took {elapsed:?}"
);
}
#[test]
fn broadcast_targets_excludes_abandoned_preferred_peer() {
use super::broadcast_targets;
let peers = vec![
"http://peer-a".to_string(),
"http://peer-b".to_string(),
"http://peer-c".to_string(),
];
assert_eq!(broadcast_targets(&peers, None), peers);
let targets = broadcast_targets(&peers, Some("http://peer-b"));
assert_eq!(
targets,
vec!["http://peer-a".to_string(), "http://peer-c".to_string()]
);
assert!(
!targets.contains(&"http://peer-b".to_string()),
"abandoned preferred peer must not be re-offered the same task_id"
);
assert_eq!(broadcast_targets(&peers, Some("http://peer-z")), peers);
}
}
#[cfg(test)]
mod runtime_tests {
use super::build_broker_runtime;
#[test]
fn runtime_honors_explicit_worker_thread_count() {
let rt = build_broker_runtime(Some(2)).unwrap();
assert_eq!(rt.metrics().num_workers(), 2);
}
#[test]
fn runtime_defaults_to_multi_thread_when_unset() {
let rt = build_broker_runtime(None).unwrap();
assert!(rt.metrics().num_workers() >= 1);
}
#[test]
fn runtime_zero_workers_is_current_thread() {
let rt = build_broker_runtime(Some(0)).unwrap();
assert_eq!(rt.metrics().num_workers(), 1);
}
}
#[cfg(test)]
mod panic_fix_tests {
use super::header_or;
use super::{check_node_sig, worker_to_info, BrokerRequest};
#[test]
fn worker_to_info_client_hides_ip_peer_keeps_it() {
use crate::broker::node_identity::NodeKey;
use crate::broker::worker::{Worker, WorkerRegistry};
let reg = WorkerRegistry::new();
let key = NodeKey::generate();
let mut w = Worker::new(
"id1".into(),
"worker-xyz".into(),
"http://10.13.13.21:3960".into(),
"zakuro".into(),
);
w.source_node = Some("node-i9".into());
w.node_fp = "abc123".into();
w.slot = "3960".into();
let ci = worker_to_info(&w, ®, &key, true);
assert_eq!(ci.uri, "zc://worker-abc123-3960");
assert_eq!(ci.node.as_deref(), Some("zc://node-i9"));
assert_eq!(ci.label.as_deref(), Some("worker-xyz"));
let cj = serde_json::to_string(&ci).unwrap();
assert!(!cj.contains("10.13.13.") && !cj.contains(":3960") && !cj.contains("http://"));
let pi = worker_to_info(&w, ®, &key, false);
assert_eq!(pi.uri, "http://10.13.13.21:3960");
}
#[test]
fn worker_to_info_carries_provider_fields_for_peers() {
use crate::broker::node_identity::NodeKey;
use crate::broker::worker::{ProviderType, Worker, WorkerRegistry};
let reg = WorkerRegistry::new();
let key = NodeKey::generate();
let mut w = Worker::new(
"id3".into(),
"zc-serve-116c1f8b".into(),
"http://192.168.0.167:8600".into(),
"llm".into(),
);
w.provider_type = ProviderType::Specialized;
w.served_models = vec!["116c1f8b-7a22-452a-867c-f8b20afdb355".into()];
w.price_per_mtok = 1.5;
let pi = worker_to_info(&w, ®, &key, false);
assert_eq!(pi.provider_type, ProviderType::Specialized);
assert_eq!(pi.served_models, w.served_models);
assert_eq!(pi.price_per_mtok, 1.5);
let j = serde_json::to_value(&pi).unwrap();
assert_eq!(j["provider_type"], "specialized");
assert_eq!(
j["served_models"][0],
"116c1f8b-7a22-452a-867c-f8b20afdb355"
);
}
#[test]
fn no_client_surface_emits_ip() {
use crate::broker::node_identity::NodeKey;
use crate::broker::worker::{Worker, WorkerRegistry};
let reg = WorkerRegistry::new();
let key = NodeKey::generate();
let mut w = Worker::new(
"id2".into(),
"worker-abc".into(),
"http://10.13.13.5:3960".into(),
"zakuro".into(),
);
w.node_fp = key.fingerprint();
w.slot = "3960".into();
let ci = worker_to_info(&w, ®, &key, true);
let body = serde_json::to_string(&super::WorkerListResponse {
total: 1,
workers: vec![ci],
})
.unwrap();
assert!(!body.contains("10.13.13."));
assert!(body.contains("zc://worker-"));
}
#[test]
fn client_response_uses_key_derived_node_uri() {
use crate::broker::node_identity::NodeKey;
use crate::broker::worker::{Worker, WorkerRegistry};
let key = NodeKey::generate();
let handle = super::node_handle_for(&key, Some("lxd"));
assert_eq!(handle, key.node_uri());
assert!(handle.starts_with("zc://node-"));
assert!(!handle.contains("lxd"));
let reg = WorkerRegistry::new();
let w = Worker::new(
"id2".into(),
"local-worker".into(),
"http://127.0.0.1:3960".into(),
"zakuro".into(),
);
let ci = worker_to_info(&w, ®, &key, true);
assert_eq!(ci.node.as_deref(), Some(key.node_uri().as_str()));
}
#[test]
fn capped_charge_never_exceeds_reserved() {
use super::capped_charge;
assert!((capped_charge(0.05, 0.02) - 0.02).abs() < 1e-9); assert!((capped_charge(0.01, 0.02) - 0.01).abs() < 1e-9); assert_eq!(capped_charge(f64::NAN, 0.02), 0.0); assert_eq!(capped_charge(-1.0, 0.02), 0.0); }
#[test]
fn header_or_accepts_valid_value() {
let (k, v) = header_or("X-Zakuro-Worker", "node-01");
assert_eq!(k, "X-Zakuro-Worker");
assert_eq!(v, "node-01");
}
#[test]
fn check_node_sig_accepts_rostered_signer_rejects_others() {
use crate::broker::node_identity::{NodeKey, ReplayGuard};
use dashmap::DashMap;
use tiny_http::Method;
let key = NodeKey::generate();
let path = "/peer/tasks/offer";
let body = br#"{"fn":"greet"}"#.to_vec();
let hdrs = key.sign_headers("POST", path, &body);
let build = |headers: Vec<(String, String)>, body: Vec<u8>| BrokerRequest {
method: Method::Post,
path: path.to_string(),
query: String::new(),
headers,
body,
peer_addr: None,
};
let roster = DashMap::new();
roster.insert(key.public_b64(), false); let guard = ReplayGuard::new();
let now = crate::broker::node_identity::now_secs();
let req = build(hdrs.clone(), body.clone());
assert_eq!(
check_node_sig(&roster, &guard, &req, now).unwrap(),
key.public_b64()
);
roster.insert(key.public_b64(), true);
let g2 = ReplayGuard::new();
let req2 = build(hdrs.clone(), body.clone());
assert!(check_node_sig(&roster, &g2, &req2, now).is_err());
roster.insert(key.public_b64(), false);
let g3 = ReplayGuard::new();
let req3 = build(hdrs.clone(), b"{}".to_vec());
assert!(check_node_sig(&roster, &g3, &req3, now).is_err());
let g4 = ReplayGuard::new();
let req4 = build(vec![], body);
assert!(check_node_sig(&roster, &g4, &req4, now).is_err());
}
#[test]
fn peer_peers_requires_valid_node_signature_and_signs_response() {
use crate::broker::node_identity::{verify_sig, NodeKey};
use crate::broker::BrokerState;
use std::sync::Arc;
use tiny_http::Method;
let state = BrokerState::new();
let signer = NodeKey::generate();
state.node_roster.insert(signer.public_b64(), false); let state = Arc::new(state);
let unsigned = BrokerRequest {
method: Method::Get,
path: "/peer/peers".to_string(),
query: String::new(),
headers: vec![],
body: vec![],
peer_addr: None,
};
let resp = super::handle(Arc::clone(&state), &unsigned, false);
assert_eq!(resp.status, 401);
let hdrs = signer.sign_headers("GET", "/peer/peers", b"");
let signed = BrokerRequest {
method: Method::Get,
path: "/peer/peers".to_string(),
query: String::new(),
headers: hdrs,
body: vec![],
peer_addr: None,
};
let resp = super::handle(Arc::clone(&state), &signed, false);
assert_eq!(resp.status, 200);
let get = |name: &str| -> Option<String> {
resp.headers
.iter()
.find(|(k, _)| k.eq_ignore_ascii_case(name))
.map(|(_, v)| v.clone())
};
let node_id = get("X-Node-Id").expect("response signed with X-Node-Id");
let sig = get("X-Node-Sig").expect("response signed with X-Node-Sig");
assert_eq!(node_id, state.node_key.public_b64());
assert!(verify_sig(&node_id, &resp.body, &sig));
let parsed: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert!(parsed["peers"].as_array().is_some());
}
#[test]
fn header_or_does_not_panic_on_peer_supplied_garbage() {
let (k, v) = header_or("X-Zakuro-Worker", "wörker\u{00ff}");
assert_eq!(k, "X-Zakuro-Invalid");
assert_eq!(v, "1");
}
#[test]
fn health_endpoint_body_has_no_mesh_ip() {
use crate::broker::BrokerState;
use std::sync::Arc;
use tiny_http::Method;
let mut state = BrokerState::new();
state.own_wireguard_ip = Some("10.13.13.42".to_string());
let state = Arc::new(state);
let req = BrokerRequest {
method: Method::Get,
path: "/health".to_string(),
query: String::new(),
headers: vec![],
body: vec![],
peer_addr: None,
};
let resp = super::handle(state, &req, false);
let body = String::from_utf8(resp.body).unwrap();
assert!(
!body.contains("10.13.13."),
"health body leaked mesh IP: {body}"
);
assert!(
!body.contains("127.0.0.1"),
"health body leaked loopback: {body}"
);
assert!(body.contains("wireguard_connected"));
assert!(!body.contains("\"wireguard_ip\""));
}
#[test]
fn health_wireguard_connected_and_mesh_tunnel_never_disagree() {
use crate::broker::BrokerState;
use std::sync::Arc;
use tiny_http::Method;
let mut state = BrokerState::new();
state.own_wireguard_ip = Some("10.13.13.42".to_string());
let state = Arc::new(state);
let req = BrokerRequest {
method: Method::Get,
path: "/health".to_string(),
query: String::new(),
headers: vec![],
body: vec![],
peer_addr: None,
};
let resp = super::handle(state, &req, false);
let body: serde_json::Value =
serde_json::from_slice(&resp.body).expect("/health returns JSON");
let connected = body["wireguard_connected"]
.as_bool()
.expect("wireguard_connected is a bool");
let tunnel = body["mesh_tunnel"]
.as_str()
.expect("mesh_tunnel is a string");
assert!(tunnel == "up" || tunnel == "down", "mesh_tunnel={tunnel}");
assert_eq!(
connected,
tunnel == "up",
"/health contradicts itself: wireguard_connected={connected} mesh_tunnel={tunnel}"
);
}
#[test]
fn stats_worker_and_response_serialization_has_no_mesh_ip() {
use crate::broker::stats::{StatsCollector, StatsResponse, WorkerStats};
let ws = WorkerStats {
id: "id1".into(),
name: "worker-abc".into(),
uri: "zc://worker-abc123-3960".into(), status: "healthy".into(),
cpus_available: 1.0,
memory_available_gib: 1.0,
gpus_available: 0,
price_per_hour: 0.0,
active_requests: 0,
avg_latency_ms: 0.0,
};
let ws_json = serde_json::to_string(&ws).unwrap();
assert!(!ws_json.contains("10.13.13."));
assert!(!ws_json.contains("127.0.0.1"));
assert!(ws_json.contains("zc://worker-"));
let resp = StatsResponse {
host: "0.0.0.0".into(),
port: 8080,
transactions: vec![],
task_offers: vec![],
workers: vec![ws],
metrics: StatsCollector::new().metrics(0, 0, true, false, None),
rps_history: vec![],
wireguard_ip: None,
wireguard_connected: true,
};
let resp_json = serde_json::to_string(&resp).unwrap();
assert!(
!resp_json.contains("10.13.13."),
"stats body leaked mesh IP: {resp_json}"
);
assert!(
!resp_json.contains("127.0.0.1"),
"stats body leaked loopback: {resp_json}"
);
}
#[test]
fn mesh_endpoint_is_ip_plus_bound_port() {
assert_eq!(
super::mesh_endpoint_for(Some("10.13.13.6".to_string()), 9100).as_deref(),
Some("10.13.13.6:9100")
);
assert_eq!(super::mesh_endpoint_for(None, 9000), None);
}
#[test]
fn mesh_endpoint_log_gate_logs_on_first_absence() {
let gate = super::MeshEndpointLogGate::new();
assert_eq!(
gate.on_tick(false, 1_000),
super::MeshEndpointLogAction::NoEndpoint
);
}
#[test]
fn mesh_endpoint_log_gate_suppresses_repeat_within_interval() {
let gate = super::MeshEndpointLogGate::new();
assert_eq!(
gate.on_tick(false, 1_000),
super::MeshEndpointLogAction::NoEndpoint
);
assert_eq!(
gate.on_tick(false, 1_060),
super::MeshEndpointLogAction::None
);
}
#[test]
fn mesh_endpoint_log_gate_logs_again_after_interval_elapses() {
let gate = super::MeshEndpointLogGate::new();
assert_eq!(
gate.on_tick(false, 1_000),
super::MeshEndpointLogAction::NoEndpoint
);
assert_eq!(
gate.on_tick(false, 1_000 + 30 * 60),
super::MeshEndpointLogAction::NoEndpoint
);
}
#[test]
fn mesh_endpoint_log_gate_logs_recovery_on_transition_back() {
let gate = super::MeshEndpointLogGate::new();
assert_eq!(
gate.on_tick(false, 1_000),
super::MeshEndpointLogAction::NoEndpoint
);
assert_eq!(
gate.on_tick(true, 1_005),
super::MeshEndpointLogAction::Recovered
);
}
#[test]
fn mesh_endpoint_log_gate_stays_silent_while_connected() {
let gate = super::MeshEndpointLogGate::new();
assert_eq!(
gate.on_tick(true, 1_000),
super::MeshEndpointLogAction::None
);
assert_eq!(
gate.on_tick(true, 1_000 + 30 * 60),
super::MeshEndpointLogAction::None
);
}
fn urls_of(entries: Vec<super::RosterPeer>) -> Vec<String> {
entries.into_iter().map(|e| e.url).collect()
}
#[test]
fn roster_peer_entries_excludes_self_by_pubkey_when_tunnel_is_down() {
use crate::broker::node_identity::NodeKey;
use crate::broker::roster_cache::RosterCache;
let self_key = NodeKey::generate();
let peer_key = NodeKey::generate();
let mut cache = RosterCache::from_entries(vec![
(self_key.public_b64(), false),
(peer_key.public_b64(), false),
]);
cache.set_endpoints(vec![
(self_key.public_b64(), "10.13.13.5:9000".to_string()),
(peer_key.public_b64(), "10.13.13.6:9000".to_string()),
]);
let urls = urls_of(super::roster_peer_entries(
&cache,
Some(&self_key.public_b64()),
None,
));
assert_eq!(
urls,
vec!["http://10.13.13.6:9000".to_string()],
"broker registered itself as a peer while its tunnel was down"
);
}
#[test]
fn roster_peer_entries_excludes_self_even_when_its_address_changed() {
use crate::broker::node_identity::NodeKey;
use crate::broker::roster_cache::RosterCache;
let self_key = NodeKey::generate();
let mut cache = RosterCache::from_entries(vec![(self_key.public_b64(), false)]);
cache.set_endpoints(vec![(
self_key.public_b64(),
"10.13.13.99:9000".to_string(),
)]);
let urls = urls_of(super::roster_peer_entries(
&cache,
Some(&self_key.public_b64()),
Some("10.13.13.5:9000"),
));
assert!(
urls.is_empty(),
"self row survived the identity filter: {urls:?}"
);
}
#[test]
fn roster_peer_entries_still_drops_revoked_rows_with_identity_filter_on() {
use crate::broker::node_identity::NodeKey;
use crate::broker::roster_cache::RosterCache;
let self_key = NodeKey::generate();
let revoked = NodeKey::generate();
let mut cache = RosterCache::from_entries(vec![
(self_key.public_b64(), false),
(revoked.public_b64(), true),
]);
cache.set_endpoints(vec![(revoked.public_b64(), "10.13.13.7:9000".to_string())]);
assert!(super::roster_peer_entries(&cache, Some(&self_key.public_b64()), None).is_empty());
}
#[test]
fn roster_peer_entries_carry_the_nodes_fingerprint() {
use crate::broker::node_identity::NodeKey;
use crate::broker::roster_cache::RosterCache;
let peer_key = NodeKey::generate();
let mut cache = RosterCache::from_entries(vec![(peer_key.public_b64(), false)]);
cache.set_endpoints(vec![(peer_key.public_b64(), "10.13.13.6:9000".to_string())]);
let entries = super::roster_peer_entries(&cache, None, None);
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].fp, peer_key.fingerprint());
assert_eq!(entries[0].url, "http://10.13.13.6:9000");
}
#[test]
fn roster_peer_urls_skips_revoked_entries() {
use crate::broker::node_identity::NodeKey;
use crate::broker::roster_cache::RosterCache;
let key = NodeKey::generate();
let mut cache = RosterCache::from_entries(vec![(key.public_b64(), true)]);
cache.set_endpoints(vec![(key.public_b64(), "10.13.13.6:9000".to_string())]);
assert!(super::roster_peer_entries(&cache, None, None).is_empty());
}
#[test]
fn roster_peer_urls_excludes_self_endpoint() {
use crate::broker::node_identity::NodeKey;
use crate::broker::roster_cache::RosterCache;
let self_key = NodeKey::generate();
let peer_key = NodeKey::generate();
let mut cache = RosterCache::from_entries(vec![
(self_key.public_b64(), false),
(peer_key.public_b64(), false),
]);
cache.set_endpoints(vec![
(self_key.public_b64(), "10.13.13.5:9000".to_string()),
(peer_key.public_b64(), "10.13.13.6:9000".to_string()),
]);
let urls = urls_of(super::roster_peer_entries(
&cache,
None,
Some("10.13.13.5:9000"),
));
assert_eq!(urls, vec!["http://10.13.13.6:9000".to_string()]);
}
#[test]
fn roster_peer_urls_prefixes_bare_addresses_with_http() {
use crate::broker::node_identity::NodeKey;
use crate::broker::roster_cache::RosterCache;
let key = NodeKey::generate();
let mut cache = RosterCache::from_entries(vec![(key.public_b64(), false)]);
cache.set_endpoints(vec![(key.public_b64(), "10.13.13.6:9000".to_string())]);
assert_eq!(
urls_of(super::roster_peer_entries(&cache, None, None)),
vec!["http://10.13.13.6:9000".to_string()]
);
}
#[test]
fn roster_peer_urls_empty_roster_yields_empty_vec() {
use crate::broker::roster_cache::RosterCache;
let cache = RosterCache::from_entries(vec![]);
assert!(super::roster_peer_entries(&cache, None, None).is_empty());
}
#[test]
fn roster_peer_urls_skips_empty_and_whitespace_only_endpoints() {
use crate::broker::node_identity::NodeKey;
use crate::broker::roster_cache::RosterCache;
let blank_key = NodeKey::generate();
let good_key = NodeKey::generate();
let mut cache = RosterCache::from_entries(vec![
(blank_key.public_b64(), false),
(good_key.public_b64(), false),
]);
cache.set_endpoints(vec![
(blank_key.public_b64(), " ".to_string()),
(good_key.public_b64(), "10.13.13.6:9000".to_string()),
]);
assert_eq!(
urls_of(super::roster_peer_entries(&cache, None, None)),
vec!["http://10.13.13.6:9000".to_string()]
);
}
}
#[cfg(test)]
mod two_phase_tests {
use super::{
cancel_offer, earn_refund_amount, offer_commission_c, reserve_offer, split_commission,
sweep_expired_offers, take_committed, voucher_commission_valid,
};
use crate::broker::{task_board, BrokerState, PendingOffer};
fn mk_offer(task_id: &str) -> task_board::TaskOffer {
serde_json::from_str(&format!(
r#"{{"task_id":"{}","payload_b64":"","max_price_per_hour":100.0,"estimated_duration_secs":1.0,"timeout_secs":10.0,"requester_user_id":"u","source_broker":"n"}}"#,
task_id
)).unwrap()
}
#[test]
fn split_commission_ten_percent_exact() {
let (payout, commission) = split_commission(100.0, 0.1);
assert_eq!(commission, 10.0);
assert_eq!(payout, 90.0);
assert_eq!(payout + commission, 100.0);
}
#[test]
fn split_commission_min_charge_no_leak() {
let gross = 0.001;
let (payout, commission) = split_commission(gross, 0.1);
assert_eq!(payout + commission, gross);
}
#[test]
fn split_commission_zero_is_identity() {
let (payout, commission) = split_commission(42.5, 0.0);
assert_eq!(payout, 42.5);
assert_eq!(commission, 0.0);
}
#[test]
fn settlement_conserves_for_all_commission_rates() {
for &cost in &[0.001_f64, 1.0, 0.036668, 100.0, 0.0] {
for &c in &[0.0_f64, 0.1, 0.5, 0.9, 0.999] {
let (earn, commission) = if c > 0.0 {
split_commission(cost, c)
} else {
(cost, 0.0)
};
assert_eq!(
earn + commission,
cost,
"conservation broken: cost={cost} c={c} earn={earn} commission={commission}"
);
}
}
}
#[test]
fn settle_payout_conserves_at_named_commission_rates() {
for &cost in &[0.001_f64, 1.0, 100.0] {
for &c in &[0.0_f64, 0.1, 0.5] {
let (payout, commission) = split_commission(cost, c);
assert_eq!(payout + commission, cost, "cost={cost} c={c}");
if c == 0.0 {
assert_eq!(commission, 0.0);
assert_eq!(payout, cost);
} else {
assert!(commission > 0.0, "commission must be > 0 for c={c}");
}
}
}
}
#[test]
fn earn_refund_conserves_on_success_and_failure() {
for &cost in &[0.001_f64, 1.0, 0.036668, 100.0] {
for &c in &[0.0_f64, 0.1, 0.5, 0.9] {
let (payout, commission) = split_commission(cost, c);
let refund_ok = earn_refund_amount(payout, true);
assert_eq!(refund_ok, 0.0);
let requester_net_ok = -cost + refund_ok;
let executor_delta_ok = payout;
let platform_delta_ok = commission;
let sum_ok = requester_net_ok + executor_delta_ok + platform_delta_ok;
assert!(
sum_ok.abs() < 1e-9,
"earn-success conservation broken: cost={cost} c={c} sum={sum_ok}"
);
assert!(
(requester_net_ok - (-cost)).abs() < 1e-9,
"cost={cost} c={c}"
);
let refund_fail = earn_refund_amount(payout, false);
assert_eq!(refund_fail, payout);
let requester_net_fail = -cost + refund_fail;
let executor_delta_fail = 0.0;
let platform_delta_fail = commission;
let sum_fail = requester_net_fail + executor_delta_fail + platform_delta_fail;
assert!(
sum_fail.abs() < 1e-9,
"earn-fail conservation broken: cost={cost} c={c} sum={sum_fail}"
);
assert!(
(requester_net_fail - (-commission)).abs() < 1e-9,
"requester should net exactly -commission when earn fails: cost={cost} c={c}"
);
}
}
}
#[test]
fn unresolved_peer_full_refund_makes_requester_whole() {
for &actual_cost in &[0.001_f64, 1.0, 0.036668, 100.0] {
let requester_net = -actual_cost + actual_cost;
let executor_delta = 0.0;
let platform_delta = 0.0; let sum = requester_net + executor_delta + platform_delta;
assert!(
sum.abs() < 1e-9,
"unresolved-peer conservation broken: actual_cost={actual_cost} sum={sum}"
);
assert!(
requester_net.abs() < 1e-9,
"requester must be made whole (net 0): actual_cost={actual_cost}"
);
}
}
#[test]
fn voucher_commission_valid_gate() {
assert!(voucher_commission_valid(0.1));
assert!(voucher_commission_valid(0.5));
assert!(voucher_commission_valid(0.999));
assert!(
!voucher_commission_valid(0.0),
"zero commission fails closed"
);
assert!(!voucher_commission_valid(1.0), "out of [0,1) fails closed");
assert!(!voucher_commission_valid(1.5), "out of range fails closed");
assert!(!voucher_commission_valid(-0.1), "negative fails closed");
assert!(!voucher_commission_valid(f64::NAN), "NaN fails closed");
}
#[test]
fn mint_and_verify_voucher_rejects_without_api_creds() {
let state = std::sync::Arc::new(BrokerState::new());
let resp = super::mint_and_verify_voucher(&state, 5.0, None)
.expect_err("no api creds → fail-closed reject");
assert_eq!(resp.status, 503, "reject must be a 503");
}
#[test]
fn mint_and_verify_voucher_fails_closed_even_with_cached_pubkey() {
let state = std::sync::Arc::new(BrokerState::new());
if let Ok(mut g) = state.dash_voucher_pubkey.write() {
*g = Some(crate::broker::node_identity::NodeKey::generate().public_b64());
}
assert!(
super::mint_and_verify_voucher(&state, 5.0, None).is_err(),
"missing api creds must reject regardless of cached pubkey"
);
}
fn mk_signed_voucher(
commission_c: f64,
) -> (
String,
String,
String,
crate::broker::node_identity::NodeKey,
) {
use crate::broker::node_identity::NodeKey;
let dash = NodeKey::generate();
let signed = format!(
r#"{{"v":1,"requester_user":"u","budget_credits":1000000.0,"task_nonce":"n","exp":9999999999,"commission_c":{}}}"#,
commission_c
);
let sig = dash.sign(signed.as_bytes());
let pubkey = dash.public_b64();
(signed, sig, pubkey, dash)
}
#[test]
fn offer_commission_c_reads_validly_signed_voucher() {
let (signed, sig, pubkey, _dash) = mk_signed_voucher(0.1);
let mut offer = mk_offer("t-c");
offer.voucher_signed_json = signed;
offer.voucher_sig = sig;
assert_eq!(offer_commission_c(&offer, Some(&pubkey)), 0.1);
}
#[test]
fn offer_commission_c_defaults_zero_without_voucher() {
let offer = mk_offer("t-nc");
assert!(offer.voucher_signed_json.is_empty());
assert_eq!(offer_commission_c(&offer, Some("irrelevant")), 0.0);
}
#[test]
fn offer_commission_c_parse_failure_falls_back_to_zero() {
let mut offer = mk_offer("t-bad");
offer.voucher_signed_json = "{not valid json".to_string();
offer.voucher_sig = "somesig".to_string();
assert_eq!(offer_commission_c(&offer, Some("irrelevant")), 0.0);
}
#[test]
fn offer_commission_c_out_of_range_falls_back_to_zero() {
let (signed, sig, pubkey, _dash) = mk_signed_voucher(1.0);
let mut offer = mk_offer("t-oob");
offer.voucher_signed_json = signed;
offer.voucher_sig = sig;
assert_eq!(offer_commission_c(&offer, Some(&pubkey)), 0.0);
}
#[test]
fn offer_commission_c_zero_when_signature_missing() {
let mut offer = mk_offer("t-nosig");
offer.voucher_signed_json = r#"{"v":1,"requester_user":"u","budget_credits":5.0,"task_nonce":"n","exp":9999999999,"commission_c":0.5}"#.to_string();
offer.voucher_sig = String::new();
assert_eq!(offer_commission_c(&offer, Some("some-pubkey")), 0.0);
}
#[test]
fn offer_commission_c_zero_when_signature_invalid() {
let (signed, _sig, pubkey, _dash) = mk_signed_voucher(0.5);
let mut offer = mk_offer("t-badsig");
offer.voucher_signed_json = signed;
offer.voucher_sig = "not-a-real-signature".to_string();
assert_eq!(offer_commission_c(&offer, Some(&pubkey)), 0.0);
}
#[test]
fn offer_commission_c_zero_when_signed_by_wrong_key() {
let (signed, sig, _their_pubkey, _their_dash) = mk_signed_voucher(0.5);
let (_real_signed, _real_sig, real_pubkey, _real_dash) = mk_signed_voucher(0.1);
let mut offer = mk_offer("t-wrongkey");
offer.voucher_signed_json = signed;
offer.voucher_sig = sig;
assert_eq!(offer_commission_c(&offer, Some(&real_pubkey)), 0.0);
}
#[test]
fn offer_commission_c_zero_when_no_pubkey_cached() {
let (signed, sig, _pubkey, _dash) = mk_signed_voucher(0.5);
let mut offer = mk_offer("t-nopubkey");
offer.voucher_signed_json = signed;
offer.voucher_sig = sig;
assert_eq!(offer_commission_c(&offer, None), 0.0);
}
#[test]
fn offer_commission_c_zero_when_tampered_after_signing() {
let (signed, sig, pubkey, _dash) = mk_signed_voucher(0.5);
let tampered = signed.replace("0.5", "0.0000001");
let mut offer = mk_offer("t-tampered");
offer.voucher_signed_json = tampered;
offer.voucher_sig = sig;
assert_eq!(offer_commission_c(&offer, Some(&pubkey)), 0.0);
}
#[test]
fn reserve_offer_rejects_when_no_local_worker() {
let state = BrokerState::new();
let r = reserve_offer(&state, &mk_offer("t1"));
assert!(r.is_err(), "no local worker → TaskReject");
}
#[test]
fn verify_offer_voucher_gate() {
use crate::broker::node_identity::NodeKey;
static LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
let _g = LOCK.lock().unwrap();
let state = std::sync::Arc::new(BrokerState::new());
let dash = NodeKey::generate();
*state.dash_voucher_pubkey.write().unwrap() = Some(dash.public_b64());
let signed = r#"{"v":1,"requester_user":"u","budget_credits":1000000.0,"task_nonce":"n","exp":9999999999}"#;
let sig = dash.sign(signed.as_bytes());
let mut offer = task_board::TaskOffer {
voucher_signed_json: signed.to_string(),
voucher_sig: sig,
..mk_offer("t-v")
};
std::env::remove_var("ZAKURO_REQUIRE_VOUCHER");
assert!(super::verify_offer_voucher(&state, &offer).is_ok());
std::env::set_var("ZAKURO_REQUIRE_VOUCHER", "1");
assert!(super::verify_offer_voucher(&state, &offer).is_ok());
let bare = mk_offer("t-bare");
assert!(super::verify_offer_voucher(&state, &bare).is_err());
offer.voucher_signed_json = offer.voucher_signed_json.replace("1000000.0", "0.0000001");
assert!(super::verify_offer_voucher(&state, &offer).is_err());
std::env::remove_var("ZAKURO_REQUIRE_VOUCHER");
}
#[test]
fn take_committed_removes_once() {
let state = BrokerState::new();
state.pending_offers.insert(
"t1".into(),
PendingOffer {
offer: mk_offer("t1"),
worker_id: "w1".into(),
created: std::time::Instant::now(),
},
);
assert!(take_committed(&state, "t1").is_some());
assert!(
take_committed(&state, "t1").is_none(),
"second commit → None"
);
}
#[test]
fn cancel_is_idempotent() {
let state = BrokerState::new();
state.pending_offers.insert(
"t1".into(),
PendingOffer {
offer: mk_offer("t1"),
worker_id: "w1".into(),
created: std::time::Instant::now(),
},
);
cancel_offer(&state, "t1");
assert!(!state.pending_offers.contains_key("t1"));
cancel_offer(&state, "t1"); cancel_offer(&state, "never"); }
#[test]
fn sweep_releases_stale_keeps_fresh() {
let state = BrokerState::new();
state.pending_offers.insert(
"t2".into(),
PendingOffer {
offer: mk_offer("t2"),
worker_id: "w".into(),
created: std::time::Instant::now(),
},
);
sweep_expired_offers(&state, std::time::Duration::from_secs(3600)); assert!(state.pending_offers.contains_key("t2"));
sweep_expired_offers(&state, std::time::Duration::ZERO); assert!(!state.pending_offers.contains_key("t2"));
}
}
#[cfg(test)]
mod forward_async_tests {
use super::forward_to_worker_async;
#[tokio::test]
async fn forward_async_round_trips_status_body_and_headers() {
let server = std::sync::Arc::new(tiny_http::Server::http("127.0.0.1:0").unwrap());
let port = server.server_addr().to_ip().unwrap().port();
std::thread::spawn(move || {
for req in server.incoming_requests() {
let r = tiny_http::Response::from_data(b"pong".to_vec())
.with_status_code(200)
.with_header(tiny_http::Header::from_bytes("X-Zakuro-Worker", "w1").unwrap())
.with_header(tiny_http::Header::from_bytes("Connection", "close").unwrap());
let _ = req.respond(r);
}
});
let client = reqwest::Client::new();
let uri = format!("http://127.0.0.1:{}/execute", port);
let reply = forward_to_worker_async(&client, &uri, b"ping".to_vec(), "rid-1", 0.0)
.await
.expect("forward ok");
assert_eq!(reply.status, 200);
assert_eq!(reply.body, b"pong");
assert!(reply
.headers
.iter()
.any(|(k, v)| k == "x-zakuro-worker" && v == "w1"));
}
}
#[cfg(test)]
mod brokers_endpoint_tests {
use super::handle_brokers;
use crate::broker::node_identity::NodeKey;
use crate::broker::{BrokerState, PeerManager};
use std::sync::Arc;
fn spawn_fake_peer_health(broker_id: &str) -> u16 {
let server = Arc::new(tiny_http::Server::http("127.0.0.1:0").unwrap());
let port = server.server_addr().to_ip().unwrap().port();
let broker_id = broker_id.to_string();
std::thread::spawn(move || {
for req in server.incoming_requests() {
let body = format!(r#"{{"status":"healthy","broker_id":"{}"}}"#, broker_id);
let resp = tiny_http::Response::from_data(body.into_bytes())
.with_status_code(200)
.with_header(
tiny_http::Header::from_bytes("Content-Type", "application/json").unwrap(),
);
let _ = req.respond(resp);
}
});
port
}
fn state_with_peers(peer_urls: &[String]) -> Arc<BrokerState> {
let mut state = BrokerState::new();
let key = Arc::new(NodeKey::generate());
state.peer_manager = PeerManager::new_with_node_key(
Some("127.0.0.2"),
peer_urls,
state.config.port,
String::new(), true,
Some(key.clone()),
);
state.node_key = key;
Arc::new(state)
}
#[test]
fn self_id_is_key_derived_and_matches_node_key() {
let state = state_with_peers(&[]);
let resp = handle_brokers(&state);
let body: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
assert_eq!(body["self"], state.node_key.node_uri());
assert!(body["self"].as_str().unwrap().starts_with("zc://node-"));
assert!(body["brokers"].as_array().unwrap().is_empty());
}
#[test]
fn known_peer_fingerprint_appears_as_zc_id_no_ip_leaked() {
let fp_id = "zc://node-deadbeefcafebabe";
let port = spawn_fake_peer_health(fp_id);
let peer_url = format!("http://127.0.0.1:{}", port);
let state = state_with_peers(std::slice::from_ref(&peer_url));
let resp = handle_brokers(&state);
let raw = String::from_utf8(resp.body.clone()).unwrap();
assert!(!raw.contains("127.0.0.1"), "leaked loopback IP: {raw}");
assert!(!raw.contains(&port.to_string()), "leaked peer port: {raw}");
assert!(!raw.contains(&peer_url), "leaked raw peer_url: {raw}");
assert!(!regex_ish_ip_shaped(&raw), "body looks IP-shaped: {raw}");
let body: serde_json::Value = serde_json::from_str(&raw).unwrap();
let brokers = body["brokers"].as_array().unwrap();
assert_eq!(brokers.len(), 1);
assert_eq!(brokers[0]["id"], fp_id);
assert_eq!(brokers[0]["reachable"], true);
}
#[test]
fn one_broker_registered_under_two_spellings_is_listed_once() {
let fp_id = "zc://node-deadbeefcafebabe";
let port_a = spawn_fake_peer_health(fp_id);
let port_b = spawn_fake_peer_health(fp_id);
let url_a = format!("http://127.0.0.1:{port_a}");
let url_b = format!("http://127.0.0.1:{port_b}");
let state = state_with_peers(std::slice::from_ref(&url_a));
state.peer_manager.register_peer(url_b);
assert_eq!(
state.peer_manager.peer_count(),
2,
"test setup: two spellings"
);
let resp = handle_brokers(&state);
let body: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
let brokers = body["brokers"].as_array().unwrap();
assert_eq!(
brokers.len(),
1,
"one broker listed twice under two spellings: {body}"
);
assert_eq!(brokers[0]["id"], fp_id);
assert_eq!(brokers[0]["reachable"], true);
}
#[test]
fn distinct_brokers_are_each_listed() {
let port_a = spawn_fake_peer_health("zc://node-aaaaaaaaaaaaaaaa");
let port_b = spawn_fake_peer_health("zc://node-bbbbbbbbbbbbbbbb");
let state = state_with_peers(&[format!("http://127.0.0.1:{port_a}")]);
state
.peer_manager
.register_peer(format!("http://127.0.0.1:{port_b}"));
assert_eq!(state.peer_manager.peer_count(), 2, "test setup: two peers");
let resp = handle_brokers(&state);
let body: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
let mut ids: Vec<&str> = body["brokers"]
.as_array()
.unwrap()
.iter()
.map(|b| b["id"].as_str().unwrap())
.collect();
ids.sort();
assert_eq!(
ids,
vec!["zc://node-aaaaaaaaaaaaaaaa", "zc://node-bbbbbbbbbbbbbbbb"]
);
}
#[test]
fn unknown_peer_fingerprint_is_omitted_not_leaked_as_ip() {
let dead_peer = "http://127.0.0.1:1".to_string(); let state = state_with_peers(std::slice::from_ref(&dead_peer));
let resp = handle_brokers(&state);
let raw = String::from_utf8(resp.body.clone()).unwrap();
assert!(!raw.contains(&dead_peer), "leaked raw peer_url: {raw}");
let body: serde_json::Value = serde_json::from_str(&raw).unwrap();
assert!(
body["brokers"].as_array().unwrap().is_empty(),
"unreachable/unknown peer must be omitted: {raw}"
);
}
fn regex_ish_ip_shaped(s: &str) -> bool {
for window in s.split(|c: char| !c.is_ascii_digit() && c != '.') {
let parts: Vec<&str> = window.split('.').collect();
if parts.len() == 4
&& parts
.iter()
.all(|p| !p.is_empty() && p.parse::<u8>().is_ok())
{
return true;
}
}
false
}
#[test]
fn peers_endpoint_is_ip_free_and_key_derived() {
use tiny_http::Method;
let fp_id = "zc://node-feedfacecafebabe";
let port = spawn_fake_peer_health(fp_id);
let peer_url = format!("http://127.0.0.1:{}", port);
let state = state_with_peers(std::slice::from_ref(&peer_url));
let self_uri = state.node_key.node_uri();
let req = crate::broker::http_adapter::BrokerRequest {
method: Method::Get,
path: "/peers".to_string(),
query: String::new(),
headers: vec![],
body: vec![],
peer_addr: None, };
let resp = super::handle(state, &req, false);
let raw = String::from_utf8(resp.body.clone()).unwrap();
assert!(!raw.contains("10.13.13."), "leaked mesh IP: {raw}");
assert!(!raw.contains("127.0.0.1"), "leaked loopback IP: {raw}");
assert!(!raw.contains(&port.to_string()), "leaked peer port: {raw}");
assert!(!raw.contains(&peer_url), "leaked raw peer_url: {raw}");
assert!(!raw.starts_with("http://"), "leaked raw http url: {raw}");
assert!(
!raw.contains("http://"),
"leaked raw http url anywhere: {raw}"
);
assert!(!regex_ish_ip_shaped(&raw), "body looks IP-shaped: {raw}");
let body: serde_json::Value = serde_json::from_str(&raw).unwrap();
assert_eq!(body["node"], self_uri);
assert!(body["node"].as_str().unwrap().starts_with("zc://node-"));
let peers = body["peers"].as_array().unwrap();
assert_eq!(peers.len(), 1);
assert_eq!(peers[0], fp_id);
assert_eq!(body["total"], 1);
}
}
#[cfg(test)]
mod marketplace_preselection_guard_tests {
use super::marketplace_select_enabled;
static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
#[test]
fn marketplace_select_escape_hatch() {
let _g = ENV_LOCK.lock().unwrap();
std::env::remove_var("ZAKURO_MARKETPLACE_SELECT");
assert!(marketplace_select_enabled(), "default (unset) must be ON");
std::env::set_var("ZAKURO_MARKETPLACE_SELECT", "0");
assert!(
!marketplace_select_enabled(),
"\"0\" must disable pre-selection"
);
std::env::set_var("ZAKURO_MARKETPLACE_SELECT", "1");
assert!(
marketplace_select_enabled(),
"any non-\"0\" value must be ON"
);
std::env::remove_var("ZAKURO_MARKETPLACE_SELECT");
}
}
#[cfg(test)]
mod executor_identity_tests {
use crate::broker::peer::PeerManager;
use crate::broker::worker::Worker;
#[test]
fn executor_fp_prefers_node_identity_and_refuses_unstamped_workers() {
let mut w = Worker::new(
"w1".into(),
"worker-a".into(),
"http://10.13.13.7:3960".into(),
"zakuro".into(),
);
assert_eq!(super::executor_fp(&w), None);
w.node_fp = "aaaaaaaaaaaaaaaa".into();
assert_eq!(super::executor_fp(&w), Some("aaaaaaaaaaaaaaaa"));
}
#[test]
fn proxied_executor_resolves_by_fingerprint_where_ip_match_fails() {
let pm = PeerManager::new(
Some("10.13.13.2"),
&["127.0.0.1:19001".to_string()],
9000,
"k".to_string(),
true,
);
pm.set_peer_fingerprint("http://127.0.0.1:19001", "cccccccccccccccc");
let mut w = Worker::new(
"w2".into(),
"worker-proxied".into(),
"http://10.13.13.7:3960".into(),
"zakuro".into(),
);
w.node_fp = "cccccccccccccccc".into();
assert_eq!(
pm.get_url_for_ip("10.13.13.7"),
None,
"precondition: the old IP-keyed lookup must miss for a proxied peer"
);
assert_eq!(
super::executor_fp(&w)
.and_then(|fp| pm.get_url_for_fingerprint(fp))
.as_deref(),
Some("http://127.0.0.1:19001")
);
}
#[test]
fn direct_but_unstamped_executor_now_routes_to_full_refund_arm() {
let pm = PeerManager::new(
Some("10.13.13.2"),
&["10.13.13.7:9000".to_string()],
9000,
"k".to_string(),
true,
);
pm.set_peer_fingerprint("http://10.13.13.7:9000", "dddddddddddddddd");
let w = Worker::new(
"w3".into(),
"worker-direct".into(),
"http://10.13.13.7:3960".into(),
"zakuro".into(),
);
assert_eq!(
pm.get_url_for_ip("10.13.13.7").as_deref(),
Some("http://10.13.13.7:9000")
);
assert_eq!(super::executor_fp(&w), None);
assert_eq!(
super::executor_fp(&w).and_then(|fp| pm.get_url_for_fingerprint(fp)),
None
);
}
}
#[cfg(test)]
mod task_result_producer_tests {
use super::*;
use crate::broker::worker::Worker;
use crate::broker::{BrokerConfig, BrokerState};
use std::io::{Read, Write};
use std::net::TcpListener;
use std::sync::atomic::{AtomicU16, Ordering};
use std::sync::Arc;
fn reserve_port() -> TcpListener {
static NEXT: AtomicU16 = AtomicU16::new(50100);
loop {
let port = NEXT.fetch_add(1, Ordering::Relaxed);
assert!(port < 51000, "reserve_port exhausted the test port range");
if let Ok(l) = TcpListener::bind(("127.0.0.1", port)) {
return l;
}
}
}
fn spawn_stub_worker(listener: TcpListener) {
std::thread::spawn(move || {
for stream in listener.incoming() {
let mut stream = match stream {
Ok(s) => s,
Err(_) => continue,
};
let mut buf = [0u8; 2048];
let _ = stream.read(&mut buf);
let body = "ok";
let resp = format!(
"HTTP/1.1 200 OK\r\nContent-Length: {}\r\nContent-Type: application/octet-stream\r\nX-Zakuro-Pid: 4242\r\nX-Zakuro-IP: 10.13.13.7\r\n\r\n{}",
body.len(),
body
);
let _ = stream.write_all(resp.as_bytes());
let _ = stream.flush();
}
});
}
fn mk_offer(task_id: &str) -> task_board::TaskOffer {
serde_json::from_str(&format!(
r#"{{"task_id":"{}","payload_b64":"","max_price_per_hour":100.0,"estimated_duration_secs":1.0,"timeout_secs":10.0,"requester_user_id":"u","source_broker":"n"}}"#,
task_id
))
.unwrap()
}
fn produce_real_result() -> task_board::TaskResult {
let listener = reserve_port();
let port = listener.local_addr().unwrap().port();
spawn_stub_worker(listener);
let st = Arc::new(BrokerState::with_config(BrokerConfig::default()));
let mut w = Worker::new(
"wid".into(),
"worker-exec".into(),
format!("http://127.0.0.1:{}", port),
"zakuro".into(),
);
w.node_fp = "aaaaaaaaaaaaaaaa".into();
w.slot = "3960".into();
super::run_offer_on(&st, &mk_offer("t-privacy"), &w)
.expect("stub worker must answer so the producer builds a TaskResult")
}
#[test]
fn producer_emits_no_address_in_task_result() {
let r = produce_real_result();
let json = serde_json::to_string(&r).unwrap();
assert!(
!json.contains("http://"),
"producer put a scheme+host on the wire: {json}"
);
assert!(
!json.contains("worker_ip"),
"worker_ip must be off the wire: {json}"
);
assert!(
!contains_dotted_quad(&json),
"producer put an IP address on the wire: {json}"
);
assert_eq!(
r.worker_uri, "zc://worker-aaaaaaaaaaaaaaaa-3960",
"producer must emit Worker::zc_uri(), not the dialable worker.uri"
);
assert_eq!(r.executor_node_uri, st_node_uri(&r));
assert!(
r.executor_node_uri.starts_with("zc://node-"),
"executor_node_uri must be the key-derived node identity: {:?}",
r.executor_node_uri
);
}
fn produce_real_peer_identity() -> task_board::PeerIdentity {
use crate::broker::http_adapter::BrokerRequest;
use crate::broker::worker::WorkerRegistration;
let config = BrokerConfig {
owner_user_id: Some("1000000001".into()),
node_name: Some("node-a".into()),
..Default::default()
};
let st = Arc::new(BrokerState::with_config(config));
st.workers.register(WorkerRegistration {
name: "worker-a".into(),
uri: "http://127.0.0.1:39999".into(),
worker_type: "zakuro".into(),
resources: Default::default(),
pricing: Default::default(),
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,
});
let req = BrokerRequest {
method: tiny_http::Method::Get,
path: "/peer/identity".into(),
query: "".into(),
headers: vec![],
body: vec![],
peer_addr: None, };
let resp = super::handle_peer_identity(&st, &req);
assert_eq!(resp.status, 200, "handle_peer_identity must succeed");
serde_json::from_slice(&resp.body).expect("handle_peer_identity must emit valid JSON")
}
#[test]
fn producer_emits_no_address_in_peer_identity() {
let id = produce_real_peer_identity();
let json = serde_json::to_string(&id).unwrap();
assert!(
!json.contains("http://"),
"producer put a scheme+host on the wire: {json}"
);
assert!(
!contains_dotted_quad(&json),
"producer put an IP address on the wire: {json}"
);
assert_eq!(id.workers.len(), 1, "the registered worker must be listed");
}
fn st_node_uri(r: &task_board::TaskResult) -> String {
let node_id = r
.executor_node_id
.as_ref()
.expect("producer signs the receipt, so node_id is present");
format!(
"zc://node-{}",
crate::broker::node_identity::fingerprint_of_pubkey_b64(node_id).unwrap()
)
}
fn contains_dotted_quad(s: &str) -> bool {
let b: Vec<char> = s.chars().collect();
let mut i = 0;
while i < b.len() {
let mut j = i;
let mut groups = 0;
loop {
let start = j;
while j < b.len() && b[j].is_ascii_digit() {
j += 1;
}
if j == start {
break;
}
groups += 1;
if groups == 4 {
return true;
}
if j < b.len() && b[j] == '.' {
j += 1;
} else {
break;
}
}
i = if j > i { j } else { i + 1 };
}
false
}
#[test]
fn quic_producer_emits_identity_not_worker_uri() {
let src = include_str!("quic.rs");
assert!(
src.contains("let wuri = worker.zc_uri();"),
"quic.rs producer must derive its mesh-facing handle from Worker::zc_uri()"
);
assert!(
!src.contains("let wuri = worker.uri"),
"quic.rs producer must not put the dialable worker.uri on the wire"
);
assert!(
!src.contains("X-Zakuro-IP"),
"quic.rs must not capture the worker IP header for the mesh payload"
);
}
#[test]
fn dotted_quad_detector_is_sound() {
assert!(contains_dotted_quad("http://10.13.13.7:3960"));
assert!(contains_dotted_quad("127.0.0.1"));
assert!(contains_dotted_quad(r#"{"worker_ip":"192.168.0.23"}"#));
assert!(!contains_dotted_quad("zc://worker-aaaaaaaaaaaaaaaa-3960"));
assert!(!contains_dotted_quad("zc://node-13a9551e881f7f8c"));
assert!(!contains_dotted_quad("1234.5")); assert!(!contains_dotted_quad("0.789")); }
}
#[cfg(test)]
mod commit_credits_arm_tests {
use super::*;
use crate::broker::peer::PeerManager;
use crate::broker::worker::Worker;
use crate::broker::{BrokerConfig, BrokerState};
use std::sync::Arc;
const PEER_ADDR: &str = "127.0.0.1:19099";
const PEER_URL: &str = "http://127.0.0.1:19099";
const PEER_FP: &str = "cccccccccccccccc";
fn state_with_seeded_peer() -> Arc<BrokerState> {
state_with_seeded_peer_api("http://127.0.0.1:1")
}
fn state_with_seeded_peer_api(api_url: &str) -> Arc<BrokerState> {
let config = BrokerConfig {
api_url: Some(api_url.into()),
api_key: Some("k".into()),
..Default::default()
};
let mut st = BrokerState::with_config(config);
st.peer_manager = PeerManager::new(
Some("10.13.13.2"),
&[PEER_ADDR.to_string()],
9000,
"k".to_string(),
true,
);
st.peer_manager.set_peer_fingerprint(PEER_URL, PEER_FP);
let st = Arc::new(st);
st.ledger
.authoritative_balances
.insert("u1".to_string(), 100.0);
st.ledger.local_credits.insert("u1".to_string(), 100.0);
assert!(st.is_billing_enabled());
st
}
fn worker(node_fp: Option<&str>) -> Worker {
let mut w = Worker::new(
"wid".into(),
"worker-exec".into(),
"http://127.0.0.1:3960".into(),
"zakuro".into(),
);
if let Some(fp) = node_fp {
w.node_fp = fp.into();
}
w
}
fn run_commit(st: &Arc<BrokerState>, w: &Worker) {
run_commit_with_nonce(st, w, "");
}
fn run_commit_with_nonce(st: &Arc<BrokerState>, w: &Worker, nonce: &str) {
let (reservation_id, _) = st
.ledger
.local_reserve("u1", 10.0, "ref")
.expect("reserve must succeed against the seeded balance");
super::commit_credits(
st,
&Authority::Local,
false, &reservation_id,
4.0, 100.0, "req-1",
"u1",
w,
1000.0, 0.10, nonce, );
}
fn reserve_port() -> std::net::TcpListener {
use std::sync::atomic::{AtomicU16, Ordering};
static NEXT: AtomicU16 = AtomicU16::new(51100);
loop {
let port = NEXT.fetch_add(1, Ordering::Relaxed);
assert!(port < 52000, "reserve_port exhausted the test port range");
if let Ok(l) = std::net::TcpListener::bind(("127.0.0.1", port)) {
return l;
}
}
}
fn spawn_redeem_stub() -> (Arc<std::sync::atomic::AtomicUsize>, String) {
use std::io::{Read, Write};
use std::sync::atomic::{AtomicUsize, Ordering};
let listener = reserve_port();
let url = format!("http://127.0.0.1:{}", listener.local_addr().unwrap().port());
let hits = Arc::new(AtomicUsize::new(0));
let hits_srv = Arc::clone(&hits);
std::thread::spawn(move || {
for stream in listener.incoming() {
let mut stream = match stream {
Ok(s) => s,
Err(_) => continue,
};
let mut buf = [0u8; 4096];
let n = stream.read(&mut buf).unwrap_or(0);
let req = String::from_utf8_lossy(&buf[..n]);
if req.contains("/api/broker/voucher/redeem") {
hits_srv.fetch_add(1, Ordering::SeqCst);
}
let body = r#"{"redeemed":true}"#;
let resp = format!(
"HTTP/1.1 200 OK\r\nContent-Length: {}\r\nContent-Type: application/json\r\nConnection: close\r\n\r\n{}",
body.len(),
body
);
let _ = stream.write_all(resp.as_bytes());
let _ = stream.flush();
}
});
(hits, url)
}
#[test]
fn some_arm_redeems_the_voucher_exactly_once() {
use std::sync::atomic::Ordering;
let (hits, api_url) = spawn_redeem_stub();
let st = state_with_seeded_peer_api(&api_url);
run_commit_with_nonce(&st, &worker(Some(PEER_FP)), "nonce-abc");
assert_eq!(
hits.load(Ordering::SeqCst),
1,
"the resolvable-executor arm must redeem the voucher exactly once \
(all three money legs fire together)"
);
}
#[test]
fn none_arm_redeems_nothing() {
use std::sync::atomic::Ordering;
let (hits, api_url) = spawn_redeem_stub();
let st = state_with_seeded_peer_api(&api_url);
run_commit_with_nonce(&st, &worker(None), "nonce-abc");
assert_eq!(
hits.load(Ordering::SeqCst),
0,
"the unresolvable-executor arm must redeem nothing"
);
assert_eq!(bal(&st, "u1"), Some(100.0));
}
fn bal(st: &Arc<BrokerState>, uid: &str) -> Option<f64> {
st.ledger.authoritative_balances.get(uid).map(|v| *v)
}
#[test]
fn unstamped_worker_takes_full_refund_arm_via_commit_credits() {
let st = state_with_seeded_peer();
run_commit(&st, &worker(None));
assert_eq!(
bal(&st, "u1"),
Some(100.0),
"unstamped executor must leave the requester whole (full refund)"
);
assert_eq!(
bal(&st, "platform-house"),
None,
"None arm credited the platform commission — a money leg escaped the Some arm"
);
assert_eq!(
st.ledger.local_credits.get("platform-house").map(|v| *v),
None
);
}
#[test]
fn stamped_resolvable_worker_takes_payout_arm_via_commit_credits() {
let st = state_with_seeded_peer();
run_commit(&st, &worker(Some(PEER_FP)));
assert_eq!(
bal(&st, "platform-house"),
Some(0.4),
"stamped, resolvable executor must credit the platform commission"
);
let got = bal(&st, "u1").unwrap();
assert!(
(got - 99.6).abs() < 1e-9,
"requester should net -commission on the payout arm; got {got}"
);
}
}
#[cfg(test)]
mod billing_price_pin_tests {
use super::{execute_offer_locally, reserve_offer};
use crate::broker::worker::test_registration;
use crate::broker::{task_board, BrokerState};
pub(super) fn slow_mock_worker(delay_ms: u64) -> u16 {
let (server, port) = crate::integration_tests::integration::bind_loopback_http();
std::thread::spawn(move || {
for req in server.incoming_requests() {
std::thread::sleep(std::time::Duration::from_millis(delay_ms));
let _ = req.respond(tiny_http::Response::from_data(b"ok".to_vec()));
}
});
port
}
pub(super) fn state_with_local_worker(name: &str, port: u16, price: f64) -> BrokerState {
let state = BrokerState::new();
state.workers.register(test_registration(
name,
format!("http://127.0.0.1:{port}"),
price,
));
state
}
pub(super) fn pin_offer(task_id: &str) -> task_board::TaskOffer {
serde_json::from_str(&format!(
r#"{{"task_id":"{task_id}","payload_b64":"","max_price_per_hour":1000000000.0,"estimated_duration_secs":10.0,"timeout_secs":10.0,"requester_user_id":"u","source_broker":"n"}}"#
))
.unwrap()
}
#[test]
fn billing_reads_the_worker_price_not_the_broker_price() {
let port = slow_mock_worker(200);
let state = state_with_local_worker("pin-w0", port, 3600.0);
state.price.set(0.0);
let accept = reserve_offer(&state, &pin_offer("pin-reserve")).unwrap();
assert!(
(accept.estimated_cost - 10.0).abs() < 1e-9,
"estimate must use the worker price (10 s at 1 cr/s), got {}",
accept.estimated_cost
);
assert_eq!(accept.price_per_hour, 3600.0);
let result = execute_offer_locally(&state, &pin_offer("pin-exec")).unwrap();
assert!(
result.actual_cost >= 0.15,
"a ~200 ms job at 1 cr/s must cost >= 0.15, got {} (BrokerPrice would give 0.001)",
result.actual_cost
);
assert_eq!(result.price_per_hour, 3600.0);
}
fn mock_sync_hub(reply: &'static str) -> (String, std::sync::mpsc::Receiver<String>) {
let (server, port) = crate::integration_tests::integration::bind_loopback_http();
let (tx, rx) = std::sync::mpsc::channel();
std::thread::spawn(move || {
for mut req in server.incoming_requests() {
let mut body = String::new();
let _ = std::io::Read::read_to_string(req.as_reader(), &mut body);
let _ = tx.send(body);
let _ = req.respond(tiny_http::Response::from_string(reply).with_header(
tiny_http::Header::from_bytes("Content-Type", "application/json").unwrap(),
));
}
});
(format!("http://127.0.0.1:{port}"), rx)
}
#[test]
fn hub_price_from_sync_reply_is_what_a_job_is_billed() {
use crate::broker::ledger::Ledger;
let port = slow_mock_worker(200);
let state = state_with_local_worker("pin-w1", port, 3600.0);
let (hub, bodies) = mock_sync_hub(
r#"{"success":true,"synced":1,"failed":0,"errors":[],"prices":{"pin-w1":7200.0}}"#,
);
let sync = || {
Ledger::sync_workers_via_api(
"u1",
&state.workers.list(),
&hub,
"zk_u1_test",
Some("n"),
None,
Some("pk"),
)
.unwrap()
};
let outcome = sync();
assert_eq!(outcome.prices.get("pin-w1"), Some(&7200.0));
assert_eq!(state.workers.apply_hub_prices(&outcome.prices), 1);
let sent: serde_json::Value = serde_json::from_str(&bodies.recv().unwrap()).unwrap();
assert_eq!(
sent[0]["price_per_hour"], 3600.0,
"first sync reports the worker's own price"
);
let result = execute_offer_locally(&state, &pin_offer("hub-exec")).unwrap();
assert_eq!(result.price_per_hour, 7200.0);
assert!(
result.actual_cost >= 0.3,
"~200 ms at 2 cr/s must cost >= 0.3, got {}",
result.actual_cost
);
sync();
let sent: serde_json::Value = serde_json::from_str(&bodies.recv().unwrap()).unwrap();
assert_eq!(
sent[0]["price_per_hour"], 3600.0,
"sync keeps reporting the non-override price"
);
}
}
#[cfg(test)]
mod drain_route_tests {
use super::*;
use crate::broker::worker::test_registration;
use tiny_http::Method;
fn post(path: &str, peer: Option<std::net::SocketAddr>) -> BrokerRequest {
BrokerRequest {
method: Method::Post,
path: path.to_string(),
query: String::new(),
headers: vec![],
body: vec![],
peer_addr: peer,
}
}
#[test]
fn drain_is_loopback_only_and_marks_the_worker_draining() {
let state = Arc::new(BrokerState::new());
let w = state.workers.register(test_registration(
"fp-w0",
"http://127.0.0.1:3960".to_string(),
3.6,
));
let path = format!("/workers/{}/drain", w.id);
let remote: std::net::SocketAddr = "10.13.13.9:5555".parse().unwrap();
assert_eq!(
handle(state.clone(), &post(&path, Some(remote)), false).status,
403
);
assert_eq!(
state.workers.get(&w.id).unwrap().status,
WorkerStatus::Healthy
);
let ok = handle(state.clone(), &post(&path, None), false);
assert_eq!(ok.status, 200, "{}", String::from_utf8_lossy(&ok.body));
let body: serde_json::Value = serde_json::from_slice(&ok.body).unwrap();
assert_eq!(body["status"], "draining");
assert_eq!(
state.workers.get(&w.id).unwrap().status,
WorkerStatus::Draining
);
let list = handle(
state.clone(),
&BrokerRequest {
method: Method::Get,
..post("/workers", None)
},
false,
);
let list: serde_json::Value = serde_json::from_slice(&list.body).unwrap();
assert_eq!(list["workers"][0]["status"], "draining");
assert_eq!(
handle(state, &post("/workers/nope/drain", None), false).status,
404
);
}
}
#[cfg(test)]
mod p2p_active_requests_tests {
use super::billing_price_pin_tests::{pin_offer, slow_mock_worker, state_with_local_worker};
use super::*;
#[test]
fn a_slow_p2p_offer_holds_active_requests_at_one_until_it_completes() {
let port = slow_mock_worker(200);
let state = Arc::new(state_with_local_worker("fp-w0", port, 3.6));
let worker_id = state.workers.list()[0].id.clone();
let state_for_thread = state.clone();
let handle_thread = std::thread::spawn(move || {
execute_offer_locally(&state_for_thread, &pin_offer("p2p-active"))
});
std::thread::sleep(std::time::Duration::from_millis(50));
let mid = state.workers.get(&worker_id).unwrap();
assert_eq!(
mid.active_requests, 1,
"a P2P offer in flight must be counted, exactly like /execute"
);
handle_thread.join().unwrap().unwrap();
let after = state.workers.get(&worker_id).unwrap();
assert_eq!(
after.active_requests, 0,
"released once the offer completes"
);
}
}
#[cfg(test)]
mod reservation_activity_tests {
use super::billing_price_pin_tests::{pin_offer, slow_mock_worker, state_with_local_worker};
use super::*;
use tiny_http::Method;
fn listed_active_requests(state: &Arc<BrokerState>) -> u64 {
let resp = handle(
state.clone(),
&BrokerRequest {
method: Method::Get,
path: "/workers".to_string(),
query: String::new(),
headers: vec![],
body: vec![],
peer_addr: None,
},
false,
);
let list: serde_json::Value = serde_json::from_slice(&resp.body).unwrap();
list["workers"][0]["active_requests"].as_u64().unwrap()
}
#[test]
fn a_reserved_offer_counts_as_activity_until_commit_cancel_or_expiry() {
let port = slow_mock_worker(0);
let state = Arc::new(state_with_local_worker("fp-w0", port, 3.6));
assert_eq!(listed_active_requests(&state), 0);
reserve_offer(&state, &pin_offer("r-commit")).unwrap();
assert_eq!(
listed_active_requests(&state),
1,
"reserved, not yet committed"
);
let pending = take_committed(&state, "r-commit").unwrap();
let worker = state.workers.get(&pending.worker_id).unwrap();
run_offer_on(&state, &pending.offer, &worker).unwrap();
assert_eq!(
listed_active_requests(&state),
0,
"released once the commit ran"
);
reserve_offer(&state, &pin_offer("r-cancel")).unwrap();
assert_eq!(listed_active_requests(&state), 1);
cancel_offer(&state, "r-cancel");
assert_eq!(listed_active_requests(&state), 0, "released by a cancel");
reserve_offer(&state, &pin_offer("r-expire")).unwrap();
assert_eq!(listed_active_requests(&state), 1);
std::thread::sleep(std::time::Duration::from_millis(5));
sweep_expired_offers(&state, std::time::Duration::ZERO);
assert_eq!(
listed_active_requests(&state),
0,
"released once the reservation expired"
);
}
}
#[cfg(test)]
mod earn_prefix_contract_tests {
use super::*;
use crate::broker::flush::dashboard_type;
use crate::broker::worker::test_registration;
use tiny_http::Method;
fn owner_state(owner: Option<&str>, api_configured: bool) -> Arc<BrokerState> {
Arc::new(BrokerState::with_config(crate::broker::BrokerConfig {
owner_user_id: owner.map(|s| s.to_string()),
api_url: api_configured.then(|| "http://127.0.0.1:1".to_string()),
api_key: api_configured.then(|| "k".to_string()),
peer_key: None,
enable_p2p: false,
..crate::broker::BrokerConfig::default()
}))
}
#[test]
fn peer_earn_books_an_earn_prefixed_credit_purchase_for_the_owner() {
let state = Arc::new(BrokerState::with_config(crate::broker::BrokerConfig {
peer_key: Some("pk".to_string()),
enable_p2p: true,
owner_user_id: Some("owner".to_string()),
api_url: None,
api_key: None,
discovery: crate::broker::DiscoveryConfig {
peers: vec!["127.0.0.1:1".to_string()],
..crate::broker::DiscoveryConfig::default()
},
..crate::broker::BrokerConfig::default()
}));
let body = serde_json::json!({
"amount": 1.5,
"duration_ms": 10.0,
"worker_id": "a1b2c3d4-w0",
"requesting_user": "someone",
"request_id": "r-42"
});
let req = BrokerRequest {
method: Method::Post,
path: "/peer/earn".to_string(),
query: String::new(),
headers: vec![("X-Peer-Key".to_string(), "pk".to_string())],
body: body.to_string().into_bytes(),
peer_addr: None,
};
let resp = handle(state.clone(), &req, false);
assert_eq!(resp.status, 200, "{}", String::from_utf8_lossy(&resp.body));
let queued = state.tx_buffer.pending_snapshot();
let tx = queued
.iter()
.find(|t| t.request_id == "earn-r-42")
.expect("the earn row must be queued under the earn- prefix");
assert_eq!(tx.user_id, "owner");
assert_eq!(dashboard_type(&tx.tx_type), "credit_purchase");
}
#[test]
fn dashboard_types_are_stable() {
assert_eq!(dashboard_type("commit"), "job_execution");
assert_eq!(dashboard_type("cancel"), "job_execution");
assert_eq!(dashboard_type("credit"), "credit_purchase");
assert_eq!(dashboard_type("serve"), "serve");
}
#[test]
fn settle_own_provider_payout_books_an_earn_prefixed_credit_for_the_owner() {
let state = owner_state(Some("owner"), false);
let worker = state.workers.register(test_registration(
"w-own",
"http://127.0.0.1:1".to_string(),
3600.0,
));
settle_own_provider_payout(&state, "requester", &worker, 10.0, 0.0, "", "r-own-1", 5.0);
let queued = state.tx_buffer.pending_snapshot();
let tx = queued
.iter()
.find(|t| t.request_id == "earn-r-own-1")
.expect("the own-provider payout must be queued under the earn- prefix");
assert_eq!(tx.user_id, "owner");
assert_eq!(tx.tx_type, "credit");
assert_eq!(dashboard_type(&tx.tx_type), "credit_purchase");
}
#[test]
fn settle_self_payout_books_an_earn_prefixed_credit_for_the_owner() {
let state = owner_state(Some("owner"), true);
let worker = state.workers.register(test_registration(
"w-self",
"http://127.0.0.1:1".to_string(),
3600.0,
));
settle_self_payout(&state, "requester", 10.0, 0.0, &worker, "r-self-1", 5.0);
let queued = state.tx_buffer.pending_snapshot();
let tx = queued
.iter()
.find(|t| t.request_id == "earn-r-self-1")
.expect("the self-executed payout must be queued under the earn- prefix");
assert_eq!(tx.user_id, "owner");
assert_eq!(tx.tx_type, "credit");
assert_eq!(dashboard_type(&tx.tx_type), "credit_purchase");
}
#[test]
fn settle_own_provider_payout_refund_id_does_not_start_with_earn_prefix() {
let state = owner_state(None, true);
let worker = state.workers.register(test_registration(
"w-refund",
"http://127.0.0.1:1".to_string(),
3600.0,
));
settle_own_provider_payout(
&state,
"requester",
&worker,
10.0,
0.0,
"",
"r-refund-1",
5.0,
);
let queued = state.tx_buffer.pending_snapshot();
let tx = queued
.iter()
.find(|t| t.user_id == "requester")
.expect("the no-owner refund must be queued");
assert_eq!(tx.request_id, "r-refund-1-earn-refund");
assert!(
!tx.request_id.starts_with("earn-"),
"a refund id must never start with `earn-`: {}",
tx.request_id
);
assert_eq!(tx.tx_type, "refund");
}
#[test]
fn settle_executor_payout_refund_id_does_not_start_with_earn_prefix() {
let state = owner_state(None, true);
settle_executor_payout(
&state,
"requester",
10.0,
0.0,
"http://127.0.0.1:2",
"unreachable-worker",
"r-refund-2",
5.0,
);
let queued = state.tx_buffer.pending_snapshot();
let tx = queued
.iter()
.find(|t| t.user_id == "requester")
.expect("the unreachable-executor refund must be queued");
assert_eq!(tx.request_id, "r-refund-2-earn-refund");
assert!(
!tx.request_id.starts_with("earn-"),
"a refund id must never start with `earn-`: {}",
tx.request_id
);
assert_eq!(tx.tx_type, "refund");
}
#[test]
fn settle_self_payout_refund_id_does_not_start_with_earn_prefix() {
let state = owner_state(None, true);
let worker = state.workers.register(test_registration(
"w-self-refund",
"http://127.0.0.1:1".to_string(),
3600.0,
));
settle_self_payout(&state, "requester", 10.0, 0.0, &worker, "r-refund-3", 5.0);
let queued = state.tx_buffer.pending_snapshot();
let tx = queued
.iter()
.find(|t| t.user_id == "requester")
.expect("the self-executed no-owner refund must be queued");
assert_eq!(tx.request_id, "r-refund-3-earn-refund");
assert!(
!tx.request_id.starts_with("earn-"),
"a refund id must never start with `earn-`: {}",
tx.request_id
);
assert_eq!(tx.tx_type, "refund");
}
}