use std::sync::Arc;
use spectra_core::{
install_config, set_sink, LoggingKind, SchemaRegistry, SharedEventBackend,
SharedMetricsBackend, SpectraConfig, SpectraRouter, SpectraSink,
};
use crate::persist_config::PersistConfig;
use crate::persist_sink::{PersistHandle, StoragePersistSink};
#[cfg(feature = "telemetry-console")]
use std::path::Path;
#[cfg(feature = "telemetry-console")]
use crate::async_writer::OffThreadSpectraSink;
#[cfg(feature = "telemetry-console")]
use spectra_core::NdjsonFileSink;
pub struct Spectra {
router: Arc<SpectraRouter>,
persist: Option<PersistHandle>,
embedded: bool,
}
impl Spectra {
pub fn router(&self) -> Arc<SpectraRouter> {
Arc::clone(&self.router)
}
pub fn is_embedded(&self) -> bool {
self.embedded
}
pub async fn flush_persist(&self) -> spectra_core::Result<()> {
match &self.persist {
Some(handle) => handle.flush().await,
None => Ok(()),
}
}
pub fn builder() -> SpectraBuilder {
SpectraBuilder::new()
}
}
pub struct SpectraBuilder {
metrics: Option<SharedMetricsBackend>,
events: Option<SharedEventBackend>,
config: Option<SpectraConfig>,
transport_sink: Option<Arc<dyn SpectraSink>>,
embedded: bool,
persist: bool,
persist_config: PersistConfig,
}
impl Default for SpectraBuilder {
fn default() -> Self {
Self::new()
}
}
impl SpectraBuilder {
pub fn new() -> Self {
Self {
metrics: None,
events: None,
config: None,
transport_sink: None,
embedded: false,
persist: true,
persist_config: PersistConfig::default(),
}
}
pub fn metrics_backend(mut self, backend: SharedMetricsBackend) -> Self {
self.metrics = Some(backend);
self
}
pub fn events_backend(mut self, backend: SharedEventBackend) -> Self {
self.events = Some(backend);
self
}
pub fn config(mut self, config: SpectraConfig) -> Self {
self.config = Some(config);
self
}
pub fn persist(mut self, config: PersistConfig) -> Self {
self.persist_config = config;
self
}
pub fn sink(mut self, sink: Arc<dyn SpectraSink>) -> Self {
self.transport_sink = Some(sink);
self
}
pub fn embedded(mut self) -> Self {
self.embedded = true;
self
}
pub fn persist_disabled(mut self) -> Self {
self.persist = false;
self
}
#[cfg(feature = "telemetry-console")]
pub fn telemetry_ndjson(mut self, dir: impl AsRef<Path>) -> spectra_core::Result<Self> {
let dir = dir.as_ref();
let ndjson = NdjsonFileSink::new(dir.join("metrics.ndjson"), dir.join("events.ndjson"))?;
let sink = Arc::new(OffThreadSpectraSink::new(ndjson));
self.transport_sink = Some(sink);
Ok(self)
}
pub fn build(self) -> spectra_core::Result<Spectra> {
let metrics = self
.metrics
.ok_or_else(|| spectra_core::Error::Internal("metrics_backend is required".into()))?;
let events = self
.events
.ok_or_else(|| spectra_core::Error::Internal("events_backend is required".into()))?;
let router = build_router(metrics, events);
let router = Arc::new(router);
SpectraRouter::set_global(Arc::clone(&router));
let config = self.config.unwrap_or_else(SpectraConfig::from_env);
install_config(config);
let persist_config = self.persist_config;
let (installed, persist_handle): (Arc<dyn SpectraSink>, Option<PersistHandle>) =
match (self.persist, self.transport_sink) {
(true, Some(inner)) => {
let sink = StoragePersistSink::with_config(
Arc::clone(&router),
Some(inner),
persist_config,
);
let handle = sink.handle();
(Arc::new(sink), Some(handle))
}
(true, None) => {
let sink =
StoragePersistSink::new_with_config(Arc::clone(&router), persist_config);
let handle = sink.handle();
(Arc::new(sink), Some(handle))
}
(false, Some(inner)) => (inner, None),
(false, None) => {
return Err(spectra_core::Error::Internal(
"SpectraBuilder: persist is disabled but no sink was configured; \
call .sink(...) or enable persist"
.into(),
));
}
};
set_sink(installed);
Ok(Spectra {
router,
persist: persist_handle,
embedded: self.embedded,
})
}
}
fn build_router(
metrics: SharedMetricsBackend,
events: SharedEventBackend,
) -> SpectraRouter {
let router = SpectraRouter::with_defaults(Arc::clone(&metrics), Arc::clone(&events));
for name in SchemaRegistry::global().list_schemas() {
let Some(meta) = SchemaRegistry::global().get_schema(name) else {
continue;
};
match meta.logging_kind {
LoggingKind::Event => {
router.register_event_backend(name, Arc::clone(&events));
}
LoggingKind::Metric => {
router.register_metrics_backend(name, Arc::clone(&metrics));
}
}
}
router
}
#[cfg(test)]
mod tests {
use super::*;
use spectra_backend_mem::{MemEventsBackend, MemMetricsBackend};
use spectra_core::{try_record_counter_now, NoOpSink, RecordingSink, SpectraConfig};
static RUNTIME_TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
fn mem_backends() -> (SharedMetricsBackend, SharedEventBackend) {
(
Arc::new(MemMetricsBackend::new()),
Arc::new(MemEventsBackend::new()),
)
}
async fn with_isolated_runtime<F, Fut>(f: F)
where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = ()>,
{
let _g = RUNTIME_TEST_LOCK.lock().await;
spectra_core::install_config(SpectraConfig {
enabled: false,
..Default::default()
});
f().await;
spectra_core::set_sink(Arc::new(NoOpSink));
}
#[tokio::test]
async fn embedded_flag_retained_on_handle() {
with_isolated_runtime(|| async {
let (metrics, events) = mem_backends();
let spectra = Spectra::builder()
.metrics_backend(metrics)
.events_backend(events)
.embedded()
.build()
.expect("build");
assert!(spectra.is_embedded());
})
.await;
}
#[tokio::test]
async fn builder_installs_persist_sink() {
with_isolated_runtime(|| async {
let (metrics, events) = mem_backends();
let spectra = Spectra::builder()
.metrics_backend(Arc::clone(&metrics))
.events_backend(Arc::clone(&events))
.embedded()
.build()
.expect("build");
try_record_counter_now("test_counter", &[], 1);
spectra.flush_persist().await.expect("flush");
let points = spectra
.router()
.query_metrics(spectra_core::MetricsQueryRange {
metric_name: "test_counter".into(),
start: chrono::Utc::now() - chrono::Duration::seconds(5),
end: chrono::Utc::now() + chrono::Duration::seconds(1),
label_matchers: vec![],
})
.await
.expect("query");
assert_eq!(points.len(), 1);
})
.await;
}
#[tokio::test]
async fn transport_and_persist_both_receive_emits() {
with_isolated_runtime(|| async {
let (metrics, events) = mem_backends();
let transport = Arc::new(RecordingSink::new());
let spectra = Spectra::builder()
.metrics_backend(Arc::clone(&metrics))
.events_backend(Arc::clone(&events))
.sink(Arc::clone(&transport) as Arc<dyn SpectraSink>)
.embedded()
.build()
.expect("build");
try_record_counter_now("dual_path_counter", &[], 1);
spectra.flush_persist().await.expect("flush");
assert_eq!(transport.counters().len(), 1);
let points = spectra
.router()
.query_metrics(spectra_core::MetricsQueryRange {
metric_name: "dual_path_counter".into(),
start: chrono::Utc::now() - chrono::Duration::seconds(5),
end: chrono::Utc::now() + chrono::Duration::seconds(1),
label_matchers: vec![],
})
.await
.expect("query");
assert_eq!(points.len(), 1);
})
.await;
}
#[tokio::test]
async fn transport_only_skips_storage() {
with_isolated_runtime(|| async {
let (metrics, events) = mem_backends();
let transport = Arc::new(RecordingSink::new());
let spectra = Spectra::builder()
.metrics_backend(Arc::clone(&metrics))
.events_backend(Arc::clone(&events))
.sink(Arc::clone(&transport) as Arc<dyn SpectraSink>)
.persist_disabled()
.build()
.expect("build");
try_record_counter_now("transport_only_counter", &[], 1);
spectra.flush_persist().await.expect("flush noop");
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
assert_eq!(transport.counters().len(), 1);
let points = spectra
.router()
.query_metrics(spectra_core::MetricsQueryRange {
metric_name: "transport_only_counter".into(),
start: chrono::Utc::now() - chrono::Duration::seconds(5),
end: chrono::Utc::now() + chrono::Duration::seconds(1),
label_matchers: vec![],
})
.await
.expect("query");
assert!(points.is_empty());
})
.await;
}
#[tokio::test]
async fn persist_disabled_without_sink_errors() {
let _g = RUNTIME_TEST_LOCK.lock().await;
let (metrics, events) = mem_backends();
let result = Spectra::builder()
.metrics_backend(metrics)
.events_backend(events)
.persist_disabled()
.build();
assert!(result.is_err());
assert!(
result
.err()
.expect("err")
.to_string()
.contains("no sink")
);
}
#[tokio::test]
async fn persist_config_batch_and_flush() {
with_isolated_runtime(|| async {
let (metrics, events) = mem_backends();
let spectra = Spectra::builder()
.metrics_backend(Arc::clone(&metrics))
.events_backend(Arc::clone(&events))
.persist(PersistConfig {
batch_max: 16,
batch_enabled: true,
..PersistConfig::default()
})
.embedded()
.build()
.expect("build");
for i in 0..8 {
try_record_counter_now(&format!("cfg_batch_{i}"), &[], 1);
}
spectra.flush_persist().await.expect("flush");
for i in 0..8 {
let points = spectra
.router()
.query_metrics(spectra_core::MetricsQueryRange {
metric_name: format!("cfg_batch_{i}"),
start: chrono::Utc::now() - chrono::Duration::seconds(5),
end: chrono::Utc::now() + chrono::Duration::seconds(1),
label_matchers: vec![],
})
.await
.expect("query");
assert_eq!(points.len(), 1, "counter {i}");
}
})
.await;
}
}