micromegas_telemetry_sink/
composite_event_sink.rs1use 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
12pub struct CompositeSink {
17 sinks: Vec<(LevelFilter, BoxedEventSink)>,
18 target_level_filters: Vec<(String, LevelFilter)>,
19}
20
21impl CompositeSink {
22 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 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 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 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}