use crate::storage::Storage;
use crate::{Region, RegionMut, Ring};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::{error, fmt, io};
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),
data_available_cond: Condvar::new(),
space_available_cond: Condvar::new(),
});
let reader = Reader {
state: state.clone(),
data_available_mutex: Mutex::new(()),
};
let writer = Writer {
state,
space_available_mutex: Mutex::new(()),
};
(reader, writer)
}
#[cfg(feature = "heap-buffer")]
mod heap_constructors {
use crate::blocking::{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,
data_available_cond: Condvar,
space_available_cond: Condvar,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WriteError {
ReaderClosed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReadExactError {
WriterClosed,
}
#[derive(Debug)]
pub struct Reader<S> {
state: Arc<State<S>>,
data_available_mutex: Mutex<()>,
}
#[derive(Debug)]
pub struct Writer<S> {
state: Arc<State<S>>,
space_available_mutex: Mutex<()>,
}
impl<S> State<S> {
fn close_reader(&self) {
let was_open = self.is_reader_open.swap(false, Ordering::AcqRel);
if was_open {
self.space_available_cond.notify_all();
}
}
fn close_writer(&self) {
let was_open = self.is_writer_open.swap(false, Ordering::AcqRel);
if was_open {
self.data_available_cond.notify_all();
}
}
}
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()
}
pub fn fill_buf<T>(&mut self) -> Region<T>
where
S: Storage<T>,
{
if self.has_data() {
return self.buf();
}
if !self.is_writer_open() {
return Default::default();
}
let mut lock = self.data_available_mutex.lock().unwrap();
loop {
lock = self.state.data_available_cond.wait(lock).unwrap();
if self.has_data() {
return self.buf();
}
if !self.is_writer_open() {
return Default::default();
}
}
}
pub fn consume(&mut self, amt: usize) {
self.state.ring.advance_left(amt);
self.state.space_available_cond.notify_all();
}
pub fn read<T>(&mut self, buf: &mut [T]) -> usize
where
S: Storage<T>,
T: Clone,
{
let src_buf = self.fill_buf();
if src_buf.is_empty() {
return 0;
}
let len = src_buf.len().min(buf.len());
src_buf.slice(..len).clone_to_slice(&mut buf[..len]);
self.consume(len);
len
}
pub fn read_exact<T>(&mut self, buf: &mut [T]) -> Result<usize, ReadExactError>
where
S: Storage<T>,
T: Clone,
{
let len = buf.len();
let src_buf = loop {
let src_buf = self.fill_buf();
if src_buf.len() >= len {
break src_buf;
}
if !self.is_writer_open() {
return Err(ReadExactError::WriterClosed);
}
};
src_buf.slice(..len).clone_to_slice(buf);
self.consume(len);
Ok(len)
}
#[inline]
pub fn close(&mut self) {
self.state.close_reader();
}
#[inline]
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),
)
}
}
impl<S> io::Read for Reader<S>
where
S: Storage<u8>,
{
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
let src_buf = self.fill_buf();
let len = src_buf.len().min(buf.len());
src_buf.slice(..len).copy_to_slice(&mut buf[..len]);
self.consume(len);
Ok(len)
}
}
impl<S> io::BufRead for Reader<S>
where
S: Storage<u8>,
{
fn fill_buf(&mut self) -> io::Result<&[u8]> {
Ok(self.fill_buf().contiguous())
}
fn consume(&mut self, amt: usize) {
self.consume(amt);
}
}
impl<S> Drop for Reader<S> {
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 (r0, r1) = self.state.ring.right_ranges();
!r0.is_empty() || !r1.is_empty()
}
#[inline]
pub fn is_flushed(&self) -> bool {
let (r0, r1) = self.state.ring.left_ranges();
r0.is_empty() && r1.is_empty()
}
fn get_flush_state(&self) -> Option<Result<(), WriteError>> {
if self.is_flushed() {
return Some(Ok(()));
}
if !self.is_reader_open() {
return Some(Err(WriteError::ReaderClosed));
}
None
}
#[inline]
fn buf<T>(&mut self) -> RegionMut<T>
where
S: Storage<T>,
{
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) },
)
}
pub fn empty_buf<T>(&mut self) -> RegionMut<T>
where
S: Storage<T>,
{
if !self.is_reader_open() {
return Default::default();
}
if self.has_space() {
return self.buf();
}
{
let mut lock = self.space_available_mutex.lock().unwrap();
loop {
lock = self.state.space_available_cond.wait(lock).unwrap();
if !self.is_reader_open() {
return Default::default();
}
if self.has_space() {
break;
}
}
}
self.buf()
}
pub fn feed(&mut self, len: usize) {
self.state.ring.advance_right(len);
self.state.data_available_cond.notify_all();
}
pub fn write<T>(&mut self, buf: &[T]) -> usize
where
S: Storage<T>,
T: Clone,
{
let mut dest_buf = self.empty_buf();
if dest_buf.is_empty() {
return 0;
}
let len = dest_buf.len().min(buf.len());
dest_buf.slice_mut(..len).clone_from_slice(&buf[..len]);
self.feed(len);
len
}
pub fn write_all<T>(&mut self, buf: &[T]) -> Result<usize, WriteError>
where
S: Storage<T>,
T: Clone,
{
let len = buf.len();
let mut dest_buf = loop {
let dest_buf = self.empty_buf();
if dest_buf.is_empty() {
return Err(WriteError::ReaderClosed);
}
if dest_buf.len() >= len {
break dest_buf;
}
};
dest_buf.slice_mut(..len).clone_from_slice(buf);
self.feed(len);
Ok(len)
}
pub fn flush(&mut self) -> Result<(), WriteError> {
if let Some(flush_state) = self.get_flush_state() {
return flush_state;
}
let mut lock = self.space_available_mutex.lock().unwrap();
loop {
lock = self.state.space_available_cond.wait(lock).unwrap();
if let Some(flush_state) = self.get_flush_state() {
return flush_state;
}
}
}
pub fn close(&mut self) -> Result<(), WriteError> {
self.flush()?;
self.state.close_writer();
Ok(())
}
}
impl<S> io::Write for Writer<S>
where
S: Storage<u8>,
{
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
let mut dest_buf = self.empty_buf();
let len = dest_buf.len().min(buf.len());
dest_buf.slice_mut(..len).copy_from_slice(&buf[..len]);
self.feed(len);
Ok(len)
}
fn flush(&mut self) -> io::Result<()> {
self.flush().map_err(Into::into)
}
}
impl<S> Drop for Writer<S> {
#[inline]
fn drop(&mut self) {
self.state.close_writer();
}
}
impl error::Error for WriteError {}
impl fmt::Display for WriteError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
WriteError::ReaderClosed => write!(f, "reader closed"),
}
}
}
impl From<WriteError> for io::Error {
fn from(err: WriteError) -> Self {
match err {
WriteError::ReaderClosed => io::ErrorKind::UnexpectedEof.into(),
}
}
}
impl error::Error for ReadExactError {}
impl fmt::Display for ReadExactError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
ReadExactError::WriterClosed => write!(f, "writer closed"),
}
}
}
impl From<ReadExactError> for io::Error {
fn from(err: ReadExactError) -> Self {
match err {
ReadExactError::WriterClosed => io::ErrorKind::UnexpectedEof.into(),
}
}
}