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