use crate::config;
use crate::shutdown;
use lazy_static::lazy_static;
use prometheus::{gather, Encoder, IntCounterVec, IntGaugeVec, Opts, Registry, TextEncoder};
use std::net::{Ipv4Addr, Ipv6Addr, SocketAddr};
use tokio::sync::mpsc;
use tracing::{error, info};
use warp::{Filter, Rejection, Reply};
const DEFAULT_PORT: u16 = 8000;
lazy_static! {
pub static ref REGISTRY: Registry = Registry::new();
pub static ref VERSION_GAUGE: IntGaugeVec =
IntGaugeVec::new(
Opts::new("version", "Version info of the service.").namespace(config::SERVICE_NAME).subsystem(config::NAME),
&["major", "minor", "git_version", "git_commit", "platform", "build_time"]
).expect("metric can be created");
pub static ref DOWNLOAD_PEER_COUNT: IntCounterVec =
IntCounterVec::new(
Opts::new("download_peer_total", "Counter of the number of the download peer.").namespace(config::SERVICE_NAME).subsystem(config::NAME),
&["task_type"]
).expect("metric can be created");
}
#[derive(Debug)]
pub struct Metrics {
addr: SocketAddr,
shutdown: shutdown::Shutdown,
_shutdown_complete: mpsc::UnboundedSender<()>,
}
impl Metrics {
pub fn new(
enable_ipv6: bool,
shutdown: shutdown::Shutdown,
shutdown_complete_tx: mpsc::UnboundedSender<()>,
) -> Self {
let addr = if enable_ipv6 {
SocketAddr::new(Ipv6Addr::UNSPECIFIED.into(), DEFAULT_PORT)
} else {
SocketAddr::new(Ipv4Addr::UNSPECIFIED.into(), DEFAULT_PORT)
};
Self {
addr,
shutdown,
_shutdown_complete: shutdown_complete_tx,
}
}
pub async fn run(&mut self) {
self.register_custom_metrics();
let metrics_route = warp::path!("metrics")
.and(warp::get())
.and(warp::path::end())
.and_then(Self::metrics_handler);
tokio::select! {
_ = warp::serve(metrics_route).run(self.addr) => {
info!("metrics server ended");
}
_ = self.shutdown.recv() => {
info!("metrics server shutting down");
}
}
}
fn register_custom_metrics(&self) {
REGISTRY
.register(Box::new(VERSION_GAUGE.clone()))
.expect("metric can be registered");
REGISTRY
.register(Box::new(DOWNLOAD_PEER_COUNT.clone()))
.expect("metric can be registered");
}
async fn metrics_handler() -> Result<impl Reply, Rejection> {
let encoder = TextEncoder::new();
let mut buf = Vec::new();
if let Err(err) = encoder.encode(®ISTRY.gather(), &mut buf) {
error!("could not encode custom metrics: {}", err);
};
let mut res = match String::from_utf8(buf.clone()) {
Ok(v) => v,
Err(err) => {
error!("custom metrics could not be from_utf8'd: {}", err);
String::default()
}
};
buf.clear();
let mut buf = Vec::new();
if let Err(err) = encoder.encode(&gather(), &mut buf) {
error!("could not encode prometheus metrics: {}", err);
};
let res_custom = match String::from_utf8(buf.clone()) {
Ok(v) => v,
Err(err) => {
error!("prometheus metrics could not be from_utf8'd: {}", err);
String::default()
}
};
buf.clear();
res.push_str(&res_custom);
Ok(res)
}
}