use std::time::Instant;
use jsonrpsee::{
server::middleware::rpc::{layer::ResponseFuture, RpcServiceT},
MethodResponse,
};
#[derive(Clone)]
pub struct RpcMetricsMiddleware<S> {
service: S,
}
impl<S> RpcMetricsMiddleware<S> {
pub fn new(service: S) -> Self {
Self { service }
}
}
impl<'a, S> RpcServiceT<'a> for RpcMetricsMiddleware<S>
where
S: RpcServiceT<'a> + Send + Sync + Clone + 'static,
{
type Future = ResponseFuture<futures::future::BoxFuture<'a, MethodResponse>>;
fn call(&self, request: jsonrpsee::types::Request<'a>) -> Self::Future {
let service = self.service.clone();
let method = request.method_name().to_owned();
let start = Instant::now();
metrics::gauge!("rpc.active_requests").increment(1.0);
ResponseFuture::future(Box::pin(async move {
let response = service.call(request).await;
let duration = start.elapsed().as_secs_f64();
let status = if response.is_error() {
"error"
} else {
"success"
};
metrics::counter!("rpc.requests.total", "method" => method.clone(), "status" => status)
.increment(1);
metrics::histogram!("rpc.request.duration_seconds", "method" => method.clone())
.record(duration);
if response.is_error() {
let error_code = response
.as_error_code()
.map(|c| c.to_string())
.unwrap_or_else(|| "unknown".to_string());
metrics::counter!(
"rpc.errors.total",
"method" => method,
"error_code" => error_code
)
.increment(1);
}
metrics::gauge!("rpc.active_requests").decrement(1.0);
response
}))
}
}