Skip to main content

moirai_async/sync/
mod.rs

1//! Advanced async synchronization primitives for Moirai
2//!
3//! This module provides async-aware synchronization that integrates with
4//! Moirai's unified runtime. Following SLAP principle, each synchronization
5//! primitive is implemented in its own focused module.
6//!
7//! # Wake discipline
8//!
9//! Every primitive here takes the waker — or the wakers — out from under its
10//! state lock and wakes only after the guard is released. `Waker::wake` may poll
11//! the task inline on the calling thread, and that poll re-locks the same state,
12//! so waking under the lock is a self-deadlock. `hybrid::notify` states the rule
13//! for its registries, and `timer`'s driver follows it by dropping its guard
14//! before waking. A new release/notify path belongs to this rule rather than to
15//! a local judgement — the sites that got it wrong were the ones that decided
16//! per-site.
17
18#![cfg_attr(test, allow(clippy::unwrap_used, reason = "test scope"))]
19
20pub mod broadcast;
21/// Async condition variable.
22pub mod condvar;
23/// Bounded async multi-producer single-consumer channel.
24pub mod mpsc;
25/// Async mutual-exclusion lock.
26pub mod mutex;
27pub mod notify;
28/// Single-value channel completed by one send.
29pub mod oneshot;
30pub mod rwlock;
31pub mod semaphore;
32pub(crate) mod subscribers;
33pub(crate) mod wait_queue;
34pub mod watch;
35
36// Re-export public types for convenience
37pub use broadcast::{Broadcast, BroadcastError, BroadcastReceiver, BroadcastRecv, BroadcastSender};
38pub use condvar::Condvar;
39pub use mpsc::{Receiver as MpscReceiver, Sender as MpscSender, channel as mpsc_channel};
40pub use mutex::{Mutex, MutexGuard, MutexLockFuture};
41pub use notify::{Notify, NotifyFuture};
42pub use oneshot::{
43    Receiver as OneshotReceiver, Sender as OneshotSender, channel as oneshot_channel,
44};
45pub use rwlock::{RwLock, RwLockReadFuture, RwLockWriteFuture};
46pub use semaphore::{Semaphore, SemaphoreAcquire, SemaphorePermit};
47pub use watch::{Watch, WatchChanged, WatchError, WatchReceiver, WatchSender};
48
49#[cfg(all(test, not(target_arch = "wasm32")))]
50mod tests {
51    use super::*;
52    use crate::executor::AsyncExecutor;
53    use std::future::Future;
54    use std::sync::Arc;
55    use std::task::{Context, Poll};
56    use std::time::Duration;
57
58    fn run_executor_to_completion_with_limit(executor: &AsyncExecutor, limit: usize) {
59        for _ in 0..limit {
60            executor.process_pending_tasks();
61            executor
62                .reactor()
63                .run_iteration(Some(Duration::from_millis(0)))
64                .ok();
65            if executor.stats().tasks_pending == 0 {
66                break;
67            }
68        }
69    }
70
71    #[test]
72    fn test_semaphore_basic() {
73        let sem = Semaphore::new(2);
74        assert_eq!(sem.available_permits(), 2);
75    }
76
77    #[test]
78    fn test_semaphore_async() {
79        let executor = AsyncExecutor::new().unwrap();
80        let sem = Arc::new(Semaphore::new(2));
81        let notify = Arc::new(Notify::new());
82
83        let sem_clone1 = sem.clone();
84        let notify_clone1 = notify.clone();
85        let handle1 = executor.spawn(async move {
86            let _permit = sem_clone1.acquire().await;
87            notify_clone1.notified().await;
88            1
89        });
90
91        let sem_clone2 = sem.clone();
92        let notify_clone2 = notify.clone();
93        let handle2 = executor.spawn(async move {
94            let _permit = sem_clone2.acquire().await;
95            notify_clone2.notified().await;
96            2
97        });
98
99        let sem_clone3 = sem.clone();
100        let notify_clone3 = notify.clone();
101        let handle3 = executor.spawn(async move {
102            let _permit = sem_clone3.acquire().await;
103            notify_clone3.notified().await;
104            3
105        });
106
107        run_executor_to_completion_with_limit(&executor, 100);
108
109        let waker = futures::task::noop_waker();
110        let mut context = Context::from_waker(&waker);
111
112        let mut handle1 = Box::pin(handle1);
113        let mut handle2 = Box::pin(handle2);
114        let mut handle3 = Box::pin(handle3);
115
116        assert!(handle1.as_mut().poll(&mut context).is_pending());
117        assert!(handle2.as_mut().poll(&mut context).is_pending());
118        assert!(handle3.as_mut().poll(&mut context).is_pending());
119
120        // Notify one task to let it complete and release its permit
121        notify.notify_one();
122
123        run_executor_to_completion_with_limit(&executor, 100);
124
125        let p1 = handle1.as_mut().poll(&mut context);
126        let p2 = handle2.as_mut().poll(&mut context);
127        let p3 = handle3.as_mut().poll(&mut context);
128
129        let r1 = p1.is_ready();
130        let r2 = p2.is_ready();
131        assert_eq!(if r1 { 1 } else { 0 } + if r2 { 1 } else { 0 }, 1);
132        assert!(p3.is_pending());
133
134        // Notify all remaining tasks
135        notify.notify_waiters();
136
137        run_executor_to_completion_with_limit(&executor, 100);
138
139        if !r1 {
140            assert!(handle1.as_mut().poll(&mut context).is_ready());
141        }
142        if !r2 {
143            assert!(handle2.as_mut().poll(&mut context).is_ready());
144        }
145        assert!(handle3.as_mut().poll(&mut context).is_ready());
146    }
147
148    #[test]
149    fn test_notify_async() {
150        let executor = AsyncExecutor::new().unwrap();
151        let notify = Arc::new(Notify::new());
152
153        let notify_clone = notify.clone();
154        let handle = executor.spawn(async move {
155            notify_clone.notified().await;
156            42
157        });
158
159        let mut handle = Box::pin(handle);
160        let waker = futures::task::noop_waker();
161        let mut context = Context::from_waker(&waker);
162
163        run_executor_to_completion_with_limit(&executor, 5);
164        assert!(handle.as_mut().poll(&mut context).is_pending());
165
166        notify.notify_one();
167
168        run_executor_to_completion_with_limit(&executor, 5);
169        assert!(matches!(
170            handle.as_mut().poll(&mut context),
171            Poll::Ready(42)
172        ));
173    }
174
175    #[test]
176    fn test_rwlock_async() {
177        let executor = AsyncExecutor::new().unwrap();
178        let lock = Arc::new(RwLock::new(100));
179
180        let lock_clone1 = lock.clone();
181        let handle_read1 = executor.spawn(async move {
182            let guard = lock_clone1.read().await;
183            *guard
184        });
185
186        let lock_clone2 = lock.clone();
187        let handle_read2 = executor.spawn(async move {
188            let guard = lock_clone2.read().await;
189            *guard
190        });
191
192        run_executor_to_completion_with_limit(&executor, 5);
193
194        let waker = futures::task::noop_waker();
195        let mut context = Context::from_waker(&waker);
196        let mut h_r1 = Box::pin(handle_read1);
197        let mut h_r2 = Box::pin(handle_read2);
198
199        assert!(matches!(h_r1.as_mut().poll(&mut context), Poll::Ready(100)));
200        assert!(matches!(h_r2.as_mut().poll(&mut context), Poll::Ready(100)));
201
202        // Drops the read handles, unlocking RwLock
203        drop(h_r1);
204        drop(h_r2);
205
206        let lock_clone3 = lock.clone();
207        let handle_write = executor.spawn(async move {
208            let mut guard = lock_clone3.write().await;
209            *guard += 50;
210            *guard
211        });
212
213        run_executor_to_completion_with_limit(&executor, 5);
214        let mut h_w = Box::pin(handle_write);
215        assert!(matches!(h_w.as_mut().poll(&mut context), Poll::Ready(150)));
216    }
217
218    #[test]
219    fn test_watch_async() {
220        let executor = AsyncExecutor::new().unwrap();
221        let (tx, rx) = Watch::new(10);
222
223        let tx = Arc::new(tx);
224        let rx_clone = rx.clone();
225        let handle = executor.spawn(async move {
226            let mut rx = rx_clone;
227            rx.changed().await.unwrap();
228            rx.borrow()
229        });
230
231        run_executor_to_completion_with_limit(&executor, 5);
232        let waker = futures::task::noop_waker();
233        let mut context = Context::from_waker(&waker);
234        let mut handle = Box::pin(handle);
235        assert!(handle.as_mut().poll(&mut context).is_pending());
236
237        tx.send(20).unwrap();
238
239        run_executor_to_completion_with_limit(&executor, 5);
240        assert!(matches!(
241            handle.as_mut().poll(&mut context),
242            Poll::Ready(20)
243        ));
244    }
245
246    #[test]
247    fn test_broadcast_async() {
248        let executor = AsyncExecutor::new().unwrap();
249        let (tx, rx1) = Broadcast::new(2);
250        let rx2 = rx1.clone();
251
252        let handle1 = executor.spawn(async move {
253            let mut rx = rx1;
254            let m1 = rx.recv().await.unwrap();
255            let m2 = rx.recv().await.unwrap();
256            (m1, m2)
257        });
258
259        let handle2 = executor.spawn(async move {
260            let mut rx = rx2;
261            let m1 = rx.recv().await.unwrap();
262            let m2 = rx.recv().await.unwrap();
263            (m1, m2)
264        });
265
266        run_executor_to_completion_with_limit(&executor, 5);
267
268        tx.send(10).unwrap();
269        tx.send(20).unwrap();
270
271        run_executor_to_completion_with_limit(&executor, 10);
272
273        let waker = futures::task::noop_waker();
274        let mut context = Context::from_waker(&waker);
275        let mut h1 = Box::pin(handle1);
276        let mut h2 = Box::pin(handle2);
277
278        assert!(matches!(
279            h1.as_mut().poll(&mut context),
280            Poll::Ready((10, 20))
281        ));
282        assert!(matches!(
283            h2.as_mut().poll(&mut context),
284            Poll::Ready((10, 20))
285        ));
286    }
287
288    #[test]
289    fn test_broadcast_waker_registration() {
290        let (tx, mut rx) = Broadcast::new(2);
291        let waker = futures::task::noop_waker();
292        let mut context = Context::from_waker(&waker);
293
294        // Initially empty, should return Pending and register the waker
295        assert!(rx.poll_recv(&mut context).is_pending());
296
297        // Send a message
298        tx.send(42).unwrap();
299
300        // Now it should return Ready(Ok(42))
301        assert!(matches!(rx.poll_recv(&mut context), Poll::Ready(Ok(42))));
302    }
303
304    #[test]
305    fn test_notify_cancellation_safety() {
306        use crate::sync::Notify;
307        use futures::task::noop_waker;
308        use std::future::Future;
309        use std::task::Context;
310
311        let notify = Notify::new();
312        let waker = noop_waker();
313        let mut context = Context::from_waker(&waker);
314
315        let mut f1 = Box::pin(notify.notified());
316        let mut f2 = Box::pin(notify.notified());
317
318        // Poll both to register them as waiters
319        assert!(f1.as_mut().poll(&mut context).is_pending());
320        assert!(f2.as_mut().poll(&mut context).is_pending());
321
322        // Notify once, which grants permit to f1
323        notify.notify_one();
324
325        // Drop f1 before it is polled. This should transfer the permit to f2.
326        drop(f1);
327
328        // Now f2 should be ready
329        assert!(f2.as_mut().poll(&mut context).is_ready());
330
331        // Second case: no other pending waiters, permit is restored to state
332        let notify = Notify::new();
333        let mut f1 = Box::pin(notify.notified());
334
335        // Poll to register f1 as a waiter
336        assert!(f1.as_mut().poll(&mut context).is_pending());
337
338        // Notify once, granting permit to f1
339        notify.notify_one();
340
341        // Drop f1. This should restore the permit back to notify's state.
342        drop(f1);
343
344        // Now a new future should be ready immediately
345        let mut f3 = Box::pin(notify.notified());
346        assert!(f3.as_mut().poll(&mut context).is_ready());
347    }
348
349    #[test]
350    fn notify_waiters_preserves_stored_notify_one_permit() {
351        // A `notify_one` issued with no waiters stores a single permit. An
352        // unrelated `notify_waiters` (which only wakes currently-registered
353        // waiters) must not destroy that stored permit, or the next
354        // `notified()` would block forever.
355        let waker = futures::task::noop_waker();
356        let mut context = Context::from_waker(&waker);
357
358        let notify = Notify::new();
359        notify.notify_one(); // store a permit (no waiters registered)
360        notify.notify_waiters(); // must leave the stored permit intact
361
362        let mut fut = Box::pin(notify.notified());
363        assert!(
364            fut.as_mut().poll(&mut context).is_ready(),
365            "notify_one permit must survive notify_waiters when no waiters are registered"
366        );
367    }
368}