Skip to main content

spectra_runtime/
persist_sink.rs

1//! Generic sink that persists emits to [`SpectraRouter`] storage backends.
2
3use std::collections::HashMap;
4use std::sync::Arc;
5use std::time::Instant;
6
7use chrono::{DateTime, Utc};
8use serde_json::{Map, Value};
9use spectra_core::{
10    current_emit_ts, record_persist_queue_drop, record_storage_batch_write_events,
11    record_storage_batch_write_metrics, record_storage_write_events, record_storage_write_metrics,
12    EventWriteRow, MetricWriteRow, SpectraRouter, SpectraSink,
13};
14use tokio::sync::{mpsc, oneshot};
15
16use crate::persist_config::{PersistConfig, PersistOverflow};
17
18enum PersistJob {
19    Counter {
20        name: String,
21        labels: Value,
22        delta: i64,
23        ts: DateTime<Utc>,
24    },
25    Gauge {
26        name: String,
27        labels: Value,
28        value: f64,
29        ts: DateTime<Utc>,
30    },
31    Event {
32        table: String,
33        fields: Value,
34        ts: DateTime<Utc>,
35    },
36    /// Barrier: completes after all prior jobs in the queue have been processed.
37    Flush(oneshot::Sender<()>),
38}
39
40/// Handle to wait until the persist queue has drained prior emits.
41#[derive(Clone)]
42pub struct PersistHandle {
43    tx: mpsc::Sender<PersistJob>,
44}
45
46impl PersistHandle {
47    /// Wait until every job enqueued before this call has been flushed to storage.
48    pub async fn flush(&self) -> spectra_core::Result<()> {
49        let (ack_tx, ack_rx) = oneshot::channel();
50        self.tx
51            .send(PersistJob::Flush(ack_tx))
52            .await
53            .map_err(|_| {
54                spectra_core::Error::Internal("persist queue closed during flush".into())
55            })?;
56        ack_rx
57            .await
58            .map_err(|_| spectra_core::Error::Internal("persist flush ack dropped".into()))
59    }
60}
61
62/// [`SpectraSink`] that queues telemetry for asynchronous storage through a [`SpectraRouter`].
63///
64/// Queue capacity and batching are configured with [`PersistConfig`] (via
65/// [`crate::SpectraBuilder::persist`]), not environment variables.
66///
67/// Default overflow policy is [`PersistOverflow::Drop`] (non-blocking; drops counted via
68/// `persist_queue_drops`). Use [`PersistOverflow::Block`] for backpressure.
69/// Call [`PersistHandle::flush`] (or [`crate::Spectra::flush_persist`]) to wait for durability.
70pub struct StoragePersistSink {
71    inner: Option<Arc<dyn SpectraSink>>,
72    tx: mpsc::Sender<PersistJob>,
73    handle: PersistHandle,
74    overflow: PersistOverflow,
75}
76
77impl StoragePersistSink {
78    /// Persist-only sink with default [`PersistConfig`].
79    pub fn new(router: Arc<SpectraRouter>) -> Self {
80        Self::with_config(router, None, PersistConfig::default())
81    }
82
83    /// Persist-only sink with explicit config.
84    pub fn new_with_config(router: Arc<SpectraRouter>, config: PersistConfig) -> Self {
85        Self::with_config(router, None, config)
86    }
87
88    /// Invoke `inner` on the hot path, then enqueue the same emit for async storage persist.
89    pub fn with_inner(router: Arc<SpectraRouter>, inner: Option<Arc<dyn SpectraSink>>) -> Self {
90        Self::with_config(router, inner, PersistConfig::default())
91    }
92
93    /// Like [`Self::with_inner`] with explicit [`PersistConfig`].
94    pub fn with_config(
95        router: Arc<SpectraRouter>,
96        inner: Option<Arc<dyn SpectraSink>>,
97        config: PersistConfig,
98    ) -> Self {
99        let config = config.normalized();
100        let overflow = config.overflow;
101        let (tx, mut rx) = mpsc::channel(config.queue_max);
102        let handle = PersistHandle { tx: tx.clone() };
103        let router_worker = Arc::clone(&router);
104        let batch_max = config.batch_max;
105        let batch_wait = config.batch_wait;
106        let batch_enabled = config.batch_enabled;
107
108        tokio::spawn(async move {
109            while let Some(first) = rx.recv().await {
110                if let PersistJob::Flush(ack) = first {
111                    let _ = ack.send(());
112                    continue;
113                }
114
115                let mut batch = vec![first];
116                let mut pending_flush: Option<oneshot::Sender<()>> = None;
117
118                while batch.len() < batch_max {
119                    match rx.try_recv() {
120                        Ok(PersistJob::Flush(ack)) => {
121                            pending_flush = Some(ack);
122                            break;
123                        }
124                        Ok(job) => batch.push(job),
125                        Err(mpsc::error::TryRecvError::Empty) => {
126                            if batch.len() == 1 && batch_enabled {
127                                tokio::time::sleep(batch_wait).await;
128                                match rx.try_recv() {
129                                    Ok(PersistJob::Flush(ack)) => {
130                                        pending_flush = Some(ack);
131                                    }
132                                    Ok(job) => {
133                                        batch.push(job);
134                                        continue;
135                                    }
136                                    Err(_) => {}
137                                }
138                            }
139                            break;
140                        }
141                        Err(mpsc::error::TryRecvError::Disconnected) => break,
142                    }
143                }
144
145                if batch_enabled {
146                    if let Err(e) = flush_batch(&router_worker, batch).await {
147                        log::warn!("[spectra:persist] batch flush: {e}");
148                    }
149                } else {
150                    for job in batch {
151                        if let Err(e) = run_job(&router_worker, job).await {
152                            log::warn!("[spectra:persist] {e}");
153                        }
154                    }
155                }
156
157                if let Some(ack) = pending_flush {
158                    let _ = ack.send(());
159                }
160            }
161        });
162
163        Self {
164            inner,
165            tx,
166            handle,
167            overflow,
168        }
169    }
170
171    /// Handle for [`PersistHandle::flush`].
172    pub fn handle(&self) -> PersistHandle {
173        self.handle.clone()
174    }
175}
176
177impl SpectraSink for StoragePersistSink {
178    fn record_counter(&self, name: &str, labels: &[(&str, &str)], delta: i64) {
179        if let Some(inner) = &self.inner {
180            inner.record_counter(name, labels, delta);
181        }
182        enqueue(
183            &self.tx,
184            PersistJob::Counter {
185                name: name.to_string(),
186                labels: labels_to_value(labels),
187                delta,
188                ts: emit_ts(),
189            },
190            self.overflow,
191        );
192    }
193
194    fn record_gauge(&self, name: &str, labels: &[(&str, &str)], value: f64) {
195        if let Some(inner) = &self.inner {
196            inner.record_gauge(name, labels, value);
197        }
198        enqueue(
199            &self.tx,
200            PersistJob::Gauge {
201                name: name.to_string(),
202                labels: labels_to_value(labels),
203                value,
204                ts: emit_ts(),
205            },
206            self.overflow,
207        );
208    }
209
210    fn log_event(&self, table: &str, fields: &Value) {
211        if let Some(inner) = &self.inner {
212            inner.log_event(table, fields);
213        }
214        enqueue(
215            &self.tx,
216            PersistJob::Event {
217                table: table.to_string(),
218                fields: fields.clone(),
219                ts: emit_ts(),
220            },
221            self.overflow,
222        );
223    }
224}
225
226fn emit_ts() -> DateTime<Utc> {
227    current_emit_ts()
228}
229
230fn labels_to_value(labels: &[(&str, &str)]) -> Value {
231    let mut map = Map::new();
232    for (k, v) in labels {
233        map.insert((*k).to_string(), Value::String((*v).to_string()));
234    }
235    Value::Object(map)
236}
237
238fn enqueue(tx: &mpsc::Sender<PersistJob>, job: PersistJob, overflow: PersistOverflow) {
239    match overflow {
240        PersistOverflow::Drop => {
241            if tx.try_send(job).is_err() {
242                record_persist_queue_drop();
243                log::warn!("[spectra:persist] queue full; dropping job");
244            }
245        }
246        PersistOverflow::Block => {
247            let send_result = if tokio::runtime::Handle::try_current().is_ok() {
248                tokio::task::block_in_place(|| {
249                    let handle = tokio::runtime::Handle::current();
250                    handle.block_on(tx.send(job))
251                })
252            } else {
253                tx.blocking_send(job)
254            };
255            if send_result.is_err() {
256                record_persist_queue_drop();
257                log::warn!("[spectra:persist] queue closed; dropping job");
258            }
259        }
260    }
261}
262
263async fn run_job(router: &SpectraRouter, job: PersistJob) -> spectra_core::Result<()> {
264    match job {
265        PersistJob::Flush(_) => Ok(()),
266        PersistJob::Counter {
267            name,
268            labels,
269            delta,
270            ts,
271        } => {
272            let started = Instant::now();
273            let backend = router.resolve_metrics(&name);
274            let result = backend.record_counter(&name, &labels, delta, ts).await;
275            if result.is_ok() {
276                record_storage_write_metrics(started.elapsed());
277            }
278            result
279        }
280        PersistJob::Gauge {
281            name,
282            labels,
283            value,
284            ts,
285        } => {
286            let started = Instant::now();
287            let backend = router.resolve_metrics(&name);
288            let result = backend.record_gauge(&name, &labels, value, ts).await;
289            if result.is_ok() {
290                record_storage_write_metrics(started.elapsed());
291            }
292            result
293        }
294        PersistJob::Event { table, fields, ts } => {
295            let started = Instant::now();
296            let backend = router.resolve_event(&table);
297            let result = backend.append_row(&table, &fields, ts, None).await;
298            if result.is_ok() {
299                record_storage_write_events(started.elapsed());
300            }
301            result
302        }
303    }
304}
305
306async fn flush_batch(router: &SpectraRouter, batch: Vec<PersistJob>) -> spectra_core::Result<()> {
307    fn arc_key<T: ?Sized>(arc: &Arc<T>) -> usize {
308        Arc::as_ptr(arc) as *const () as usize
309    }
310
311    let mut metrics_buckets: HashMap<usize, (spectra_core::SharedMetricsBackend, Vec<MetricWriteRow>)> =
312        HashMap::new();
313    let mut event_buckets: HashMap<usize, (spectra_core::SharedEventBackend, Vec<EventWriteRow>)> =
314        HashMap::new();
315
316    for job in batch {
317        match job {
318            PersistJob::Flush(_) => {}
319            PersistJob::Counter {
320                name,
321                labels,
322                delta,
323                ts,
324            } => {
325                let backend = router.resolve_metrics(&name);
326                let key = arc_key(&backend);
327                metrics_buckets
328                    .entry(key)
329                    .or_insert_with(|| (Arc::clone(&backend), Vec::new()))
330                    .1
331                    .push(MetricWriteRow {
332                        name,
333                        kind: "counter",
334                        value: Value::from(delta),
335                        labels,
336                        ts,
337                        correlation_id: None,
338                    });
339            }
340            PersistJob::Gauge {
341                name,
342                labels,
343                value,
344                ts,
345            } => {
346                let backend = router.resolve_metrics(&name);
347                let key = arc_key(&backend);
348                metrics_buckets
349                    .entry(key)
350                    .or_insert_with(|| (Arc::clone(&backend), Vec::new()))
351                    .1
352                    .push(MetricWriteRow {
353                        name,
354                        kind: "gauge",
355                        value: serde_json::json!(value),
356                        labels,
357                        ts,
358                        correlation_id: None,
359                    });
360            }
361            PersistJob::Event { table, fields, ts } => {
362                let backend = router.resolve_event(&table);
363                let key = arc_key(&backend);
364                event_buckets
365                    .entry(key)
366                    .or_insert_with(|| (Arc::clone(&backend), Vec::new()))
367                    .1
368                    .push(EventWriteRow {
369                        table,
370                        fields,
371                        ts,
372                        correlation_id: None,
373                    });
374            }
375        }
376    }
377
378    for (_, (backend, rows)) in metrics_buckets {
379        if !rows.is_empty() {
380            let started = Instant::now();
381            let row_count = rows.len() as u64;
382            backend.record_metrics_batch(&rows).await?;
383            record_storage_batch_write_metrics(started.elapsed(), row_count);
384        }
385    }
386    for (_, (backend, rows)) in event_buckets {
387        if !rows.is_empty() {
388            let started = Instant::now();
389            let row_count = rows.len() as u64;
390            backend.append_rows_batch(&rows).await?;
391            record_storage_batch_write_events(started.elapsed(), row_count);
392        }
393    }
394    Ok(())
395}
396
397#[cfg(test)]
398mod tests {
399    use super::*;
400    use spectra_backend_mem::{MemEventsBackend, MemMetricsBackend};
401    use spectra_core::{
402        try_record_counter_now, NoOpSink, SharedEventBackend, SharedMetricsBackend, SpectraConfig,
403    };
404
405    static PERSIST_TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
406
407    #[tokio::test]
408    async fn batch_flush_persists_multiple_counters() {
409        let _g = PERSIST_TEST_LOCK.lock().await;
410        spectra_core::install_config(SpectraConfig {
411            enabled: false,
412            ..Default::default()
413        });
414
415        let metrics: SharedMetricsBackend = Arc::new(MemMetricsBackend::new());
416        let events: SharedEventBackend = Arc::new(MemEventsBackend::new());
417        let router = Arc::new(SpectraRouter::with_defaults(
418            Arc::clone(&metrics),
419            Arc::clone(&events),
420        ));
421        let sink = StoragePersistSink::new_with_config(
422            Arc::clone(&router),
423            PersistConfig {
424                batch_max: 8,
425                batch_enabled: true,
426                ..PersistConfig::default()
427            },
428        );
429        let handle = sink.handle();
430        spectra_core::set_sink(Arc::new(sink));
431
432        for i in 0..4 {
433            try_record_counter_now(&format!("batch_counter_{i}"), &[], 1);
434        }
435        handle.flush().await.expect("flush");
436
437        for i in 0..4 {
438            let points = router
439                .query_metrics(spectra_core::MetricsQueryRange {
440                    metric_name: format!("batch_counter_{i}"),
441                    start: Utc::now() - chrono::Duration::seconds(5),
442                    end: Utc::now() + chrono::Duration::seconds(1),
443                    label_matchers: vec![],
444                })
445                .await
446                .expect("query");
447            assert_eq!(points.len(), 1, "counter {i}");
448        }
449
450        spectra_core::set_sink(Arc::new(NoOpSink));
451    }
452}