Skip to main content

optirs_core/streaming/streaming_metrics/
accumulator.rs

1// Real metric accumulation for the streaming metrics collector (finding M1).
2//
3// The four `update_*_metrics` methods used to be `Ok(())` with a comment. They
4// now derive every metric from the rolling raw observations kept here, and
5// leave an `Option` at `None` whenever the underlying measurement was never
6// supplied.
7
8use super::*;
9use std::collections::VecDeque;
10
11/// OS-level resource measurements the collector cannot take itself.
12#[derive(Debug, Clone, Default)]
13pub struct ResourceProbe {
14    /// CPU utilization percentage (0-100)
15    pub cpu_utilization: Option<f64>,
16    /// GPU utilization percentage (0-100)
17    pub gpu_utilization: Option<f64>,
18    /// Network bandwidth in MB/s
19    pub network_bandwidth_mbps: Option<f64>,
20    /// Disk I/O in MB/s
21    pub disk_io_mbps: Option<f64>,
22    /// Fraction of available threads busy (0-1)
23    pub thread_utilization: Option<f64>,
24    /// Bytes handed out by the allocator
25    pub total_allocated_bytes: Option<u64>,
26    /// Allocator fragmentation ratio (0-1)
27    pub fragmentation_ratio: Option<f64>,
28}
29
30/// Robustness measurements obtained from deliberate perturbation experiments.
31#[derive(Debug, Clone)]
32pub struct RobustnessProbe<A: Float + Send + Sync> {
33    /// Retained accuracy under input noise
34    pub noise_tolerance: Option<A>,
35    /// Retained accuracy under adversarial perturbation
36    pub adversarial_robustness: Option<A>,
37    /// Output sensitivity to a unit input perturbation
38    pub perturbation_sensitivity: Option<A>,
39    /// Fraction of performance recovered after a shock
40    pub recovery_capability: Option<A>,
41    /// Fraction of injected faults survived
42    pub fault_tolerance: Option<A>,
43}
44
45impl<A: Float + Send + Sync> Default for RobustnessProbe<A> {
46    fn default() -> Self {
47        Self {
48            noise_tolerance: None,
49            adversarial_robustness: None,
50            perturbation_sensitivity: None,
51            recovery_capability: None,
52            fault_tolerance: None,
53        }
54    }
55}
56
57/// Most recent externally reported drift event.
58#[derive(Debug, Clone)]
59struct DriftReport<A: Float + Send + Sync> {
60    magnitude: A,
61    confidence: A,
62    detection_latency: Duration,
63    adaptation_effectiveness: Option<A>,
64}
65
66/// Rolling raw observations every aggregate metric is derived from.
67#[derive(Debug, Clone)]
68pub struct MetricsAccumulator<A: Float + Send + Sync> {
69    window: usize,
70
71    losses: VecDeque<A>,
72    gradients: VecDeque<A>,
73    latencies: VecDeque<Duration>,
74    inter_arrival: VecDeque<Duration>,
75    memory: VecDeque<u64>,
76    gradient_times: VecDeque<Duration>,
77    update_times: VecDeque<Duration>,
78    communication_times: VecDeque<Duration>,
79    queue_times: VecDeque<Duration>,
80
81    first_timestamp: Option<SystemTime>,
82    last_timestamp: Option<SystemTime>,
83    first_loss: Option<A>,
84
85    sample_count: u64,
86    valid_sample_count: u64,
87    peak_memory: u64,
88    peak_rate: Option<f64>,
89    min_rate: Option<f64>,
90
91    // Welford accumulator over the loss stream, used for the anomaly z-score.
92    loss_n: u64,
93    loss_mean: f64,
94    loss_m2: f64,
95    anomaly_count: u64,
96
97    slo_evaluated: u64,
98    slo_met: u64,
99
100    downtime: Duration,
101    drift_events: u64,
102    last_drift: Option<DriftReport<A>>,
103    energy_joules: Option<f64>,
104
105    resource_probe: Option<ResourceProbe>,
106    robustness_probe: Option<RobustnessProbe<A>>,
107}
108
109fn push_capped<T>(queue: &mut VecDeque<T>, value: T, window: usize) {
110    queue.push_back(value);
111    while queue.len() > window {
112        queue.pop_front();
113    }
114}
115
116pub(crate) fn to_scalar<A: Float>(value: f64) -> A {
117    A::from(value).unwrap_or_else(A::zero)
118}
119
120pub(crate) fn from_scalar<A: Float>(value: A) -> f64 {
121    value.to_f64().unwrap_or(0.0)
122}
123
124pub(crate) fn mean_f64(values: &[f64]) -> f64 {
125    if values.is_empty() {
126        return 0.0;
127    }
128    values.iter().sum::<f64>() / values.len() as f64
129}
130
131pub(crate) fn variance_f64(values: &[f64]) -> f64 {
132    if values.len() < 2 {
133        return 0.0;
134    }
135    let mean = mean_f64(values);
136    values.iter().map(|v| (v - mean).powi(2)).sum::<f64>() / values.len() as f64
137}
138
139/// Ordinary-least-squares slope of `values` against their index.
140pub(crate) fn ols_slope(values: &[f64]) -> f64 {
141    let n = values.len();
142    if n < 2 {
143        return 0.0;
144    }
145    let x_mean = (n - 1) as f64 / 2.0;
146    let y_mean = mean_f64(values);
147    let mut numerator = 0.0;
148    let mut denominator = 0.0;
149    for (index, value) in values.iter().enumerate() {
150        let dx = index as f64 - x_mean;
151        numerator += dx * (value - y_mean);
152        denominator += dx * dx;
153    }
154    if denominator == 0.0 {
155        0.0
156    } else {
157        numerator / denominator
158    }
159}
160
161/// Full latency statistics over a sample window, or `None` when empty.
162pub(crate) fn duration_stats(samples: &VecDeque<Duration>) -> Option<LatencyStats> {
163    if samples.is_empty() {
164        return None;
165    }
166    let mut sorted: Vec<Duration> = samples.iter().copied().collect();
167    sorted.sort();
168    let last = sorted.len() - 1;
169    let quantile = |q: f64| sorted[((sorted.len() as f64 * q) as usize).min(last)];
170
171    let seconds: Vec<f64> = sorted.iter().map(|d| d.as_secs_f64()).collect();
172    let mean_seconds = mean_f64(&seconds);
173    let std_seconds = variance_f64(&seconds).sqrt();
174
175    Some(LatencyStats {
176        mean: Duration::from_secs_f64(mean_seconds.max(0.0)),
177        median: quantile(0.50),
178        p95: quantile(0.95),
179        p99: quantile(0.99),
180        p999: quantile(0.999),
181        max: sorted[last],
182        min: sorted[0],
183        std_dev: Duration::from_secs_f64(std_seconds.max(0.0)),
184    })
185}
186
187impl<A: Float + Send + Sync> MetricsAccumulator<A> {
188    /// Create an accumulator retaining `window` observations per series.
189    pub fn new(window: usize) -> Self {
190        let window = window.max(2);
191        Self {
192            window,
193            losses: VecDeque::with_capacity(window),
194            gradients: VecDeque::with_capacity(window),
195            latencies: VecDeque::with_capacity(window),
196            inter_arrival: VecDeque::with_capacity(window),
197            memory: VecDeque::with_capacity(window),
198            gradient_times: VecDeque::with_capacity(window),
199            update_times: VecDeque::with_capacity(window),
200            communication_times: VecDeque::with_capacity(window),
201            queue_times: VecDeque::with_capacity(window),
202            first_timestamp: None,
203            last_timestamp: None,
204            first_loss: None,
205            sample_count: 0,
206            valid_sample_count: 0,
207            peak_memory: 0,
208            peak_rate: None,
209            min_rate: None,
210            loss_n: 0,
211            loss_mean: 0.0,
212            loss_m2: 0.0,
213            anomaly_count: 0,
214            slo_evaluated: 0,
215            slo_met: 0,
216            downtime: Duration::ZERO,
217            drift_events: 0,
218            last_drift: None,
219            energy_joules: None,
220            resource_probe: None,
221            robustness_probe: None,
222        }
223    }
224
225    /// Fold one sample into the rolling state.
226    pub fn ingest(&mut self, sample: &MetricsSample<A>) {
227        self.sample_count += 1;
228
229        let loss = from_scalar(sample.loss);
230        let gradient = from_scalar(sample.gradient_magnitude);
231        if loss.is_finite() && gradient.is_finite() && gradient >= 0.0 {
232            self.valid_sample_count += 1;
233        }
234
235        if self.first_timestamp.is_none() {
236            self.first_timestamp = Some(sample.timestamp);
237        }
238        if self.first_loss.is_none() && loss.is_finite() {
239            self.first_loss = Some(sample.loss);
240        }
241
242        if let Some(previous) = self.last_timestamp {
243            let gap = saturating_elapsed(sample.timestamp, previous);
244            push_capped(&mut self.inter_arrival, gap, self.window);
245            let seconds = gap.as_secs_f64();
246            if seconds > 0.0 {
247                let rate = 1.0 / seconds;
248                self.peak_rate = Some(self.peak_rate.map_or(rate, |p| p.max(rate)));
249                self.min_rate = Some(self.min_rate.map_or(rate, |p| p.min(rate)));
250            }
251        }
252        self.last_timestamp = Some(sample.timestamp);
253
254        push_capped(&mut self.losses, sample.loss, self.window);
255        push_capped(&mut self.gradients, sample.gradient_magnitude, self.window);
256        push_capped(&mut self.latencies, sample.processing_time, self.window);
257        push_capped(&mut self.memory, sample.memory_usage, self.window);
258        self.peak_memory = self.peak_memory.max(sample.memory_usage);
259
260        if let Some(value) = sample.gradient_computation_time {
261            push_capped(&mut self.gradient_times, value, self.window);
262        }
263        if let Some(value) = sample.update_application_time {
264            push_capped(&mut self.update_times, value, self.window);
265        }
266        if let Some(value) = sample.communication_time {
267            push_capped(&mut self.communication_times, value, self.window);
268        }
269        if let Some(value) = sample.queue_wait_time {
270            push_capped(&mut self.queue_times, value, self.window);
271        }
272
273        // Welford update plus anomaly counting against the *previous*
274        // statistics, so a sample never suppresses its own anomaly score.
275        if loss.is_finite() {
276            let previous_std = self.loss_std();
277            if self.loss_n >= 2 && previous_std > 0.0 {
278                let z = (loss - self.loss_mean).abs() / previous_std;
279                if z >= 3.0 {
280                    self.anomaly_count += 1;
281                }
282            }
283            self.loss_n += 1;
284            let delta = loss - self.loss_mean;
285            self.loss_mean += delta / self.loss_n as f64;
286            self.loss_m2 += delta * (loss - self.loss_mean);
287        }
288    }
289
290    fn loss_std(&self) -> f64 {
291        if self.loss_n < 2 {
292            0.0
293        } else {
294            (self.loss_m2 / self.loss_n as f64).sqrt()
295        }
296    }
297
298    /// Robust z-score of the most recent loss against the running statistics.
299    pub(crate) fn anomaly_score(&self) -> f64 {
300        let Some(&latest) = self.losses.back() else {
301            return 0.0;
302        };
303        let latest = from_scalar(latest);
304        let std = self.loss_std();
305        if !latest.is_finite() || std <= 0.0 {
306            0.0
307        } else {
308            (latest - self.loss_mean).abs() / std
309        }
310    }
311
312    pub(crate) fn anomaly_frequency(&self) -> f64 {
313        if self.loss_n == 0 {
314            0.0
315        } else {
316            self.anomaly_count as f64 / self.loss_n as f64
317        }
318    }
319
320    pub(crate) fn observed_span(&self) -> Duration {
321        match (self.first_timestamp, self.last_timestamp) {
322            (Some(first), Some(last)) => saturating_elapsed(last, first),
323            _ => Duration::ZERO,
324        }
325    }
326
327    pub(crate) fn losses_f64(&self) -> Vec<f64> {
328        self.losses.iter().map(|v| from_scalar(*v)).collect()
329    }
330
331    pub(crate) fn gradients_f64(&self) -> Vec<f64> {
332        self.gradients.iter().map(|v| from_scalar(*v)).collect()
333    }
334
335    pub(crate) fn rates_f64(&self) -> Vec<f64> {
336        self.inter_arrival
337            .iter()
338            .map(|d| d.as_secs_f64())
339            .filter(|s| *s > 0.0)
340            .map(|s| 1.0 / s)
341            .collect()
342    }
343
344    pub(crate) fn total_processing_time(&self) -> f64 {
345        self.latencies.iter().map(|d| d.as_secs_f64()).sum()
346    }
347
348    pub(crate) fn total_communication_time(&self) -> f64 {
349        self.communication_times
350            .iter()
351            .map(|d| d.as_secs_f64())
352            .sum()
353    }
354
355    /// Record an SLO evaluation outcome for the sample just ingested.
356    pub(crate) fn record_slo_outcome(&mut self, met: bool) {
357        self.slo_evaluated += 1;
358        if met {
359            self.slo_met += 1;
360        }
361    }
362
363    pub(crate) fn slo_compliance(&self) -> Option<f64> {
364        if self.slo_evaluated == 0 {
365            None
366        } else {
367            Some(self.slo_met as f64 / self.slo_evaluated as f64)
368        }
369    }
370
371    /// Report observed downtime.
372    pub fn record_outage(&mut self, downtime: Duration) {
373        self.downtime = self.downtime.saturating_add(downtime);
374    }
375
376    pub(crate) fn availability(&self) -> Option<f64> {
377        if self.downtime.is_zero() {
378            // Without a single reported outage there is no evidence about
379            // availability; reporting 100% would be an invention.
380            return None;
381        }
382        let span = self.observed_span().as_secs_f64() + self.downtime.as_secs_f64();
383        if span <= 0.0 {
384            return None;
385        }
386        Some((span - self.downtime.as_secs_f64()) / span)
387    }
388
389    /// Report a detected concept-drift event.
390    pub fn record_drift_event(
391        &mut self,
392        magnitude: A,
393        confidence: A,
394        detection_latency: Duration,
395        adaptation_effectiveness: Option<A>,
396    ) {
397        self.drift_events += 1;
398        self.last_drift = Some(DriftReport {
399            magnitude,
400            confidence,
401            detection_latency,
402            adaptation_effectiveness,
403        });
404    }
405
406    /// Report measured energy consumption.
407    pub fn record_energy(&mut self, joules: f64) {
408        self.energy_joules = Some(self.energy_joules.unwrap_or(0.0) + joules);
409    }
410
411    pub fn record_resource_probe(&mut self, probe: ResourceProbe) {
412        self.resource_probe = Some(probe);
413    }
414
415    pub fn record_robustness_probe(&mut self, probe: RobustnessProbe<A>) {
416        self.robustness_probe = Some(probe);
417    }
418}
419
420impl<A: Float + Default + Clone + std::fmt::Debug + Send + Sync> StreamingMetricsCollector<A> {
421    pub(crate) fn update_performance_metrics(&mut self, sample: &MetricsSample<A>) -> Result<()> {
422        let rates = self.accumulator.rates_f64();
423        let losses = self.accumulator.losses_f64();
424        let gradients = self.accumulator.gradients_f64();
425
426        // ---- throughput -------------------------------------------------
427        let throughput = &mut self.performance_metrics.throughput;
428        if !rates.is_empty() {
429            let mean_rate = mean_f64(&rates);
430            throughput.samples_per_second = mean_rate;
431            throughput.updates_per_second = mean_rate;
432            throughput.gradients_per_second = mean_rate;
433            throughput.throughput_variance = variance_f64(&rates);
434            throughput.throughput_trend = ols_slope(&rates);
435        }
436        if let Some(peak) = self.accumulator.peak_rate {
437            throughput.peak_throughput = peak;
438        }
439        if let Some(min) = self.accumulator.min_rate {
440            throughput.min_throughput = min;
441        }
442
443        // ---- latency ----------------------------------------------------
444        let latency = &mut self.performance_metrics.latency;
445        if let Some(stats) = duration_stats(&self.accumulator.latencies) {
446            latency.end_to_end = stats;
447        }
448        latency.gradient_computation = duration_stats(&self.accumulator.gradient_times);
449        latency.update_application = duration_stats(&self.accumulator.update_times);
450        latency.communication = duration_stats(&self.accumulator.communication_times);
451        latency.queue_wait_time = duration_stats(&self.accumulator.queue_times);
452        latency.jitter = {
453            let seconds: Vec<f64> = self
454                .accumulator
455                .latencies
456                .iter()
457                .map(|d| d.as_secs_f64())
458                .collect();
459            if seconds.len() < 2 {
460                0.0
461            } else {
462                seconds
463                    .windows(2)
464                    .map(|pair| (pair[1] - pair[0]).abs())
465                    .sum::<f64>()
466                    / (seconds.len() - 1) as f64
467            }
468        };
469
470        // ---- accuracy / convergence --------------------------------------
471        let span_seconds = self.accumulator.observed_span().as_secs_f64();
472        let first_loss = self.accumulator.first_loss.map(from_scalar);
473        let current_loss = from_scalar(sample.loss);
474        let loss_slope = ols_slope(&losses);
475
476        let accuracy = &mut self.performance_metrics.accuracy;
477        accuracy.current_loss = sample.loss;
478        accuracy.gradient_magnitude = sample.gradient_magnitude;
479        accuracy.loss_reduction_rate = match (first_loss, span_seconds > 0.0) {
480            (Some(first), true) => to_scalar((first - current_loss) / span_seconds),
481            _ => A::zero(),
482        };
483        // Positive when the loss is trending down.
484        accuracy.convergence_rate = to_scalar(-loss_slope);
485        accuracy.prediction_accuracy = sample.custom_metrics.get("accuracy").copied();
486        accuracy.parameter_stability = {
487            let spread = variance_f64(&gradients).sqrt();
488            to_scalar(1.0 / (1.0 + spread))
489        };
490        accuracy.learning_progress = match first_loss {
491            Some(first) if first.abs() > f64::EPSILON => {
492                to_scalar(((first - current_loss) / first.abs()).clamp(-1.0, 1.0))
493            }
494            _ => A::zero(),
495        };
496
497        // ---- stability ---------------------------------------------------
498        let increases = losses.windows(2).filter(|pair| pair[1] > pair[0]).count() as f64;
499        let sign_changes = losses
500            .windows(3)
501            .filter(|triple| {
502                let first = triple[1] - triple[0];
503                let second = triple[2] - triple[1];
504                first * second < 0.0
505            })
506            .count() as f64;
507
508        let stability = &mut self.performance_metrics.stability;
509        stability.loss_variance = to_scalar(variance_f64(&losses));
510        stability.gradient_variance = to_scalar(variance_f64(&gradients));
511        stability.parameter_drift = to_scalar(mean_f64(&gradients));
512        stability.oscillation_score = if losses.len() >= 3 {
513            to_scalar(sign_changes / (losses.len() - 2) as f64)
514        } else {
515            A::zero()
516        };
517        stability.divergence_probability = if losses.len() >= 2 {
518            to_scalar(increases / (losses.len() - 1) as f64)
519        } else {
520            A::zero()
521        };
522        stability.stability_confidence = {
523            let mean_loss = mean_f64(&losses).abs();
524            let spread = variance_f64(&losses).sqrt();
525            if mean_loss > f64::EPSILON {
526                to_scalar((1.0 - (spread / mean_loss)).clamp(0.0, 1.0))
527            } else {
528                A::zero()
529            }
530        };
531
532        // ---- efficiency ---------------------------------------------------
533        let processing_seconds = self.accumulator.total_processing_time();
534        let communication_seconds = self.accumulator.total_communication_time();
535        let memory_values: Vec<f64> = self
536            .accumulator
537            .memory
538            .iter()
539            .map(|bytes| *bytes as f64)
540            .collect();
541        let peak_memory = self.accumulator.peak_memory as f64;
542
543        let efficiency = &mut self.performance_metrics.efficiency;
544        efficiency.computational_efficiency = match (first_loss, processing_seconds > 0.0) {
545            (Some(first), true) => Some(to_scalar((first - current_loss) / processing_seconds)),
546            _ => None,
547        };
548        efficiency.memory_efficiency = if peak_memory > 0.0 && !memory_values.is_empty() {
549            Some(to_scalar(mean_f64(&memory_values) / peak_memory))
550        } else {
551            None
552        };
553        efficiency.communication_efficiency =
554            if !self.accumulator.communication_times.is_empty() && processing_seconds > 0.0 {
555                Some(to_scalar(
556                    (1.0 - communication_seconds / processing_seconds).clamp(0.0, 1.0),
557                ))
558            } else {
559                None
560            };
561        efficiency.energy_efficiency = match (self.accumulator.energy_joules, first_loss) {
562            (Some(joules), Some(first)) if joules > 0.0 => {
563                Some(to_scalar((first - current_loss) / joules))
564            }
565            _ => None,
566        };
567        efficiency.resource_utilization = if span_seconds > 0.0 {
568            to_scalar((processing_seconds / span_seconds).clamp(0.0, 1.0))
569        } else {
570            A::zero()
571        };
572
573        Ok(())
574    }
575
576    pub(crate) fn update_resource_metrics(&mut self, sample: &MetricsSample<A>) -> Result<()> {
577        let memory_values: Vec<f64> = self
578            .accumulator
579            .memory
580            .iter()
581            .map(|bytes| *bytes as f64)
582            .collect();
583        let peak = self.accumulator.peak_memory;
584
585        self.resource_metrics.memory_usage = MemoryUsage {
586            total_allocated: self
587                .accumulator
588                .resource_probe
589                .as_ref()
590                .and_then(|probe| probe.total_allocated_bytes),
591            current_used: sample.memory_usage,
592            peak_usage: peak,
593            fragmentation_ratio: self
594                .accumulator
595                .resource_probe
596                .as_ref()
597                .and_then(|probe| probe.fragmentation_ratio),
598            // Rust is not garbage collected; there is no overhead to report.
599            gc_overhead: None,
600            efficiency: if peak > 0 && !memory_values.is_empty() {
601                Some(mean_f64(&memory_values) / peak as f64)
602            } else {
603                None
604            },
605        };
606
607        if let Some(probe) = self.accumulator.resource_probe.as_ref() {
608            self.resource_metrics.cpu_utilization = probe.cpu_utilization;
609            self.resource_metrics.gpu_utilization = probe.gpu_utilization;
610            self.resource_metrics.network_bandwidth = probe.network_bandwidth_mbps;
611            self.resource_metrics.disk_io = probe.disk_io_mbps;
612            self.resource_metrics.thread_utilization = probe.thread_utilization;
613        }
614
615        Ok(())
616    }
617
618    pub(crate) fn update_quality_metrics(&mut self, sample: &MetricsSample<A>) -> Result<()> {
619        let losses = self.accumulator.losses_f64();
620        let current_loss = from_scalar(sample.loss);
621        let first_loss = self.accumulator.first_loss.map(from_scalar);
622
623        self.quality_metrics.data_quality = if self.accumulator.sample_count == 0 {
624            A::zero()
625        } else {
626            to_scalar(
627                self.accumulator.valid_sample_count as f64 / self.accumulator.sample_count as f64,
628            )
629        };
630
631        let validation_loss = sample
632            .custom_metrics
633            .get("val_loss")
634            .copied()
635            .map(from_scalar);
636        let model_quality = &mut self.quality_metrics.model_quality;
637        model_quality.training_quality = match first_loss {
638            Some(first) if first.abs() > f64::EPSILON => {
639                to_scalar(((first - current_loss) / first.abs()).clamp(0.0, 1.0))
640            }
641            _ => A::zero(),
642        };
643        model_quality.generalization_score = validation_loss.map(|validation| {
644            // 1 when validation matches training loss, decaying as the gap grows.
645            let gap = (validation - current_loss).abs();
646            to_scalar(1.0 / (1.0 + gap))
647        });
648        model_quality.overfitting_score = validation_loss.map(|validation| {
649            let denominator = current_loss.abs().max(f64::EPSILON);
650            to_scalar(((validation - current_loss) / denominator).max(0.0))
651        });
652        model_quality.underfitting_score = validation_loss.map(|validation| {
653            // Both losses high and close together indicates underfitting.
654            let denominator = current_loss.abs().max(f64::EPSILON);
655            let gap = ((validation - current_loss) / denominator).abs();
656            match first_loss {
657                Some(first) if first.abs() > f64::EPSILON => {
658                    to_scalar(((current_loss / first.abs()) * (1.0 - gap)).clamp(0.0, 1.0))
659                }
660                _ => to_scalar((1.0 - gap).clamp(0.0, 1.0)),
661            }
662        });
663        model_quality.complexity_score =
664            sample
665                .custom_metrics
666                .get("parameter_count")
667                .copied()
668                .map(|count| {
669                    let count = from_scalar(count).max(1.0);
670                    to_scalar(count.ln() / (1.0 + count.ln()))
671                });
672
673        let span_seconds = self.accumulator.observed_span().as_secs_f64();
674        let drift = &mut self.quality_metrics.concept_drift;
675        drift.drift_frequency = if span_seconds > 0.0 {
676            self.accumulator.drift_events as f64 / span_seconds
677        } else {
678            0.0
679        };
680        if let Some(report) = self.accumulator.last_drift.as_ref() {
681            drift.drift_confidence = Some(report.confidence);
682            drift.drift_magnitude = Some(report.magnitude);
683            drift.detection_latency = Some(report.detection_latency);
684            drift.adaptation_effectiveness = report.adaptation_effectiveness;
685        }
686
687        let anomaly = &mut self.quality_metrics.anomaly_detection;
688        anomaly.anomaly_score = to_scalar(self.accumulator.anomaly_score());
689        anomaly.anomaly_frequency = self.accumulator.anomaly_frequency();
690        // The remaining anomaly metrics need labelled ground truth, which a
691        // streaming collector never observes; they stay `None`.
692
693        if let Some(probe) = self.accumulator.robustness_probe.as_ref() {
694            self.quality_metrics.robustness = RobustnessMetrics {
695                noise_tolerance: probe.noise_tolerance,
696                adversarial_robustness: probe.adversarial_robustness,
697                perturbation_sensitivity: probe.perturbation_sensitivity,
698                recovery_capability: probe.recovery_capability,
699                fault_tolerance: probe.fault_tolerance,
700            };
701        }
702
703        let _ = losses;
704        Ok(())
705    }
706
707    pub(crate) fn update_business_metrics(&mut self, sample: &MetricsSample<A>) -> Result<()> {
708        if let Some(targets) = self.slo.clone() {
709            let mut met = true;
710            if let Some(limit) = targets.max_processing_time {
711                met &= sample.processing_time <= limit;
712            }
713            if let Some(limit) = targets.max_loss {
714                met &= from_scalar(sample.loss) <= limit;
715            }
716            if let Some(limit) = targets.max_memory_bytes {
717                met &= sample.memory_usage <= limit;
718            }
719            self.accumulator.record_slo_outcome(met);
720        }
721
722        self.business_metrics.slo_compliance = self.accumulator.slo_compliance();
723        self.business_metrics.availability = self.accumulator.availability();
724        self.business_metrics.user_satisfaction =
725            sample.custom_metrics.get("user_satisfaction").copied();
726
727        if let Some(model) = self.cost_model.clone() {
728            let compute_seconds = self.accumulator.total_processing_time();
729            let memory_gb_hours = {
730                let bytes: Vec<f64> = self
731                    .accumulator
732                    .memory
733                    .iter()
734                    .map(|value| *value as f64)
735                    .collect();
736                let mean_gb = mean_f64(&bytes) / (1024.0 * 1024.0 * 1024.0);
737                mean_gb * (self.accumulator.observed_span().as_secs_f64() / 3600.0)
738            };
739            let energy = self.accumulator.energy_joules;
740
741            let computational_cost =
742                to_scalar::<A>(compute_seconds) * model.compute_cost_per_second;
743            let infrastructure_cost =
744                to_scalar::<A>(memory_gb_hours) * model.memory_cost_per_gb_hour;
745            let energy_cost = energy.map(|j| to_scalar::<A>(j) * model.energy_cost_per_joule);
746
747            let first_loss = self.accumulator.first_loss.map(from_scalar);
748            let loss_reduction = match first_loss {
749                Some(first) => (first - from_scalar(sample.loss)).max(0.0),
750                None => 0.0,
751            };
752            let business_value = to_scalar::<A>(loss_reduction) * model.value_per_loss_unit;
753            let total =
754                computational_cost + infrastructure_cost + energy_cost.unwrap_or_else(A::zero);
755
756            self.business_metrics.cost_metrics = CostMetrics {
757                computational_cost: Some(computational_cost),
758                infrastructure_cost: Some(infrastructure_cost),
759                energy_cost,
760                // The value forgone by spending this compute elsewhere is the
761                // value it produced here, at the margin.
762                opportunity_cost: Some(business_value),
763                total_cost: Some(total),
764            };
765            self.business_metrics.business_value = Some(business_value);
766            self.performance_metrics.efficiency.cost_efficiency = if total > A::zero() {
767                Some(business_value / total)
768            } else {
769                None
770            };
771        }
772
773        Ok(())
774    }
775}