mirror of
https://github.com/pezkuwichain/pezkuwi-subxt.git
synced 2026-06-13 07:01:05 +00:00
rpc server with HTTP/WS on the same socket (#12663)
* jsonrpsee v0.16 add backwards compatibility run old http server on http only * cargo fmt * update jsonrpsee 0.16.1 * less verbose cors log * fix nit in log: WS -> HTTP * revert needless changes in Cargo.lock * remove unused features in tower * fix nits; add client-core feature * jsonrpsee v0.16.2
This commit is contained in:
@@ -16,9 +16,9 @@
|
||||
// You should have received a copy of the GNU General Public License
|
||||
// along with this program. If not, see <https://www.gnu.org/licenses/>.
|
||||
|
||||
//! RPC middlware to collect prometheus metrics on RPC calls.
|
||||
//! RPC middleware to collect prometheus metrics on RPC calls.
|
||||
|
||||
use jsonrpsee::core::middleware::{Headers, HttpMiddleware, MethodKind, Params, WsMiddleware};
|
||||
use jsonrpsee::server::logger::{HttpRequest, Logger, MethodKind, Params, TransportProtocol};
|
||||
use prometheus_endpoint::{
|
||||
register, Counter, CounterVec, HistogramOpts, HistogramVec, Opts, PrometheusError, Registry,
|
||||
U64,
|
||||
@@ -54,9 +54,9 @@ pub struct RpcMetrics {
|
||||
calls_started: CounterVec<U64>,
|
||||
/// Number of calls completed.
|
||||
calls_finished: CounterVec<U64>,
|
||||
/// Number of Websocket sessions opened (Websocket only).
|
||||
/// Number of Websocket sessions opened.
|
||||
ws_sessions_opened: Option<Counter<U64>>,
|
||||
/// Number of Websocket sessions closed (Websocket only).
|
||||
/// Number of Websocket sessions closed.
|
||||
ws_sessions_closed: Option<Counter<U64>>,
|
||||
}
|
||||
|
||||
@@ -139,62 +139,61 @@ impl RpcMetrics {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
/// Middleware for RPC calls
|
||||
pub struct RpcMiddleware {
|
||||
metrics: RpcMetrics,
|
||||
transport_label: &'static str,
|
||||
}
|
||||
impl Logger for RpcMetrics {
|
||||
type Instant = std::time::Instant;
|
||||
|
||||
impl RpcMiddleware {
|
||||
/// Create a new [`RpcMiddleware`] with the provided [`RpcMetrics`].
|
||||
pub fn new(metrics: RpcMetrics, transport_label: &'static str) -> Self {
|
||||
Self { metrics, transport_label }
|
||||
fn on_connect(
|
||||
&self,
|
||||
_remote_addr: SocketAddr,
|
||||
_request: &HttpRequest,
|
||||
transport: TransportProtocol,
|
||||
) {
|
||||
if let TransportProtocol::WebSocket = transport {
|
||||
self.ws_sessions_opened.as_ref().map(|counter| counter.inc());
|
||||
}
|
||||
}
|
||||
|
||||
/// Called when a new JSON-RPC request comes to the server.
|
||||
fn on_request(&self) -> std::time::Instant {
|
||||
fn on_request(&self, transport: TransportProtocol) -> Self::Instant {
|
||||
let transport_label = transport_label_str(transport);
|
||||
let now = std::time::Instant::now();
|
||||
self.metrics.requests_started.with_label_values(&[self.transport_label]).inc();
|
||||
self.requests_started.with_label_values(&[transport_label]).inc();
|
||||
now
|
||||
}
|
||||
|
||||
/// Called on each JSON-RPC method call, batch requests will trigger `on_call` multiple times.
|
||||
fn on_call(&self, name: &str, params: Params, kind: MethodKind) {
|
||||
fn on_call(&self, name: &str, params: Params, kind: MethodKind, transport: TransportProtocol) {
|
||||
let transport_label = transport_label_str(transport);
|
||||
log::trace!(
|
||||
target: "rpc_metrics",
|
||||
"[{}] on_call name={} params={:?} kind={}",
|
||||
self.transport_label,
|
||||
transport_label,
|
||||
name,
|
||||
params,
|
||||
kind,
|
||||
);
|
||||
self.metrics
|
||||
.calls_started
|
||||
.with_label_values(&[self.transport_label, name])
|
||||
.inc();
|
||||
self.calls_started.with_label_values(&[transport_label, name]).inc();
|
||||
}
|
||||
|
||||
/// Called on each JSON-RPC method completion, batch requests will trigger `on_result` multiple
|
||||
/// times.
|
||||
fn on_result(&self, name: &str, success: bool, started_at: std::time::Instant) {
|
||||
fn on_result(
|
||||
&self,
|
||||
name: &str,
|
||||
success: bool,
|
||||
started_at: Self::Instant,
|
||||
transport: TransportProtocol,
|
||||
) {
|
||||
let transport_label = transport_label_str(transport);
|
||||
let micros = started_at.elapsed().as_micros();
|
||||
log::debug!(
|
||||
target: "rpc_metrics",
|
||||
"[{}] {} call took {} μs",
|
||||
self.transport_label,
|
||||
transport_label,
|
||||
name,
|
||||
micros,
|
||||
);
|
||||
self.metrics
|
||||
.calls_time
|
||||
.with_label_values(&[self.transport_label, name])
|
||||
.observe(micros as _);
|
||||
self.calls_time.with_label_values(&[transport_label, name]).observe(micros as _);
|
||||
|
||||
self.metrics
|
||||
.calls_finished
|
||||
self.calls_finished
|
||||
.with_label_values(&[
|
||||
self.transport_label,
|
||||
transport_label,
|
||||
name,
|
||||
// the label "is_error", so `success` should be regarded as false
|
||||
// and vice-versa to be registrered correctly.
|
||||
@@ -203,58 +202,23 @@ impl RpcMiddleware {
|
||||
.inc();
|
||||
}
|
||||
|
||||
/// Called once the JSON-RPC request is finished and response is sent to the output buffer.
|
||||
fn on_response(&self, result: &str, started_at: std::time::Instant) {
|
||||
log::trace!(target: "rpc_metrics", "[{}] on_response started_at={:?}", self.transport_label, started_at);
|
||||
log::trace!(target: "rpc_metrics::extra", "[{}] result={:?}", self.transport_label, result);
|
||||
self.metrics.requests_finished.with_label_values(&[self.transport_label]).inc();
|
||||
fn on_response(&self, result: &str, started_at: Self::Instant, transport: TransportProtocol) {
|
||||
let transport_label = transport_label_str(transport);
|
||||
log::trace!(target: "rpc_metrics", "[{}] on_response started_at={:?}", transport_label, started_at);
|
||||
log::trace!(target: "rpc_metrics::extra", "[{}] result={:?}", transport_label, result);
|
||||
self.requests_finished.with_label_values(&[transport_label]).inc();
|
||||
}
|
||||
|
||||
fn on_disconnect(&self, _remote_addr: SocketAddr, transport: TransportProtocol) {
|
||||
if let TransportProtocol::WebSocket = transport {
|
||||
self.ws_sessions_closed.as_ref().map(|counter| counter.inc());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl WsMiddleware for RpcMiddleware {
|
||||
type Instant = std::time::Instant;
|
||||
|
||||
fn on_connect(&self, _remote_addr: SocketAddr, _headers: &Headers) {
|
||||
self.metrics.ws_sessions_opened.as_ref().map(|counter| counter.inc());
|
||||
}
|
||||
|
||||
fn on_request(&self) -> Self::Instant {
|
||||
self.on_request()
|
||||
}
|
||||
|
||||
fn on_call(&self, name: &str, params: Params, kind: MethodKind) {
|
||||
self.on_call(name, params, kind)
|
||||
}
|
||||
|
||||
fn on_result(&self, name: &str, success: bool, started_at: Self::Instant) {
|
||||
self.on_result(name, success, started_at)
|
||||
}
|
||||
|
||||
fn on_response(&self, _result: &str, started_at: Self::Instant) {
|
||||
self.on_response(_result, started_at)
|
||||
}
|
||||
|
||||
fn on_disconnect(&self, _remote_addr: SocketAddr) {
|
||||
self.metrics.ws_sessions_closed.as_ref().map(|counter| counter.inc());
|
||||
}
|
||||
}
|
||||
|
||||
impl HttpMiddleware for RpcMiddleware {
|
||||
type Instant = std::time::Instant;
|
||||
|
||||
fn on_request(&self, _remote_addr: SocketAddr, _headers: &Headers) -> Self::Instant {
|
||||
self.on_request()
|
||||
}
|
||||
|
||||
fn on_call(&self, name: &str, params: Params, kind: MethodKind) {
|
||||
self.on_call(name, params, kind)
|
||||
}
|
||||
|
||||
fn on_result(&self, name: &str, success: bool, started_at: Self::Instant) {
|
||||
self.on_result(name, success, started_at)
|
||||
}
|
||||
|
||||
fn on_response(&self, _result: &str, started_at: Self::Instant) {
|
||||
self.on_response(_result, started_at)
|
||||
fn transport_label_str(t: TransportProtocol) -> &'static str {
|
||||
match t {
|
||||
TransportProtocol::Http => "http",
|
||||
TransportProtocol::WebSocket => "ws",
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user