Skip to main content

spectra_core/
router.rs

1//! Resolves event tables and metric names to storage backends for reads and writes.
2
3use std::collections::HashMap;
4use std::sync::{Arc, OnceLock};
5
6use parking_lot::RwLock;
7
8use crate::error::Result;
9use crate::query::EventAggregateResult;
10use crate::storage::{
11    EventsAggregateFilter, EventsQueryFilter, MetricsQueryRange, NoOpEventBackend,
12    NoOpMetricsBackend, SharedEventBackend, SharedMetricsBackend,
13};
14
15static GLOBAL_ROUTER: OnceLock<Arc<SpectraRouter>> = OnceLock::new();
16
17/// Resolves event tables and metric names to storage backends.
18///
19/// A runtime installs default metrics and events backends, then registers schema-specific
20/// routes. Queries use a named route when present and otherwise fall back to the corresponding
21/// default backend. Most applications access this through `Spectra::router()`.
22///
23/// # Examples
24///
25/// ```no_run
26/// use chrono::{Duration, Utc};
27/// use spectra_core::{MetricsQueryRange, SpectraRouter};
28///
29/// # async fn example() -> spectra_core::Result<()> {
30/// let router = SpectraRouter::new();
31/// let now = Utc::now();
32/// let points = router.query_metrics(MetricsQueryRange {
33///     metric_name: "cache_hits".into(),
34///     start: now - Duration::minutes(5),
35///     end: now,
36///     label_matchers: vec![],
37/// }).await?;
38///
39/// // A new router uses no-op defaults, so no rows are returned.
40/// assert!(points.is_empty());
41/// # Ok(())
42/// # }
43/// ```
44pub struct SpectraRouter {
45    events: RwLock<HashMap<String, SharedEventBackend>>,
46    metrics: RwLock<HashMap<String, SharedMetricsBackend>>,
47    default_events: SharedEventBackend,
48    default_metrics: SharedMetricsBackend,
49}
50
51impl Default for SpectraRouter {
52    fn default() -> Self {
53        Self::new()
54    }
55}
56
57impl SpectraRouter {
58    /// Creates a router with no-op default backends.
59    pub fn new() -> Self {
60        Self {
61            events: RwLock::new(HashMap::new()),
62            metrics: RwLock::new(HashMap::new()),
63            default_events: Arc::new(NoOpEventBackend),
64            default_metrics: Arc::new(NoOpMetricsBackend),
65        }
66    }
67
68    /// Create a router with default backends for unregistered schema names.
69    pub fn with_defaults(
70        default_metrics: SharedMetricsBackend,
71        default_events: SharedEventBackend,
72    ) -> Self {
73        Self {
74            events: RwLock::new(HashMap::new()),
75            metrics: RwLock::new(HashMap::new()),
76            default_events,
77            default_metrics,
78        }
79    }
80
81    /// Registers a storage backend for an event table.
82    pub fn register_event_backend(&self, table: impl Into<String>, backend: SharedEventBackend) {
83        self.events.write().insert(table.into(), backend);
84    }
85
86    /// Registers a storage backend for a metric family.
87    pub fn register_metrics_backend(&self, name: impl Into<String>, backend: SharedMetricsBackend) {
88        self.metrics.write().insert(name.into(), backend);
89    }
90
91    /// Resolves the backend for an event table, falling back to the default.
92    pub fn resolve_event(&self, table: &str) -> SharedEventBackend {
93        self.events
94            .read()
95            .get(table)
96            .cloned()
97            .unwrap_or_else(|| Arc::clone(&self.default_events))
98    }
99
100    /// Resolves the backend for a metric family, falling back to the default.
101    pub fn resolve_metrics(&self, name: &str) -> SharedMetricsBackend {
102        self.metrics
103            .read()
104            .get(name)
105            .cloned()
106            .unwrap_or_else(|| Arc::clone(&self.default_metrics))
107    }
108
109    /// Queries event rows through the resolved backend.
110    ///
111    /// # Errors
112    ///
113    /// Returns [`crate::Error::Config`] when the table name, sort field, or any grid filter
114    /// field fails [`crate::validate_spectra_ident`].
115    pub async fn query_events(
116        &self,
117        filter: EventsQueryFilter,
118    ) -> Result<Vec<crate::storage::EventRow>> {
119        validate_events_query(&filter)?;
120        let backend = self.resolve_event(&filter.table);
121        backend.query_rows(filter).await
122    }
123
124    /// Queries metric points through the resolved backend.
125    ///
126    /// # Errors
127    ///
128    /// Returns [`crate::Error::Config`] when `metric_name` is not a valid Spectra identifier.
129    pub async fn query_metrics(
130        &self,
131        query: MetricsQueryRange,
132    ) -> Result<Vec<crate::storage::MetricPoint>> {
133        crate::validate_spectra_ident(&query.metric_name)?;
134        let backend = self.resolve_metrics(&query.metric_name);
135        backend.query_range(query).await
136    }
137
138    /// Queries aggregated chart data through the resolved backend.
139    pub async fn query_event_aggregate(
140        &self,
141        filter: EventsAggregateFilter,
142    ) -> Result<EventAggregateResult> {
143        crate::validate_spectra_ident(&filter.table)?;
144        if let Some(ref field) = filter.group_by_field {
145            crate::validate_spectra_ident(field)?;
146        }
147        for item in &filter.filter.items {
148            if item.field != "ts" {
149                crate::validate_spectra_ident(&item.field)?;
150            }
151        }
152        let table = filter.table.clone();
153        let backend = self.resolve_event(&table);
154        backend.query_aggregate(filter).await
155    }
156
157    /// Installs a router as the process-global instance (call once).
158    pub fn set_global(router: Arc<Self>) {
159        let _ = GLOBAL_ROUTER.set(router);
160    }
161
162    /// Returns the process-global router (panics if not installed).
163    ///
164    /// **Do not call from library code.** Prefer [`Self::try_global`], or hold an
165    /// [`Arc`] from [`Self::set_global`] / `Spectra::router()`. This panicking accessor is
166    /// reserved for hosts and examples after a successful install.
167    ///
168    /// # Panics
169    ///
170    /// Panics if [`Self::set_global`] was not called.
171    pub fn global() -> Arc<Self> {
172        // Host/example convenience after install — libraries must use `try_global`.
173        #[allow(clippy::expect_used)]
174        GLOBAL_ROUTER
175            .get()
176            .cloned()
177            .expect("SpectraRouter::set_global was not called")
178    }
179
180    /// Returns the process-global router if installed (panic-free; prefer in libraries).
181    pub fn try_global() -> Option<Arc<Self>> {
182        GLOBAL_ROUTER.get().cloned()
183    }
184}
185
186fn validate_events_query(filter: &EventsQueryFilter) -> Result<()> {
187    crate::validate_spectra_ident(&filter.table)?;
188    if let Some(ref field) = filter.sort_field {
189        if field != "ts" {
190            crate::validate_spectra_ident(field)?;
191        }
192    }
193    for item in &filter.filter.items {
194        if item.field != "ts" {
195            crate::validate_spectra_ident(&item.field)?;
196        }
197    }
198    Ok(())
199}
200
201#[cfg(test)]
202mod tests {
203    use super::*;
204    use crate::storage::{EventStorageBackend, EventsQueryFilter};
205    use async_trait::async_trait;
206    use chrono::Utc;
207    use serde_json::json;
208    use std::sync::atomic::{AtomicU32, Ordering};
209
210    struct CountingEventBackend {
211        appends: AtomicU32,
212    }
213
214    #[async_trait]
215    impl EventStorageBackend for CountingEventBackend {
216        fn engine_type(&self) -> crate::storage::StorageEngineType {
217            crate::storage::StorageEngineType::NoOp
218        }
219
220        async fn append_row(
221            &self,
222            _: &str,
223            _: &serde_json::Value,
224            _: chrono::DateTime<Utc>,
225            _: Option<&str>,
226        ) -> crate::error::Result<()> {
227            self.appends.fetch_add(1, Ordering::SeqCst);
228            Ok(())
229        }
230    }
231
232    #[test]
233    fn router_noop_default() {
234        let router = SpectraRouter::new();
235        let rt = tokio::runtime::Runtime::new().expect("runtime");
236        let rows = rt
237            .block_on(router.query_events(EventsQueryFilter {
238                table: "missing".into(),
239                ..Default::default()
240            }))
241            .expect("query");
242        assert!(rows.is_empty());
243    }
244
245    #[test]
246    fn router_resolve_event_backend() {
247        let router = SpectraRouter::new();
248        let counting = Arc::new(CountingEventBackend {
249            appends: AtomicU32::new(0),
250        });
251        let backend: SharedEventBackend = Arc::clone(&counting) as SharedEventBackend;
252        router.register_event_backend("t1", backend);
253        let resolved = router.resolve_event("t1");
254        let rt = tokio::runtime::Runtime::new().expect("runtime");
255        rt.block_on(async {
256            resolved
257                .append_row("t1", &json!({}), Utc::now(), None)
258                .await
259                .expect("append");
260        });
261        assert_eq!(counting.appends.load(Ordering::SeqCst), 1);
262    }
263
264    #[test]
265    fn query_events_rejects_bad_filter_field() {
266        use crate::error::Error;
267        use crate::{GridFilterItem, GridFilterModel, GridFilterOperator};
268
269        let router = SpectraRouter::new();
270        let rt = tokio::runtime::Runtime::new().expect("runtime");
271        let err = rt
272            .block_on(router.query_events(EventsQueryFilter {
273                table: "req_log".into(),
274                filter: GridFilterModel {
275                    items: vec![GridFilterItem {
276                        field: "msg; DROP".into(),
277                        operator: GridFilterOperator::Equals,
278                        value: json!("x"),
279                    }],
280                    ..Default::default()
281                },
282                ..Default::default()
283            }))
284            .expect_err("invalid filter field");
285        assert!(matches!(err, Error::Config(_)));
286    }
287}