mirror of
https://github.com/hl-archive-node/nanoreth.git
synced 2025-12-06 19:09:54 +00:00
feat: add metrics for tx channel (#3345)
This commit is contained in:
@ -1,9 +1,13 @@
|
||||
//! Support for metering senders. Facilitates debugging by exposing metrics for number of messages
|
||||
//! sent, number of errors, etc.
|
||||
|
||||
use futures::Stream;
|
||||
use metrics::Counter;
|
||||
use reth_metrics_derive::Metrics;
|
||||
use std::task::{ready, Context, Poll};
|
||||
use std::{
|
||||
pin::Pin,
|
||||
task::{ready, Context, Poll},
|
||||
};
|
||||
use tokio::sync::mpsc::{
|
||||
self,
|
||||
error::{SendError, TryRecvError, TrySendError},
|
||||
@ -115,6 +119,14 @@ impl<T> UnboundedMeteredReceiver<T> {
|
||||
}
|
||||
}
|
||||
|
||||
impl<T> Stream for UnboundedMeteredReceiver<T> {
|
||||
type Item = T;
|
||||
|
||||
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
|
||||
self.receiver.poll_recv(cx)
|
||||
}
|
||||
}
|
||||
|
||||
/// A wrapper type around [Sender](mpsc::Sender) that updates metrics on send.
|
||||
#[derive(Debug)]
|
||||
pub struct MeteredSender<T> {
|
||||
@ -215,6 +227,14 @@ impl<T> MeteredReceiver<T> {
|
||||
}
|
||||
}
|
||||
|
||||
impl<T> Stream for MeteredReceiver<T> {
|
||||
type Item = T;
|
||||
|
||||
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
|
||||
self.receiver.poll_recv(cx)
|
||||
}
|
||||
}
|
||||
|
||||
/// Throughput metrics for [MeteredSender]
|
||||
#[derive(Clone, Metrics)]
|
||||
#[metrics(dynamic = true)]
|
||||
|
||||
Reference in New Issue
Block a user