Skip to main content

kmp_adapter_embedded/adapter/telemetry/
sqlite_quality_telemetry_writer.rs

1use 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/// Multi-process SQLite writer for the bounded local quality journal.
16#[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}