1use std::sync::Arc;
4
5use spectra_core::{
6 install_config, set_sink, LoggingKind, SchemaRegistry, SharedEventBackend,
7 SharedMetricsBackend, SpectraConfig, SpectraRouter, SpectraSink,
8};
9
10use crate::persist_config::PersistConfig;
11use crate::persist_sink::{PersistHandle, StoragePersistSink};
12
13#[cfg(feature = "telemetry-console")]
14use std::path::Path;
15#[cfg(feature = "telemetry-console")]
16use crate::async_writer::OffThreadSpectraSink;
17#[cfg(feature = "telemetry-console")]
18use spectra_core::NdjsonFileSink;
19
20pub struct Spectra {
45 router: Arc<SpectraRouter>,
46 persist: Option<PersistHandle>,
47 embedded: bool,
49}
50
51impl Spectra {
52 pub fn router(&self) -> Arc<SpectraRouter> {
54 Arc::clone(&self.router)
55 }
56
57 pub fn is_embedded(&self) -> bool {
62 self.embedded
63 }
64
65 pub async fn flush_persist(&self) -> spectra_core::Result<()> {
70 match &self.persist {
71 Some(handle) => handle.flush().await,
72 None => Ok(()),
73 }
74 }
75
76 pub fn builder() -> SpectraBuilder {
95 SpectraBuilder::new()
96 }
97}
98
99pub struct SpectraBuilder {
134 metrics: Option<SharedMetricsBackend>,
135 events: Option<SharedEventBackend>,
136 config: Option<SpectraConfig>,
137 transport_sink: Option<Arc<dyn SpectraSink>>,
138 embedded: bool,
139 persist: bool,
140 persist_config: PersistConfig,
141}
142
143impl Default for SpectraBuilder {
144 fn default() -> Self {
145 Self::new()
146 }
147}
148
149impl SpectraBuilder {
150 pub fn new() -> Self {
152 Self {
153 metrics: None,
154 events: None,
155 config: None,
156 transport_sink: None,
157 embedded: false,
158 persist: true,
159 persist_config: PersistConfig::default(),
160 }
161 }
162
163 pub fn metrics_backend(mut self, backend: SharedMetricsBackend) -> Self {
165 self.metrics = Some(backend);
166 self
167 }
168
169 pub fn events_backend(mut self, backend: SharedEventBackend) -> Self {
171 self.events = Some(backend);
172 self
173 }
174
175 pub fn config(mut self, config: SpectraConfig) -> Self {
177 self.config = Some(config);
178 self
179 }
180
181 pub fn persist(mut self, config: PersistConfig) -> Self {
186 self.persist_config = config;
187 self
188 }
189
190 pub fn sink(mut self, sink: Arc<dyn SpectraSink>) -> Self {
203 self.transport_sink = Some(sink);
204 self
205 }
206
207 pub fn embedded(mut self) -> Self {
213 self.embedded = true;
214 self
215 }
216
217 pub fn persist_disabled(mut self) -> Self {
223 self.persist = false;
224 self
225 }
226
227 #[cfg(feature = "telemetry-console")]
231 pub fn telemetry_ndjson(mut self, dir: impl AsRef<Path>) -> spectra_core::Result<Self> {
232 let dir = dir.as_ref();
233 let ndjson = NdjsonFileSink::new(dir.join("metrics.ndjson"), dir.join("events.ndjson"))?;
234 let sink = Arc::new(OffThreadSpectraSink::new(ndjson));
235 self.transport_sink = Some(sink);
236 Ok(self)
237 }
238
239 pub fn build(self) -> spectra_core::Result<Spectra> {
241 let metrics = self
242 .metrics
243 .ok_or_else(|| spectra_core::Error::Internal("metrics_backend is required".into()))?;
244 let events = self
245 .events
246 .ok_or_else(|| spectra_core::Error::Internal("events_backend is required".into()))?;
247
248 let router = build_router(metrics, events);
249 let router = Arc::new(router);
250 SpectraRouter::set_global(Arc::clone(&router));
251
252 let config = self.config.unwrap_or_else(SpectraConfig::from_env);
253 install_config(config);
254
255 let persist_config = self.persist_config;
256 let (installed, persist_handle): (Arc<dyn SpectraSink>, Option<PersistHandle>) =
257 match (self.persist, self.transport_sink) {
258 (true, Some(inner)) => {
259 let sink = StoragePersistSink::with_config(
260 Arc::clone(&router),
261 Some(inner),
262 persist_config,
263 );
264 let handle = sink.handle();
265 (Arc::new(sink), Some(handle))
266 }
267 (true, None) => {
268 let sink =
269 StoragePersistSink::new_with_config(Arc::clone(&router), persist_config);
270 let handle = sink.handle();
271 (Arc::new(sink), Some(handle))
272 }
273 (false, Some(inner)) => (inner, None),
274 (false, None) => {
275 return Err(spectra_core::Error::Internal(
276 "SpectraBuilder: persist is disabled but no sink was configured; \
277 call .sink(...) or enable persist"
278 .into(),
279 ));
280 }
281 };
282 set_sink(installed);
283
284 Ok(Spectra {
285 router,
286 persist: persist_handle,
287 embedded: self.embedded,
288 })
289 }
290}
291
292fn build_router(
293 metrics: SharedMetricsBackend,
294 events: SharedEventBackend,
295) -> SpectraRouter {
296 let router = SpectraRouter::with_defaults(Arc::clone(&metrics), Arc::clone(&events));
297 for name in SchemaRegistry::global().list_schemas() {
298 let Some(meta) = SchemaRegistry::global().get_schema(name) else {
299 continue;
300 };
301 match meta.logging_kind {
302 LoggingKind::Event => {
303 router.register_event_backend(name, Arc::clone(&events));
304 }
305 LoggingKind::Metric => {
306 router.register_metrics_backend(name, Arc::clone(&metrics));
307 }
308 }
309 }
310 router
311}
312
313#[cfg(test)]
314mod tests {
315 use super::*;
316 use spectra_backend_mem::{MemEventsBackend, MemMetricsBackend};
317 use spectra_core::{try_record_counter_now, NoOpSink, RecordingSink, SpectraConfig};
318
319 static RUNTIME_TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
320
321 fn mem_backends() -> (SharedMetricsBackend, SharedEventBackend) {
322 (
323 Arc::new(MemMetricsBackend::new()),
324 Arc::new(MemEventsBackend::new()),
325 )
326 }
327
328 async fn with_isolated_runtime<F, Fut>(f: F)
329 where
330 F: FnOnce() -> Fut,
331 Fut: std::future::Future<Output = ()>,
332 {
333 let _g = RUNTIME_TEST_LOCK.lock().await;
334 spectra_core::install_config(SpectraConfig {
335 enabled: false,
336 ..Default::default()
337 });
338 f().await;
339 spectra_core::set_sink(Arc::new(NoOpSink));
340 }
341
342 #[tokio::test]
343 async fn embedded_flag_retained_on_handle() {
344 with_isolated_runtime(|| async {
345 let (metrics, events) = mem_backends();
346 let spectra = Spectra::builder()
347 .metrics_backend(metrics)
348 .events_backend(events)
349 .embedded()
350 .build()
351 .expect("build");
352 assert!(spectra.is_embedded());
353 })
354 .await;
355 }
356
357 #[tokio::test]
358 async fn builder_installs_persist_sink() {
359 with_isolated_runtime(|| async {
360 let (metrics, events) = mem_backends();
361
362 let spectra = Spectra::builder()
363 .metrics_backend(Arc::clone(&metrics))
364 .events_backend(Arc::clone(&events))
365 .embedded()
366 .build()
367 .expect("build");
368
369 try_record_counter_now("test_counter", &[], 1);
370 spectra.flush_persist().await.expect("flush");
371
372 let points = spectra
373 .router()
374 .query_metrics(spectra_core::MetricsQueryRange {
375 metric_name: "test_counter".into(),
376 start: chrono::Utc::now() - chrono::Duration::seconds(5),
377 end: chrono::Utc::now() + chrono::Duration::seconds(1),
378 label_matchers: vec![],
379 })
380 .await
381 .expect("query");
382 assert_eq!(points.len(), 1);
383 })
384 .await;
385 }
386
387 #[tokio::test]
388 async fn transport_and_persist_both_receive_emits() {
389 with_isolated_runtime(|| async {
390 let (metrics, events) = mem_backends();
391 let transport = Arc::new(RecordingSink::new());
392
393 let spectra = Spectra::builder()
394 .metrics_backend(Arc::clone(&metrics))
395 .events_backend(Arc::clone(&events))
396 .sink(Arc::clone(&transport) as Arc<dyn SpectraSink>)
397 .embedded()
398 .build()
399 .expect("build");
400
401 try_record_counter_now("dual_path_counter", &[], 1);
402 spectra.flush_persist().await.expect("flush");
403
404 assert_eq!(transport.counters().len(), 1);
405 let points = spectra
406 .router()
407 .query_metrics(spectra_core::MetricsQueryRange {
408 metric_name: "dual_path_counter".into(),
409 start: chrono::Utc::now() - chrono::Duration::seconds(5),
410 end: chrono::Utc::now() + chrono::Duration::seconds(1),
411 label_matchers: vec![],
412 })
413 .await
414 .expect("query");
415 assert_eq!(points.len(), 1);
416 })
417 .await;
418 }
419
420 #[tokio::test]
421 async fn transport_only_skips_storage() {
422 with_isolated_runtime(|| async {
423 let (metrics, events) = mem_backends();
424 let transport = Arc::new(RecordingSink::new());
425
426 let spectra = Spectra::builder()
427 .metrics_backend(Arc::clone(&metrics))
428 .events_backend(Arc::clone(&events))
429 .sink(Arc::clone(&transport) as Arc<dyn SpectraSink>)
430 .persist_disabled()
431 .build()
432 .expect("build");
433
434 try_record_counter_now("transport_only_counter", &[], 1);
435 spectra.flush_persist().await.expect("flush noop");
436 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
437
438 assert_eq!(transport.counters().len(), 1);
439 let points = spectra
440 .router()
441 .query_metrics(spectra_core::MetricsQueryRange {
442 metric_name: "transport_only_counter".into(),
443 start: chrono::Utc::now() - chrono::Duration::seconds(5),
444 end: chrono::Utc::now() + chrono::Duration::seconds(1),
445 label_matchers: vec![],
446 })
447 .await
448 .expect("query");
449 assert!(points.is_empty());
450 })
451 .await;
452 }
453
454 #[tokio::test]
455 async fn persist_disabled_without_sink_errors() {
456 let _g = RUNTIME_TEST_LOCK.lock().await;
457 let (metrics, events) = mem_backends();
458 let result = Spectra::builder()
459 .metrics_backend(metrics)
460 .events_backend(events)
461 .persist_disabled()
462 .build();
463 assert!(result.is_err());
464 assert!(
465 result
466 .err()
467 .expect("err")
468 .to_string()
469 .contains("no sink")
470 );
471 }
472
473 #[tokio::test]
474 async fn persist_config_batch_and_flush() {
475 with_isolated_runtime(|| async {
476 let (metrics, events) = mem_backends();
477
478 let spectra = Spectra::builder()
479 .metrics_backend(Arc::clone(&metrics))
480 .events_backend(Arc::clone(&events))
481 .persist(PersistConfig {
482 batch_max: 16,
483 batch_enabled: true,
484 ..PersistConfig::default()
485 })
486 .embedded()
487 .build()
488 .expect("build");
489
490 for i in 0..8 {
491 try_record_counter_now(&format!("cfg_batch_{i}"), &[], 1);
492 }
493 spectra.flush_persist().await.expect("flush");
494
495 for i in 0..8 {
496 let points = spectra
497 .router()
498 .query_metrics(spectra_core::MetricsQueryRange {
499 metric_name: format!("cfg_batch_{i}"),
500 start: chrono::Utc::now() - chrono::Duration::seconds(5),
501 end: chrono::Utc::now() + chrono::Duration::seconds(1),
502 label_matchers: vec![],
503 })
504 .await
505 .expect("query");
506 assert_eq!(points.len(), 1, "counter {i}");
507 }
508 })
509 .await;
510 }
511}