kmp_observability/
buffered_quality_metrics_observer.rs1use 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
10pub 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 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}