1use std::cell::{Cell, RefCell};
5use std::collections::VecDeque;
6use std::future::Future;
7use std::pin::Pin;
8use std::rc::Rc;
9use std::task::{Context, Poll, Waker};
10
11struct EventInner {
16 is_set: Cell<bool>,
17 waiters: RefCell<Vec<Waker>>,
18}
19
20#[derive(Clone)]
22pub struct Event {
23 inner: Rc<EventInner>,
24}
25
26impl Event {
27 #[allow(clippy::new_without_default)]
28 pub fn new() -> Event {
29 Event { inner: Rc::new(EventInner { is_set: Cell::new(false), waiters: RefCell::new(Vec::new()) }) }
30 }
31
32 pub fn set(&self) {
33 self.inner.is_set.set(true);
34 for w in self.inner.waiters.borrow_mut().drain(..) {
35 w.wake();
36 }
37 }
38
39 pub fn clear(&self) {
40 self.inner.is_set.set(false);
41 }
42
43 pub fn is_set(&self) -> bool {
44 self.inner.is_set.get()
45 }
46
47 pub fn wait(&self) -> EventWait {
48 EventWait { inner: self.inner.clone() }
49 }
50}
51
52pub struct EventWait {
53 inner: Rc<EventInner>,
54}
55
56impl Future for EventWait {
57 type Output = ();
58 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> {
59 if self.inner.is_set.get() {
60 Poll::Ready(())
61 } else {
62 self.inner.waiters.borrow_mut().push(cx.waker().clone());
63 Poll::Pending
64 }
65 }
66}
67
68struct LockWaiter {
73 granted: Cell<bool>,
74 waker: RefCell<Option<Waker>>,
75}
76
77struct LockInner {
78 locked: Cell<bool>,
79 waiters: RefCell<VecDeque<Rc<LockWaiter>>>,
80}
81
82impl LockInner {
83 fn release(&self) {
85 loop {
86 let next = self.waiters.borrow_mut().pop_front();
87 match next {
88 Some(w) => {
89 w.granted.set(true);
90 if let Some(waker) = w.waker.borrow_mut().take() {
91 waker.wake();
92 }
93 return;
95 }
96 None => {
97 self.locked.set(false);
98 return;
99 }
100 }
101 }
102 }
103}
104
105#[derive(Clone)]
106pub struct Lock {
107 inner: Rc<LockInner>,
108}
109
110impl Lock {
111 #[allow(clippy::new_without_default)]
112 pub fn new() -> Lock {
113 Lock { inner: Rc::new(LockInner { locked: Cell::new(false), waiters: RefCell::new(VecDeque::new()) }) }
114 }
115
116 pub fn locked(&self) -> bool {
117 self.inner.locked.get()
118 }
119
120 pub fn acquire(&self) -> Acquire {
121 Acquire { inner: self.inner.clone(), waiter: None, acquired: false }
122 }
123}
124
125pub struct Acquire {
126 inner: Rc<LockInner>,
127 waiter: Option<Rc<LockWaiter>>,
128 acquired: bool,
129}
130
131impl Future for Acquire {
132 type Output = LockGuard;
133 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<LockGuard> {
134 match &self.waiter {
135 None => {
136 if !self.inner.locked.get() {
137 self.inner.locked.set(true);
138 self.acquired = true;
139 Poll::Ready(LockGuard { inner: self.inner.clone() })
140 } else {
141 let w = Rc::new(LockWaiter {
142 granted: Cell::new(false),
143 waker: RefCell::new(Some(cx.waker().clone())),
144 });
145 self.inner.waiters.borrow_mut().push_back(w.clone());
146 self.waiter = Some(w);
147 Poll::Pending
148 }
149 }
150 Some(w) => {
151 if w.granted.get() {
152 self.acquired = true;
153 Poll::Ready(LockGuard { inner: self.inner.clone() })
154 } else {
155 *w.waker.borrow_mut() = Some(cx.waker().clone());
156 Poll::Pending
157 }
158 }
159 }
160 }
161}
162
163impl Drop for Acquire {
164 fn drop(&mut self) {
165 if self.acquired {
166 return; }
168 if let Some(w) = &self.waiter {
169 if w.granted.get() {
170 self.inner.release();
172 } else {
173 self.inner.waiters.borrow_mut().retain(|x| !Rc::ptr_eq(x, w));
175 }
176 }
177 }
178}
179
180pub struct LockGuard {
182 inner: Rc<LockInner>,
183}
184
185impl Drop for LockGuard {
186 fn drop(&mut self) {
187 self.inner.release();
188 }
189}
190
191#[cfg(test)]
196mod tests {
197 use super::*;
198 use crate::executor;
199 use crate::testing::{assert_pending, block_on};
200 use std::cell::RefCell;
201 use std::rc::Rc;
202
203 #[test]
204 fn wait_before_set_wakes() {
205 block_on(async {
206 let ev = Event::new();
207 let setter = ev.clone();
208 executor::spawn(async move {
209 setter.set();
210 });
211 ev.wait().await; });
213 }
214
215 #[test]
221 fn set_latches_until_cleared() {
222 block_on(async {
223 let ev = Event::new();
224 ev.set();
225 ev.wait().await; assert!(ev.is_set());
227 });
228 }
229
230 #[test]
231 fn clear_makes_a_later_wait_block_again() {
232 let ev = Event::new();
233 ev.set();
234 ev.clear();
235 assert!(!ev.is_set());
236 assert_pending(async move { ev.wait().await });
237 }
238
239 #[test]
240 fn one_set_wakes_every_waiter() {
241 block_on(async {
242 let ev = Event::new();
243 let woken = Rc::new(RefCell::new(0));
244 for _ in 0..3 {
245 let (e, w) = (ev.clone(), woken.clone());
246 executor::spawn(async move {
247 e.wait().await;
248 *w.borrow_mut() += 1;
249 });
250 }
251 crate::executor::current().run_until_idle();
253 ev.set();
254 crate::executor::current().run_until_idle();
255 assert_eq!(*woken.borrow(), 3, "all three saw one set");
256 });
257 }
258
259 #[test]
260 fn lock_is_fifo_fair() {
261 block_on(async {
262 let lock = Lock::new();
263 let order = Rc::new(RefCell::new(Vec::new()));
264
265 let held = lock.acquire().await; for who in 1..=3 {
268 let (l, o) = (lock.clone(), order.clone());
269 executor::spawn(async move {
270 let _g = l.acquire().await;
271 o.borrow_mut().push(who);
272 });
273 crate::executor::current().run_until_idle();
275 }
276
277 drop(held);
278 crate::executor::current().run_until_idle();
279 assert_eq!(*order.borrow(), vec![1, 2, 3], "granted in request order");
280 });
281 }
282}