kmp_adapter_embedded/adapter/telemetry/
sqlite_quality_telemetry_writer.rs1use std::path::Path;
2use std::sync::atomic::{AtomicU64, Ordering};
3use std::sync::{Arc, Mutex};
4
5use kmp_domain::PortError;
6use kmp_observability::QualityTelemetryObservation;
7use rusqlite::{Connection, TransactionBehavior};
8
9use super::quality_telemetry_retention::QualityTelemetryRetention;
10use super::sqlite_quality_telemetry_reader::SqliteQualityTelemetryReader;
11use super::storage::{enforce_retention, insert_observation, open_quality_connection};
12
13const DEFAULT_DURABLE_EVERY_BATCHES: u64 = 16;
14
15#[derive(Debug)]
17pub struct SqliteQualityTelemetryWriter {
18 connection: Arc<Mutex<Connection>>,
19 retention: QualityTelemetryRetention,
20 batch_number: AtomicU64,
21 durable_every_batches: u64,
22 write_failures: AtomicU64,
23}
24
25impl SqliteQualityTelemetryWriter {
26 pub fn reader(&self) -> SqliteQualityTelemetryReader {
27 SqliteQualityTelemetryReader::from_connection(Arc::clone(&self.connection))
28 }
29
30 pub fn open(data_dir: &Path, retention: QualityTelemetryRetention) -> Result<Self, PortError> {
31 Self::open_with_durable_interval(data_dir, retention, DEFAULT_DURABLE_EVERY_BATCHES)
32 }
33
34 pub fn open_with_durable_interval(
35 data_dir: &Path,
36 retention: QualityTelemetryRetention,
37 durable_every_batches: u64,
38 ) -> Result<Self, PortError> {
39 if durable_every_batches == 0 {
40 return Err(PortError::Unavailable(
41 "quality telemetry durable interval must be greater than zero".to_string(),
42 ));
43 }
44 let connection = open_quality_connection(data_dir)?;
45 Ok(Self {
46 connection: Arc::new(Mutex::new(connection)),
47 retention,
48 batch_number: AtomicU64::new(0),
49 durable_every_batches,
50 write_failures: AtomicU64::new(0),
51 })
52 }
53
54 pub fn write_batch(
55 &self,
56 observations: &[QualityTelemetryObservation],
57 ) -> Result<(), PortError> {
58 if observations.is_empty() {
59 return Ok(());
60 }
61 let result = self.write_batch_inner(observations);
62 if result.is_err() {
63 self.write_failures.fetch_add(1, Ordering::Relaxed);
64 }
65 result
66 }
67
68 pub fn flush_durable(&self) -> Result<(), PortError> {
69 let result = self.flush_durable_inner();
70 if result.is_err() {
71 self.write_failures.fetch_add(1, Ordering::Relaxed);
72 }
73 result
74 }
75
76 pub fn write_failures(&self) -> u64 {
77 self.write_failures.load(Ordering::Relaxed)
78 }
79
80 fn write_batch_inner(
81 &self,
82 observations: &[QualityTelemetryObservation],
83 ) -> Result<(), PortError> {
84 let current_batch = self.batch_number.fetch_add(1, Ordering::Relaxed) + 1;
85 let synchronous = if current_batch.is_multiple_of(self.durable_every_batches) {
86 "FULL"
87 } else {
88 "NORMAL"
89 };
90 let mut connection = self.connection.lock().map_err(|_| poisoned())?;
91 connection
92 .pragma_update(None, "synchronous", synchronous)
93 .map_err(write_error)?;
94 let transaction = connection
95 .transaction_with_behavior(TransactionBehavior::Immediate)
96 .map_err(write_error)?;
97 for observation in observations {
98 insert_observation(&transaction, observation)?;
99 }
100 enforce_retention(&transaction, self.retention)?;
101 transaction.commit().map_err(write_error)
102 }
103
104 fn flush_durable_inner(&self) -> Result<(), PortError> {
105 let connection = self.connection.lock().map_err(|_| poisoned())?;
106 connection
107 .pragma_update(None, "synchronous", "FULL")
108 .map_err(write_error)?;
109 connection
110 .execute_batch("BEGIN IMMEDIATE; COMMIT; PRAGMA wal_checkpoint(FULL);")
111 .map_err(write_error)?;
112 connection
113 .pragma_update(None, "synchronous", "NORMAL")
114 .map_err(write_error)
115 }
116}
117
118fn poisoned() -> PortError {
119 PortError::Unavailable("quality telemetry connection lock is poisoned".to_string())
120}
121
122fn write_error(error: rusqlite::Error) -> PortError {
123 PortError::Unavailable(format!("quality telemetry write failed: {error}"))
124}