use dashmap::DashMap;
use serde::{Deserialize, Serialize};
use std::collections::BTreeSet;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerReserveRequest {
pub user_id: String,
pub amount: f64,
pub request_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerReserveResponse {
pub reservation_id: String,
pub balance_before: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerCommitRequest {
pub reservation_id: String,
pub actual_cost: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerCommitResponse {
pub balance_after: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerCancelRequest {
pub reservation_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerEarnRequest {
pub amount: f64,
pub duration_ms: f64,
pub worker_id: String,
pub requesting_user: String,
pub request_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerBalanceResponse {
pub user_id: String,
pub balance: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerHealthResponse {
pub status: String,
pub broker_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerErrorResponse {
pub error: String,
pub code: String,
}
#[derive(Debug, Clone)]
pub enum Authority {
Local,
Peer(String), Standalone,
}
pub struct PeerClient {
agent: ureq::Agent,
base_url: String,
peer_key: String,
reachable: AtomicBool,
}
impl PeerClient {
pub fn new(base_url: String, peer_key: String) -> Self {
let agent = ureq::AgentBuilder::new()
.timeout_connect(Duration::from_millis(500))
.timeout_read(Duration::from_secs(5))
.timeout_write(Duration::from_secs(2))
.build();
Self {
agent,
base_url,
peer_key,
reachable: AtomicBool::new(false),
}
}
pub fn base_url(&self) -> &str {
&self.base_url
}
pub fn is_reachable(&self) -> bool {
self.reachable.load(Ordering::Relaxed)
}
pub fn check_health(&self) -> bool {
let url = format!("{}/peer/health", self.base_url);
let result = self.agent.get(&url)
.set("X-Peer-Key", &self.peer_key)
.call();
let ok = matches!(result, Ok(r) if r.status() == 200);
self.reachable.store(ok, Ordering::Relaxed);
ok
}
pub fn reserve(
&self,
user_id: &str,
amount: f64,
request_id: &str,
) -> Result<PeerReserveResponse, String> {
let url = format!("{}/peer/reserve", self.base_url);
let req = PeerReserveRequest {
user_id: user_id.to_string(),
amount,
request_id: request_id.to_string(),
};
let body = serde_json::to_vec(&req).map_err(|e| e.to_string())?;
match self.agent.post(&url)
.set("X-Peer-Key", &self.peer_key)
.set("Content-Type", "application/json")
.send_bytes(&body)
{
Ok(resp) => {
self.reachable.store(true, Ordering::Relaxed);
let text = resp.into_string().map_err(|e| e.to_string())?;
serde_json::from_str(&text).map_err(|e| e.to_string())
}
Err(ureq::Error::Status(code, resp)) => {
self.reachable.store(true, Ordering::Relaxed);
let text = resp.into_string().unwrap_or_default();
if let Ok(err) = serde_json::from_str::<PeerErrorResponse>(&text) {
Err(format!("[{}] {}", code, err.error))
} else {
Err(format!("Peer returned {}: {}", code, text))
}
}
Err(e) => {
self.reachable.store(false, Ordering::Relaxed);
Err(format!("Peer unreachable: {}", e))
}
}
}
pub fn commit(
&self,
reservation_id: &str,
actual_cost: f64,
) -> Result<PeerCommitResponse, String> {
let url = format!("{}/peer/commit", self.base_url);
let req = PeerCommitRequest {
reservation_id: reservation_id.to_string(),
actual_cost,
};
let body = serde_json::to_vec(&req).map_err(|e| e.to_string())?;
match self.agent.post(&url)
.set("X-Peer-Key", &self.peer_key)
.set("Content-Type", "application/json")
.send_bytes(&body)
{
Ok(resp) => {
self.reachable.store(true, Ordering::Relaxed);
let text = resp.into_string().map_err(|e| e.to_string())?;
serde_json::from_str(&text).map_err(|e| e.to_string())
}
Err(ureq::Error::Status(code, resp)) => {
self.reachable.store(true, Ordering::Relaxed);
let text = resp.into_string().unwrap_or_default();
Err(format!("Peer commit failed ({}): {}", code, text))
}
Err(e) => {
self.reachable.store(false, Ordering::Relaxed);
Err(format!("Peer unreachable: {}", e))
}
}
}
pub fn cancel(&self, reservation_id: &str) -> Result<(), String> {
let url = format!("{}/peer/cancel", self.base_url);
let req = PeerCancelRequest {
reservation_id: reservation_id.to_string(),
};
let body = serde_json::to_vec(&req).map_err(|e| e.to_string())?;
match self.agent.post(&url)
.set("X-Peer-Key", &self.peer_key)
.set("Content-Type", "application/json")
.send_bytes(&body)
{
Ok(_) => {
self.reachable.store(true, Ordering::Relaxed);
Ok(())
}
Err(ureq::Error::Status(_, _)) => {
self.reachable.store(true, Ordering::Relaxed);
Ok(()) }
Err(e) => {
self.reachable.store(false, Ordering::Relaxed);
Err(format!("Peer unreachable: {}", e))
}
}
}
pub fn earn(
&self,
amount: f64,
duration_ms: f64,
worker_id: &str,
requesting_user: &str,
request_id: &str,
) -> Result<(), String> {
let url = format!("{}/peer/earn", self.base_url);
let req = PeerEarnRequest {
amount,
duration_ms,
worker_id: worker_id.to_string(),
requesting_user: requesting_user.to_string(),
request_id: request_id.to_string(),
};
let body = serde_json::to_vec(&req).map_err(|e| e.to_string())?;
match self.agent.post(&url)
.set("X-Peer-Key", &self.peer_key)
.set("Content-Type", "application/json")
.send_bytes(&body)
{
Ok(_) => {
self.reachable.store(true, Ordering::Relaxed);
Ok(())
}
Err(ureq::Error::Status(_, _)) => {
self.reachable.store(true, Ordering::Relaxed);
Ok(()) }
Err(e) => {
self.reachable.store(false, Ordering::Relaxed);
Err(format!("Peer unreachable: {}", e))
}
}
}
pub fn handshake(&self) -> Result<super::task_board::PeerIdentity, String> {
let url = format!("{}/peer/identity", self.base_url);
match self.agent.get(&url)
.set("X-Peer-Key", &self.peer_key)
.call()
{
Ok(resp) => {
self.reachable.store(true, Ordering::Relaxed);
let text = resp.into_string().map_err(|e| e.to_string())?;
serde_json::from_str(&text).map_err(|e| e.to_string())
}
Err(e) => {
self.reachable.store(false, Ordering::Relaxed);
Err(format!("Peer handshake failed: {}", e))
}
}
}
pub fn offer_task(
&self,
offer: &super::task_board::TaskOffer,
) -> Result<super::task_board::TaskResult, String> {
let url = format!("{}/peer/tasks/offer", self.base_url);
let body = serde_json::to_vec(offer).map_err(|e| e.to_string())?;
let timeout = if offer.timeout_secs > 0.0 { offer.timeout_secs + 10.0 } else { 310.0 };
let task_agent = ureq::AgentBuilder::new()
.timeout_connect(Duration::from_millis(2000))
.timeout_read(Duration::from_secs_f64(timeout))
.timeout_write(Duration::from_secs(5))
.build();
match task_agent.post(&url)
.set("X-Peer-Key", &self.peer_key)
.set("Content-Type", "application/json")
.send_bytes(&body)
{
Ok(resp) => {
self.reachable.store(true, Ordering::Relaxed);
let text = resp.into_string().map_err(|e| e.to_string())?;
serde_json::from_str(&text).map_err(|e| e.to_string())
}
Err(ureq::Error::Status(404, resp)) => {
self.reachable.store(true, Ordering::Relaxed);
let text = resp.into_string().unwrap_or_default();
Err(format!("Peer rejected task: {}", text))
}
Err(ureq::Error::Status(code, resp)) => {
self.reachable.store(true, Ordering::Relaxed);
let text = resp.into_string().unwrap_or_default();
Err(format!("Peer error ({}): {}", code, text))
}
Err(e) => {
self.reachable.store(false, Ordering::Relaxed);
Err(format!("Peer unreachable: {}", e))
}
}
}
pub fn get_balance(&self, user_id: &str) -> Result<f64, String> {
let url = format!("{}/peer/balance?user_id={}", self.base_url, user_id);
match self.agent.get(&url)
.set("X-Peer-Key", &self.peer_key)
.call()
{
Ok(resp) => {
self.reachable.store(true, Ordering::Relaxed);
let text = resp.into_string().map_err(|e| e.to_string())?;
let parsed: PeerBalanceResponse =
serde_json::from_str(&text).map_err(|e| e.to_string())?;
Ok(parsed.balance)
}
Err(e) => {
self.reachable.store(false, Ordering::Relaxed);
Err(format!("Peer unreachable: {}", e))
}
}
}
}
impl std::fmt::Debug for PeerClient {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PeerClient")
.field("base_url", &self.base_url)
.field("reachable", &self.is_reachable())
.finish()
}
}
pub struct PeerManager {
peers: DashMap<String, PeerClient>,
own_broker_index: u8,
total_brokers: u8,
peer_key: String,
enabled: bool,
quic_ports: DashMap<String, u16>,
}
impl PeerManager {
pub fn new(
own_ip: Option<&str>,
peer_addresses: &[String],
broker_port: u16,
peer_key: String,
enabled: bool,
) -> Self {
let peers = DashMap::new();
let mut all_addresses = BTreeSet::new();
let own_addr = if let Some(ip) = own_ip {
format!("{}:{}", ip, broker_port)
} else {
format!("127.0.0.1:{}", broker_port)
};
all_addresses.insert(own_addr.clone());
let mut peer_url_map: std::collections::HashMap<String, String> = std::collections::HashMap::new();
for peer in peer_addresses {
if peer.is_empty() {
continue;
}
let raw = peer
.strip_prefix("http://")
.or_else(|| peer.strip_prefix("https://"))
.unwrap_or(peer);
let actual_url = if raw.contains("://") {
peer.to_string()
} else {
format!("http://{}", raw)
};
let (ip, _) = raw.rsplit_once(':').unwrap_or((raw, ""));
let norm_addr = format!("{}:{}", ip, broker_port);
all_addresses.insert(norm_addr.clone());
peer_url_map.insert(norm_addr, actual_url);
}
let sorted: Vec<String> = all_addresses.into_iter().collect();
let own_index = sorted.iter().position(|a| *a == own_addr).unwrap_or(0) as u8;
let total = sorted.len() as u8;
if enabled {
for addr in &sorted {
if *addr != own_addr {
let base_url = peer_url_map.get(addr)
.cloned()
.unwrap_or_else(|| format!("http://{}", addr));
let client = PeerClient::new(base_url.clone(), peer_key.clone());
eprintln!(" [P2P] Registered peer broker at {}", base_url);
peers.insert(base_url, client);
}
}
}
eprintln!(
" [P2P] Broker index={}/{} (p2p={})",
own_index, total, enabled
);
Self {
peers,
own_broker_index: own_index,
total_brokers: total,
peer_key,
enabled,
quic_ports: DashMap::new(),
}
}
pub fn determine_authority(&self, user_id: &str) -> Authority {
if !self.enabled || self.total_brokers <= 1 {
return Authority::Standalone;
}
let hash = simple_hash(user_id);
let owner_index = (hash % self.total_brokers as u64) as u8;
if owner_index == self.own_broker_index {
Authority::Local
} else {
let mut peer_urls: Vec<String> = self.peers.iter().map(|e| e.key().clone()).collect();
peer_urls.sort();
let peer_idx = if owner_index < self.own_broker_index {
owner_index as usize
} else {
(owner_index - 1) as usize
};
if peer_idx < peer_urls.len() {
let url = &peer_urls[peer_idx];
if let Some(client) = self.peers.get(url) {
if client.is_reachable() {
return Authority::Peer(url.clone());
}
}
}
Authority::Standalone
}
}
pub fn get_client(&self, url: &str) -> Option<dashmap::mapref::one::Ref<'_, String, PeerClient>> {
self.peers.get(url)
}
pub fn get_url_for_ip(&self, ip: &str) -> Option<String> {
self.peers.iter()
.find(|e| e.key().contains(ip))
.map(|e| e.key().clone())
}
pub fn register_peer(&self, base_url: String) {
if !self.peers.contains_key(&base_url) {
let client = PeerClient::new(base_url.clone(), self.peer_key.clone());
self.peers.insert(base_url, client);
}
}
pub fn health_check_all(&self) {
for entry in self.peers.iter() {
entry.value().check_health();
}
}
pub fn is_enabled(&self) -> bool {
self.enabled
}
pub fn peer_key(&self) -> &str {
&self.peer_key
}
pub fn peer_count(&self) -> usize {
self.peers.len()
}
pub fn peer_urls(&self) -> Vec<String> {
self.peers.iter().map(|e| e.key().clone()).collect()
}
pub fn set_quic_port(&self, peer_url: &str, port: u16) {
if port > 0 {
self.quic_ports.insert(peer_url.to_string(), port);
}
}
pub fn get_quic_port(&self, peer_url: &str) -> u16 {
self.quic_ports.get(peer_url).map(|v| *v).unwrap_or(0)
}
pub fn discover_quic_ports(&self) {
for entry in self.peers.iter() {
if let Ok(identity) = entry.value().handshake() {
if identity.quic_port > 0 {
self.quic_ports.insert(entry.key().clone(), identity.quic_port);
eprintln!(" [QUIC] Discovered peer {} QUIC port: {}", entry.key(), identity.quic_port);
}
}
}
}
}
impl std::fmt::Debug for PeerManager {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PeerManager")
.field("own_broker_index", &self.own_broker_index)
.field("total_brokers", &self.total_brokers)
.field("enabled", &self.enabled)
.field("peer_count", &self.peers.len())
.finish()
}
}
pub fn verify_peer_key(request: &tiny_http::Request, expected_key: &str) -> bool {
if expected_key.is_empty() {
return false; }
request.headers().iter()
.find(|h| h.field.as_str().as_str().eq_ignore_ascii_case("X-Peer-Key"))
.map(|h| h.value.as_str() == expected_key)
.unwrap_or(false)
}
pub fn simple_hash(s: &str) -> u64 {
let mut hash: u64 = 0xcbf29ce484222325; for byte in s.as_bytes() {
hash ^= *byte as u64;
hash = hash.wrapping_mul(0x100000001b3); }
hash
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_simple_hash_deterministic() {
let h1 = simple_hash("9000000001");
let h2 = simple_hash("9000000001");
assert_eq!(h1, h2);
}
#[test]
fn test_simple_hash_distribution() {
let mut counts = [0u32; 2];
for i in 9000000001..9000000101u64 {
let h = simple_hash(&i.to_string());
counts[(h % 2) as usize] += 1;
}
assert!(counts[0] >= 20, "Broker 0 got {} users", counts[0]);
assert!(counts[1] >= 20, "Broker 1 got {} users", counts[1]);
}
#[test]
fn test_determine_authority_single_broker() {
let pm = PeerManager::new(Some("10.0.0.1"), &[], 9000, "key".into(), true);
match pm.determine_authority("9000000001") {
Authority::Standalone => {}
other => panic!("Expected Standalone, got {:?}", other),
}
}
#[test]
fn test_determine_authority_disabled() {
let pm = PeerManager::new(
Some("10.0.0.1"),
&["10.0.0.2:3960".to_string()],
9000,
"key".into(),
false,
);
match pm.determine_authority("9000000001") {
Authority::Standalone => {}
other => panic!("Expected Standalone, got {:?}", other),
}
}
}