1use std::{
4 convert::Infallible,
5 num::NonZeroUsize,
6 sync::{Arc, LazyLock},
7};
8
9use async_trait::async_trait;
10use axum::{extract::Request, Router};
11use ddsketch::DDSketch;
12use http::{Response, StatusCode};
13use prometheus_exposition::{MetricType, PrometheusRenderer};
14use saluki_common::{collections::FastIndexMap, iter::ReusableDeduplicator};
15use saluki_core::{
16 accounting::{MemoryBounds, MemoryBoundsBuilder},
17 components::{destinations::*, BuildContext},
18 data_model::{
19 event::{
20 metric::{context::Context, Histogram, Metric, MetricValues},
21 EventType,
22 },
23 tags::Tag,
24 },
25 runtime,
26};
27use saluki_error::GenericError;
28use saluki_io::net::{server::http::HttpServer, ListenAddress};
29use serde::Deserialize;
30use stringtheory::{
31 interning::{FixedSizeInterner, Interner as _},
32 MetaString,
33};
34use tokio::{select, sync::RwLock};
35use tower::util::service_fn;
36use tracing::debug;
37
38const CONTEXT_LIMIT: usize = 10_000;
39const PAYLOAD_SIZE_LIMIT_BYTES: usize = 1024 * 1024;
40const TAGS_BUFFER_SIZE_LIMIT_BYTES: usize = 2048;
41const RAW_METRICS_PATH: &str = "/metrics";
42const LEGACY_RAW_METRICS_PATH: &str = "/";
43
44const TIME_HISTOGRAM_BUCKET_COUNT: usize = 30;
46static TIME_HISTOGRAM_BUCKETS: LazyLock<[(f64, &'static str); TIME_HISTOGRAM_BUCKET_COUNT]> =
47 LazyLock::new(|| histogram_buckets::<TIME_HISTOGRAM_BUCKET_COUNT>(0.000000128, 4.0));
48
49const NON_TIME_HISTOGRAM_BUCKET_COUNT: usize = 30;
50static NON_TIME_HISTOGRAM_BUCKETS: LazyLock<[(f64, &'static str); NON_TIME_HISTOGRAM_BUCKET_COUNT]> =
51 LazyLock::new(|| histogram_buckets::<NON_TIME_HISTOGRAM_BUCKET_COUNT>(1.0, 2.0));
52
53const METRIC_NAME_STRING_INTERNER_BYTES: NonZeroUsize = NonZeroUsize::new(65536).unwrap();
55
56pub trait PrometheusPayloadProvider: Send + Sync {
58 fn render_payload(&self) -> String;
60}
61
62impl<F> PrometheusPayloadProvider for F
63where
64 F: Fn() -> String + Send + Sync,
65{
66 fn render_payload(&self) -> String {
67 self()
68 }
69}
70
71#[derive(Clone)]
72struct PrometheusAdditionalRoute {
73 path: String,
74 provider: Arc<dyn PrometheusPayloadProvider>,
75}
76
77#[derive(Deserialize)]
95pub struct PrometheusConfiguration {
96 #[serde(rename = "prometheus_listen_addr")]
97 listen_addr: ListenAddress,
98
99 #[serde(skip)]
100 additional_routes: Vec<PrometheusAdditionalRoute>,
101}
102
103impl PrometheusConfiguration {
104 pub fn from_listen_address(listen_addr: ListenAddress) -> Self {
106 Self {
107 listen_addr,
108 additional_routes: Vec::new(),
109 }
110 }
111
112 pub fn with_additional_route(
114 mut self, path: impl Into<String>, provider: Arc<dyn PrometheusPayloadProvider>,
115 ) -> Self {
116 self.additional_routes.push(PrometheusAdditionalRoute {
117 path: path.into(),
118 provider,
119 });
120 self
121 }
122}
123
124#[async_trait]
125impl DestinationBuilder for PrometheusConfiguration {
126 fn input_event_type(&self) -> EventType {
127 EventType::Metric
128 }
129
130 async fn build(&self, _context: BuildContext) -> Result<Box<dyn Destination + Send>, GenericError> {
131 Ok(Box::new(Prometheus {
132 listen_addr: self.listen_addr.clone(),
133 additional_routes: self.additional_routes.clone(),
134 metrics: FastIndexMap::default(),
135 payload: Arc::new(RwLock::new(String::new())),
136 renderer: PrometheusRenderer::new(),
137 interner: FixedSizeInterner::new(METRIC_NAME_STRING_INTERNER_BYTES),
138 }))
139 }
140}
141
142impl MemoryBounds for PrometheusConfiguration {
143 fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
144 builder
145 .minimum()
146 .with_single_value::<Prometheus>("component struct");
148
149 builder
150 .firm()
151 .with_map::<Context, PrometheusValue>("state map", CONTEXT_LIMIT)
155 .with_fixed_amount("payload size", PAYLOAD_SIZE_LIMIT_BYTES)
156 .with_fixed_amount("tags buffer", TAGS_BUFFER_SIZE_LIMIT_BYTES);
157 }
158}
159
160struct Prometheus {
161 listen_addr: ListenAddress,
162 additional_routes: Vec<PrometheusAdditionalRoute>,
163 metrics: FastIndexMap<PrometheusContext, FastIndexMap<Context, PrometheusValue>>,
164 payload: Arc<RwLock<String>>,
165 renderer: PrometheusRenderer,
166 interner: FixedSizeInterner<1>,
167}
168
169#[async_trait]
170impl Destination for Prometheus {
171 async fn run(mut self: Box<Self>, mut context: DestinationContext) -> Result<(), GenericError> {
172 let Self {
173 listen_addr,
174 additional_routes,
175 mut metrics,
176 payload,
177 mut renderer,
178 interner,
179 } = *self;
180
181 let mut health = context.take_health_handle();
182
183 runtime::nested_supervisor(
186 build_scrape_server(listen_addr, Arc::clone(&payload), additional_routes)
187 .with_worker_pool(context.topology_context().global_thread_pool().clone())
188 .into_supervisor(),
189 )
190 .spawn();
191
192 health.mark_ready();
193
194 debug!("Prometheus destination started.");
195
196 let mut contexts = 0;
197 let mut tags_deduplicator = ReusableDeduplicator::new();
198
199 loop {
200 select! {
201 _ = health.live() => continue,
202 maybe_events = context.events().next() => match maybe_events {
203 Some(events) => {
204 for event in events {
207 if let Some(metric) = event.try_into_metric() {
208 let prom_context = match into_prometheus_metric(&metric, &mut renderer, &interner) {
212 Some(prom_context) => prom_context,
213 None => continue,
214 };
215
216 let (context, values, _) = metric.into_parts();
217
218 let existing_contexts = metrics.entry(prom_context.clone()).or_default();
220 match existing_contexts.get_mut(&context) {
221 Some(existing_prom_value) => merge_metric_values_with_prom_value(values, existing_prom_value),
222 None => {
223 if contexts >= CONTEXT_LIMIT {
224 debug!("Prometheus destination reached context limit. Skipping metric '{}'.", context.name());
225 continue
226 }
227
228 let mut new_prom_value = get_prom_value_for_prom_context(&prom_context);
229 merge_metric_values_with_prom_value(values, &mut new_prom_value);
230
231 existing_contexts.insert(context, new_prom_value);
232 contexts += 1;
233 }
234 }
235 }
236 }
237
238 regenerate_payload(&metrics, &payload, &mut renderer, &mut tags_deduplicator).await;
240 },
241 None => break,
242 },
243 }
244 }
245
246 debug!("Prometheus destination stopped.");
247
248 Ok(())
249 }
250}
251
252fn build_scrape_server(
257 listen_addr: ListenAddress, payload: Arc<RwLock<String>>, additional_routes: Vec<PrometheusAdditionalRoute>,
258) -> HttpServer {
259 let additional_routes = Arc::new(additional_routes);
260 let service = service_fn(move |req: Request| {
261 let payload = Arc::clone(&payload);
262 let additional_routes = Arc::clone(&additional_routes);
263 async move {
264 Ok::<_, Infallible>(build_scrape_response(req.uri().path(), &payload, additional_routes.as_ref()).await)
265 }
266 });
267
268 HttpServer::from_listen_address(listen_addr).with_routes(Router::new().fallback_service(service))
269}
270
271async fn build_scrape_response(
272 path: &str, payload: &Arc<RwLock<String>>, additional_routes: &[PrometheusAdditionalRoute],
273) -> Response<axum::body::Body> {
274 if path == RAW_METRICS_PATH || path == LEGACY_RAW_METRICS_PATH {
275 let payload = payload.read().await;
276 return Response::new(axum::body::Body::from(payload.to_string()));
277 }
278
279 if let Some(route) = additional_routes.iter().find(|route| route.path == path) {
280 return Response::new(axum::body::Body::from(route.provider.render_payload()));
281 }
282
283 Response::builder()
284 .status(StatusCode::NOT_FOUND)
285 .body(axum::body::Body::empty())
286 .expect("response builder should accept static status and empty body")
287}
288
289#[allow(clippy::mutable_key_type)]
290async fn regenerate_payload(
291 metrics: &FastIndexMap<PrometheusContext, FastIndexMap<Context, PrometheusValue>>, payload: &Arc<RwLock<String>>,
292 renderer: &mut PrometheusRenderer, tags_deduplicator: &mut ReusableDeduplicator<Tag>,
293) {
294 renderer.clear();
295
296 for (prom_context, contexts) in metrics {
297 if !write_metrics(renderer, prom_context, contexts, tags_deduplicator) {
298 debug!("Failed to write metric to payload. Continuing...");
299 continue;
300 }
301
302 if renderer.output().len() > PAYLOAD_SIZE_LIMIT_BYTES {
303 debug!(
304 payload_len = renderer.output().len(),
305 "Payload size limit exceeded. Skipping remaining metrics."
306 );
307 break;
308 }
309 }
310
311 let mut payload = payload.write().await;
312 payload.clear();
313 payload.push_str(renderer.output());
314}
315
316fn write_metrics(
317 renderer: &mut PrometheusRenderer, prom_context: &PrometheusContext,
318 contexts: &FastIndexMap<Context, PrometheusValue>, tags_deduplicator: &mut ReusableDeduplicator<Tag>,
319) -> bool {
320 if contexts.is_empty() {
321 debug!("No contexts for metric '{}'. Skipping.", prom_context.metric_name);
322 return true;
323 }
324
325 renderer.begin_group(&prom_context.metric_name, prom_context.metric_type, None);
326
327 for (context, values) in contexts {
328 let labels = match collect_tags(context, tags_deduplicator) {
329 Some(labels) => labels,
330 None => return false,
331 };
332
333 match values {
334 PrometheusValue::Counter(value) | PrometheusValue::Gauge(value) => {
335 renderer.write_gauge_or_counter_series(labels, *value);
336 }
337 PrometheusValue::Histogram(histogram) => {
338 renderer.write_histogram_series(labels, histogram.buckets(), histogram.sum, histogram.count);
339 }
340 PrometheusValue::Summary(sketch) => {
341 let quantiles = [0.1, 0.25, 0.5, 0.95, 0.99, 0.999]
342 .into_iter()
343 .map(|q| (q, sketch.quantile(q).unwrap_or_default()));
344
345 renderer.write_summary_series(labels, quantiles, sketch.sum().unwrap_or_default(), sketch.count());
346 }
347 }
348 }
349
350 renderer.finish_group();
351 true
352}
353
354fn collect_tags<'a>(
356 context: &'a Context, tags_deduplicator: &mut ReusableDeduplicator<Tag>,
357) -> Option<Vec<(&'a str, &'a str)>> {
358 let mut labels = Vec::new();
359 let mut total_bytes = 0;
360
361 let chained_tags = context.tags().into_iter().chain(context.origin_tags());
362 let deduplicated_tags = tags_deduplicator.deduplicated(chained_tags);
363
364 for tag in deduplicated_tags {
365 let tag_name = tag.name();
366 let tag_value = match tag.value() {
367 Some(value) => value,
368 None => {
369 debug!("Skipping bare tag.");
370 continue;
371 }
372 };
373
374 total_bytes += tag_name.len() + tag_value.len() + 4;
377 if total_bytes > TAGS_BUFFER_SIZE_LIMIT_BYTES {
378 debug!("Tags buffer size limit exceeded. Tags may be missing from this metric.");
379 return None;
380 }
381
382 labels.push((tag_name, tag_value));
383 }
384
385 Some(labels)
386}
387
388#[derive(Clone, Eq, Hash, Ord, PartialEq, PartialOrd)]
389struct PrometheusContext {
390 metric_name: MetaString,
391 metric_type: MetricType,
392}
393
394enum PrometheusValue {
395 Counter(f64),
396 Gauge(f64),
397 Histogram(PrometheusHistogram),
398 Summary(DDSketch),
399}
400
401fn into_prometheus_metric(
402 metric: &Metric, renderer: &mut PrometheusRenderer, interner: &FixedSizeInterner<1>,
403) -> Option<PrometheusContext> {
404 let normalized = renderer.normalize_metric_name(metric.context().name());
406 let metric_name = match interner.try_intern(normalized).map(MetaString::from) {
407 Some(name) => name,
408 None => {
409 debug!(
410 "Failed to intern normalized metric name. Skipping metric '{}'.",
411 metric.context().name()
412 );
413 return None;
414 }
415 };
416
417 let metric_type = match metric.values() {
418 MetricValues::Counter(_) => MetricType::Counter,
419 MetricValues::Gauge(_) | MetricValues::Set(_) => MetricType::Gauge,
420 MetricValues::Histogram(_) => MetricType::Histogram,
421 MetricValues::Distribution(_) => MetricType::Summary,
422 _ => return None,
423 };
424
425 Some(PrometheusContext {
426 metric_name,
427 metric_type,
428 })
429}
430
431fn get_prom_value_for_prom_context(prom_context: &PrometheusContext) -> PrometheusValue {
432 match prom_context.metric_type {
433 MetricType::Counter => PrometheusValue::Counter(0.0),
434 MetricType::Gauge => PrometheusValue::Gauge(0.0),
435 MetricType::Histogram => PrometheusValue::Histogram(PrometheusHistogram::new(&prom_context.metric_name)),
436 MetricType::Summary => PrometheusValue::Summary(DDSketch::default()),
437 }
438}
439
440fn merge_metric_values_with_prom_value(values: MetricValues, prom_value: &mut PrometheusValue) {
441 match (values, prom_value) {
442 (MetricValues::Counter(counter_values), PrometheusValue::Counter(prom_counter)) => {
443 for (_, value) in counter_values {
444 *prom_counter += value;
445 }
446 }
447 (MetricValues::Gauge(gauge_values), PrometheusValue::Gauge(prom_gauge)) => {
448 let latest_value = gauge_values
449 .into_iter()
450 .max_by_key(|(ts, _)| ts.map(|v| v.get()).unwrap_or_default())
451 .map(|(_, value)| value)
452 .unwrap_or_default();
453 *prom_gauge = latest_value;
454 }
455 (MetricValues::Set(set_values), PrometheusValue::Gauge(prom_gauge)) => {
456 let latest_value = set_values
457 .into_iter()
458 .max_by_key(|(ts, _)| ts.map(|v| v.get()).unwrap_or_default())
459 .map(|(_, value)| value)
460 .unwrap_or_default();
461 *prom_gauge = latest_value;
462 }
463 (MetricValues::Histogram(histogram_values), PrometheusValue::Histogram(prom_histogram)) => {
464 for (_, value) in histogram_values {
465 prom_histogram.merge_histogram(&value);
466 }
467 }
468 (MetricValues::Distribution(distribution_values), PrometheusValue::Summary(prom_summary)) => {
469 for (_, value) in distribution_values {
470 prom_summary.merge(&value);
471 }
472 }
473 _ => panic!("Mismatched metric types"),
474 }
475}
476
477#[derive(Clone)]
478struct PrometheusHistogram {
479 sum: f64,
480 count: u64,
481 buckets: Vec<(f64, &'static str, u64)>,
482}
483
484impl PrometheusHistogram {
485 fn new(metric_name: &str) -> Self {
486 let base_buckets = if metric_name.ends_with("_seconds") {
488 &TIME_HISTOGRAM_BUCKETS[..]
489 } else {
490 &NON_TIME_HISTOGRAM_BUCKETS[..]
491 };
492
493 let buckets = base_buckets
494 .iter()
495 .map(|(upper_bound, upper_bound_str)| (*upper_bound, *upper_bound_str, 0))
496 .collect();
497
498 Self {
499 sum: 0.0,
500 count: 0,
501 buckets,
502 }
503 }
504
505 fn merge_histogram(&mut self, histogram: &Histogram) {
506 for sample in histogram.samples() {
507 self.add_sample(sample.value.into_inner(), sample.weight.0 as u64);
508 }
509 }
510
511 fn add_sample(&mut self, value: f64, weight: u64) {
512 self.sum += value * weight as f64;
513 self.count += weight;
514
515 for (upper_bound, _, count) in &mut self.buckets {
517 if value <= *upper_bound {
518 *count += weight;
519 }
520 }
521 }
522
523 fn buckets(&self) -> impl Iterator<Item = (&'static str, u64)> + '_ {
524 self.buckets
525 .iter()
526 .map(|(_, upper_bound_str, count)| (*upper_bound_str, *count))
527 }
528}
529
530fn histogram_buckets<const N: usize>(base: f64, scale: f64) -> [(f64, &'static str); N] {
531 let mut buckets = [(0.0, ""); N];
539
540 let log_linear_buckets = std::iter::repeat(base).enumerate().flat_map(|(i, base)| {
541 let pow = scale.powf(i as f64);
542 let value = base * pow;
543
544 let next_pow = scale.powf((i + 1) as f64);
545 let next_value = base * next_pow;
546 let midpoint = (value + next_value) / 2.0;
547
548 [value, midpoint]
549 });
550
551 for (i, current_le) in log_linear_buckets.enumerate().take(N) {
552 let (bucket_le, bucket_le_str) = &mut buckets[i];
553 let current_le_str = format!("{}", current_le);
554
555 *bucket_le = current_le;
556 *bucket_le_str = current_le_str.leak();
557 }
558
559 buckets
560}
561
562#[cfg(test)]
563mod tests {
564 use std::collections::BTreeSet;
565
566 use http_body_util::BodyExt as _;
567 use saluki_core::data_model::tags::TagSet;
568
569 use super::*;
570
571 fn prom_context(metric_type: MetricType) -> PrometheusContext {
572 PrometheusContext {
573 metric_name: MetaString::from("test.metric"),
574 metric_type,
575 }
576 }
577
578 #[test]
579 fn histogram_buckets_are_monotonic_finite_and_labeled() {
580 fn assert_bucket_table(buckets: &[(f64, &'static str)], base: f64) {
586 assert!(!buckets.is_empty(), "bucket table must not be empty");
587 assert_eq!(
588 buckets[0].0, base,
589 "first bucket upper bound must equal the configured base"
590 );
591
592 let mut previous = f64::NEG_INFINITY;
593 for (upper_bound, label) in buckets {
594 assert!(
595 upper_bound.is_finite(),
596 "bucket upper bound must be finite, got {upper_bound}"
597 );
598 assert!(
599 *upper_bound > 0.0,
600 "bucket upper bound must be positive, got {upper_bound}"
601 );
602 assert!(
603 *upper_bound > previous,
604 "bucket upper bounds must be strictly increasing ({upper_bound} !> {previous})"
605 );
606 previous = *upper_bound;
607
608 let parsed: f64 = label.parse().expect("bucket label must parse as an f64");
609 assert_eq!(
610 parsed, *upper_bound,
611 "bucket label {label:?} must render its own upper bound {upper_bound}"
612 );
613 }
614 }
615
616 assert_eq!(TIME_HISTOGRAM_BUCKETS.len(), TIME_HISTOGRAM_BUCKET_COUNT);
617 assert_eq!(NON_TIME_HISTOGRAM_BUCKETS.len(), NON_TIME_HISTOGRAM_BUCKET_COUNT);
618 assert_bucket_table(&TIME_HISTOGRAM_BUCKETS[..], 0.000000128);
619 assert_bucket_table(&NON_TIME_HISTOGRAM_BUCKETS[..], 1.0);
620 }
621
622 #[test]
623 fn prom_histogram_add_sample() {
624 let sample1 = (0.25, 1);
625 let sample2 = (1.0, 2);
626 let sample3 = (2.0, 3);
627
628 let mut histogram = PrometheusHistogram::new("time_metric_seconds");
629 histogram.add_sample(sample1.0, sample1.1);
630 histogram.add_sample(sample2.0, sample2.1);
631 histogram.add_sample(sample3.0, sample3.1);
632
633 let sample1_weighted_value = sample1.0 * sample1.1 as f64;
634 let sample2_weighted_value = sample2.0 * sample2.1 as f64;
635 let sample3_weighted_value = sample3.0 * sample3.1 as f64;
636 let expected_sum = sample1_weighted_value + sample2_weighted_value + sample3_weighted_value;
637 let expected_count = sample1.1 + sample2.1 + sample3.1;
638 assert_eq!(histogram.sum, expected_sum);
639 assert_eq!(histogram.count, expected_count);
640
641 let mut expected_bucket_count = 0;
643 for sample in [sample1, sample2, sample3] {
644 for bucket in &histogram.buckets {
645 if sample.0 <= bucket.0 {
648 assert!(bucket.2 >= expected_bucket_count + sample.1);
649 }
650 }
651
652 expected_bucket_count += sample.1;
654 }
655 }
656
657 #[tokio::test]
658 async fn scrape_routes_serve_raw_compat_and_404() {
659 let payload = Arc::new(RwLock::new("raw".to_string()));
660 let routes = vec![PrometheusAdditionalRoute {
661 path: "/compat/metrics".to_string(),
662 provider: Arc::new(|| "compat".to_string()),
663 }];
664
665 let raw_response = build_scrape_response("/metrics", &payload, &routes).await;
666 assert_eq!(raw_response.status(), StatusCode::OK);
667 let raw_body = raw_response
668 .into_body()
669 .collect()
670 .await
671 .expect("body should collect")
672 .to_bytes();
673 assert_eq!(&raw_body[..], b"raw");
674
675 let legacy_response = build_scrape_response("/", &payload, &routes).await;
676 assert_eq!(legacy_response.status(), StatusCode::OK);
677 let legacy_body = legacy_response
678 .into_body()
679 .collect()
680 .await
681 .expect("body should collect")
682 .to_bytes();
683 assert_eq!(&legacy_body[..], b"raw");
684
685 let compat_response = build_scrape_response("/compat/metrics", &payload, &routes).await;
686 assert_eq!(compat_response.status(), StatusCode::OK);
687 let compat_body = compat_response
688 .into_body()
689 .collect()
690 .await
691 .expect("body should collect")
692 .to_bytes();
693 assert_eq!(&compat_body[..], b"compat");
694
695 let missing_response = build_scrape_response("/missing", &payload, &routes).await;
696 assert_eq!(missing_response.status(), StatusCode::NOT_FOUND);
697 }
698
699 #[test]
700 fn into_prometheus_metric_maps_value_type_and_normalizes_name() {
701 let mut renderer = PrometheusRenderer::new();
702 let interner = FixedSizeInterner::<1>::new(METRIC_NAME_STRING_INTERNER_BYTES);
703
704 let counter = into_prometheus_metric(&Metric::counter("my.counter", 1.0), &mut renderer, &interner)
707 .expect("counter should translate");
708 assert_eq!("my__counter", counter.metric_name.as_ref());
709 assert!(matches!(counter.metric_type, MetricType::Counter));
710
711 let gauge =
713 into_prometheus_metric(&Metric::gauge("g", 1.0), &mut renderer, &interner).expect("gauge should translate");
714 assert!(matches!(gauge.metric_type, MetricType::Gauge));
715 let set =
716 into_prometheus_metric(&Metric::set("s", "a"), &mut renderer, &interner).expect("set should translate");
717 assert!(matches!(set.metric_type, MetricType::Gauge));
718
719 let histogram = into_prometheus_metric(&Metric::histogram("h", [1.0]), &mut renderer, &interner)
721 .expect("histogram should translate");
722 assert!(matches!(histogram.metric_type, MetricType::Histogram));
723
724 let distribution = into_prometheus_metric(&Metric::distribution("d", [1.0]), &mut renderer, &interner)
727 .expect("distribution should translate");
728 assert!(matches!(distribution.metric_type, MetricType::Summary));
729 }
730
731 #[test]
732 fn merge_counter_sums_all_points() {
733 let mut value = get_prom_value_for_prom_context(&prom_context(MetricType::Counter));
734 let (_, values, _) = Metric::counter("c", [(1, 1.0), (2, 2.0), (3, 4.0)]).into_parts();
735 merge_metric_values_with_prom_value(values, &mut value);
736 match value {
737 PrometheusValue::Counter(sum) => assert_eq!(7.0, sum),
738 _ => panic!("counter values should stay a counter"),
739 }
740 }
741
742 #[test]
743 fn merge_gauge_keeps_latest_by_timestamp() {
744 let mut value = get_prom_value_for_prom_context(&prom_context(MetricType::Gauge));
745 let (_, values, _) = Metric::gauge("g", [(10, 1.0), (30, 3.0), (20, 2.0)]).into_parts();
747 merge_metric_values_with_prom_value(values, &mut value);
748 match value {
749 PrometheusValue::Gauge(latest) => assert_eq!(3.0, latest),
750 _ => panic!("gauge values should stay a gauge"),
751 }
752 }
753
754 #[test]
755 fn merge_set_maps_cardinality_into_gauge() {
756 let mut value = get_prom_value_for_prom_context(&prom_context(MetricType::Gauge));
757 let (_, values, _) = Metric::set("s", "a").into_parts();
758 merge_metric_values_with_prom_value(values, &mut value);
759 match value {
760 PrometheusValue::Gauge(cardinality) => assert_eq!(1.0, cardinality),
761 _ => panic!("set values should merge into a gauge"),
762 }
763 }
764
765 #[test]
766 fn merge_histogram_accumulates_sum_and_count() {
767 let mut value = get_prom_value_for_prom_context(&prom_context(MetricType::Histogram));
768 let (_, values, _) = Metric::histogram("h", [1.0, 2.0, 3.0]).into_parts();
769 merge_metric_values_with_prom_value(values, &mut value);
770 match value {
771 PrometheusValue::Histogram(histogram) => {
772 assert_eq!(3, histogram.count);
773 assert_eq!(6.0, histogram.sum);
774 }
775 _ => panic!("histogram values should stay a histogram"),
776 }
777 }
778
779 #[test]
780 fn merge_distribution_folds_into_ddsketch_summary() {
781 let mut value = get_prom_value_for_prom_context(&prom_context(MetricType::Summary));
782 let (_, values, _) = Metric::distribution("d", [1.0, 2.0, 3.0, 4.0, 5.0]).into_parts();
783 merge_metric_values_with_prom_value(values, &mut value);
784 match value {
785 PrometheusValue::Summary(sketch) => {
786 assert_eq!(5, sketch.count());
787 let median = sketch.quantile(0.5).expect("median should be computable");
788 assert!((median - 3.0).abs() <= 0.5, "median should be ~= 3.0, got {median}");
789 }
790 _ => panic!("distribution values should merge into a summary"),
791 }
792 }
793
794 #[test]
795 #[should_panic(expected = "Mismatched metric types")]
796 fn merge_mismatched_value_and_accumulator_types_panics() {
797 let mut value = get_prom_value_for_prom_context(&prom_context(MetricType::Counter));
799 let (_, values, _) = Metric::gauge("g", 1.0).into_parts();
800 merge_metric_values_with_prom_value(values, &mut value);
801 }
802
803 #[test]
804 fn collect_tags_skips_bare_tags_and_keeps_key_value_pairs() {
805 let mut tags_deduplicator = ReusableDeduplicator::new();
806 let context = Context::from_static_parts("m", &["env:prod", "bare", "team:core"]);
807
808 let labels = collect_tags(&context, &mut tags_deduplicator).expect("tags should collect");
809 let labels = labels.into_iter().collect::<BTreeSet<_>>();
810
811 assert_eq!(BTreeSet::from([("env", "prod"), ("team", "core")]), labels);
813 }
814
815 #[test]
816 fn collect_tags_returns_none_when_buffer_limit_exceeded() {
817 let mut tags_deduplicator = ReusableDeduplicator::new();
818 let oversized_tag = format!("big:{}", "x".repeat(TAGS_BUFFER_SIZE_LIMIT_BYTES + 1));
820 let tags = std::iter::once(Tag::from(oversized_tag)).collect::<TagSet>();
821 let context = Context::from_parts("m", tags);
822
823 assert_eq!(None, collect_tags(&context, &mut tags_deduplicator));
824 }
825}