kmp_adapter_embedded/adapter/telemetry/
redb_quality_telemetry_reader.rs1use std::path::Path;
2use std::sync::Arc;
3
4use kmp_domain::PortError;
5use kmp_observability::QualityTelemetryObservation;
6use redb::{Database, ReadableDatabase, ReadableTable, ReadableTableMetadata};
7
8use super::storage::{OBSERVATIONS, quality_telemetry_path};
9use crate::adapter::engine::redb::{range_error, table_error};
10use crate::adapter::serdes::decode;
11
12#[derive(Debug, Clone)]
14pub struct RedbQualityTelemetryReader {
15 database: Arc<Database>,
16}
17
18impl RedbQualityTelemetryReader {
19 pub fn open(data_dir: &Path) -> Result<Self, PortError> {
20 let path = quality_telemetry_path(data_dir);
21 let database = Database::open(&path).map_err(|error| {
22 PortError::Unavailable(format!(
23 "quality telemetry could not open `{}` for reading: {error}",
24 path.display()
25 ))
26 })?;
27 Ok(Self {
28 database: Arc::new(database),
29 })
30 }
31
32 pub fn count(&self) -> Result<u64, PortError> {
33 let tx = self.begin_read()?;
34 let table = tx.open_table(OBSERVATIONS).map_err(table_error)?;
35 table.len().map_err(|error| {
36 PortError::Unavailable(format!("quality telemetry count failed: {error}"))
37 })
38 }
39
40 pub fn query_since(
41 &self,
42 since_millis: u64,
43 rpc: Option<&str>,
44 limit: usize,
45 ) -> Result<Vec<QualityTelemetryObservation>, PortError> {
46 self.query_between(since_millis, u64::MAX, rpc, limit)
47 }
48
49 pub fn query_between(
50 &self,
51 since_millis: u64,
52 until_millis: u64,
53 rpc: Option<&str>,
54 limit: usize,
55 ) -> Result<Vec<QualityTelemetryObservation>, PortError> {
56 if limit == 0 || until_millis < since_millis {
57 return Ok(Vec::new());
58 }
59 let tx = self.begin_read()?;
60 let table = tx.open_table(OBSERVATIONS).map_err(table_error)?;
61 let mut observations = Vec::new();
62 for row in table
63 .range((since_millis, 0u64)..=(until_millis, u64::MAX))
64 .map_err(range_error)?
65 {
66 let (_, value) = row.map_err(range_error)?;
67 let observation: QualityTelemetryObservation =
68 decode("quality observation", value.value())?;
69 if rpc.is_some_and(|wanted| wanted != observation.rpc()) {
70 continue;
71 }
72 observations.push(observation);
73 if observations.len() >= limit {
74 break;
75 }
76 }
77 Ok(observations)
78 }
79
80 pub fn latest(&self, limit: usize) -> Result<Vec<QualityTelemetryObservation>, PortError> {
81 if limit == 0 {
82 return Ok(Vec::new());
83 }
84 let tx = self.begin_read()?;
85 let table = tx.open_table(OBSERVATIONS).map_err(table_error)?;
86 let mut observations = Vec::new();
87 for row in table.iter().map_err(range_error)?.rev().take(limit) {
88 let (_, value) = row.map_err(range_error)?;
89 observations.push(decode("quality observation", value.value())?);
90 }
91 observations.reverse();
92 Ok(observations)
93 }
94
95 fn begin_read(&self) -> Result<redb::ReadTransaction, PortError> {
96 self.database.begin_read().map_err(|error| {
97 PortError::Unavailable(format!(
98 "quality telemetry read transaction failed: {error}"
99 ))
100 })
101 }
102}