saluki_core/observability/metrics/
processor.rs1use std::collections::BTreeMap;
12
13use prometheus_exposition::{MetricType, PrometheusRenderer};
14use stringtheory::MetaString;
15
16use super::aggregated::{AggregatedMetricValue, AggregatedMetricsState};
17use super::remapper::RemapperRule;
18
19pub struct TelemetryProcessor {
24 renderer: PrometheusRenderer,
25 rules: Vec<RemapperRule>,
26}
27
28impl TelemetryProcessor {
29 pub fn new() -> Self {
33 Self {
34 renderer: PrometheusRenderer::new(),
35 rules: Vec::new(),
36 }
37 }
38
39 pub fn with_remapper_rules(mut self, rules: Vec<RemapperRule>) -> Self {
45 self.rules = rules;
46 self
47 }
48
49 pub fn process(&mut self, state: &AggregatedMetricsState) -> String {
51 self.renderer.clear();
52
53 if self.rules.is_empty() {
54 self.process_all(state);
55 } else {
56 self.process_remapped(state);
57 }
58
59 self.renderer.output().to_string()
60 }
61
62 fn process_all(&mut self, state: &AggregatedMetricsState) {
63 let mut groups: BTreeMap<MetaString, BTreeMap<Vec<MetaString>, AggregatedMetricValue>> = BTreeMap::new();
66
67 state.visit_metrics(|context, value| {
68 let mut tags: Vec<MetaString> = context.tags().into_iter().map(|tag| tag.clone().into_inner()).collect();
69 tags.sort();
70
71 let series = groups.entry(context.name().clone()).or_default();
72 series
73 .entry(tags)
74 .and_modify(|existing| existing.merge(value))
75 .or_insert_with(|| value.clone());
76 });
77
78 for (name, series) in &groups {
79 render_group(&mut self.renderer, name, None, series);
80 }
81 }
82
83 fn process_remapped(&mut self, state: &AggregatedMetricsState) {
84 let mut groups: BTreeMap<&'static str, BTreeMap<Vec<MetaString>, AggregatedMetricValue>> = BTreeMap::new();
90 let mut help_text: BTreeMap<&'static str, &'static str> = BTreeMap::new();
91
92 state.visit_metrics(|context, value| {
93 for rule in &self.rules {
94 if let Some(mut remapped) = rule.try_match_no_context(context) {
95 remapped.tags.sort();
96
97 let continue_matching = rule.should_continue_matching();
98 if let Some(text) = rule.help_text() {
99 help_text.entry(remapped.name).or_insert(text);
100 }
101 let series = groups.entry(remapped.name).or_default();
102 series
103 .entry(remapped.tags)
104 .and_modify(|existing| existing.merge(value))
105 .or_insert_with(|| value.clone());
106
107 if !continue_matching {
108 return;
109 }
110 }
111 }
112 });
113
114 for (name, series) in &groups {
115 render_group(&mut self.renderer, name, help_text.get(name).copied(), series);
116 }
117 }
118}
119
120impl Default for TelemetryProcessor {
121 fn default() -> Self {
122 Self::new()
123 }
124}
125
126fn render_group(
127 renderer: &mut PrometheusRenderer, name: &str, help_text: Option<&str>,
128 series: &BTreeMap<Vec<MetaString>, AggregatedMetricValue>,
129) {
130 let metric_type = match series.values().next() {
132 Some(AggregatedMetricValue::Counter(_)) => MetricType::Counter,
133 Some(AggregatedMetricValue::Gauge(_)) => MetricType::Gauge,
134 Some(AggregatedMetricValue::Histogram(_)) => MetricType::Histogram,
135 None => return,
136 };
137
138 match metric_type {
139 MetricType::Counter | MetricType::Gauge => {
140 let rendered = series.iter().map(|(tags, value)| (split_tags(tags), value.value()));
141 renderer.render_scalar_group(name, metric_type, help_text, rendered);
142 }
143 MetricType::Histogram => {
144 renderer.begin_group(name, metric_type, help_text);
145 for (tags, value) in series {
146 if let AggregatedMetricValue::Histogram(histogram) = value {
147 renderer.write_histogram_series(
148 split_tags(tags),
149 histogram.buckets(),
150 histogram.sum(),
151 histogram.count(),
152 );
153 }
154 }
155 renderer.finish_group();
156 }
157 MetricType::Summary => {}
158 }
159}
160
161fn split_tags(tags: &[MetaString]) -> impl Iterator<Item = (&str, &str)> {
163 tags.iter().filter_map(|tag| tag.as_ref().split_once(':'))
164}
165
166#[cfg(test)]
167mod tests {
168 use super::super::aggregate_upserts;
169 use super::*;
170 use crate::data_model::event::{
171 metric::{context::Context, Metric},
172 Event,
173 };
174
175 #[test]
176 fn renders_counter_and_gauge_groups_without_rules() {
177 let state = aggregate_upserts(vec![
178 Event::Metric(Metric::counter(
179 Context::from_static_parts("adp.requests_total", &["method:get"]),
180 10.0,
181 )),
182 Event::Metric(Metric::counter(
183 Context::from_static_parts("adp.requests_total", &["method:post"]),
184 3.0,
185 )),
186 Event::Metric(Metric::gauge(
187 Context::from_static_parts("adp.queue_depth", &["queue:work"]),
188 5.0,
189 )),
190 ]);
191
192 let output = TelemetryProcessor::new().process(&state);
193
194 assert!(output.contains("# TYPE adp__requests_total counter"));
195 assert!(output.contains("adp__requests_total{method=\"get\"} 10"));
196 assert!(output.contains("adp__requests_total{method=\"post\"} 3"));
197 assert!(output.contains("# TYPE adp__queue_depth gauge"));
198 assert!(output.contains("adp__queue_depth{queue=\"work\"} 5"));
199 }
200
201 #[test]
202 fn renders_histogram_groups_without_rules() {
203 let state = aggregate_upserts(vec![
204 Event::Metric(Metric::histogram(
205 Context::from_static_parts("adp.latency_seconds", &["op:read"]),
206 [0.001, 0.002, 0.5],
207 )),
208 Event::Metric(Metric::histogram(
209 Context::from_static_parts("adp.latency_seconds", &["op:write"]),
210 [0.01, 0.02],
211 )),
212 ]);
213
214 let output = TelemetryProcessor::new().process(&state);
215
216 assert!(output.contains("# TYPE adp__latency_seconds histogram"));
217 assert!(output.contains("adp__latency_seconds_bucket{op=\"read\","));
218 assert!(output.contains("adp__latency_seconds_bucket{op=\"write\","));
219 assert!(output.contains("le=\"+Inf\""));
220 assert!(output.contains("adp__latency_seconds_sum{op=\"read\"}"));
221 assert!(output.contains("adp__latency_seconds_count{op=\"read\"} 3"));
222 assert!(output.contains("adp__latency_seconds_count{op=\"write\"} 2"));
223 }
224
225 #[test]
226 fn renders_only_matched_metrics_with_rules() {
227 let state = aggregate_upserts(vec![
228 Event::Metric(Metric::counter(
229 Context::from_static_parts("src.matched", &["component_id:x"]),
230 42.0,
231 )),
232 Event::Metric(Metric::counter(
233 Context::from_static_parts("src.unmatched", &["component_id:x"]),
234 100.0,
235 )),
236 ]);
237
238 let rules =
239 vec![RemapperRule::by_name("src.matched", "dst.renamed").with_help_text("Renamed counter help text")];
240
241 let output = TelemetryProcessor::new().with_remapper_rules(rules).process(&state);
242
243 assert!(output.contains("# HELP dst__renamed Renamed counter help text"));
244 assert!(output.contains("# TYPE dst__renamed counter"));
245 assert!(output.contains("dst__renamed 42"));
246 assert!(!output.contains("unmatched"));
248 }
249
250 #[test]
251 fn renders_histograms_through_rules() {
252 let state = aggregate_upserts(vec![Event::Metric(Metric::histogram(
253 Context::from_static_parts("src.latency_seconds", &["op:read"]),
254 [0.001, 0.5],
255 ))]);
256
257 let rules = vec![RemapperRule::by_name("src.latency_seconds", "dst.latency_seconds")
258 .with_original_tags(["op"])
259 .with_help_text("Remapped latency")];
260
261 let output = TelemetryProcessor::new().with_remapper_rules(rules).process(&state);
262
263 assert!(output.contains("# HELP dst__latency_seconds Remapped latency"));
264 assert!(output.contains("# TYPE dst__latency_seconds histogram"));
265 assert!(output.contains("dst__latency_seconds_bucket{op=\"read\","));
266 assert!(output.contains("dst__latency_seconds_count{op=\"read\"} 2"));
267 }
268}