1use std::sync::Arc;
18
19use receiver::MassaReceiver;
20use sender::MassaSender;
21
22pub mod receiver;
23pub mod sender;
24
25#[derive(Clone)]
26pub struct MassaChannel {}
27
28impl MassaChannel {
29 #[allow(clippy::new_ret_no_self)]
30 pub fn new<T>(name: String, capacity: Option<usize>) -> (MassaSender<T>, MassaReceiver<T>) {
31 use prometheus::{Counter, Gauge};
32
33 let (s, r) = if let Some(capacity) = capacity {
34 crossbeam::channel::bounded::<T>(capacity)
35 } else {
36 crossbeam::channel::unbounded::<T>()
37 };
38
39 let actual_len = Gauge::new(
42 format!("{}_channel_actual_size", name),
43 "Actual length of channel",
44 )
45 .expect("Failed to create gauge");
46
47 let received = Counter::new(
49 format!("{}_channel_total_receive", name),
50 "Total received messages",
51 )
52 .expect("Failed to create counter");
53
54 #[cfg(not(feature = "test-exports"))]
58 {
59 use tracing::debug;
60 if let Err(e) = prometheus::register(Box::new(actual_len.clone())) {
61 debug!("Failed to register actual_len gauge for {} : {}", name, e);
62 }
63
64 if let Err(e) = prometheus::register(Box::new(received.clone())) {
65 debug!("Failed to register received counter for {} : {}", name, e);
66 }
67 }
68
69 let sender = MassaSender {
70 sender: s,
71 name: name.clone(),
72 actual_len: actual_len.clone(),
73 };
74
75 let receiver = MassaReceiver {
76 receiver: r,
77 name,
78 actual_len,
79 received,
80 ref_counter: Arc::new(()),
81 };
82
83 (sender, receiver)
84 }
85}