Skip to main content

kmp_adapter_embedded/adapter/telemetry/
redb_quality_telemetry_reader.rs

1use 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/// Read-only query adapter for the local quality journal.
13#[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}