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}