Skip to main content

kmp_observability/
buffered_quality_metrics_observer.rs

1//! Runtime-neutral, non-blocking quality-telemetry buffering.
2
3use std::sync::atomic::{AtomicU64, Ordering};
4use std::sync::mpsc::{Receiver, SyncSender, sync_channel};
5
6use kmp_domain::{BundleQualityMetrics, QualityMetricsObserver, QualityObservationContext};
7
8use crate::QualityTelemetryObservation;
9
10/// Implements `QualityMetricsObserver` over a bounded fail-open channel.
11pub struct BufferedQualityMetricsObserver {
12    sender: SyncSender<QualityTelemetryObservation>,
13    dropped: AtomicU64,
14}
15
16impl BufferedQualityMetricsObserver {
17    pub fn with_capacity(capacity: usize) -> (Self, Receiver<QualityTelemetryObservation>) {
18        let (sender, receiver) = sync_channel(capacity);
19        (
20            Self {
21                sender,
22                dropped: AtomicU64::new(0),
23            },
24            receiver,
25        )
26    }
27
28    /// Observations discarded because the buffer was full or disconnected.
29    pub fn dropped_observations(&self) -> u64 {
30        self.dropped.load(Ordering::Relaxed)
31    }
32}
33
34impl QualityMetricsObserver for BufferedQualityMetricsObserver {
35    fn observe(&self, metrics: &BundleQualityMetrics, context: &QualityObservationContext) {
36        let observation = QualityTelemetryObservation::capture(metrics, context);
37        if self.sender.try_send(observation).is_err() {
38            let dropped = self.dropped.fetch_add(1, Ordering::Relaxed) + 1;
39            if dropped == 1 || dropped.is_multiple_of(1000) {
40                tracing::warn!(
41                    dropped,
42                    "quality telemetry buffer unavailable; dropping observations"
43                );
44            }
45        }
46    }
47}
48
49#[cfg(test)]
50mod tests {
51    use kmp_domain::{BundleQualityMetrics, QualityMetricsObserver, QualityObservationContext};
52
53    use super::BufferedQualityMetricsObserver;
54
55    fn sample_metrics() -> BundleQualityMetrics {
56        BundleQualityMetrics::new(120, 2.5, 0.4, 0.1, 0.8).expect("valid metrics")
57    }
58
59    fn sample_context() -> QualityObservationContext {
60        QualityObservationContext {
61            rpc: "kernel_wake".to_string(),
62            root_node_id: "question:t".to_string(),
63            role: "resumer".to_string(),
64        }
65    }
66
67    #[test]
68    fn observations_flow_through_the_bounded_channel() {
69        let (observer, receiver) = BufferedQualityMetricsObserver::with_capacity(4);
70
71        observer.observe(&sample_metrics(), &sample_context());
72
73        let observation = receiver.try_recv().expect("observation buffered");
74        assert_eq!(observation.rpc(), "kernel_wake");
75        assert_eq!(observation.root_node_id(), "question:t");
76        assert_eq!(observation.raw_equivalent_tokens(), 120);
77        assert!(observation.observed_at_millis() > 0);
78        assert_eq!(observer.dropped_observations(), 0);
79    }
80
81    #[test]
82    fn overflow_drops_and_counts_without_blocking() {
83        let (observer, receiver) = BufferedQualityMetricsObserver::with_capacity(1);
84
85        observer.observe(&sample_metrics(), &sample_context());
86        observer.observe(&sample_metrics(), &sample_context());
87        observer.observe(&sample_metrics(), &sample_context());
88
89        assert_eq!(observer.dropped_observations(), 2);
90        assert!(receiver.try_recv().is_ok());
91        assert!(receiver.try_recv().is_err());
92    }
93
94    #[test]
95    fn disconnected_worker_never_fails_the_kernel() {
96        let (observer, receiver) = BufferedQualityMetricsObserver::with_capacity(4);
97        drop(receiver);
98
99        observer.observe(&sample_metrics(), &sample_context());
100
101        assert_eq!(observer.dropped_observations(), 1);
102    }
103}