Skip to main content

kmp_observability/
quality_observers.rs

1//! Adapters for the [`QualityMetricsObserver`] domain port.
2//!
3//! - [`OTelQualityObserver`]: OpenTelemetry histograms (Prometheus/Grafana via OTLP)
4//! - [`TracingQualityObserver`]: Structured tracing logs (Loki/Grafana via Promtail)
5//! - [`CompositeQualityObserver`]: Fan-out to multiple observers
6
7use std::sync::Arc;
8
9use kmp_domain::{BundleQualityMetrics, QualityMetricsObserver, QualityObservationContext};
10use opentelemetry::metrics::{Histogram, Meter};
11
12// ── OTel adapter ────────────────────────────────────────────────────────
13
14/// Emits quality metrics as OpenTelemetry histograms.
15///
16/// Requires `OTEL_EXPORTER_OTLP_ENDPOINT` to export; otherwise instruments
17/// discard data silently (noop meter behavior).
18pub struct OTelQualityObserver {
19    raw_equivalent_tokens: Histogram<u64>,
20    compression_ratio: Histogram<f64>,
21    causal_density: Histogram<f64>,
22    noise_ratio: Histogram<f64>,
23    detail_coverage: Histogram<f64>,
24}
25
26impl OTelQualityObserver {
27    pub fn new(meter: &Meter) -> Self {
28        Self {
29            raw_equivalent_tokens: meter
30                .u64_histogram("rehydration.quality.raw_equivalent_tokens")
31                .with_description("Flat text token count baseline")
32                .build(),
33            compression_ratio: meter
34                .f64_histogram("rehydration.quality.compression_ratio")
35                .with_description("Raw / rendered token ratio")
36                .build(),
37            causal_density: meter
38                .f64_histogram("rehydration.quality.causal_density")
39                .with_description("Fraction of explanatory relationships")
40                .build(),
41            noise_ratio: meter
42                .f64_histogram("rehydration.quality.noise_ratio")
43                .with_description("Fraction of noise/distractor nodes")
44                .build(),
45            detail_coverage: meter
46                .f64_histogram("rehydration.quality.detail_coverage")
47                .with_description("Fraction of nodes with detail")
48                .build(),
49        }
50    }
51}
52
53impl QualityMetricsObserver for OTelQualityObserver {
54    fn observe(&self, metrics: &BundleQualityMetrics, context: &QualityObservationContext) {
55        let attrs = &[opentelemetry::KeyValue::new("rpc", context.rpc.clone())];
56        self.raw_equivalent_tokens
57            .record(metrics.raw_equivalent_tokens() as u64, attrs);
58        self.compression_ratio
59            .record(metrics.compression_ratio(), attrs);
60        self.causal_density.record(metrics.causal_density(), attrs);
61        self.noise_ratio.record(metrics.noise_ratio(), attrs);
62        self.detail_coverage
63            .record(metrics.detail_coverage(), attrs);
64    }
65}
66
67// ── Tracing / Loki adapter ──────────────────────────────────────────────
68
69/// Emits quality metrics as structured tracing log events.
70///
71/// When the kernel runs with `KMP_LOG_FORMAT=json`, these become
72/// structured JSON log lines that Promtail / Grafana Agent collect and
73/// push to Loki. Grafana can then query via LogQL:
74///
75/// ```logql
76/// {job="kmp"} | json | quality_compression_ratio > 1.5
77/// ```
78pub struct TracingQualityObserver;
79
80impl QualityMetricsObserver for TracingQualityObserver {
81    fn observe(&self, metrics: &BundleQualityMetrics, context: &QualityObservationContext) {
82        tracing::info!(
83            target: "rehydration.quality",
84            rpc = %context.rpc,
85            root_node_id = %context.root_node_id,
86            role = %context.role,
87            quality_raw_equivalent_tokens = metrics.raw_equivalent_tokens(),
88            quality_compression_ratio = metrics.compression_ratio(),
89            quality_causal_density = metrics.causal_density(),
90            quality_noise_ratio = metrics.noise_ratio(),
91            quality_detail_coverage = metrics.detail_coverage(),
92            "bundle quality metrics"
93        );
94    }
95}
96
97// ── Composite adapter ───────────────────────────────────────────────────
98
99/// Fan-out observer that delegates to multiple backends.
100///
101/// ```rust,ignore
102/// use kmp_observability::quality_observers::*;
103/// let meter = opentelemetry::global::meter("example");
104/// let observer = CompositeQualityObserver::new(vec![
105///     Box::new(OTelQualityObserver::new(&meter)),
106///     Box::new(TracingQualityObserver),
107/// ]);
108/// ```
109/// Fan-out observer that spawns each adapter on a background task.
110///
111/// Individual adapters run fire-and-forget via `tokio::spawn`, keeping
112/// the gRPC handler hot path free from observer I/O latency.
113pub struct CompositeQualityObserver {
114    observers: Arc<Vec<Box<dyn QualityMetricsObserver>>>,
115}
116
117impl CompositeQualityObserver {
118    pub fn new(observers: Vec<Box<dyn QualityMetricsObserver>>) -> Self {
119        Self {
120            observers: Arc::new(observers),
121        }
122    }
123}
124
125impl QualityMetricsObserver for CompositeQualityObserver {
126    fn observe(&self, metrics: &BundleQualityMetrics, context: &QualityObservationContext) {
127        let observers = Arc::clone(&self.observers);
128        let metrics = metrics.clone();
129        let context = context.clone();
130        tokio::spawn(async move {
131            for observer in observers.iter() {
132                observer.observe(&metrics, &context);
133            }
134        });
135    }
136}
137
138// ── Noop adapter (for tests / when observability is disabled) ───────────
139
140/// No-op observer that discards all metrics.
141pub struct NoopQualityObserver;
142
143impl QualityMetricsObserver for NoopQualityObserver {
144    fn observe(&self, _metrics: &BundleQualityMetrics, _context: &QualityObservationContext) {}
145}
146
147#[cfg(test)]
148mod tests {
149    use std::sync::{Arc, Mutex};
150
151    use kmp_domain::{BundleQualityMetrics, QualityMetricsObserver, QualityObservationContext};
152
153    use super::{
154        CompositeQualityObserver, NoopQualityObserver, OTelQualityObserver, TracingQualityObserver,
155    };
156
157    fn sample_metrics() -> BundleQualityMetrics {
158        BundleQualityMetrics::new(200, 1.5, 0.6, 0.1, 0.8).expect("valid")
159    }
160
161    fn sample_context() -> QualityObservationContext {
162        QualityObservationContext {
163            rpc: "GetContext".to_string(),
164            root_node_id: "node:case:123".to_string(),
165            role: "developer".to_string(),
166        }
167    }
168
169    #[test]
170    fn noop_observer_does_not_panic() {
171        NoopQualityObserver.observe(&sample_metrics(), &sample_context());
172    }
173
174    #[test]
175    fn otel_observer_does_not_panic_with_noop_meter() {
176        let meter = opentelemetry::global::meter("test");
177        let observer = OTelQualityObserver::new(&meter);
178        observer.observe(&sample_metrics(), &sample_context());
179    }
180
181    #[test]
182    fn tracing_observer_does_not_panic() {
183        TracingQualityObserver.observe(&sample_metrics(), &sample_context());
184    }
185
186    /// Recording spy for verifying composite fan-out.
187    struct SpyObserver {
188        count: Arc<Mutex<u32>>,
189    }
190
191    impl SpyObserver {
192        fn new(count: Arc<Mutex<u32>>) -> Self {
193            Self { count }
194        }
195    }
196
197    impl QualityMetricsObserver for SpyObserver {
198        fn observe(&self, _metrics: &BundleQualityMetrics, _context: &QualityObservationContext) {
199            *self.count.lock().expect("mutex not poisoned") += 1;
200        }
201    }
202
203    #[tokio::test]
204    async fn composite_fans_out_to_all_observers() {
205        let count = Arc::new(Mutex::new(0u32));
206        let observer = CompositeQualityObserver::new(vec![
207            Box::new(SpyObserver::new(Arc::clone(&count))),
208            Box::new(SpyObserver::new(Arc::clone(&count))),
209            Box::new(SpyObserver::new(Arc::clone(&count))),
210        ]);
211        observer.observe(&sample_metrics(), &sample_context());
212        // Yield to let the spawned task run
213        tokio::task::yield_now().await;
214        assert_eq!(*count.lock().expect("mutex not poisoned"), 3);
215    }
216}