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