1use 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
21const HTTP_REQUESTS_TOTAL: &str = "praxis_http_requests_total";
27
28const HTTP_REQUEST_DURATION_SECONDS: &str = "praxis_http_request_duration_seconds";
30
31const HTTP_REQUEST_BODY_BYTES: &str = "praxis_http_request_body_bytes";
33
34const HTTP_RESPONSE_BODY_BYTES: &str = "praxis_http_response_body_bytes";
36
37const HTTP_ACTIVE_REQUESTS: &str = "praxis_http_active_requests";
44
45const TCP_ACTIVE_CONNECTIONS: &str = "praxis_tcp_active_connections";
47
48const OVERLOAD_REJECTS_TOTAL: &str = "praxis_overload_rejects_total";
50
51const UPSTREAM_REQUESTS_TOTAL: &str = "praxis_upstream_requests_total";
53
54const UPSTREAM_CONNECT_DURATION_SECONDS: &str = "praxis_upstream_connect_duration_seconds";
56
57const UPSTREAM_CONNECT_FAILURES_TOTAL: &str = "praxis_upstream_connect_failures_total";
59
60const UPSTREAM_RETRIES_TOTAL: &str = "praxis_upstream_retries_total";
62
63const UPSTREAM_HEALTHY_ENDPOINTS: &str = "praxis_upstream_healthy_endpoints";
65
66const UPSTREAM_TOTAL_ENDPOINTS: &str = "praxis_upstream_total_endpoints";
68
69const UPSTREAM_HEALTH_TRANSITIONS_TOTAL: &str = "praxis_upstream_health_transitions_total";
71
72const CONFIG_RELOAD_TOTAL: &str = "praxis_config_reload_total";
74
75const CONFIG_RELOAD_LAST_SUCCESS_TIMESTAMP: &str = "praxis_config_reload_last_success_timestamp";
77
78const ERRORS_TOTAL: &str = "praxis_errors_total";
80
81pub(crate) const ERROR_TYPE_FILTER_REJECT: &str = "filter_reject";
87
88pub(crate) const ERROR_TYPE_TIMEOUT: &str = "timeout";
90
91pub(crate) const ERROR_TYPE_UPSTREAM_UNAVAILABLE: &str = "upstream_unavailable";
93
94pub(crate) const ERROR_TYPE_UPSTREAM_PROTOCOL: &str = "upstream_protocol";
96
97pub(crate) const ERROR_TYPE_DOWNSTREAM: &str = "downstream";
99
100pub(crate) const ERROR_TYPE_INTERNAL: &str = "internal";
102
103pub(crate) const OVERLOAD_REASON_MEMORY: &str = "memory";
105
106pub(crate) const OVERLOAD_REASON_GLOBAL_CONNECTIONS: &str = "global_connections";
108
109pub(crate) const OVERLOAD_REASON_LISTENER_CONNECTIONS: &str = "listener_connections";
111
112pub(crate) const RETRY_RESULT_SUCCESS: &str = "success";
114
115pub(crate) const RETRY_RESULT_EXHAUSTED: &str = "exhausted";
117
118pub(crate) const HEALTH_RESULT_HEALTHY: &str = "healthy";
120
121pub(crate) const HEALTH_RESULT_UNHEALTHY: &str = "unhealthy";
123
124pub(crate) const RELOAD_RESULT_SUCCESS: &str = "success";
126
127pub(crate) const RELOAD_RESULT_FAILURE: &str = "failure";
129
130const 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
152static PROMETHEUS_HANDLE: OnceLock<PrometheusHandle> = OnceLock::new();
158
159static LABEL_CONFIG: OnceLock<MetricLabelsConfig> = OnceLock::new();
161
162static ALL_LABELS: OnceLock<MetricLabelsConfig> = OnceLock::new();
164
165pub fn install_metric_labels(labels: MetricLabelsConfig) {
171 let _existing = LABEL_CONFIG.set(labels);
172}
173
174pub(crate) fn metric_labels() -> &'static MetricLabelsConfig {
176 LABEL_CONFIG
177 .get()
178 .unwrap_or_else(|| ALL_LABELS.get_or_init(MetricLabelsConfig::default))
179}
180
181fn 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
192fn label_if(enabled: bool, value: SharedString) -> Option<SharedString> {
194 enabled.then_some(value)
195}
196
197pub 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
228pub fn render_prometheus() -> Option<String> {
232 PROMETHEUS_HANDLE.get().map(PrometheusHandle::render)
233}
234
235pub(crate) fn is_recorder_installed() -> bool {
237 PROMETHEUS_HANDLE.get().is_some()
238}
239
240#[derive(Clone, Debug, Default, Eq, PartialEq)]
246pub struct StatsMetricsSnapshot {
247 pub http_active_by_listener: std::collections::HashMap<String, u64>,
249 pub http_active_aggregate: Option<u64>,
251 pub tcp_active_by_listener: std::collections::HashMap<String, u64>,
253 pub tcp_active_aggregate: Option<u64>,
255 pub upstream_requests_by_cluster: std::collections::HashMap<String, u64>,
257 pub upstream_requests_aggregate: Option<u64>,
259 pub connect_failures_by_cluster: std::collections::HashMap<String, u64>,
261 pub connect_failures_aggregate: Option<u64>,
263}
264
265#[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
317type PrometheusSample<'a> = (&'a str, std::collections::HashMap<String, String>, u64);
323
324fn 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
345fn 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
362fn 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
380fn 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
387pub 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
414pub 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
442pub(crate) struct RequestMetricLabels {
453 pub cluster: SharedString,
455 pub method: &'static str,
457 pub route: SharedString,
459 pub status_class: &'static str,
461}
462
463fn 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
490pub(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
504fn 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
529fn 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
556fn 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
580pub(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
606pub struct ActiveRequestGuard {
617 listener: SharedString,
619}
620
621impl ActiveRequestGuard {
622 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
647pub(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
663pub(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
686fn 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
696fn 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
706pub(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
718pub(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
734pub(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
772pub(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
784pub(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
796pub(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
820pub 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
839pub 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
853pub(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
871pub(crate) fn count_healthy_endpoints(health: &praxis_core::health::ClusterHealthEntry) -> (usize, usize) {
873 health.endpoint_counts()
874}
875
876pub 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
892pub 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
900pub(crate) fn cluster_none() -> SharedString {
908 SharedString::const_str("none")
909}
910
911pub(crate) fn route_unknown() -> SharedString {
915 SharedString::const_str("unknown")
916}
917
918#[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(®istry);
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}