1use 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#[derive(Clone)]
122pub struct SqliteMetricsBackend {
123 conn: Arc<Mutex<Connection>>,
124 path: PathBuf,
125}
126
127impl SqliteMetricsBackend {
128 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 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#[derive(Clone)]
303pub struct SqliteEventsBackend {
304 conn: Arc<Mutex<Connection>>,
305 path: PathBuf,
306}
307
308impl SqliteEventsBackend {
309 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 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 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}