massa_channel/
sender.rs

1use std::{
2    ops::Deref,
3    time::{Duration, Instant},
4};
5
6use crossbeam::channel::{SendError, SendTimeoutError, Sender, TrySendError};
7use prometheus::Gauge;
8
9#[derive(Clone, Debug)]
10pub struct MassaSender<T> {
11    pub(crate) sender: Sender<T>,
12    #[allow(dead_code)]
13    pub(crate) name: String,
14    /// channel size
15    pub(crate) actual_len: Gauge,
16}
17
18impl<T> MassaSender<T> {
19    /// Send a message to the channel
20    pub fn send(&self, msg: T) -> Result<(), SendError<T>> {
21        match self.sender.send(msg) {
22            Ok(()) => {
23                self.actual_len.inc();
24                Ok(())
25            }
26            Err(e) => Err(e),
27        }
28    }
29
30    pub fn send_timeout(&self, msg: T, duration: Duration) -> Result<(), SendTimeoutError<T>> {
31        match self.sender.send_timeout(msg, duration) {
32            Ok(()) => {
33                self.actual_len.inc();
34                Ok(())
35            }
36            Err(e) => Err(e),
37        }
38    }
39
40    pub fn send_deadline(&self, msg: T, deadline: Instant) -> Result<(), SendTimeoutError<T>> {
41        match self.sender.send_deadline(msg, deadline) {
42            Ok(()) => {
43                self.actual_len.inc();
44                Ok(())
45            }
46            Err(e) => Err(e),
47        }
48    }
49
50    pub fn try_send(&self, msg: T) -> Result<(), TrySendError<T>> {
51        match self.sender.try_send(msg) {
52            Ok(()) => {
53                self.actual_len.inc();
54                Ok(())
55            }
56            Err(e) => Err(e),
57        }
58    }
59}
60
61impl<T> Deref for MassaSender<T> {
62    type Target = Sender<T>;
63
64    fn deref(&self) -> &Self::Target {
65        &self.sender
66    }
67}