use crate::error::StreamError;
use dashmap::DashMap;
use std::sync::Arc;
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub enum HealthStatus {
Healthy,
Stale,
Unknown,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct FeedHealth {
pub feed_id: String,
pub status: HealthStatus,
pub last_tick_ms: Option<u64>,
pub stale_threshold_ms: u64,
pub tick_count: u64,
pub consecutive_stale: u32,
}
impl FeedHealth {
pub fn elapsed_ms(&self, now_ms: u64) -> Option<u64> {
self.last_tick_ms.map(|t| now_ms.saturating_sub(t))
}
pub fn is_healthy(&self) -> bool {
self.status == HealthStatus::Healthy
}
pub fn is_stale(&self) -> bool {
self.status == HealthStatus::Stale
}
}
pub struct HealthMonitor {
feeds: Arc<DashMap<String, FeedHealth>>,
default_stale_threshold_ms: u64,
circuit_breaker_threshold: u32,
}
impl HealthMonitor {
pub fn new(default_stale_threshold_ms: u64) -> Self {
Self {
feeds: Arc::new(DashMap::new()),
default_stale_threshold_ms,
circuit_breaker_threshold: 3,
}
}
pub fn with_circuit_breaker_threshold(mut self, threshold: u32) -> Self {
self.circuit_breaker_threshold = threshold;
self
}
pub fn is_circuit_open(&self, feed_id: &str) -> bool {
if self.circuit_breaker_threshold == 0 {
return false;
}
self.feeds
.get(feed_id)
.map_or(false, |e| e.consecutive_stale >= self.circuit_breaker_threshold)
}
pub fn register_many(&self, ids: &[&str], stale_threshold_ms: Option<u64>) {
for id in ids {
self.register(*id, stale_threshold_ms);
}
}
pub fn register(&self, feed_id: impl Into<String>, stale_threshold_ms: Option<u64>) {
let id = feed_id.into();
let threshold = stale_threshold_ms.unwrap_or(self.default_stale_threshold_ms);
self.feeds.insert(
id.clone(),
FeedHealth {
feed_id: id,
status: HealthStatus::Unknown,
last_tick_ms: None,
stale_threshold_ms: threshold,
tick_count: 0,
consecutive_stale: 0,
},
);
}
pub fn deregister(&self, feed_id: &str) -> Option<FeedHealth> {
self.feeds.remove(feed_id).map(|(_, v)| v)
}
pub fn heartbeat(&self, feed_id: &str, ts_ms: u64) -> Result<(), StreamError> {
let mut entry = self
.feeds
.get_mut(feed_id)
.ok_or_else(|| StreamError::UnknownFeed {
feed_id: feed_id.to_string(),
})?;
entry.last_tick_ms = Some(ts_ms);
entry.tick_count += 1;
entry.status = HealthStatus::Healthy;
entry.consecutive_stale = 0;
Ok(())
}
pub fn check_all(&self, now_ms: u64) -> Vec<(String, StreamError)> {
let mut errors = Vec::new();
for mut entry in self.feeds.iter_mut() {
let elapsed = entry.elapsed_ms(now_ms);
if let Some(elapsed) = elapsed {
if elapsed > entry.stale_threshold_ms {
entry.status = HealthStatus::Stale;
entry.consecutive_stale += 1;
let feed_id = entry.feed_id.clone();
errors.push((
feed_id.clone(),
StreamError::StaleFeed {
feed_id,
elapsed_ms: elapsed,
threshold_ms: entry.stale_threshold_ms,
},
));
}
}
}
errors
}
pub fn get(&self, feed_id: &str) -> Option<FeedHealth> {
self.feeds.get(feed_id).map(|e| e.value().clone())
}
pub fn all_feeds(&self) -> Vec<FeedHealth> {
self.feeds.iter().map(|e| e.value().clone()).collect()
}
pub fn feed_count(&self) -> usize {
self.feeds.len()
}
pub fn check_one(
&self,
feed_id: &str,
now_ms: u64,
) -> Result<Option<StreamError>, StreamError> {
let mut entry = self
.feeds
.get_mut(feed_id)
.ok_or_else(|| StreamError::UnknownFeed {
feed_id: feed_id.to_string(),
})?;
let elapsed = match entry.last_tick_ms {
Some(t) => now_ms.saturating_sub(t),
None => return Ok(None),
};
if elapsed > entry.stale_threshold_ms {
entry.status = HealthStatus::Stale;
entry.consecutive_stale += 1;
Ok(Some(StreamError::StaleFeed {
feed_id: entry.feed_id.clone(),
elapsed_ms: elapsed,
threshold_ms: entry.stale_threshold_ms,
}))
} else {
Ok(None)
}
}
pub fn reset_feed(&self, feed_id: &str) -> Result<(), StreamError> {
let mut entry = self
.feeds
.get_mut(feed_id)
.ok_or_else(|| StreamError::UnknownFeed {
feed_id: feed_id.to_string(),
})?;
entry.status = HealthStatus::Unknown;
entry.last_tick_ms = None;
entry.tick_count = 0;
entry.consecutive_stale = 0;
Ok(())
}
pub fn feed_ids(&self) -> Vec<String> {
let mut ids: Vec<String> = self.feeds.iter().map(|e| e.feed_id.clone()).collect();
ids.sort();
ids
}
pub fn healthy_count(&self) -> usize {
self.feeds_by_status(HealthStatus::Healthy).len()
}
pub fn stale_ratio(&self) -> f64 {
let total = self.feed_count();
if total == 0 {
return 0.0;
}
self.stale_count() as f64 / total as f64
}
pub fn is_any_stale(&self) -> bool {
self.stale_count() > 0
}
pub fn stale_count(&self) -> usize {
self.feeds_by_status(HealthStatus::Stale).len()
}
pub fn stale_feeds(&self) -> Vec<FeedHealth> {
self.feeds_by_status(HealthStatus::Stale)
}
pub fn oldest_tick_ms(&self) -> Option<u64> {
self.feeds
.iter()
.filter_map(|e| e.last_tick_ms)
.min()
}
pub fn newest_tick_ms(&self) -> Option<u64> {
self.feeds
.iter()
.filter_map(|e| e.last_tick_ms)
.max()
}
pub fn total_tick_count(&self) -> u64 {
self.feeds.iter().map(|e| e.tick_count).sum()
}
pub fn lag_ms(&self) -> Option<u64> {
let newest = self.newest_tick_ms()?;
let oldest = self.oldest_tick_ms()?;
Some(newest.saturating_sub(oldest))
}
pub fn most_stale_feed(&self) -> Option<FeedHealth> {
self.feeds.iter().fold(None, |acc: Option<FeedHealth>, entry| {
let feed = entry.clone();
match acc {
None => Some(feed),
Some(current) => {
let more_stale = match (feed.last_tick_ms, current.last_tick_ms) {
(None, _) => true,
(Some(_), None) => false,
(Some(a), Some(b)) => a < b,
};
if more_stale { Some(feed) } else { Some(current) }
}
}
})
}
pub fn stale_ratio_excluding_unknown(&self) -> f64 {
let known: Vec<_> = self.feeds.iter()
.filter(|e| e.status != HealthStatus::Unknown)
.collect();
if known.is_empty() { return 0.0; }
let stale = known.iter().filter(|e| e.status == HealthStatus::Stale).count();
stale as f64 / known.len() as f64
}
pub fn healthy_feeds(&self) -> Vec<String> {
let mut ids: Vec<String> = self
.feeds
.iter()
.filter(|e| e.status == HealthStatus::Healthy)
.map(|e| e.feed_id.clone())
.collect();
ids.sort();
ids
}
pub fn unhealthy_feeds(&self) -> Vec<String> {
let mut ids: Vec<String> = self
.feeds
.iter()
.filter(|e| e.status != HealthStatus::Healthy)
.map(|e| e.feed_id.clone())
.collect();
ids.sort();
ids
}
pub fn reset_all(&self) {
for mut entry in self.feeds.iter_mut() {
entry.status = HealthStatus::Unknown;
entry.last_tick_ms = None;
entry.consecutive_stale = 0;
}
}
pub fn all_healthy(&self) -> bool {
self.feeds.iter().all(|e| e.status == HealthStatus::Healthy)
}
pub fn ratio_healthy(&self) -> f64 {
let total = self.feed_count();
if total == 0 {
return 0.0;
}
self.healthy_count() as f64 / total as f64
}
pub fn last_updated_feed_id(&self) -> Option<String> {
self.feeds
.iter()
.filter_map(|e| e.last_tick_ms.map(|t| (t, e.feed_id.clone())))
.max_by_key(|(t, _)| *t)
.map(|(_, id)| id)
}
pub fn unknown_feed_ids(&self) -> Vec<String> {
self.feeds
.iter()
.filter(|e| e.status == HealthStatus::Unknown)
.map(|e| e.feed_id.clone())
.collect()
}
pub fn status_summary(&self) -> (usize, usize, usize) {
let (mut healthy, mut stale, mut unknown) = (0, 0, 0);
for e in self.feeds.iter() {
match e.status {
HealthStatus::Healthy => healthy += 1,
HealthStatus::Stale => stale += 1,
HealthStatus::Unknown => unknown += 1,
}
}
(healthy, stale, unknown)
}
pub fn stale_feed_ids(&self) -> Vec<String> {
let mut ids: Vec<String> = self
.feeds
.iter()
.filter(|e| e.status == HealthStatus::Stale)
.map(|e| e.feed_id.clone())
.collect();
ids.sort();
ids
}
#[deprecated(since = "2.2.0", note = "Use `stale_count` instead")]
pub fn total_stale_count(&self) -> usize {
self.stale_count()
}
#[deprecated(since = "2.2.0", note = "Use `unhealthy_feeds` instead")]
pub fn feeds_needing_check(&self) -> Vec<String> {
self.unhealthy_feeds()
}
pub fn avg_feed_age_ms(&self, now_ms: u64) -> Option<f64> {
let ages: Vec<u64> = self
.feeds
.iter()
.filter_map(|e| e.last_tick_ms)
.map(|t| now_ms.saturating_sub(t))
.collect();
if ages.is_empty() {
return None;
}
Some(ages.iter().sum::<u64>() as f64 / ages.len() as f64)
}
pub fn unknown_count(&self) -> usize {
self.feeds_by_status(HealthStatus::Unknown).len()
}
pub fn avg_tick_count(&self) -> f64 {
let count = self.feed_count();
if count == 0 {
return 0.0;
}
self.total_tick_count() as f64 / count as f64
}
pub fn max_consecutive_stale(&self) -> u32 {
self.feeds
.iter()
.map(|e| e.consecutive_stale)
.max()
.unwrap_or(0)
}
#[deprecated(since = "2.2.0", note = "Use `unknown_count` instead")]
pub fn unknown_feed_count(&self) -> usize {
self.unknown_count()
}
pub fn feeds_by_status(&self, status: HealthStatus) -> Vec<FeedHealth> {
self.feeds
.iter()
.filter(|e| e.value().status == status)
.map(|e| e.value().clone())
.collect()
}
pub fn oldest_stale_feed(&self) -> Option<FeedHealth> {
self.feeds
.iter()
.filter(|e| e.value().status == HealthStatus::Stale)
.min_by_key(|e| e.value().last_tick_ms.unwrap_or(u64::MAX))
.map(|e| e.value().clone())
}
#[deprecated(since = "2.2.0", note = "Use `ratio_healthy` instead")]
pub fn healthy_ratio(&self) -> f64 {
self.ratio_healthy()
}
pub fn most_reliable_feed(&self) -> Option<FeedHealth> {
self.feeds.iter()
.max_by_key(|e| e.value().tick_count)
.map(|e| e.value().clone())
}
pub fn feeds_never_seen(&self) -> Vec<FeedHealth> {
self.feeds.iter()
.filter(|e| e.value().last_tick_ms.is_none())
.map(|e| e.value().clone())
.collect()
}
#[deprecated(since = "2.2.0", note = "Use `is_any_stale` instead")]
pub fn is_any_feed_stale(&self) -> bool {
self.is_any_stale()
}
pub fn all_feeds_seen(&self) -> bool {
self.feeds.iter().all(|e| e.last_tick_ms.is_some())
}
pub fn tick_count_for(&self, feed_id: &str) -> Option<u64> {
self.feeds.iter()
.find(|e| e.feed_id == feed_id)
.map(|e| e.tick_count)
}
#[deprecated(since = "2.2.0", note = "Use `avg_tick_count` instead")]
pub fn average_tick_count(&self) -> f64 {
self.avg_tick_count()
}
pub fn feeds_above_tick_count(&self, threshold: u64) -> usize {
self.feeds.iter().filter(|e| e.tick_count > threshold).count()
}
pub fn oldest_feed_age_ms(&self, now_ms: u64) -> Option<u64> {
self.feeds
.iter()
.filter_map(|e| e.last_tick_ms)
.map(|t| now_ms.saturating_sub(t))
.max()
}
pub fn has_any_unknown(&self) -> bool {
self.feeds.iter().any(|e| e.status == HealthStatus::Unknown)
}
pub fn is_degraded(&self) -> bool {
let total = self.feed_count();
if total == 0 {
return false;
}
let healthy = self.healthy_count();
healthy > 0 && healthy < total
}
pub fn unhealthy_count(&self) -> usize {
self.feed_count().saturating_sub(self.healthy_count())
}
pub fn feed_exists(&self, feed_id: &str) -> bool {
self.feeds.iter().any(|e| e.feed_id == feed_id)
}
#[deprecated(since = "2.2.0", note = "Use `has_any_unknown` instead")]
pub fn any_unknown(&self) -> bool {
self.has_any_unknown()
}
#[deprecated(since = "2.2.0", note = "Use `stale_count` instead")]
pub fn degraded_count(&self) -> usize {
self.stale_count()
}
pub fn time_since_last_heartbeat(&self, feed_id: &str, now_ms: u64) -> Option<u64> {
self.feeds
.iter()
.find(|e| e.feed_id == feed_id)?
.last_tick_ms
.map(|t| now_ms.saturating_sub(t))
}
pub fn healthy_feed_ids(&self) -> Vec<String> {
self.feeds
.iter()
.filter(|e| e.status == HealthStatus::Healthy)
.map(|e| e.feed_id.clone())
.collect()
}
pub fn register_batch(&self, feeds: &[(&str, u64)]) {
for (id, threshold) in feeds {
self.register(*id, Some(*threshold));
}
}
pub fn min_healthy_age_ms(&self, now_ms: u64) -> Option<u64> {
self.feeds
.iter()
.filter(|e| e.status == HealthStatus::Healthy)
.filter_map(|e| e.last_tick_ms)
.map(|t| now_ms.saturating_sub(t))
.min()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn monitor() -> HealthMonitor {
HealthMonitor::new(5_000)
}
#[test]
fn test_register_creates_unknown_feed() {
let m = monitor();
m.register("BTC-USD", None);
let h = m.get("BTC-USD").unwrap();
assert_eq!(h.status, HealthStatus::Unknown);
assert!(h.last_tick_ms.is_none());
}
#[test]
fn test_heartbeat_marks_feed_healthy() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
let h = m.get("BTC-USD").unwrap();
assert_eq!(h.status, HealthStatus::Healthy);
assert_eq!(h.last_tick_ms, Some(1_000_000));
}
#[test]
fn test_heartbeat_increments_tick_count() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1000).unwrap();
m.heartbeat("BTC-USD", 2000).unwrap();
m.heartbeat("BTC-USD", 3000).unwrap();
assert_eq!(m.get("BTC-USD").unwrap().tick_count, 3);
}
#[test]
fn test_heartbeat_unknown_feed_returns_unknown_feed_error() {
let m = monitor();
let result = m.heartbeat("ghost", 1000);
assert!(matches!(result, Err(StreamError::UnknownFeed { .. })));
}
#[test]
fn test_check_all_healthy_feed_no_errors() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
let errors = m.check_all(1_003_000); assert!(errors.is_empty());
}
#[test]
fn test_check_all_stale_feed_returns_error() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
let errors = m.check_all(1_010_000); assert_eq!(errors.len(), 1);
assert_eq!(errors[0].0, "BTC-USD");
assert!(matches!(&errors[0].1, StreamError::StaleFeed { feed_id, .. } if feed_id == "BTC-USD"));
}
#[test]
fn test_check_all_marks_stale_in_state() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
m.check_all(1_010_000);
assert_eq!(m.get("BTC-USD").unwrap().status, HealthStatus::Stale);
}
#[test]
fn test_check_all_unknown_feed_not_counted_as_stale() {
let m = monitor();
m.register("BTC-USD", None);
let errors = m.check_all(9_999_999);
assert!(errors.is_empty());
}
#[test]
fn test_custom_threshold_per_feed() {
let m = monitor();
m.register("BTC-USD", Some(1_000)); m.heartbeat("BTC-USD", 1_000_000).unwrap();
let errors = m.check_all(1_002_000); assert!(!errors.is_empty());
}
#[test]
fn test_feed_count() {
let m = monitor();
m.register("BTC-USD", None);
m.register("ETH-USD", None);
assert_eq!(m.feed_count(), 2);
}
#[test]
fn test_healthy_count_and_stale_count() {
let m = monitor();
m.register("BTC-USD", None);
m.register("ETH-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
m.heartbeat("ETH-USD", 1_000_000).unwrap();
m.check_all(1_010_000); assert_eq!(m.stale_count(), 2);
assert_eq!(m.healthy_count(), 0);
}
#[test]
fn test_feed_health_elapsed_ms() {
let h = FeedHealth {
feed_id: "BTC-USD".into(),
status: HealthStatus::Healthy,
last_tick_ms: Some(1_000_000),
stale_threshold_ms: 5_000,
tick_count: 1,
consecutive_stale: 0,
};
assert_eq!(h.elapsed_ms(1_003_000), Some(3_000));
}
#[test]
fn test_feed_health_elapsed_ms_none_when_no_last_tick() {
let h = FeedHealth {
feed_id: "X".into(),
status: HealthStatus::Unknown,
last_tick_ms: None,
stale_threshold_ms: 5_000,
tick_count: 0,
consecutive_stale: 0,
};
assert!(h.elapsed_ms(9_999_999).is_none());
}
#[test]
fn test_all_feeds_returns_all() {
let m = monitor();
m.register("A", None);
m.register("B", None);
let feeds = m.all_feeds();
assert_eq!(feeds.len(), 2);
}
#[test]
fn test_circuit_not_open_before_threshold() {
let m = monitor(); m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
m.check_all(1_010_000); m.check_all(1_020_000); assert!(!m.is_circuit_open("BTC-USD"));
}
#[test]
fn test_circuit_opens_at_threshold() {
let m = monitor(); m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
m.check_all(1_010_000); m.check_all(1_020_000); m.check_all(1_030_000); assert!(m.is_circuit_open("BTC-USD"));
}
#[test]
fn test_circuit_resets_on_heartbeat() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
m.check_all(1_010_000);
m.check_all(1_020_000);
m.check_all(1_030_000);
assert!(m.is_circuit_open("BTC-USD"));
m.heartbeat("BTC-USD", 1_040_000).unwrap();
assert!(!m.is_circuit_open("BTC-USD"));
}
#[test]
fn test_circuit_disabled_when_threshold_zero() {
let m = HealthMonitor::new(5_000).with_circuit_breaker_threshold(0);
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
for i in 0..10 {
m.check_all(1_010_000 + i * 10_000);
}
assert!(!m.is_circuit_open("BTC-USD"));
}
#[test]
fn test_circuit_open_returns_false_for_unknown_feed() {
let m = monitor();
assert!(!m.is_circuit_open("ghost"));
}
#[test]
fn test_deregister_removes_feed() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
let removed = m.deregister("BTC-USD");
assert!(removed.is_some());
assert_eq!(removed.unwrap().feed_id, "BTC-USD");
assert!(m.get("BTC-USD").is_none());
assert_eq!(m.feed_count(), 0);
}
#[test]
fn test_deregister_unknown_feed_returns_none() {
let m = monitor();
assert!(m.deregister("ghost").is_none());
}
#[test]
fn test_feed_ids_returns_sorted_ids() {
let m = monitor();
m.register("ETH-USD", None);
m.register("BTC-USD", None);
m.register("SOL-USD", None);
let ids = m.feed_ids();
assert_eq!(ids, vec!["BTC-USD", "ETH-USD", "SOL-USD"]);
}
#[test]
fn test_feed_ids_empty_when_no_feeds() {
let m = monitor();
assert!(m.feed_ids().is_empty());
}
#[test]
fn test_feed_ids_updates_after_deregister() {
let m = monitor();
m.register("BTC-USD", None);
m.register("ETH-USD", None);
m.deregister("BTC-USD");
assert_eq!(m.feed_ids(), vec!["ETH-USD"]);
}
#[test]
fn test_reset_feed_clears_state() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
m.check_all(1_010_000); assert_eq!(m.get("BTC-USD").unwrap().status, HealthStatus::Stale);
m.reset_feed("BTC-USD").unwrap();
let h = m.get("BTC-USD").unwrap();
assert_eq!(h.status, HealthStatus::Unknown);
assert_eq!(h.tick_count, 0);
assert_eq!(h.consecutive_stale, 0);
assert!(h.last_tick_ms.is_none());
}
#[test]
fn test_reset_feed_unknown_returns_error() {
let m = monitor();
assert!(matches!(
m.reset_feed("ghost"),
Err(StreamError::UnknownFeed { .. })
));
}
#[test]
fn test_check_one_healthy_feed_returns_none() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
assert!(m.check_one("BTC-USD", 1_003_000).unwrap().is_none());
}
#[test]
fn test_check_one_stale_feed_returns_some_error() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
let result = m.check_one("BTC-USD", 1_010_000).unwrap();
assert!(matches!(result, Some(StreamError::StaleFeed { .. })));
}
#[test]
fn test_check_one_unknown_feed_returns_err() {
let m = monitor();
assert!(matches!(
m.check_one("ghost", 0),
Err(StreamError::UnknownFeed { .. })
));
}
#[test]
fn test_heartbeat_after_deregister_returns_unknown_feed_error() {
let m = monitor();
m.register("BTC-USD", None);
m.deregister("BTC-USD");
let result = m.heartbeat("BTC-USD", 1_000_000);
assert!(matches!(result, Err(StreamError::UnknownFeed { .. })));
}
#[test]
fn test_unhealthy_feeds_returns_non_healthy_sorted() {
let m = HealthMonitor::new(5_000);
m.register("BTC-USD", None);
m.register("ETH-USD", None);
m.register("SOL-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
let unhealthy = m.unhealthy_feeds();
assert!(unhealthy.contains(&"ETH-USD".to_string()));
assert!(unhealthy.contains(&"SOL-USD".to_string()));
assert!(!unhealthy.contains(&"BTC-USD".to_string()));
assert_eq!(unhealthy[0], "ETH-USD");
assert_eq!(unhealthy[1], "SOL-USD");
}
#[test]
fn test_oldest_tick_ms_returns_minimum() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000_000).unwrap();
m.heartbeat("B", 2_000_000).unwrap();
assert_eq!(m.oldest_tick_ms(), Some(1_000_000));
}
#[test]
fn test_newest_tick_ms_returns_maximum() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000_000).unwrap();
m.heartbeat("B", 2_000_000).unwrap();
assert_eq!(m.newest_tick_ms(), Some(2_000_000));
}
#[test]
fn test_oldest_newest_tick_ms_none_when_no_ticks() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
assert!(m.oldest_tick_ms().is_none());
assert!(m.newest_tick_ms().is_none());
}
#[test]
fn test_unhealthy_feeds_empty_when_all_healthy() {
let m = HealthMonitor::new(5_000);
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
assert!(m.unhealthy_feeds().is_empty());
}
#[test]
fn test_unhealthy_feeds_includes_stale_feeds() {
let m = HealthMonitor::new(1_000); m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
m.check_all(1_002_000); let unhealthy = m.unhealthy_feeds();
assert!(unhealthy.contains(&"BTC-USD".to_string()));
}
#[test]
fn test_all_healthy_vacuously_true_no_feeds() {
let m = HealthMonitor::new(5_000);
assert!(m.all_healthy());
}
#[test]
fn test_all_healthy_false_when_unknown() {
let m = HealthMonitor::new(5_000);
m.register("BTC-USD", None);
assert!(!m.all_healthy());
}
#[test]
fn test_all_healthy_true_after_heartbeats() {
let m = HealthMonitor::new(5_000);
m.register("BTC-USD", None);
m.register("ETH-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
m.heartbeat("ETH-USD", 1_000_000).unwrap();
assert!(m.all_healthy());
}
#[test]
fn test_all_healthy_false_when_one_stale() {
let m = HealthMonitor::new(1_000);
m.register("BTC-USD", None);
m.register("ETH-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
m.heartbeat("ETH-USD", 1_000_000).unwrap();
m.check_all(1_002_000); assert!(!m.all_healthy());
}
#[test]
fn test_stale_feed_ids_empty_when_none_stale() {
let m = HealthMonitor::new(5_000);
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
assert!(m.stale_feed_ids().is_empty());
}
#[test]
fn test_stale_feed_ids_excludes_unknown() {
let m = HealthMonitor::new(5_000);
m.register("BTC-USD", None); assert!(m.stale_feed_ids().is_empty()); }
#[test]
fn test_stale_feed_ids_returns_only_stale() {
let m = HealthMonitor::new(1_000);
m.register("BTC-USD", None);
m.register("ETH-USD", Some(10_000)); m.heartbeat("BTC-USD", 1_000_000).unwrap();
m.heartbeat("ETH-USD", 1_000_000).unwrap();
m.check_all(1_002_000); let stale = m.stale_feed_ids();
assert_eq!(stale, vec!["BTC-USD".to_string()]);
}
#[test]
fn test_lag_ms_none_when_no_ticks() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
assert!(m.lag_ms().is_none());
}
#[test]
fn test_lag_ms_zero_when_one_feed_or_equal_timestamps() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000_000).unwrap();
m.heartbeat("B", 1_000_000).unwrap();
assert_eq!(m.lag_ms(), Some(0));
}
#[test]
fn test_lag_ms_returns_spread() {
let m = HealthMonitor::new(5_000);
m.register("fast", None);
m.register("slow", None);
m.heartbeat("fast", 2_000_000).unwrap();
m.heartbeat("slow", 1_000_000).unwrap();
assert_eq!(m.lag_ms(), Some(1_000_000));
}
#[test]
fn test_total_tick_count_zero_initially() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
assert_eq!(m.total_tick_count(), 0);
}
#[test]
fn test_total_tick_count_sums_across_feeds() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000_000).unwrap();
m.heartbeat("A", 1_001_000).unwrap();
m.heartbeat("B", 1_000_000).unwrap();
assert_eq!(m.total_tick_count(), 3);
}
#[test]
fn test_healthy_feeds_empty_initially() {
let m = HealthMonitor::new(5_000);
m.register("A", None); assert!(m.healthy_feeds().is_empty());
}
#[test]
fn test_healthy_feeds_after_heartbeat() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000_000).unwrap();
m.check_all(1_001_000); let healthy = m.healthy_feeds();
assert_eq!(healthy, vec!["A".to_string()]);
}
#[test]
fn test_healthy_feeds_sorted() {
let m = HealthMonitor::new(5_000);
m.register("C", None);
m.register("A", None);
m.register("B", None);
m.heartbeat("C", 1_000_000).unwrap();
m.heartbeat("A", 1_000_000).unwrap();
m.heartbeat("B", 1_000_000).unwrap();
m.check_all(1_001_000);
let healthy = m.healthy_feeds();
assert_eq!(healthy, vec!["A".to_string(), "B".to_string(), "C".to_string()]);
}
#[test]
fn test_unknown_count_all_new_feeds() {
let m = monitor();
m.register("A", None);
m.register("B", None);
assert_eq!(m.unknown_count(), 2);
}
#[test]
fn test_unknown_count_decreases_after_heartbeat() {
let m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
assert_eq!(m.unknown_count(), 1); }
#[test]
fn test_unknown_count_zero_when_all_healthy() {
let m = monitor();
m.register("A", None);
m.heartbeat("A", 1_000).unwrap();
assert_eq!(m.unknown_count(), 0);
}
#[test]
fn test_feed_health_is_healthy_true_after_heartbeat() {
let m = monitor();
m.register("X", None);
m.heartbeat("X", 1_000).unwrap();
let fh = m.get("X").unwrap();
assert!(fh.is_healthy());
}
#[test]
fn test_feed_health_is_healthy_false_when_unknown() {
let m = monitor();
m.register("X", None);
let fh = m.get("X").unwrap();
assert!(!fh.is_healthy());
}
#[test]
fn test_most_stale_feed_none_when_no_feeds() {
let m = HealthMonitor::new(5_000);
assert!(m.most_stale_feed().is_none());
}
#[test]
fn test_most_stale_feed_returns_unticked_feed_first() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000_000).unwrap();
let stale = m.most_stale_feed().unwrap();
assert_eq!(stale.feed_id, "B");
}
#[test]
fn test_most_stale_feed_returns_oldest_last_tick() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 5_000).unwrap();
let stale = m.most_stale_feed().unwrap();
assert_eq!(stale.feed_id, "A");
}
#[test]
fn test_stale_feeds_empty_when_all_healthy() {
let m = monitor();
m.register("A", None);
m.heartbeat("A", 1_000).unwrap();
assert!(m.stale_feeds().is_empty());
}
#[test]
fn test_stale_feeds_returns_all_stale() {
let m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 9_500).unwrap();
m.check_all(10_000);
let stale = m.stale_feeds();
assert_eq!(stale.len(), 1);
assert_eq!(stale[0].feed_id, "A");
}
#[test]
fn test_register_many_creates_all_feeds() {
let m = monitor();
m.register_many(&["BTC-USD", "ETH-USD", "SOL-USD"], None);
assert_eq!(m.feed_count(), 3);
}
#[test]
fn test_register_many_custom_threshold_applies() {
let m = monitor();
m.register_many(&["A", "B"], Some(1_000));
m.heartbeat("A", 1_000_000).unwrap();
let errors = m.check_all(1_002_000); assert!(!errors.is_empty());
}
#[test]
fn test_register_many_empty_slice_is_noop() {
let m = monitor();
m.register_many(&[], None);
assert_eq!(m.feed_count(), 0);
}
#[test]
fn test_feed_health_is_stale_true_when_stale() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
m.check_all(1_010_000); assert!(m.get("BTC-USD").unwrap().is_stale());
}
#[test]
fn test_feed_health_is_stale_false_when_healthy() {
let m = monitor();
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000_000).unwrap();
assert!(!m.get("BTC-USD").unwrap().is_stale());
}
#[test]
fn test_avg_tick_count_zero_with_no_feeds() {
let m = HealthMonitor::new(5_000);
assert!((m.avg_tick_count() - 0.0).abs() < 1e-9);
}
#[test]
fn test_avg_tick_count_with_equal_ticks() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 1_000).unwrap();
assert!((m.avg_tick_count() - 1.0).abs() < 1e-9);
}
#[test]
fn test_avg_tick_count_with_different_ticks() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("A", 2_000).unwrap();
m.heartbeat("A", 3_000).unwrap();
m.heartbeat("B", 1_000).unwrap();
assert!((m.avg_tick_count() - 2.0).abs() < 1e-9);
}
#[test]
fn test_max_consecutive_stale_zero_with_no_feeds() {
let m = HealthMonitor::new(5_000);
assert_eq!(m.max_consecutive_stale(), 0);
}
#[test]
fn test_max_consecutive_stale_picks_highest() {
let m = HealthMonitor::new(1_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 1_000).unwrap();
m.check_all(4_000);
m.check_all(5_000);
assert!(m.max_consecutive_stale() >= 2);
}
#[test]
fn test_max_consecutive_stale_zero_after_heartbeat() {
let m = HealthMonitor::new(1_000);
m.register("A", None);
m.check_all(10_000);
m.heartbeat("A", 10_001).unwrap();
assert_eq!(m.max_consecutive_stale(), 0);
}
#[test]
fn test_status_summary_all_unknown() {
let m = monitor();
m.register("A", None);
m.register("B", None);
let (healthy, stale, unknown) = m.status_summary();
assert_eq!(healthy, 0);
assert_eq!(stale, 0);
assert_eq!(unknown, 2);
}
#[test]
fn test_status_summary_mixed() {
let m = monitor();
m.register("A", None);
m.register("B", None);
m.register("C", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 5_000).unwrap();
m.check_all(10_000);
let (healthy, stale, unknown) = m.status_summary();
assert_eq!(stale, 1);
assert_eq!(unknown, 1);
assert_eq!(healthy, 1);
}
#[test]
fn test_feeds_by_status_returns_only_matching_status() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 9_000).unwrap();
m.check_all(10_000); let stale = m.feeds_by_status(HealthStatus::Stale);
assert_eq!(stale.len(), 1);
assert_eq!(stale[0].feed_id, "A");
}
#[test]
fn test_feeds_by_status_empty_when_none_match() {
let mut m = monitor();
m.register("A", None);
m.heartbeat("A", 1_000).unwrap();
m.check_all(2_000);
let unknown = m.feeds_by_status(HealthStatus::Unknown);
assert!(unknown.is_empty());
}
#[test]
fn test_feeds_by_status_all_start_as_unknown() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
let unknown = m.feeds_by_status(HealthStatus::Unknown);
assert_eq!(unknown.len(), 2);
}
#[test]
fn test_unknown_feed_count_all_unknown_at_start() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
assert_eq!(m.unknown_feed_count(), 2);
}
#[test]
fn test_unknown_feed_count_decreases_after_heartbeat() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
assert_eq!(m.unknown_feed_count(), 1);
}
#[test]
fn test_unknown_feed_count_zero_with_no_feeds() {
let m = HealthMonitor::new(5_000);
assert_eq!(m.unknown_feed_count(), 0);
}
#[test]
fn test_all_healthy_vacuously_true_with_no_feeds() {
let m = HealthMonitor::new(5_000);
assert!(m.all_healthy());
}
#[test]
fn test_all_healthy_true_when_all_feeds_healthy() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 1_000).unwrap();
assert!(m.all_healthy());
}
#[test]
fn test_all_healthy_false_when_one_feed_unknown() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
assert!(!m.all_healthy());
}
#[test]
fn test_oldest_stale_feed_returns_feed_with_smallest_last_tick() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap(); m.heartbeat("B", 3_000).unwrap(); m.check_all(10_000); let oldest = m.oldest_stale_feed().unwrap();
assert_eq!(oldest.feed_id, "A");
}
#[test]
fn test_oldest_stale_feed_none_when_no_stale_feeds() {
let mut m = monitor();
m.register("A", None);
m.heartbeat("A", 9_000).unwrap();
m.check_all(10_000); assert!(m.oldest_stale_feed().is_none());
}
#[test]
fn test_healthy_ratio_zero_when_no_feeds() {
let m = monitor();
assert_eq!(m.healthy_ratio(), 0.0);
}
#[test]
fn test_healthy_ratio_one_when_all_healthy() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 9_000).unwrap();
m.heartbeat("B", 9_500).unwrap();
m.check_all(10_000); assert!((m.healthy_ratio() - 1.0).abs() < 1e-10);
}
#[test]
fn test_healthy_ratio_half_when_one_of_two_healthy() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap(); m.heartbeat("B", 9_500).unwrap(); m.check_all(10_000);
assert!((m.healthy_ratio() - 0.5).abs() < 1e-10);
}
#[test]
fn test_most_reliable_feed_returns_highest_tick_count() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 1_100).unwrap();
m.heartbeat("B", 1_200).unwrap();
let best = m.most_reliable_feed().unwrap();
assert_eq!(best.feed_id, "B");
}
#[test]
fn test_most_reliable_feed_none_when_no_feeds() {
let m = monitor();
assert!(m.most_reliable_feed().is_none());
}
#[test]
fn test_feeds_never_seen_returns_feeds_with_no_heartbeat() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
let never_seen = m.feeds_never_seen();
assert_eq!(never_seen.len(), 1);
assert_eq!(never_seen[0].feed_id, "B");
}
#[test]
fn test_feeds_never_seen_empty_when_all_have_heartbeat() {
let mut m = monitor();
m.register("A", None);
m.heartbeat("A", 1_000).unwrap();
assert!(m.feeds_never_seen().is_empty());
}
#[test]
fn test_is_any_feed_stale_true_when_stale_feed_exists() {
let mut m = monitor();
m.register("A", None);
m.heartbeat("A", 1_000).unwrap();
m.check_all(10_000); assert!(m.is_any_feed_stale());
}
#[test]
fn test_is_any_feed_stale_false_when_all_healthy() {
let mut m = monitor();
m.register("A", None);
m.heartbeat("A", 9_500).unwrap();
m.check_all(10_000); assert!(!m.is_any_feed_stale());
}
#[test]
fn test_is_any_feed_stale_false_when_no_feeds() {
let m = monitor();
assert!(!m.is_any_feed_stale());
}
#[test]
fn test_all_feeds_seen_true_when_all_have_heartbeat() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 2_000).unwrap();
assert!(m.all_feeds_seen());
}
#[test]
fn test_all_feeds_seen_false_when_one_never_seen() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
assert!(!m.all_feeds_seen());
}
#[test]
fn test_all_feeds_seen_true_vacuously_when_no_feeds() {
let m = monitor();
assert!(m.all_feeds_seen());
}
#[test]
fn test_tick_count_for_returns_correct_count() {
let mut m = monitor();
m.register("A", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("A", 2_000).unwrap();
assert_eq!(m.tick_count_for("A"), Some(2));
}
#[test]
fn test_tick_count_for_none_when_not_registered() {
let m = monitor();
assert!(m.tick_count_for("nonexistent").is_none());
}
#[test]
fn test_tick_count_for_zero_when_no_heartbeats() {
let mut m = monitor();
m.register("A", None);
assert_eq!(m.tick_count_for("A"), Some(0));
}
#[test]
fn test_average_tick_count_zero_when_no_feeds() {
let m = monitor();
assert_eq!(m.average_tick_count(), 0.0);
}
#[test]
fn test_average_tick_count_correct_value() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("A", 2_000).unwrap(); m.heartbeat("B", 1_000).unwrap(); assert!((m.average_tick_count() - 1.5).abs() < 1e-10);
}
#[test]
fn test_feeds_above_tick_count_correct() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.register("C", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("A", 2_000).unwrap();
m.heartbeat("A", 3_000).unwrap(); m.heartbeat("B", 1_000).unwrap(); assert_eq!(m.feeds_above_tick_count(1), 1);
assert_eq!(m.feeds_above_tick_count(0), 2);
assert_eq!(m.feeds_above_tick_count(5), 0);
}
#[test]
fn test_feeds_above_tick_count_zero_when_no_feeds() {
let m = HealthMonitor::new(5_000);
assert_eq!(m.feeds_above_tick_count(0), 0);
}
#[test]
fn test_oldest_feed_age_ms_returns_max_age() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 5_000).unwrap();
m.heartbeat("B", 8_000).unwrap();
assert_eq!(m.oldest_feed_age_ms(10_000), Some(5_000));
}
#[test]
fn test_oldest_feed_age_ms_none_when_no_ticks() {
let mut m = monitor();
m.register("A", None);
assert!(m.oldest_feed_age_ms(10_000).is_none());
}
#[test]
fn test_total_stale_count_zero_when_all_healthy() {
let mut m = HealthMonitor::new(5_000);
m.register("A", None);
m.heartbeat("A", 9_500).unwrap();
let _ = m.check_all(10_000);
assert_eq!(m.total_stale_count(), 0);
}
#[test]
fn test_total_stale_count_correct_when_stale() {
let mut m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 1_000).unwrap();
let _ = m.check_all(10_000);
assert_eq!(m.total_stale_count(), 2);
}
#[test]
fn test_avg_feed_age_ms_none_when_no_ticks() {
let mut m = monitor();
m.register("A", None);
assert!(m.avg_feed_age_ms(10_000).is_none());
}
#[test]
fn test_avg_feed_age_ms_correct_average() {
let mut m = monitor();
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 5_000).unwrap(); m.heartbeat("B", 8_000).unwrap(); let avg = m.avg_feed_age_ms(10_000).unwrap();
assert!((avg - 3500.0).abs() < 1e-10, "got {avg}");
}
#[test]
fn test_stale_ratio_zero_with_no_feeds() {
let m = HealthMonitor::new(5_000);
assert_eq!(m.stale_ratio(), 0.0);
}
#[test]
fn test_stale_ratio_one_when_all_stale() {
let m = HealthMonitor::new(1_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 1_000).unwrap();
m.check_all(5_000);
assert!((m.stale_ratio() - 1.0).abs() < 1e-10);
}
#[test]
fn test_stale_ratio_half_when_one_of_two_stale() {
let m = HealthMonitor::new(1_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
m.heartbeat("B", 9_500).unwrap();
m.check_all(10_000);
assert!((m.stale_ratio() - 0.5).abs() < 1e-10);
}
#[test]
fn test_has_any_unknown_true_for_fresh_feed() {
let m = HealthMonitor::new(5_000);
m.register("feed", None);
assert!(m.has_any_unknown());
}
#[test]
fn test_has_any_unknown_false_after_heartbeat() {
let m = HealthMonitor::new(5_000);
m.register("feed", None);
m.heartbeat("feed", 1_000).unwrap();
assert!(!m.has_any_unknown());
}
#[test]
fn test_has_any_unknown_false_with_no_feeds() {
let m = HealthMonitor::new(5_000);
assert!(!m.has_any_unknown());
}
#[test]
fn test_is_degraded_true_when_some_unhealthy() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 9_500).unwrap(); m.heartbeat("B", 1_000).unwrap(); m.check_all(10_000);
assert!(m.is_degraded());
}
#[test]
fn test_is_degraded_false_when_all_healthy() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.heartbeat("A", 1_000).unwrap();
assert!(!m.is_degraded());
}
#[test]
fn test_is_degraded_false_with_no_feeds() {
let m = HealthMonitor::new(5_000);
assert!(!m.is_degraded());
}
#[test]
fn test_unhealthy_count_all_unknown() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
assert_eq!(m.unhealthy_count(), 2);
}
#[test]
fn test_unhealthy_count_one_healthy() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap();
assert_eq!(m.unhealthy_count(), 1);
}
#[test]
fn test_unhealthy_count_zero_when_empty() {
let m = HealthMonitor::new(5_000);
assert_eq!(m.unhealthy_count(), 0);
}
#[test]
fn test_feed_exists_true_after_register() {
let m = HealthMonitor::new(5_000);
m.register("BTC-USD", None);
assert!(m.feed_exists("BTC-USD"));
}
#[test]
fn test_feed_exists_false_for_unknown_feed() {
let m = HealthMonitor::new(5_000);
assert!(!m.feed_exists("ETH-USD"));
}
#[test]
fn test_most_stale_feed_returns_oldest_feed() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 1_000).unwrap(); m.heartbeat("B", 9_000).unwrap(); let stale = m.most_stale_feed().unwrap();
assert_eq!(stale.feed_id, "A");
}
#[test]
fn test_most_stale_feed_some_with_single_feed() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.heartbeat("A", 1_000).unwrap();
assert!(m.most_stale_feed().is_some());
}
#[test]
fn test_most_stale_feed_none_when_empty() {
let m = HealthMonitor::new(5_000);
assert!(m.most_stale_feed().is_none());
}
#[test]
fn test_stale_ratio_excl_unknown_zero_when_empty() {
let m = HealthMonitor::new(5_000);
assert_eq!(m.stale_ratio_excluding_unknown(), 0.0);
}
#[test]
fn test_stale_ratio_excl_unknown_zero_when_all_unknown() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
assert_eq!(m.stale_ratio_excluding_unknown(), 0.0);
}
#[test]
fn test_stale_ratio_excl_unknown_half_when_one_stale() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 9_500).unwrap(); m.heartbeat("B", 1_000).unwrap(); m.check_all(10_000);
assert!((m.stale_ratio_excluding_unknown() - 0.5).abs() < 1e-10);
}
#[test]
fn test_any_unknown_false_when_empty() {
let m = HealthMonitor::new(5_000);
assert!(!m.any_unknown());
}
#[test]
fn test_any_unknown_true_when_feed_registered_but_no_heartbeat() {
let m = HealthMonitor::new(5_000);
m.register("BTC-USD", None);
assert!(m.any_unknown());
}
#[test]
fn test_any_unknown_false_when_all_have_heartbeats() {
let m = HealthMonitor::new(5_000);
m.register("BTC-USD", None);
m.heartbeat("BTC-USD", 1_000).unwrap();
assert!(!m.any_unknown());
}
#[test]
fn test_degraded_count_zero_when_all_healthy() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.heartbeat("A", 9_500).unwrap();
m.check_all(10_000);
assert_eq!(m.degraded_count(), 0);
}
#[test]
fn test_degraded_count_one_when_one_stale() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 9_500).unwrap(); m.heartbeat("B", 1_000).unwrap(); m.check_all(10_000);
assert_eq!(m.degraded_count(), 1);
}
#[test]
fn test_degraded_count_zero_when_empty() {
let m = HealthMonitor::new(5_000);
assert_eq!(m.degraded_count(), 0);
}
#[test]
fn test_min_healthy_age_ms_none_when_no_healthy_feeds() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
assert!(m.min_healthy_age_ms(10_000).is_none());
}
#[test]
fn test_min_healthy_age_ms_returns_most_recent() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 8_000).unwrap(); m.heartbeat("B", 9_000).unwrap(); m.check_all(10_000);
assert_eq!(m.min_healthy_age_ms(10_000), Some(1_000));
}
#[test]
fn test_healthy_feed_ids_empty_when_no_feeds() {
let m = HealthMonitor::new(5_000);
assert!(m.healthy_feed_ids().is_empty());
}
#[test]
fn test_healthy_feed_ids_returns_healthy_only() {
let m = HealthMonitor::new(5_000);
m.register("A", None);
m.register("B", None);
m.heartbeat("A", 9_500).unwrap(); m.check_all(10_000);
let ids = m.healthy_feed_ids();
assert_eq!(ids, vec!["A".to_string()]);
}
#[test]
fn test_time_since_last_heartbeat_correct() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
m.heartbeat("BTC", 9_000).unwrap();
assert_eq!(m.time_since_last_heartbeat("BTC", 10_000), Some(1_000));
}
#[test]
fn test_time_since_last_heartbeat_none_when_no_tick() {
let m = HealthMonitor::new(5_000);
m.register("ETH", None);
assert!(m.time_since_last_heartbeat("ETH", 10_000).is_none());
}
#[test]
fn test_time_since_last_heartbeat_none_when_unknown_feed() {
let m = HealthMonitor::new(5_000);
assert!(m.time_since_last_heartbeat("MISSING", 10_000).is_none());
}
#[test]
fn test_register_batch_registers_all_feeds() {
let m = HealthMonitor::new(5_000);
m.register_batch(&[("BTC", 1_000), ("ETH", 2_000), ("SOL", 3_000)]);
assert!(m.feed_exists("BTC"));
assert!(m.feed_exists("ETH"));
assert!(m.feed_exists("SOL"));
}
#[test]
fn test_register_batch_uses_custom_thresholds() {
let m = HealthMonitor::new(10_000);
m.register_batch(&[("BTC", 500)]);
m.heartbeat("BTC", 0).unwrap();
m.check_all(600);
assert_eq!(m.stale_count(), 1);
}
#[test]
fn test_register_batch_empty_slice_is_noop() {
let m = HealthMonitor::new(5_000);
m.register_batch(&[]);
assert_eq!(m.feed_count(), 0);
}
#[test]
fn test_unknown_feed_ids_all_new_feeds_are_unknown() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
m.register("ETH", None);
let ids = m.unknown_feed_ids();
assert!(ids.contains(&"BTC".to_string()));
assert!(ids.contains(&"ETH".to_string()));
}
#[test]
fn test_unknown_feed_ids_empty_after_heartbeat_and_check() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
m.heartbeat("BTC", 0).unwrap();
m.check_all(100);
assert!(m.unknown_feed_ids().is_empty());
}
#[test]
fn test_unknown_feed_ids_empty_when_no_feeds() {
let m = HealthMonitor::new(5_000);
assert!(m.unknown_feed_ids().is_empty());
}
#[test]
fn test_feeds_needing_check_returns_stale_and_unknown() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None); m.register("ETH", None);
m.heartbeat("ETH", 0).unwrap();
m.check_all(100); let needing = m.feeds_needing_check();
assert!(needing.contains(&"BTC".to_string()));
assert!(!needing.contains(&"ETH".to_string()));
}
#[test]
fn test_feeds_needing_check_empty_when_all_healthy() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
m.heartbeat("BTC", 0).unwrap();
m.check_all(100);
assert!(m.feeds_needing_check().is_empty());
}
#[test]
fn test_feeds_needing_check_sorted() {
let m = HealthMonitor::new(5_000);
m.register("ZZZ", None);
m.register("AAA", None);
let needing = m.feeds_needing_check();
assert_eq!(needing, vec!["AAA".to_string(), "ZZZ".to_string()]);
}
#[test]
fn test_ratio_healthy_zero_when_no_feeds() {
let m = HealthMonitor::new(5_000);
assert_eq!(m.ratio_healthy(), 0.0);
}
#[test]
fn test_ratio_healthy_zero_when_all_unknown() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
m.register("ETH", None);
assert_eq!(m.ratio_healthy(), 0.0);
}
#[test]
fn test_ratio_healthy_one_when_all_healthy() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
m.heartbeat("BTC", 0).unwrap();
m.check_all(100);
assert_eq!(m.ratio_healthy(), 1.0);
}
#[test]
fn test_ratio_healthy_half_when_one_of_two() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
m.register("ETH", None);
m.heartbeat("BTC", 0).unwrap();
m.check_all(100);
let ratio = m.ratio_healthy();
assert!((ratio - 0.5).abs() < 1e-10);
}
#[test]
fn test_total_tick_count_zero_when_no_feeds() {
let m = HealthMonitor::new(5_000);
assert_eq!(m.total_tick_count(), 0);
}
#[test]
fn test_total_tick_count_sums_all_feeds() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
m.register("ETH", None);
m.heartbeat("BTC", 0).unwrap();
m.heartbeat("BTC", 1).unwrap();
m.heartbeat("ETH", 0).unwrap();
assert_eq!(m.total_tick_count(), 3);
}
#[test]
fn test_last_updated_feed_id_none_when_no_ticks() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
assert!(m.last_updated_feed_id().is_none());
}
#[test]
fn test_last_updated_feed_id_returns_most_recent() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
m.register("ETH", None);
m.heartbeat("BTC", 100).unwrap();
m.heartbeat("ETH", 200).unwrap(); assert_eq!(m.last_updated_feed_id(), Some("ETH".to_string()));
}
#[test]
fn test_is_any_stale_false_when_no_feeds() {
let m = HealthMonitor::new(5_000);
assert!(!m.is_any_stale());
}
#[test]
fn test_is_any_stale_false_when_all_healthy() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
m.heartbeat("BTC", 0).unwrap();
m.check_all(100);
assert!(!m.is_any_stale());
}
#[test]
fn test_is_any_stale_true_when_stale_feed() {
let m = HealthMonitor::new(5_000);
m.register("BTC", None);
m.heartbeat("BTC", 0).unwrap();
m.check_all(10_000); assert!(m.is_any_stale());
}
}