1use 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
13const HTTP_REQUESTS_TOTAL: &str = "praxis_http_requests_total";
19
20const HTTP_REQUEST_DURATION_SECONDS: &str = "praxis_http_request_duration_seconds";
22
23const HTTP_REQUEST_BODY_BYTES: &str = "praxis_http_request_body_bytes";
25
26const HTTP_RESPONSE_BODY_BYTES: &str = "praxis_http_response_body_bytes";
28
29const HTTP_ACTIVE_REQUESTS: &str = "praxis_http_active_requests";
36
37const TCP_ACTIVE_CONNECTIONS: &str = "praxis_tcp_active_connections";
39
40const OVERLOAD_REJECTS_TOTAL: &str = "praxis_overload_rejects_total";
42
43const UPSTREAM_REQUESTS_TOTAL: &str = "praxis_upstream_requests_total";
45
46const UPSTREAM_CONNECT_DURATION_SECONDS: &str = "praxis_upstream_connect_duration_seconds";
48
49const UPSTREAM_CONNECT_FAILURES_TOTAL: &str = "praxis_upstream_connect_failures_total";
51
52const UPSTREAM_RETRIES_TOTAL: &str = "praxis_upstream_retries_total";
54
55const UPSTREAM_HEALTHY_ENDPOINTS: &str = "praxis_upstream_healthy_endpoints";
57
58const UPSTREAM_TOTAL_ENDPOINTS: &str = "praxis_upstream_total_endpoints";
60
61const UPSTREAM_HEALTH_TRANSITIONS_TOTAL: &str = "praxis_upstream_health_transitions_total";
63
64const CONFIG_RELOAD_TOTAL: &str = "praxis_config_reload_total";
66
67const CONFIG_RELOAD_LAST_SUCCESS_TIMESTAMP: &str = "praxis_config_reload_last_success_timestamp";
69
70const ERRORS_TOTAL: &str = "praxis_errors_total";
72
73pub(crate) const ERROR_TYPE_FILTER_REJECT: &str = "filter_reject";
75
76pub(crate) const ERROR_TYPE_TIMEOUT: &str = "timeout";
78
79pub(crate) const ERROR_TYPE_UPSTREAM_UNAVAILABLE: &str = "upstream_unavailable";
81
82pub(crate) const ERROR_TYPE_UPSTREAM_PROTOCOL: &str = "upstream_protocol";
84
85pub(crate) const ERROR_TYPE_DOWNSTREAM: &str = "downstream";
87
88pub(crate) const ERROR_TYPE_INTERNAL: &str = "internal";
90
91pub(crate) const OVERLOAD_REASON_MEMORY: &str = "memory";
93
94pub(crate) const OVERLOAD_REASON_GLOBAL_CONNECTIONS: &str = "global_connections";
96
97pub(crate) const OVERLOAD_REASON_LISTENER_CONNECTIONS: &str = "listener_connections";
99
100pub(crate) const RETRY_RESULT_SUCCESS: &str = "success";
102
103pub(crate) const RETRY_RESULT_EXHAUSTED: &str = "exhausted";
105
106pub(crate) const HEALTH_RESULT_HEALTHY: &str = "healthy";
108
109pub(crate) const HEALTH_RESULT_UNHEALTHY: &str = "unhealthy";
111
112pub(crate) const RELOAD_RESULT_SUCCESS: &str = "success";
114
115pub(crate) const RELOAD_RESULT_FAILURE: &str = "failure";
117
118const 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
134static PROMETHEUS_HANDLE: OnceLock<PrometheusHandle> = OnceLock::new();
140
141static LABEL_CONFIG: OnceLock<MetricLabelsConfig> = OnceLock::new();
143
144static ALL_LABELS: OnceLock<MetricLabelsConfig> = OnceLock::new();
146
147pub fn install_metric_labels(labels: MetricLabelsConfig) {
153 let _existing = LABEL_CONFIG.set(labels);
154}
155
156pub(crate) fn metric_labels() -> &'static MetricLabelsConfig {
158 LABEL_CONFIG
159 .get()
160 .unwrap_or_else(|| ALL_LABELS.get_or_init(MetricLabelsConfig::default))
161}
162
163fn 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
173fn label_if(enabled: bool, value: SharedString) -> Option<SharedString> {
175 enabled.then_some(value)
176}
177
178pub 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
209pub fn render_prometheus() -> Option<String> {
213 PROMETHEUS_HANDLE.get().map(PrometheusHandle::render)
214}
215
216pub(crate) fn is_recorder_installed() -> bool {
218 PROMETHEUS_HANDLE.get().is_some()
219}
220
221#[derive(Clone, Debug, Default, Eq, PartialEq)]
227pub struct StatsMetricsSnapshot {
228 pub http_active_by_listener: std::collections::HashMap<String, u64>,
230 pub http_active_aggregate: Option<u64>,
232 pub tcp_active_by_listener: std::collections::HashMap<String, u64>,
234 pub tcp_active_aggregate: Option<u64>,
236 pub upstream_requests_by_cluster: std::collections::HashMap<String, u64>,
238 pub upstream_requests_aggregate: Option<u64>,
240 pub connect_failures_by_cluster: std::collections::HashMap<String, u64>,
242 pub connect_failures_aggregate: Option<u64>,
244}
245
246#[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
298type PrometheusSample<'a> = (&'a str, std::collections::HashMap<String, String>, u64);
300
301fn 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
322fn 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
339fn 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
357fn 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
364pub 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
391pub 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
419pub(crate) struct RequestMetricLabels {
430 pub cluster: SharedString,
432 pub method: &'static str,
434 pub route: SharedString,
436 pub status_class: &'static str,
438}
439
440fn 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
467pub(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
481fn 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
506fn 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
529fn 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
553pub(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
579pub struct ActiveRequestGuard {
586 listener: SharedString,
588}
589
590impl ActiveRequestGuard {
591 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
616pub(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
628pub(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
651fn 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
661fn 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
671pub(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
679pub(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
691pub(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
729pub(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
741pub(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
753pub(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
773pub 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
792pub 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
806pub(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
824pub(crate) fn count_healthy_endpoints(health: &praxis_core::health::ClusterHealthEntry) -> (usize, usize) {
826 health.endpoint_counts()
827}
828
829pub 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
841pub 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
849pub(crate) fn cluster_none() -> SharedString {
853 SharedString::const_str("none")
854}
855
856pub(crate) fn route_unknown() -> SharedString {
860 SharedString::const_str("unknown")
861}
862
863#[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(®istry);
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}