saluki_components/destinations/prometheus/
mod.rs

1//! Prometheus destination.
2
3use 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
44// Histogram-related constants and pre-calculated buckets.
45const 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
53// SAFETY: This is obviously not zero.
54const METRIC_NAME_STRING_INTERNER_BYTES: NonZeroUsize = NonZeroUsize::new(65536).unwrap();
55
56/// Provides a Prometheus scrape payload for an additional route.
57pub trait PrometheusPayloadProvider: Send + Sync {
58    /// Renders the current Prometheus text payload.
59    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/// Prometheus destination.
78///
79/// Exposes a Prometheus scrape endpoint that emits metrics in the Prometheus exposition format.
80///
81/// # Limits
82///
83/// - Number of contexts (unique series) is limited to 10,000.
84/// - Maximum size of scrape payload response is ~1MiB.
85///
86/// # Missing
87///
88/// - no support for expiring metrics (which we don't really need because the only use for this destination at the
89///   moment is internal metrics, which aren't dynamic since we don't use dynamic tags or have dynamic topology support,
90///   but... you know, we'll eventually need this)
91/// - full support for distributions (we can't convert a distribution to an aggregated histogram, and native histogram
92///   support is still too fresh for most clients, so we simply expose aggregated summaries as a stopgap)
93///
94#[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    /// Creates a new `PrometheusConfiguration` for the given listen address.
105    pub fn from_listen_address(listen_addr: ListenAddress) -> Self {
106        Self {
107            listen_addr,
108            additional_routes: Vec::new(),
109        }
110    }
111
112    /// Adds an additional scrape route backed by the given payload provider.
113    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            // Capture the size of the heap allocation when the component is built.
147            .with_single_value::<Prometheus>("component struct");
148
149        builder
150            .firm()
151            // Even though our context map is really the Prometheus context to a map of context/value pairs, we're just
152            // simplifying things here because the ratio of true "contexts" to Prometheus contexts should be very high,
153            // high enough to make this a reasonable approximation.
154            .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        // The scrape endpoint runs as a supervised worker of its own rather than as part of this component: the
184        // component's supervisor is what stops it, and drains its in-flight connections, once this component is done.
185        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                        // Process each metric event in the batch, either merging it with the existing value or
205                        // inserting it for the first time.
206                        for event in events {
207                            if let Some(metric) = event.try_into_metric() {
208                                // Break apart our metric into its constituent parts, and then normalize it for
209                                // Prometheus: adjust the name if necessary, figuring out the equivalent Prometheus
210                                // metric type, and so on.
211                                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                                // Create an entry for the context if we don't already have one, obeying our configured context limit.
219                                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 the scrape payload.
239                        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
252/// Builds the server that answers scrape requests.
253///
254/// Every path is answered by the same service rather than being routed: the raw metrics paths take precedence over the
255/// additional routes, so a path claimed by both is answered with the raw payload rather than being a conflict.
256fn 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
354/// Collects tags from a context into key-value pairs suitable for the renderer.
355fn 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        // Can't exceed the tags buffer size limit: we calculate the addition as tag name/value length plus three bytes
375        // to account for having to format it as `name="value",`.
376        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    // Normalize the metric name using the renderer, then intern it.
405    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        // Super hacky but effective way to decide when to switch to the time-oriented buckets.
487        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        // Add the value to each bucket that it falls into, up to the maximum number of buckets.
516        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    // We generate a set of "log-linear" buckets: logarithmically spaced values which are then subdivided linearly.
532    //
533    // As an example, with base=2 and scale=4, we would get: 2, 5, 8, 20, 32, 80, 128, 320, 512, and so on.
534    //
535    // We calculate buckets in pairs, where the n-th pair is `i` and `j`, such that `i` is `base * scale^n` and `j` is
536    // the midpoint between `i` and the next `i` (`base * scale^(n+1)`).
537
538    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        // The pre-computed log-linear bucket tables back every Prometheus histogram we expose, so their invariants are
581        // load-bearing: exactly N upper bounds, strictly increasing, finite and positive, starting at the configured
582        // base, each paired with a string label that renders (and parses back to) its own float value. A regression
583        // producing NaN/inf bounds, non-monotonic spacing, or a mismatched label would silently corrupt the exposition
584        // output, so assert the contract rather than just printing the tables.
585        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        // Go through and make sure we have things in the right buckets.
642        let mut expected_bucket_count = 0;
643        for sample in [sample1, sample2, sample3] {
644            for bucket in &histogram.buckets {
645                // If we've finally hit a bucket that includes our sample value, it's count should be equal to or
646                // greater than our expected bucket count when we account for the current sample.
647                if sample.0 <= bucket.0 {
648                    assert!(bucket.2 >= expected_bucket_count + sample.1);
649                }
650            }
651
652            // Adjust the expected bucket count to fully account for the current sample before moving on.
653            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        // The metric name is normalized to the Prometheus character set: the `.` separator is not a valid name
705        // character, so it is replaced (the renderer yields `my__counter`).
706        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        // Gauges and sets both map to a Prometheus gauge.
712        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        // Histograms map to a native Prometheus histogram.
720        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        // Distributions are exposed as a DDSketch-backed summary: the documented stopgap, since a distribution can't
725        // be converted to an aggregated histogram.
726        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        // Points are supplied out of timestamp order; the highest timestamp's value wins.
746        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        // Feeding gauge values into a counter accumulator is an invariant violation and panics.
798        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        // Key/value tags become labels; the bare `bare` tag (no value) is dropped.
812        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        // A single tag larger than the 2048-byte buffer trips the size-limit bail-out.
819        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}