#[macro_use]
extern crate crossbeam_channel as channel;
use channel::internal::select::{RecvArgument, SendArgument};
use channel::{Receiver, Sender};
use std::option;
pub fn bounded<T, U>(cap: usize) -> (BiChannel<T, U>, BiChannel<U, T>) {
let (tx1, rx1) = channel::bounded(cap);
let (tx2, rx2) = channel::bounded(cap);
(BiChannel::new(tx1, rx2), BiChannel::new(tx2, rx1))
}
pub fn unbounded<T, U>() -> (BiChannel<T, U>, BiChannel<U, T>) {
let (tx1, rx1) = channel::unbounded();
let (tx2, rx2) = channel::unbounded();
(BiChannel::new(tx1, rx2), BiChannel::new(tx2, rx1))
}
pub struct BiChannel<T, U> {
pub rx: Receiver<T>,
pub tx: Sender<U>,
}
impl<T, U> BiChannel<T, U> {
pub fn new(tx: Sender<U>, rx: Receiver<T>) -> Self {
BiChannel { rx, tx }
}
pub fn send(&self, msg: U) {
self.tx.send(msg)
}
pub fn recv(&self) -> Option<T> {
self.rx.recv()
}
}
impl<'a, T, U> RecvArgument<'a, T> for &'a BiChannel<T, U> {
type Iter = option::IntoIter<&'a Receiver<T>>;
fn _as_recv_argument(&'a self) -> Self::Iter {
Some(&self.rx).into_iter()
}
}
impl<'a, T, U> SendArgument<'a, T> for &'a BiChannel<U, T> {
type Iter = option::IntoIter<&'a Sender<T>>;
fn _as_send_argument(&'a self) -> Self::Iter {
Some(&self.tx).into_iter()
}
}
#[cfg(test)]
mod tests {
use channel;
use std::thread;
use std::time::Duration;
#[test]
fn simultaneous_handover() {
let (left, right) = super::bounded(1);
left.send(10);
right.send(20);
assert_eq!(left.recv(), Some(20));
assert_eq!(right.recv(), Some(10));
}
#[test]
fn rendezvous_recv_select() {
let (left, right) = super::bounded::<(), ()>(0);
let timeout = Duration::from_millis(10);
thread::spawn(move || {
left.send(());
});
select! {
recv(right, _msg) => {},
recv(channel::after(timeout)) => {
panic!("timeout waiting for rendezvous");
},
}
}
#[test]
fn rendezvous_send_select() {
let (left, right) = super::bounded::<(), ()>(0);
let timeout = Duration::from_millis(10);
thread::spawn(move || {
left.recv();
});
select! {
send(right, ()) => {},
recv(channel::after(timeout)) => {
panic!("timeout waiting for rendezvous");
},
}
}
#[test]
fn asymmetric_message_types() {
let (left, right) = super::unbounded::<u8, i16>();
left.send(0i16);
assert_eq!(right.recv().unwrap(), 0i16);
right.send(0u8);
assert_eq!(left.recv().unwrap(), 0u8);
}
}