use super::ring::{blocking, SpscChannel};
use crate::channel::error::{Channel, Result};
use crate::channel::roles::{Consumer, Producer};
use std::cell::Cell;
use std::sync::atomic::Ordering;
use std::sync::Arc;
impl<T> SpscChannel<T> {
pub fn channel(capacity: usize) -> (SpscSender<T>, SpscReceiver<T>) {
let channel = Arc::new(Self::new(capacity));
(
SpscSender {
channel: channel.clone(),
cached_tail: Cell::new(0),
},
SpscReceiver {
channel,
cached_head: Cell::new(0),
},
)
}
}
pub struct SpscSender<T> {
pub(super) channel: Arc<SpscChannel<T>>,
cached_tail: Cell<usize>,
}
impl<T: Send> SpscSender<T> {
pub fn send(&self, value: T) -> Result<()> {
self.channel.send_cached(value, &self.cached_tail)
}
pub fn try_send(&self, value: T) -> Result<()> {
self.channel.try_send_cached(value, &self.cached_tail)
}
}
pub struct SpscReceiver<T> {
pub(super) channel: Arc<SpscChannel<T>>,
cached_head: Cell<usize>,
}
impl<T: Send> SpscReceiver<T> {
pub fn recv(&self) -> Result<T> {
blocking(|| self.channel.try_recv_cached(&self.cached_head))
}
pub fn try_recv(&self) -> Result<T> {
self.channel.try_recv_cached(&self.cached_head)
}
}
impl<T: Send> Producer<T> for SpscSender<T> {
#[inline]
fn send(&self, value: T) -> Result<()> {
SpscSender::send(self, value)
}
#[inline]
fn try_send(&self, value: T) -> Result<()> {
SpscSender::try_send(self, value)
}
#[inline]
fn is_full(&self) -> bool {
Channel::is_full(&*self.channel)
}
#[inline]
fn capacity(&self) -> Option<usize> {
Channel::capacity(&*self.channel)
}
}
impl<T: Send> Consumer<T> for SpscReceiver<T> {
#[inline]
fn recv(&self) -> Result<T> {
SpscReceiver::recv(self)
}
#[inline]
fn try_recv(&self) -> Result<T> {
SpscReceiver::try_recv(self)
}
#[inline]
fn is_empty(&self) -> bool {
Channel::is_empty(&*self.channel)
}
}
impl<T> Drop for SpscSender<T> {
fn drop(&mut self) {
self.channel.closed.store(true, Ordering::Release);
}
}
impl<T> Drop for SpscReceiver<T> {
fn drop(&mut self) {
self.channel.closed.store(true, Ordering::Release);
}
}
#[cfg(test)]
mod auto_traits {
use super::{SpscReceiver, SpscSender};
use static_assertions::{assert_impl_all, assert_not_impl_any};
assert_impl_all!(SpscSender<u64>: Send);
assert_impl_all!(SpscReceiver<u64>: Send);
#[allow(dead_code)]
fn halves_are_not_sync() {
assert_not_impl_any!(SpscSender<u64>: Sync);
assert_not_impl_any!(SpscReceiver<u64>: Sync);
}
}