use crate::storage::Storage;
use crate::{Region, RegionMut, Ring};
use alloc::sync::Arc;
use core::sync::atomic::{AtomicBool, Ordering};
pub fn queue_from<T, S>(storage: S) -> (Reader<S>, Writer<S>)
where
S: Storage<T>,
{
let ring = Ring::new(storage.capacity());
queue_from_parts(ring, storage)
}
pub fn queue_from_parts<S>(ring: Ring, storage: S) -> (Reader<S>, Writer<S>) {
let state = Arc::new(State {
ring,
storage,
is_reader_open: AtomicBool::new(true),
is_writer_open: AtomicBool::new(true),
});
let reader = Reader {
state: state.clone(),
};
let writer = Writer { state };
(reader, writer)
}
#[cfg(feature = "heap-buffer")]
mod heap_constructors {
use crate::nonblocking::{queue_from_parts, Reader, Writer};
use crate::storage::HeapBuffer;
use crate::Ring;
#[cfg_attr(docsrs, doc(cfg(feature = "heap-buffer")))]
pub fn queue<T>(capacity: usize) -> (Reader<HeapBuffer<T>>, Writer<HeapBuffer<T>>)
where
T: Default,
{
let ring = Ring::new(capacity);
let buffer = HeapBuffer::new(capacity);
queue_from_parts(ring, buffer)
}
}
#[cfg(feature = "heap-buffer")]
pub use self::heap_constructors::*;
#[derive(Debug)]
struct State<S> {
ring: Ring,
storage: S,
is_reader_open: AtomicBool,
is_writer_open: AtomicBool,
}
#[derive(Debug)]
pub struct Reader<S> {
state: Arc<State<S>>,
}
#[derive(Debug)]
pub struct Writer<S> {
state: Arc<State<S>>,
}
impl<S> State<S> {
fn close_reader(&self) {
self.is_reader_open.store(false, Ordering::Release);
}
fn close_writer(&self) {
self.is_writer_open.store(false, Ordering::Release);
}
}
impl<S> Reader<S> {
#[inline]
pub fn is_writer_open(&self) -> bool {
self.state.is_writer_open.load(Ordering::Acquire)
}
#[inline]
pub fn has_data(&self) -> bool {
let (r1, r2) = self.state.ring.left_ranges();
!r1.is_empty() || !r2.is_empty()
}
#[inline]
pub fn is_full(&self) -> bool {
let (r1, r2) = self.state.ring.right_ranges();
r1.is_empty() && r2.is_empty()
}
#[inline]
pub fn buf<T>(&self) -> Region<T>
where
S: Storage<T>,
{
let (range_0, range_1) = self.state.ring.left_ranges();
Region::new(
self.state.storage.slice(range_0),
self.state.storage.slice(range_1),
)
}
#[inline]
pub fn consume(&mut self, amt: usize) {
self.state.ring.advance_left(amt);
}
pub fn read<T>(&mut self, buf: &mut [T]) -> usize
where
S: Storage<T>,
T: Clone,
{
let src_buf = self.buf();
let len = src_buf.len().min(buf.len());
src_buf.slice(..len).clone_to_slice(&mut buf[..len]);
self.consume(len);
len
}
#[inline]
pub fn close(&mut self) {
self.state.close_reader();
}
}
impl<S> Drop for Reader<S> {
#[inline]
fn drop(&mut self) {
self.state.close_reader();
}
}
impl<S> Writer<S> {
#[inline]
pub fn is_reader_open(&self) -> bool {
self.state.is_reader_open.load(Ordering::Acquire)
}
#[inline]
pub fn has_space(&self) -> bool {
let (r1, r2) = self.state.ring.right_ranges();
!r1.is_empty() || !r2.is_empty()
}
#[inline]
pub fn is_flushed(&self) -> bool {
let (r1, r2) = self.state.ring.left_ranges();
r1.is_empty() && r2.is_empty()
}
pub fn buf<T>(&mut self) -> RegionMut<T>
where
S: Storage<T>,
{
if !self.is_reader_open() {
return Default::default();
}
let (range_0, range_1) = self.state.ring.right_ranges();
RegionMut::new(
unsafe { self.state.storage.slice_mut_unchecked(range_0) },
unsafe { self.state.storage.slice_mut_unchecked(range_1) },
)
}
#[inline]
pub fn feed(&mut self, amt: usize) {
self.state.ring.advance_right(amt);
}
pub fn write<T>(&mut self, buf: &[T]) -> usize
where
S: Storage<T>,
T: Clone,
{
let mut dest_buf = self.buf();
let len = dest_buf.len().min(buf.len());
dest_buf.slice_mut(..len).clone_from_slice(&buf[..len]);
self.feed(len);
len
}
#[inline]
pub fn close(&mut self) {
self.state.close_writer();
}
}
impl<S> Drop for Writer<S> {
#[inline]
fn drop(&mut self) {
self.state.close_writer();
}
}
#[cfg(feature = "std-io")]
mod io_impls {
use crate::nonblocking::{Reader, Writer};
use crate::storage::Storage;
use std::io::{BufRead, ErrorKind, Read, Result, Write};
#[cfg_attr(docsrs, doc(cfg(feature = "std-io")))]
impl<S> Read for Reader<S>
where
S: Storage<u8>,
{
fn read(&mut self, buf: &mut [u8]) -> Result<usize> {
let src_buf = self.buf();
if src_buf.is_empty() {
return if self.is_writer_open() {
Err(ErrorKind::WouldBlock.into())
} else {
Ok(0)
};
}
let len = src_buf.len().min(buf.len());
src_buf.slice(..len).copy_to_slice(&mut buf[..len]);
self.consume(len);
Ok(len)
}
}
#[cfg_attr(docsrs, doc(cfg(feature = "std-io")))]
impl<S> BufRead for Reader<S>
where
S: Storage<u8>,
{
fn fill_buf(&mut self) -> Result<&[u8]> {
let buf = self.buf().contiguous();
if !buf.is_empty() {
return Ok(buf);
}
if self.is_writer_open() {
return Err(ErrorKind::WouldBlock.into());
}
Ok(Default::default())
}
fn consume(&mut self, amt: usize) {
self.consume(amt);
}
}
#[cfg_attr(docsrs, doc(cfg(feature = "std-io")))]
impl<S> Write for Writer<S>
where
S: Storage<u8>,
{
fn write(&mut self, buf: &[u8]) -> Result<usize> {
let mut dest_buf = self.buf();
if !dest_buf.is_empty() {
let len = dest_buf.len().min(buf.len());
dest_buf.slice_mut(..len).copy_from_slice(&buf[..len]);
self.feed(len);
return Ok(len);
}
if !self.is_reader_open() {
return Ok(Default::default());
}
Err(ErrorKind::WouldBlock.into())
}
fn flush(&mut self) -> Result<()> {
if self.is_flushed() {
return Ok(());
}
if self.is_reader_open() {
return Err(ErrorKind::WouldBlock.into());
}
Err(ErrorKind::UnexpectedEof.into())
}
}
}
#[cfg(feature = "std-io")]
pub use self::io_impls::*;