moirai_async/sync/notify.rs
1//! Notification primitive for efficient async task coordination
2//!
3//! Provides notification mechanisms for waking up waiting async tasks,
4//! following SLAP principle with focused responsibility. Waiter-queue
5//! mechanics live in `WaitQueue`; this module keeps only the notify
6//! admission state (the stored single permit) and the grant-restoration
7//! policy for cancelled futures.
8
9#![expect(
10 clippy::unwrap_used,
11 reason = "ratchet MOIRAI-UNWRAP-1: pre-existing debt"
12)]
13
14use std::future::Future;
15use std::pin::Pin;
16use std::sync::Mutex;
17use std::task::{Context, Poll};
18
19use crate::sync::wait_queue::{WaitQueue, WaiterPoll};
20
21/// Grant payload distinguishing how a waiter was notified: a `notify_one`
22/// grant is a transferable single permit (restored to the next waiter or the
23/// stored-permit slot when the granted future is cancelled), while a
24/// `notify_waiters` grant is a broadcast wakeup that is not restored.
25#[derive(Clone, Copy, PartialEq, Eq, Debug)]
26enum NotifyGrant {
27 One,
28 All,
29}
30
31/// Notification primitive for efficient task coordination
32pub struct Notify {
33 state: Mutex<NotifyState>,
34}
35
36struct NotifyState {
37 /// Single stored permit from a `notify_one` issued with no waiters.
38 notified: bool,
39 waiters: WaitQueue<NotifyGrant>,
40}
41
42impl Notify {
43 /// Create a new notification primitive
44 pub fn new() -> Self {
45 Self {
46 state: Mutex::new(NotifyState {
47 notified: false,
48 waiters: WaitQueue::new(),
49 }),
50 }
51 }
52
53 /// Wait for a notification
54 pub fn notified(&self) -> NotifyFuture<'_> {
55 NotifyFuture {
56 notify: self,
57 id: None,
58 }
59 }
60
61 /// Notify one waiting task
62 ///
63 /// The waker is taken under the state lock and woken after it is released:
64 /// `Waker::wake` may poll the task inline on this thread, and that poll
65 /// re-locks this state — waking under the lock would self-deadlock. Same
66 /// discipline as `notify_waiters` below and `hybrid::notify`.
67 pub fn notify_one(&self) {
68 let waker = {
69 let mut state = self.state.lock().unwrap();
70 let waker = state.waiters.grant_oldest(NotifyGrant::One);
71 if waker.is_none() {
72 state.notified = true;
73 }
74 waker
75 };
76 if let Some(waker) = waker {
77 waker.wake();
78 }
79 }
80
81 /// Notify all waiting tasks.
82 ///
83 /// Wakes every currently-registered waiter. This is independent of the
84 /// single-permit `notify_one` mechanism: a permit stored by a prior
85 /// `notify_one` (issued with no waiters present) is left intact, so a
86 /// subsequent `notified()` still observes it.
87 pub fn notify_waiters(&self) {
88 let mut state = self.state.lock().unwrap();
89 let wakers = state.waiters.grant_all(NotifyGrant::All);
90 drop(state);
91
92 for waker in wakers {
93 waker.wake();
94 }
95 }
96}
97
98impl Default for Notify {
99 fn default() -> Self {
100 Self::new()
101 }
102}
103
104/// Future for waiting on notifications
105pub struct NotifyFuture<'a> {
106 notify: &'a Notify,
107 id: Option<u64>,
108}
109
110impl<'a> Future for NotifyFuture<'a> {
111 type Output = ();
112
113 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
114 let mut state = self.notify.state.lock().unwrap();
115
116 // 1. Check if a stored permit is available
117 if state.notified {
118 state.notified = false;
119 if let Some(id) = self.id.take() {
120 // Parity with the pre-consolidation state machine: the entry
121 // is removed regardless of grant state, so a `One` grant on
122 // our own entry is consumed together with the stored permit.
123 let _removed_grant = state.waiters.deregister(id);
124 }
125 return Poll::Ready(());
126 }
127
128 // 2. Check if already registered
129 if let Some(id) = self.id {
130 match state.waiters.poll_waiter(id, cx.waker()) {
131 WaiterPoll::Granted(_) => {
132 self.id = None;
133 return Poll::Ready(());
134 }
135 WaiterPoll::Pending => return Poll::Pending,
136 // Our registration was lost; re-register below.
137 WaiterPoll::NotRegistered => {}
138 }
139 }
140
141 // 3. Register a (new) waiter.
142 self.id = Some(state.waiters.register(cx.waker().clone()));
143 Poll::Pending
144 }
145}
146
147impl<'a> Drop for NotifyFuture<'a> {
148 fn drop(&mut self) {
149 // The restored permit's waker leaves the state lock before it is woken;
150 // `Waker::wake` may poll the task inline on this thread, and that poll
151 // re-locks this state. Same discipline as `notify_one`.
152 let waker = if let Some(id) = self.id {
153 let mut state = match self.notify.state.lock() {
154 Ok(state) => state,
155 Err(_) => return,
156 };
157 // If we were holding a single-task permit but never observed
158 // it, hand it to the next pending waiter (or store it) so it
159 // is not lost. Broadcast (`All`) grants are not restored.
160 if state.waiters.deregister(id) == Some(NotifyGrant::One) {
161 let waker = state.waiters.grant_oldest(NotifyGrant::One);
162 if waker.is_none() {
163 state.notified = true;
164 }
165 waker
166 } else {
167 None
168 }
169 } else {
170 None
171 };
172 if let Some(waker) = waker {
173 waker.wake();
174 }
175 }
176}