Skip to main content

praxis_protocol/http/pingora/
metrics.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 Praxis Contributors
3
4//! Prometheus metrics: recorder installation, metric recording for the HTTP
5//! and upstream paths, config-reload counters, and scrape rendering.
6//!
7//! Every recorder guards on `is_recorder_installed`, so recording stays a
8//! cheap no-op until the exporter is installed. Label emission then takes one
9//! of two shapes: the default, all-enabled configuration keeps the
10//! static-label macro form so it emits exactly the series it did before label
11//! selection existed, while disabling a dimension drops it from the emitted
12//! labels (`selected_labels` builds the reduced set for the multi-label
13//! metrics; the single-label recorders check the flag inline).
14
15use std::sync::OnceLock;
16
17use metrics::{Label, SharedString, counter, gauge, histogram};
18use metrics_exporter_prometheus::{Matcher, PrometheusBuilder, PrometheusHandle};
19use praxis_core::config::{MetricLabel, MetricLabelsConfig};
20
21// -----------------------------------------------------------------------------
22// Metric Names
23// -----------------------------------------------------------------------------
24
25/// Counter for completed HTTP requests.
26const HTTP_REQUESTS_TOTAL: &str = "praxis_http_requests_total";
27
28/// Histogram for HTTP request duration in seconds.
29const HTTP_REQUEST_DURATION_SECONDS: &str = "praxis_http_request_duration_seconds";
30
31/// Histogram for HTTP request body size in bytes.
32const HTTP_REQUEST_BODY_BYTES: &str = "praxis_http_request_body_bytes";
33
34/// Histogram for HTTP response body size in bytes.
35const HTTP_RESPONSE_BODY_BYTES: &str = "praxis_http_response_body_bytes";
36
37/// Gauge for in-flight HTTP requests per listener.
38///
39/// Each active request holds an RAII guard for its lifetime, so the
40/// decrement runs on every terminal path, including client aborts and
41/// reset HTTP/2 streams. TCP connections are tracked separately by
42/// [`TCP_ACTIVE_CONNECTIONS`].
43const HTTP_ACTIVE_REQUESTS: &str = "praxis_http_active_requests";
44
45/// Gauge for open TCP sessions per listener.
46const TCP_ACTIVE_CONNECTIONS: &str = "praxis_tcp_active_connections";
47
48/// Counter for connections rejected by overload protection.
49const OVERLOAD_REJECTS_TOTAL: &str = "praxis_overload_rejects_total";
50
51/// Counter for requests sent to upstream endpoints.
52const UPSTREAM_REQUESTS_TOTAL: &str = "praxis_upstream_requests_total";
53
54/// Histogram for upstream connect duration in seconds.
55const UPSTREAM_CONNECT_DURATION_SECONDS: &str = "praxis_upstream_connect_duration_seconds";
56
57/// Counter for upstream connect failures.
58const UPSTREAM_CONNECT_FAILURES_TOTAL: &str = "praxis_upstream_connect_failures_total";
59
60/// Counter for upstream connect-failure retries.
61const UPSTREAM_RETRIES_TOTAL: &str = "praxis_upstream_retries_total";
62
63/// Gauge for healthy endpoints per cluster.
64const UPSTREAM_HEALTHY_ENDPOINTS: &str = "praxis_upstream_healthy_endpoints";
65
66/// Gauge for total endpoints per cluster.
67const UPSTREAM_TOTAL_ENDPOINTS: &str = "praxis_upstream_total_endpoints";
68
69/// Counter for endpoint health state transitions.
70const UPSTREAM_HEALTH_TRANSITIONS_TOTAL: &str = "praxis_upstream_health_transitions_total";
71
72/// Counter for config reload attempts.
73const CONFIG_RELOAD_TOTAL: &str = "praxis_config_reload_total";
74
75/// Gauge for unix timestamp of last successful config reload.
76const CONFIG_RELOAD_LAST_SUCCESS_TIMESTAMP: &str = "praxis_config_reload_last_success_timestamp";
77
78/// Counter for proxy errors not already counted by a dedicated metric.
79const ERRORS_TOTAL: &str = "praxis_errors_total";
80
81// -----------------------------------------------------------------------------
82// Label Values
83// -----------------------------------------------------------------------------
84
85/// Error type: a filter rejected the request.
86pub(crate) const ERROR_TYPE_FILTER_REJECT: &str = "filter_reject";
87
88/// Error type: an upstream connect, read or write timed out.
89pub(crate) const ERROR_TYPE_TIMEOUT: &str = "timeout";
90
91/// Error type: the upstream could not be reached.
92pub(crate) const ERROR_TYPE_UPSTREAM_UNAVAILABLE: &str = "upstream_unavailable";
93
94/// Error type: the upstream was reached but the exchange failed.
95pub(crate) const ERROR_TYPE_UPSTREAM_PROTOCOL: &str = "upstream_protocol";
96
97/// Error type: the downstream client connection failed.
98pub(crate) const ERROR_TYPE_DOWNSTREAM: &str = "downstream";
99
100/// Error type: an internal proxy fault.
101pub(crate) const ERROR_TYPE_INTERNAL: &str = "internal";
102
103/// Overload reject reason: process memory pressure.
104pub(crate) const OVERLOAD_REASON_MEMORY: &str = "memory";
105
106/// Overload reject reason: process-wide connection limit.
107pub(crate) const OVERLOAD_REASON_GLOBAL_CONNECTIONS: &str = "global_connections";
108
109/// Overload reject reason: per-listener connection limit.
110pub(crate) const OVERLOAD_REASON_LISTENER_CONNECTIONS: &str = "listener_connections";
111
112/// Retry result: connect eventually succeeded after at least one retry.
113pub(crate) const RETRY_RESULT_SUCCESS: &str = "success";
114
115/// Retry result: retries gave up or were exhausted.
116pub(crate) const RETRY_RESULT_EXHAUSTED: &str = "exhausted";
117
118/// Health transition result: endpoint became healthy.
119pub(crate) const HEALTH_RESULT_HEALTHY: &str = "healthy";
120
121/// Health transition result: endpoint became unhealthy.
122pub(crate) const HEALTH_RESULT_UNHEALTHY: &str = "unhealthy";
123
124/// Config reload result: success.
125pub(crate) const RELOAD_RESULT_SUCCESS: &str = "success";
126
127/// Config reload result: failure.
128pub(crate) const RELOAD_RESULT_FAILURE: &str = "failure";
129
130// -----------------------------------------------------------------------------
131// Histogram Buckets
132// -----------------------------------------------------------------------------
133
134/// Histogram bucket upper bounds for HTTP body sizes in bytes.
135///
136/// `PrometheusBuilder`'s default buckets target request durations in seconds
137/// (`0.005` to `10`), so without an override almost every body would fall into
138/// the `+Inf` bucket. These byte-scaled bounds keep the body-size histograms
139/// meaningful.
140const BODY_SIZE_BUCKETS_BYTES: &[f64] = &[
141    64.0,
142    256.0,
143    1_024.0,
144    4_096.0,
145    16_384.0,
146    65_536.0,
147    262_144.0,
148    1_048_576.0,
149    10_485_760.0,
150];
151
152// -----------------------------------------------------------------------------
153// Recorder Installation
154// -----------------------------------------------------------------------------
155
156/// Global handle to the Prometheus exporter.
157static PROMETHEUS_HANDLE: OnceLock<PrometheusHandle> = OnceLock::new();
158
159/// Label dimensions to emit, installed once at startup.
160static LABEL_CONFIG: OnceLock<MetricLabelsConfig> = OnceLock::new();
161
162/// Every dimension enabled: the default, and the fallback before install.
163static ALL_LABELS: OnceLock<MetricLabelsConfig> = OnceLock::new();
164
165/// Install the label dimensions to emit.
166///
167/// Must be called once during startup, before any metric is recorded.
168/// Later calls are ignored: a gauge guard acquired before a change and
169/// released after it would increment one series and decrement another.
170pub fn install_metric_labels(labels: MetricLabelsConfig) {
171    let _existing = LABEL_CONFIG.set(labels);
172}
173
174/// The installed label dimensions, defaulting to all enabled.
175pub(crate) fn metric_labels() -> &'static MetricLabelsConfig {
176    LABEL_CONFIG
177        .get()
178        .unwrap_or_else(|| ALL_LABELS.get_or_init(MetricLabelsConfig::default))
179}
180
181/// Build a label set, dropping the dimensions that are disabled.
182///
183/// Only reached when at least one dimension is off; the all-enabled path
184/// uses the static-label macro form.
185fn selected_labels(pairs: &[(&'static str, Option<SharedString>)]) -> Vec<Label> {
186    pairs
187        .iter()
188        .filter_map(|(name, value)| value.clone().map(|value| Label::new(*name, value)))
189        .collect()
190}
191
192/// The value for a dimension, or `None` when that dimension is disabled.
193fn label_if(enabled: bool, value: SharedString) -> Option<SharedString> {
194    enabled.then_some(value)
195}
196
197/// Install the global Prometheus metrics recorder.
198///
199/// Must be called exactly once during server startup. Subsequent
200/// calls are no-ops and return the existing handle.
201///
202/// # Panics
203///
204/// Panics if the global recorder cannot be installed (another
205/// recorder was already set by a different subsystem).
206pub fn install_prometheus_recorder() -> &'static PrometheusHandle {
207    #[expect(
208        clippy::expect_used,
209        reason = "recorder installation is a one-time startup operation"
210    )]
211    PROMETHEUS_HANDLE.get_or_init(|| {
212        PrometheusBuilder::new()
213            .set_buckets_for_metric(
214                Matcher::Full(HTTP_REQUEST_BODY_BYTES.to_owned()),
215                BODY_SIZE_BUCKETS_BYTES,
216            )
217            .expect("body request histogram buckets must be non-empty")
218            .set_buckets_for_metric(
219                Matcher::Full(HTTP_RESPONSE_BODY_BYTES.to_owned()),
220                BODY_SIZE_BUCKETS_BYTES,
221            )
222            .expect("body response histogram buckets must be non-empty")
223            .install_recorder()
224            .expect("failed to install Prometheus recorder")
225    })
226}
227
228/// Render all collected metrics in Prometheus text exposition format.
229///
230/// Returns `None` if the recorder has not been installed.
231pub fn render_prometheus() -> Option<String> {
232    PROMETHEUS_HANDLE.get().map(PrometheusHandle::render)
233}
234
235/// Returns `true` if the Prometheus recorder has been installed.
236pub(crate) fn is_recorder_installed() -> bool {
237    PROMETHEUS_HANDLE.get().is_some()
238}
239
240// -----------------------------------------------------------------------------
241// Stats snapshot (`GET /api/stats`)
242// -----------------------------------------------------------------------------
243
244/// Parsed operational counters for [`collect_stats_metrics`].
245#[derive(Clone, Debug, Default, Eq, PartialEq)]
246pub struct StatsMetricsSnapshot {
247    /// In-flight HTTP requests per listener (`praxis_http_active_requests`).
248    pub http_active_by_listener: std::collections::HashMap<String, u64>,
249    /// Aggregate in-flight HTTP requests when the listener label is disabled.
250    pub http_active_aggregate: Option<u64>,
251    /// Open TCP sessions per listener (`praxis_tcp_active_connections`).
252    pub tcp_active_by_listener: std::collections::HashMap<String, u64>,
253    /// Aggregate open TCP sessions when the listener label is disabled.
254    pub tcp_active_aggregate: Option<u64>,
255    /// Upstream requests grouped by cluster (`praxis_upstream_requests_total`).
256    pub upstream_requests_by_cluster: std::collections::HashMap<String, u64>,
257    /// Aggregate upstream requests when the cluster label is disabled.
258    pub upstream_requests_aggregate: Option<u64>,
259    /// Upstream connect failures grouped by cluster.
260    pub connect_failures_by_cluster: std::collections::HashMap<String, u64>,
261    /// Aggregate upstream connect failures when the cluster label is disabled.
262    pub connect_failures_aggregate: Option<u64>,
263}
264
265/// Extract operational counters needed by `/api/stats` from Prometheus text.
266pub fn collect_stats_metrics(prometheus_text: &str) -> StatsMetricsSnapshot {
267    let mut snapshot = StatsMetricsSnapshot::default();
268    for line in prometheus_text.lines() {
269        let line = line.trim();
270        if line.is_empty() || line.starts_with('#') {
271            continue;
272        }
273        if let Some((name, labels, value)) = parse_prometheus_sample(line) {
274            match name {
275                HTTP_ACTIVE_REQUESTS => snapshot.record_http_active(&labels, value),
276                TCP_ACTIVE_CONNECTIONS => snapshot.record_tcp_active(&labels, value),
277                UPSTREAM_REQUESTS_TOTAL => snapshot.add_upstream_requests(&labels, value),
278                UPSTREAM_CONNECT_FAILURES_TOTAL => snapshot.add_connect_failures(&labels, value),
279                _ => {},
280            }
281        }
282    }
283    snapshot
284}
285
286impl StatsMetricsSnapshot {
287    /// Record an HTTP active-requests gauge sample.
288    fn record_http_active(&mut self, labels: &std::collections::HashMap<String, String>, value: u64) {
289        record_by_listener(
290            &mut self.http_active_by_listener,
291            &mut self.http_active_aggregate,
292            labels,
293            value,
294        );
295    }
296
297    /// Record a TCP active-connections gauge sample.
298    fn record_tcp_active(&mut self, labels: &std::collections::HashMap<String, String>, value: u64) {
299        record_by_listener(
300            &mut self.tcp_active_by_listener,
301            &mut self.tcp_active_aggregate,
302            labels,
303            value,
304        );
305    }
306
307    /// Accumulate an upstream-requests counter sample.
308    fn add_upstream_requests(&mut self, labels: &std::collections::HashMap<String, String>, value: u64) {
309        accumulate_by_cluster(
310            &mut self.upstream_requests_by_cluster,
311            &mut self.upstream_requests_aggregate,
312            labels,
313            value,
314        );
315    }
316
317    /// Accumulate an upstream connect-failures counter sample.
318    fn add_connect_failures(&mut self, labels: &std::collections::HashMap<String, String>, value: u64) {
319        accumulate_by_cluster(
320            &mut self.connect_failures_by_cluster,
321            &mut self.connect_failures_aggregate,
322            labels,
323            value,
324        );
325    }
326}
327
328/// Record a by-listener gauge sample: keyed by the `listener` label, or the
329/// aggregate when the label is absent.
330fn record_by_listener(
331    by_listener: &mut std::collections::HashMap<String, u64>,
332    aggregate: &mut Option<u64>,
333    labels: &std::collections::HashMap<String, String>,
334    value: u64,
335) {
336    match labels.get("listener") {
337        Some(listener) => {
338            by_listener.insert(listener.clone(), value);
339        },
340        None => *aggregate = Some(value),
341    }
342}
343
344/// Accumulate a by-cluster counter sample: added to the `cluster` label's
345/// running total, or to the aggregate when the label is absent.
346fn accumulate_by_cluster(
347    by_cluster: &mut std::collections::HashMap<String, u64>,
348    aggregate: &mut Option<u64>,
349    labels: &std::collections::HashMap<String, String>,
350    value: u64,
351) {
352    match labels.get("cluster") {
353        Some(cluster) => {
354            *by_cluster.entry(cluster.clone()).or_insert(0) += value;
355        },
356        None => *aggregate = Some(aggregate.unwrap_or(0) + value),
357    }
358}
359
360// -----------------------------------------------------------------------------
361// Prometheus Text Parsing
362// -----------------------------------------------------------------------------
363
364/// Parsed Prometheus sample: metric name, label map, and integer value.
365type PrometheusSample<'a> = (&'a str, std::collections::HashMap<String, String>, u64);
366
367/// Parse one Prometheus text sample line into metric name, labels, and value.
368fn parse_prometheus_sample(line: &str) -> Option<PrometheusSample<'_>> {
369    let (name_and_labels, value_str) = line.rsplit_once(' ')?;
370    #[expect(
371        clippy::cast_possible_truncation,
372        clippy::cast_sign_loss,
373        reason = "Prometheus counter/gauge values are non-negative integers"
374    )]
375    let value = {
376        let parsed = value_str.parse::<f64>().ok()?;
377        parsed.round() as u64
378    };
379    let (name, labels) = if let Some((name, label_blob)) = name_and_labels.split_once('{') {
380        let label_blob = label_blob.strip_suffix('}')?;
381        (name, parse_prometheus_labels(label_blob))
382    } else {
383        (name_and_labels, std::collections::HashMap::new())
384    };
385    Some((name, labels, value))
386}
387
388/// Parse `{key="value",...}` label sets from Prometheus text exposition.
389///
390/// Splits on commas only outside double-quoted values so embedded commas in
391/// label values parse correctly.
392fn parse_prometheus_labels(input: &str) -> std::collections::HashMap<String, String> {
393    let mut labels = std::collections::HashMap::new();
394    let mut rest = input.trim();
395    while !rest.is_empty() {
396        let (pair, tail) = split_prometheus_label_pair(rest);
397        if let Some((key, value)) = parse_prometheus_label_pair(pair) {
398            labels.insert(key, value);
399        }
400        rest = tail;
401    }
402    labels
403}
404
405/// Split one `key="value"` pair from the head of a label-set fragment.
406fn split_prometheus_label_pair(input: &str) -> (&str, &str) {
407    let bytes = input.as_bytes();
408    let mut in_quotes = false;
409    for (index, byte) in bytes.iter().enumerate() {
410        match *byte {
411            b'"' => in_quotes = !in_quotes,
412            b',' if !in_quotes => {
413                let head = input.get(..index).unwrap_or(input);
414                let tail = input.get(index + 1..).unwrap_or("").trim_start();
415                return (head, tail);
416            },
417            _ => {},
418        }
419    }
420    (input, "")
421}
422
423/// Parse one `key="value"` pair from a Prometheus label set fragment.
424fn parse_prometheus_label_pair(pair: &str) -> Option<(String, String)> {
425    let (key, value) = pair.split_once('=')?;
426    let value = value.strip_prefix('"')?.strip_suffix('"')?;
427    Some((key.to_owned(), value.to_owned()))
428}
429
430// -----------------------------------------------------------------------------
431// Status & Method Labels
432// -----------------------------------------------------------------------------
433
434/// Map an HTTP status code to its class label (`"1xx"`, `"2xx"`, etc.).
435///
436/// Returns `"unknown"` for zero (no response written) or codes
437/// outside the 100–599 range.
438///
439/// ```
440/// use praxis_protocol::http::pingora::metrics::status_class;
441///
442/// assert_eq!(status_class(200), "2xx");
443/// assert_eq!(status_class(404), "4xx");
444/// assert_eq!(status_class(0), "unknown");
445/// ```
446pub fn status_class(code: u16) -> &'static str {
447    match code {
448        100..=199 => "1xx",
449        200..=299 => "2xx",
450        300..=399 => "3xx",
451        400..=499 => "4xx",
452        500..=599 => "5xx",
453        _ => "unknown",
454    }
455}
456
457/// Map an HTTP method to a bounded label value.
458///
459/// Returns the method string for the nine standard methods
460/// defined in [RFC 9110]; all others collapse to `"OTHER"`.
461///
462/// ```
463/// use praxis_protocol::http::pingora::metrics::method_label;
464///
465/// assert_eq!(method_label("GET"), "GET");
466/// assert_eq!(method_label("PURGE"), "OTHER");
467/// ```
468///
469/// [RFC 9110]: https://datatracker.ietf.org/doc/html/rfc9110#section-9.1
470pub fn method_label(method: &str) -> &'static str {
471    match method {
472        "GET" => "GET",
473        "POST" => "POST",
474        "PUT" => "PUT",
475        "DELETE" => "DELETE",
476        "PATCH" => "PATCH",
477        "HEAD" => "HEAD",
478        "OPTIONS" => "OPTIONS",
479        "TRACE" => "TRACE",
480        "CONNECT" => "CONNECT",
481        _ => "OTHER",
482    }
483}
484
485// -----------------------------------------------------------------------------
486// Request Metrics
487// -----------------------------------------------------------------------------
488
489/// Labels for a completed HTTP request.
490///
491/// Static labels (`method`, `status_class`) use `&'static str`;
492/// `cluster` and `route` are dynamic [`SharedString`] values.
493///
494/// [`SharedString`]: ::metrics::SharedString
495pub(crate) struct RequestMetricLabels {
496    /// Cluster name or `"none"`.
497    pub cluster: SharedString,
498    /// HTTP method (e.g. `"GET"`).
499    pub method: &'static str,
500    /// Route path-match pattern or `"unknown"`.
501    pub route: SharedString,
502    /// Status class (e.g. `"2xx"`).
503    pub status_class: &'static str,
504}
505
506/// Build the enabled subset of the request-metric labels.
507fn selected_request_labels(labels: RequestMetricLabels) -> Vec<Label> {
508    let selected = metric_labels();
509    let pairs = [
510        (
511            "method",
512            label_if(
513                selected.is_enabled(MetricLabel::Method),
514                SharedString::const_str(labels.method),
515            ),
516        ),
517        (
518            "status_class",
519            label_if(
520                selected.is_enabled(MetricLabel::StatusClass),
521                SharedString::const_str(labels.status_class),
522            ),
523        ),
524        ("route", label_if(selected.is_enabled(MetricLabel::Route), labels.route)),
525        (
526            "cluster",
527            label_if(selected.is_enabled(MetricLabel::Cluster), labels.cluster),
528        ),
529    ];
530    selected_labels(&pairs)
531}
532
533/// Record HTTP request metrics for a completed request.
534pub(crate) fn record_request_metrics(labels: RequestMetricLabels, duration_secs: f64) {
535    if !is_recorder_installed() {
536        return;
537    }
538    if !metric_labels().all_enabled() {
539        let emitted = selected_request_labels(labels);
540        counter!(HTTP_REQUESTS_TOTAL, emitted.clone()).increment(1);
541        histogram!(HTTP_REQUEST_DURATION_SECONDS, emitted).record(duration_secs);
542        return;
543    }
544    record_request_metrics_all_labels(labels, duration_secs);
545}
546
547/// Record request metrics with the full default label set.
548///
549/// Kept on the static-label macro form so the default configuration emits
550/// exactly the series it did before label selection existed.
551fn record_request_metrics_all_labels(labels: RequestMetricLabels, duration_secs: f64) {
552    let cluster = labels.cluster;
553    let route = labels.route;
554    counter!(
555        HTTP_REQUESTS_TOTAL,
556        "method" => labels.method,
557        "status_class" => labels.status_class,
558        "route" => route.clone(),
559        "cluster" => cluster.clone()
560    )
561    .increment(1);
562    histogram!(
563        HTTP_REQUEST_DURATION_SECONDS,
564        "method" => labels.method,
565        "status_class" => labels.status_class,
566        "route" => route,
567        "cluster" => cluster
568    )
569    .record(duration_secs);
570}
571
572// -----------------------------------------------------------------------------
573// Body-Size Metrics
574// -----------------------------------------------------------------------------
575
576/// Build the enabled subset of the body-size histogram labels.
577fn selected_body_labels(method: &'static str, status_class: &'static str, cluster: SharedString) -> Vec<Label> {
578    let selected = metric_labels();
579    let pairs = [
580        (
581            "method",
582            label_if(
583                selected.is_enabled(MetricLabel::Method),
584                SharedString::const_str(method),
585            ),
586        ),
587        (
588            "status_class",
589            label_if(
590                selected.is_enabled(MetricLabel::StatusClass),
591                SharedString::const_str(status_class),
592            ),
593        ),
594        ("cluster", label_if(selected.is_enabled(MetricLabel::Cluster), cluster)),
595    ];
596    selected_labels(&pairs)
597}
598
599/// Record body-size histograms with the full default label set.
600fn record_body_size_all_labels(
601    method: &'static str,
602    status_class: &'static str,
603    cluster: SharedString,
604    request_bytes: f64,
605    response_bytes: f64,
606) {
607    histogram!(
608        HTTP_REQUEST_BODY_BYTES,
609        "method" => method,
610        "status_class" => status_class,
611        "cluster" => cluster.clone()
612    )
613    .record(request_bytes);
614    histogram!(
615        HTTP_RESPONSE_BODY_BYTES,
616        "method" => method,
617        "status_class" => status_class,
618        "cluster" => cluster
619    )
620    .record(response_bytes);
621}
622
623/// Record HTTP request and response body size histograms.
624pub(crate) fn record_body_size_metrics(
625    method: &'static str,
626    status_class: &'static str,
627    cluster: SharedString,
628    request_body_bytes: u64,
629    response_body_bytes: u64,
630) {
631    if !is_recorder_installed() {
632        return;
633    }
634    #[expect(
635        clippy::cast_precision_loss,
636        reason = "body byte counts as histogram observations; exact integer precision not required"
637    )]
638    let (request_bytes, response_bytes) = (request_body_bytes as f64, response_body_bytes as f64);
639
640    if !metric_labels().all_enabled() {
641        let emitted = selected_body_labels(method, status_class, cluster);
642        histogram!(HTTP_REQUEST_BODY_BYTES, emitted.clone()).record(request_bytes);
643        histogram!(HTTP_RESPONSE_BODY_BYTES, emitted).record(response_bytes);
644        return;
645    }
646    record_body_size_all_labels(method, status_class, cluster, request_bytes, response_bytes);
647}
648
649// -----------------------------------------------------------------------------
650// Active Request Tracking
651// -----------------------------------------------------------------------------
652
653/// RAII guard that decrements `praxis_http_active_requests` on drop.
654///
655/// Acquired once per HTTP request. Pingora owns the request context by
656/// value, so the drop runs on every terminal path (including a client
657/// abort or an HTTP/2 stream reset, which skip the `logging` callback
658/// entirely).
659pub struct ActiveRequestGuard {
660    /// Listener name label.
661    listener: SharedString,
662}
663
664impl ActiveRequestGuard {
665    /// Increment the gauge and return a guard that decrements on drop.
666    pub(crate) fn acquire(listener: SharedString) -> Self {
667        if is_recorder_installed() {
668            if metric_labels().is_enabled(MetricLabel::Listener) {
669                gauge!(HTTP_ACTIVE_REQUESTS, "listener" => listener.clone()).increment(1.0);
670            } else {
671                gauge!(HTTP_ACTIVE_REQUESTS).increment(1.0);
672            }
673        }
674        Self { listener }
675    }
676}
677
678impl Drop for ActiveRequestGuard {
679    fn drop(&mut self) {
680        if is_recorder_installed() {
681            if metric_labels().is_enabled(MetricLabel::Listener) {
682                gauge!(HTTP_ACTIVE_REQUESTS, "listener" => self.listener.clone()).decrement(1.0);
683            } else {
684                gauge!(HTTP_ACTIVE_REQUESTS).decrement(1.0);
685            }
686        }
687    }
688}
689
690// -----------------------------------------------------------------------------
691// Error Classification
692// -----------------------------------------------------------------------------
693
694/// Record a proxy error.
695///
696/// Counted once per request, from the logging hook. Overload rejections and
697/// upstream connect failures have dedicated counters and are not repeated
698/// here; this counter covers the causes those miss.
699pub(crate) fn record_error(error_type: &'static str) {
700    if !is_recorder_installed() {
701        return;
702    }
703    counter!(ERRORS_TOTAL, "type" => error_type).increment(1);
704}
705
706/// Classify a Pingora error into a bounded `type` label value.
707///
708/// Connect failures map to `upstream_unavailable` and are also counted by
709/// `praxis_upstream_connect_failures_total`; the overlap is deliberate so
710/// that `praxis_errors_total` is a complete error denominator on its own.
711pub(crate) fn error_type_for(etype: &::pingora_core::ErrorType, source: &::pingora_core::ErrorSource) -> &'static str {
712    use ::pingora_core::ErrorSource::{Downstream, Internal, Unset};
713
714    if matches!(source, Downstream) {
715        return ERROR_TYPE_DOWNSTREAM;
716    }
717    if is_timeout(etype) {
718        return ERROR_TYPE_TIMEOUT;
719    }
720    if is_unreachable(etype) {
721        return ERROR_TYPE_UPSTREAM_UNAVAILABLE;
722    }
723    if matches!(source, Internal | Unset) {
724        return ERROR_TYPE_INTERNAL;
725    }
726    ERROR_TYPE_UPSTREAM_PROTOCOL
727}
728
729/// Whether the error is a connect, handshake, read or write timeout.
730fn is_timeout(etype: &::pingora_core::ErrorType) -> bool {
731    use ::pingora_core::ErrorType::{ConnectTimedout, ReadTimedout, TLSHandshakeTimedout, WriteTimedout};
732
733    matches!(
734        etype,
735        ConnectTimedout | TLSHandshakeTimedout | ReadTimedout | WriteTimedout
736    )
737}
738
739/// Whether the error means the upstream was never reached.
740fn is_unreachable(etype: &::pingora_core::ErrorType) -> bool {
741    use ::pingora_core::ErrorType::{BindError, ConnectError, ConnectNoRoute, ConnectRefused, SocketError};
742
743    matches!(
744        etype,
745        ConnectRefused | ConnectNoRoute | ConnectError | BindError | SocketError
746    )
747}
748
749// -----------------------------------------------------------------------------
750// Overload Rejections
751// -----------------------------------------------------------------------------
752
753/// Record an overload rejection.
754pub(crate) fn record_overload_reject(reason: &'static str) {
755    if !is_recorder_installed() {
756        return;
757    }
758    counter!(OVERLOAD_REJECTS_TOTAL, "reason" => reason).increment(1);
759}
760
761// -----------------------------------------------------------------------------
762// Upstream Metrics
763// -----------------------------------------------------------------------------
764
765/// Record upstream connect duration for a cluster.
766pub(crate) fn record_upstream_connect_duration(cluster: SharedString, duration_secs: f64) {
767    if !is_recorder_installed() {
768        return;
769    }
770    if metric_labels().is_enabled(MetricLabel::Cluster) {
771        histogram!(UPSTREAM_CONNECT_DURATION_SECONDS, "cluster" => cluster).record(duration_secs);
772    } else {
773        histogram!(UPSTREAM_CONNECT_DURATION_SECONDS).record(duration_secs);
774    }
775}
776
777/// Record a request that reached an upstream endpoint.
778///
779/// Counted once per request, from the logging hook, so a request retried
780/// across endpoints increments once against the endpoint that answered
781/// rather than once per attempt. Requests that never reached an upstream
782/// (filter rejections, connect failures) are not counted here.
783pub(crate) fn record_upstream_request(cluster: SharedString, endpoint: SharedString, status_class: &'static str) {
784    if !is_recorder_installed() {
785        return;
786    }
787    let selected = metric_labels();
788    if !selected.all_enabled() {
789        let pairs = [
790            ("cluster", label_if(selected.is_enabled(MetricLabel::Cluster), cluster)),
791            (
792                "endpoint",
793                label_if(selected.is_enabled(MetricLabel::Endpoint), endpoint),
794            ),
795            (
796                "status_class",
797                label_if(
798                    selected.is_enabled(MetricLabel::StatusClass),
799                    SharedString::const_str(status_class),
800                ),
801            ),
802        ];
803        counter!(UPSTREAM_REQUESTS_TOTAL, selected_labels(&pairs)).increment(1);
804        return;
805    }
806    counter!(
807        UPSTREAM_REQUESTS_TOTAL,
808        "cluster" => cluster,
809        "endpoint" => endpoint,
810        "status_class" => status_class
811    )
812    .increment(1);
813}
814
815/// Record an upstream connect failure.
816pub(crate) fn record_upstream_connect_failure(cluster: SharedString) {
817    if !is_recorder_installed() {
818        return;
819    }
820    if metric_labels().is_enabled(MetricLabel::Cluster) {
821        counter!(UPSTREAM_CONNECT_FAILURES_TOTAL, "cluster" => cluster).increment(1);
822    } else {
823        counter!(UPSTREAM_CONNECT_FAILURES_TOTAL).increment(1);
824    }
825}
826
827/// Record an upstream connect-failure retry outcome.
828pub(crate) fn record_upstream_retry(cluster: SharedString, result: &'static str) {
829    if !is_recorder_installed() {
830        return;
831    }
832    if metric_labels().is_enabled(MetricLabel::Cluster) {
833        counter!(UPSTREAM_RETRIES_TOTAL, "cluster" => cluster, "result" => result).increment(1);
834    } else {
835        counter!(UPSTREAM_RETRIES_TOTAL, "result" => result).increment(1);
836    }
837}
838
839// -----------------------------------------------------------------------------
840// Upstream Health
841// -----------------------------------------------------------------------------
842
843/// Refresh cluster endpoint health gauges.
844///
845/// The `cluster` label is structural here and is kept even when the
846/// `cluster` dimension is disabled: these gauges are keyed by cluster and
847/// set (not incremented), so dropping the label would collapse every
848/// cluster onto one last-writer-wins series rather than merely lowering
849/// cardinality. Disabling `cluster` therefore drops it from the additive
850/// cluster metrics (counters and the connect-duration histogram) but not
851/// from the per-cluster health gauges.
852pub(crate) fn set_upstream_endpoint_gauges(cluster: SharedString, healthy: usize, total: usize) {
853    if !is_recorder_installed() {
854        return;
855    }
856    #[expect(clippy::cast_precision_loss, reason = "endpoint counts fit f64 exactly below 2^53")]
857    {
858        gauge!(UPSTREAM_HEALTHY_ENDPOINTS, "cluster" => cluster.clone()).set(healthy as f64);
859        gauge!(UPSTREAM_TOTAL_ENDPOINTS, "cluster" => cluster).set(total as f64);
860    }
861}
862
863/// Zero health gauges for clusters that lost active health checks on reload.
864///
865/// Prometheus does not drop series automatically; clearing removed clusters
866/// prevents stale `healthy`/`total` values from lingering after config change.
867pub fn clear_stale_upstream_health_gauges<'a, P: IntoIterator<Item = &'a str>, C: IntoIterator<Item = &'a str>>(
868    previous_health_clusters: P,
869    current_health_clusters: C,
870) {
871    if !is_recorder_installed() {
872        return;
873    }
874    let current: std::collections::HashSet<&str> = current_health_clusters.into_iter().collect();
875    for name in previous_health_clusters {
876        if !current.contains(name) {
877            set_upstream_endpoint_gauges(SharedString::from(name.to_owned()), 0, 0);
878        }
879    }
880}
881
882/// Publish current healthy/total gauges for every cluster in a health registry.
883///
884/// Called on reload so scrapes reflect the new registry immediately instead of
885/// waiting for the first probe round.
886pub fn seed_upstream_health_gauges(registry: &praxis_core::health::HealthRegistry) {
887    if !is_recorder_installed() {
888        return;
889    }
890    for (name, state) in registry.iter() {
891        let (healthy, total) = state.endpoint_counts();
892        set_upstream_endpoint_gauges(SharedString::from(name.as_ref().to_owned()), healthy, total);
893    }
894}
895
896/// Record an endpoint health state transition and refresh gauges.
897pub(crate) fn record_health_transition(cluster: SharedString, result: &'static str, healthy: usize, total: usize) {
898    if !is_recorder_installed() {
899        return;
900    }
901    if metric_labels().is_enabled(MetricLabel::Cluster) {
902        counter!(
903            UPSTREAM_HEALTH_TRANSITIONS_TOTAL,
904            "cluster" => cluster.clone(),
905            "result" => result
906        )
907        .increment(1);
908    } else {
909        counter!(UPSTREAM_HEALTH_TRANSITIONS_TOTAL, "result" => result).increment(1);
910    }
911    set_upstream_endpoint_gauges(cluster, healthy, total);
912}
913
914/// Count healthy endpoints in a cluster health entry.
915pub(crate) fn count_healthy_endpoints(health: &praxis_core::health::ClusterHealthEntry) -> (usize, usize) {
916    health.endpoint_counts()
917}
918
919// -----------------------------------------------------------------------------
920// Config Reload
921// -----------------------------------------------------------------------------
922
923/// Record a successful config reload.
924pub fn record_config_reload_success() {
925    if !is_recorder_installed() {
926        return;
927    }
928    counter!(CONFIG_RELOAD_TOTAL, "result" => RELOAD_RESULT_SUCCESS).increment(1);
929    let ts = std::time::SystemTime::now()
930        .duration_since(std::time::UNIX_EPOCH)
931        .map_or(0.0, |d| d.as_secs_f64());
932    gauge!(CONFIG_RELOAD_LAST_SUCCESS_TIMESTAMP).set(ts);
933}
934
935/// Record a failed config reload.
936pub fn record_config_reload_failure() {
937    if !is_recorder_installed() {
938        return;
939    }
940    counter!(CONFIG_RELOAD_TOTAL, "result" => RELOAD_RESULT_FAILURE).increment(1);
941}
942
943// -----------------------------------------------------------------------------
944// Shared Label Values
945// -----------------------------------------------------------------------------
946
947/// [`SharedString`] for the `"none"` cluster label.
948///
949/// [`SharedString`]: ::metrics::SharedString
950pub(crate) fn cluster_none() -> SharedString {
951    SharedString::const_str("none")
952}
953
954/// [`SharedString`] for the `"unknown"` route label.
955///
956/// [`SharedString`]: ::metrics::SharedString
957pub(crate) fn route_unknown() -> SharedString {
958    SharedString::const_str("unknown")
959}
960
961// -----------------------------------------------------------------------------
962// Tests
963// -----------------------------------------------------------------------------
964
965#[cfg(test)]
966#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
967#[allow(clippy::unwrap_used, clippy::expect_used, clippy::indexing_slicing, reason = "tests")]
968mod tests {
969    use super::*;
970
971    #[test]
972    fn status_class_1xx() {
973        assert_eq!(status_class(100), "1xx", "100 should be 1xx");
974        assert_eq!(status_class(199), "1xx", "199 should be 1xx");
975    }
976
977    #[test]
978    fn status_class_2xx() {
979        assert_eq!(status_class(200), "2xx", "200 should be 2xx");
980        assert_eq!(status_class(204), "2xx", "204 should be 2xx");
981        assert_eq!(status_class(299), "2xx", "299 should be 2xx");
982    }
983
984    #[test]
985    fn status_class_3xx() {
986        assert_eq!(status_class(301), "3xx", "301 should be 3xx");
987        assert_eq!(status_class(399), "3xx", "399 should be 3xx");
988    }
989
990    #[test]
991    fn status_class_4xx() {
992        assert_eq!(status_class(400), "4xx", "400 should be 4xx");
993        assert_eq!(status_class(404), "4xx", "404 should be 4xx");
994        assert_eq!(status_class(499), "4xx", "499 should be 4xx");
995    }
996
997    #[test]
998    fn status_class_5xx() {
999        assert_eq!(status_class(500), "5xx", "500 should be 5xx");
1000        assert_eq!(status_class(503), "5xx", "503 should be 5xx");
1001        assert_eq!(status_class(599), "5xx", "599 should be 5xx");
1002    }
1003
1004    #[test]
1005    fn status_class_zero_is_unknown() {
1006        assert_eq!(status_class(0), "unknown", "0 should be unknown");
1007    }
1008
1009    #[test]
1010    fn status_class_out_of_range_is_unknown() {
1011        assert_eq!(status_class(600), "unknown", "600 should be unknown");
1012        assert_eq!(status_class(99), "unknown", "99 should be unknown");
1013    }
1014
1015    #[test]
1016    fn method_label_standard_methods() {
1017        for m in [
1018            "GET", "POST", "PUT", "DELETE", "PATCH", "HEAD", "OPTIONS", "TRACE", "CONNECT",
1019        ] {
1020            assert_eq!(method_label(m), m, "{m} should pass through");
1021        }
1022    }
1023
1024    #[test]
1025    fn method_label_custom_methods_collapse_to_other() {
1026        assert_eq!(method_label("PURGE"), "OTHER", "PURGE should be OTHER");
1027        assert_eq!(method_label("FOOBAR"), "OTHER", "FOOBAR should be OTHER");
1028        assert_eq!(method_label(""), "OTHER", "empty should be OTHER");
1029    }
1030
1031    #[test]
1032    fn record_utilities_noop_without_recorder() {
1033        record_overload_reject(OVERLOAD_REASON_MEMORY);
1034        record_upstream_connect_failure(cluster_none());
1035        record_error(ERROR_TYPE_INTERNAL);
1036        record_upstream_request(cluster_none(), SharedString::const_str("10.0.0.1:80"), "2xx");
1037        record_upstream_retry(cluster_none(), RETRY_RESULT_SUCCESS);
1038        record_upstream_connect_duration(cluster_none(), 0.01);
1039        set_upstream_endpoint_gauges(cluster_none(), 1, 2);
1040        record_health_transition(cluster_none(), HEALTH_RESULT_HEALTHY, 1, 2);
1041        record_config_reload_success();
1042        record_config_reload_failure();
1043        clear_stale_upstream_health_gauges(["gone"], std::iter::empty::<&str>());
1044        let _request_guard = ActiveRequestGuard::acquire(SharedString::const_str("test"));
1045    }
1046
1047    #[test]
1048    fn active_request_guard_returns_to_zero_on_drop() {
1049        install_prometheus_recorder();
1050        let listener = SharedString::const_str("active-request-guard-listener");
1051        let guard = ActiveRequestGuard::acquire(listener.clone());
1052        let held = render_prometheus().expect("recorder should render");
1053        assert!(
1054            held.contains("praxis_http_active_requests{listener=\"active-request-guard-listener\"} 1"),
1055            "gauge should read 1 while the guard is held:\n{held}"
1056        );
1057        drop(guard);
1058        let released = render_prometheus().expect("recorder should render");
1059        assert!(
1060            released.contains("praxis_http_active_requests{listener=\"active-request-guard-listener\"} 0"),
1061            "gauge should return to 0 once the guard drops:\n{released}"
1062        );
1063    }
1064
1065    #[test]
1066    fn overload_reject_reasons_appear_in_scrape() {
1067        install_prometheus_recorder();
1068        record_overload_reject(OVERLOAD_REASON_MEMORY);
1069        record_overload_reject(OVERLOAD_REASON_GLOBAL_CONNECTIONS);
1070        record_overload_reject(OVERLOAD_REASON_LISTENER_CONNECTIONS);
1071        let body = render_prometheus().expect("recorder should render");
1072        for reason in [
1073            OVERLOAD_REASON_MEMORY,
1074            OVERLOAD_REASON_GLOBAL_CONNECTIONS,
1075            OVERLOAD_REASON_LISTENER_CONNECTIONS,
1076        ] {
1077            let needle = format!("praxis_overload_rejects_total{{reason=\"{reason}\"}}");
1078            assert!(body.contains(&needle), "expected `{needle}` in scrape:\n{body}");
1079        }
1080    }
1081
1082    #[test]
1083    fn upstream_requests_carry_cluster_endpoint_and_status_class() {
1084        install_prometheus_recorder();
1085        record_upstream_request(
1086            SharedString::const_str("api"),
1087            SharedString::const_str("10.0.0.7:8080"),
1088            "5xx",
1089        );
1090        let body = render_prometheus().expect("recorder should render");
1091        assert!(
1092            body.contains(
1093                "praxis_upstream_requests_total{cluster=\"api\",endpoint=\"10.0.0.7:8080\",status_class=\"5xx\"} 1"
1094            ),
1095            "counter should carry all three labels:\n{body}"
1096        );
1097    }
1098
1099    #[test]
1100    fn selected_labels_drops_disabled_dimensions() {
1101        let pairs = [
1102            ("method", Some(SharedString::const_str("GET"))),
1103            ("route", None),
1104            ("cluster", Some(SharedString::const_str("api"))),
1105        ];
1106        let emitted = selected_labels(&pairs);
1107        let names: Vec<&str> = emitted.iter().map(Label::key).collect();
1108        assert_eq!(names, vec!["method", "cluster"], "a disabled dimension must be absent");
1109    }
1110
1111    #[test]
1112    fn selected_labels_preserves_order_and_values() {
1113        let pairs = [
1114            ("cluster", Some(SharedString::const_str("api"))),
1115            ("endpoint", Some(SharedString::const_str("10.0.0.1:80"))),
1116        ];
1117        let emitted = selected_labels(&pairs);
1118        let rendered: Vec<(&str, &str)> = emitted.iter().map(|l| (l.key(), l.value())).collect();
1119        assert_eq!(
1120            rendered,
1121            vec![("cluster", "api"), ("endpoint", "10.0.0.1:80")],
1122            "enabled dimensions keep their order and values"
1123        );
1124    }
1125
1126    #[test]
1127    fn label_if_gates_on_the_flag() {
1128        assert_eq!(
1129            label_if(true, SharedString::const_str("x")).as_deref(),
1130            Some("x"),
1131            "an enabled dimension keeps its value"
1132        );
1133        assert_eq!(
1134            label_if(false, SharedString::const_str("x")),
1135            None,
1136            "a disabled dimension yields no value"
1137        );
1138    }
1139
1140    #[test]
1141    fn metric_labels_default_to_all_enabled() {
1142        assert!(
1143            metric_labels().all_enabled(),
1144            "without an explicit install every dimension must stay on, so the \
1145             recorders keep their allocation-free fast path"
1146        );
1147    }
1148
1149    #[test]
1150    fn error_types_appear_in_scrape() {
1151        install_prometheus_recorder();
1152        for error_type in [
1153            ERROR_TYPE_FILTER_REJECT,
1154            ERROR_TYPE_TIMEOUT,
1155            ERROR_TYPE_UPSTREAM_UNAVAILABLE,
1156            ERROR_TYPE_UPSTREAM_PROTOCOL,
1157            ERROR_TYPE_DOWNSTREAM,
1158            ERROR_TYPE_INTERNAL,
1159        ] {
1160            record_error(error_type);
1161            let body = render_prometheus().expect("recorder should render");
1162            let needle = format!("praxis_errors_total{{type=\"{error_type}\"}}");
1163            assert!(body.contains(&needle), "expected `{needle}` in scrape:\n{body}");
1164        }
1165    }
1166
1167    #[test]
1168    fn error_type_for_maps_pingora_errors_to_bounded_values() {
1169        use ::pingora_core::{ErrorSource, ErrorType};
1170
1171        assert_eq!(
1172            error_type_for(&ErrorType::ConnectTimedout, &ErrorSource::Upstream),
1173            ERROR_TYPE_TIMEOUT,
1174            "connect timeout is a timeout"
1175        );
1176        assert_eq!(
1177            error_type_for(&ErrorType::ConnectRefused, &ErrorSource::Upstream),
1178            ERROR_TYPE_UPSTREAM_UNAVAILABLE,
1179            "a refused connect means the upstream was unreachable"
1180        );
1181        assert_eq!(
1182            error_type_for(&ErrorType::ReadError, &ErrorSource::Upstream),
1183            ERROR_TYPE_UPSTREAM_PROTOCOL,
1184            "a mid-exchange read error is a protocol failure"
1185        );
1186        assert_eq!(
1187            error_type_for(&ErrorType::ReadTimedout, &ErrorSource::Downstream),
1188            ERROR_TYPE_DOWNSTREAM,
1189            "downstream source wins over the error kind"
1190        );
1191        assert_eq!(
1192            error_type_for(&ErrorType::InternalError, &ErrorSource::Internal),
1193            ERROR_TYPE_INTERNAL,
1194            "internal source is an internal fault"
1195        );
1196    }
1197
1198    #[test]
1199    fn body_size_histograms_use_byte_buckets() {
1200        install_prometheus_recorder();
1201        record_body_size_metrics("GET", "2xx", cluster_none(), 500, 4_000);
1202        let body = render_prometheus().expect("recorder should render");
1203        assert!(
1204            body.contains("praxis_http_request_body_bytes_bucket") && body.contains("le=\"1024\""),
1205            "request body histogram should use byte buckets, not duration defaults:\n{body}"
1206        );
1207        assert!(
1208            body.contains("praxis_http_response_body_bytes_bucket") && body.contains("le=\"4096\""),
1209            "response body histogram should use byte buckets:\n{body}"
1210        );
1211        assert!(
1212            !body.contains("praxis_http_request_body_bytes_bucket{le=\"0.005\"}")
1213                && !body.contains("praxis_http_request_body_bytes_bucket{method=\"GET\",status_class=\"2xx\",cluster=\"\",le=\"0.005\"}"),
1214            "request body histogram must not use duration default buckets:\n{body}"
1215        );
1216    }
1217
1218    #[test]
1219    fn clear_stale_upstream_health_gauges_zeros_removed_clusters() {
1220        install_prometheus_recorder();
1221        set_upstream_endpoint_gauges(SharedString::from("old-cluster".to_owned()), 2, 3);
1222        set_upstream_endpoint_gauges(SharedString::from("kept-cluster".to_owned()), 1, 1);
1223        clear_stale_upstream_health_gauges(["old-cluster", "kept-cluster"], ["kept-cluster"]);
1224        let body = render_prometheus().expect("recorder should render");
1225        assert!(
1226            body.contains("praxis_upstream_healthy_endpoints{cluster=\"old-cluster\"} 0"),
1227            "removed cluster healthy gauge should be zeroed:\n{body}"
1228        );
1229        assert!(
1230            body.contains("praxis_upstream_total_endpoints{cluster=\"old-cluster\"} 0"),
1231            "removed cluster total gauge should be zeroed:\n{body}"
1232        );
1233        assert!(
1234            body.contains("praxis_upstream_healthy_endpoints{cluster=\"kept-cluster\"} 1"),
1235            "kept cluster should retain its value:\n{body}"
1236        );
1237    }
1238
1239    #[test]
1240    fn collect_stats_metrics_sums_cluster_counters() {
1241        let text = r#"
1242praxis_http_active_requests{listener="web"} 2
1243praxis_tcp_active_connections{listener="tcp-in"} 1
1244praxis_upstream_requests_total{cluster="backend",endpoint="127.0.0.1:1",status_class="2xx"} 3
1245praxis_upstream_requests_total{cluster="backend",endpoint="127.0.0.1:2",status_class="5xx"} 1
1246praxis_upstream_connect_failures_total{cluster="backend"} 2
1247"#;
1248        let snap = collect_stats_metrics(text);
1249        assert_eq!(
1250            snap.http_active_by_listener.get("web"),
1251            Some(&2),
1252            "HTTP active per listener should parse"
1253        );
1254        assert_eq!(
1255            snap.tcp_active_by_listener.get("tcp-in"),
1256            Some(&1),
1257            "TCP active per listener should parse"
1258        );
1259        assert_eq!(
1260            snap.upstream_requests_by_cluster.get("backend"),
1261            Some(&4),
1262            "upstream requests should sum by cluster"
1263        );
1264        assert_eq!(
1265            snap.connect_failures_by_cluster.get("backend"),
1266            Some(&2),
1267            "connect failures should parse by cluster"
1268        );
1269    }
1270
1271    #[test]
1272    fn parse_prometheus_labels_handles_commas_inside_quoted_values() {
1273        let labels = parse_prometheus_labels(r#"tag="a,b",listener="web""#);
1274        assert_eq!(labels.get("tag"), Some(&"a,b".to_owned()), "comma inside quotes");
1275        assert_eq!(labels.get("listener"), Some(&"web".to_owned()), "second label");
1276    }
1277
1278    #[test]
1279    fn collect_stats_metrics_parses_unlabeled_upstream_counters() {
1280        let text = r#"
1281praxis_upstream_requests_total{status_class="2xx"} 5
1282praxis_upstream_connect_failures_total 2
1283"#;
1284        let snap = collect_stats_metrics(text);
1285        assert_eq!(
1286            snap.upstream_requests_aggregate,
1287            Some(5),
1288            "unlabeled upstream requests should aggregate"
1289        );
1290        assert_eq!(
1291            snap.connect_failures_aggregate,
1292            Some(2),
1293            "unlabeled connect failures should aggregate"
1294        );
1295    }
1296
1297    #[test]
1298    fn seed_upstream_health_gauges_publishes_registry_counts() {
1299        use std::sync::Arc;
1300
1301        use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
1302
1303        install_prometheus_recorder();
1304        let endpoints = vec![EndpointHealth::new(), EndpointHealth::new()];
1305        endpoints[0].mark_unhealthy();
1306        let entry = Arc::new(ClusterHealthEntry::new(
1307            endpoints,
1308            vec![Arc::from("a:1"), Arc::from("b:1")],
1309            None,
1310            None,
1311        ));
1312        let registry = Arc::new([(Arc::from("backend"), entry)].into_iter().collect());
1313        seed_upstream_health_gauges(&registry);
1314        let body = render_prometheus().expect("recorder should render");
1315        assert!(
1316            body.contains("praxis_upstream_healthy_endpoints{cluster=\"backend\"} 1"),
1317            "seed should publish healthy count:\n{body}"
1318        );
1319        assert!(
1320            body.contains("praxis_upstream_total_endpoints{cluster=\"backend\"} 2"),
1321            "seed should publish total count:\n{body}"
1322        );
1323    }
1324
1325    #[test]
1326    fn count_healthy_endpoints_counts_correctly() {
1327        use std::sync::Arc;
1328
1329        use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
1330
1331        let endpoints = vec![EndpointHealth::new(), EndpointHealth::new(), EndpointHealth::new()];
1332        endpoints[1].mark_unhealthy();
1333        let entry = ClusterHealthEntry::new(
1334            endpoints,
1335            vec![Arc::from("a:1"), Arc::from("b:1"), Arc::from("c:1")],
1336            None,
1337            None,
1338        );
1339        let (healthy, total) = count_healthy_endpoints(&entry);
1340        assert_eq!(total, 3, "total should be 3");
1341        assert_eq!(healthy, 2, "two endpoints should be healthy");
1342    }
1343}