kmp_adapter_embedded/adapter/telemetry/
storage.rs1use 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}