use std::io::Read;
use std::net::SocketAddr;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
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, Request, Response, Server};
use super::{
discovery::{detect_discovery_mode, Discovery, DiscoveryMode},
flush::BufferedTransaction,
ledger,
peer::{self, Authority},
recovery,
router::{ResourceRequirements, RoutingDecision},
stats::{StatsResponse, TaskOfferRecord, TransactionRecord, TransactionStatus, WorkerStats},
task_board,
tui,
wal::{WalEntry, WalStatus},
worker::{self, Worker, WorkerHeartbeat, WorkerRegistration, WorkerStatus},
BrokerConfig, BrokerState,
};
static REQUEST_COUNTER: AtomicU64 = AtomicU64::new(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()
}
}
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);
}
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 || state.config.tui_mode {
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 tailscale_ip: Option<String>,
pub is_docker: Option<bool>,
}
#[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) -> Response<std::io::Cursor<Vec<u8>>> {
let body = serde_json::to_vec(data).unwrap_or_default();
Response::from_data(body)
.with_status_code(status)
.with_header(Header::from_bytes("Content-Type", "application/json").unwrap())
}
fn error_response(error: &str, code: &str, status: u16) -> Response<std::io::Cursor<Vec<u8>>> {
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)
}
fn read_body(request: &mut Request) -> Vec<u8> {
let mut body = Vec::new();
let _ = request.as_reader().read_to_end(&mut body);
body
}
fn get_header(request: &Request, name: &str) -> Option<String> {
request.headers().iter()
.find(|h| h.field.as_str().as_str().eq_ignore_ascii_case(name))
.map(|h| h.value.to_string())
}
fn authenticate_bearer(request: &Request, state: &BrokerState) -> Result<String, Response<std::io::Cursor<Vec<u8>>>> {
match get_header(request, "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: &Request, state: &BrokerState) -> bool {
match &state.config.worker_key {
None => true, Some(expected) => {
get_header(request, "X-Worker-Key")
.map(|k| k == *expected)
.unwrap_or(false)
}
}
}
fn worker_to_info(w: &worker::Worker, registry: &worker::WorkerRegistry) -> 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);
WorkerInfo {
id: w.id.clone(),
name: w.name.clone(),
uri: w.uri.clone(),
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,
tailscale_ip: w.tailscale_ip.clone(),
is_docker: w.is_docker,
}
}
fn handle_request(state: Arc<BrokerState>, mut request: Request, verbose: bool) {
let path = request.url().to_string();
let method = request.method().clone();
let response = match (method, path.as_str()) {
(Method::Get, "/health") => {
let live_ts_ip = super::discovery::get_tailscale_ip()
.or_else(|| state.own_tailscale_ip.clone());
let ts_connected = live_ts_ip.as_ref().map(|ip| ip.starts_with("100.")).unwrap_or(false);
json_response(&serde_json::json!({
"status": "healthy",
"service": "zakuro-broker",
"tailscale_ip": live_ts_ip,
"tailscale_connected": ts_connected,
"node_name": state.config.node_name,
}), 200)
}
(Method::Get, path) if path == "/stats" || path.starts_with("/stats?") => {
match authenticate_bearer(&request, &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.uri.clone(),
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 tailscale_ip = state.own_tailscale_ip.clone();
let metrics = state.stats.metrics(
active,
workers.len(),
state.is_local_mode(),
state.is_billing_enabled(),
tailscale_ip.clone(),
);
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(),
tailscale_ip: tailscale_ip.clone(),
tailscale_connected: tailscale_ip.is_some(),
};
json_response(&stats_resp, 200)
}
}
}
(Method::Get, "/workers") => {
let infos: Vec<WorkerInfo> = state.workers.list().iter()
.map(|w| worker_to_info(w, &state.workers))
.collect();
json_response(&WorkerListResponse { total: infos.len(), workers: infos }, 200)
}
(Method::Post, "/workers") => {
if !verify_worker_key(&request, &state) {
error_response("Worker key required. Set X-Worker-Key header", "UNAUTHORIZED", 401)
} else {
let body = read_body(&mut request);
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 worker = state.workers.register(registration);
if let Some(ref owner_id) = state.config.owner_user_id {
let node_name = state.config.node_name.as_deref();
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,
&vec![worker.clone()],
api_url,
api_key,
node_name,
state.own_tailscale_ip.as_deref(),
) {
Ok(()) => 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(&request, &state) {
error_response("Worker key required. Set X-Worker-Key header", "UNAUTHORIZED", 401)
} else {
let body = read_body(&mut request);
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") => {
let body = read_body(&mut request);
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, "/peer/health") => {
handle_peer_health(&state, &request)
}
(Method::Get, "/peer/workers") => {
handle_peer_workers(&state, &request)
}
(Method::Post, "/peer/reserve") => {
handle_peer_reserve(state.clone(), &mut request)
}
(Method::Post, "/peer/commit") => {
handle_peer_commit(state.clone(), &mut request)
}
(Method::Post, "/peer/cancel") => {
handle_peer_cancel(state.clone(), &mut request)
}
(Method::Get, p) if p.starts_with("/peer/balance") => {
handle_peer_balance(&state, &request)
}
(Method::Post, "/peer/earn") => {
handle_peer_earn(state.clone(), &mut request)
}
(Method::Get, "/peer/identity") => {
handle_peer_identity(&state, &request)
}
(Method::Post, "/peer/tasks/offer") => {
handle_peer_task_offer(state.clone(), &mut request)
}
(Method::Post, "/execute") => {
handle_execute(state, &mut request, verbose)
}
(Method::Get, "/instances") => {
match authenticate_bearer(&request, &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(&request, &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(&request, &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) = get_header(&request, "X-Api-Key") {
let body = read_body(&mut request);
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(&request, &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::Delete, path) if path.starts_with("/workers/") => {
if !verify_worker_key(&request, &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(&request, &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) = get_header(&request, "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)
};
let _ = request.respond(response);
}
fn handle_peer_health(
state: &Arc<BrokerState>,
request: &Request,
) -> Response<std::io::Cursor<Vec<u8>>> {
let peer_key = state.peer_manager.peer_key();
if !peer_key.is_empty() && !peer::verify_peer_key(request, peer_key) {
return error_response("Invalid peer key", "UNAUTHORIZED", 401);
}
json_response(
&peer::PeerHealthResponse {
status: "healthy".to_string(),
broker_id: super::discovery::get_tailscale_ip()
.or_else(|| state.own_tailscale_ip.clone())
.unwrap_or_else(|| "local".to_string()),
},
200,
)
}
fn handle_peer_workers(
state: &Arc<BrokerState>,
request: &Request,
) -> Response<std::io::Cursor<Vec<u8>>> {
let peer_key = state.peer_manager.peer_key();
if !peer_key.is_empty() && !peer::verify_peer_key(request, peer_key) {
return error_response("Invalid peer key", "UNAUTHORIZED", 401);
}
let worker_infos: Vec<WorkerInfo> = state.workers.list().iter()
.filter(|w| w.uri.contains("127.0.0.1") || w.uri.contains("localhost"))
.map(|w| worker_to_info(w, &state.workers))
.collect();
json_response(&WorkerListResponse { total: worker_infos.len(), workers: worker_infos }, 200)
}
fn handle_peer_reserve(
state: Arc<BrokerState>,
request: &mut Request,
) -> Response<std::io::Cursor<Vec<u8>>> {
let peer_key = state.peer_manager.peer_key();
if !peer_key.is_empty() && !peer::verify_peer_key(request, peer_key) {
return error_response("Invalid peer key", "UNAUTHORIZED", 401);
}
let body = read_body(request);
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) => error_response(&e.to_string(), "LEDGER_ERROR", 500),
}
}
fn handle_peer_commit(
state: Arc<BrokerState>,
request: &mut Request,
) -> Response<std::io::Cursor<Vec<u8>>> {
let peer_key = state.peer_manager.peer_key();
if !peer_key.is_empty() && !peer::verify_peer_key(request, peer_key) {
return error_response("Invalid peer key", "UNAUTHORIZED", 401);
}
let body = read_body(request);
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: &mut Request,
) -> Response<std::io::Cursor<Vec<u8>>> {
let peer_key = state.peer_manager.peer_key();
if !peer_key.is_empty() && !peer::verify_peer_key(request, peer_key) {
return error_response("Invalid peer key", "UNAUTHORIZED", 401);
}
let body = read_body(request);
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: &Request,
) -> Response<std::io::Cursor<Vec<u8>>> {
let peer_key = state.peer_manager.peer_key();
if !peer_key.is_empty() && !peer::verify_peer_key(request, peer_key) {
return error_response("Invalid peer key", "UNAUTHORIZED", 401);
}
let path = request.url();
let user_id = path.find('?')
.and_then(|qpos| {
let query = &path[qpos + 1..];
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: &mut Request,
) -> Response<std::io::Cursor<Vec<u8>>> {
let peer_key = state.peer_manager.peer_key();
if !peer_key.is_empty() && !peer::verify_peer_key(request, peer_key) {
return error_response("Invalid peer key", "UNAUTHORIZED", 401);
}
let body = read_body(request);
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 balance_after = state.ledger.local_add_credits(&owner_user_id, req.amount);
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: format!("earn-{}", req.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: &Request,
) -> Response<std::io::Cursor<Vec<u8>>> {
let peer_key = state.peer_manager.peer_key();
if !peer_key.is_empty() && !peer::verify_peer_key(request, peer_key) {
return error_response("Invalid peer key", "UNAUTHORIZED", 401);
}
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),
};
json_response(&identity, 200)
}
pub fn execute_offer_locally(
state: &BrokerState,
offer: &task_board::TaskOffer,
) -> Result<task_board::TaskResult, task_board::TaskReject> {
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().map_or(true, |wt| w.worker_type.as_str() == wt.as_str())
})
.collect();
if local_workers.is_empty() {
return Err(task_board::TaskReject {
task_id: offer.task_id.clone(),
reason: "No matching local worker".to_string(),
});
}
let idx = (REQUEST_COUNTER.fetch_add(1, Ordering::SeqCst) as usize) % local_workers.len();
let worker = local_workers[idx].clone();
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::AgentBuilder::new()
.timeout(std::time::Duration::from_secs_f64(timeout + 5.0))
.build();
let forward_result = agent.post(&worker_uri)
.set("Content-Type", "application/octet-stream")
.set("X-Zakuro-Request-Id", &offer.task_id)
.set("X-Zakuro-Timeout-Secs", &format!("{:.1}", timeout))
.send_bytes(&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.header("X-Zakuro-Pid").map(|s| s.to_string());
let worker_ip = response.header("X-Zakuro-IP").map(|s| s.to_string());
let mut response_body = Vec::new();
let _ = response.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();
Ok(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.uri.clone(),
price_per_hour: worker.pricing.price_per_hour,
executor_owner,
worker_pid,
worker_ip,
})
}
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: &mut Request,
) -> Response<std::io::Cursor<Vec<u8>>> {
let peer_key = state.peer_manager.peer_key();
if !peer_key.is_empty() && !peer::verify_peer_key(request, peer_key) {
return error_response("Invalid peer key", "UNAUTHORIZED", 401);
}
let body = read_body(request);
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),
};
match execute_offer_locally(&state, &offer) {
Ok(result) => json_response(&result, 200),
Err(reject) => json_response(&reject, 404),
}
}
fn forward_to_worker(
worker_uri: &str,
body: &[u8],
request_id: &str,
effective_timeout: f64,
) -> Result<ureq::Response, ureq::Error> {
if effective_timeout > 0.0 {
let agent = ureq::AgentBuilder::new()
.timeout(std::time::Duration::from_secs_f64(effective_timeout + 5.0))
.build();
agent.post(worker_uri)
.set("Content-Type", "application/octet-stream")
.set("X-Zakuro-Request-Id", request_id)
.set("X-Zakuro-Timeout-Secs", &format!("{:.1}", effective_timeout))
.send_bytes(body)
} else {
ureq::post(worker_uri)
.set("Content-Type", "application/octet-stream")
.set("X-Zakuro-Request-Id", request_id)
.send_bytes(body)
}
}
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, Response<std::io::Cursor<Vec<u8>>>> {
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) => 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, Response<std::io::Cursor<Vec<u8>>>> {
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) => fail(&e.to_string(), "LEDGER_ERROR", 503),
}
}
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),
);
}
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,
) -> 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);
let dur_secs = duration_ms / 1000.0;
let earn = worker.pricing.price_per_hour / 3600.0 * dur_secs;
let worker_ip = worker.uri
.strip_prefix("http://").unwrap_or(&worker.uri)
.split(':').next().unwrap_or("");
if let Some(peer_url) = state.peer_manager.get_url_for_ip(worker_ip) {
if let Some(client) = state.peer_manager.get_client(&peer_url) {
if let Err(e) = client.earn(earn, duration_ms, &worker.id, user_id, request_id) {
eprintln!(" [EARN] Failed to notify peer {}: {}", peer_url, e);
}
}
}
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,
});
balance
}
}
}
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,
) -> Response<std::io::Cursor<Vec<u8>>> {
let peer_urls: Vec<String> = state.peer_manager.peer_urls();
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());
}
}
if winning_peer.is_none() && !peer_urls.is_empty() {
let (tx, rx) = std::sync::mpsc::channel::<Option<(String, task_board::TaskResult, String)>>();
for peer_url in &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);
for _ in 0..peer_urls.len() {
if let Ok(Some((url, res, trans))) = rx.recv() {
winning_peer = Some(url);
winning_result = Some(res);
winning_transport = Some(trans);
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(),
result.worker_ip.as_deref(),
Some(request_id),
);
if state.is_billing_enabled() && cost > 0.0 {
state.credits.set_balance(user_id, balance_after);
state.ledger.local_credits.entry(user_id.to_string())
.and_modify(|b| *b = (*b - cost).max(0.0));
state.ledger.authoritative_balances.entry(user_id.to_string())
.and_modify(|b| *b = (*b - cost).max(0.0));
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 {
if let Some(client) = state.peer_manager.get_client(&peer_url) {
let dur_secs = duration_ms / 1000.0;
let earn_amount = result.price_per_hour / 3600.0 * dur_secs;
let _ = client.earn(earn_amount, duration_ms, &result.worker_name, user_id, request_id);
}
}
let response_body = base64::Engine::decode(
&base64::engine::general_purpose::STANDARD,
&result.payload_b64,
).unwrap_or_default();
return Response::from_data(response_body)
.with_status_code(200)
.with_header(Header::from_bytes("Content-Type", "application/octet-stream").unwrap())
.with_header(Header::from_bytes("X-Zakuro-Request-Id", request_id).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Worker", result.worker_name.as_str()).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Cost", format!("{:.6}", cost)).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Credits-Remaining", format!("{:.6}", balance_after)).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Duration-Ms", format!("{:.2}", duration_ms)).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Transport", transport_used.as_str()).unwrap());
}
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)
}
fn handle_execute(state: Arc<BrokerState>, request: &mut Request, verbose: bool) -> Response<std::io::Cursor<Vec<u8>>> {
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 = get_header(request, "X-Zakuro-Instance-Action");
let instance_id = get_header(request, "X-Zakuro-Instance-Id");
let user_id = if let Some(auth) = get_header(request, "Authorization") {
match auth.strip_prefix("Bearer ") {
Some(token) => {
match state.ledger.resolve_user_from_api_key(token) {
Ok(uid) => uid,
Err(e) => return error_response(&e.to_string(), "UNAUTHORIZED", 401),
}
}
None => return error_response("Invalid Authorization header format, expected Bearer token", "UNAUTHORIZED", 401),
}
} else if state.is_local_mode() {
get_header(request, "X-Zakuro-User").unwrap_or_else(|| "anonymous".to_string())
} else {
return error_response("API key required. Set Authorization: Bearer <key>", "UNAUTHORIZED", 401)
};
let requirements: ResourceRequirements = get_header(request, "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 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 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 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 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 try_local: Option<RoutingDecision> = if requirements.remote_only {
None
} else {
let live_ts_ip = super::discovery::get_tailscale_ip()
.or_else(|| state.own_tailscale_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 = read_body(request);
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(),
};
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(),
});
return dispatch_to_peers(
&state, offer, start, &user_id, is_self_owned,
&request_id, tx_num, verbose, balance_before,
);
}
}
} 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 error_response(&e.message, &e.code, status);
}
let healthy = state.workers.healthy();
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 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 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()
|| state.is_local_worker(&routing.worker.uri)
|| is_self_owned;
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 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 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,
});
}
state.active_requests.insert(request_id.clone(), routing.worker.id.clone());
state.workers.increment_active(&routing.worker.id);
let body = read_body(request);
let worker_uri = format!("{}/execute", routing.worker.uri.trim_end_matches('/'));
let worker_name = routing.worker.name.clone();
let forward_result = forward_to_worker(&worker_uri, &body, &request_id, effective_timeout);
let duration_ms = start.elapsed().as_secs_f64() * 1000.0;
match forward_result {
Ok(response) => {
let actual_cost = if is_local {
0.0
} else {
let actual_duration_secs = duration_ms / 1000.0;
routing.worker.pricing.estimate_cost(actual_duration_secs)
};
if !is_local {
let _ = state.wal.update_status(
&request_id,
WalStatus::Executed,
Some(actual_cost),
Some(duration_ms),
);
}
let credits_remaining = commit_credits(
&state, &authority, is_local, &reservation_id, actual_cost,
balance_before, &request_id, &user_id, &routing.worker, duration_ms,
);
if !is_local {
let _ = state.wal.update_status(
&request_id,
WalStatus::Committed,
Some(actual_cost),
Some(duration_ms),
);
}
state.workers.record_request(&routing.worker.id, duration_ms, true);
state.active_requests.remove(&request_id);
let worker_pid = response.header("X-Zakuro-Pid").map(|s| s.to_string());
let worker_ip = response.header("X-Zakuro-IP").map(|s| s.to_string())
.or_else(|| {
let u = routing.worker.uri.strip_prefix("http://").unwrap_or(&routing.worker.uri);
u.split(':').next().map(|s| s.to_string())
});
let mut response_body = Vec::new();
let _ = response.into_reader().read_to_end(&mut response_body);
if instance_action.as_deref() == Some("create_instance") {
if let Some(ref iid) = instance_id {
state.instance_registry.insert(iid.clone(), routing.worker.id.clone());
}
}
let action_label = match instance_action.as_deref() {
Some("create_instance") => "CREATE_INST",
Some("call_method") => "CALL_METHOD",
_ => "EXECUTE",
};
log_transaction_with_worker_info(
&state, verbose, tx_num, &user_id, action_label,
actual_cost, credits_remaining, Some(&worker_name), duration_ms, "OK",
Some(routing.worker.pricing.price_per_hour),
state.config.owner_user_id.as_deref(),
worker_pid.as_deref(),
worker_ip.as_deref(),
Some(&request_id),
);
Response::from_data(response_body)
.with_status_code(200)
.with_header(Header::from_bytes("Content-Type", "application/octet-stream").unwrap())
.with_header(Header::from_bytes("X-Zakuro-Request-Id", request_id).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Worker", worker_name.as_str()).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Cost", format!("{:.6}", actual_cost)).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Credits-Remaining", format!("{:.6}", credits_remaining)).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Duration-Ms", format!("{:.2}", duration_ms)).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Transport", "local").unwrap())
}
Err(e) => {
let is_timeout = e.to_string().contains("timed out")
|| e.to_string().contains("Timeout");
let timeout_cost = if is_timeout && !is_local {
let elapsed_secs = duration_ms / 1000.0;
routing.worker.pricing.estimate_cost(elapsed_secs)
} else {
0.0
};
if is_timeout {
if !is_local && timeout_cost > 0.0 {
match &authority {
Authority::Local => {
if let Err(e) = state.ledger.local_commit(&reservation_id, timeout_cost) {
eprintln!(" [BILLING] local_commit failed for {}: {}", request_id, e);
}
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: request_id.clone(),
user_id: user_id.clone(),
tx_type: "commit".to_string(),
amount: timeout_cost,
balance_after: balance_before - timeout_cost,
worker_id: routing.worker.id.clone(),
duration_ms,
source_node: state.config.node_name.clone(),
worker_name: Some(routing.worker.name.clone()),
worker_uri: Some(routing.worker.uri.clone()),
price_per_hour: routing.worker.pricing.price_per_hour,
});
}
Authority::Peer(url) => {
if let Some(client) = state.peer_manager.get_client(url) {
let _ = client.commit(&reservation_id, timeout_cost);
} else {
let _ = state.ledger.commit(&reservation_id, timeout_cost);
}
state.ledger.publish_transaction(
&request_id, &user_id, "commit", timeout_cost,
balance_before - timeout_cost, &routing.worker.id, duration_ms,
state.config.node_name.as_deref(),
);
}
Authority::Standalone => {
let _ = state.ledger.commit(&reservation_id, timeout_cost);
state.ledger.publish_transaction(
&request_id, &user_id, "commit", timeout_cost,
balance_before - timeout_cost, &routing.worker.id, duration_ms,
state.config.node_name.as_deref(),
);
}
}
let _ = state.wal.update_status(
&request_id,
WalStatus::Committed,
Some(timeout_cost),
Some(duration_ms),
);
}
state.active_requests.remove(&request_id);
log_transaction(&state, verbose, tx_num, &user_id, "EXECUTE", timeout_cost, balance_before - timeout_cost, Some(&worker_name), duration_ms, "FAIL");
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(&routing.worker.id);
state.workers.cancel_quota_reservation(&routing.worker.id);
if !is_local {
cancel_credits(&state, &authority, &reservation_id, &request_id, &user_id, &routing.worker.id, balance_before, duration_ms);
}
state.active_requests.remove(&request_id);
const MAX_RETRIES: usize = 2;
let mut tried_ids: Vec<String> = vec![routing.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, tx_num, &user_id, "EXECUTE", 0.0, 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
);
state.workers.increment_active(&next.id);
state.active_requests.insert(request_id.clone(), next.id.clone());
let retry_start = Instant::now();
let retry_result = forward_to_worker(&retry_uri, &body, &request_id, effective_timeout);
let retry_duration_ms = retry_start.elapsed().as_secs_f64() * 1000.0;
match retry_result {
Ok(response) => {
let actual_cost = if is_local {
0.0
} else {
next.pricing.estimate_cost(retry_duration_ms / 1000.0)
};
let retry_reservation_id = if !is_local {
let res_id = format!("{}-retry{}", request_id, attempt);
match &authority {
Authority::Local => {
let _ = state.ledger.local_reserve(&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(&user_id, actual_cost, &res_id).is_ok() {
let _ = client.commit(&res_id, actual_cost);
}
}
}
Authority::Standalone => {
if state.ledger.reserve(&user_id, actual_cost, &res_id).is_ok() {
let _ = state.ledger.commit(&res_id, actual_cost);
}
}
}
res_id
} else {
String::new()
};
if !is_local {
let _ = state.wal.update_status(
&request_id,
WalStatus::Executed,
Some(actual_cost),
Some(retry_duration_ms),
);
let _ = state.wal.update_status(
&request_id,
WalStatus::Committed,
Some(actual_cost),
Some(retry_duration_ms),
);
}
state.workers.record_request(&next.id, retry_duration_ms, true);
state.active_requests.remove(&request_id);
let credits_remaining = if is_local {
balance_before
} else {
balance_before - actual_cost
};
state.tx_buffer.push_transaction(BufferedTransaction {
request_id: format!("{}-retry{}", request_id, attempt),
user_id: 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 instance_action.as_deref() {
Some("create_instance") => "CREATE_INST",
Some("call_method") => "CALL_METHOD",
_ => "EXECUTE",
};
log_transaction(&state, verbose, tx_num, &user_id, action_label, actual_cost, credits_remaining, Some(&retry_name), retry_duration_ms, "OK");
let mut response_body = Vec::new();
let _ = response.into_reader().read_to_end(&mut response_body);
let _ = retry_reservation_id; return Response::from_data(response_body)
.with_status_code(200)
.with_header(Header::from_bytes("Content-Type", "application/octet-stream").unwrap())
.with_header(Header::from_bytes("X-Zakuro-Request-Id", request_id).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Cost", format!("{:.6}", actual_cost)).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Credits-Remaining", format!("{:.6}", credits_remaining)).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Duration-Ms", format!("{:.2}", retry_duration_ms)).unwrap())
.with_header(Header::from_bytes("X-Zakuro-Retries", format!("{}", attempt)).unwrap());
}
Err(_retry_err) => {
state.workers.mark_unhealthy(&next.id);
state.workers.cancel_quota_reservation(&next.id);
state.active_requests.remove(&request_id);
}
}
}
log_transaction(&state, verbose, tx_num, &user_id, "EXECUTE", 0.0, 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,
)
}
}
}
pub fn start_server(config: BrokerConfig) -> std::io::Result<()> {
crate::async_exec::spawn_detached(|| {
if let Err(e) = super::config::fetch_broker_config() {
eprintln!(" {} Failed to fetch broker config from API: {}", "Warning:".yellow(), e);
}
});
let addr = format!("{}:{}", config.host, config.port);
let server = Server::http(&addr).map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e))?;
let verbose = config.verbose;
let daemon = config.daemon;
let tui_mode = config.tui_mode;
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);
if state.peer_manager.is_enabled() {
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_tailscale_ip.as_deref().unwrap_or("127.0.0.1"),
config.port
);
super::quic::run_subscription_client(
state.clone(),
our_url,
std::sync::Arc::new(|s, o| execute_offer_locally(s, o)),
);
}
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);
let running = Arc::new(AtomicBool::new(true));
let state_clone = state.clone();
let running_cleanup = running.clone();
async_exec::spawn_detached(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());
}
state_clone.stats.tick_rps();
if compaction_counter % 12 == 0 {
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.tailscale_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();
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(),
) {
Ok(()) => {
println!(" [WORKER_SYNC] Synced {} worker(s) to {} via API", workers.len(), api_url);
}
Err(e) => {
eprintln!(" [WORKER_SYNC] API sync 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);
}
if let Err(e) = state_clone.wal.flush_buffer() {
eprintln!(" [WAL] flush_buffer failed: {}", e);
}
if compaction_counter % 12 == 0 && state_clone.peer_manager.is_enabled() {
state_clone.peer_manager.health_check_all();
}
if (compaction_counter <= 1 || compaction_counter % 12 == 0)
&& state_clone.peer_manager.is_enabled()
&& state_clone.quic.get().is_some()
{
state_clone.peer_manager.discover_quic_ports();
}
if compaction_counter % 6 == 0 && 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 % 60 == 0 {
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);
}
if let Err(e) = state_clone.wal.flush_buffer() {
eprintln!(" [WAL] flush_buffer failed: {}", e);
}
});
if config.enable_discovery {
let state_for_discovery = state.clone();
let discovery_config = config.discovery.clone();
let discovery_verbose = verbose && !tui_mode;
let mode = detect_discovery_mode(&discovery_config.subnet);
let is_local = matches!(mode, DiscoveryMode::Local);
state.set_local_mode(is_local);
if verbose && !daemon && !tui_mode {
match &mode {
DiscoveryMode::Tailscale { subnet } => {
println!(" {} Tailscale 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 tui_mode {
let default_panic = std::panic::take_hook();
std::panic::set_hook(Box::new(move |info| {
tui::cleanup_terminal();
default_panic(info);
}));
let state_for_server = state.clone();
let running_for_server = running.clone();
async_exec::spawn_detached(move || {
for request in server.incoming_requests() {
if !running_for_server.load(Ordering::Relaxed) {
break;
}
let state = state_for_server.clone();
async_exec::spawn_detached(move || {
handle_request(state, request, false); });
}
});
let stats = state.stats.clone();
let result = tui::run_tui(state, stats, running.clone());
tui::cleanup_terminal();
return result;
}
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));
}
}
for request in server.incoming_requests() {
if !running.load(Ordering::Relaxed) {
break;
}
let state = state.clone();
async_exec::spawn_detached(move || {
handle_request(state, request, verbose);
});
}
Ok(())
}