use std::sync::PoisonError;
use std::time::Instant;
use crate::sync::{AtomicUsize, Condvar, Mutex, MutexGuard, Ordering::Relaxed};
use crate::utils::CachePadded;
#[derive(Debug)]
pub(crate) struct WaitQueue {
waiters: CachePadded<AtomicUsize>,
lock: Mutex<()>,
cv: Condvar,
}
impl WaitQueue {
pub(crate) fn new() -> Self {
Self {
waiters: CachePadded::new(AtomicUsize::new(0)),
lock: Mutex::new(()),
cv: Condvar::new(),
}
}
fn guard(&self) -> MutexGuard<'_, ()> {
self.lock.lock().unwrap_or_else(PoisonError::into_inner)
}
#[inline]
pub(crate) fn notify_one(&self) {
if self.waiters.load(Relaxed) == 0 {
return;
}
let _guard = self.guard();
self.cv.notify_one();
}
#[inline]
pub(crate) fn has_waiters(&self) -> bool {
self.waiters.load(Relaxed) != 0
}
pub(crate) fn notify_all(&self) {
let _guard = self.guard();
self.cv.notify_all();
}
pub(crate) fn wait_until(
&self,
mut ready: impl FnMut() -> bool,
deadline: Option<Instant>,
) -> bool {
let mut guard = self.guard();
self.waiters.fetch_add(1, Relaxed);
let observed = loop {
if ready() {
break true;
}
match deadline {
None => {
guard = self.cv.wait(guard).unwrap_or_else(PoisonError::into_inner);
}
Some(deadline) => {
let now = Instant::now();
if now >= deadline {
break false;
}
guard = self
.cv
.wait_timeout(guard, deadline - now)
.unwrap_or_else(PoisonError::into_inner)
.0;
}
}
};
self.waiters.fetch_sub(1, Relaxed);
drop(guard);
observed
}
}