Skip to main content

kmp_adapter_embedded/adapter/telemetry/
storage.rs

1use std::path::{Path, PathBuf};
2use std::time::Duration;
3
4use kmp_domain::PortError;
5use kmp_observability::QualityTelemetryObservation;
6use rusqlite::{Connection, Transaction, config::DbConfig, params};
7
8use super::QualityTelemetryRetention;
9use crate::adapter::serdes::encode;
10
11const BUSY_TIMEOUT: Duration = Duration::from_secs(10);
12
13pub fn quality_telemetry_path(data_dir: &Path) -> PathBuf {
14    data_dir.join("telemetry").join("quality.sqlite3")
15}
16
17pub(super) fn open_quality_connection(data_dir: &Path) -> Result<Connection, PortError> {
18    let path = quality_telemetry_path(data_dir);
19    let parent = path.parent().expect("quality telemetry path has a parent");
20    std::fs::create_dir_all(parent).map_err(|error| {
21        PortError::Unavailable(format!(
22            "quality telemetry could not create `{}`: {error}",
23            parent.display()
24        ))
25    })?;
26    let connection = Connection::open(&path).map_err(|error| {
27        PortError::Unavailable(format!(
28            "quality telemetry could not open `{}`: {error}",
29            path.display()
30        ))
31    })?;
32    harden_connection(&connection, &path)?;
33    connection
34        .busy_timeout(BUSY_TIMEOUT)
35        .map_err(|error| sqlite_error(&path, "configure busy timeout", error))?;
36    enter_wal(&connection, &path)?;
37    connection
38        .pragma_update(None, "synchronous", "NORMAL")
39        .map_err(|error| sqlite_error(&path, "configure durability", error))?;
40    initialize_schema(&connection, &path)?;
41    Ok(connection)
42}
43
44pub(super) fn insert_observation(
45    transaction: &Transaction<'_>,
46    observation: &QualityTelemetryObservation,
47) -> Result<(), PortError> {
48    let observed_at = i64::try_from(observation.observed_at_millis()).map_err(|_| {
49        PortError::InvalidState("quality observation timestamp exceeds SQLite range".to_string())
50    })?;
51    let payload = encode("quality observation", observation)?;
52    transaction
53        .execute(
54            "INSERT INTO quality_observations (observed_at_millis, payload) VALUES (?1, ?2)",
55            params![observed_at, payload],
56        )
57        .map_err(|error| {
58            PortError::Unavailable(format!(
59                "quality telemetry could not persist an observation: {error}"
60            ))
61        })?;
62    Ok(())
63}
64
65pub(super) fn enforce_retention(
66    transaction: &Transaction<'_>,
67    retention: QualityTelemetryRetention,
68) -> Result<(), PortError> {
69    let total: i64 = transaction
70        .query_row("SELECT COUNT(*) FROM quality_observations", [], |row| {
71            row.get(0)
72        })
73        .map_err(|error| {
74            PortError::Unavailable(format!("quality telemetry count failed: {error}"))
75        })?;
76    let total = u64::try_from(total)
77        .map_err(|_| PortError::InvalidState("quality telemetry count is negative".to_string()))?;
78    let excess = retention.excess(total);
79    if excess > 0 {
80        transaction
81            .execute(
82                "DELETE FROM quality_observations WHERE id IN (\
83                 SELECT id FROM quality_observations \
84                 ORDER BY observed_at_millis ASC, id ASC LIMIT ?1)",
85                params![i64::try_from(excess).unwrap_or(i64::MAX)],
86            )
87            .map_err(|error| {
88                PortError::Unavailable(format!(
89                    "quality telemetry retention cleanup failed: {error}"
90                ))
91            })?;
92    }
93    Ok(())
94}
95
96fn initialize_schema(connection: &Connection, path: &Path) -> Result<(), PortError> {
97    let integrity: String = connection
98        .query_row("PRAGMA quick_check(1)", [], |row| row.get(0))
99        .map_err(|error| sqlite_error(path, "run quick_check", error))?;
100    if integrity != "ok" {
101        return Err(PortError::InvalidState(format!(
102            "quality telemetry SQLite journal `{}` failed quick_check: {integrity}",
103            path.display()
104        )));
105    }
106    connection
107        .execute_batch(
108            "CREATE TABLE IF NOT EXISTS quality_observations (\
109               id INTEGER PRIMARY KEY AUTOINCREMENT,\
110               observed_at_millis INTEGER NOT NULL,\
111               payload BLOB NOT NULL\
112             );\
113             CREATE INDEX IF NOT EXISTS quality_observations_by_time \
114               ON quality_observations (observed_at_millis, id);\
115             CREATE TABLE IF NOT EXISTS quality_metadata (\
116               key TEXT PRIMARY KEY,\
117               value TEXT NOT NULL\
118             ) WITHOUT ROWID;",
119        )
120        .map_err(|error| sqlite_error(path, "initialize schema", error))
121}
122
123fn harden_connection(connection: &Connection, path: &Path) -> Result<(), PortError> {
124    for (config, enabled, name) in [
125        (DbConfig::SQLITE_DBCONFIG_DEFENSIVE, true, "defensive mode"),
126        (
127            DbConfig::SQLITE_DBCONFIG_TRUSTED_SCHEMA,
128            false,
129            "trusted schema",
130        ),
131        (DbConfig::SQLITE_DBCONFIG_ENABLE_TRIGGER, false, "triggers"),
132        (DbConfig::SQLITE_DBCONFIG_ENABLE_VIEW, false, "views"),
133    ] {
134        connection
135            .set_db_config(config, enabled)
136            .map_err(|error| sqlite_error(path, name, error))?;
137    }
138    connection
139        .execute_batch("PRAGMA cell_size_check=ON; PRAGMA mmap_size=0;")
140        .map_err(|error| sqlite_error(path, "configure safe page access", error))
141}
142
143fn enter_wal(connection: &Connection, path: &Path) -> Result<(), PortError> {
144    let deadline = std::time::Instant::now() + BUSY_TIMEOUT;
145    let mut backoff = Duration::from_millis(2);
146    loop {
147        match connection.pragma_update(None, "journal_mode", "WAL") {
148            Ok(()) => return Ok(()),
149            Err(error) if is_busy(&error) && std::time::Instant::now() < deadline => {
150                std::thread::sleep(backoff);
151                backoff = (backoff * 2).min(Duration::from_millis(64));
152            }
153            Err(error) => return Err(sqlite_error(path, "enter WAL mode", error)),
154        }
155    }
156}
157
158fn is_busy(error: &rusqlite::Error) -> bool {
159    matches!(
160        error.sqlite_error_code(),
161        Some(rusqlite::ErrorCode::DatabaseBusy | rusqlite::ErrorCode::DatabaseLocked)
162    )
163}
164
165fn sqlite_error(path: &Path, action: &str, error: rusqlite::Error) -> PortError {
166    PortError::Unavailable(format!(
167        "quality telemetry could not {action} at `{}`: {error}",
168        path.display()
169    ))
170}