mod error;
use std::fmt;
use std::time::Duration;
pub use error::{
RecvError, RecvTimeoutError, SendError, SendTimeoutError, TryRecvError, TrySendError,
};
use crate::LockFreeQueue;
use crate::error::{PopTimeoutError, PushTimeoutError, TryPopError, TryPushError};
use crate::sync::{Arc, AtomicUsize, Ordering};
#[must_use]
pub fn bounded<T>(capacity: usize) -> (Sender<T>, Receiver<T>) {
let shared = Arc::new(Shared {
queue: LockFreeQueue::new(capacity),
senders: AtomicUsize::new(1),
receivers: AtomicUsize::new(1),
});
(
Sender {
shared: Arc::clone(&shared),
},
Receiver { shared },
)
}
struct Shared<T> {
queue: LockFreeQueue<T>,
senders: AtomicUsize,
receivers: AtomicUsize,
}
impl<T> Shared<T> {
fn release(&self, count: &AtomicUsize) {
let previous = count.fetch_sub(1, Ordering::AcqRel);
#[cfg(not(parkring_mutant = "channel_no_disconnect"))]
if previous == 1 {
self.queue.close();
}
#[cfg(parkring_mutant = "channel_no_disconnect")]
let _ = previous;
}
}
pub struct Sender<T> {
shared: Arc<Shared<T>>,
}
pub struct Receiver<T> {
shared: Arc<Shared<T>>,
}
impl<T> Sender<T> {
pub fn send(&self, msg: T) -> Result<(), SendError<T>> {
self.shared
.queue
.push(msg)
.map_err(|e| SendError(e.into_inner()))
}
pub fn try_send(&self, msg: T) -> Result<(), TrySendError<T>> {
self.shared.queue.try_push(msg).map_err(|e| match e {
TryPushError::Full(m) => TrySendError::Full(m),
TryPushError::Closed(m) => TrySendError::Disconnected(m),
})
}
pub fn send_timeout(&self, msg: T, timeout: Duration) -> Result<(), SendTimeoutError<T>> {
self.shared
.queue
.push_timeout(msg, timeout)
.map_err(|e| match e {
PushTimeoutError::Timeout(m) => SendTimeoutError::Timeout(m),
PushTimeoutError::Closed(m) => SendTimeoutError::Disconnected(m),
})
}
#[must_use]
pub fn is_disconnected(&self) -> bool {
self.shared.queue.is_closed()
}
#[must_use]
pub fn capacity(&self) -> usize {
self.shared.queue.capacity()
}
#[must_use]
pub fn len(&self) -> usize {
self.shared.queue.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.shared.queue.is_empty()
}
#[must_use]
pub fn is_full(&self) -> bool {
self.shared.queue.is_full()
}
}
impl<T> Receiver<T> {
pub fn recv(&self) -> Result<T, RecvError> {
self.shared.queue.pop().map_err(|_| RecvError)
}
pub fn try_recv(&self) -> Result<T, TryRecvError> {
self.shared.queue.try_pop().map_err(|e| match e {
TryPopError::Empty => TryRecvError::Empty,
TryPopError::Closed => TryRecvError::Disconnected,
})
}
pub fn recv_timeout(&self, timeout: Duration) -> Result<T, RecvTimeoutError> {
self.shared.queue.pop_timeout(timeout).map_err(|e| match e {
PopTimeoutError::Timeout => RecvTimeoutError::Timeout,
PopTimeoutError::Closed => RecvTimeoutError::Disconnected,
})
}
#[must_use]
pub fn iter(&self) -> Iter<'_, T> {
Iter { rx: self }
}
#[must_use]
pub fn try_iter(&self) -> TryIter<'_, T> {
TryIter { rx: self }
}
#[must_use]
pub fn is_disconnected(&self) -> bool {
self.shared.queue.is_closed()
}
#[must_use]
pub fn capacity(&self) -> usize {
self.shared.queue.capacity()
}
#[must_use]
pub fn len(&self) -> usize {
self.shared.queue.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.shared.queue.is_empty()
}
#[must_use]
pub fn is_full(&self) -> bool {
self.shared.queue.is_full()
}
}
impl<T> Clone for Sender<T> {
fn clone(&self) -> Self {
self.shared.senders.fetch_add(1, Ordering::Relaxed);
Self {
shared: Arc::clone(&self.shared),
}
}
}
impl<T> Clone for Receiver<T> {
fn clone(&self) -> Self {
self.shared.receivers.fetch_add(1, Ordering::Relaxed);
Self {
shared: Arc::clone(&self.shared),
}
}
}
impl<T> Drop for Sender<T> {
fn drop(&mut self) {
self.shared.release(&self.shared.senders);
}
}
impl<T> Drop for Receiver<T> {
fn drop(&mut self) {
self.shared.release(&self.shared.receivers);
}
}
impl<T> fmt::Debug for Sender<T> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Sender")
.field("len", &self.len())
.field("capacity", &self.capacity())
.field("disconnected", &self.is_disconnected())
.finish()
}
}
impl<T> fmt::Debug for Receiver<T> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Receiver")
.field("len", &self.len())
.field("capacity", &self.capacity())
.field("disconnected", &self.is_disconnected())
.finish()
}
}
#[derive(Debug)]
pub struct Iter<'a, T> {
rx: &'a Receiver<T>,
}
impl<T> Iterator for Iter<'_, T> {
type Item = T;
fn next(&mut self) -> Option<T> {
self.rx.recv().ok()
}
}
#[derive(Debug)]
pub struct TryIter<'a, T> {
rx: &'a Receiver<T>,
}
impl<T> Iterator for TryIter<'_, T> {
type Item = T;
fn next(&mut self) -> Option<T> {
self.rx.try_recv().ok()
}
}
#[derive(Debug)]
pub struct IntoIter<T> {
rx: Receiver<T>,
}
impl<T> Iterator for IntoIter<T> {
type Item = T;
fn next(&mut self) -> Option<T> {
self.rx.recv().ok()
}
}
impl<T> IntoIterator for Receiver<T> {
type Item = T;
type IntoIter = IntoIter<T>;
fn into_iter(self) -> IntoIter<T> {
IntoIter { rx: self }
}
}
impl<'a, T> IntoIterator for &'a Receiver<T> {
type Item = T;
type IntoIter = Iter<'a, T>;
fn into_iter(self) -> Iter<'a, T> {
self.iter()
}
}