#![allow(dead_code)]
use std::{
collections::VecDeque,
pin::Pin,
sync::{
atomic::{AtomicBool, Ordering},
Arc,
},
task::{Context, Poll, Waker},
};
use futures::stream::Stream;
use parking_lot::Mutex;
#[derive(Debug)]
struct ReceiverNotifier {
handle: Waker,
awake: Arc<AtomicBool>,
}
struct RawDeque<T> {
front_values: VecDeque<T>,
back_values: VecDeque<T>,
rx_notifiers: VecDeque<ReceiverNotifier>,
}
impl<T> RawDeque<T> {
const fn new() -> Self {
Self {
front_values: VecDeque::new(),
back_values: VecDeque::new(),
rx_notifiers: VecDeque::new(),
}
}
}
impl<T> RawDeque<T> {
fn notify_rx(&mut self) {
if let Some(n) = self.rx_notifiers.pop_front() {
n.handle.wake();
n.awake.store(true, Ordering::Relaxed);
}
}
}
pub struct StreamableDeque<T> {
inner: Mutex<RawDeque<T>>,
}
impl<T> std::fmt::Debug for StreamableDeque<T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("StreamableDeque { ... }").finish()
}
}
impl<T> Default for StreamableDeque<T> {
fn default() -> Self {
Self {
inner: Mutex::new(RawDeque::new()),
}
}
}
impl<T> StreamableDeque<T> {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn push_front(&self, item: T) {
let mut inner = self.inner.lock();
inner.front_values.push_back(item);
inner.notify_rx();
}
pub fn push_back(&self, item: T) {
let mut inner = self.inner.lock();
inner.back_values.push_back(item);
inner.notify_rx();
}
pub const fn stream(&self) -> StreamReceiver<T> {
StreamReceiver {
queue: self,
awake: None,
}
}
pub fn pop_front(&self) -> Option<T> {
let mut inner = self.inner.lock();
inner
.front_values
.pop_front()
.or_else(|| inner.back_values.pop_front())
}
#[cfg(test)]
pub(crate) fn pop_back(&self) -> Option<T> {
let mut inner = self.inner.lock();
inner
.back_values
.pop_back()
.or_else(|| inner.front_values.pop_back())
}
}
pub struct StreamReceiver<'a, T> {
queue: &'a StreamableDeque<T>,
awake: Option<Arc<AtomicBool>>,
}
impl<T> Stream for StreamReceiver<'_, T> {
type Item = T;
fn poll_next(mut self: Pin<&mut Self>, ctx: &mut Context) -> Poll<Option<Self::Item>> {
let mut inner = self.queue.inner.lock();
let value = inner
.front_values
.pop_front()
.or_else(|| inner.back_values.pop_front());
if let Some(v) = value {
self.awake = None;
Poll::Ready(Some(v))
} else {
let awake = Arc::new(AtomicBool::new(false));
inner.rx_notifiers.push_back(ReceiverNotifier {
handle: ctx.waker().clone(),
awake: awake.clone(),
});
self.awake = Some(awake);
drop(inner);
Poll::Pending
}
}
}
impl<T> Drop for StreamReceiver<'_, T> {
fn drop(&mut self) {
let awake = self.awake.take().map(|w| w.load(Ordering::Relaxed));
if awake == Some(true) {
let mut queue_wakers = self.queue.inner.lock();
if let Some(n) = queue_wakers.rx_notifiers.pop_front() {
n.awake.store(true, Ordering::Relaxed);
n.handle.wake();
}
}
}
}
#[cfg(test)]
mod tests {
use futures::stream::StreamExt;
use super::*;
#[tokio::test]
async fn streamable_deque() {
let queue = Arc::new(StreamableDeque::<i32>::new());
let pos_queue = queue.clone();
tokio::spawn(async move {
for i in 0..=10 {
pos_queue.push_back(i);
}
});
let neg_queue = queue.clone();
tokio::spawn(async move {
for i in -10..=-1 {
neg_queue.push_front(i);
}
});
let mut rx_vec = vec![];
let mut stream = queue.stream().enumerate();
while let Some((i, v)) = stream.next().await {
rx_vec.push(v);
if i >= 20 {
break;
}
}
let expected_vec: Vec<i32> = (-10..=10).collect();
assert_eq!(expected_vec, rx_vec);
}
}