kmp_adapter_embedded/adapter/telemetry/
redb_quality_telemetry_writer.rs1use 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::serdes::encode;
12use crate::adapter::store::{commit_error, range_error, storage_error, table_error};
13
14const DEFAULT_DURABLE_EVERY_BATCHES: u64 = 16;
15
16#[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}