use std::sync::Arc;
use serde_json::Value;
use crate::sink::SpectraSink;
#[derive(Default)]
pub struct ChainedSink {
sinks: Vec<Arc<dyn SpectraSink>>,
}
impl ChainedSink {
pub fn new() -> Self {
Self::default()
}
pub fn push(mut self, sink: Arc<dyn SpectraSink>) -> Self {
self.sinks.push(sink);
self
}
pub fn len(&self) -> usize {
self.sinks.len()
}
pub fn is_empty(&self) -> bool {
self.sinks.is_empty()
}
}
impl SpectraSink for ChainedSink {
fn record_counter(&self, name: &str, labels: &[(&str, &str)], delta: i64) {
for sink in &self.sinks {
sink.record_counter(name, labels, delta);
}
}
fn record_gauge(&self, name: &str, labels: &[(&str, &str)], value: f64) {
for sink in &self.sinks {
sink.record_gauge(name, labels, value);
}
}
fn log_event(&self, table: &str, fields: &Value) {
for sink in &self.sinks {
sink.log_event(table, fields);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::sinks::RecordingSink;
use serde_json::json;
#[test]
fn fans_out_to_all_sinks() {
let a = Arc::new(RecordingSink::new());
let b = Arc::new(RecordingSink::new());
let chain = ChainedSink::new()
.push(Arc::clone(&a) as Arc<dyn SpectraSink>)
.push(Arc::clone(&b) as Arc<dyn SpectraSink>);
chain.record_counter("hits", &[("region", "us")], 2);
chain.record_gauge("load", &[("host", "a")], 0.5);
chain.log_event("request_log", &json!({"status": 200}));
assert_eq!(a.counters().len(), 1);
assert_eq!(b.counters().len(), 1);
assert_eq!(a.gauges().len(), 1);
assert_eq!(b.gauges().len(), 1);
assert_eq!(a.events().len(), 1);
assert_eq!(b.events().len(), 1);
}
}