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
265pub fn collect_stats_metrics(prometheus_text: &str) -> StatsMetricsSnapshot {
267 let mut snapshot = StatsMetricsSnapshot::default();
268 for line in prometheus_text.lines() {
269 let line = line.trim();
270 if line.is_empty() || line.starts_with('#') {
271 continue;
272 }
273 if let Some((name, labels, value)) = parse_prometheus_sample(line) {
274 match name {
275 HTTP_ACTIVE_REQUESTS => snapshot.record_http_active(&labels, value),
276 TCP_ACTIVE_CONNECTIONS => snapshot.record_tcp_active(&labels, value),
277 UPSTREAM_REQUESTS_TOTAL => snapshot.add_upstream_requests(&labels, value),
278 UPSTREAM_CONNECT_FAILURES_TOTAL => snapshot.add_connect_failures(&labels, value),
279 _ => {},
280 }
281 }
282 }
283 snapshot
284}
285
286impl StatsMetricsSnapshot {
287 fn record_http_active(&mut self, labels: &std::collections::HashMap<String, String>, value: u64) {
289 record_by_listener(
290 &mut self.http_active_by_listener,
291 &mut self.http_active_aggregate,
292 labels,
293 value,
294 );
295 }
296
297 fn record_tcp_active(&mut self, labels: &std::collections::HashMap<String, String>, value: u64) {
299 record_by_listener(
300 &mut self.tcp_active_by_listener,
301 &mut self.tcp_active_aggregate,
302 labels,
303 value,
304 );
305 }
306
307 fn add_upstream_requests(&mut self, labels: &std::collections::HashMap<String, String>, value: u64) {
309 accumulate_by_cluster(
310 &mut self.upstream_requests_by_cluster,
311 &mut self.upstream_requests_aggregate,
312 labels,
313 value,
314 );
315 }
316
317 fn add_connect_failures(&mut self, labels: &std::collections::HashMap<String, String>, value: u64) {
319 accumulate_by_cluster(
320 &mut self.connect_failures_by_cluster,
321 &mut self.connect_failures_aggregate,
322 labels,
323 value,
324 );
325 }
326}
327
328fn record_by_listener(
331 by_listener: &mut std::collections::HashMap<String, u64>,
332 aggregate: &mut Option<u64>,
333 labels: &std::collections::HashMap<String, String>,
334 value: u64,
335) {
336 match labels.get("listener") {
337 Some(listener) => {
338 by_listener.insert(listener.clone(), value);
339 },
340 None => *aggregate = Some(value),
341 }
342}
343
344fn accumulate_by_cluster(
347 by_cluster: &mut std::collections::HashMap<String, u64>,
348 aggregate: &mut Option<u64>,
349 labels: &std::collections::HashMap<String, String>,
350 value: u64,
351) {
352 match labels.get("cluster") {
353 Some(cluster) => {
354 *by_cluster.entry(cluster.clone()).or_insert(0) += value;
355 },
356 None => *aggregate = Some(aggregate.unwrap_or(0) + value),
357 }
358}
359
360type PrometheusSample<'a> = (&'a str, std::collections::HashMap<String, String>, u64);
366
367fn parse_prometheus_sample(line: &str) -> Option<PrometheusSample<'_>> {
369 let (name_and_labels, value_str) = line.rsplit_once(' ')?;
370 #[expect(
371 clippy::cast_possible_truncation,
372 clippy::cast_sign_loss,
373 reason = "Prometheus counter/gauge values are non-negative integers"
374 )]
375 let value = {
376 let parsed = value_str.parse::<f64>().ok()?;
377 parsed.round() as u64
378 };
379 let (name, labels) = if let Some((name, label_blob)) = name_and_labels.split_once('{') {
380 let label_blob = label_blob.strip_suffix('}')?;
381 (name, parse_prometheus_labels(label_blob))
382 } else {
383 (name_and_labels, std::collections::HashMap::new())
384 };
385 Some((name, labels, value))
386}
387
388fn parse_prometheus_labels(input: &str) -> std::collections::HashMap<String, String> {
393 let mut labels = std::collections::HashMap::new();
394 let mut rest = input.trim();
395 while !rest.is_empty() {
396 let (pair, tail) = split_prometheus_label_pair(rest);
397 if let Some((key, value)) = parse_prometheus_label_pair(pair) {
398 labels.insert(key, value);
399 }
400 rest = tail;
401 }
402 labels
403}
404
405fn split_prometheus_label_pair(input: &str) -> (&str, &str) {
407 let bytes = input.as_bytes();
408 let mut in_quotes = false;
409 for (index, byte) in bytes.iter().enumerate() {
410 match *byte {
411 b'"' => in_quotes = !in_quotes,
412 b',' if !in_quotes => {
413 let head = input.get(..index).unwrap_or(input);
414 let tail = input.get(index + 1..).unwrap_or("").trim_start();
415 return (head, tail);
416 },
417 _ => {},
418 }
419 }
420 (input, "")
421}
422
423fn parse_prometheus_label_pair(pair: &str) -> Option<(String, String)> {
425 let (key, value) = pair.split_once('=')?;
426 let value = value.strip_prefix('"')?.strip_suffix('"')?;
427 Some((key.to_owned(), value.to_owned()))
428}
429
430pub fn status_class(code: u16) -> &'static str {
447 match code {
448 100..=199 => "1xx",
449 200..=299 => "2xx",
450 300..=399 => "3xx",
451 400..=499 => "4xx",
452 500..=599 => "5xx",
453 _ => "unknown",
454 }
455}
456
457pub fn method_label(method: &str) -> &'static str {
471 match method {
472 "GET" => "GET",
473 "POST" => "POST",
474 "PUT" => "PUT",
475 "DELETE" => "DELETE",
476 "PATCH" => "PATCH",
477 "HEAD" => "HEAD",
478 "OPTIONS" => "OPTIONS",
479 "TRACE" => "TRACE",
480 "CONNECT" => "CONNECT",
481 _ => "OTHER",
482 }
483}
484
485pub(crate) struct RequestMetricLabels {
496 pub cluster: SharedString,
498 pub method: &'static str,
500 pub route: SharedString,
502 pub status_class: &'static str,
504}
505
506fn selected_request_labels(labels: RequestMetricLabels) -> Vec<Label> {
508 let selected = metric_labels();
509 let pairs = [
510 (
511 "method",
512 label_if(
513 selected.is_enabled(MetricLabel::Method),
514 SharedString::const_str(labels.method),
515 ),
516 ),
517 (
518 "status_class",
519 label_if(
520 selected.is_enabled(MetricLabel::StatusClass),
521 SharedString::const_str(labels.status_class),
522 ),
523 ),
524 ("route", label_if(selected.is_enabled(MetricLabel::Route), labels.route)),
525 (
526 "cluster",
527 label_if(selected.is_enabled(MetricLabel::Cluster), labels.cluster),
528 ),
529 ];
530 selected_labels(&pairs)
531}
532
533pub(crate) fn record_request_metrics(labels: RequestMetricLabels, duration_secs: f64) {
535 if !is_recorder_installed() {
536 return;
537 }
538 if !metric_labels().all_enabled() {
539 let emitted = selected_request_labels(labels);
540 counter!(HTTP_REQUESTS_TOTAL, emitted.clone()).increment(1);
541 histogram!(HTTP_REQUEST_DURATION_SECONDS, emitted).record(duration_secs);
542 return;
543 }
544 record_request_metrics_all_labels(labels, duration_secs);
545}
546
547fn record_request_metrics_all_labels(labels: RequestMetricLabels, duration_secs: f64) {
552 let cluster = labels.cluster;
553 let route = labels.route;
554 counter!(
555 HTTP_REQUESTS_TOTAL,
556 "method" => labels.method,
557 "status_class" => labels.status_class,
558 "route" => route.clone(),
559 "cluster" => cluster.clone()
560 )
561 .increment(1);
562 histogram!(
563 HTTP_REQUEST_DURATION_SECONDS,
564 "method" => labels.method,
565 "status_class" => labels.status_class,
566 "route" => route,
567 "cluster" => cluster
568 )
569 .record(duration_secs);
570}
571
572fn selected_body_labels(method: &'static str, status_class: &'static str, cluster: SharedString) -> Vec<Label> {
578 let selected = metric_labels();
579 let pairs = [
580 (
581 "method",
582 label_if(
583 selected.is_enabled(MetricLabel::Method),
584 SharedString::const_str(method),
585 ),
586 ),
587 (
588 "status_class",
589 label_if(
590 selected.is_enabled(MetricLabel::StatusClass),
591 SharedString::const_str(status_class),
592 ),
593 ),
594 ("cluster", label_if(selected.is_enabled(MetricLabel::Cluster), cluster)),
595 ];
596 selected_labels(&pairs)
597}
598
599fn record_body_size_all_labels(
601 method: &'static str,
602 status_class: &'static str,
603 cluster: SharedString,
604 request_bytes: f64,
605 response_bytes: f64,
606) {
607 histogram!(
608 HTTP_REQUEST_BODY_BYTES,
609 "method" => method,
610 "status_class" => status_class,
611 "cluster" => cluster.clone()
612 )
613 .record(request_bytes);
614 histogram!(
615 HTTP_RESPONSE_BODY_BYTES,
616 "method" => method,
617 "status_class" => status_class,
618 "cluster" => cluster
619 )
620 .record(response_bytes);
621}
622
623pub(crate) fn record_body_size_metrics(
625 method: &'static str,
626 status_class: &'static str,
627 cluster: SharedString,
628 request_body_bytes: u64,
629 response_body_bytes: u64,
630) {
631 if !is_recorder_installed() {
632 return;
633 }
634 #[expect(
635 clippy::cast_precision_loss,
636 reason = "body byte counts as histogram observations; exact integer precision not required"
637 )]
638 let (request_bytes, response_bytes) = (request_body_bytes as f64, response_body_bytes as f64);
639
640 if !metric_labels().all_enabled() {
641 let emitted = selected_body_labels(method, status_class, cluster);
642 histogram!(HTTP_REQUEST_BODY_BYTES, emitted.clone()).record(request_bytes);
643 histogram!(HTTP_RESPONSE_BODY_BYTES, emitted).record(response_bytes);
644 return;
645 }
646 record_body_size_all_labels(method, status_class, cluster, request_bytes, response_bytes);
647}
648
649pub struct ActiveRequestGuard {
660 listener: SharedString,
662}
663
664impl ActiveRequestGuard {
665 pub(crate) fn acquire(listener: SharedString) -> Self {
667 if is_recorder_installed() {
668 if metric_labels().is_enabled(MetricLabel::Listener) {
669 gauge!(HTTP_ACTIVE_REQUESTS, "listener" => listener.clone()).increment(1.0);
670 } else {
671 gauge!(HTTP_ACTIVE_REQUESTS).increment(1.0);
672 }
673 }
674 Self { listener }
675 }
676}
677
678impl Drop for ActiveRequestGuard {
679 fn drop(&mut self) {
680 if is_recorder_installed() {
681 if metric_labels().is_enabled(MetricLabel::Listener) {
682 gauge!(HTTP_ACTIVE_REQUESTS, "listener" => self.listener.clone()).decrement(1.0);
683 } else {
684 gauge!(HTTP_ACTIVE_REQUESTS).decrement(1.0);
685 }
686 }
687 }
688}
689
690pub(crate) fn record_error(error_type: &'static str) {
700 if !is_recorder_installed() {
701 return;
702 }
703 counter!(ERRORS_TOTAL, "type" => error_type).increment(1);
704}
705
706pub(crate) fn error_type_for(etype: &::pingora_core::ErrorType, source: &::pingora_core::ErrorSource) -> &'static str {
712 use ::pingora_core::ErrorSource::{Downstream, Internal, Unset};
713
714 if matches!(source, Downstream) {
715 return ERROR_TYPE_DOWNSTREAM;
716 }
717 if is_timeout(etype) {
718 return ERROR_TYPE_TIMEOUT;
719 }
720 if is_unreachable(etype) {
721 return ERROR_TYPE_UPSTREAM_UNAVAILABLE;
722 }
723 if matches!(source, Internal | Unset) {
724 return ERROR_TYPE_INTERNAL;
725 }
726 ERROR_TYPE_UPSTREAM_PROTOCOL
727}
728
729fn is_timeout(etype: &::pingora_core::ErrorType) -> bool {
731 use ::pingora_core::ErrorType::{ConnectTimedout, ReadTimedout, TLSHandshakeTimedout, WriteTimedout};
732
733 matches!(
734 etype,
735 ConnectTimedout | TLSHandshakeTimedout | ReadTimedout | WriteTimedout
736 )
737}
738
739fn is_unreachable(etype: &::pingora_core::ErrorType) -> bool {
741 use ::pingora_core::ErrorType::{BindError, ConnectError, ConnectNoRoute, ConnectRefused, SocketError};
742
743 matches!(
744 etype,
745 ConnectRefused | ConnectNoRoute | ConnectError | BindError | SocketError
746 )
747}
748
749pub(crate) fn record_overload_reject(reason: &'static str) {
755 if !is_recorder_installed() {
756 return;
757 }
758 counter!(OVERLOAD_REJECTS_TOTAL, "reason" => reason).increment(1);
759}
760
761pub(crate) fn record_upstream_connect_duration(cluster: SharedString, duration_secs: f64) {
767 if !is_recorder_installed() {
768 return;
769 }
770 if metric_labels().is_enabled(MetricLabel::Cluster) {
771 histogram!(UPSTREAM_CONNECT_DURATION_SECONDS, "cluster" => cluster).record(duration_secs);
772 } else {
773 histogram!(UPSTREAM_CONNECT_DURATION_SECONDS).record(duration_secs);
774 }
775}
776
777pub(crate) fn record_upstream_request(cluster: SharedString, endpoint: SharedString, status_class: &'static str) {
784 if !is_recorder_installed() {
785 return;
786 }
787 let selected = metric_labels();
788 if !selected.all_enabled() {
789 let pairs = [
790 ("cluster", label_if(selected.is_enabled(MetricLabel::Cluster), cluster)),
791 (
792 "endpoint",
793 label_if(selected.is_enabled(MetricLabel::Endpoint), endpoint),
794 ),
795 (
796 "status_class",
797 label_if(
798 selected.is_enabled(MetricLabel::StatusClass),
799 SharedString::const_str(status_class),
800 ),
801 ),
802 ];
803 counter!(UPSTREAM_REQUESTS_TOTAL, selected_labels(&pairs)).increment(1);
804 return;
805 }
806 counter!(
807 UPSTREAM_REQUESTS_TOTAL,
808 "cluster" => cluster,
809 "endpoint" => endpoint,
810 "status_class" => status_class
811 )
812 .increment(1);
813}
814
815pub(crate) fn record_upstream_connect_failure(cluster: SharedString) {
817 if !is_recorder_installed() {
818 return;
819 }
820 if metric_labels().is_enabled(MetricLabel::Cluster) {
821 counter!(UPSTREAM_CONNECT_FAILURES_TOTAL, "cluster" => cluster).increment(1);
822 } else {
823 counter!(UPSTREAM_CONNECT_FAILURES_TOTAL).increment(1);
824 }
825}
826
827pub(crate) fn record_upstream_retry(cluster: SharedString, result: &'static str) {
829 if !is_recorder_installed() {
830 return;
831 }
832 if metric_labels().is_enabled(MetricLabel::Cluster) {
833 counter!(UPSTREAM_RETRIES_TOTAL, "cluster" => cluster, "result" => result).increment(1);
834 } else {
835 counter!(UPSTREAM_RETRIES_TOTAL, "result" => result).increment(1);
836 }
837}
838
839pub(crate) fn set_upstream_endpoint_gauges(cluster: SharedString, healthy: usize, total: usize) {
853 if !is_recorder_installed() {
854 return;
855 }
856 #[expect(clippy::cast_precision_loss, reason = "endpoint counts fit f64 exactly below 2^53")]
857 {
858 gauge!(UPSTREAM_HEALTHY_ENDPOINTS, "cluster" => cluster.clone()).set(healthy as f64);
859 gauge!(UPSTREAM_TOTAL_ENDPOINTS, "cluster" => cluster).set(total as f64);
860 }
861}
862
863pub fn clear_stale_upstream_health_gauges<'a, P: IntoIterator<Item = &'a str>, C: IntoIterator<Item = &'a str>>(
868 previous_health_clusters: P,
869 current_health_clusters: C,
870) {
871 if !is_recorder_installed() {
872 return;
873 }
874 let current: std::collections::HashSet<&str> = current_health_clusters.into_iter().collect();
875 for name in previous_health_clusters {
876 if !current.contains(name) {
877 set_upstream_endpoint_gauges(SharedString::from(name.to_owned()), 0, 0);
878 }
879 }
880}
881
882pub fn seed_upstream_health_gauges(registry: &praxis_core::health::HealthRegistry) {
887 if !is_recorder_installed() {
888 return;
889 }
890 for (name, state) in registry.iter() {
891 let (healthy, total) = state.endpoint_counts();
892 set_upstream_endpoint_gauges(SharedString::from(name.as_ref().to_owned()), healthy, total);
893 }
894}
895
896pub(crate) fn record_health_transition(cluster: SharedString, result: &'static str, healthy: usize, total: usize) {
898 if !is_recorder_installed() {
899 return;
900 }
901 if metric_labels().is_enabled(MetricLabel::Cluster) {
902 counter!(
903 UPSTREAM_HEALTH_TRANSITIONS_TOTAL,
904 "cluster" => cluster.clone(),
905 "result" => result
906 )
907 .increment(1);
908 } else {
909 counter!(UPSTREAM_HEALTH_TRANSITIONS_TOTAL, "result" => result).increment(1);
910 }
911 set_upstream_endpoint_gauges(cluster, healthy, total);
912}
913
914pub(crate) fn count_healthy_endpoints(health: &praxis_core::health::ClusterHealthEntry) -> (usize, usize) {
916 health.endpoint_counts()
917}
918
919pub fn record_config_reload_success() {
925 if !is_recorder_installed() {
926 return;
927 }
928 counter!(CONFIG_RELOAD_TOTAL, "result" => RELOAD_RESULT_SUCCESS).increment(1);
929 let ts = std::time::SystemTime::now()
930 .duration_since(std::time::UNIX_EPOCH)
931 .map_or(0.0, |d| d.as_secs_f64());
932 gauge!(CONFIG_RELOAD_LAST_SUCCESS_TIMESTAMP).set(ts);
933}
934
935pub fn record_config_reload_failure() {
937 if !is_recorder_installed() {
938 return;
939 }
940 counter!(CONFIG_RELOAD_TOTAL, "result" => RELOAD_RESULT_FAILURE).increment(1);
941}
942
943pub(crate) fn cluster_none() -> SharedString {
951 SharedString::const_str("none")
952}
953
954pub(crate) fn route_unknown() -> SharedString {
958 SharedString::const_str("unknown")
959}
960
961#[cfg(test)]
966#[expect(clippy::allow_attributes, reason = "blanket test suppressions")]
967#[allow(clippy::unwrap_used, clippy::expect_used, clippy::indexing_slicing, reason = "tests")]
968mod tests {
969 use super::*;
970
971 #[test]
972 fn status_class_1xx() {
973 assert_eq!(status_class(100), "1xx", "100 should be 1xx");
974 assert_eq!(status_class(199), "1xx", "199 should be 1xx");
975 }
976
977 #[test]
978 fn status_class_2xx() {
979 assert_eq!(status_class(200), "2xx", "200 should be 2xx");
980 assert_eq!(status_class(204), "2xx", "204 should be 2xx");
981 assert_eq!(status_class(299), "2xx", "299 should be 2xx");
982 }
983
984 #[test]
985 fn status_class_3xx() {
986 assert_eq!(status_class(301), "3xx", "301 should be 3xx");
987 assert_eq!(status_class(399), "3xx", "399 should be 3xx");
988 }
989
990 #[test]
991 fn status_class_4xx() {
992 assert_eq!(status_class(400), "4xx", "400 should be 4xx");
993 assert_eq!(status_class(404), "4xx", "404 should be 4xx");
994 assert_eq!(status_class(499), "4xx", "499 should be 4xx");
995 }
996
997 #[test]
998 fn status_class_5xx() {
999 assert_eq!(status_class(500), "5xx", "500 should be 5xx");
1000 assert_eq!(status_class(503), "5xx", "503 should be 5xx");
1001 assert_eq!(status_class(599), "5xx", "599 should be 5xx");
1002 }
1003
1004 #[test]
1005 fn status_class_zero_is_unknown() {
1006 assert_eq!(status_class(0), "unknown", "0 should be unknown");
1007 }
1008
1009 #[test]
1010 fn status_class_out_of_range_is_unknown() {
1011 assert_eq!(status_class(600), "unknown", "600 should be unknown");
1012 assert_eq!(status_class(99), "unknown", "99 should be unknown");
1013 }
1014
1015 #[test]
1016 fn method_label_standard_methods() {
1017 for m in [
1018 "GET", "POST", "PUT", "DELETE", "PATCH", "HEAD", "OPTIONS", "TRACE", "CONNECT",
1019 ] {
1020 assert_eq!(method_label(m), m, "{m} should pass through");
1021 }
1022 }
1023
1024 #[test]
1025 fn method_label_custom_methods_collapse_to_other() {
1026 assert_eq!(method_label("PURGE"), "OTHER", "PURGE should be OTHER");
1027 assert_eq!(method_label("FOOBAR"), "OTHER", "FOOBAR should be OTHER");
1028 assert_eq!(method_label(""), "OTHER", "empty should be OTHER");
1029 }
1030
1031 #[test]
1032 fn record_utilities_noop_without_recorder() {
1033 record_overload_reject(OVERLOAD_REASON_MEMORY);
1034 record_upstream_connect_failure(cluster_none());
1035 record_error(ERROR_TYPE_INTERNAL);
1036 record_upstream_request(cluster_none(), SharedString::const_str("10.0.0.1:80"), "2xx");
1037 record_upstream_retry(cluster_none(), RETRY_RESULT_SUCCESS);
1038 record_upstream_connect_duration(cluster_none(), 0.01);
1039 set_upstream_endpoint_gauges(cluster_none(), 1, 2);
1040 record_health_transition(cluster_none(), HEALTH_RESULT_HEALTHY, 1, 2);
1041 record_config_reload_success();
1042 record_config_reload_failure();
1043 clear_stale_upstream_health_gauges(["gone"], std::iter::empty::<&str>());
1044 let _request_guard = ActiveRequestGuard::acquire(SharedString::const_str("test"));
1045 }
1046
1047 #[test]
1048 fn active_request_guard_returns_to_zero_on_drop() {
1049 install_prometheus_recorder();
1050 let listener = SharedString::const_str("active-request-guard-listener");
1051 let guard = ActiveRequestGuard::acquire(listener.clone());
1052 let held = render_prometheus().expect("recorder should render");
1053 assert!(
1054 held.contains("praxis_http_active_requests{listener=\"active-request-guard-listener\"} 1"),
1055 "gauge should read 1 while the guard is held:\n{held}"
1056 );
1057 drop(guard);
1058 let released = render_prometheus().expect("recorder should render");
1059 assert!(
1060 released.contains("praxis_http_active_requests{listener=\"active-request-guard-listener\"} 0"),
1061 "gauge should return to 0 once the guard drops:\n{released}"
1062 );
1063 }
1064
1065 #[test]
1066 fn overload_reject_reasons_appear_in_scrape() {
1067 install_prometheus_recorder();
1068 record_overload_reject(OVERLOAD_REASON_MEMORY);
1069 record_overload_reject(OVERLOAD_REASON_GLOBAL_CONNECTIONS);
1070 record_overload_reject(OVERLOAD_REASON_LISTENER_CONNECTIONS);
1071 let body = render_prometheus().expect("recorder should render");
1072 for reason in [
1073 OVERLOAD_REASON_MEMORY,
1074 OVERLOAD_REASON_GLOBAL_CONNECTIONS,
1075 OVERLOAD_REASON_LISTENER_CONNECTIONS,
1076 ] {
1077 let needle = format!("praxis_overload_rejects_total{{reason=\"{reason}\"}}");
1078 assert!(body.contains(&needle), "expected `{needle}` in scrape:\n{body}");
1079 }
1080 }
1081
1082 #[test]
1083 fn upstream_requests_carry_cluster_endpoint_and_status_class() {
1084 install_prometheus_recorder();
1085 record_upstream_request(
1086 SharedString::const_str("api"),
1087 SharedString::const_str("10.0.0.7:8080"),
1088 "5xx",
1089 );
1090 let body = render_prometheus().expect("recorder should render");
1091 assert!(
1092 body.contains(
1093 "praxis_upstream_requests_total{cluster=\"api\",endpoint=\"10.0.0.7:8080\",status_class=\"5xx\"} 1"
1094 ),
1095 "counter should carry all three labels:\n{body}"
1096 );
1097 }
1098
1099 #[test]
1100 fn selected_labels_drops_disabled_dimensions() {
1101 let pairs = [
1102 ("method", Some(SharedString::const_str("GET"))),
1103 ("route", None),
1104 ("cluster", Some(SharedString::const_str("api"))),
1105 ];
1106 let emitted = selected_labels(&pairs);
1107 let names: Vec<&str> = emitted.iter().map(Label::key).collect();
1108 assert_eq!(names, vec!["method", "cluster"], "a disabled dimension must be absent");
1109 }
1110
1111 #[test]
1112 fn selected_labels_preserves_order_and_values() {
1113 let pairs = [
1114 ("cluster", Some(SharedString::const_str("api"))),
1115 ("endpoint", Some(SharedString::const_str("10.0.0.1:80"))),
1116 ];
1117 let emitted = selected_labels(&pairs);
1118 let rendered: Vec<(&str, &str)> = emitted.iter().map(|l| (l.key(), l.value())).collect();
1119 assert_eq!(
1120 rendered,
1121 vec![("cluster", "api"), ("endpoint", "10.0.0.1:80")],
1122 "enabled dimensions keep their order and values"
1123 );
1124 }
1125
1126 #[test]
1127 fn label_if_gates_on_the_flag() {
1128 assert_eq!(
1129 label_if(true, SharedString::const_str("x")).as_deref(),
1130 Some("x"),
1131 "an enabled dimension keeps its value"
1132 );
1133 assert_eq!(
1134 label_if(false, SharedString::const_str("x")),
1135 None,
1136 "a disabled dimension yields no value"
1137 );
1138 }
1139
1140 #[test]
1141 fn metric_labels_default_to_all_enabled() {
1142 assert!(
1143 metric_labels().all_enabled(),
1144 "without an explicit install every dimension must stay on, so the \
1145 recorders keep their allocation-free fast path"
1146 );
1147 }
1148
1149 #[test]
1150 fn error_types_appear_in_scrape() {
1151 install_prometheus_recorder();
1152 for error_type in [
1153 ERROR_TYPE_FILTER_REJECT,
1154 ERROR_TYPE_TIMEOUT,
1155 ERROR_TYPE_UPSTREAM_UNAVAILABLE,
1156 ERROR_TYPE_UPSTREAM_PROTOCOL,
1157 ERROR_TYPE_DOWNSTREAM,
1158 ERROR_TYPE_INTERNAL,
1159 ] {
1160 record_error(error_type);
1161 let body = render_prometheus().expect("recorder should render");
1162 let needle = format!("praxis_errors_total{{type=\"{error_type}\"}}");
1163 assert!(body.contains(&needle), "expected `{needle}` in scrape:\n{body}");
1164 }
1165 }
1166
1167 #[test]
1168 fn error_type_for_maps_pingora_errors_to_bounded_values() {
1169 use ::pingora_core::{ErrorSource, ErrorType};
1170
1171 assert_eq!(
1172 error_type_for(&ErrorType::ConnectTimedout, &ErrorSource::Upstream),
1173 ERROR_TYPE_TIMEOUT,
1174 "connect timeout is a timeout"
1175 );
1176 assert_eq!(
1177 error_type_for(&ErrorType::ConnectRefused, &ErrorSource::Upstream),
1178 ERROR_TYPE_UPSTREAM_UNAVAILABLE,
1179 "a refused connect means the upstream was unreachable"
1180 );
1181 assert_eq!(
1182 error_type_for(&ErrorType::ReadError, &ErrorSource::Upstream),
1183 ERROR_TYPE_UPSTREAM_PROTOCOL,
1184 "a mid-exchange read error is a protocol failure"
1185 );
1186 assert_eq!(
1187 error_type_for(&ErrorType::ReadTimedout, &ErrorSource::Downstream),
1188 ERROR_TYPE_DOWNSTREAM,
1189 "downstream source wins over the error kind"
1190 );
1191 assert_eq!(
1192 error_type_for(&ErrorType::InternalError, &ErrorSource::Internal),
1193 ERROR_TYPE_INTERNAL,
1194 "internal source is an internal fault"
1195 );
1196 }
1197
1198 #[test]
1199 fn body_size_histograms_use_byte_buckets() {
1200 install_prometheus_recorder();
1201 record_body_size_metrics("GET", "2xx", cluster_none(), 500, 4_000);
1202 let body = render_prometheus().expect("recorder should render");
1203 assert!(
1204 body.contains("praxis_http_request_body_bytes_bucket") && body.contains("le=\"1024\""),
1205 "request body histogram should use byte buckets, not duration defaults:\n{body}"
1206 );
1207 assert!(
1208 body.contains("praxis_http_response_body_bytes_bucket") && body.contains("le=\"4096\""),
1209 "response body histogram should use byte buckets:\n{body}"
1210 );
1211 assert!(
1212 !body.contains("praxis_http_request_body_bytes_bucket{le=\"0.005\"}")
1213 && !body.contains("praxis_http_request_body_bytes_bucket{method=\"GET\",status_class=\"2xx\",cluster=\"\",le=\"0.005\"}"),
1214 "request body histogram must not use duration default buckets:\n{body}"
1215 );
1216 }
1217
1218 #[test]
1219 fn clear_stale_upstream_health_gauges_zeros_removed_clusters() {
1220 install_prometheus_recorder();
1221 set_upstream_endpoint_gauges(SharedString::from("old-cluster".to_owned()), 2, 3);
1222 set_upstream_endpoint_gauges(SharedString::from("kept-cluster".to_owned()), 1, 1);
1223 clear_stale_upstream_health_gauges(["old-cluster", "kept-cluster"], ["kept-cluster"]);
1224 let body = render_prometheus().expect("recorder should render");
1225 assert!(
1226 body.contains("praxis_upstream_healthy_endpoints{cluster=\"old-cluster\"} 0"),
1227 "removed cluster healthy gauge should be zeroed:\n{body}"
1228 );
1229 assert!(
1230 body.contains("praxis_upstream_total_endpoints{cluster=\"old-cluster\"} 0"),
1231 "removed cluster total gauge should be zeroed:\n{body}"
1232 );
1233 assert!(
1234 body.contains("praxis_upstream_healthy_endpoints{cluster=\"kept-cluster\"} 1"),
1235 "kept cluster should retain its value:\n{body}"
1236 );
1237 }
1238
1239 #[test]
1240 fn collect_stats_metrics_sums_cluster_counters() {
1241 let text = r#"
1242praxis_http_active_requests{listener="web"} 2
1243praxis_tcp_active_connections{listener="tcp-in"} 1
1244praxis_upstream_requests_total{cluster="backend",endpoint="127.0.0.1:1",status_class="2xx"} 3
1245praxis_upstream_requests_total{cluster="backend",endpoint="127.0.0.1:2",status_class="5xx"} 1
1246praxis_upstream_connect_failures_total{cluster="backend"} 2
1247"#;
1248 let snap = collect_stats_metrics(text);
1249 assert_eq!(
1250 snap.http_active_by_listener.get("web"),
1251 Some(&2),
1252 "HTTP active per listener should parse"
1253 );
1254 assert_eq!(
1255 snap.tcp_active_by_listener.get("tcp-in"),
1256 Some(&1),
1257 "TCP active per listener should parse"
1258 );
1259 assert_eq!(
1260 snap.upstream_requests_by_cluster.get("backend"),
1261 Some(&4),
1262 "upstream requests should sum by cluster"
1263 );
1264 assert_eq!(
1265 snap.connect_failures_by_cluster.get("backend"),
1266 Some(&2),
1267 "connect failures should parse by cluster"
1268 );
1269 }
1270
1271 #[test]
1272 fn parse_prometheus_labels_handles_commas_inside_quoted_values() {
1273 let labels = parse_prometheus_labels(r#"tag="a,b",listener="web""#);
1274 assert_eq!(labels.get("tag"), Some(&"a,b".to_owned()), "comma inside quotes");
1275 assert_eq!(labels.get("listener"), Some(&"web".to_owned()), "second label");
1276 }
1277
1278 #[test]
1279 fn collect_stats_metrics_parses_unlabeled_upstream_counters() {
1280 let text = r#"
1281praxis_upstream_requests_total{status_class="2xx"} 5
1282praxis_upstream_connect_failures_total 2
1283"#;
1284 let snap = collect_stats_metrics(text);
1285 assert_eq!(
1286 snap.upstream_requests_aggregate,
1287 Some(5),
1288 "unlabeled upstream requests should aggregate"
1289 );
1290 assert_eq!(
1291 snap.connect_failures_aggregate,
1292 Some(2),
1293 "unlabeled connect failures should aggregate"
1294 );
1295 }
1296
1297 #[test]
1298 fn seed_upstream_health_gauges_publishes_registry_counts() {
1299 use std::sync::Arc;
1300
1301 use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
1302
1303 install_prometheus_recorder();
1304 let endpoints = vec![EndpointHealth::new(), EndpointHealth::new()];
1305 endpoints[0].mark_unhealthy();
1306 let entry = Arc::new(ClusterHealthEntry::new(
1307 endpoints,
1308 vec![Arc::from("a:1"), Arc::from("b:1")],
1309 None,
1310 None,
1311 ));
1312 let registry = Arc::new([(Arc::from("backend"), entry)].into_iter().collect());
1313 seed_upstream_health_gauges(®istry);
1314 let body = render_prometheus().expect("recorder should render");
1315 assert!(
1316 body.contains("praxis_upstream_healthy_endpoints{cluster=\"backend\"} 1"),
1317 "seed should publish healthy count:\n{body}"
1318 );
1319 assert!(
1320 body.contains("praxis_upstream_total_endpoints{cluster=\"backend\"} 2"),
1321 "seed should publish total count:\n{body}"
1322 );
1323 }
1324
1325 #[test]
1326 fn count_healthy_endpoints_counts_correctly() {
1327 use std::sync::Arc;
1328
1329 use praxis_core::health::{ClusterHealthEntry, EndpointHealth};
1330
1331 let endpoints = vec![EndpointHealth::new(), EndpointHealth::new(), EndpointHealth::new()];
1332 endpoints[1].mark_unhealthy();
1333 let entry = ClusterHealthEntry::new(
1334 endpoints,
1335 vec![Arc::from("a:1"), Arc::from("b:1"), Arc::from("c:1")],
1336 None,
1337 None,
1338 );
1339 let (healthy, total) = count_healthy_endpoints(&entry);
1340 assert_eq!(total, 3, "total should be 3");
1341 assert_eq!(healthy, 2, "two endpoints should be healthy");
1342 }
1343}