1use 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 Flush(oneshot::Sender<()>),
38}
39
40#[derive(Clone)]
42pub struct PersistHandle {
43 tx: mpsc::Sender<PersistJob>,
44}
45
46impl PersistHandle {
47 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
62pub struct StoragePersistSink {
71 inner: Option<Arc<dyn SpectraSink>>,
72 tx: mpsc::Sender<PersistJob>,
73 handle: PersistHandle,
74 overflow: PersistOverflow,
75}
76
77impl StoragePersistSink {
78 pub fn new(router: Arc<SpectraRouter>) -> Self {
80 Self::with_config(router, None, PersistConfig::default())
81 }
82
83 pub fn new_with_config(router: Arc<SpectraRouter>, config: PersistConfig) -> Self {
85 Self::with_config(router, None, config)
86 }
87
88 pub fn with_inner(router: Arc<SpectraRouter>, inner: Option<Arc<dyn SpectraSink>>) -> Self {
90 Self::with_config(router, inner, PersistConfig::default())
91 }
92
93 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 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}