Skip to main content

spectra_backend_mem/
lib.rs

1//! In-memory [`MetricsStorageBackend`] and [`EventStorageBackend`] for tests and quick start.
2//!
3//! Inject through `SpectraBuilder::metrics_backend` / `events_backend`, or use the
4//! re-exports from the `spectra` crate (`MemMetricsBackend`, `MemEventsBackend`).
5//!
6//! - Data is process-local and lost on exit; not suitable for production durability.
7//! - `query_aggregate` supports `Count` measure only; other measures return empty series.
8//! - Uses `RwLock`; contended writes may block the async runtime thread briefly.
9
10use std::collections::HashMap;
11use std::sync::RwLock;
12
13use async_trait::async_trait;
14use chrono::{DateTime, Utc};
15use serde_json::{json, Value};
16use spectra_core::{
17    EventAggregateResult, EventMeasure, EventRow, EventStorageBackend,
18    EventsAggregateFilter, EventsQueryFilter, LabelMatcher, MetricPoint, MetricPointDto,
19    MetricsQueryRange, MetricsStorageBackend, Result, StorageEngineType, TimeSeriesDto,
20};
21
22/// Non-durable in-memory metrics storage.
23///
24/// Counter and gauge writes are appended as points and can be queried immediately. This backend
25/// is the default `spectra` backend and is useful for local development and tests.
26///
27/// # Examples
28///
29/// Facade wiring through `Spectra::builder()` (requires the `spectra` crate):
30///
31/// ```ignore
32/// use std::sync::Arc;
33/// use spectra::{MemEventsBackend, MemMetricsBackend, Spectra};
34///
35/// let spectra = Spectra::builder()
36///     .metrics_backend(Arc::new(MemMetricsBackend::new()))
37///     .events_backend(Arc::new(MemEventsBackend::new()))
38///     .embedded()
39///     .build()?;
40/// ```
41///
42/// Direct backend usage:
43///
44/// ```no_run
45/// use chrono::{Duration, Utc};
46/// use serde_json::json;
47/// use spectra_backend_mem::MemMetricsBackend;
48/// use spectra_core::{MetricsQueryRange, MetricsStorageBackend};
49///
50/// # async fn example() -> spectra_core::Result<()> {
51/// let backend = MemMetricsBackend::new();
52/// let now = Utc::now();
53/// backend.record_counter("cache_hits", &json!({"region": "us"}), 1, now).await?;
54///
55/// let points = backend.query_range(MetricsQueryRange {
56///     metric_name: "cache_hits".into(),
57///     start: now - Duration::seconds(1),
58///     end: now + Duration::seconds(1),
59///     label_matchers: vec![],
60/// }).await?;
61/// assert_eq!(points.len(), 1);
62/// # Ok(())
63/// # }
64/// ```
65#[derive(Default)]
66pub struct MemMetricsBackend {
67    points: RwLock<Vec<StoredMetric>>,
68}
69
70#[derive(Clone)]
71struct StoredMetric {
72    name: String,
73    value: f64,
74    labels: Value,
75    ts: DateTime<Utc>,
76}
77
78impl MemMetricsBackend {
79    /// Create an empty in-memory metrics backend.
80    ///
81    /// Writes are process-local and discarded when the process exits. Inject the result into
82    /// `Spectra::builder().metrics_backend(...)` for a full runtime, or call storage trait
83    /// methods directly in tests.
84    ///
85    /// # Examples
86    ///
87    /// ```
88    /// use spectra_backend_mem::MemMetricsBackend;
89    /// use spectra_core::{MetricsStorageBackend, StorageEngineType};
90    ///
91    /// let backend = MemMetricsBackend::new();
92    /// assert_eq!(backend.engine_type(), StorageEngineType::Mem);
93    /// ```
94    pub fn new() -> Self {
95        Self::default()
96    }
97}
98
99#[async_trait]
100impl MetricsStorageBackend for MemMetricsBackend {
101    fn engine_type(&self) -> StorageEngineType {
102        StorageEngineType::Mem
103    }
104
105    async fn record_counter(
106        &self,
107        name: &str,
108        labels: &Value,
109        delta: i64,
110        ts: DateTime<Utc>,
111    ) -> Result<()> {
112        let mut guard = self.points.write().expect("mem metrics lock");
113        guard.push(StoredMetric {
114            name: name.to_string(),
115            value: delta as f64,
116            labels: labels.clone(),
117            ts,
118        });
119        Ok(())
120    }
121
122    async fn record_gauge(
123        &self,
124        name: &str,
125        labels: &Value,
126        value: f64,
127        ts: DateTime<Utc>,
128    ) -> Result<()> {
129        let mut guard = self.points.write().expect("mem metrics lock");
130        guard.push(StoredMetric {
131            name: name.to_string(),
132            value,
133            labels: labels.clone(),
134            ts,
135        });
136        Ok(())
137    }
138
139    async fn query_range(&self, query: MetricsQueryRange) -> Result<Vec<MetricPoint>> {
140        let guard = self.points.read().expect("mem metrics lock");
141        Ok(guard
142            .iter()
143            .filter(|p| p.name == query.metric_name)
144            .filter(|p| p.ts >= query.start && p.ts <= query.end)
145            .filter(|p| labels_match(&p.labels, &query.label_matchers))
146            .map(|p| MetricPoint {
147                ts: p.ts,
148                value: p.value,
149                labels: p.labels.clone(),
150            })
151            .collect())
152    }
153}
154
155/// Non-durable in-memory structured-event storage keyed by logical table name.
156///
157/// Rows are available immediately after append, which makes this backend useful for tests and
158/// local development.
159///
160/// # Examples
161///
162/// Facade wiring through `Spectra::builder()` (requires the `spectra` crate):
163///
164/// ```ignore
165/// use std::sync::Arc;
166/// use spectra::{MemEventsBackend, MemMetricsBackend, Spectra};
167///
168/// let spectra = Spectra::builder()
169///     .metrics_backend(Arc::new(MemMetricsBackend::new()))
170///     .events_backend(Arc::new(MemEventsBackend::new()))
171///     .embedded()
172///     .build()?;
173/// ```
174///
175/// Direct backend usage:
176///
177/// ```no_run
178/// use chrono::Utc;
179/// use serde_json::json;
180/// use spectra_backend_mem::MemEventsBackend;
181/// use spectra_core::{EventStorageBackend, EventsQueryFilter};
182///
183/// # async fn example() -> spectra_core::Result<()> {
184/// let backend = MemEventsBackend::new();
185/// backend.append_row(
186///     "request_log",
187///     &json!({"message": "handled"}),
188///     Utc::now(),
189///     None,
190/// ).await?;
191///
192/// let rows = backend.query_rows(EventsQueryFilter {
193///     table: "request_log".into(),
194///     ..Default::default()
195/// }).await?;
196/// assert_eq!(rows.len(), 1);
197/// # Ok(())
198/// # }
199/// ```
200#[derive(Default)]
201pub struct MemEventsBackend {
202    rows: RwLock<HashMap<String, Vec<StoredEvent>>>,
203}
204
205#[derive(Clone)]
206struct StoredEvent {
207    fields: Value,
208    ts: DateTime<Utc>,
209}
210
211impl MemEventsBackend {
212    /// Create an empty in-memory events backend.
213    ///
214    /// Rows are keyed by logical table name and discarded when the process exits. Inject the
215    /// result into `Spectra::builder().events_backend(...)` for a full runtime.
216    ///
217    /// # Examples
218    ///
219    /// ```
220    /// use spectra_backend_mem::MemEventsBackend;
221    /// use spectra_core::{EventStorageBackend, StorageEngineType};
222    ///
223    /// let backend = MemEventsBackend::new();
224    /// assert_eq!(backend.engine_type(), StorageEngineType::Mem);
225    /// ```
226    pub fn new() -> Self {
227        Self::default()
228    }
229}
230
231#[async_trait]
232impl EventStorageBackend for MemEventsBackend {
233    fn engine_type(&self) -> StorageEngineType {
234        StorageEngineType::Mem
235    }
236
237    async fn append_row(
238        &self,
239        table: &str,
240        fields: &Value,
241        ts: DateTime<Utc>,
242        _correlation_id: Option<&str>,
243    ) -> Result<()> {
244        let mut guard = self.rows.write().expect("mem events lock");
245        guard
246            .entry(table.to_string())
247            .or_default()
248            .push(StoredEvent {
249                fields: fields.clone(),
250                ts,
251            });
252        Ok(())
253    }
254
255    async fn query_rows(&self, filter: EventsQueryFilter) -> Result<Vec<EventRow>> {
256        let guard = self.rows.read().expect("mem events lock");
257        let Some(rows) = guard.get(&filter.table) else {
258            return Ok(Vec::new());
259        };
260        let out: Vec<EventRow> = rows
261            .iter()
262            .filter(|r| filter.start.map(|s| r.ts >= s).unwrap_or(true))
263            .filter(|r| filter.end.map(|e| r.ts <= e).unwrap_or(true))
264            .map(|r| EventRow {
265                ts: r.ts,
266                fields: r.fields.clone(),
267            })
268            .collect();
269        Ok(spectra_core::finalize_event_rows(out, &filter))
270    }
271
272    async fn query_aggregate(&self, filter: EventsAggregateFilter) -> Result<EventAggregateResult> {
273        let rows = self
274            .query_rows(EventsQueryFilter {
275                table: filter.table.clone(),
276                start: Some(filter.start),
277                end: Some(filter.end),
278                partition: filter.partition.clone(),
279                limit: None,
280                offset: None,
281                sort_field: None,
282                sort_desc: false,
283                filter: filter.filter.clone(),
284            })
285            .await?;
286        let count = rows.len() as u64;
287        match filter.measure {
288            EventMeasure::Count => Ok(EventAggregateResult::TimeSeries {
289                series: vec![TimeSeriesDto {
290                    labels: json!({}),
291                    points: vec![MetricPointDto {
292                        ts: filter.end,
293                        value: count as f64,
294                    }],
295                }],
296                headline: vec![],
297            }),
298            _ => Ok(EventAggregateResult::TimeSeries {
299                series: vec![],
300                headline: vec![],
301            }),
302        }
303    }
304}
305
306fn labels_match(labels: &Value, matchers: &[LabelMatcher]) -> bool {
307    matchers.iter().all(|m| {
308        labels
309            .get(&m.key)
310            .and_then(|v| v.as_str())
311            .map(|v| v == m.value)
312            .unwrap_or(false)
313    })
314}
315
316#[cfg(test)]
317mod tests {
318    use super::*;
319    use chrono::Duration;
320    use serde_json::json;
321
322    #[tokio::test]
323    async fn metrics_roundtrip() {
324        let backend = MemMetricsBackend::new();
325        let ts = Utc::now();
326        backend
327            .record_counter("hits", &json!({"k": "v"}), 3, ts)
328            .await
329            .expect("write");
330        let points = backend
331            .query_range(MetricsQueryRange {
332                metric_name: "hits".into(),
333                start: ts - Duration::seconds(1),
334                end: ts + Duration::seconds(1),
335                label_matchers: vec![],
336            })
337            .await
338            .expect("query");
339        assert_eq!(points.len(), 1);
340        assert!((points[0].value - 3.0).abs() < f64::EPSILON);
341    }
342
343    #[tokio::test]
344    async fn events_roundtrip() {
345        let backend = MemEventsBackend::new();
346        let ts = Utc::now();
347        backend
348            .append_row("req_log", &json!({"msg": "ok"}), ts, None)
349            .await
350            .expect("write");
351        let rows = backend
352            .query_rows(EventsQueryFilter {
353                table: "req_log".into(),
354                start: Some(ts - Duration::seconds(1)),
355                end: Some(ts + Duration::seconds(1)),
356                ..Default::default()
357            })
358            .await
359            .expect("query");
360        assert_eq!(rows.len(), 1);
361    }
362}