use core::{convert::Infallible, fmt};
use std::rc::Rc;
use either::Either::{self};
use crate::prelude::*;
use crate::queues::Queue;
use super::advanced::sssr::{Receiver as AdvancedReceiver, Sender as AdvancedSender, State};
pub fn new_sssr<Q: Queue, F>(queue: Q) -> (Sender<Q, F>, Receiver<Q, F>) where {
let state = Rc::new(State::new(queue));
let (s, r) = super::advanced::new_sssr(state);
(Sender(s), Receiver(r))
}
pub struct Sender<Q, F>(AdvancedSender<Rc<State<Q, F>>, Q, F>);
impl<Q, F> fmt::Debug for Sender<Q, F>
where
Q: fmt::Debug,
F: fmt::Debug,
{
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.0.fmt(f)
}
}
impl<Q, F> Sender<Q, F>
where
Q: Queue,
{
pub fn len(&self) -> usize {
self.0.len()
}
pub fn is_empty(&self) -> bool {
self.0.is_empty()
}
pub fn is_receiver_dropped(&self) -> bool {
self.0.is_receiver_dropped()
}
}
impl<Q: Queue, F> Consumer for Sender<Q, F> {
type Item = Q::Item;
type Final = F;
type Error = Infallible;
async fn consume(&mut self, val: Either<Self::Item, Self::Final>) -> Result<(), Self::Error> {
self.0.consume(val).await
}
async fn flush(&mut self) -> Result<(), Self::Error> {
self.0.flush().await
}
}
impl<Q: Queue, F> BulkConsumer for Sender<Q, F> {
async fn expose_slots_gracefully<Fun, Ret>(&mut self, f: Fun) -> Result<Ret, (Fun, Self::Error)>
where
Fun: AsyncFnOnce(&mut [Self::Item]) -> (usize, Ret),
{
self.0.expose_slots_gracefully(f).await
}
}
pub struct Receiver<Q, F>(AdvancedReceiver<Rc<State<Q, F>>, Q, F>);
impl<Q, F> fmt::Debug for Receiver<Q, F>
where
Q: fmt::Debug,
F: fmt::Debug,
{
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.0.fmt(f)
}
}
impl<Q: Queue, F> Receiver<Q, F> {
pub fn len(&self) -> usize {
self.0.len()
}
pub fn is_empty(&self) -> bool {
self.0.is_empty()
}
pub fn is_sender_dropped(&self) -> bool {
self.0.is_sender_dropped()
}
}
impl<Q: Queue, F> Producer for Receiver<Q, F> {
type Item = Q::Item;
type Final = F;
type Error = Infallible;
async fn produce(&mut self) -> Result<Either<Self::Item, Self::Final>, Self::Error> {
self.0.produce().await
}
async fn slurp(&mut self) -> Result<(), Self::Error> {
self.0.slurp().await
}
}
impl<Q: Queue, F> BulkProducer for Receiver<Q, F> {
async fn expose_items_gracefully<Fun, Ret>(
&mut self,
f: Fun,
) -> Result<Either<Ret, (Fun, Self::Final)>, (Fun, Self::Error)>
where
Fun: AsyncFnOnce(&[Self::Item]) -> (usize, Ret),
{
self.0.expose_items_gracefully(f).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use futures::join;
use crate::queues::new_static;
#[test]
fn test_spsc_sufficient_capacity() {
let (mut sender, mut receiver) = new_sssr(new_static::<i16, 99>());
pollster::block_on(async {
assert!(sender.consume_item(300).await.is_ok());
assert!(sender.consume_item(400).await.is_ok());
assert!(sender.consume_item(500).await.is_ok());
assert!(sender.consume_final(-17).await.is_ok());
assert_eq!(300, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(400, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(500, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(-17, receiver.produce().await.unwrap().unwrap_right());
});
}
#[test]
fn test_spsc_low_capacity() {
pollster::block_on(async {
let (mut sender, mut receiver) = new_sssr(new_static::<i16, 2>());
let send_things = async {
assert!(sender.consume_item(300).await.is_ok());
assert!(sender.consume_item(400).await.is_ok());
assert!(sender.consume_item(500).await.is_ok());
assert!(sender.consume_final(-17).await.is_ok());
};
let receive_things = async {
assert_eq!(300, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(400, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(500, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(-17, receiver.produce().await.unwrap().unwrap_right());
};
join!(send_things, receive_things);
});
}
#[test]
fn test_spsc_immediate_final() {
pollster::block_on(async {
let (mut sender, mut receiver) = new_sssr(new_static::<i16, 3>());
let send_things = async {
assert!(sender.consume_final(-17).await.is_ok());
};
let receive_things = async {
assert_eq!(-17, receiver.produce().await.unwrap().unwrap_right());
};
join!(send_things, receive_things);
});
}
#[test]
fn test_spsc_receive_then_send_concurrently() {
pollster::block_on(async {
let (mut sender, mut receiver) = new_sssr(new_static::<i16, 2>());
let send_things = async {
assert!(sender.consume_item(300).await.is_ok());
assert!(sender.consume_item(400).await.is_ok());
assert!(sender.consume_item(500).await.is_ok());
assert!(sender.consume_final(-17).await.is_ok());
};
let receive_things = async {
assert_eq!(300, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(400, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(500, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(-17, receiver.produce().await.unwrap().unwrap_right());
};
join!(receive_things, send_things);
});
}
#[test]
fn test_spsc_capacity_1() {
pollster::block_on(async {
let (mut sender, mut receiver) = new_sssr(new_static::<i16, 1>());
let send_things = async {
assert!(sender.consume_item(300).await.is_ok());
assert!(sender.consume_item(400).await.is_ok());
assert!(sender.consume_item(500).await.is_ok());
assert!(sender.consume_final(-17).await.is_ok());
};
let receive_things = async {
assert_eq!(300, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(400, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(500, receiver.produce().await.unwrap().unwrap_left());
assert_eq!(-17, receiver.produce().await.unwrap().unwrap_right());
};
join!(receive_things, send_things);
});
}
}