1use std::net::SocketAddr;
2
3use hyper::service::make_service_fn;
4use hyper::{body::Body, header::CONTENT_TYPE, service::service_fn, Request, Response};
5use prometheus::{Encoder, TextEncoder};
6use tracing::{error, info};
7
8use crate::MetricsStopper;
9
10#[allow(dead_code)]
11pub(crate) fn bind_metrics(addr: SocketAddr) -> MetricsStopper {
12 let (tx, rx) = tokio::sync::oneshot::channel::<()>();
13 let handle = std::thread::spawn(move || {
14 let rt = tokio::runtime::Builder::new_current_thread()
15 .enable_all()
16 .build()
17 .expect("error on build tokio runtime for metrics server");
18
19 rt.block_on(async {
20 let server = hyper::Server::bind(&addr).serve(make_service_fn(|_| async {
21 Ok::<_, hyper::Error>(service_fn(serve_req))
22 }));
23
24 let graceful_server = server.with_graceful_shutdown(async {
25 rx.await.ok();
26 });
27 info!("METRICS | listening on http://{}", addr);
28 if let Err(e) = graceful_server.await {
29 error!("metrics server error: {}", e);
30 }
31 info!("METRICS | server stopped");
32 });
33 });
34 MetricsStopper {
35 stopper: Some(tx),
36 stop_handle: Some(handle),
37 }
38}
39
40#[allow(dead_code)]
41async fn serve_req(req: Request<Body>) -> Result<Response<Body>, hyper::Error> {
42 if req.uri().path() != "/metrics" {
43 Ok(Response::builder()
45 .status(404)
46 .body(Body::from("Not Found"))
47 .unwrap())
48 } else {
49 let encoder = TextEncoder::new();
50 let mut buffer = vec![];
51 encoder
52 .encode(&prometheus::gather(), &mut buffer)
53 .expect("Failed to encode metrics");
54
55 let response = Response::builder()
56 .status(200)
57 .header(CONTENT_TYPE, encoder.format_type())
58 .body(Body::from(buffer))
59 .unwrap();
60
61 Ok(response)
62 }
63}