use std::sync::{Arc, Mutex};
use crate::cell::{Computed, Source};
use crate::queue::{QueuePopError, QueuePushError, QueueStorage, VecDequeStorage};
use crate::thread_safe::ThreadSafeContext;
struct Inner<T, S> {
storage: Arc<Mutex<S>>,
capacity: Option<usize>,
head: Computed<Option<T>>,
len: Computed<usize>,
is_empty: Computed<bool>,
is_full: Computed<bool>,
closed: Source<bool>,
}
pub struct ThreadSafeQueueCell<T, S = VecDequeStorage<T>> {
inner: Arc<Inner<T, S>>,
}
impl<T, S> Clone for ThreadSafeQueueCell<T, S> {
fn clone(&self) -> Self {
Self {
inner: Arc::clone(&self.inner),
}
}
}
impl<T> ThreadSafeQueueCell<T, VecDequeStorage<T>>
where
T: PartialEq + Clone + Send + Sync + 'static,
{
pub fn new(ctx: &ThreadSafeContext) -> Self {
Self::with_storage(ctx, VecDequeStorage::unbounded())
}
pub fn with_capacity(ctx: &ThreadSafeContext, capacity: usize) -> Self {
Self::with_storage(ctx, VecDequeStorage::with_capacity(capacity))
}
}
impl<T, S> ThreadSafeQueueCell<T, S>
where
T: PartialEq + Clone + Send + Sync + 'static,
S: QueueStorage<T> + Send + 'static,
{
pub fn with_storage(ctx: &ThreadSafeContext, storage: S) -> Self {
let closed_val = storage.is_closed();
let capacity = storage.capacity();
let storage = Arc::new(Mutex::new(storage));
let head = {
let storage = Arc::clone(&storage);
ctx.computed(move |_| storage.lock().expect("queue storage").peek().cloned())
};
let len = {
let storage = Arc::clone(&storage);
ctx.computed(move |_| storage.lock().expect("queue storage").len())
};
let is_empty = {
let storage = Arc::clone(&storage);
ctx.computed(move |_| storage.lock().expect("queue storage").len() == 0)
};
let is_full = {
let storage = Arc::clone(&storage);
ctx.computed(move |_| {
let s = storage.lock().expect("queue storage");
match s.capacity() {
Some(cap) => s.len() >= cap,
None => false,
}
})
};
Self {
inner: Arc::new(Inner {
storage,
capacity,
head,
len,
is_empty,
is_full,
closed: ctx.source(closed_val),
}),
}
}
fn invalidate_readers(
&self,
ctx: &ThreadSafeContext,
len_before: usize,
len_after: usize,
head_changed: bool,
) {
let is_empty_changed = (len_before == 0) != (len_after == 0);
let is_full_changed = self
.inner
.capacity
.map(|c| (len_before >= c) != (len_after >= c))
.unwrap_or(false);
ctx.batch(|_| {
ctx.clear(&self.inner.len);
if is_empty_changed {
ctx.clear(&self.inner.is_empty);
}
if is_full_changed {
ctx.clear(&self.inner.is_full);
}
if head_changed {
ctx.clear(&self.inner.head);
}
});
}
pub fn try_push(&self, ctx: &ThreadSafeContext, value: T) -> Result<(), QueuePushError> {
let (result, len_before) = {
let mut s = self.inner.storage.lock().expect("queue storage");
let len_before = s.len();
(s.try_push(value), len_before)
};
if result.is_ok() {
self.invalidate_readers(ctx, len_before, len_before + 1, len_before == 0);
}
result
}
pub fn try_pop(&self, ctx: &ThreadSafeContext) -> Result<T, QueuePopError> {
let (result, len_before) = {
let mut s = self.inner.storage.lock().expect("queue storage");
let len_before = s.len();
(s.try_pop(), len_before)
};
if result.is_ok() {
self.invalidate_readers(ctx, len_before, len_before - 1, true);
}
result
}
pub fn close(&self, ctx: &ThreadSafeContext) {
let newly_closed = {
let mut s = self.inner.storage.lock().expect("queue storage");
let was = s.is_closed();
s.close();
!was
};
if newly_closed {
ctx.set(&self.inner.closed, true);
}
}
pub fn head(&self, ctx: &ThreadSafeContext) -> Option<T> {
ctx.get(&self.inner.head)
}
pub fn len(&self, ctx: &ThreadSafeContext) -> usize {
ctx.get(&self.inner.len)
}
pub fn is_empty(&self, ctx: &ThreadSafeContext) -> bool {
ctx.get(&self.inner.is_empty)
}
pub fn is_full(&self, ctx: &ThreadSafeContext) -> bool {
ctx.get(&self.inner.is_full)
}
pub fn closed(&self, ctx: &ThreadSafeContext) -> bool {
ctx.get(&self.inner.closed)
}
pub fn capacity(&self) -> Option<usize> {
self.inner.capacity
}
}