#![allow(dead_code)]
use std::collections::VecDeque;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use crate::broker::stats::{BrokerMetrics, StatsResponse};
use crate::tui::views::snapshot::{Connection, Snapshot, SourceMode, WorkerKind, WorkerRow};
use super::DataSource;
#[derive(Default)]
pub struct RemoteShared {
pub stats: Option<StatsResponse>,
pub last_error: Option<String>,
pub event_log: VecDeque<String>,
}
pub struct RemoteSource {
pub broker_url: String,
pub shared: Arc<Mutex<RemoteShared>>,
pub refresh_requested: Arc<AtomicBool>,
}
impl RemoteSource {
pub fn new(broker_url: String) -> Self {
Self {
broker_url,
shared: Arc::new(Mutex::new(RemoteShared::default())),
refresh_requested: Arc::new(AtomicBool::new(false)),
}
}
}
impl DataSource for RemoteSource {
fn snapshot(&self) -> Snapshot {
let (stats, last_error, events) = {
let g = self.shared.lock().unwrap();
(
g.stats.clone(),
g.last_error.clone(),
g.event_log.iter().cloned().collect::<Vec<_>>(),
)
};
if let Some(s) = stats {
let workers: Vec<WorkerRow> = s
.workers
.iter()
.map(|w| WorkerRow {
status: WorkerKind::from_str_status(&w.status),
name: w.name.clone(),
cpus_available: w.cpus_available,
memory_gib: w.memory_available_gib,
avg_latency_ms: w.avg_latency_ms,
})
.collect();
Snapshot {
title: "Zakuro Compute Broker".to_string(),
mode: SourceMode::Remote,
host_port: format!("{}:{}", s.host, s.port),
wireguard_ip: s.wireguard_ip.clone(),
ledger_connected: s.metrics.ledger_connected,
connection: Connection::Connected,
metrics: s.metrics.clone(),
workers,
transactions: s.transactions.clone(),
task_offers: s.task_offers.clone(),
rps_history: s.rps_history.clone(),
events,
}
} else {
let connection = match last_error {
Some(e) => Connection::Error(e),
None => Connection::Connecting,
};
Snapshot {
title: "Zakuro Compute Broker".to_string(),
mode: SourceMode::Remote,
host_port: self.broker_url.clone(),
wireguard_ip: None,
ledger_connected: false,
connection,
metrics: BrokerMetrics::default(),
workers: Vec::new(),
transactions: Vec::new(),
task_offers: Vec::new(),
rps_history: Vec::new(),
events,
}
}
}
fn request_refresh(&self) {
self.refresh_requested.store(true, Ordering::Relaxed);
}
}