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 | |
| 6 | use crate::{ |
| 7 | prelude::*, |
| 8 | //bot::CtrlMsg, |
| 9 | time::wait_for_true, |
| 10 | }; |
| 11 | |
| 12 | use std::{ |
| 13 | fmt::Debug, |
| 14 | sync::{ |
| 15 | Arc, |
| 16 | RwLock, |
| 17 | }, |
| 18 | time::{ |
| 19 | Duration, |
| 20 | Instant, |
| 21 | }, |
| 22 | }; |
| 23 | |
| 24 | pub use flume::{ |
| 25 | unbounded, |
| 26 | Sender, |
| 27 | Receiver, |
| 28 | TryRecvError, |
| 29 | RecvTimeoutError, |
| 30 | }; |
| 31 | |
| 32 | pub fn full_duplex<M>() -> FullDuplex<M> { |
| 33 | FullDuplex ( |
| 34 | simplex(), |
| 35 | simplex(), |
| 36 | ) |
| 37 | } |
| 38 | |
| 39 | pub 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 | /// ``` |
| 59 | pub 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 | |
| 65 | impl<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 | |
| 75 | impl<M> Default for Simplex<M> { |
| 76 | fn default() -> Self { |
| 77 | simplex::<M>() |
| 78 | } |
| 79 | } |
| 80 | |
| 81 | impl<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)] |
| 89 | pub enum Recv<M> { |
| 90 | Result(Outcome<M>), |
| 91 | Empty, |
| 92 | } |
| 93 | |
| 94 | impl<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 | /// ``` |
| 202 | pub struct FullDuplex<M> ( |
| 203 | Simplex<M>, |
| 204 | Simplex<M>, |
| 205 | ); |
| 206 | |
| 207 | impl<M> Clone for FullDuplex<M> { |
| 208 | fn clone(&self) -> Self { |
| 209 | Self ( |
| 210 | self.0.clone(), |
| 211 | self.1.clone(), |
| 212 | ) |
| 213 | } |
| 214 | } |
| 215 | |
| 216 | impl<M> FullDuplex<M> { |
| 217 | pub fn fwd(&self) -> &Simplex<M> { &self.0 } |
| 218 | pub fn rev(&self) -> &Simplex<M> { &self.1 } |
| 219 | } |
| 220 | |
| 221 | impl<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 | |
| 234 | impl<M> Default for FullDuplex<M> { |
| 235 | fn default() -> Self { |
| 236 | full_duplex::<M>() |
| 237 | } |
| 238 | } |