Skip to main content

spectra_runtime/
builder.rs

1//! Build and install a Spectra runtime with injected storage backends.
2
3use 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
20/// Handle to an installed process-wide Spectra runtime.
21///
22/// Keep this value to access the configured [`SpectraRouter`]. Building installs the emit
23/// configuration, router, and sink globally; applications should build once during startup.
24///
25/// # Examples
26///
27/// ```no_run
28/// use std::sync::Arc;
29/// use spectra_backend_mem::{MemEventsBackend, MemMetricsBackend};
30/// use spectra_runtime::Spectra;
31///
32/// # fn example() -> spectra_core::Result<()> {
33/// let spectra = Spectra::builder()
34///     .metrics_backend(Arc::new(MemMetricsBackend::new()))
35///     .events_backend(Arc::new(MemEventsBackend::new()))
36///     .embedded()
37///     .build()?;
38///
39/// let router = spectra.router();
40/// # let _ = router;
41/// # Ok(())
42/// # }
43/// ```
44pub struct Spectra {
45    router: Arc<SpectraRouter>,
46    persist: Option<PersistHandle>,
47    /// Topology marker from [`SpectraBuilder::embedded`] (does not select backends).
48    embedded: bool,
49}
50
51impl Spectra {
52    /// Returns the installed global router (metrics/events query and backend resolution).
53    pub fn router(&self) -> Arc<SpectraRouter> {
54        Arc::clone(&self.router)
55    }
56
57    /// Whether this runtime was marked as in-process embedded topology.
58    ///
59    /// Set by [`SpectraBuilder::embedded`]. This is a topology marker only — it does not
60    /// select, validate, or change storage backends.
61    pub fn is_embedded(&self) -> bool {
62        self.embedded
63    }
64
65    /// Wait until every persist job enqueued before this call has been written to storage.
66    ///
67    /// Use after `try_record_*_now` / generated helpers when a script must exit only after
68    /// durable writes complete. No-op when persist was disabled at build time.
69    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    /// Start building a process-scoped Spectra runtime with explicit storage injection.
77    ///
78    /// # Examples
79    ///
80    /// ```no_run
81    /// # use std::sync::Arc;
82    /// # fn demo() -> spectra_core::Result<()> {
83    /// use spectra_backend_mem::{MemEventsBackend, MemMetricsBackend};
84    /// use spectra_runtime::Spectra;
85    ///
86    /// let _spectra = Spectra::builder()
87    ///     .metrics_backend(Arc::new(MemMetricsBackend::new()))
88    ///     .events_backend(Arc::new(MemEventsBackend::new()))
89    ///     .embedded()
90    ///     .build()?;
91    /// # Ok(())
92    /// # }
93    /// ```
94    pub fn builder() -> SpectraBuilder {
95        SpectraBuilder::new()
96    }
97}
98
99/// Configures and installs [`Spectra`] with explicit storage adapter injection.
100///
101/// Both metrics and events backends are required. Storage persistence is enabled by default;
102/// an optional [`SpectraSink`] can receive the same emits before they are queued for persistence.
103///
104/// | Mode | Calls |
105/// |------|-------|
106/// | Direct persist | backends + `.build()` |
107/// | Dual-path | `.sink(transport).build()` |
108/// | **Publisher** (distributed) | `.sink(transport).persist_disabled().build()` |
109///
110/// Publisher processes publish through the sink; consumers subscribe on the host bus and write
111/// storage. See the `spectra` crate **Getting started → Mode 2**.
112///
113/// # Examples
114///
115/// ```no_run
116/// use std::sync::Arc;
117/// use spectra_backend_mem::{MemEventsBackend, MemMetricsBackend};
118/// use spectra_core::{RecordingSink, SpectraSink};
119/// use spectra_runtime::Spectra;
120///
121/// # fn example() -> spectra_core::Result<()> {
122/// let transport = Arc::new(RecordingSink::new());
123/// let spectra = Spectra::builder()
124///     .metrics_backend(Arc::new(MemMetricsBackend::new()))
125///     .events_backend(Arc::new(MemEventsBackend::new()))
126///     .sink(Arc::clone(&transport) as Arc<dyn SpectraSink>)
127///     .embedded()
128///     .build()?;
129/// # let _ = spectra;
130/// # Ok(())
131/// # }
132/// ```
133pub 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    /// New builder with async storage persist enabled and no backends configured yet.
151    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    /// Register the metrics storage backend (required before [`Self::build`]).
164    pub fn metrics_backend(mut self, backend: SharedMetricsBackend) -> Self {
165        self.metrics = Some(backend);
166        self
167    }
168
169    /// Register the events storage backend (required before [`Self::build`]).
170    pub fn events_backend(mut self, backend: SharedEventBackend) -> Self {
171        self.events = Some(backend);
172        self
173    }
174
175    /// Override emit policy and feature flags (defaults to [`SpectraConfig::from_env`]).
176    pub fn config(mut self, config: SpectraConfig) -> Self {
177        self.config = Some(config);
178        self
179    }
180
181    /// Configure the L2 persist queue and batch settings (ignored when persist is disabled).
182    ///
183    /// Defaults: `queue_max=8192`, `batch_max=32`, `batch_wait=5ms`, `batch_enabled=true`.
184    /// Raise `batch_max` for high-throughput DW ingest.
185    pub fn persist(mut self, config: PersistConfig) -> Self {
186        self.persist_config = config;
187        self
188    }
189
190    /// Optional transport or telemetry sink (invoked before async storage persist when enabled).
191    ///
192    /// Use this for the **publisher** side of distributed ingest: your [`SpectraSink`]
193    /// publishes emits onto a bus (for example Photon). Combine with
194    /// [`Self::persist_disabled`] when the publisher must not write storage.
195    ///
196    /// Without `persist_disabled`, the runtime installs a persist wrapper that calls your
197    /// sink first, then queues async storage writes (dual-path).
198    ///
199    /// Implement [`SpectraSink`](spectra_core::SpectraSink) in your application or use
200    /// [`RecordingSink`](spectra_core::RecordingSink) in tests. See the `spectra` crate
201    /// **Getting started → Mode 2**.
202    pub fn sink(mut self, sink: Arc<dyn SpectraSink>) -> Self {
203        self.transport_sink = Some(sink);
204        self
205    }
206
207    /// Mark in-process embedded topology (in-host storage backends).
208    ///
209    /// This sets a topology flag retained on the returned [`Spectra`] handle
210    /// ([`Spectra::is_embedded`]). It does **not** select or validate backends — you still
211    /// inject storage via [`Self::metrics_backend`] / [`Self::events_backend`].
212    pub fn embedded(mut self) -> Self {
213        self.embedded = true;
214        self
215    }
216
217    /// Disable async storage persist (only the transport sink receives emits).
218    ///
219    /// **Publisher / distributed mode.** Pair with [`Self::sink`]: writers publish through the
220    /// sink; separate consumer processes subscribe and persist. Requires a sink — building
221    /// with persist disabled and no sink returns an error.
222    pub fn persist_disabled(mut self) -> Self {
223        self.persist = false;
224        self
225    }
226
227    /// Attach off-thread NDJSON + optional console mirror telemetry (`telemetry-console` feature).
228    ///
229    /// Writes `{dir}/metrics.ndjson` and `{dir}/events.ndjson`.
230    #[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    /// Install global config, router, and sink; returns a handle to the running runtime.
240    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}