use crate::{
RuntimeError,
futures::channel::{
channel_task::{ChannelRecvTask, ChannelSendTask},
core::{Core, Refused},
},
};
use std::{fmt, sync::Arc};
pub struct Sender<T> {
core: Arc<Core<T>>,
}
impl<T> Sender<T>
where
T: Send + 'static,
{
pub(crate) fn new(core: Arc<Core<T>>) -> Self {
Self { core }
}
pub fn send(&self, value: T) -> Result<(), RuntimeError> {
match self.core.push(value) {
Ok(()) => Ok(()),
Err(_) => Err(RuntimeError::Closed),
}
}
pub fn len(&self) -> usize {
self.core.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
}
pub struct BoundedSender<T> {
core: Arc<Core<T>>,
}
impl<T> BoundedSender<T>
where
T: Send + 'static,
{
pub(crate) fn new(core: Arc<Core<T>>) -> Self {
Self { core }
}
pub fn send(&self, value: T) -> ChannelSendTask<T> {
ChannelSendTask::new(self.clone(), value)
}
pub fn try_send(&self, value: T) -> Result<Result<(), T>, RuntimeError> {
match self.core.push(value) {
Ok(()) => Ok(Ok(())),
Err(Refused::Full(value)) => Ok(Err(value)),
Err(Refused::Closed) => Err(RuntimeError::Closed),
}
}
pub fn len(&self) -> usize {
self.core.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub(crate) fn core(&self) -> &Core<T> {
&self.core
}
}
pub struct Receiver<T> {
core: Arc<Core<T>>,
}
impl<T> Receiver<T>
where
T: Send + 'static,
{
pub(crate) fn new(core: Arc<Core<T>>) -> Self {
Self { core }
}
pub fn recv(&self) -> ChannelRecvTask<T> {
ChannelRecvTask::new(self.clone())
}
pub fn try_recv(&self) -> Result<T, RuntimeError> {
self.core.pop()
}
pub fn len(&self) -> usize {
self.core.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub(crate) fn core(&self) -> &Core<T> {
&self.core
}
}
impl<T> Clone for Sender<T> {
fn clone(&self) -> Self {
self.core.add_sender();
Self {
core: Arc::clone(&self.core),
}
}
}
impl<T> Clone for BoundedSender<T> {
fn clone(&self) -> Self {
self.core.add_sender();
Self {
core: Arc::clone(&self.core),
}
}
}
impl<T> Clone for Receiver<T> {
fn clone(&self) -> Self {
self.core.add_receiver();
Self {
core: Arc::clone(&self.core),
}
}
}
impl<T> Drop for Sender<T> {
fn drop(&mut self) {
self.core.drop_sender();
}
}
impl<T> Drop for BoundedSender<T> {
fn drop(&mut self) {
self.core.drop_sender();
}
}
impl<T> Drop for Receiver<T> {
fn drop(&mut self) {
self.core.drop_receiver();
}
}
impl<T> fmt::Debug for Sender<T> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.debug_struct("Sender").finish_non_exhaustive()
}
}
impl<T> fmt::Debug for BoundedSender<T> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("BoundedSender")
.finish_non_exhaustive()
}
}
impl<T> fmt::Debug for Receiver<T> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.debug_struct("Receiver").finish_non_exhaustive()
}
}