Skip to main content

micromegas_telemetry_sink/
composite_event_sink.rs

1use micromegas_tracing::{
2    event::{BoxedEventSink, EventSink},
3    images::{ImageBlock, ImageStream},
4    logs::{LogBlock, LogMetadata, LogStream},
5    metrics::{MetricsBlock, MetricsStream},
6    prelude::*,
7    property_set::Property,
8    spans::{ThreadBlock, ThreadStream},
9};
10use std::{fmt, sync::Arc};
11
12/// An `EventSink` that dispatches events to multiple other `EventSink`s.
13///
14/// This allows for fanning out telemetry events to different sinks, each with its own
15/// filtering and processing logic.
16pub struct CompositeSink {
17    sinks: Vec<(LevelFilter, BoxedEventSink)>,
18    target_level_filters: Vec<(String, LevelFilter)>,
19}
20
21impl CompositeSink {
22    /// Creates a new `CompositeSink`.
23    pub fn new(
24        sinks: Vec<(LevelFilter, BoxedEventSink)>,
25        target_max_level: Vec<(String, LevelFilter)>,
26        max_level_override: Option<LevelFilter>,
27    ) -> Self {
28        if let Some(max_level) = max_level_override {
29            micromegas_tracing::levels::set_max_level(max_level);
30        } else {
31            let mut max_level = LevelFilter::Off;
32            for (_, level_filter) in &target_max_level {
33                max_level = max_level.max(*level_filter);
34            }
35            for (level_filter, _) in &sinks {
36                max_level = max_level.max(*level_filter);
37            }
38            micromegas_tracing::levels::set_max_level(max_level);
39        }
40
41        let mut target_max_level = target_max_level;
42        target_max_level.sort_by_key(|(name, _)| name.len().wrapping_neg());
43
44        Self {
45            sinks,
46            target_level_filters: target_max_level,
47        }
48    }
49
50    fn target_max_level(&self, metadata: &LogMetadata) -> Option<LevelFilter> {
51        const GENERATION: u16 = 1;
52        // At this point we would have already tested the max level on the macro
53        match metadata.level_filter(GENERATION) {
54            micromegas_tracing::logs::FilterState::Outdated => {
55                let level_filter =
56                    Self::find_max_match(metadata.target, &self.target_level_filters);
57                metadata.set_level_filter(GENERATION, level_filter);
58                level_filter
59            }
60            micromegas_tracing::logs::FilterState::NotSet => None,
61            micromegas_tracing::logs::FilterState::Set(level_filter) => Some(level_filter),
62        }
63    }
64
65    /// This needs to be optimized
66    fn find_max_match(
67        target: &str,
68        level_filters: &[(String, LevelFilter)],
69    ) -> Option<LevelFilter> {
70        for (t, l) in level_filters.iter() {
71            if target.starts_with(t) {
72                return Some(*l);
73            }
74        }
75        None
76    }
77}
78
79impl EventSink for CompositeSink {
80    fn on_startup(&self, process_info: Arc<ProcessInfo>) {
81        if self.sinks.len() == 1 {
82            self.sinks[0].1.on_startup(process_info);
83        } else {
84            self.sinks
85                .iter()
86                .for_each(|(_, sink)| sink.on_startup(process_info.clone()));
87        }
88    }
89
90    fn on_shutdown(&self) {
91        self.sinks.iter().for_each(|(_, sink)| sink.on_shutdown());
92    }
93
94    fn on_log_enabled(&self, metadata: &LogMetadata) -> bool {
95        // The log is enabled if any of the sinks are enabled
96        // If the sinks vec is empty `false` will be returned
97        let target_max_level = self.target_max_level(metadata);
98        self.sinks.iter().any(|(max_level, sink)| {
99            metadata.level <= target_max_level.unwrap_or(*max_level)
100                && sink.on_log_enabled(metadata)
101        })
102    }
103
104    fn on_log(
105        &self,
106        metadata: &LogMetadata,
107        properties: &[Property],
108        time: i64,
109        args: fmt::Arguments<'_>,
110    ) {
111        let target_max_level = self.target_max_level(metadata);
112        self.sinks.iter().for_each(|(max_level, sink)| {
113            if metadata.level <= target_max_level.unwrap_or(*max_level)
114                && sink.on_log_enabled(metadata)
115            {
116                sink.on_log(metadata, properties, time, args);
117            }
118        });
119    }
120
121    fn on_init_log_stream(&self, log_stream: &LogStream) {
122        self.sinks
123            .iter()
124            .for_each(|(_, sink)| sink.on_init_log_stream(log_stream));
125    }
126
127    fn on_process_log_block(&self, old_event_block: Arc<LogBlock>) {
128        self.sinks
129            .iter()
130            .for_each(|(_, sink)| sink.on_process_log_block(old_event_block.clone()));
131    }
132
133    fn on_init_metrics_stream(&self, metrics_stream: &MetricsStream) {
134        self.sinks
135            .iter()
136            .for_each(|(_, sink)| sink.on_init_metrics_stream(metrics_stream));
137    }
138
139    fn on_process_metrics_block(&self, old_event_block: Arc<MetricsBlock>) {
140        self.sinks
141            .iter()
142            .for_each(|(_, sink)| sink.on_process_metrics_block(old_event_block.clone()));
143    }
144
145    fn on_init_thread_stream(&self, thread_stream: &ThreadStream) {
146        self.sinks
147            .iter()
148            .for_each(|(_, sink)| sink.on_init_thread_stream(thread_stream));
149    }
150
151    fn on_process_thread_block(&self, old_event_block: Arc<ThreadBlock>) {
152        self.sinks
153            .iter()
154            .for_each(|(_, sink)| sink.on_process_thread_block(old_event_block.clone()));
155    }
156
157    fn on_init_image_stream(&self, stream: &ImageStream) {
158        self.sinks
159            .iter()
160            .for_each(|(_, sink)| sink.on_init_image_stream(stream));
161    }
162
163    fn on_process_image_block(&self, block: Arc<ImageBlock>) {
164        self.sinks
165            .iter()
166            .for_each(|(_, sink)| sink.on_process_image_block(block.clone()));
167    }
168
169    fn is_busy(&self) -> bool {
170        for (_, sink) in &self.sinks {
171            if sink.is_busy() {
172                return true;
173            }
174        }
175        false
176    }
177}