massa_channel/
receiver.rs1use std::{
2 ops::{Deref, DerefMut},
3 sync::Arc,
4 time::{Duration, Instant},
5};
6
7use crossbeam::channel::{Receiver, RecvError, RecvTimeoutError, TryRecvError};
8use prometheus::{Counter, Gauge};
9use tracing::trace;
10
11#[derive(Clone)]
12pub struct MassaReceiver<T> {
13 pub(crate) receiver: Receiver<T>,
14 #[allow(dead_code)]
15 pub(crate) name: String,
16 pub(crate) actual_len: Gauge,
18 pub(crate) received: Counter,
20 pub(crate) ref_counter: Arc<()>,
22}
23
24impl<T> Drop for MassaReceiver<T> {
26 fn drop(&mut self) {
27 let ref_count = Arc::strong_count(&self.ref_counter);
28 if ref_count == 1 {
29 self.unregister_metrics();
31 }
32 }
33}
34
35impl<T> MassaReceiver<T> {
36 pub fn update_metrics(&self) {
40 self.actual_len.set(self.receiver.len() as f64);
43
44 self.received.inc();
45 }
46
47 fn unregister_metrics(&self) {
58 if let Err(e) = prometheus::unregister(Box::new(self.actual_len.clone())) {
59 trace!(
60 "promethetus error unregister actual_len for {} : {}",
61 self.name,
62 e
63 );
64 }
65
66 if let Err(e) = prometheus::unregister(Box::new(self.received.clone())) {
67 trace!(
68 "promethetus error unregister received for {} : {}",
69 self.name,
70 e
71 );
72 }
73 }
74
75 pub fn try_recv(&self) -> Result<T, TryRecvError> {
77 match self.receiver.try_recv() {
78 Ok(msg) => {
79 self.update_metrics();
80 Ok(msg)
81 }
82 Err(crossbeam::channel::TryRecvError::Empty) => Err(TryRecvError::Empty),
83 Err(crossbeam::channel::TryRecvError::Disconnected) => {
84 self.unregister_metrics();
85 Err(TryRecvError::Disconnected)
86 }
87 }
88 }
89
90 pub fn recv_deadline(&self, deadline: Instant) -> Result<T, RecvTimeoutError> {
91 match self.receiver.recv_deadline(deadline) {
92 Ok(msg) => {
93 self.update_metrics();
94 Ok(msg)
95 }
96 Err(RecvTimeoutError::Timeout) => Err(RecvTimeoutError::Timeout),
97 Err(RecvTimeoutError::Disconnected) => {
98 self.unregister_metrics();
99 Err(RecvTimeoutError::Disconnected)
100 }
101 }
102 }
103
104 pub fn recv_timeout(&self, timeout: Duration) -> Result<T, RecvTimeoutError> {
105 match self.receiver.recv_timeout(timeout) {
106 Ok(msg) => {
107 self.update_metrics();
108 Ok(msg)
109 }
110 Err(RecvTimeoutError::Timeout) => Err(RecvTimeoutError::Timeout),
111 Err(RecvTimeoutError::Disconnected) => {
112 self.unregister_metrics();
113 Err(RecvTimeoutError::Disconnected)
114 }
115 }
116 }
117
118 pub fn recv(&self) -> Result<T, RecvError> {
119 match self.receiver.recv() {
120 Ok(msg) => {
121 self.update_metrics();
122 Ok(msg)
123 }
124 Err(e) => {
125 self.unregister_metrics();
126 Err(e)
127 }
128 }
129 }
130}
131
132impl<T> Deref for MassaReceiver<T> {
133 type Target = Receiver<T>;
134
135 fn deref(&self) -> &Self::Target {
136 &self.receiver
137 }
138}
139
140impl<T> DerefMut for MassaReceiver<T> {
141 fn deref_mut(&mut self) -> &mut Self::Target {
142 &mut self.receiver
143 }
144}