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.
266#[expect(clippy::too_many_lines, reason = "metric name dispatch table")]
267pub fn collect_stats_metrics(prometheus_text: &str) -> StatsMetricsSnapshot {
268    let mut snapshot = StatsMetricsSnapshot::default();
269    for line in prometheus_text.lines() {
270        let line = line.trim();
271        if line.is_empty() || line.starts_with('#') {
272            continue;
273        }
274        let Some((name, labels, value)) = parse_prometheus_sample(line) else {
275            continue;
276        };
277        match name {
278            HTTP_ACTIVE_REQUESTS => {
279                if let Some(listener) = labels.get("listener") {
280                    snapshot.http_active_by_listener.insert(listener.clone(), value);
281                } else {
282                    snapshot.http_active_aggregate = Some(value);
283                }
284            },
285            TCP_ACTIVE_CONNECTIONS => {
286                if let Some(listener) = labels.get("listener") {
287                    snapshot.tcp_active_by_listener.insert(listener.clone(), value);
288                } else {
289                    snapshot.tcp_active_aggregate = Some(value);
290                }
291            },
292            UPSTREAM_REQUESTS_TOTAL => {
293                if let Some(cluster) = labels.get("cluster") {
294                    *snapshot
295                        .upstream_requests_by_cluster
296                        .entry(cluster.clone())
297                        .or_insert(0) += value;
298                } else {
299                    snapshot.upstream_requests_aggregate =
300                        Some(snapshot.upstream_requests_aggregate.unwrap_or(0) + value);
301                }
302            },
303            UPSTREAM_CONNECT_FAILURES_TOTAL => {
304                if let Some(cluster) = labels.get("cluster") {
305                    *snapshot.connect_failures_by_cluster.entry(cluster.clone()).or_insert(0) += value;
306                } else {
307                    snapshot.connect_failures_aggregate =
308                        Some(snapshot.connect_failures_aggregate.unwrap_or(0) + value);
309                }
310            },
311            _ => {},
312        }
313    }
314    snapshot
315}
316
317// -----------------------------------------------------------------------------
318// Prometheus Text Parsing
319// -----------------------------------------------------------------------------
320
321/// Parsed Prometheus sample: metric name, label map, and integer value.
322type PrometheusSample<'a> = (&'a str, std::collections::HashMap<String, String>, u64);
323
324/// Parse one Prometheus text sample line into metric name, labels, and value.
325fn parse_prometheus_sample(line: &str) -> Option<PrometheusSample<'_>> {
326    let (name_and_labels, value_str) = line.rsplit_once(' ')?;
327    #[expect(
328        clippy::cast_possible_truncation,
329        clippy::cast_sign_loss,
330        reason = "Prometheus counter/gauge values are non-negative integers"
331    )]
332    let value = {
333        let parsed = value_str.parse::<f64>().ok()?;
334        parsed.round() as u64
335    };
336    let (name, labels) = if let Some((name, label_blob)) = name_and_labels.split_once('{') {
337        let label_blob = label_blob.strip_suffix('}')?;
338        (name, parse_prometheus_labels(label_blob))
339    } else {
340        (name_and_labels, std::collections::HashMap::new())
341    };
342    Some((name, labels, value))
343}
344
345/// Parse `{key="value",...}` label sets from Prometheus text exposition.
346///
347/// Splits on commas only outside double-quoted values so embedded commas in
348/// label values parse correctly.
349fn parse_prometheus_labels(input: &str) -> std::collections::HashMap<String, String> {
350    let mut labels = std::collections::HashMap::new();
351    let mut rest = input.trim();
352    while !rest.is_empty() {
353        let (pair, tail) = split_prometheus_label_pair(rest);
354        if let Some((key, value)) = parse_prometheus_label_pair(pair) {
355            labels.insert(key, value);
356        }
357        rest = tail;
358    }
359    labels
360}
361
362/// Split one `key="value"` pair from the head of a label-set fragment.
363fn split_prometheus_label_pair(input: &str) -> (&str, &str) {
364    let bytes = input.as_bytes();
365    let mut in_quotes = false;
366    for (index, byte) in bytes.iter().enumerate() {
367        match *byte {
368            b'"' => in_quotes = !in_quotes,
369            b',' if !in_quotes => {
370                let head = input.get(..index).unwrap_or(input);
371                let tail = input.get(index + 1..).unwrap_or("").trim_start();
372                return (head, tail);
373            },
374            _ => {},
375        }
376    }
377    (input, "")
378}
379
380/// Parse one `key="value"` pair from a Prometheus label set fragment.
381fn parse_prometheus_label_pair(pair: &str) -> Option<(String, String)> {
382    let (key, value) = pair.split_once('=')?;
383    let value = value.strip_prefix('"')?.strip_suffix('"')?;
384    Some((key.to_owned(), value.to_owned()))
385}
386
387// -----------------------------------------------------------------------------
388// Status & Method Labels
389// -----------------------------------------------------------------------------
390
391/// Map an HTTP status code to its class label (`"1xx"`, `"2xx"`, etc.).
392///
393/// Returns `"unknown"` for zero (no response written) or codes
394/// outside the 100–599 range.
395///
396/// ```
397/// use praxis_protocol::http::pingora::metrics::status_class;
398///
399/// assert_eq!(status_class(200), "2xx");
400/// assert_eq!(status_class(404), "4xx");
401/// assert_eq!(status_class(0), "unknown");
402/// ```
403pub fn status_class(code: u16) -> &'static str {
404    match code {
405        100..=199 => "1xx",
406        200..=299 => "2xx",
407        300..=399 => "3xx",
408        400..=499 => "4xx",
409        500..=599 => "5xx",
410        _ => "unknown",
411    }
412}
413
414/// Map an HTTP method to a bounded label value.
415///
416/// Returns the method string for the nine standard methods
417/// defined in [RFC 9110]; all others collapse to `"OTHER"`.
418///
419/// ```
420/// use praxis_protocol::http::pingora::metrics::method_label;
421///
422/// assert_eq!(method_label("GET"), "GET");
423/// assert_eq!(method_label("PURGE"), "OTHER");
424/// ```
425///
426/// [RFC 9110]: https://datatracker.ietf.org/doc/html/rfc9110#section-9.1
427pub fn method_label(method: &str) -> &'static str {
428    match method {
429        "GET" => "GET",
430        "POST" => "POST",
431        "PUT" => "PUT",
432        "DELETE" => "DELETE",
433        "PATCH" => "PATCH",
434        "HEAD" => "HEAD",
435        "OPTIONS" => "OPTIONS",
436        "TRACE" => "TRACE",
437        "CONNECT" => "CONNECT",
438        _ => "OTHER",
439    }
440}
441
442// -----------------------------------------------------------------------------
443// Request Metrics
444// -----------------------------------------------------------------------------
445
446/// Labels for a completed HTTP request.
447///
448/// Static labels (`method`, `status_class`) use `&'static str`;
449/// `cluster` and `route` are dynamic [`SharedString`] values.
450///
451/// [`SharedString`]: ::metrics::SharedString
452pub(crate) struct RequestMetricLabels {
453    /// Cluster name or `"none"`.
454    pub cluster: SharedString,
455    /// HTTP method (e.g. `"GET"`).
456    pub method: &'static str,
457    /// Route path-match pattern or `"unknown"`.
458    pub route: SharedString,
459    /// Status class (e.g. `"2xx"`).
460    pub status_class: &'static str,
461}
462
463/// Build the enabled subset of the request-metric labels.
464fn selected_request_labels(labels: RequestMetricLabels) -> Vec<Label> {
465    let selected = metric_labels();
466    let pairs = [
467        (
468            "method",
469            label_if(
470                selected.is_enabled(MetricLabel::Method),
471                SharedString::const_str(labels.method),
472            ),
473        ),
474        (
475            "status_class",
476            label_if(
477                selected.is_enabled(MetricLabel::StatusClass),
478                SharedString::const_str(labels.status_class),
479            ),
480        ),
481        ("route", label_if(selected.is_enabled(MetricLabel::Route), labels.route)),
482        (
483            "cluster",
484            label_if(selected.is_enabled(MetricLabel::Cluster), labels.cluster),
485        ),
486    ];
487    selected_labels(&pairs)
488}
489
490/// Record HTTP request metrics for a completed request.
491pub(crate) fn record_request_metrics(labels: RequestMetricLabels, duration_secs: f64) {
492    if !is_recorder_installed() {
493        return;
494    }
495    if !metric_labels().all_enabled() {
496        let emitted = selected_request_labels(labels);
497        counter!(HTTP_REQUESTS_TOTAL, emitted.clone()).increment(1);
498        histogram!(HTTP_REQUEST_DURATION_SECONDS, emitted).record(duration_secs);
499        return;
500    }
501    record_request_metrics_all_labels(labels, duration_secs);
502}
503
504/// Record request metrics with the full default label set.
505///
506/// Kept on the static-label macro form so the default configuration emits
507/// exactly the series it did before label selection existed.
508fn record_request_metrics_all_labels(labels: RequestMetricLabels, duration_secs: f64) {
509    let cluster = labels.cluster;
510    let route = labels.route;
511    counter!(
512        HTTP_REQUESTS_TOTAL,
513        "method" => labels.method,
514        "status_class" => labels.status_class,
515        "route" => route.clone(),
516        "cluster" => cluster.clone()
517    )
518    .increment(1);
519    histogram!(
520        HTTP_REQUEST_DURATION_SECONDS,
521        "method" => labels.method,
522        "status_class" => labels.status_class,
523        "route" => route,
524        "cluster" => cluster
525    )
526    .record(duration_secs);
527}
528
529// -----------------------------------------------------------------------------
530// Body-Size Metrics
531// -----------------------------------------------------------------------------
532
533/// Build the enabled subset of the body-size histogram labels.
534fn selected_body_labels(method: &'static str, status_class: &'static str, cluster: SharedString) -> Vec<Label> {
535    let selected = metric_labels();
536    let pairs = [
537        (
538            "method",
539            label_if(
540                selected.is_enabled(MetricLabel::Method),
541                SharedString::const_str(method),
542            ),
543        ),
544        (
545            "status_class",
546            label_if(
547                selected.is_enabled(MetricLabel::StatusClass),
548                SharedString::const_str(status_class),
549            ),
550        ),
551        ("cluster", label_if(selected.is_enabled(MetricLabel::Cluster), cluster)),
552    ];
553    selected_labels(&pairs)
554}
555
556/// Record body-size histograms with the full default label set.
557fn record_body_size_all_labels(
558    method: &'static str,
559    status_class: &'static str,
560    cluster: SharedString,
561    request_bytes: f64,
562    response_bytes: f64,
563) {
564    histogram!(
565        HTTP_REQUEST_BODY_BYTES,
566        "method" => method,
567        "status_class" => status_class,
568        "cluster" => cluster.clone()
569    )
570    .record(request_bytes);
571    histogram!(
572        HTTP_RESPONSE_BODY_BYTES,
573        "method" => method,
574        "status_class" => status_class,
575        "cluster" => cluster
576    )
577    .record(response_bytes);
578}
579
580/// Record HTTP request and response body size histograms.
581pub(crate) fn record_body_size_metrics(
582    method: &'static str,
583    status_class: &'static str,
584    cluster: SharedString,
585    request_body_bytes: u64,
586    response_body_bytes: u64,
587) {
588    if !is_recorder_installed() {
589        return;
590    }
591    #[expect(
592        clippy::cast_precision_loss,
593        reason = "body byte counts as histogram observations; exact integer precision not required"
594    )]
595    let (request_bytes, response_bytes) = (request_body_bytes as f64, response_body_bytes as f64);
596
597    if !metric_labels().all_enabled() {
598        let emitted = selected_body_labels(method, status_class, cluster);
599        histogram!(HTTP_REQUEST_BODY_BYTES, emitted.clone()).record(request_bytes);
600        histogram!(HTTP_RESPONSE_BODY_BYTES, emitted).record(response_bytes);
601        return;
602    }
603    record_body_size_all_labels(method, status_class, cluster, request_bytes, response_bytes);
604}
605
606// -----------------------------------------------------------------------------
607// Active Request Tracking
608// -----------------------------------------------------------------------------
609
610/// RAII guard that decrements `praxis_http_active_requests` on drop.
611///
612/// Acquired once per HTTP request. Pingora owns the request context by
613/// value, so the drop runs on every terminal path (including a client
614/// abort or an HTTP/2 stream reset, which skip the `logging` callback
615/// entirely).
616pub struct ActiveRequestGuard {
617    /// Listener name label.
618    listener: SharedString,
619}
620
621impl ActiveRequestGuard {
622    /// Increment the gauge and return a guard that decrements on drop.
623    pub(crate) fn acquire(listener: SharedString) -> Self {
624        if is_recorder_installed() {
625            if metric_labels().is_enabled(MetricLabel::Listener) {
626                gauge!(HTTP_ACTIVE_REQUESTS, "listener" => listener.clone()).increment(1.0);
627            } else {
628                gauge!(HTTP_ACTIVE_REQUESTS).increment(1.0);
629            }
630        }
631        Self { listener }
632    }
633}
634
635impl Drop for ActiveRequestGuard {
636    fn drop(&mut self) {
637        if is_recorder_installed() {
638            if metric_labels().is_enabled(MetricLabel::Listener) {
639                gauge!(HTTP_ACTIVE_REQUESTS, "listener" => self.listener.clone()).decrement(1.0);
640            } else {
641                gauge!(HTTP_ACTIVE_REQUESTS).decrement(1.0);
642            }
643        }
644    }
645}
646
647// -----------------------------------------------------------------------------
648// Error Classification
649// -----------------------------------------------------------------------------
650
651/// Record a proxy error.
652///
653/// Counted once per request, from the logging hook. Overload rejections and
654/// upstream connect failures have dedicated counters and are not repeated
655/// here; this counter covers the causes those miss.
656pub(crate) fn record_error(error_type: &'static str) {
657    if !is_recorder_installed() {
658        return;
659    }
660    counter!(ERRORS_TOTAL, "type" => error_type).increment(1);
661}
662
663/// Classify a Pingora error into a bounded `type` label value.
664///
665/// Connect failures map to `upstream_unavailable` and are also counted by
666/// `praxis_upstream_connect_failures_total`; the overlap is deliberate so
667/// that `praxis_errors_total` is a complete error denominator on its own.
668pub(crate) fn error_type_for(etype: &::pingora_core::ErrorType, source: &::pingora_core::ErrorSource) -> &'static str {
669    use ::pingora_core::ErrorSource::{Downstream, Internal, Unset};
670
671    if matches!(source, Downstream) {
672        return ERROR_TYPE_DOWNSTREAM;
673    }
674    if is_timeout(etype) {
675        return ERROR_TYPE_TIMEOUT;
676    }
677    if is_unreachable(etype) {
678        return ERROR_TYPE_UPSTREAM_UNAVAILABLE;
679    }
680    if matches!(source, Internal | Unset) {
681        return ERROR_TYPE_INTERNAL;
682    }
683    ERROR_TYPE_UPSTREAM_PROTOCOL
684}
685
686/// Whether the error is a connect, handshake, read or write timeout.
687fn is_timeout(etype: &::pingora_core::ErrorType) -> bool {
688    use ::pingora_core::ErrorType::{ConnectTimedout, ReadTimedout, TLSHandshakeTimedout, WriteTimedout};
689
690    matches!(
691        etype,
692        ConnectTimedout | TLSHandshakeTimedout | ReadTimedout | WriteTimedout
693    )
694}
695
696/// Whether the error means the upstream was never reached.
697fn is_unreachable(etype: &::pingora_core::ErrorType) -> bool {
698    use ::pingora_core::ErrorType::{BindError, ConnectError, ConnectNoRoute, ConnectRefused, SocketError};
699
700    matches!(
701        etype,
702        ConnectRefused | ConnectNoRoute | ConnectError | BindError | SocketError
703    )
704}
705
706// -----------------------------------------------------------------------------
707// Overload Rejections
708// -----------------------------------------------------------------------------
709
710/// Record an overload rejection.
711pub(crate) fn record_overload_reject(reason: &'static str) {
712    if !is_recorder_installed() {
713        return;
714    }
715    counter!(OVERLOAD_REJECTS_TOTAL, "reason" => reason).increment(1);
716}
717
718// -----------------------------------------------------------------------------
719// Upstream Metrics
720// -----------------------------------------------------------------------------
721
722/// Record upstream connect duration for a cluster.
723pub(crate) fn record_upstream_connect_duration(cluster: SharedString, duration_secs: f64) {
724    if !is_recorder_installed() {
725        return;
726    }
727    if metric_labels().is_enabled(MetricLabel::Cluster) {
728        histogram!(UPSTREAM_CONNECT_DURATION_SECONDS, "cluster" => cluster).record(duration_secs);
729    } else {
730        histogram!(UPSTREAM_CONNECT_DURATION_SECONDS).record(duration_secs);
731    }
732}
733
734/// Record a request that reached an upstream endpoint.
735///
736/// Counted once per request, from the logging hook, so a request retried
737/// across endpoints increments once against the endpoint that answered
738/// rather than once per attempt. Requests that never reached an upstream
739/// (filter rejections, connect failures) are not counted here.
740pub(crate) fn record_upstream_request(cluster: SharedString, endpoint: SharedString, status_class: &'static str) {
741    if !is_recorder_installed() {
742        return;
743    }
744    let selected = metric_labels();
745    if !selected.all_enabled() {
746        let pairs = [
747            ("cluster", label_if(selected.is_enabled(MetricLabel::Cluster), cluster)),
748            (
749                "endpoint",
750                label_if(selected.is_enabled(MetricLabel::Endpoint), endpoint),
751            ),
752            (
753                "status_class",
754                label_if(
755                    selected.is_enabled(MetricLabel::StatusClass),
756                    SharedString::const_str(status_class),
757                ),
758            ),
759        ];
760        counter!(UPSTREAM_REQUESTS_TOTAL, selected_labels(&pairs)).increment(1);
761        return;
762    }
763    counter!(
764        UPSTREAM_REQUESTS_TOTAL,
765        "cluster" => cluster,
766        "endpoint" => endpoint,
767        "status_class" => status_class
768    )
769    .increment(1);
770}
771
772/// Record an upstream connect failure.
773pub(crate) fn record_upstream_connect_failure(cluster: SharedString) {
774    if !is_recorder_installed() {
775        return;
776    }
777    if metric_labels().is_enabled(MetricLabel::Cluster) {
778        counter!(UPSTREAM_CONNECT_FAILURES_TOTAL, "cluster" => cluster).increment(1);
779    } else {
780        counter!(UPSTREAM_CONNECT_FAILURES_TOTAL).increment(1);
781    }
782}
783
784/// Record an upstream connect-failure retry outcome.
785pub(crate) fn record_upstream_retry(cluster: SharedString, result: &'static str) {
786    if !is_recorder_installed() {
787        return;
788    }
789    if metric_labels().is_enabled(MetricLabel::Cluster) {
790        counter!(UPSTREAM_RETRIES_TOTAL, "cluster" => cluster, "result" => result).increment(1);
791    } else {
792        counter!(UPSTREAM_RETRIES_TOTAL, "result" => result).increment(1);
793    }
794}
795
796// -----------------------------------------------------------------------------
797// Upstream Health
798// -----------------------------------------------------------------------------
799
800/// Refresh cluster endpoint health gauges.
801///
802/// The `cluster` label is structural here and is kept even when the
803/// `cluster` dimension is disabled: these gauges are keyed by cluster and
804/// set (not incremented), so dropping the label would collapse every
805/// cluster onto one last-writer-wins series rather than merely lowering
806/// cardinality. Disabling `cluster` therefore drops it from the additive
807/// cluster metrics (counters and the connect-duration histogram) but not
808/// from the per-cluster health gauges.
809pub(crate) fn set_upstream_endpoint_gauges(cluster: SharedString, healthy: usize, total: usize) {
810    if !is_recorder_installed() {
811        return;
812    }
813    #[expect(clippy::cast_precision_loss, reason = "endpoint counts fit f64 exactly below 2^53")]
814    {
815        gauge!(UPSTREAM_HEALTHY_ENDPOINTS, "cluster" => cluster.clone()).set(healthy as f64);
816        gauge!(UPSTREAM_TOTAL_ENDPOINTS, "cluster" => cluster).set(total as f64);
817    }
818}
819
820/// Zero health gauges for clusters that lost active health checks on reload.
821///
822/// Prometheus does not drop series automatically; clearing removed clusters
823/// prevents stale `healthy`/`total` values from lingering after config change.
824pub fn clear_stale_upstream_health_gauges<'a, P: IntoIterator<Item = &'a str>, C: IntoIterator<Item = &'a str>>(
825    previous_health_clusters: P,
826    current_health_clusters: C,
827) {
828    if !is_recorder_installed() {
829        return;
830    }
831    let current: std::collections::HashSet<&str> = current_health_clusters.into_iter().collect();
832    for name in previous_health_clusters {
833        if !current.contains(name) {
834            set_upstream_endpoint_gauges(SharedString::from(name.to_owned()), 0, 0);
835        }
836    }
837}
838
839/// Publish current healthy/total gauges for every cluster in a health registry.
840///
841/// Called on reload so scrapes reflect the new registry immediately instead of
842/// waiting for the first probe round.
843pub fn seed_upstream_health_gauges(registry: &praxis_core::health::HealthRegistry) {
844    if !is_recorder_installed() {
845        return;
846    }
847    for (name, state) in registry.iter() {
848        let (healthy, total) = state.endpoint_counts();
849        set_upstream_endpoint_gauges(SharedString::from(name.as_ref().to_owned()), healthy, total);
850    }
851}
852
853/// Record an endpoint health state transition and refresh gauges.
854pub(crate) fn record_health_transition(cluster: SharedString, result: &'static str, healthy: usize, total: usize) {
855    if !is_recorder_installed() {
856        return;
857    }
858    if metric_labels().is_enabled(MetricLabel::Cluster) {
859        counter!(
860            UPSTREAM_HEALTH_TRANSITIONS_TOTAL,
861            "cluster" => cluster.clone(),
862            "result" => result
863        )
864        .increment(1);
865    } else {
866        counter!(UPSTREAM_HEALTH_TRANSITIONS_TOTAL, "result" => result).increment(1);
867    }
868    set_upstream_endpoint_gauges(cluster, healthy, total);
869}
870
871/// Count healthy endpoints in a cluster health entry.
872pub(crate) fn count_healthy_endpoints(health: &praxis_core::health::ClusterHealthEntry) -> (usize, usize) {
873    health.endpoint_counts()
874}
875
876// -----------------------------------------------------------------------------
877// Config Reload
878// -----------------------------------------------------------------------------
879
880/// Record a successful config reload.
881pub fn record_config_reload_success() {
882    if !is_recorder_installed() {
883        return;
884    }
885    counter!(CONFIG_RELOAD_TOTAL, "result" => RELOAD_RESULT_SUCCESS).increment(1);
886    let ts = std::time::SystemTime::now()
887        .duration_since(std::time::UNIX_EPOCH)
888        .map_or(0.0, |d| d.as_secs_f64());
889    gauge!(CONFIG_RELOAD_LAST_SUCCESS_TIMESTAMP).set(ts);
890}
891
892/// Record a failed config reload.
893pub fn record_config_reload_failure() {
894    if !is_recorder_installed() {
895        return;
896    }
897    counter!(CONFIG_RELOAD_TOTAL, "result" => RELOAD_RESULT_FAILURE).increment(1);
898}
899
900// -----------------------------------------------------------------------------
901// Shared Label Values
902// -----------------------------------------------------------------------------
903
904/// [`SharedString`] for the `"none"` cluster label.
905///
906/// [`SharedString`]: ::metrics::SharedString
907pub(crate) fn cluster_none() -> SharedString {
908    SharedString::const_str("none")
909}
910
911/// [`SharedString`] for the `"unknown"` route label.
912///
913/// [`SharedString`]: ::metrics::SharedString
914pub(crate) fn route_unknown() -> SharedString {
915    SharedString::const_str("unknown")
916}
917
918// -----------------------------------------------------------------------------
919// Tests
920// -----------------------------------------------------------------------------
921
922#[cfg(test)]
923#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
924#[allow(clippy::unwrap_used, clippy::expect_used, clippy::indexing_slicing, reason = "tests")]
925mod tests {
926    use super::*;
927
928    #[test]
929    fn status_class_1xx() {
930        assert_eq!(status_class(100), "1xx", "100 should be 1xx");
931        assert_eq!(status_class(199), "1xx", "199 should be 1xx");
932    }
933
934    #[test]
935    fn status_class_2xx() {
936        assert_eq!(status_class(200), "2xx", "200 should be 2xx");
937        assert_eq!(status_class(204), "2xx", "204 should be 2xx");
938        assert_eq!(status_class(299), "2xx", "299 should be 2xx");
939    }
940
941    #[test]
942    fn status_class_3xx() {
943        assert_eq!(status_class(301), "3xx", "301 should be 3xx");
944        assert_eq!(status_class(399), "3xx", "399 should be 3xx");
945    }
946
947    #[test]
948    fn status_class_4xx() {
949        assert_eq!(status_class(400), "4xx", "400 should be 4xx");
950        assert_eq!(status_class(404), "4xx", "404 should be 4xx");
951        assert_eq!(status_class(499), "4xx", "499 should be 4xx");
952    }
953
954    #[test]
955    fn status_class_5xx() {
956        assert_eq!(status_class(500), "5xx", "500 should be 5xx");
957        assert_eq!(status_class(503), "5xx", "503 should be 5xx");
958        assert_eq!(status_class(599), "5xx", "599 should be 5xx");
959    }
960
961    #[test]
962    fn status_class_zero_is_unknown() {
963        assert_eq!(status_class(0), "unknown", "0 should be unknown");
964    }
965
966    #[test]
967    fn status_class_out_of_range_is_unknown() {
968        assert_eq!(status_class(600), "unknown", "600 should be unknown");
969        assert_eq!(status_class(99), "unknown", "99 should be unknown");
970    }
971
972    #[test]
973    fn method_label_standard_methods() {
974        for m in [
975            "GET", "POST", "PUT", "DELETE", "PATCH", "HEAD", "OPTIONS", "TRACE", "CONNECT",
976        ] {
977            assert_eq!(method_label(m), m, "{m} should pass through");
978        }
979    }
980
981    #[test]
982    fn method_label_custom_methods_collapse_to_other() {
983        assert_eq!(method_label("PURGE"), "OTHER", "PURGE should be OTHER");
984        assert_eq!(method_label("FOOBAR"), "OTHER", "FOOBAR should be OTHER");
985        assert_eq!(method_label(""), "OTHER", "empty should be OTHER");
986    }
987
988    #[test]
989    fn record_utilities_noop_without_recorder() {
990        record_overload_reject(OVERLOAD_REASON_MEMORY);
991        record_upstream_connect_failure(cluster_none());
992        record_error(ERROR_TYPE_INTERNAL);
993        record_upstream_request(cluster_none(), SharedString::const_str("10.0.0.1:80"), "2xx");
994        record_upstream_retry(cluster_none(), RETRY_RESULT_SUCCESS);
995        record_upstream_connect_duration(cluster_none(), 0.01);
996        set_upstream_endpoint_gauges(cluster_none(), 1, 2);
997        record_health_transition(cluster_none(), HEALTH_RESULT_HEALTHY, 1, 2);
998        record_config_reload_success();
999        record_config_reload_failure();
1000        clear_stale_upstream_health_gauges(["gone"], std::iter::empty::<&str>());
1001        let _request_guard = ActiveRequestGuard::acquire(SharedString::const_str("test"));
1002    }
1003
1004    #[test]
1005    fn active_request_guard_returns_to_zero_on_drop() {
1006        install_prometheus_recorder();
1007        let listener = SharedString::const_str("active-request-guard-listener");
1008        let guard = ActiveRequestGuard::acquire(listener.clone());
1009        let held = render_prometheus().expect("recorder should render");
1010        assert!(
1011            held.contains("praxis_http_active_requests{listener=\"active-request-guard-listener\"} 1"),
1012            "gauge should read 1 while the guard is held:\n{held}"
1013        );
1014        drop(guard);
1015        let released = render_prometheus().expect("recorder should render");
1016        assert!(
1017            released.contains("praxis_http_active_requests{listener=\"active-request-guard-listener\"} 0"),
1018            "gauge should return to 0 once the guard drops:\n{released}"
1019        );
1020    }
1021
1022    #[test]
1023    fn overload_reject_reasons_appear_in_scrape() {
1024        install_prometheus_recorder();
1025        record_overload_reject(OVERLOAD_REASON_MEMORY);
1026        record_overload_reject(OVERLOAD_REASON_GLOBAL_CONNECTIONS);
1027        record_overload_reject(OVERLOAD_REASON_LISTENER_CONNECTIONS);
1028        let body = render_prometheus().expect("recorder should render");
1029        for reason in [
1030            OVERLOAD_REASON_MEMORY,
1031            OVERLOAD_REASON_GLOBAL_CONNECTIONS,
1032            OVERLOAD_REASON_LISTENER_CONNECTIONS,
1033        ] {
1034            let needle = format!("praxis_overload_rejects_total{{reason=\"{reason}\"}}");
1035            assert!(body.contains(&needle), "expected `{needle}` in scrape:\n{body}");
1036        }
1037    }
1038
1039    #[test]
1040    fn upstream_requests_carry_cluster_endpoint_and_status_class() {
1041        install_prometheus_recorder();
1042        record_upstream_request(
1043            SharedString::const_str("api"),
1044            SharedString::const_str("10.0.0.7:8080"),
1045            "5xx",
1046        );
1047        let body = render_prometheus().expect("recorder should render");
1048        assert!(
1049            body.contains(
1050                "praxis_upstream_requests_total{cluster=\"api\",endpoint=\"10.0.0.7:8080\",status_class=\"5xx\"} 1"
1051            ),
1052            "counter should carry all three labels:\n{body}"
1053        );
1054    }
1055
1056    #[test]
1057    fn selected_labels_drops_disabled_dimensions() {
1058        let pairs = [
1059            ("method", Some(SharedString::const_str("GET"))),
1060            ("route", None),
1061            ("cluster", Some(SharedString::const_str("api"))),
1062        ];
1063        let emitted = selected_labels(&pairs);
1064        let names: Vec<&str> = emitted.iter().map(Label::key).collect();
1065        assert_eq!(names, vec!["method", "cluster"], "a disabled dimension must be absent");
1066    }
1067
1068    #[test]
1069    fn selected_labels_preserves_order_and_values() {
1070        let pairs = [
1071            ("cluster", Some(SharedString::const_str("api"))),
1072            ("endpoint", Some(SharedString::const_str("10.0.0.1:80"))),
1073        ];
1074        let emitted = selected_labels(&pairs);
1075        let rendered: Vec<(&str, &str)> = emitted.iter().map(|l| (l.key(), l.value())).collect();
1076        assert_eq!(
1077            rendered,
1078            vec![("cluster", "api"), ("endpoint", "10.0.0.1:80")],
1079            "enabled dimensions keep their order and values"
1080        );
1081    }
1082
1083    #[test]
1084    fn label_if_gates_on_the_flag() {
1085        assert_eq!(
1086            label_if(true, SharedString::const_str("x")).as_deref(),
1087            Some("x"),
1088            "an enabled dimension keeps its value"
1089        );
1090        assert_eq!(
1091            label_if(false, SharedString::const_str("x")),
1092            None,
1093            "a disabled dimension yields no value"
1094        );
1095    }
1096
1097    #[test]
1098    fn metric_labels_default_to_all_enabled() {
1099        assert!(
1100            metric_labels().all_enabled(),
1101            "without an explicit install every dimension must stay on, so the \
1102             recorders keep their allocation-free fast path"
1103        );
1104    }
1105
1106    #[test]
1107    fn error_types_appear_in_scrape() {
1108        install_prometheus_recorder();
1109        for error_type in [
1110            ERROR_TYPE_FILTER_REJECT,
1111            ERROR_TYPE_TIMEOUT,
1112            ERROR_TYPE_UPSTREAM_UNAVAILABLE,
1113            ERROR_TYPE_UPSTREAM_PROTOCOL,
1114            ERROR_TYPE_DOWNSTREAM,
1115            ERROR_TYPE_INTERNAL,
1116        ] {
1117            record_error(error_type);
1118            let body = render_prometheus().expect("recorder should render");
1119            let needle = format!("praxis_errors_total{{type=\"{error_type}\"}}");
1120            assert!(body.contains(&needle), "expected `{needle}` in scrape:\n{body}");
1121        }
1122    }
1123
1124    #[test]
1125    fn error_type_for_maps_pingora_errors_to_bounded_values() {
1126        use ::pingora_core::{ErrorSource, ErrorType};
1127
1128        assert_eq!(
1129            error_type_for(&ErrorType::ConnectTimedout, &ErrorSource::Upstream),
1130            ERROR_TYPE_TIMEOUT,
1131            "connect timeout is a timeout"
1132        );
1133        assert_eq!(
1134            error_type_for(&ErrorType::ConnectRefused, &ErrorSource::Upstream),
1135            ERROR_TYPE_UPSTREAM_UNAVAILABLE,
1136            "a refused connect means the upstream was unreachable"
1137        );
1138        assert_eq!(
1139            error_type_for(&ErrorType::ReadError, &ErrorSource::Upstream),
1140            ERROR_TYPE_UPSTREAM_PROTOCOL,
1141            "a mid-exchange read error is a protocol failure"
1142        );
1143        assert_eq!(
1144            error_type_for(&ErrorType::ReadTimedout, &ErrorSource::Downstream),
1145            ERROR_TYPE_DOWNSTREAM,
1146            "downstream source wins over the error kind"
1147        );
1148        assert_eq!(
1149            error_type_for(&ErrorType::InternalError, &ErrorSource::Internal),
1150            ERROR_TYPE_INTERNAL,
1151            "internal source is an internal fault"
1152        );
1153    }
1154
1155    #[test]
1156    fn body_size_histograms_use_byte_buckets() {
1157        install_prometheus_recorder();
1158        record_body_size_metrics("GET", "2xx", cluster_none(), 500, 4_000);
1159        let body = render_prometheus().expect("recorder should render");
1160        assert!(
1161            body.contains("praxis_http_request_body_bytes_bucket") && body.contains("le=\"1024\""),
1162            "request body histogram should use byte buckets, not duration defaults:\n{body}"
1163        );
1164        assert!(
1165            body.contains("praxis_http_response_body_bytes_bucket") && body.contains("le=\"4096\""),
1166            "response body histogram should use byte buckets:\n{body}"
1167        );
1168        assert!(
1169            !body.contains("praxis_http_request_body_bytes_bucket{le=\"0.005\"}")
1170                && !body.contains("praxis_http_request_body_bytes_bucket{method=\"GET\",status_class=\"2xx\",cluster=\"\",le=\"0.005\"}"),
1171            "request body histogram must not use duration default buckets:\n{body}"
1172        );
1173    }
1174
1175    #[test]
1176    fn clear_stale_upstream_health_gauges_zeros_removed_clusters() {
1177        install_prometheus_recorder();
1178        set_upstream_endpoint_gauges(SharedString::from("old-cluster".to_owned()), 2, 3);
1179        set_upstream_endpoint_gauges(SharedString::from("kept-cluster".to_owned()), 1, 1);
1180        clear_stale_upstream_health_gauges(["old-cluster", "kept-cluster"], ["kept-cluster"]);
1181        let body = render_prometheus().expect("recorder should render");
1182        assert!(
1183            body.contains("praxis_upstream_healthy_endpoints{cluster=\"old-cluster\"} 0"),
1184            "removed cluster healthy gauge should be zeroed:\n{body}"
1185        );
1186        assert!(
1187            body.contains("praxis_upstream_total_endpoints{cluster=\"old-cluster\"} 0"),
1188            "removed cluster total gauge should be zeroed:\n{body}"
1189        );
1190        assert!(
1191            body.contains("praxis_upstream_healthy_endpoints{cluster=\"kept-cluster\"} 1"),
1192            "kept cluster should retain its value:\n{body}"
1193        );
1194    }
1195
1196    #[test]
1197    fn collect_stats_metrics_sums_cluster_counters() {
1198        let text = r#"
1199praxis_http_active_requests{listener="web"} 2
1200praxis_tcp_active_connections{listener="tcp-in"} 1
1201praxis_upstream_requests_total{cluster="backend",endpoint="127.0.0.1:1",status_class="2xx"} 3
1202praxis_upstream_requests_total{cluster="backend",endpoint="127.0.0.1:2",status_class="5xx"} 1
1203praxis_upstream_connect_failures_total{cluster="backend"} 2
1204"#;
1205        let snap = collect_stats_metrics(text);
1206        assert_eq!(
1207            snap.http_active_by_listener.get("web"),
1208            Some(&2),
1209            "HTTP active per listener should parse"
1210        );
1211        assert_eq!(
1212            snap.tcp_active_by_listener.get("tcp-in"),
1213            Some(&1),
1214            "TCP active per listener should parse"
1215        );
1216        assert_eq!(
1217            snap.upstream_requests_by_cluster.get("backend"),
1218            Some(&4),
1219            "upstream requests should sum by cluster"
1220        );
1221        assert_eq!(
1222            snap.connect_failures_by_cluster.get("backend"),
1223            Some(&2),
1224            "connect failures should parse by cluster"
1225        );
1226    }
1227
1228    #[test]
1229    fn parse_prometheus_labels_handles_commas_inside_quoted_values() {
1230        let labels = parse_prometheus_labels(r#"tag="a,b",listener="web""#);
1231        assert_eq!(labels.get("tag"), Some(&"a,b".to_owned()), "comma inside quotes");
1232        assert_eq!(labels.get("listener"), Some(&"web".to_owned()), "second label");
1233    }
1234
1235    #[test]
1236    fn collect_stats_metrics_parses_unlabeled_upstream_counters() {
1237        let text = r#"
1238praxis_upstream_requests_total{status_class="2xx"} 5
1239praxis_upstream_connect_failures_total 2
1240"#;
1241        let snap = collect_stats_metrics(text);
1242        assert_eq!(
1243            snap.upstream_requests_aggregate,
1244            Some(5),
1245            "unlabeled upstream requests should aggregate"
1246        );
1247        assert_eq!(
1248            snap.connect_failures_aggregate,
1249            Some(2),
1250            "unlabeled connect failures should aggregate"
1251        );
1252    }
1253
1254    #[test]
1255    fn seed_upstream_health_gauges_publishes_registry_counts() {
1256        use std::sync::Arc;
1257
1258        use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
1259
1260        install_prometheus_recorder();
1261        let endpoints = vec![EndpointHealth::new(), EndpointHealth::new()];
1262        endpoints[0].mark_unhealthy();
1263        let entry = Arc::new(ClusterHealthEntry::new(
1264            endpoints,
1265            vec![Arc::from("a:1"), Arc::from("b:1")],
1266            None,
1267            None,
1268        ));
1269        let registry = Arc::new([(Arc::from("backend"), entry)].into_iter().collect());
1270        seed_upstream_health_gauges(&registry);
1271        let body = render_prometheus().expect("recorder should render");
1272        assert!(
1273            body.contains("praxis_upstream_healthy_endpoints{cluster=\"backend\"} 1"),
1274            "seed should publish healthy count:\n{body}"
1275        );
1276        assert!(
1277            body.contains("praxis_upstream_total_endpoints{cluster=\"backend\"} 2"),
1278            "seed should publish total count:\n{body}"
1279        );
1280    }
1281
1282    #[test]
1283    fn count_healthy_endpoints_counts_correctly() {
1284        use std::sync::Arc;
1285
1286        use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
1287
1288        let endpoints = vec![EndpointHealth::new(), EndpointHealth::new(), EndpointHealth::new()];
1289        endpoints[1].mark_unhealthy();
1290        let entry = ClusterHealthEntry::new(
1291            endpoints,
1292            vec![Arc::from("a:1"), Arc::from("b:1"), Arc::from("c:1")],
1293            None,
1294            None,
1295        );
1296        let (healthy, total) = count_healthy_endpoints(&entry);
1297        assert_eq!(total, 3, "total should be 3");
1298        assert_eq!(healthy, 2, "two endpoints should be healthy");
1299    }
1300}