massa_channel/
receiver.rs

1use 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    /// channel size
17    pub(crate) actual_len: Gauge,
18    /// total received messages
19    pub(crate) received: Counter,
20    /// reference counter to know how many receiver are cloned
21    pub(crate) ref_counter: Arc<()>,
22}
23
24/// implement drop on MassaReceiver
25impl<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            // this is the last ref so we can unregister metrics
30            self.unregister_metrics();
31        }
32    }
33}
34
35impl<T> MassaReceiver<T> {
36    /// increment manually the metrics
37    /// Should be used when using the receiver with select! macro
38    /// select! does not call recv()
39    pub fn update_metrics(&self) {
40        // use the len of the channel for actual_len instead of actual_len.dec()
41        // because for each send we call recv more than one time
42        self.actual_len.set(self.receiver.len() as f64);
43
44        self.received.inc();
45    }
46
47    /// Unregister the channel metrics.
48    ///
49    /// Called from every `recv*`/`try_recv` disconnect branch on purpose, even
50    /// though `MassaReceiver` is `Clone`: a crossbeam disconnect means all senders
51    /// are dropped and the buffer is empty, so the channel is dead for every clone
52    /// at once and the metrics are frozen. Unregistering right away also frees the
53    /// metric names for the relaunch (`NeedSync`) that re-creates the channels
54    /// with the same names; `MassaChannel::new` only logs registration failures
55    /// at debug level. Later calls (from other clones or from `Drop`) fail with
56    /// "not registered" and are swallowed at trace level.
57    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    /// attempt to receive a message from the channel
76    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}