use std::sync::Arc;
use std::time::Duration;
use tokio::sync::RwLock;
use tracing::{debug, info, warn};
use trusty_common::console_metrics::ConsoleMetricsReport;
use crate::mcp_handle::McpServiceHandle;
#[derive(Clone, Debug)]
pub struct MetricsCache {
inner: Arc<RwLock<Option<ConsoleMetricsReport>>>,
}
impl Default for MetricsCache {
fn default() -> Self {
Self::new()
}
}
impl MetricsCache {
pub fn new() -> Self {
Self {
inner: Arc::new(RwLock::new(None)),
}
}
pub async fn get(&self) -> Option<ConsoleMetricsReport> {
self.inner.read().await.clone()
}
pub async fn set(&self, report: ConsoleMetricsReport) {
*self.inner.write().await = Some(report);
}
}
pub(crate) async fn poll_once(handle: &McpServiceHandle, cache: &MetricsCache) {
match handle.poll_metrics().await {
Ok(report) => {
debug!(
service_id = %report.service_id,
status = ?report.status,
"metrics_poller: poll succeeded"
);
cache.set(report).await;
}
Err(e) => {
warn!(error = %e, "metrics_poller: poll failed — retaining previous cache");
}
}
}
pub fn start(handle: Arc<McpServiceHandle>, cache: MetricsCache, interval: Duration) {
tokio::spawn(async move {
info!(
"metrics_poller: starting (interval={}s)",
interval.as_secs()
);
loop {
poll_once(&handle, &cache).await;
tokio::time::sleep(interval).await;
}
});
}
#[cfg(test)]
mod tests {
use super::*;
use trusty_common::console_metrics::{ServiceHealth, make_report};
#[tokio::test]
async fn test_metrics_cache_initialises_empty() {
let cache = MetricsCache::new();
assert!(cache.get().await.is_none(), "cache must start empty");
}
#[tokio::test]
async fn test_metrics_cache_write_read_roundtrip() {
let cache = MetricsCache::new();
let report = make_report(
"trusty-analyze",
"Trusty Analyze",
"0.7.0",
ServiceHealth::Ok,
serde_json::json!({ "search_reachable": true }),
1,
);
cache.set(report.clone()).await;
let got = cache.get().await.expect("must have report after set");
assert_eq!(got.service_id, "trusty-analyze");
assert_eq!(got.display_name, "Trusty Analyze");
assert_eq!(got.version, "0.7.0");
assert_eq!(got.status, ServiceHealth::Ok);
assert_eq!(got.metrics["search_reachable"], true);
assert_eq!(got.metrics_schema_version, 1);
}
#[tokio::test]
async fn test_metrics_cache_overwrite() {
let cache = MetricsCache::new();
cache
.set(make_report(
"trusty-analyze",
"Trusty Analyze",
"0.6.0",
ServiceHealth::Degraded,
serde_json::json!({}),
1,
))
.await;
cache
.set(make_report(
"trusty-analyze",
"Trusty Analyze",
"0.7.0",
ServiceHealth::Ok,
serde_json::json!({ "search_reachable": true }),
1,
))
.await;
let got = cache.get().await.expect("must have report");
assert_eq!(got.version, "0.7.0");
assert_eq!(got.status, ServiceHealth::Ok);
}
}