use prometheus::{Encoder, IntCounter, TextEncoder};
use std::net::SocketAddr;
use warp::Filter;
use crate::metrics_helper::MetricsCounter;
pub async fn start_metrics_server(addr: impl Into<SocketAddr>) {
let metrics_route = metrics_filter();
warp::serve(metrics_route).run(addr).await;
}
fn metrics_filter() -> impl Filter<Extract = impl warp::Reply, Error = warp::Rejection> + Clone {
warp::path("metrics").map(|| {
let encoder = TextEncoder::new();
let metric_families = prometheus::gather();
let mut buffer = Vec::new();
encoder
.encode(&metric_families, &mut buffer)
.expect("Failed to encode metrics");
warp::reply::with_header(buffer, "Content-Type", encoder.format_type())
})
}
impl MetricsCounter for IntCounter {
fn inc_by(&self, amount: u64) {
prometheus::IntCounter::inc_by(self, amount);
}
}
impl<T: MetricsCounter> MetricsCounter for &T {
fn inc_by(&self, amount: u64) {
(**self).inc_by(amount)
}
}
#[cfg(test)]
mod tests {
use super::*;
use prometheus::{register_int_counter, TextEncoder};
use warp::http::StatusCode;
#[tokio::test]
async fn test_metrics_endpoint() {
let counter = register_int_counter!(
"connections_accepted_total",
"Total number of connections accepted"
)
.expect("Failed to register counter");
counter.inc();
let filter = metrics_filter();
let response = warp::test::request().path("/metrics").reply(&filter).await;
assert_eq!(response.status(), StatusCode::OK);
let content_type = response
.headers()
.get("Content-Type")
.expect("Missing Content-Type header")
.to_str()
.expect("Invalid Content-Type header");
assert_eq!(content_type, TextEncoder::new().format_type());
let body = std::str::from_utf8(response.body()).expect("Response body is not valid UTF-8");
assert!(
body.contains("connections_accepted_total"),
"Metrics output did not contain expected counter"
);
}
}