Skip to main content

spectra_core/sinks/
chained.rs

1//! Synchronous fan-out to multiple [`SpectraSink`] implementations.
2
3use std::sync::Arc;
4
5use serde_json::Value;
6
7use crate::sink::SpectraSink;
8
9/// Forwards each emit to every sink in the chain, synchronously and in registration order.
10///
11/// Use a chain when one emit must reach multiple transports or telemetry destinations. A slow
12/// child sink delays later children and the caller.
13///
14/// # Examples
15///
16/// ```
17/// use std::sync::Arc;
18/// use spectra_core::{ChainedSink, RecordingSink, SpectraSink};
19///
20/// let first = Arc::new(RecordingSink::new());
21/// let second = Arc::new(RecordingSink::new());
22/// let chain = ChainedSink::new()
23///     .push(Arc::clone(&first) as Arc<dyn SpectraSink>)
24///     .push(Arc::clone(&second) as Arc<dyn SpectraSink>);
25///
26/// chain.record_counter("cache_hits", &[("region", "us")], 1);
27/// assert_eq!(first.counters().len(), 1);
28/// assert_eq!(second.counters().len(), 1);
29/// ```
30#[derive(Default)]
31pub struct ChainedSink {
32    sinks: Vec<Arc<dyn SpectraSink>>,
33}
34
35impl ChainedSink {
36    /// Creates an empty sink chain.
37    pub fn new() -> Self {
38        Self::default()
39    }
40
41    /// Appends a sink to the end of the chain.
42    pub fn push(mut self, sink: Arc<dyn SpectraSink>) -> Self {
43        self.sinks.push(sink);
44        self
45    }
46
47    /// Returns the number of sinks in the chain.
48    pub fn len(&self) -> usize {
49        self.sinks.len()
50    }
51
52    /// Returns whether the chain has no sinks.
53    pub fn is_empty(&self) -> bool {
54        self.sinks.is_empty()
55    }
56}
57
58impl SpectraSink for ChainedSink {
59    fn record_counter(&self, name: &str, labels: &[(&str, &str)], delta: i64) {
60        for sink in &self.sinks {
61            sink.record_counter(name, labels, delta);
62        }
63    }
64
65    fn record_gauge(&self, name: &str, labels: &[(&str, &str)], value: f64) {
66        for sink in &self.sinks {
67            sink.record_gauge(name, labels, value);
68        }
69    }
70
71    fn log_event(&self, table: &str, fields: &Value) {
72        for sink in &self.sinks {
73            sink.log_event(table, fields);
74        }
75    }
76}
77
78#[cfg(test)]
79mod tests {
80    use super::*;
81    use crate::sinks::RecordingSink;
82    use serde_json::json;
83
84    #[test]
85    fn fans_out_to_all_sinks() {
86        let a = Arc::new(RecordingSink::new());
87        let b = Arc::new(RecordingSink::new());
88        let chain = ChainedSink::new()
89            .push(Arc::clone(&a) as Arc<dyn SpectraSink>)
90            .push(Arc::clone(&b) as Arc<dyn SpectraSink>);
91
92        chain.record_counter("hits", &[("region", "us")], 2);
93        chain.record_gauge("load", &[("host", "a")], 0.5);
94        chain.log_event("request_log", &json!({"status": 200}));
95
96        assert_eq!(a.counters().len(), 1);
97        assert_eq!(b.counters().len(), 1);
98        assert_eq!(a.gauges().len(), 1);
99        assert_eq!(b.gauges().len(), 1);
100        assert_eq!(a.events().len(), 1);
101        assert_eq!(b.events().len(), 1);
102    }
103}