use scirs2_core::numeric::Float;
use std::collections::{BTreeMap, HashMap};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use crate::error::Result;
mod accumulator;
mod aggregation;
mod alerts;
mod export;
#[cfg(test)]
mod regression_tests;
pub use accumulator::{MetricsAccumulator, ResourceProbe, RobustnessProbe};
pub use aggregation::AggregatedSeries;
pub use alerts::KNOWN_METRIC_PATHS;
pub(crate) fn unix_timestamp(time: SystemTime) -> u64 {
time.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
}
pub(crate) const MICROS_PER_SEC: u64 = 1_000_000;
pub(crate) fn unix_timestamp_micros(time: SystemTime) -> u64 {
let elapsed = time.duration_since(UNIX_EPOCH).unwrap_or_default();
elapsed
.as_secs()
.saturating_mul(MICROS_PER_SEC)
.saturating_add(u64::from(elapsed.subsec_micros()))
}
pub(crate) fn saturating_elapsed(later: SystemTime, earlier: SystemTime) -> Duration {
later.duration_since(earlier).unwrap_or_default()
}
#[derive(Debug)]
pub struct StreamingMetricsCollector<A: Float + Send + Sync> {
performance_metrics: PerformanceMetrics<A>,
resource_metrics: ResourceMetrics,
quality_metrics: QualityMetrics<A>,
business_metrics: BusinessMetrics<A>,
historical_data: HistoricalMetrics<A>,
dashboards: Vec<Dashboard>,
alert_system: AlertSystem<A>,
aggregation_config: AggregationConfig,
export_config: ExportConfig,
accumulator: MetricsAccumulator<A>,
slo: Option<SloTargets>,
cost_model: Option<CostModel<A>>,
last_export: Option<SystemTime>,
}
#[derive(Debug, Clone, Default)]
pub struct SloTargets {
pub max_processing_time: Option<Duration>,
pub max_loss: Option<f64>,
pub max_memory_bytes: Option<u64>,
}
impl SloTargets {
pub fn is_empty(&self) -> bool {
self.max_processing_time.is_none()
&& self.max_loss.is_none()
&& self.max_memory_bytes.is_none()
}
}
#[derive(Debug, Clone)]
pub struct CostModel<A: Float + Send + Sync> {
pub compute_cost_per_second: A,
pub memory_cost_per_gb_hour: A,
pub energy_cost_per_joule: A,
pub value_per_loss_unit: A,
}
#[derive(Debug, Clone)]
pub struct PerformanceMetrics<A: Float + Send + Sync> {
pub throughput: ThroughputMetrics,
pub latency: LatencyMetrics,
pub accuracy: AccuracyMetrics<A>,
pub stability: StabilityMetrics<A>,
pub efficiency: EfficiencyMetrics<A>,
}
#[derive(Debug, Clone)]
pub struct ThroughputMetrics {
pub samples_per_second: f64,
pub updates_per_second: f64,
pub gradients_per_second: f64,
pub peak_throughput: f64,
pub min_throughput: f64,
pub throughput_variance: f64,
pub throughput_trend: f64,
}
#[derive(Debug, Clone)]
pub struct LatencyMetrics {
pub end_to_end: LatencyStats,
pub gradient_computation: Option<LatencyStats>,
pub update_application: Option<LatencyStats>,
pub communication: Option<LatencyStats>,
pub queue_wait_time: Option<LatencyStats>,
pub jitter: f64,
}
#[derive(Debug, Clone)]
pub struct LatencyStats {
pub mean: Duration,
pub median: Duration,
pub p95: Duration,
pub p99: Duration,
pub p999: Duration,
pub max: Duration,
pub min: Duration,
pub std_dev: Duration,
}
#[derive(Debug, Clone)]
pub struct AccuracyMetrics<A: Float + Send + Sync> {
pub current_loss: A,
pub loss_reduction_rate: A,
pub convergence_rate: A,
pub prediction_accuracy: Option<A>,
pub gradient_magnitude: A,
pub parameter_stability: A,
pub learning_progress: A,
}
#[derive(Debug, Clone)]
pub struct StabilityMetrics<A: Float + Send + Sync> {
pub loss_variance: A,
pub gradient_variance: A,
pub parameter_drift: A,
pub oscillation_score: A,
pub divergence_probability: A,
pub stability_confidence: A,
}
#[derive(Debug, Clone)]
pub struct EfficiencyMetrics<A: Float + Send + Sync> {
pub computational_efficiency: Option<A>,
pub memory_efficiency: Option<A>,
pub communication_efficiency: Option<A>,
pub energy_efficiency: Option<A>,
pub resource_utilization: A,
pub cost_efficiency: Option<A>,
}
#[derive(Debug, Clone, Default)]
pub struct ResourceMetrics {
pub cpu_utilization: Option<f64>,
pub memory_usage: MemoryUsage,
pub gpu_utilization: Option<f64>,
pub network_bandwidth: Option<f64>,
pub disk_io: Option<f64>,
pub thread_utilization: Option<f64>,
}
#[derive(Debug, Clone, Default)]
pub struct MemoryUsage {
pub total_allocated: Option<u64>,
pub current_used: u64,
pub peak_usage: u64,
pub fragmentation_ratio: Option<f64>,
pub gc_overhead: Option<f64>,
pub efficiency: Option<f64>,
}
#[derive(Debug, Clone)]
pub struct QualityMetrics<A: Float + Send + Sync> {
pub data_quality: A,
pub model_quality: ModelQuality<A>,
pub concept_drift: ConceptDriftMetrics<A>,
pub anomaly_detection: AnomalyMetrics<A>,
pub robustness: RobustnessMetrics<A>,
}
#[derive(Debug, Clone)]
pub struct ModelQuality<A: Float + Send + Sync> {
pub training_quality: A,
pub generalization_score: Option<A>,
pub overfitting_score: Option<A>,
pub underfitting_score: Option<A>,
pub complexity_score: Option<A>,
}
#[derive(Debug, Clone)]
pub struct ConceptDriftMetrics<A: Float + Send + Sync> {
pub drift_confidence: Option<A>,
pub drift_magnitude: Option<A>,
pub drift_frequency: f64,
pub adaptation_effectiveness: Option<A>,
pub detection_latency: Option<Duration>,
}
#[derive(Debug, Clone)]
pub struct AnomalyMetrics<A: Float + Send + Sync> {
pub anomaly_score: A,
pub false_positive_rate: Option<A>,
pub false_negative_rate: Option<A>,
pub detection_accuracy: Option<A>,
pub anomaly_frequency: f64,
}
#[derive(Debug, Clone, Default)]
pub struct RobustnessMetrics<A: Float + Send + Sync> {
pub noise_tolerance: Option<A>,
pub adversarial_robustness: Option<A>,
pub perturbation_sensitivity: Option<A>,
pub recovery_capability: Option<A>,
pub fault_tolerance: Option<A>,
}
#[derive(Debug, Clone)]
pub struct BusinessMetrics<A: Float + Send + Sync> {
pub availability: Option<f64>,
pub slo_compliance: Option<f64>,
pub cost_metrics: CostMetrics<A>,
pub user_satisfaction: Option<A>,
pub business_value: Option<A>,
}
#[derive(Debug, Clone, Default)]
pub struct CostMetrics<A: Float + Send + Sync> {
pub computational_cost: Option<A>,
pub infrastructure_cost: Option<A>,
pub energy_cost: Option<A>,
pub opportunity_cost: Option<A>,
pub total_cost: Option<A>,
}
#[derive(Debug)]
pub struct HistoricalMetrics<A: Float + Send + Sync> {
pub(crate) time_series: BTreeMap<u64, MetricsSnapshot<A>>,
pub(crate) retention_policy: RetentionPolicy,
pub(crate) compression_config: CompressionConfig,
}
#[derive(Debug, Clone)]
pub struct MetricsSnapshot<A: Float + Send + Sync> {
pub timestamp: u64,
pub timestamp_micros: u64,
pub performance: PerformanceMetrics<A>,
pub resource: ResourceMetrics,
pub quality: QualityMetrics<A>,
pub business: BusinessMetrics<A>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum AggregationPeriod {
Minute,
Hour,
Day,
Week,
Month,
}
impl AggregationPeriod {
pub fn seconds(self) -> u64 {
match self {
AggregationPeriod::Minute => 60,
AggregationPeriod::Hour => 3_600,
AggregationPeriod::Day => 86_400,
AggregationPeriod::Week => 604_800,
AggregationPeriod::Month => 2_592_000, }
}
}
#[derive(Debug, Clone)]
pub struct AggregatedMetrics {
pub period: AggregationPeriod,
pub period_start: u64,
pub period_end: u64,
pub sample_count: usize,
pub series: BTreeMap<String, AggregatedSeries>,
}
#[derive(Debug, Clone)]
pub struct RetentionPolicy {
pub raw_data_retention: u64,
pub aggregated_retention: HashMap<AggregationPeriod, u64>,
pub auto_cleanup: bool,
pub max_storage_size: u64,
}
#[derive(Debug, Clone)]
pub struct CompressionConfig {
pub enabled: bool,
pub algorithm: CompressionAlgorithm,
pub target_ratio: f64,
pub lossy_tolerance: f64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CompressionAlgorithm {
None,
Gzip,
Lz4,
Zstd,
Custom,
}
#[derive(Debug)]
pub struct Dashboard {
pub name: String,
pub widgets: Vec<Widget>,
pub update_frequency: Duration,
pub auto_refresh: bool,
}
#[derive(Debug)]
pub struct Widget {
pub widget_type: WidgetType,
pub metrics: Vec<String>,
pub config: WidgetConfig,
}
#[derive(Debug, Clone)]
pub enum WidgetType {
LineChart,
BarChart,
Gauge,
Table,
Heatmap,
Histogram,
ScatterPlot,
TextDisplay,
}
#[derive(Debug, Clone)]
pub struct WidgetConfig {
pub title: String,
pub time_range: Duration,
pub refresh_rate: Duration,
pub color_scheme: String,
pub layout: WidgetLayout,
}
#[derive(Debug, Clone)]
pub struct WidgetLayout {
pub x: u32,
pub y: u32,
pub width: u32,
pub height: u32,
}
#[derive(Debug)]
pub struct AlertSystem<A: Float + Send + Sync> {
pub rules: Vec<AlertRule<A>>,
pub active_alerts: Vec<Alert<A>>,
pub alert_history: Vec<Alert<A>>,
pub notification_channels: Vec<NotificationChannel>,
pub(crate) rule_state: HashMap<String, alerts::RuleState>,
pub(crate) next_alert_id: u64,
pub(crate) max_history: usize,
}
#[derive(Debug, Clone)]
pub struct AlertRule<A: Float + Send + Sync> {
pub name: String,
pub metric_path: String,
pub condition: AlertCondition<A>,
pub severity: AlertSeverity,
pub evaluation_frequency: Duration,
pub notifications: Vec<String>,
}
#[derive(Debug, Clone)]
pub enum AlertCondition<A: Float + Send + Sync> {
Threshold {
operator: ComparisonOperator,
value: A,
},
RateOfChange { threshold: A, time_window: Duration },
Anomaly { sensitivity: A },
Custom { expression: String },
}
#[derive(Debug, Clone, Copy)]
pub enum ComparisonOperator {
GreaterThan,
LessThan,
GreaterThanOrEqual,
LessThanOrEqual,
Equal,
NotEqual,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AlertSeverity {
Critical,
Warning,
Info,
}
#[derive(Debug, Clone)]
pub struct Alert<A: Float + Send + Sync> {
pub id: String,
pub rule_name: String,
pub triggered_at: SystemTime,
pub resolved_at: Option<SystemTime>,
pub current_value: A,
pub threshold: A,
pub severity: AlertSeverity,
pub message: String,
}
#[derive(Debug, Clone)]
pub enum NotificationChannel {
Email {
addresses: Vec<String>,
},
Webhook {
url: String,
headers: HashMap<String, String>,
},
Slack {
webhook_url: String,
channel: String,
},
PagerDuty {
integration_key: String,
},
Custom {
config: HashMap<String, String>,
},
}
#[derive(Debug, Clone)]
pub struct AggregationConfig {
pub default_functions: Vec<AggregationFunction>,
pub custom_aggregations: HashMap<String, Vec<AggregationFunction>>,
pub intervals: Vec<Duration>,
pub max_window: Duration,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AggregationFunction {
Mean,
Median,
Min,
Max,
Sum,
Count,
StdDev,
Percentile(u8), }
#[derive(Debug, Clone)]
pub struct ExportConfig {
pub formats: Vec<ExportFormat>,
pub destinations: Vec<ExportDestination>,
pub frequency: Duration,
pub batch_size: usize,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ExportFormat {
Json,
Csv,
Parquet,
Prometheus,
InfluxDB,
Custom { format: String },
}
#[derive(Debug, Clone)]
pub enum ExportDestination {
File {
path: String,
},
Database {
connection_string: String,
},
S3 {
bucket: String,
prefix: String,
},
Http {
endpoint: String,
headers: HashMap<String, String>,
},
Kafka {
topic: String,
brokers: Vec<String>,
},
}
impl<A: Float + Default + Clone + std::fmt::Debug + Send + Sync> StreamingMetricsCollector<A> {
pub fn new() -> Self {
Self::with_window(512)
}
pub fn with_window(window: usize) -> Self {
Self {
performance_metrics: PerformanceMetrics::default(),
resource_metrics: ResourceMetrics::default(),
quality_metrics: QualityMetrics::default(),
business_metrics: BusinessMetrics::default(),
historical_data: HistoricalMetrics::new(),
dashboards: Vec::new(),
alert_system: AlertSystem::new(),
aggregation_config: AggregationConfig::default(),
export_config: ExportConfig::default(),
accumulator: MetricsAccumulator::new(window),
slo: None,
cost_model: None,
last_export: None,
}
}
pub fn set_slo_targets(&mut self, targets: SloTargets) {
self.slo = if targets.is_empty() {
None
} else {
Some(targets)
};
}
pub fn set_cost_model(&mut self, model: CostModel<A>) {
self.cost_model = Some(model);
}
pub fn record_resource_probe(&mut self, probe: ResourceProbe) {
self.accumulator.record_resource_probe(probe);
}
pub fn record_robustness_probe(&mut self, probe: RobustnessProbe<A>) {
self.accumulator.record_robustness_probe(probe);
}
pub fn record_outage(&mut self, downtime: Duration) {
self.accumulator.record_outage(downtime);
}
pub fn record_drift_event(
&mut self,
magnitude: A,
confidence: A,
detection_latency: Duration,
adaptation_effectiveness: Option<A>,
) {
self.accumulator.record_drift_event(
magnitude,
confidence,
detection_latency,
adaptation_effectiveness,
);
}
pub fn record_energy(&mut self, joules: f64) {
self.accumulator.record_energy(joules);
}
pub fn record_sample(&mut self, sample: MetricsSample<A>) -> Result<()> {
self.accumulator.ingest(&sample);
self.update_performance_metrics(&sample)?;
self.update_resource_metrics(&sample)?;
self.update_quality_metrics(&sample)?;
self.update_business_metrics(&sample)?;
let timestamp = unix_timestamp(sample.timestamp);
let timestamp_micros = unix_timestamp_micros(sample.timestamp);
let snapshot = MetricsSnapshot {
timestamp,
timestamp_micros,
performance: self.performance_metrics.clone(),
resource: self.resource_metrics.clone(),
quality: self.quality_metrics.clone(),
business: self.business_metrics.clone(),
};
self.historical_data.store_snapshot(snapshot)?;
let summary = self.current_summary();
self.alert_system.evaluate_rules(&sample, &summary)?;
Ok(())
}
pub fn get_current_metrics(&self) -> MetricsSummary<A> {
self.current_summary()
}
pub(crate) fn current_summary(&self) -> MetricsSummary<A> {
MetricsSummary {
performance: self.performance_metrics.clone(),
resource: self.resource_metrics.clone(),
quality: self.quality_metrics.clone(),
business: self.business_metrics.clone(),
timestamp: SystemTime::now(),
}
}
pub fn get_historical_metrics(
&self,
start_time: SystemTime,
end_time: SystemTime,
) -> Result<Vec<MetricsSnapshot<A>>> {
self.historical_data.get_range(start_time, end_time)
}
pub fn retained_snapshot_count(&self) -> usize {
self.historical_data.time_series.len()
}
pub fn add_alert_rule(&mut self, rule: AlertRule<A>) -> Result<()> {
self.alert_system.add_rule(rule)
}
pub fn active_alerts(&self) -> &[Alert<A>] {
&self.alert_system.active_alerts
}
pub fn alert_history(&self) -> &[Alert<A>] {
&self.alert_system.alert_history
}
pub fn add_dashboard(&mut self, dashboard: Dashboard) {
self.dashboards.push(dashboard);
}
pub fn render_dashboard(&self, name: &str) -> Option<Vec<(String, Option<f64>)>> {
let dashboard = self.dashboards.iter().find(|d| d.name == name)?;
let summary = self.current_summary();
let mut resolved = Vec::new();
for widget in &dashboard.widgets {
for metric in &widget.metrics {
resolved.push((metric.clone(), alerts::resolve_metric(&summary, metric)));
}
}
Some(resolved)
}
pub fn aggregation_config(&self) -> &AggregationConfig {
&self.aggregation_config
}
pub fn set_aggregation_config(&mut self, config: AggregationConfig) {
self.aggregation_config = config;
}
pub fn set_export_config(&mut self, config: ExportConfig) {
self.export_config = config;
}
pub fn set_retention_policy(&mut self, policy: RetentionPolicy) {
self.historical_data.retention_policy = policy;
self.historical_data.prune();
}
pub fn set_compression_config(&mut self, config: CompressionConfig) {
self.historical_data.compression_config = config;
}
}
impl<A: Float + Default + Clone + std::fmt::Debug + Send + Sync> Default
for StreamingMetricsCollector<A>
{
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Clone)]
pub struct MetricsSample<A: Float + Send + Sync> {
pub timestamp: SystemTime,
pub loss: A,
pub gradient_magnitude: A,
pub processing_time: Duration,
pub memory_usage: u64,
pub gradient_computation_time: Option<Duration>,
pub update_application_time: Option<Duration>,
pub communication_time: Option<Duration>,
pub queue_wait_time: Option<Duration>,
pub custom_metrics: HashMap<String, A>,
}
impl<A: Float + Send + Sync> MetricsSample<A> {
pub fn new(
timestamp: SystemTime,
loss: A,
gradient_magnitude: A,
processing_time: Duration,
memory_usage: u64,
) -> Self {
Self {
timestamp,
loss,
gradient_magnitude,
processing_time,
memory_usage,
gradient_computation_time: None,
update_application_time: None,
communication_time: None,
queue_wait_time: None,
custom_metrics: HashMap::new(),
}
}
}
#[derive(Debug, Clone)]
pub struct MetricsSummary<A: Float + Send + Sync> {
pub performance: PerformanceMetrics<A>,
pub resource: ResourceMetrics,
pub quality: QualityMetrics<A>,
pub business: BusinessMetrics<A>,
pub timestamp: SystemTime,
}
impl<A: Float + Default + Send + Sync> Default for PerformanceMetrics<A> {
fn default() -> Self {
Self {
throughput: ThroughputMetrics::default(),
latency: LatencyMetrics::default(),
accuracy: AccuracyMetrics::default(),
stability: StabilityMetrics::default(),
efficiency: EfficiencyMetrics::default(),
}
}
}
impl Default for ThroughputMetrics {
fn default() -> Self {
Self {
samples_per_second: 0.0,
updates_per_second: 0.0,
gradients_per_second: 0.0,
peak_throughput: 0.0,
min_throughput: f64::MAX,
throughput_variance: 0.0,
throughput_trend: 0.0,
}
}
}
impl Default for LatencyMetrics {
fn default() -> Self {
Self {
end_to_end: LatencyStats::default(),
gradient_computation: None,
update_application: None,
communication: None,
queue_wait_time: None,
jitter: 0.0,
}
}
}
impl Default for LatencyStats {
fn default() -> Self {
Self {
mean: Duration::from_micros(0),
median: Duration::from_micros(0),
p95: Duration::from_micros(0),
p99: Duration::from_micros(0),
p999: Duration::from_micros(0),
max: Duration::from_micros(0),
min: Duration::from_micros(u64::MAX),
std_dev: Duration::from_micros(0),
}
}
}
impl<A: Float + Default + Send + Sync> Default for AccuracyMetrics<A> {
fn default() -> Self {
Self {
current_loss: A::default(),
loss_reduction_rate: A::default(),
convergence_rate: A::default(),
prediction_accuracy: None,
gradient_magnitude: A::default(),
parameter_stability: A::default(),
learning_progress: A::default(),
}
}
}
impl<A: Float + Default + Send + Sync> Default for StabilityMetrics<A> {
fn default() -> Self {
Self {
loss_variance: A::default(),
gradient_variance: A::default(),
parameter_drift: A::default(),
oscillation_score: A::default(),
divergence_probability: A::default(),
stability_confidence: A::default(),
}
}
}
impl<A: Float + Default + Send + Sync> Default for EfficiencyMetrics<A> {
fn default() -> Self {
Self {
computational_efficiency: None,
memory_efficiency: None,
communication_efficiency: None,
energy_efficiency: None,
resource_utilization: A::default(),
cost_efficiency: None,
}
}
}
impl<A: Float + Default + Send + Sync> Default for QualityMetrics<A> {
fn default() -> Self {
Self {
data_quality: A::default(),
model_quality: ModelQuality::default(),
concept_drift: ConceptDriftMetrics::default(),
anomaly_detection: AnomalyMetrics::default(),
robustness: RobustnessMetrics::default(),
}
}
}
impl<A: Float + Default + Send + Sync> Default for ModelQuality<A> {
fn default() -> Self {
Self {
training_quality: A::default(),
generalization_score: None,
overfitting_score: None,
underfitting_score: None,
complexity_score: None,
}
}
}
impl<A: Float + Default + Send + Sync> Default for ConceptDriftMetrics<A> {
fn default() -> Self {
Self {
drift_confidence: None,
drift_magnitude: None,
drift_frequency: 0.0,
adaptation_effectiveness: None,
detection_latency: None,
}
}
}
impl<A: Float + Default + Send + Sync> Default for AnomalyMetrics<A> {
fn default() -> Self {
Self {
anomaly_score: A::default(),
false_positive_rate: None,
false_negative_rate: None,
detection_accuracy: None,
anomaly_frequency: 0.0,
}
}
}
impl<A: Float + Default + Send + Sync> Default for BusinessMetrics<A> {
fn default() -> Self {
Self {
availability: None,
slo_compliance: None,
cost_metrics: CostMetrics::default(),
user_satisfaction: None,
business_value: None,
}
}
}
impl<A: Float + Send + Sync> HistoricalMetrics<A> {
fn new() -> Self {
Self {
time_series: BTreeMap::new(),
retention_policy: RetentionPolicy::default(),
compression_config: CompressionConfig::default(),
}
}
}
impl Default for RetentionPolicy {
fn default() -> Self {
let mut aggregated_retention = HashMap::new();
aggregated_retention.insert(AggregationPeriod::Minute, 3600 * 24); aggregated_retention.insert(AggregationPeriod::Hour, 3600 * 24 * 7); aggregated_retention.insert(AggregationPeriod::Day, 3600 * 24 * 30); aggregated_retention.insert(AggregationPeriod::Week, 3600 * 24 * 365); aggregated_retention.insert(AggregationPeriod::Month, 3600 * 24 * 365 * 5);
Self {
raw_data_retention: 3600 * 24, aggregated_retention,
auto_cleanup: true,
max_storage_size: 1024 * 1024 * 1024 * 10, }
}
}
impl Default for CompressionConfig {
fn default() -> Self {
Self {
enabled: false,
algorithm: CompressionAlgorithm::None,
target_ratio: 0.3,
lossy_tolerance: 0.01,
}
}
}
impl Default for AggregationConfig {
fn default() -> Self {
Self {
default_functions: vec![
AggregationFunction::Mean,
AggregationFunction::Min,
AggregationFunction::Max,
AggregationFunction::StdDev,
],
custom_aggregations: HashMap::new(),
intervals: vec![
Duration::from_secs(60), Duration::from_secs(3600), Duration::from_secs(86400), ],
max_window: Duration::from_secs(86400 * 30), }
}
}
impl Default for ExportConfig {
fn default() -> Self {
Self {
formats: vec![ExportFormat::Json],
destinations: vec![ExportDestination::File {
path: std::env::temp_dir()
.join("streaming_metrics")
.to_string_lossy()
.into_owned(),
}],
frequency: Duration::from_secs(300), batch_size: 1000,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_metrics_collector_creation() {
let collector = StreamingMetricsCollector::<f64>::new();
assert_eq!(
collector.performance_metrics.throughput.samples_per_second,
0.0
);
assert!(collector.dashboards.is_empty());
}
#[test]
fn test_metrics_sample() {
let sample = MetricsSample::new(
SystemTime::now(),
0.5f64,
0.1f64,
Duration::from_millis(10),
1024,
);
assert_eq!(sample.loss, 0.5f64);
assert_eq!(sample.gradient_magnitude, 0.1f64);
}
#[test]
fn test_latency_stats_default() {
let stats = LatencyStats::default();
assert_eq!(stats.mean, Duration::from_micros(0));
assert_eq!(stats.min, Duration::from_micros(u64::MAX));
}
#[test]
fn test_aggregation_period() {
let periods = [
AggregationPeriod::Minute,
AggregationPeriod::Hour,
AggregationPeriod::Day,
AggregationPeriod::Week,
AggregationPeriod::Month,
];
assert_eq!(periods.len(), 5);
}
#[test]
fn test_alert_severity() {
let severities = [
AlertSeverity::Critical,
AlertSeverity::Warning,
AlertSeverity::Info,
];
assert_eq!(severities.len(), 3);
}
}