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 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}