Skip to main content

spectra_backend_sqlite/
lib.rs

1//! Durable embedded SQLite metrics and events storage.
2//!
3//! Enable with the `spectra` feature `sqlite` and wire through `Spectra::builder()`.
4//!
5//! - [`SqliteMetricsBackend::new`] / [`SqliteEventsBackend::new`] — open or create database files
6//! - Parent directories are created automatically; uses `spawn_blocking` for rusqlite I/O.
7//! - `query_aggregate` is not yet implemented (returns empty series).
8//! - Default event query limit is 1000 rows when `limit` is unset.
9
10use std::path::{Path, PathBuf};
11use std::sync::{Arc, Mutex};
12
13use async_trait::async_trait;
14use chrono::{DateTime, Utc};
15use rusqlite::{params, Connection};
16use serde_json::Value;
17use spectra_core::{
18    Error, EventAggregateResult, EventRow, EventStorageBackend, EventsAggregateFilter,
19    EventsQueryFilter, LabelMatcher, MetricPoint, MetricsQueryRange, MetricsStorageBackend, Result,
20    StorageEngineType,
21};
22
23const METRICS_DDL: &str = r"
24CREATE TABLE IF NOT EXISTS spectra_metrics (
25    name TEXT NOT NULL,
26    kind TEXT NOT NULL,
27    value REAL NOT NULL,
28    labels TEXT NOT NULL,
29    ts TEXT NOT NULL,
30    correlation_id TEXT
31);
32CREATE INDEX IF NOT EXISTS idx_spectra_metrics_name_ts ON spectra_metrics(name, ts);
33";
34
35const EVENTS_DDL: &str = r"
36CREATE TABLE IF NOT EXISTS spectra_events (
37    table_name TEXT NOT NULL,
38    fields TEXT NOT NULL,
39    ts TEXT NOT NULL,
40    correlation_id TEXT
41);
42CREATE INDEX IF NOT EXISTS idx_spectra_events_table_ts ON spectra_events(table_name, ts);
43";
44
45fn open_and_migrate(path: &Path, ddl: &str) -> Result<Connection> {
46    if let Some(parent) = path.parent() {
47        std::fs::create_dir_all(parent).map_err(Error::Io)?;
48    }
49    let conn = Connection::open(path).map_err(|e| Error::Internal(e.to_string()))?;
50    conn.execute_batch(ddl)
51        .map_err(|e| Error::Internal(e.to_string()))?;
52    Ok(conn)
53}
54
55fn ts_to_rfc3339(ts: DateTime<Utc>) -> String {
56    ts.to_rfc3339()
57}
58
59fn parse_ts(s: &str) -> Result<DateTime<Utc>> {
60    DateTime::parse_from_rfc3339(s)
61        .map(|dt| dt.with_timezone(&Utc))
62        .map_err(|e| Error::Internal(e.to_string()))
63}
64
65fn labels_match(labels: &Value, matchers: &[LabelMatcher]) -> bool {
66    matchers.iter().all(|m| {
67        labels
68            .get(&m.key)
69            .and_then(|v| v.as_str())
70            .map(|v| v == m.value)
71            .unwrap_or(false)
72    })
73}
74
75/// Durable SQLite metrics storage in the `spectra_metrics` table.
76///
77/// [`new`](Self::new) opens the database file and applies the required DDL. Use a separate
78/// file from the events backend when following the standard Spectra wiring.
79///
80/// # Examples
81///
82/// Facade wiring through `Spectra::builder()` (requires the `spectra` crate with the
83/// `sqlite` feature):
84///
85/// ```ignore
86/// use std::sync::Arc;
87/// use spectra::{Spectra, SqliteEventsBackend, SqliteMetricsBackend};
88///
89/// let dir = std::env::temp_dir().join("spectra-example");
90/// let spectra = Spectra::builder()
91///     .metrics_backend(Arc::new(SqliteMetricsBackend::new(dir.join("metrics.db"))?))
92///     .events_backend(Arc::new(SqliteEventsBackend::new(dir.join("events.db"))?))
93///     .embedded()
94///     .build()?;
95/// ```
96///
97/// Direct backend usage:
98///
99/// ```no_run
100/// use chrono::{Duration, Utc};
101/// use serde_json::json;
102/// use spectra_backend_sqlite::SqliteMetricsBackend;
103/// use spectra_core::{MetricsQueryRange, MetricsStorageBackend};
104///
105/// # async fn example() -> spectra_core::Result<()> {
106/// let dir = tempfile::tempdir()?;
107/// let backend = SqliteMetricsBackend::new(dir.path().join("metrics.db"))?;
108/// let now = Utc::now();
109/// backend.record_counter("cache_hits", &json!({}), 1, now).await?;
110///
111/// let points = backend.query_range(MetricsQueryRange {
112///     metric_name: "cache_hits".into(),
113///     start: now - Duration::seconds(1),
114///     end: now + Duration::seconds(1),
115///     label_matchers: vec![],
116/// }).await?;
117/// assert_eq!(points.len(), 1);
118/// # Ok(())
119/// # }
120/// ```
121#[derive(Clone)]
122pub struct SqliteMetricsBackend {
123    conn: Arc<Mutex<Connection>>,
124    path: PathBuf,
125}
126
127impl SqliteMetricsBackend {
128    /// Open or create a SQLite database at `path` and apply metrics DDL migrations.
129    ///
130    /// The constructor is synchronous. Parent directories must already exist. Use a separate
131    /// database file from the events backend when following the standard Spectra wiring.
132    ///
133    /// # Examples
134    ///
135    /// ```no_run
136    /// use spectra_backend_sqlite::SqliteMetricsBackend;
137    ///
138    /// # fn example() -> spectra_core::Result<()> {
139    /// let backend = SqliteMetricsBackend::new("/tmp/spectra-metrics.db")?;
140    /// assert!(backend.path().ends_with("spectra-metrics.db"));
141    /// # Ok(())
142    /// # }
143    /// ```
144    pub fn new(path: impl AsRef<Path>) -> Result<Self> {
145        let path = path.as_ref().to_path_buf();
146        let conn = open_and_migrate(&path, METRICS_DDL)?;
147        Ok(Self {
148            conn: Arc::new(Mutex::new(conn)),
149            path,
150        })
151    }
152
153    /// Filesystem path of the backing SQLite database.
154    pub fn path(&self) -> &Path {
155        &self.path
156    }
157}
158
159#[async_trait]
160impl MetricsStorageBackend for SqliteMetricsBackend {
161    fn engine_type(&self) -> StorageEngineType {
162        StorageEngineType::Sqlite
163    }
164
165    async fn record_counter(
166        &self,
167        name: &str,
168        labels: &Value,
169        delta: i64,
170        ts: DateTime<Utc>,
171    ) -> Result<()> {
172        let conn = Arc::clone(&self.conn);
173        let name = name.to_string();
174        let labels = labels.to_string();
175        let ts = ts_to_rfc3339(ts);
176        tokio::task::spawn_blocking(move || {
177            let conn = conn.lock().expect("sqlite metrics lock");
178            conn.execute(
179                "INSERT INTO spectra_metrics (name, kind, value, labels, ts, correlation_id) VALUES (?1, 'counter', ?2, ?3, ?4, NULL)",
180                params![name, delta as f64, labels, ts],
181            )
182            .map_err(|e| Error::Internal(e.to_string()))?;
183            Ok(())
184        })
185        .await
186        .map_err(|e| Error::Internal(e.to_string()))?
187    }
188
189    async fn record_gauge(
190        &self,
191        name: &str,
192        labels: &Value,
193        value: f64,
194        ts: DateTime<Utc>,
195    ) -> Result<()> {
196        let conn = Arc::clone(&self.conn);
197        let name = name.to_string();
198        let labels = labels.to_string();
199        let ts = ts_to_rfc3339(ts);
200        tokio::task::spawn_blocking(move || {
201            let conn = conn.lock().expect("sqlite metrics lock");
202            conn.execute(
203                "INSERT INTO spectra_metrics (name, kind, value, labels, ts, correlation_id) VALUES (?1, 'gauge', ?2, ?3, ?4, NULL)",
204                params![name, value, labels, ts],
205            )
206            .map_err(|e| Error::Internal(e.to_string()))?;
207            Ok(())
208        })
209        .await
210        .map_err(|e| Error::Internal(e.to_string()))?
211    }
212
213    async fn query_range(&self, query: MetricsQueryRange) -> Result<Vec<MetricPoint>> {
214        let conn = Arc::clone(&self.conn);
215        let name = query.metric_name.clone();
216        let start = ts_to_rfc3339(query.start);
217        let end = ts_to_rfc3339(query.end);
218        let label_matchers = query.label_matchers.clone();
219        tokio::task::spawn_blocking(move || {
220            let conn = conn.lock().expect("sqlite metrics lock");
221            let mut stmt = conn
222                .prepare(
223                    "SELECT value, labels, ts FROM spectra_metrics WHERE name = ?1 AND ts >= ?2 AND ts <= ?3 ORDER BY ts ASC",
224                )
225                .map_err(|e| Error::Internal(e.to_string()))?;
226            let rows = stmt
227                .query_map(params![name, start, end], |row| {
228                    let value: f64 = row.get(0)?;
229                    let labels: String = row.get(1)?;
230                    let ts: String = row.get(2)?;
231                    Ok((value, labels, ts))
232                })
233                .map_err(|e| Error::Internal(e.to_string()))?;
234            let mut out = Vec::new();
235            for row in rows {
236                let (value, labels, ts) = row.map_err(|e| Error::Internal(e.to_string()))?;
237                let labels: Value = serde_json::from_str(&labels)
238                    .map_err(|e| Error::Internal(e.to_string()))?;
239                if labels_match(&labels, &label_matchers) {
240                    out.push(MetricPoint {
241                        ts: parse_ts(&ts)?,
242                        value,
243                        labels,
244                    });
245                }
246            }
247            Ok(out)
248        })
249        .await
250        .map_err(|e| Error::Internal(e.to_string()))?
251    }
252}
253
254/// Durable SQLite structured-event storage in the `spectra_events` table.
255///
256/// [`new`](Self::new) opens the database file and applies the required DDL. Event rows retain
257/// their logical table name and JSON field payload.
258///
259/// # Examples
260///
261/// Facade wiring through `Spectra::builder()` (requires the `spectra` crate with the
262/// `sqlite` feature):
263///
264/// ```ignore
265/// use std::sync::Arc;
266/// use spectra::{Spectra, SqliteEventsBackend, SqliteMetricsBackend};
267///
268/// let dir = std::env::temp_dir().join("spectra-example");
269/// let spectra = Spectra::builder()
270///     .metrics_backend(Arc::new(SqliteMetricsBackend::new(dir.join("metrics.db"))?))
271///     .events_backend(Arc::new(SqliteEventsBackend::new(dir.join("events.db"))?))
272///     .embedded()
273///     .build()?;
274/// ```
275///
276/// Direct backend usage:
277///
278/// ```no_run
279/// use chrono::Utc;
280/// use serde_json::json;
281/// use spectra_backend_sqlite::SqliteEventsBackend;
282/// use spectra_core::{EventStorageBackend, EventsQueryFilter};
283///
284/// # async fn example() -> spectra_core::Result<()> {
285/// let dir = tempfile::tempdir()?;
286/// let backend = SqliteEventsBackend::new(dir.path().join("events.db"))?;
287/// backend.append_row(
288///     "request_log",
289///     &json!({"message": "handled"}),
290///     Utc::now(),
291///     None,
292/// ).await?;
293///
294/// let rows = backend.query_rows(EventsQueryFilter {
295///     table: "request_log".into(),
296///     ..Default::default()
297/// }).await?;
298/// assert_eq!(rows.len(), 1);
299/// # Ok(())
300/// # }
301/// ```
302#[derive(Clone)]
303pub struct SqliteEventsBackend {
304    conn: Arc<Mutex<Connection>>,
305    path: PathBuf,
306}
307
308impl SqliteEventsBackend {
309    /// Open or create a SQLite database at `path` and apply events DDL migrations.
310    ///
311    /// The constructor is synchronous. Parent directories must already exist. Use a separate
312    /// database file from the metrics backend when following the standard Spectra wiring.
313    ///
314    /// # Examples
315    ///
316    /// ```no_run
317    /// use spectra_backend_sqlite::SqliteEventsBackend;
318    ///
319    /// # fn example() -> spectra_core::Result<()> {
320    /// let backend = SqliteEventsBackend::new("/tmp/spectra-events.db")?;
321    /// assert!(backend.path().ends_with("spectra-events.db"));
322    /// # Ok(())
323    /// # }
324    /// ```
325    pub fn new(path: impl AsRef<Path>) -> Result<Self> {
326        let path = path.as_ref().to_path_buf();
327        let conn = open_and_migrate(&path, EVENTS_DDL)?;
328        Ok(Self {
329            conn: Arc::new(Mutex::new(conn)),
330            path,
331        })
332    }
333
334    /// Filesystem path of the backing SQLite database.
335    pub fn path(&self) -> &Path {
336        &self.path
337    }
338}
339
340#[async_trait]
341impl EventStorageBackend for SqliteEventsBackend {
342    fn engine_type(&self) -> StorageEngineType {
343        StorageEngineType::Sqlite
344    }
345
346    async fn append_row(
347        &self,
348        table: &str,
349        fields: &Value,
350        ts: DateTime<Utc>,
351        correlation_id: Option<&str>,
352    ) -> Result<()> {
353        let conn = Arc::clone(&self.conn);
354        let table = table.to_string();
355        let fields = fields.to_string();
356        let ts = ts_to_rfc3339(ts);
357        let cid = correlation_id.map(str::to_string);
358        tokio::task::spawn_blocking(move || {
359            let conn = conn.lock().expect("sqlite events lock");
360            conn.execute(
361                "INSERT INTO spectra_events (table_name, fields, ts, correlation_id) VALUES (?1, ?2, ?3, ?4)",
362                params![table, fields, ts, cid],
363            )
364            .map_err(|e| Error::Internal(e.to_string()))?;
365            Ok(())
366        })
367        .await
368        .map_err(|e| Error::Internal(e.to_string()))?
369    }
370
371    async fn query_rows(&self, filter: EventsQueryFilter) -> Result<Vec<EventRow>> {
372        let conn = Arc::clone(&self.conn);
373        let table = filter.table.clone();
374        let start = filter.start.map(ts_to_rfc3339);
375        let end = filter.end.map(ts_to_rfc3339);
376        // Fetch table/time-scoped candidates, then apply shared filter/sort/pagination
377        // so operator semantics match mem and remote-common mem store.
378        let out = tokio::task::spawn_blocking(move || {
379            let conn = conn.lock().expect("sqlite events lock");
380            let sql = "SELECT fields, ts FROM spectra_events WHERE table_name = ?1 \
381                 AND (?2 IS NULL OR ts >= ?2) AND (?3 IS NULL OR ts <= ?3)";
382            let mut stmt = conn
383                .prepare(sql)
384                .map_err(|e| Error::Internal(e.to_string()))?;
385            let rows = stmt
386                .query_map(params![table, start, end], |row| {
387                    let fields: String = row.get(0)?;
388                    let ts: String = row.get(1)?;
389                    Ok((fields, ts))
390                })
391                .map_err(|e| Error::Internal(e.to_string()))?;
392            let mut out = Vec::new();
393            for row in rows {
394                let (fields, ts) = row.map_err(|e| Error::Internal(e.to_string()))?;
395                out.push(EventRow {
396                    ts: parse_ts(&ts)?,
397                    fields: serde_json::from_str(&fields)
398                        .map_err(|e| Error::Internal(e.to_string()))?,
399                });
400            }
401            Ok::<_, Error>(out)
402        })
403        .await
404        .map_err(|e| Error::Internal(e.to_string()))??;
405        let mut filter = filter;
406        if filter.limit.is_none() {
407            filter.limit = Some(1000);
408        }
409        Ok(spectra_core::finalize_event_rows(out, &filter))
410    }
411
412    async fn query_aggregate(&self, _filter: EventsAggregateFilter) -> Result<EventAggregateResult> {
413        Ok(EventAggregateResult::TimeSeries {
414            series: vec![],
415            headline: vec![],
416        })
417    }
418}
419
420#[cfg(test)]
421mod tests {
422    use super::*;
423    use chrono::Duration;
424    use serde_json::json;
425    use tempfile::tempdir;
426
427    #[tokio::test]
428    async fn sqlite_metrics_roundtrip() {
429        let dir = tempdir().expect("tempdir");
430        let backend = SqliteMetricsBackend::new(dir.path().join("metrics.db")).expect("open");
431        let ts = Utc::now();
432        backend
433            .record_counter("hits", &json!({}), 2, ts)
434            .await
435            .expect("write");
436        let points = backend
437            .query_range(MetricsQueryRange {
438                metric_name: "hits".into(),
439                start: ts - Duration::seconds(1),
440                end: ts + Duration::seconds(1),
441                label_matchers: vec![],
442            })
443            .await
444            .expect("query");
445        assert_eq!(points.len(), 1);
446    }
447
448    #[tokio::test]
449    async fn sqlite_events_roundtrip() {
450        let dir = tempdir().expect("tempdir");
451        let backend = SqliteEventsBackend::new(dir.path().join("events.db")).expect("open");
452        let ts = Utc::now();
453        backend
454            .append_row("req", &json!({"x": 1}), ts, None)
455            .await
456            .expect("write");
457        let rows = backend
458            .query_rows(EventsQueryFilter {
459                table: "req".into(),
460                ..Default::default()
461            })
462            .await
463            .expect("query");
464        assert_eq!(rows.len(), 1);
465    }
466}