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::{
12    enforce_retention, insert_observation, migrate_legacy_quality_telemetry,
13    open_quality_connection,
14};
15
16const DEFAULT_DURABLE_EVERY_BATCHES: u64 = 16;
17
18/// Multi-process SQLite writer for the bounded local quality journal.
19#[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}