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 redb::{Database, ReadableDatabase, ReadableTable, TableDefinition};
7use rusqlite::{
8 Connection, OptionalExtension, Transaction, TransactionBehavior, config::DbConfig, params,
9};
10
11use super::QualityTelemetryRetention;
12use crate::adapter::serdes::{decode, encode};
13
14const LEGACY_OBSERVATIONS: TableDefinition<(u64, u64), &[u8]> =
15 TableDefinition::new("quality_observations");
16const LEGACY_IMPORT_KEY: &str = "legacy-quality-redb-v1";
17const BUSY_TIMEOUT: Duration = Duration::from_secs(10);
18
19pub fn quality_telemetry_path(data_dir: &Path) -> PathBuf {
20 data_dir.join("telemetry").join("quality.sqlite3")
21}
22
23pub fn legacy_quality_telemetry_path(data_dir: &Path) -> PathBuf {
24 data_dir.join("telemetry").join("quality.redb")
25}
26
27pub(super) fn open_quality_connection(data_dir: &Path) -> Result<Connection, PortError> {
28 let path = quality_telemetry_path(data_dir);
29 let parent = path.parent().expect("quality telemetry path has a parent");
30 std::fs::create_dir_all(parent).map_err(|error| {
31 PortError::Unavailable(format!(
32 "quality telemetry could not create `{}`: {error}",
33 parent.display()
34 ))
35 })?;
36 let connection = Connection::open(&path).map_err(|error| {
37 PortError::Unavailable(format!(
38 "quality telemetry could not open `{}`: {error}",
39 path.display()
40 ))
41 })?;
42 harden_connection(&connection, &path)?;
43 connection
44 .busy_timeout(BUSY_TIMEOUT)
45 .map_err(|error| sqlite_error(&path, "configure busy timeout", error))?;
46 enter_wal(&connection, &path)?;
47 connection
48 .pragma_update(None, "synchronous", "NORMAL")
49 .map_err(|error| sqlite_error(&path, "configure durability", error))?;
50 initialize_schema(&connection, &path)?;
51 Ok(connection)
52}
53
54pub(super) fn migrate_legacy_quality_telemetry(
55 connection: &mut Connection,
56 data_dir: &Path,
57 retention: QualityTelemetryRetention,
58) -> Result<u64, PortError> {
59 let legacy_path = legacy_quality_telemetry_path(data_dir);
60 if !legacy_path.is_file() || legacy_import_complete(connection)? {
61 return Ok(0);
62 }
63
64 let observations = read_legacy_observations(&legacy_path)?;
65 let transaction = connection
66 .transaction_with_behavior(TransactionBehavior::Immediate)
67 .map_err(|error| {
68 PortError::Unavailable(format!(
69 "quality telemetry legacy import could not start: {error}"
70 ))
71 })?;
72 if legacy_import_complete_in(&transaction)? {
73 return Ok(0);
74 }
75 for observation in &observations {
76 insert_observation(&transaction, observation)?;
77 }
78 enforce_retention(&transaction, retention)?;
79 transaction
80 .execute(
81 "INSERT INTO quality_metadata (key, value) VALUES (?1, ?2)",
82 params![LEGACY_IMPORT_KEY, observations.len().to_string()],
83 )
84 .map_err(|error| {
85 PortError::Unavailable(format!(
86 "quality telemetry could not record its legacy import: {error}"
87 ))
88 })?;
89 transaction.commit().map_err(|error| {
90 PortError::Unavailable(format!(
91 "quality telemetry legacy import could not commit: {error}"
92 ))
93 })?;
94 Ok(observations.len() as u64)
95}
96
97pub(super) fn insert_observation(
98 transaction: &Transaction<'_>,
99 observation: &QualityTelemetryObservation,
100) -> Result<(), PortError> {
101 let observed_at = i64::try_from(observation.observed_at_millis()).map_err(|_| {
102 PortError::InvalidState("quality observation timestamp exceeds SQLite range".to_string())
103 })?;
104 let payload = encode("quality observation", observation)?;
105 transaction
106 .execute(
107 "INSERT INTO quality_observations (observed_at_millis, payload) VALUES (?1, ?2)",
108 params![observed_at, payload],
109 )
110 .map_err(|error| {
111 PortError::Unavailable(format!(
112 "quality telemetry could not persist an observation: {error}"
113 ))
114 })?;
115 Ok(())
116}
117
118pub(super) fn enforce_retention(
119 transaction: &Transaction<'_>,
120 retention: QualityTelemetryRetention,
121) -> Result<(), PortError> {
122 let total: i64 = transaction
123 .query_row("SELECT COUNT(*) FROM quality_observations", [], |row| {
124 row.get(0)
125 })
126 .map_err(|error| {
127 PortError::Unavailable(format!("quality telemetry count failed: {error}"))
128 })?;
129 let total = u64::try_from(total)
130 .map_err(|_| PortError::InvalidState("quality telemetry count is negative".to_string()))?;
131 let excess = retention.excess(total);
132 if excess > 0 {
133 transaction
134 .execute(
135 "DELETE FROM quality_observations WHERE id IN (\
136 SELECT id FROM quality_observations \
137 ORDER BY observed_at_millis ASC, id ASC LIMIT ?1)",
138 params![i64::try_from(excess).unwrap_or(i64::MAX)],
139 )
140 .map_err(|error| {
141 PortError::Unavailable(format!(
142 "quality telemetry retention cleanup failed: {error}"
143 ))
144 })?;
145 }
146 Ok(())
147}
148
149fn initialize_schema(connection: &Connection, path: &Path) -> Result<(), PortError> {
150 let integrity: String = connection
151 .query_row("PRAGMA quick_check(1)", [], |row| row.get(0))
152 .map_err(|error| sqlite_error(path, "run quick_check", error))?;
153 if integrity != "ok" {
154 return Err(PortError::InvalidState(format!(
155 "quality telemetry SQLite journal `{}` failed quick_check: {integrity}",
156 path.display()
157 )));
158 }
159 connection
160 .execute_batch(
161 "CREATE TABLE IF NOT EXISTS quality_observations (\
162 id INTEGER PRIMARY KEY AUTOINCREMENT,\
163 observed_at_millis INTEGER NOT NULL,\
164 payload BLOB NOT NULL\
165 );\
166 CREATE INDEX IF NOT EXISTS quality_observations_by_time \
167 ON quality_observations (observed_at_millis, id);\
168 CREATE TABLE IF NOT EXISTS quality_metadata (\
169 key TEXT PRIMARY KEY,\
170 value TEXT NOT NULL\
171 ) WITHOUT ROWID;",
172 )
173 .map_err(|error| sqlite_error(path, "initialize schema", error))
174}
175
176fn harden_connection(connection: &Connection, path: &Path) -> Result<(), PortError> {
177 for (config, enabled, name) in [
178 (DbConfig::SQLITE_DBCONFIG_DEFENSIVE, true, "defensive mode"),
179 (
180 DbConfig::SQLITE_DBCONFIG_TRUSTED_SCHEMA,
181 false,
182 "trusted schema",
183 ),
184 (DbConfig::SQLITE_DBCONFIG_ENABLE_TRIGGER, false, "triggers"),
185 (DbConfig::SQLITE_DBCONFIG_ENABLE_VIEW, false, "views"),
186 ] {
187 connection
188 .set_db_config(config, enabled)
189 .map_err(|error| sqlite_error(path, name, error))?;
190 }
191 connection
192 .execute_batch("PRAGMA cell_size_check=ON; PRAGMA mmap_size=0;")
193 .map_err(|error| sqlite_error(path, "configure safe page access", error))
194}
195
196fn enter_wal(connection: &Connection, path: &Path) -> Result<(), PortError> {
197 let deadline = std::time::Instant::now() + BUSY_TIMEOUT;
198 let mut backoff = Duration::from_millis(2);
199 loop {
200 match connection.pragma_update(None, "journal_mode", "WAL") {
201 Ok(()) => return Ok(()),
202 Err(error) if is_busy(&error) && std::time::Instant::now() < deadline => {
203 std::thread::sleep(backoff);
204 backoff = (backoff * 2).min(Duration::from_millis(64));
205 }
206 Err(error) => return Err(sqlite_error(path, "enter WAL mode", error)),
207 }
208 }
209}
210
211fn is_busy(error: &rusqlite::Error) -> bool {
212 matches!(
213 error.sqlite_error_code(),
214 Some(rusqlite::ErrorCode::DatabaseBusy | rusqlite::ErrorCode::DatabaseLocked)
215 )
216}
217
218fn legacy_import_complete(connection: &Connection) -> Result<bool, PortError> {
219 connection
220 .query_row(
221 "SELECT 1 FROM quality_metadata WHERE key = ?1",
222 [LEGACY_IMPORT_KEY],
223 |_| Ok(()),
224 )
225 .optional()
226 .map(|value| value.is_some())
227 .map_err(|error| {
228 PortError::Unavailable(format!(
229 "quality telemetry could not inspect legacy import state: {error}"
230 ))
231 })
232}
233
234fn legacy_import_complete_in(transaction: &Transaction<'_>) -> Result<bool, PortError> {
235 transaction
236 .query_row(
237 "SELECT 1 FROM quality_metadata WHERE key = ?1",
238 [LEGACY_IMPORT_KEY],
239 |_| Ok(()),
240 )
241 .optional()
242 .map(|value| value.is_some())
243 .map_err(|error| {
244 PortError::Unavailable(format!(
245 "quality telemetry could not inspect concurrent legacy import state: {error}"
246 ))
247 })
248}
249
250fn read_legacy_observations(path: &Path) -> Result<Vec<QualityTelemetryObservation>, PortError> {
251 let database = open_legacy_with_retry(path)?;
252 let transaction = database.begin_read().map_err(|error| {
253 PortError::Unavailable(format!(
254 "legacy quality telemetry migration could not start reading: {error}"
255 ))
256 })?;
257 let table = transaction
258 .open_table(LEGACY_OBSERVATIONS)
259 .map_err(|error| {
260 PortError::Unavailable(format!(
261 "legacy quality telemetry migration could not open observations: {error}"
262 ))
263 })?;
264 let mut observations = Vec::new();
265 for row in table.iter().map_err(|error| {
266 PortError::Unavailable(format!(
267 "legacy quality telemetry migration could not scan observations: {error}"
268 ))
269 })? {
270 let (_, value) = row.map_err(|error| {
271 PortError::Unavailable(format!(
272 "legacy quality telemetry migration could not read an observation: {error}"
273 ))
274 })?;
275 observations.push(decode("legacy quality observation", value.value())?);
276 }
277 Ok(observations)
278}
279
280fn open_legacy_with_retry(path: &Path) -> Result<Database, PortError> {
281 let deadline = std::time::Instant::now() + BUSY_TIMEOUT;
282 let mut backoff = Duration::from_millis(2);
283 loop {
284 match Database::open(path) {
285 Ok(database) => return Ok(database),
286 Err(error)
287 if (error.to_string().contains("Cannot acquire lock")
288 || error
289 .to_string()
290 .to_ascii_lowercase()
291 .contains("already open"))
292 && std::time::Instant::now() < deadline =>
293 {
294 std::thread::sleep(backoff);
295 backoff = (backoff * 2).min(Duration::from_millis(64));
296 }
297 Err(error) => {
298 return Err(PortError::Unavailable(format!(
299 "legacy quality telemetry could not open `{}` for one-time migration: \
300 {error}",
301 path.display()
302 )));
303 }
304 }
305 }
306}
307
308fn sqlite_error(path: &Path, action: &str, error: rusqlite::Error) -> PortError {
309 PortError::Unavailable(format!(
310 "quality telemetry could not {action} at `{}`: {error}",
311 path.display()
312 ))
313}