use parking_lot::RwLock;
use serde::Serialize;
use std::collections::HashMap;
use std::sync::Arc;
use tokio::time::Instant;
#[derive(Debug, Clone, Serialize)]
pub struct UpstreamHealthSnapshot {
pub address: String,
pub tag: Option<String>,
pub plugin: String,
pub status: String,
pub success_rate: f64,
pub avg_response_time_ms: f64,
pub queries: u64,
pub successes: u64,
pub failures: u64,
pub last_success: Option<String>,
}
struct RegistryEntry {
address: String,
tag: Option<String>,
plugin: String,
health_fn: Box<dyn Fn() -> UpstreamHealthData + Send + Sync>,
}
pub struct UpstreamHealthData {
pub queries: u64,
pub successes: u64,
pub failures: u64,
pub avg_response_time_us: u64,
pub last_success: Option<Instant>,
}
pub struct UpstreamRegistry {
entries: RwLock<HashMap<String, RegistryEntry>>,
}
impl UpstreamRegistry {
pub fn new() -> Self {
Self {
entries: RwLock::new(HashMap::new()),
}
}
pub fn register<F>(
&self,
key: String,
address: String,
tag: Option<String>,
plugin: String,
health_fn: F,
) where
F: Fn() -> UpstreamHealthData + Send + Sync + 'static,
{
let entry = RegistryEntry {
address,
tag,
plugin,
health_fn: Box::new(health_fn),
};
self.entries.write().insert(key, entry);
}
pub fn unregister(&self, key: &str) {
self.entries.write().remove(key);
}
pub fn get_all(&self) -> Vec<UpstreamHealthSnapshot> {
let entries = self.entries.read();
entries
.values()
.map(|entry| {
let data = (entry.health_fn)();
let success_rate = if data.queries > 0 {
(data.successes as f64 / data.queries as f64) * 100.0
} else {
100.0
};
let status = if data.queries == 0 {
"unknown".to_string()
} else if success_rate >= 95.0 {
"healthy".to_string()
} else if success_rate >= 50.0 {
"degraded".to_string()
} else {
"unhealthy".to_string()
};
let last_success = data.last_success.map(|instant| {
let elapsed_secs = instant.elapsed().as_secs();
if elapsed_secs < 60 {
format!("{}s ago", elapsed_secs)
} else if elapsed_secs < 3600 {
format!("{}m ago", elapsed_secs / 60)
} else if elapsed_secs < 86400 {
format!("{}h ago", elapsed_secs / 3600)
} else {
format!("{}d ago", elapsed_secs / 86400)
}
});
UpstreamHealthSnapshot {
address: entry.address.clone(),
tag: entry.tag.clone(),
plugin: entry.plugin.clone(),
status,
success_rate,
avg_response_time_ms: data.avg_response_time_us as f64 / 1000.0,
queries: data.queries,
successes: data.successes,
failures: data.failures,
last_success,
}
})
.collect()
}
}
impl Default for UpstreamRegistry {
fn default() -> Self {
Self::new()
}
}
static UPSTREAM_REGISTRY: once_cell::sync::Lazy<Arc<UpstreamRegistry>> =
once_cell::sync::Lazy::new(|| Arc::new(UpstreamRegistry::new()));
pub fn upstream_registry() -> Arc<UpstreamRegistry> {
Arc::clone(&UPSTREAM_REGISTRY)
}
pub fn register_upstream<F>(
key: String,
address: String,
tag: Option<String>,
plugin: String,
health_fn: F,
) where
F: Fn() -> UpstreamHealthData + Send + Sync + 'static,
{
UPSTREAM_REGISTRY.register(key, address, tag, plugin, health_fn);
}
pub fn unregister_upstream(key: &str) {
UPSTREAM_REGISTRY.unregister(key);
}
pub fn get_all_upstream_health() -> Vec<UpstreamHealthSnapshot> {
UPSTREAM_REGISTRY.get_all()
}