mod segment;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
#[cfg(feature = "futures")]
use futures::{Async, AsyncSink, StartSend, Poll, Stream, Sink};
#[cfg(feature = "futures")]
use futures::task::AtomicTask;
use self::segment::{Segment, SEGMENT_SIZE, Expanded};
pub use super::{SendError, ReceiveError};
pub fn channel<T: Send + Sync>() -> (Sender<T>, Receiver<T>) {
let segment = Arc::new(Segment::empty());
let shared = Arc::new(Shared{
receiver_alive: AtomicBool::new(true),
#[cfg(feature = "futures")]
task: AtomicTask::new(),
});
let sender = Sender {
current: Arc::clone(&segment),
shared: Arc::clone(&shared),
};
let receiver = Receiver {
current: segment,
read_index: 0,
shared: shared,
};
(sender, receiver)
}
#[derive(Debug)]
pub struct Sender<T> {
current: Arc<Segment<T>>,
shared: Arc<Shared>,
}
impl<T> Sender<T> {
pub fn send(&mut self, value: T) -> Result<(), SendError<T>> {
if !self.shared.receiver_alive.load(Ordering::Relaxed) {
return Err(SendError(value));
}
if let Expanded::Expanded(next_segment) = self.current.append(value) {
self.current = next_segment;
}
#[cfg(feature = "futures")]
self.shared.task.notify();
Ok(())
}
}
#[cfg(feature = "futures")]
impl<T> Sink for Sender<T> {
type SinkItem = T;
type SinkError = SendError<T>;
fn start_send(&mut self, item: T) -> StartSend<T, SendError<T>> {
match self.send(item) {
Ok(()) => Ok(AsyncSink::Ready),
Err(err) => Err(err),
}
}
fn poll_complete(&mut self) -> Poll<(), SendError<T>> {
Ok(Async::Ready(()))
}
}
impl<T> Clone for Sender<T> {
fn clone(&self) -> Sender<T> {
Sender {
current: Arc::clone(&self.current),
shared: Arc::clone(&self.shared),
}
}
}
#[cfg(feature = "futures")]
impl<T> Drop for Sender<T> {
fn drop(&mut self) {
self.shared.task.notify();
}
}
#[derive(Debug)]
pub struct Receiver<T> {
current: Arc<Segment<T>>,
read_index: usize,
shared: Arc<Shared>,
}
impl<T> Receiver<T> {
pub fn try_receive(&mut self) -> Result<T, ReceiveError> {
let index = match self.get_read_index() {
Some(index) => index,
None => {
if self.update_current_segment() {
0
} else if self.is_disconnected() {
return Err(ReceiveError::Disconnected);
} else {
return Err(ReceiveError::Empty);
}
},
};
match self.current.try_pop(index) {
Some(value) => Ok(value),
None => {
self.read_index -= 1;
if self.is_disconnected() {
Err(ReceiveError::Disconnected)
} else {
Err(ReceiveError::Empty)
}
},
}
}
fn get_read_index(&mut self) -> Option<usize> {
match self.read_index {
current if current >= SEGMENT_SIZE => None,
current => {
self.read_index += 1;
Some(current)
},
}
}
fn update_current_segment(&mut self) -> bool {
match self.current.next_segment() {
Some(next_segment) => {
self.current = next_segment;
self.read_index = 1;
true
},
None => false
}
}
fn is_disconnected(&self) -> bool {
Arc::strong_count(&self.shared) == 1
}
}
#[cfg(feature = "futures")]
impl<T> Stream for Receiver<T> {
type Item = T;
type Error = ();
fn poll(&mut self) -> Poll<Option<T>, ()> {
self.shared.task.register();
match self.try_receive() {
Ok(value) => Ok(Async::Ready(Some(value))),
Err(err) => match err {
ReceiveError::Empty => Ok(Async::NotReady),
ReceiveError::Disconnected => Ok(Async::Ready(None)),
},
}
}
}
impl<T> Drop for Receiver<T> {
fn drop(&mut self) {
self.shared.receiver_alive.store(false, Ordering::Relaxed);
}
}
#[derive(Debug)]
struct Shared {
receiver_alive: AtomicBool,
#[cfg(feature = "futures")]
task: AtomicTask,
}