Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_core/src/channels.rs

5.7 KiB, 7 runs

created by r1870400018:70, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! Provides added convenience and ergonomics when using an existing channel implementation,
2//! currently [crossbeam-channel](https://crates.io/crates/crossbeam-channel). `Simplex`
3//! packages a transmission and receiver pair for communicating in one direction, while
4//! `FullDuplex` packages two of these for simultaneous bi-directional communication.
5
6use crate::{
7 prelude::*,
8 //bot::CtrlMsg,
9 time::wait_for_true,
10};
11
12use std::{
13 fmt::Debug,
14 sync::{
15 Arc,
16 RwLock,
17 },
18 time::{
19 Duration,
20 Instant,
21 },
22};
23
24pub use flume::{
25 unbounded,
26 Sender,
27 Receiver,
28 TryRecvError,
29 RecvTimeoutError,
30};
31
32pub fn full_duplex<M>() -> FullDuplex<M> {
33 FullDuplex (
34 simplex(),
35 simplex(),
36 )
37}
38
39pub fn simplex<M>() -> Simplex<M> {
40 let (tx, rx) = unbounded();
41 Simplex {
42 tx: tx,
43 rx: rx,
44 open: Arc::new(RwLock::new(true)),
45 }
46}
47
48#[derive(Debug)]
49/// A channel for communicating in a single direction. It includes a thread-safe count of
50/// pending messages.
51///
52/// ```ignore
53///
54/// tx ----->----- rx A simplex channel is simple, just a
55/// transmitter end (tx) and a receiver
56/// end (rx).
57///
58/// ```
59pub struct Simplex<M> {
60 pub tx: Sender<M>,
61 pub rx: Receiver<M>,
62 open: Arc<RwLock<bool>>, // Is the channel open for use?
63}
64
65impl<M> Clone for Simplex<M> {
66 fn clone(&self) -> Self {
67 Self {
68 tx: self.tx.clone(),
69 rx: self.rx.clone(),
70 open: self.open.clone(),
71 }
72 }
73}
74
75impl<M> Default for Simplex<M> {
76 fn default() -> Self {
77 simplex::<M>()
78 }
79}
80
81impl<M> Simplex<M> {
82
83 pub fn tx(&self) -> &Sender<M> { &self.tx }
84 pub fn rx(&self) -> &Receiver<M> { &self.rx }
85
86}
87
88#[derive(Debug)]
89pub enum Recv<M> {
90 Result(Outcome<M>),
91 Empty,
92}
93
94impl<M: 'static + Debug + Send + Sync> Simplex<M> {
95
96 pub fn len(&self) -> usize {
97 self.tx.len()
98 }
99
100 pub fn len_non_zero(&self) -> bool {
101 self.len() > 0
102 }
103
104 /// Returns whether the value of the open flag.
105 pub fn is_open(&self) -> Outcome<bool> {
106 let open_read = lock_read!(self.open,
107 "While trying to read whether channel is open.",
108 );
109 Ok(*open_read)
110 }
111
112 /// Sets the open flag to closed, and returns the existing status of the flag. This simple
113 /// mechanism starts the channel closing process by telling others to stop sending messages.
114 pub fn close(&self) -> Outcome<bool> {
115 let mut open_write = lock_write!(self.open,
116 "While trying to close the channel.",
117 );
118 let is_open = *open_write;
119 *open_write = false;
120 Ok(is_open)
121 }
122
123 pub fn send(&self, msg: M) -> Outcome<()> {
124 res!(self.tx().send(msg));
125 Ok(())
126 }
127
128 /// Waits until a message is available.
129 pub fn recv(&self) -> Outcome<M> {
130 let msg = res!(self.rx().recv());
131 Ok(msg)
132 }
133
134 /// Captures a message but does not wait until one is present.
135 pub fn try_recv(&self) -> Recv<M> {
136 match self.rx().try_recv() {
137 Err(TryRecvError::Empty) => Recv::Empty,
138 Err(e) => Recv::Result(Err(err!(e,
139 "While trying to read channel without waiting.";
140 Channel, Read))),
141 Ok(msg) => Recv::Result(Ok(msg)),
142 }
143 }
144
145 pub fn recv_timeout(&self, sleep: Duration) -> Recv<M> {
146 match self.rx().recv_timeout(sleep) {
147 Err(RecvTimeoutError::Timeout) => Recv::Empty,
148 Err(e) => Recv::Result(Err(err!(e,
149 "While reading channel with a timeout of {:?}.", sleep;
150 Channel, Read))),
151 Ok(msg) => Recv::Result(Ok(msg)),
152 }
153 }
154
155 //pub fn send_ready(&self) -> Outcome<()> {
156 // self.send(M::ready())
157 //}
158 //pub fn send_finish(&self) -> Outcome<()> {
159 // self.send(M::finish())
160 //}
161
162 pub fn drain_messages(&self) -> Vec<String> {
163 let mut lines = Vec::new();
164 loop {
165 match self.try_recv() {
166 Recv::Empty => break,
167 Recv::Result(Err(e)) => lines.push(fmt!("<err>: {:?}", e)),
168 Recv::Result(Ok(m)) => lines.push(fmt!("{:?}", m)),
169 }
170 }
171 lines
172 }
173
174 /// Returns as soon as no more messages are detected in the channel. Returns an error if the
175 /// given `Duration`s are inconsistent.
176 pub fn wait_for_empty_channel(
177 &self,
178 check_interval: Duration,
179 max_wait: Duration,
180 )
181 -> Outcome<(Instant, bool)>
182 {
183 wait_for_true(
184 check_interval,
185 max_wait,
186 || { self.len() == 0 },
187 )
188 }
189
190}
191
192#[derive(Debug)]
193/// A channel for communicating in a two directions.
194///
195/// ```ignore
196///
197/// fwd tx1 ---->---- rx1 A full duplex contains two simplex
198/// channels. Each of these is accessed
199/// rev rx2 ----<---- tx2 via "fwd" and "rev".
200///
201/// ```
202pub struct FullDuplex<M> (
203 Simplex<M>,
204 Simplex<M>,
205);
206
207impl<M> Clone for FullDuplex<M> {
208 fn clone(&self) -> Self {
209 Self (
210 self.0.clone(),
211 self.1.clone(),
212 )
213 }
214}
215
216impl<M> FullDuplex<M> {
217 pub fn fwd(&self) -> &Simplex<M> { &self.0 }
218 pub fn rev(&self) -> &Simplex<M> { &self.1 }
219}
220
221impl<M: 'static + Debug + Send + Sync> FullDuplex<M> {
222 pub fn rx(&self) -> &Receiver<M> { &self.fwd().rx() }
223 pub fn tx(&self) -> &Sender<M> { &self.fwd().tx() }
224
225 pub fn send(&self, msg: M) -> Outcome<()> {
226 self.fwd().send(msg)
227 }
228
229 pub fn recv(&self) -> Outcome<M> {
230 self.fwd().recv()
231 }
232}
233
234impl<M> Default for FullDuplex<M> {
235 fn default() -> Self {
236 full_duplex::<M>()
237 }
238}