Skip to main content

rivet/sync/
signal.rs

1//! A one-shot, latching wake-up handoff between an external context (an
2//! ISR, another hart, another task) and exactly one waiting task.
3//!
4//! This is the building block for a peripheral driver's async completion:
5//! the driver's interrupt handler calls [`Signal::signal`] (documented
6//! ISR-safe — no critical section, no allocator, no `current_task()`
7//! requirement) when a hardware transfer finishes, and the cooperative
8//! task driving that transfer `.await`s [`Signal::wait`].
9//!
10//! `Signal` carries no value — the driver reads the peripheral's own
11//! status registers to find out what happened; this avoids an allocator
12//! and `Cell<T>` variance questions entirely. If a future version needs a
13//! value-carrying variant, that's a new type (`Signal<T>`), not a change
14//! to this one.
15//!
16//! Unlike [`crate::sync::Semaphore`] (a 32-wide per-priority waiter
17//! bitmap, because several tasks may legitimately contend for one
18//! resource) or [`crate::sync::Channel`] (two independently-driven waiter
19//! slots), `Signal` has exactly one waiter slot: a peripheral has one
20//! owning task by construction. Registering a second concurrent waiter is
21//! a caller bug, not a supported use case — see [`Wait::poll`].
22
23use core::future::Future;
24use core::pin::Pin;
25use core::task::{Context, Poll};
26
27use crate::waker;
28
29const NO_WAITER: u32 = 0xFFFF_FFFF;
30
31/// See the [module docs](self).
32pub struct Signal {
33    /// Latched by `signal()`, consumed by a successful [`Signal::try_take`].
34    /// Latching (rather than a plain wake) is what makes "the ISR fired
35    /// before `wait()`'s first poll" correct instead of a lost wakeup.
36    fired: crate::sync::atomic::AtomicBool,
37    /// The single registered cooperative waiter, or `NO_WAITER`.
38    waiter: crate::sync::atomic::AtomicU32,
39}
40
41impl Signal {
42    /// Create a new, unfired signal.
43    #[cfg(not(loom))]
44    pub const fn new() -> Self {
45        Self {
46            fired: crate::sync::atomic::AtomicBool::new(false),
47            waiter: crate::sync::atomic::AtomicU32::new(NO_WAITER),
48        }
49    }
50
51    /// Loom's atomics are not const-constructible; runtime constructor used
52    /// by the loom models.
53    #[cfg(loom)]
54    pub fn new() -> Self {
55        Self {
56            fired: crate::sync::atomic::AtomicBool::new(false),
57            waiter: crate::sync::atomic::AtomicU32::new(NO_WAITER),
58        }
59    }
60
61    /// External/ISR side: latch the signal and wake the registered waiter,
62    /// if any. Safe to call from Handler mode / an interrupt, from either
63    /// hart on an SMP build, or from plain task code.
64    pub fn signal(&self) {
65        // Store the latch *before* consuming the waiter slot: a `wait()`
66        // that re-checks `try_take()` immediately after registering (see
67        // `Wait::poll`) must see `fired = true` if this call's swap below
68        // is about to (or just did) claim its registration — Release here
69        // pairs with the Acquire in `try_take`.
70        self.fired
71            .store(true, crate::sync::atomic::Ordering::Release);
72        let w = self
73            .waiter
74            .swap(NO_WAITER, crate::sync::atomic::Ordering::AcqRel);
75        if w != NO_WAITER {
76            waker::wake_task(crate::task::TaskId::from_u16(w as u16));
77        }
78    }
79
80    /// Non-blocking consume: clears and returns whether the signal had
81    /// fired. Usable from a preemptive task, or from boot code before the
82    /// scheduler starts.
83    pub fn try_take(&self) -> bool {
84        self.fired
85            .compare_exchange(
86                true,
87                false,
88                crate::sync::atomic::Ordering::AcqRel,
89                crate::sync::atomic::Ordering::Acquire,
90            )
91            .is_ok()
92    }
93
94    /// Clear a stale latch. Drivers **must** call this before arming
95    /// hardware for a new transaction — otherwise a signal left over from
96    /// a previous, already-completed (or cancelled) transaction would
97    /// make the next `wait()` return immediately with nothing to report.
98    pub fn reset(&self) {
99        self.fired
100            .store(false, crate::sync::atomic::Ordering::Release);
101    }
102
103    /// Cooperative-tier wait: yields the task until [`Signal::signal`] is
104    /// called (or returns immediately if it already latched).
105    ///
106    /// # Panics
107    /// Panics if polled outside of a task context (i.e. not from within
108    /// the executor's poll of a `#[rivet::task]`).
109    pub fn wait(&self) -> Wait<'_> {
110        Wait {
111            sig: self,
112            registered: None,
113        }
114    }
115}
116
117impl Default for Signal {
118    fn default() -> Self {
119        Self::new()
120    }
121}
122
123/// Future returned by [`Signal::wait`].
124pub struct Wait<'a> {
125    sig: &'a Signal,
126    /// `Some(id)` while registered as the waiter; cleared on completion,
127    /// and cancelled in [`Drop`] so a dropped `wait()` never leaves a
128    /// stale registration behind. Does **not** clear `sig.fired` — a
129    /// signal that fired while this future was being cancelled must stay
130    /// observable to whatever calls `try_take()` next.
131    registered: Option<crate::task::TaskId>,
132}
133
134impl<'a> Drop for Wait<'a> {
135    fn drop(&mut self) {
136        if self.registered.take().is_some() {
137            self.sig
138                .waiter
139                .store(NO_WAITER, crate::sync::atomic::Ordering::Release);
140        }
141    }
142}
143
144impl<'a> Future for Wait<'a> {
145    type Output = ();
146
147    fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<()> {
148        // SAFETY: `Wait` holds only a `&Signal` and an `Option`; no
149        // `!Unpin` fields, so projecting is sound.
150        let this = unsafe { self.get_unchecked_mut() };
151        if this.sig.try_take() {
152            return Poll::Ready(());
153        }
154
155        let id = crate::executor::current_task()
156            .expect("Signal::wait() polled outside of a task context");
157        let prev = this
158            .sig
159            .waiter
160            .swap(id.as_u16() as u32, crate::sync::atomic::Ordering::AcqRel);
161        debug_assert_eq!(
162            prev, NO_WAITER,
163            "rivet: two tasks awaiting the same Signal concurrently — \
164             a peripheral has exactly one owning task by construction"
165        );
166        this.registered = Some(id);
167
168        // Re-check: signal() may have fired between our first try_take()
169        // and registering above.
170        if this.sig.try_take() {
171            if this.registered.take().is_some() {
172                this.sig
173                    .waiter
174                    .store(NO_WAITER, crate::sync::atomic::Ordering::Release);
175            }
176            return Poll::Ready(());
177        }
178
179        Poll::Pending
180    }
181}
182
183// Safety: Signal uses only atomics; safe to share across contexts.
184unsafe impl Sync for Signal {}
185
186#[cfg(test)]
187mod tests {
188    use super::*;
189
190    #[test]
191    fn try_take_false_when_unfired() {
192        crate::kernel_test! {
193            let sig = Signal::new();
194            assert!(!sig.try_take());
195        }
196    }
197
198    #[test]
199    fn signal_then_try_take() {
200        crate::kernel_test! {
201            let sig = Signal::new();
202            sig.signal();
203            assert!(sig.try_take());
204            assert!(!sig.try_take(), "try_take consumes the latch");
205        }
206    }
207
208    #[test]
209    fn reset_clears_stale_latch() {
210        crate::kernel_test! {
211            let sig = Signal::new();
212            sig.signal();
213            sig.reset();
214            assert!(!sig.try_take());
215        }
216    }
217
218    #[test]
219    fn wait_ready_when_already_fired() {
220        crate::kernel_test! {
221            let sig = Signal::new();
222            sig.signal();
223            let waker = crate::waker::task_waker(crate::task::TaskId::new(0, 0));
224            let mut cx = Context::from_waker(&waker);
225            let mut fut = sig.wait();
226            // SAFETY: `fut` is a local `Wait` future; `Unpin`, never moved
227            // while pinned — sound for this single poll.
228            let pinned = unsafe { Pin::new_unchecked(&mut fut) };
229            assert_eq!(pinned.poll(&mut cx), Poll::Ready(()));
230        }
231    }
232
233    #[test]
234    #[should_panic(expected = "outside of a task context")]
235    fn wait_panics_without_task_context() {
236        crate::kernel_test! {
237            let sig = Signal::new();
238            let waker = crate::waker::task_waker(crate::task::TaskId::new(0, 0));
239            let mut cx = Context::from_waker(&waker);
240            let mut fut = sig.wait();
241            // SAFETY: `fut` is a local `Wait` future; `Unpin`, never moved
242            // while pinned — sound for this single poll.
243            let pinned = unsafe { Pin::new_unchecked(&mut fut) };
244            let _ = pinned.poll(&mut cx);
245        }
246    }
247
248    #[test]
249    fn signal_after_registration_wakes() {
250        crate::kernel_test! {
251            let sig = Signal::new();
252            let id = crate::task::TaskId::new(1, 0);
253            let waker = crate::waker::task_waker(id);
254            let mut cx = Context::from_waker(&waker);
255            let mut fut = sig.wait();
256            // SAFETY: see above.
257            let pinned = unsafe { Pin::new_unchecked(&mut fut) };
258            // First poll: nothing fired yet, registers as waiter.
259            crate::executor::set_current_for_test(id.priority(), id.index());
260            assert_eq!(pinned.poll(&mut cx), Poll::Pending);
261
262            // ISR fires the signal — should mark the task ready.
263            sig.signal();
264            assert_eq!(crate::waker::next_ready(), Some(id));
265        }
266    }
267
268    #[test]
269    fn drop_clears_registration_not_latch() {
270        crate::kernel_test! {
271            let sig = Signal::new();
272            let id = crate::task::TaskId::new(2, 0);
273            {
274                let mut fut = sig.wait();
275                // SAFETY: see above.
276                let pinned = unsafe { Pin::new_unchecked(&mut fut) };
277                let waker = crate::waker::task_waker(id);
278                let mut cx = Context::from_waker(&waker);
279                crate::executor::set_current_for_test(id.priority(), id.index());
280                assert_eq!(pinned.poll(&mut cx), Poll::Pending);
281                // fut dropped here — registration must be cleared.
282            }
283            // A signal firing after the drop must not wake the old id.
284            sig.signal();
285            assert_eq!(crate::waker::next_ready(), None);
286        }
287    }
288
289    #[test]
290    fn signal_during_cancellation_stays_observable() {
291        crate::kernel_test! {
292            let sig = Signal::new();
293            let id = crate::task::TaskId::new(3, 0);
294            {
295                let mut fut = sig.wait();
296                // SAFETY: see above.
297                let pinned = unsafe { Pin::new_unchecked(&mut fut) };
298                let waker = crate::waker::task_waker(id);
299                let mut cx = Context::from_waker(&waker);
300                crate::executor::set_current_for_test(id.priority(), id.index());
301                assert_eq!(pinned.poll(&mut cx), Poll::Pending);
302                sig.signal();
303                // fut dropped here without ever observing the Ready value.
304            }
305            // The latch must still be observable by whatever asks next.
306            assert!(sig.try_take());
307        }
308    }
309}