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}