kmp_adapter_embedded/adapter/telemetry/
sqlite_quality_telemetry_reader.rs1use 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#[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}