pub mod bench;
pub mod bench_mesh;
pub mod credits;
pub mod discovery;
pub mod flush;
pub mod http_adapter;
pub mod info;
pub mod ledger;
pub mod node_health;
pub mod node_identity;
pub mod node_sync;
pub mod peer;
pub mod pricing;
pub mod quic;
pub mod recovery;
pub mod roster_cache;
pub mod router;
pub mod selection;
pub mod server;
pub mod stats;
pub mod task_board;
pub mod tui;
pub mod uri;
pub mod voucher;
pub mod wal;
pub mod worker;
pub mod worker_quic;
use std::sync::Arc;
use std::sync::OnceLock;
use dashmap::DashMap;
pub use credits::CreditManager;
pub use discovery::DiscoveryConfig;
pub use flush::TransactionBuffer;
pub use ledger::Ledger;
pub use peer::PeerManager;
pub use pricing::BrokerPrice;
pub use router::Router;
pub use server::start_server;
pub use stats::StatsCollector;
pub use tui::{cleanup_terminal, run_remote_tui};
pub use wal::Wal;
pub use worker::{WorkerRegistry, WorkerStatus};
fn detect_node_name() -> Option<String> {
std::process::Command::new("hostname")
.output()
.ok()
.and_then(|o| String::from_utf8(o.stdout).ok())
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
}
pub fn node_name_or_default() -> String {
std::env::var("ZAKURO_NODE_NAME")
.ok()
.filter(|s| !s.is_empty())
.or_else(detect_node_name)
.unwrap_or_else(|| "unknown".to_string())
}
#[derive(Debug, Clone)]
pub struct BrokerConfig {
pub host: String,
pub port: u16,
pub health_check_interval: u64,
pub worker_timeout: u64,
pub min_credits: f64,
pub daemon: bool,
pub verbose: bool,
pub tui_mode: bool,
pub enable_discovery: bool,
pub discovery: DiscoveryConfig,
pub worker_key: Option<String>,
pub owner_user_id: Option<String>,
pub node_name: Option<String>,
pub peer_key: Option<String>,
pub enable_p2p: bool,
pub api_url: Option<String>,
pub api_key: Option<String>,
pub wireguard_ip_override: Option<String>,
pub quic_port: Option<u16>,
pub runtime_worker_threads: Option<usize>,
}
impl Default for BrokerConfig {
fn default() -> Self {
Self {
host: "0.0.0.0".to_string(),
port: 9000,
health_check_interval: 5,
worker_timeout: 30,
min_credits: 0.001,
daemon: false,
verbose: true,
tui_mode: false,
enable_discovery: true,
discovery: DiscoveryConfig::default(),
worker_key: std::env::var("ZAKURO_WORKER_KEY").ok(),
owner_user_id: {
std::env::var("ZAKURO_API_KEY").ok().and_then(|k| {
k.strip_prefix("zk_")
.and_then(|rest| rest.rfind('_').map(|pos| rest[..pos].to_string()))
})
},
node_name: std::env::var("ZAKURO_NODE_NAME")
.ok()
.or_else(detect_node_name),
peer_key: std::env::var("ZAKURO_PEER_KEY")
.ok()
.filter(|k| !k.trim().is_empty())
.or_else(crate::credentials::load_mesh_peer_key),
enable_p2p: std::env::var("ZAKURO_P2P")
.map(|v| v == "true" || v == "1")
.unwrap_or(false),
api_url: std::env::var("ZAKURO_API_URL").ok(),
api_key: std::env::var("ZAKURO_API_KEY").ok(),
wireguard_ip_override: None,
quic_port: None,
runtime_worker_threads: None,
}
}
}
pub fn p2p_default() -> bool {
match std::env::var("ZAKURO_P2P") {
Ok(v) => v == "true" || v == "1",
Err(_) => discovery::get_mesh_ip().is_some(),
}
}
pub fn apply_user_broker_defaults() {
crate::credentials::load_into_env();
let set = |k: &str| {
std::env::var(k)
.map(|v| !v.trim().is_empty())
.unwrap_or(false)
};
if set("ZAKURO_API_KEY")
&& set("ZAKURO_API_URL")
&& std::env::var("ZAKURO_REQUIRE_VOUCHER").is_err()
{
std::env::set_var("ZAKURO_REQUIRE_VOUCHER", "1");
}
}
pub struct BrokerState {
pub workers: WorkerRegistry,
pub credits: CreditManager,
pub ledger: Ledger,
pub router: Router,
pub config: BrokerConfig,
pub active_requests: DashMap<String, String>,
pub instance_registry: DashMap<String, String>,
pub local_mode: std::sync::atomic::AtomicBool,
pub stats: Arc<StatsCollector>,
pub wal: Wal,
pub own_wireguard_ip: Option<String>,
pub peer_manager: PeerManager,
pub tx_buffer: TransactionBuffer,
pub quic: OnceLock<Arc<quic::QuicTransport>>,
pub quic_port: std::sync::atomic::AtomicU16,
pub http_client: ureq::Agent,
pub async_http: reqwest::Client,
pub pending_offers: DashMap<String, PendingOffer>,
pub node_key: Arc<node_identity::NodeKey>,
pub node_roster: DashMap<String, bool>,
pub node_sig_guard: node_identity::ReplayGuard,
pub dash_voucher_pubkey: std::sync::RwLock<Option<String>>,
pub price: BrokerPrice,
}
pub struct PendingOffer {
pub offer: task_board::TaskOffer,
pub worker_id: String,
pub created: std::time::Instant,
}
impl BrokerState {
pub fn new() -> Self {
Self::with_config(BrokerConfig::default())
}
pub fn with_config(config: BrokerConfig) -> Self {
let ledger = Ledger::new(config.api_url.clone(), config.api_key.clone());
let wal = {
let mut candidates: Vec<std::path::PathBuf> = vec![];
if let Ok(p) = std::env::var("ZAKURO_WAL_PATH") {
candidates.push(std::path::PathBuf::from(p));
}
if let Some(home) = std::env::var("HOME")
.ok()
.or_else(|| std::env::var("USERPROFILE").ok())
{
let dir = std::path::PathBuf::from(home).join(".zakuro");
let _ = std::fs::create_dir_all(&dir);
candidates.push(dir.join("wal.jsonl"));
}
candidates.push(std::path::PathBuf::from("/tmp/zakuro-wal.jsonl"));
candidates
.into_iter()
.find_map(|p| Wal::open(p.to_str().unwrap_or("/tmp/zakuro-wal.jsonl")).ok())
.expect("Failed to open WAL on any candidate path")
};
let own_wireguard_ip = config
.wireguard_ip_override
.clone()
.or_else(discovery::get_effective_node_ip);
let peer_addresses: Vec<String> = if config.discovery.peers.is_empty()
&& config.enable_p2p
&& std::env::var("ZAKURO_DISCOVER_BROKER_PEERS").unwrap_or_else(|_| "true".into())
!= "false"
{
let mut discovered =
discovery::discover_broker_peers_on_localhost(config.port, 9000, 9010);
if !discovered.is_empty() {
eprintln!(
" [P2P] Discovered {} broker(s) on localhost (ZAKURO_API_KEY optional)",
discovered.len()
);
}
if let Some(mesh_ip) = discovery::get_mesh_ip() {
let mesh_peers = discovery::discover_broker_peers_on_mesh_subnet(
&mesh_ip,
config.peer_key.as_deref().unwrap_or(""),
);
if !mesh_peers.is_empty() {
eprintln!(
" [P2P] Discovered {} broker(s) on the mesh subnet",
mesh_peers.len()
);
}
discovered.extend(mesh_peers);
}
discovered
} else {
config.discovery.peers.clone()
};
let node_key = Arc::new(node_identity::NodeKey::load_or_create());
let peer_key = config.peer_key.clone().unwrap_or_default();
let peer_manager = PeerManager::new_with_node_key(
own_wireguard_ip.as_deref(),
&peer_addresses,
config.port,
peer_key,
config.enable_p2p,
Some(node_key.clone()),
);
Self {
workers: WorkerRegistry::new(),
credits: CreditManager::new(),
ledger,
router: Router::new(),
config,
active_requests: DashMap::new(),
instance_registry: DashMap::new(),
local_mode: std::sync::atomic::AtomicBool::new(false),
stats: Arc::new(StatsCollector::new()),
wal,
own_wireguard_ip,
peer_manager,
tx_buffer: TransactionBuffer::new(),
quic: OnceLock::new(),
quic_port: std::sync::atomic::AtomicU16::new(0),
http_client: ureq::Agent::new_with_config(ureq::Agent::config_builder().build()),
async_http: reqwest::Client::builder()
.pool_max_idle_per_host(32)
.build()
.expect("reqwest client builds with rustls"),
pending_offers: DashMap::new(),
node_key,
node_roster: DashMap::new(),
node_sig_guard: node_identity::ReplayGuard::new(),
dash_voucher_pubkey: std::sync::RwLock::new(None),
price: BrokerPrice::from_env_or_default(),
}
}
pub fn is_local_mode(&self) -> bool {
self.local_mode.load(std::sync::atomic::Ordering::Relaxed)
}
pub fn set_local_mode(&self, local: bool) {
self.local_mode
.store(local, std::sync::atomic::Ordering::Relaxed);
}
pub fn is_billing_enabled(&self) -> bool {
(self.config.api_url.is_some() && self.config.api_key.is_some())
|| !std::env::var("ZAKURO_MASTER_KEY").unwrap_or_default().is_empty()
}
pub fn is_local_worker(&self, worker_uri: &str) -> bool {
if self.is_local_mode() {
return true; }
if worker_uri.contains("127.0.0.1") || worker_uri.contains("localhost") {
return true;
}
if let Some(ref own_ip) = self.own_wireguard_ip {
worker_uri.contains(own_ip)
} else {
false
}
}
pub fn verify_owner_with_dashboard(&mut self) -> Option<String> {
let api_url = self.config.api_url.as_ref()?;
let api_key = self.config.api_key.as_ref()?;
let url = format!("{}/api/auth/me/api-key", api_url.trim_end_matches('/'));
let agent = ureq::Agent::new_with_config(
ureq::Agent::config_builder()
.timeout_global(Some(std::time::Duration::from_secs(10)))
.build(),
);
match agent
.get(&url)
.header("Authorization", &format!("Bearer {}", api_key))
.call()
{
Ok(resp) => {
let body = resp.into_body().read_to_string().unwrap_or_default();
let parsed: serde_json::Value = serde_json::from_str(&body).unwrap_or_default();
if let Some(uid) = parsed["zakuro_user_id"].as_str() {
let derived = self.config.owner_user_id.as_deref();
if derived.is_some() && derived != Some(uid) {
eprintln!(" [HANDSHAKE] owner_user_id mismatch: key-derived={}, dashboard={}. Using key-derived value.",
derived.unwrap_or("?"), uid);
return derived.map(str::to_string);
}
if self.config.owner_user_id.is_none() {
self.config.owner_user_id = Some(uid.to_string());
}
Some(
self.config
.owner_user_id
.clone()
.unwrap_or_else(|| uid.to_string()),
)
} else {
eprintln!(" [HANDSHAKE] Dashboard did not return zakuro_user_id");
None
}
}
Err(e) => {
eprintln!(" [HANDSHAKE] Failed to verify owner with dashboard: {}", e);
None
}
}
}
}
pub type SharedBrokerState = Arc<BrokerState>;
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn broker_state_has_shared_async_client() {
let s = BrokerState::new();
let _c: reqwest::Client = s.async_http.clone();
}
}