1#![cfg_attr(test, allow(clippy::unwrap_used, reason = "test scope"))]
19
20pub mod broadcast;
21pub mod condvar;
23pub mod mpsc;
25pub mod mutex;
27pub mod notify;
28pub mod oneshot;
30pub mod rwlock;
31pub mod semaphore;
32pub(crate) mod subscribers;
33pub(crate) mod wait_queue;
34pub mod watch;
35
36pub 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.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.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 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 assert!(rx.poll_recv(&mut context).is_pending());
296
297 tx.send(42).unwrap();
299
300 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 assert!(f1.as_mut().poll(&mut context).is_pending());
320 assert!(f2.as_mut().poll(&mut context).is_pending());
321
322 notify.notify_one();
324
325 drop(f1);
327
328 assert!(f2.as_mut().poll(&mut context).is_ready());
330
331 let notify = Notify::new();
333 let mut f1 = Box::pin(notify.notified());
334
335 assert!(f1.as_mut().poll(&mut context).is_pending());
337
338 notify.notify_one();
340
341 drop(f1);
343
344 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 let waker = futures::task::noop_waker();
356 let mut context = Context::from_waker(&waker);
357
358 let notify = Notify::new();
359 notify.notify_one(); notify.notify_waiters(); 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}