1use 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
17pub 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 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 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 pub fn register_event_backend(&self, table: impl Into<String>, backend: SharedEventBackend) {
83 self.events.write().insert(table.into(), backend);
84 }
85
86 pub fn register_metrics_backend(&self, name: impl Into<String>, backend: SharedMetricsBackend) {
88 self.metrics.write().insert(name.into(), backend);
89 }
90
91 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 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 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 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 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 pub fn set_global(router: Arc<Self>) {
159 let _ = GLOBAL_ROUTER.set(router);
160 }
161
162 pub fn global() -> Arc<Self> {
172 #[allow(clippy::expect_used)]
174 GLOBAL_ROUTER
175 .get()
176 .cloned()
177 .expect("SpectraRouter::set_global was not called")
178 }
179
180 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}