Skip to main content

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}