use std::cell::{Cell, RefCell};
use std::collections::VecDeque;
use std::future::Future;
use std::pin::Pin;
use std::rc::Rc;
use std::task::{Context, Poll, Waker};
struct EventInner {
is_set: Cell<bool>,
waiters: RefCell<Vec<Waker>>,
}
#[derive(Clone)]
pub struct Event {
inner: Rc<EventInner>,
}
impl Event {
#[allow(clippy::new_without_default)]
pub fn new() -> Event {
Event { inner: Rc::new(EventInner { is_set: Cell::new(false), waiters: RefCell::new(Vec::new()) }) }
}
pub fn set(&self) {
self.inner.is_set.set(true);
for w in self.inner.waiters.borrow_mut().drain(..) {
w.wake();
}
}
pub fn clear(&self) {
self.inner.is_set.set(false);
}
pub fn is_set(&self) -> bool {
self.inner.is_set.get()
}
pub fn wait(&self) -> EventWait {
EventWait { inner: self.inner.clone() }
}
}
pub struct EventWait {
inner: Rc<EventInner>,
}
impl Future for EventWait {
type Output = ();
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
if self.inner.is_set.get() {
Poll::Ready(())
} else {
self.inner.waiters.borrow_mut().push(cx.waker().clone());
Poll::Pending
}
}
}
struct LockWaiter {
granted: Cell<bool>,
waker: RefCell<Option<Waker>>,
}
struct LockInner {
locked: Cell<bool>,
waiters: RefCell<VecDeque<Rc<LockWaiter>>>,
}
impl LockInner {
fn release(&self) {
loop {
let next = self.waiters.borrow_mut().pop_front();
match next {
Some(w) => {
w.granted.set(true);
if let Some(waker) = w.waker.borrow_mut().take() {
waker.wake();
}
return;
}
None => {
self.locked.set(false);
return;
}
}
}
}
}
#[derive(Clone)]
pub struct Lock {
inner: Rc<LockInner>,
}
impl Lock {
#[allow(clippy::new_without_default)]
pub fn new() -> Lock {
Lock { inner: Rc::new(LockInner { locked: Cell::new(false), waiters: RefCell::new(VecDeque::new()) }) }
}
pub fn locked(&self) -> bool {
self.inner.locked.get()
}
pub fn acquire(&self) -> Acquire {
Acquire { inner: self.inner.clone(), waiter: None, acquired: false }
}
}
pub struct Acquire {
inner: Rc<LockInner>,
waiter: Option<Rc<LockWaiter>>,
acquired: bool,
}
impl Future for Acquire {
type Output = LockGuard;
fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<LockGuard> {
match &self.waiter {
None => {
if !self.inner.locked.get() {
self.inner.locked.set(true);
self.acquired = true;
Poll::Ready(LockGuard { inner: self.inner.clone() })
} else {
let w = Rc::new(LockWaiter {
granted: Cell::new(false),
waker: RefCell::new(Some(cx.waker().clone())),
});
self.inner.waiters.borrow_mut().push_back(w.clone());
self.waiter = Some(w);
Poll::Pending
}
}
Some(w) => {
if w.granted.get() {
self.acquired = true;
Poll::Ready(LockGuard { inner: self.inner.clone() })
} else {
*w.waker.borrow_mut() = Some(cx.waker().clone());
Poll::Pending
}
}
}
}
}
impl Drop for Acquire {
fn drop(&mut self) {
if self.acquired {
return; }
if let Some(w) = &self.waiter {
if w.granted.get() {
self.inner.release();
} else {
self.inner.waiters.borrow_mut().retain(|x| !Rc::ptr_eq(x, w));
}
}
}
}
pub struct LockGuard {
inner: Rc<LockInner>,
}
impl Drop for LockGuard {
fn drop(&mut self) {
self.inner.release();
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::executor;
use crate::testing::{assert_pending, block_on};
use std::cell::RefCell;
use std::rc::Rc;
#[test]
fn wait_before_set_wakes() {
block_on(async {
let ev = Event::new();
let setter = ev.clone();
executor::spawn(async move {
setter.set();
});
ev.wait().await; });
}
#[test]
fn set_latches_until_cleared() {
block_on(async {
let ev = Event::new();
ev.set();
ev.wait().await; assert!(ev.is_set());
});
}
#[test]
fn clear_makes_a_later_wait_block_again() {
let ev = Event::new();
ev.set();
ev.clear();
assert!(!ev.is_set());
assert_pending(async move { ev.wait().await });
}
#[test]
fn one_set_wakes_every_waiter() {
block_on(async {
let ev = Event::new();
let woken = Rc::new(RefCell::new(0));
for _ in 0..3 {
let (e, w) = (ev.clone(), woken.clone());
executor::spawn(async move {
e.wait().await;
*w.borrow_mut() += 1;
});
}
crate::executor::current().run_until_idle();
ev.set();
crate::executor::current().run_until_idle();
assert_eq!(*woken.borrow(), 3, "all three saw one set");
});
}
#[test]
fn lock_is_fifo_fair() {
block_on(async {
let lock = Lock::new();
let order = Rc::new(RefCell::new(Vec::new()));
let held = lock.acquire().await;
for who in 1..=3 {
let (l, o) = (lock.clone(), order.clone());
executor::spawn(async move {
let _g = l.acquire().await;
o.borrow_mut().push(who);
});
crate::executor::current().run_until_idle();
}
drop(held);
crate::executor::current().run_until_idle();
assert_eq!(*order.borrow(), vec![1, 2, 3], "granted in request order");
});
}
}