Skip to main content

kmp_adapter_embedded/adapter/telemetry/
sqlite_quality_telemetry_reader.rs

1use std::path::Path;
2use std::sync::{Arc, Mutex, MutexGuard};
3
4use kmp_application::{
5    ObservabilityExemplar, ObservabilityMetricPoint, ObservabilityProjection, ObservabilityQuery,
6    ObservabilityQueryPort, ObservabilitySeries,
7};
8use kmp_domain::PortError;
9use kmp_observability::QualityTelemetryObservation;
10use rusqlite::{Connection, params};
11
12use super::storage::open_quality_connection;
13use crate::adapter::serdes::decode;
14
15/// Read-only query adapter for the shareable local quality journal.
16#[derive(Debug, Clone)]
17pub struct SqliteQualityTelemetryReader {
18    connection: Arc<Mutex<Connection>>,
19}
20
21impl ObservabilityQueryPort for SqliteQualityTelemetryReader {
22    fn available_series(&self) -> Vec<String> {
23        SUPPORTED_SERIES
24            .iter()
25            .map(|(name, _, _)| (*name).to_string())
26            .collect()
27    }
28
29    fn query<'a>(
30        &'a self,
31        query: ObservabilityQuery,
32    ) -> std::pin::Pin<
33        Box<
34            dyn std::future::Future<Output = Result<ObservabilityProjection, PortError>>
35                + Send
36                + 'a,
37        >,
38    > {
39        Box::pin(async move {
40            let requested = if query.series.is_empty() {
41                SUPPORTED_SERIES
42                    .iter()
43                    .map(|(name, _, _)| (*name).to_string())
44                    .collect::<Vec<_>>()
45            } else {
46                query.series.clone()
47            };
48            if query.to_millis < query.from_millis || query.max_points == 0 {
49                return Ok(ObservabilityProjection::empty(&query));
50            }
51            let mut observations = self.query_between_filtered(
52                query.from_millis,
53                query.to_millis,
54                None,
55                query.about.as_deref(),
56                query.max_points.saturating_add(1),
57            )?;
58            let truncated = observations.len() > query.max_points;
59            observations.truncate(query.max_points);
60            let exemplars = observations
61                .iter()
62                .enumerate()
63                .map(|(index, observation)| {
64                    let id = format!("local-quality:{}:{index}", observation.observed_at_millis());
65                    let mut attributes = std::collections::BTreeMap::new();
66                    attributes.insert("role".to_string(), observation.role().to_string());
67                    ObservabilityExemplar {
68                        id,
69                        at_millis: observation.observed_at_millis(),
70                        operation: observation.rpc().to_string(),
71                        about: (!observation.root_node_id().is_empty())
72                            .then(|| observation.root_node_id().to_string()),
73                        bundle_ref: (!observation.root_node_id().is_empty())
74                            .then(|| observation.root_node_id().to_string()),
75                        revision: observation.revision(),
76                        attributes,
77                    }
78                })
79                .collect::<Vec<_>>();
80            let mut series = Vec::new();
81            let mut missing = Vec::new();
82            for requested_name in requested {
83                let Some((name, unit, scope)) = SUPPORTED_SERIES
84                    .iter()
85                    .find(|(name, _, _)| *name == requested_name)
86                else {
87                    missing.push(requested_name);
88                    continue;
89                };
90                let points = observations
91                    .iter()
92                    .zip(&exemplars)
93                    .map(|(observation, exemplar)| ObservabilityMetricPoint {
94                        at_millis: observation.observed_at_millis(),
95                        value: quality_metric(observation, name),
96                        exemplar_id: Some(exemplar.id.clone()),
97                    })
98                    .collect();
99                series.push(ObservabilitySeries {
100                    name: (*name).to_string(),
101                    unit: (*unit).to_string(),
102                    scope: (*scope).to_string(),
103                    points,
104                });
105            }
106            Ok(ObservabilityProjection {
107                contract: "kmp.observability.projection.v1".to_string(),
108                from_millis: query.from_millis,
109                to_millis: query.to_millis,
110                series,
111                exemplars,
112                missing,
113                truncated,
114            })
115        })
116    }
117}
118
119const SUPPORTED_SERIES: &[(&str, &str, &str)] = &[
120    ("raw_equivalent_tokens", "tokens", "rendered_bundle"),
121    ("compression_ratio", "ratio", "rendered_bundle"),
122    ("causal_density", "ratio", "rendered_bundle"),
123    ("noise_ratio", "ratio", "rendered_bundle"),
124    ("detail_coverage", "ratio", "rendered_bundle"),
125];
126
127fn quality_metric(observation: &QualityTelemetryObservation, name: &str) -> f64 {
128    match name {
129        "raw_equivalent_tokens" => f64::from(observation.raw_equivalent_tokens()),
130        "compression_ratio" => observation.compression_ratio(),
131        "causal_density" => observation.causal_density(),
132        "noise_ratio" => observation.noise_ratio(),
133        "detail_coverage" => observation.detail_coverage(),
134        _ => 0.0,
135    }
136}
137
138impl SqliteQualityTelemetryReader {
139    pub(super) fn from_connection(connection: Arc<Mutex<Connection>>) -> Self {
140        Self { connection }
141    }
142
143    pub fn open(data_dir: &Path) -> Result<Self, PortError> {
144        let connection = open_quality_connection(data_dir)?;
145        Ok(Self::from_connection(Arc::new(Mutex::new(connection))))
146    }
147
148    pub fn count(&self) -> Result<u64, PortError> {
149        let count: i64 = self
150            .connection()?
151            .query_row("SELECT COUNT(*) FROM quality_observations", [], |row| {
152                row.get(0)
153            })
154            .map_err(|error| {
155                PortError::Unavailable(format!("quality telemetry count failed: {error}"))
156            })?;
157        u64::try_from(count)
158            .map_err(|_| PortError::InvalidState("quality telemetry count is negative".to_string()))
159    }
160
161    pub fn query_since(
162        &self,
163        since_millis: u64,
164        rpc: Option<&str>,
165        limit: usize,
166    ) -> Result<Vec<QualityTelemetryObservation>, PortError> {
167        self.query_between(since_millis, u64::MAX, rpc, limit)
168    }
169
170    pub fn query_between(
171        &self,
172        since_millis: u64,
173        until_millis: u64,
174        rpc: Option<&str>,
175        limit: usize,
176    ) -> Result<Vec<QualityTelemetryObservation>, PortError> {
177        self.query_between_filtered(since_millis, until_millis, rpc, None, limit)
178    }
179
180    fn query_between_filtered(
181        &self,
182        since_millis: u64,
183        until_millis: u64,
184        rpc: Option<&str>,
185        root_node_id: Option<&str>,
186        limit: usize,
187    ) -> Result<Vec<QualityTelemetryObservation>, PortError> {
188        if limit == 0 || until_millis < since_millis {
189            return Ok(Vec::new());
190        }
191        let from = i64::try_from(since_millis).unwrap_or(i64::MAX);
192        let to = i64::try_from(until_millis).unwrap_or(i64::MAX);
193        let connection = self.connection()?;
194        let mut statement = connection
195            .prepare(
196                "SELECT payload FROM quality_observations \
197                 WHERE observed_at_millis BETWEEN ?1 AND ?2 \
198                 ORDER BY observed_at_millis ASC, id ASC",
199            )
200            .map_err(query_error)?;
201        let rows = statement
202            .query_map(params![from, to], |row| row.get::<_, Vec<u8>>(0))
203            .map_err(query_error)?;
204        let mut observations = Vec::new();
205        for row in rows {
206            let payload = row.map_err(query_error)?;
207            let observation: QualityTelemetryObservation = decode("quality observation", &payload)?;
208            if rpc.is_some_and(|wanted| wanted != observation.rpc()) {
209                continue;
210            }
211            if root_node_id.is_some_and(|wanted| wanted != observation.root_node_id()) {
212                continue;
213            }
214            observations.push(observation);
215            if observations.len() >= limit {
216                break;
217            }
218        }
219        Ok(observations)
220    }
221
222    pub fn latest(&self, limit: usize) -> Result<Vec<QualityTelemetryObservation>, PortError> {
223        if limit == 0 {
224            return Ok(Vec::new());
225        }
226        let connection = self.connection()?;
227        let mut statement = connection
228            .prepare(
229                "SELECT payload FROM quality_observations \
230                 ORDER BY observed_at_millis DESC, id DESC LIMIT ?1",
231            )
232            .map_err(query_error)?;
233        let rows = statement
234            .query_map([i64::try_from(limit).unwrap_or(i64::MAX)], |row| {
235                row.get::<_, Vec<u8>>(0)
236            })
237            .map_err(query_error)?;
238        let mut observations = rows
239            .map(|row| {
240                let payload = row.map_err(query_error)?;
241                decode("quality observation", &payload)
242            })
243            .collect::<Result<Vec<_>, _>>()?;
244        observations.reverse();
245        Ok(observations)
246    }
247
248    fn connection(&self) -> Result<MutexGuard<'_, Connection>, PortError> {
249        self.connection.lock().map_err(|_| {
250            PortError::Unavailable("quality telemetry connection lock is poisoned".to_string())
251        })
252    }
253}
254
255fn query_error(error: rusqlite::Error) -> PortError {
256    PortError::Unavailable(format!("quality telemetry query failed: {error}"))
257}