Skip to main content

kmp_adapter_embedded/adapter/telemetry/
redb_quality_telemetry_writer.rs

1use std::path::Path;
2use std::sync::Arc;
3use std::sync::atomic::{AtomicU64, Ordering};
4
5use kmp_domain::PortError;
6use kmp_observability::QualityTelemetryObservation;
7use redb::{Database, Durability, ReadableDatabase, ReadableTable, ReadableTableMetadata};
8
9use super::quality_telemetry_retention::QualityTelemetryRetention;
10use super::storage::{OBSERVATIONS, quality_telemetry_path};
11use crate::adapter::engine::redb::{commit_error, range_error, storage_error, table_error};
12use crate::adapter::serdes::encode;
13
14const DEFAULT_DURABLE_EVERY_BATCHES: u64 = 16;
15
16/// Relaxed-durability writer for `telemetry/quality.redb`.
17#[derive(Debug)]
18pub struct RedbQualityTelemetryWriter {
19    database: Arc<Database>,
20    retention: QualityTelemetryRetention,
21    next_sequence: AtomicU64,
22    batch_number: AtomicU64,
23    durable_every_batches: u64,
24    write_failures: AtomicU64,
25}
26
27impl RedbQualityTelemetryWriter {
28    pub fn open(data_dir: &Path, retention: QualityTelemetryRetention) -> Result<Self, PortError> {
29        Self::open_with_durable_interval(data_dir, retention, DEFAULT_DURABLE_EVERY_BATCHES)
30    }
31
32    pub fn open_with_durable_interval(
33        data_dir: &Path,
34        retention: QualityTelemetryRetention,
35        durable_every_batches: u64,
36    ) -> Result<Self, PortError> {
37        if durable_every_batches == 0 {
38            return Err(PortError::Unavailable(
39                "quality telemetry durable interval must be greater than zero".to_string(),
40            ));
41        }
42        let path = quality_telemetry_path(data_dir);
43        let parent = path.parent().expect("quality telemetry path has a parent");
44        std::fs::create_dir_all(parent).map_err(|error| {
45            PortError::Unavailable(format!(
46                "quality telemetry could not create `{}`: {error}",
47                parent.display()
48            ))
49        })?;
50        let database = Arc::new(Database::create(&path).map_err(|error| {
51            PortError::Unavailable(format!(
52                "quality telemetry could not open `{}`: {error}",
53                path.display()
54            ))
55        })?);
56        initialize_table(&database)?;
57        let next_sequence = load_highest_sequence(&database)?;
58        Ok(Self {
59            database,
60            retention,
61            next_sequence: AtomicU64::new(next_sequence),
62            batch_number: AtomicU64::new(0),
63            durable_every_batches,
64            write_failures: AtomicU64::new(0),
65        })
66    }
67
68    pub fn write_batch(
69        &self,
70        observations: &[QualityTelemetryObservation],
71    ) -> Result<(), PortError> {
72        if observations.is_empty() {
73            return Ok(());
74        }
75        let result = self.write_batch_inner(observations);
76        if result.is_err() {
77            self.write_failures.fetch_add(1, Ordering::Relaxed);
78        }
79        result
80    }
81
82    pub fn flush_durable(&self) -> Result<(), PortError> {
83        let result = self.flush_durable_inner();
84        if result.is_err() {
85            self.write_failures.fetch_add(1, Ordering::Relaxed);
86        }
87        result
88    }
89
90    pub fn write_failures(&self) -> u64 {
91        self.write_failures.load(Ordering::Relaxed)
92    }
93
94    fn write_batch_inner(
95        &self,
96        observations: &[QualityTelemetryObservation],
97    ) -> Result<(), PortError> {
98        let current_batch = self.batch_number.fetch_add(1, Ordering::Relaxed) + 1;
99        let durability = if current_batch.is_multiple_of(self.durable_every_batches) {
100            Durability::Immediate
101        } else {
102            Durability::None
103        };
104        let mut tx = self.database.begin_write().map_err(|error| {
105            PortError::Unavailable(format!(
106                "quality telemetry write transaction failed: {error}"
107            ))
108        })?;
109        tx.set_durability(durability).map_err(|error| {
110            PortError::Unavailable(format!(
111                "quality telemetry durability configuration failed: {error}"
112            ))
113        })?;
114        {
115            let mut table = tx.open_table(OBSERVATIONS).map_err(table_error)?;
116            for observation in observations {
117                let sequence = self.next_sequence.fetch_add(1, Ordering::Relaxed) + 1;
118                let bytes = encode("quality observation", observation)?;
119                table
120                    .insert(
121                        (observation.observed_at_millis(), sequence),
122                        bytes.as_slice(),
123                    )
124                    .map_err(storage_error)?;
125            }
126            let excess = self.retention.excess(table.len().map_err(storage_error)?);
127            for _ in 0..excess {
128                table.pop_first().map_err(storage_error)?;
129            }
130        }
131        tx.commit().map_err(commit_error)
132    }
133
134    fn flush_durable_inner(&self) -> Result<(), PortError> {
135        let mut tx = self.database.begin_write().map_err(|error| {
136            PortError::Unavailable(format!(
137                "quality telemetry durable flush failed to start: {error}"
138            ))
139        })?;
140        tx.set_durability(Durability::Immediate).map_err(|error| {
141            PortError::Unavailable(format!(
142                "quality telemetry durable flush configuration failed: {error}"
143            ))
144        })?;
145        tx.commit().map_err(commit_error)
146    }
147}
148
149fn initialize_table(database: &Database) -> Result<(), PortError> {
150    let tx = database.begin_write().map_err(|error| {
151        PortError::Unavailable(format!("quality telemetry initialization failed: {error}"))
152    })?;
153    tx.open_table(OBSERVATIONS).map_err(table_error)?;
154    tx.commit().map_err(commit_error)
155}
156
157fn load_highest_sequence(database: &Database) -> Result<u64, PortError> {
158    let tx = database.begin_read().map_err(|error| {
159        PortError::Unavailable(format!("quality telemetry sequence read failed: {error}"))
160    })?;
161    let table = tx.open_table(OBSERVATIONS).map_err(table_error)?;
162    let mut highest = 0u64;
163    for row in table.iter().map_err(range_error)? {
164        let (key, _) = row.map_err(range_error)?;
165        highest = highest.max(key.value().1);
166    }
167    Ok(highest)
168}