#![allow(dead_code)]
use std::collections::HashMap;
#[derive(Debug, Clone)]
pub struct MetricPoint {
pub name: String,
pub value: f64,
pub tags: HashMap<String, String>,
pub timestamp_ms: u64,
}
impl MetricPoint {
pub fn new(name: impl Into<String>, value: f64, timestamp_ms: u64) -> Self {
Self {
name: name.into(),
value,
tags: HashMap::new(),
timestamp_ms,
}
}
pub fn with_tag(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
self.tags.insert(key.into(), value.into());
self
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AggregationFn {
Sum,
Mean,
Min,
Max,
P50,
P95,
P99,
}
#[derive(Debug, Default)]
pub struct MetricAggregator {
points: HashMap<String, Vec<MetricPoint>>,
}
impl MetricAggregator {
#[must_use]
pub fn new() -> Self {
Self {
points: HashMap::new(),
}
}
pub fn record(&mut self, point: MetricPoint) {
self.points
.entry(point.name.clone())
.or_default()
.push(point);
}
#[must_use]
pub fn query(&self, name: &str, window_ms: u64) -> Vec<f64> {
let pts = match self.points.get(name) {
Some(v) => v,
None => return Vec::new(),
};
let now_max = pts.iter().map(|p| p.timestamp_ms).max().unwrap_or(0);
let start = now_max.saturating_sub(window_ms);
pts.iter()
.filter(|p| p.timestamp_ms >= start)
.map(|p| p.value)
.collect()
}
#[must_use]
pub fn aggregate(&self, name: &str, window_ms: u64, agg: AggregationFn) -> Option<f64> {
let mut values = self.query(name, window_ms);
if values.is_empty() {
return None;
}
values.sort_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal));
let result = match agg {
AggregationFn::Sum => values.iter().sum(),
AggregationFn::Mean => values.iter().sum::<f64>() / values.len() as f64,
AggregationFn::Min => *values
.first()
.expect("values non-empty: is_empty check returned above"),
AggregationFn::Max => *values
.last()
.expect("values non-empty: is_empty check returned above"),
AggregationFn::P50 => percentile(&values, 50.0),
AggregationFn::P95 => percentile(&values, 95.0),
AggregationFn::P99 => percentile(&values, 99.0),
};
Some(result)
}
}
fn percentile(sorted: &[f64], p: f64) -> f64 {
if sorted.is_empty() {
return 0.0;
}
let idx = ((p / 100.0) * (sorted.len() - 1) as f64).round() as usize;
sorted[idx.min(sorted.len() - 1)]
}
#[derive(Debug, Clone)]
pub struct ClusterMetrics {
pub worker_count: u32,
pub total_tasks: u64,
pub running_tasks: u32,
pub failed_tasks: u64,
pub avg_throughput: f64,
}
impl ClusterMetrics {
#[must_use]
pub fn compute(aggregator: &MetricAggregator, now_ms: u64) -> Self {
let window_ms = now_ms;
let worker_count = aggregator
.aggregate("worker_count", window_ms, AggregationFn::Max)
.unwrap_or(0.0) as u32;
let total_tasks = aggregator
.aggregate("total_tasks", window_ms, AggregationFn::Sum)
.unwrap_or(0.0) as u64;
let running_tasks = aggregator
.aggregate("running_tasks", window_ms, AggregationFn::Max)
.unwrap_or(0.0) as u32;
let failed_tasks = aggregator
.aggregate("failed_tasks", window_ms, AggregationFn::Sum)
.unwrap_or(0.0) as u64;
let avg_throughput = aggregator
.aggregate("throughput", window_ms, AggregationFn::Mean)
.unwrap_or(0.0);
Self {
worker_count,
total_tasks,
running_tasks,
failed_tasks,
avg_throughput,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Comparison {
Above,
Below,
Equal,
}
#[derive(Debug, Clone)]
pub struct AlertRule {
pub metric_name: String,
pub threshold: f64,
pub comparison: Comparison,
pub window_ms: u64,
}
impl AlertRule {
pub fn new(
metric_name: impl Into<String>,
threshold: f64,
comparison: Comparison,
window_ms: u64,
) -> Self {
Self {
metric_name: metric_name.into(),
threshold,
comparison,
window_ms,
}
}
}
#[derive(Debug, Clone)]
pub struct Alert {
pub rule: AlertRule,
pub current_value: f64,
pub triggered_at_ms: u64,
}
pub struct AlertEvaluator;
impl AlertEvaluator {
#[must_use]
pub fn check(rules: &[AlertRule], aggregator: &MetricAggregator, now_ms: u64) -> Vec<Alert> {
let mut alerts = Vec::new();
for rule in rules {
if let Some(value) =
aggregator.aggregate(&rule.metric_name, rule.window_ms, AggregationFn::Mean)
{
let fires = match rule.comparison {
Comparison::Above => value > rule.threshold,
Comparison::Below => value < rule.threshold,
Comparison::Equal => (value - rule.threshold).abs() < f64::EPSILON,
};
if fires {
alerts.push(Alert {
rule: rule.clone(),
current_value: value,
triggered_at_ms: now_ms,
});
}
}
}
alerts
}
}
#[cfg(test)]
mod tests {
use super::*;
fn ts(ms: u64) -> u64 {
ms
}
fn record_values(agg: &mut MetricAggregator, name: &str, values: &[(f64, u64)]) {
for &(v, t) in values {
agg.record(MetricPoint::new(name, v, ts(t)));
}
}
#[test]
fn test_record_and_query() {
let mut agg = MetricAggregator::new();
record_values(&mut agg, "cpu", &[(0.5, 1000), (0.6, 2000), (0.7, 3000)]);
let vals = agg.query("cpu", 5000);
assert_eq!(vals.len(), 3);
}
#[test]
fn test_query_windowed() {
let mut agg = MetricAggregator::new();
record_values(&mut agg, "fps", &[(10.0, 100), (20.0, 500), (30.0, 1000)]);
let vals = agg.query("fps", 600);
assert_eq!(vals.len(), 2);
}
#[test]
fn test_aggregate_sum() {
let mut agg = MetricAggregator::new();
record_values(&mut agg, "bytes", &[(100.0, 1), (200.0, 2), (300.0, 3)]);
let sum = agg.aggregate("bytes", 100, AggregationFn::Sum);
assert!((sum.expect("metric computation should succeed") - 600.0).abs() < 1e-9);
}
#[test]
fn test_aggregate_mean() {
let mut agg = MetricAggregator::new();
record_values(&mut agg, "lat", &[(10.0, 1), (20.0, 2), (30.0, 3)]);
let mean = agg.aggregate("lat", 100, AggregationFn::Mean);
assert!((mean.expect("metric computation should succeed") - 20.0).abs() < 1e-9);
}
#[test]
fn test_aggregate_min_max() {
let mut agg = MetricAggregator::new();
record_values(&mut agg, "val", &[(5.0, 1), (1.0, 2), (9.0, 3)]);
assert!(
(agg.aggregate("val", 100, AggregationFn::Min)
.expect("aggregation should succeed")
- 1.0)
.abs()
< 1e-9
);
assert!(
(agg.aggregate("val", 100, AggregationFn::Max)
.expect("aggregation should succeed")
- 9.0)
.abs()
< 1e-9
);
}
#[test]
fn test_aggregate_p50() {
let mut agg = MetricAggregator::new();
record_values(
&mut agg,
"m",
&[(3.0, 1), (1.0, 2), (5.0, 3), (2.0, 4), (4.0, 5)],
);
let p50 = agg.aggregate("m", 100, AggregationFn::P50);
assert!((p50.expect("metric computation should succeed") - 3.0).abs() < 1e-9);
}
#[test]
fn test_aggregate_empty() {
let agg = MetricAggregator::new();
assert!(agg
.aggregate("nonexistent", 1000, AggregationFn::Mean)
.is_none());
}
#[test]
fn test_cluster_metrics_compute() {
let mut agg = MetricAggregator::new();
agg.record(MetricPoint::new("worker_count", 4.0, 1000));
agg.record(MetricPoint::new("total_tasks", 100.0, 1000));
agg.record(MetricPoint::new("running_tasks", 12.0, 1000));
agg.record(MetricPoint::new("failed_tasks", 3.0, 1000));
agg.record(MetricPoint::new("throughput", 25.0, 1000));
let metrics = ClusterMetrics::compute(&agg, 2000);
assert_eq!(metrics.worker_count, 4);
assert_eq!(metrics.running_tasks, 12);
assert!((metrics.avg_throughput - 25.0).abs() < 1e-9);
}
#[test]
fn test_alert_above_fires() {
let mut agg = MetricAggregator::new();
record_values(&mut agg, "cpu", &[(0.95, 1000)]);
let rules = vec![AlertRule::new("cpu", 0.8, Comparison::Above, 5000)];
let alerts = AlertEvaluator::check(&rules, &agg, 1000);
assert_eq!(alerts.len(), 1);
assert!((alerts[0].current_value - 0.95).abs() < 1e-9);
}
#[test]
fn test_alert_below_fires() {
let mut agg = MetricAggregator::new();
record_values(&mut agg, "workers", &[(2.0, 1000)]);
let rules = vec![AlertRule::new("workers", 5.0, Comparison::Below, 5000)];
let alerts = AlertEvaluator::check(&rules, &agg, 1000);
assert_eq!(alerts.len(), 1);
}
#[test]
fn test_alert_does_not_fire_when_ok() {
let mut agg = MetricAggregator::new();
record_values(&mut agg, "cpu", &[(0.5, 1000)]);
let rules = vec![AlertRule::new("cpu", 0.8, Comparison::Above, 5000)];
let alerts = AlertEvaluator::check(&rules, &agg, 1000);
assert!(alerts.is_empty());
}
}