moirai_utils/result_cell.rs
1//! One-shot completion cell: one result, one waiter, one hand-off.
2//!
3//! `ResultCell` carries a producer's single output to a single consumer and
4//! wakes the parked waiter across that hand-off. It is the completion path
5//! shared by `moirai-core`'s blocking `TaskResultSlot` and `moirai-async`'s
6//! `AsyncResultSlot`: both need this atomic hand-off, and neither may reintroduce
7//! a lock or a per-task waker map on the result path.
8//!
9//! # Roles
10//!
11//! The cell is reached through an `Arc` shared by exactly two owners:
12//!
13//! - the **producer** calls [`complete`](ResultCell::complete) exactly once — the tail
14//! of a spawned unit of work;
15//! - the **consumer** calls [`try_take_ready`](ResultCell::try_take_ready) and
16//! [`register`](ResultCell::register) from its own poll or park loop.
17//!
18//! Because those are the only two owners, the `result` and `waiter` cells need
19//! no lock. The type does not enforce that, so [`register`](ResultCell::register)
20//! is `unsafe` and its caller vouches for the single consumer. The producer runs
21//! once; the consumer is serialized with itself (`poll` takes `Pin<&mut Self>`
22//! on the async side, and the blocking side registers once per wait); and `Drop`
23//! runs only after the last `Arc`, so it has exclusive access and races neither
24//! side.
25//!
26//! # State machine
27//!
28//! `state: AtomicU8` is the sole synchronization variable; the `result` and
29//! `waiter` cells are touched only while a transition grants exclusive access to
30//! them (C = consumer, P = producer):
31//!
32//! ```text
33//! PENDING ──C register──▶ WAITING ──C re-register──▶ UPDATING_WAITER ──C──▶ WAITING
34//! │ │
35//! │ P complete │ P complete
36//! ▼ ▼
37//! WRITING ────────────────▶ WRITING ──P──▶ READY ──C take──▶ TAKEN
38//! ```
39//!
40//! `WRITING` is the producer's exclusive claim on `result`, and `READY`
41//! publishes it; `UPDATING_WAITER` is the consumer's exclusive claim on `waiter`
42//! while it swaps a stale one, past which the producer spins. A waiter that is
43//! never replaced (see [`Waiter::REPLACE_ON_REPEAT`]) never enters
44//! `UPDATING_WAITER`, and monomorphization removes that arm entirely.
45//!
46//! # Cell-access invariants (what the per-site `// Safety:` comments rely on)
47//!
48//! 1. **`result`: written once, read once.** Only the producer writes it, only
49//! under `WRITING` (entered by winning the producer transition, so no
50//! consumer can see it yet). It is read exactly once — by the unique
51//! `READY -> TAKEN` consumer transition, or by `Drop` at `READY` when the
52//! consumer never took it.
53//! 2. **`waiter`: written by the consumer, read once by the producer.** The
54//! consumer writes it under `PENDING` (published by the `PENDING -> WAITING`
55//! release transition) or under the `UPDATING_WAITER` claim. The producer
56//! reads it exactly once, on `WAITING -> WRITING`; otherwise `Drop` at
57//! `WAITING` drops it. A *failed* publish transition proves no producer
58//! observed the write, so the consumer drops its own value.
59//! 3. **No lost wakeup.** [`complete`](ResultCell::complete) on the `PENDING ->
60//! WRITING` path (it beat registration) deliberately does not wake: no waiter
61//! is registered yet. Liveness therefore requires the consumer to check →
62//! [`register`](ResultCell::register) → **re-check**; should `complete` land in that
63//! window, the re-check observes `READY`. That re-check is load-bearing, not
64//! defensive — dropping it reintroduces a hang.
65//!
66//! # Ordering
67//!
68//! Each cell access is ordered by a release/acquire pair on `state`: the writer
69//! releases on the publishing transition, the reader acquires on the transition
70//! that reads the cell. Thus `PENDING -> WAITING` and the `UPDATING_WAITER ->
71//! WAITING` store are `Release` (they publish `waiter`), while `WAITING ->
72//! WRITING` and `READY -> TAKEN` are `Acquire` (they read a cell). The `PENDING
73//! -> WRITING` success is `Relaxed`: that path reads neither cell before its own
74//! `store(READY, Release)` publishes `result`, so it carries no incoming edge to
75//! establish.
76//!
77//! # The waiter payload
78//!
79//! [`Waiter`] abstracts the parked handle: `thread::Thread` for the blocking
80//! side, `Waker` for the async one. The two differ in exactly one behaviour, and
81//! it is a trait constant rather than a branch —
82//! [`REPLACE_ON_REPEAT`](Waiter::REPLACE_ON_REPEAT) is `false` for a thread,
83//! which parks once per wait and must not be overwritten, and `true` for a
84//! waker, which a re-poll may legitimately replace. Each instantiation therefore
85//! compiles to its own machine, and both share this one copy of the protocol,
86//! its invariants and its ordering argument.
87//!
88//! # Layout
89//!
90//! The cell is packed by default: the state word sits beside the two cells with
91//! no padding, which is what a per-task allocation wants. A caller that wants the
92//! state in an interference sector of its own supplies that as the
93//! [`StateWord`] — `moirai-core`'s `TaskResultSlot` passes
94//! `CacheAligned<AtomicU8>` — and the machine runs identically either way, since
95//! the alignment never reaches the protocol.
96
97use core::cell::UnsafeCell;
98use core::mem::MaybeUninit;
99use core::sync::atomic::{AtomicU8, Ordering};
100
101mod state_word;
102
103pub use state_word::StateWord;
104
105const RESULT_PENDING: u8 = 0;
106const RESULT_WAITING: u8 = 1;
107const RESULT_UPDATING_WAITER: u8 = 2;
108const RESULT_WRITING: u8 = 3;
109const RESULT_READY: u8 = 4;
110const RESULT_TAKEN: u8 = 5;
111
112/// A handle that can be parked on and woken once.
113///
114/// See the [module docs](self#the-waiter-payload) for the one behaviour the two
115/// implementations differ in.
116pub trait Waiter: Send + Clone + Sized {
117 /// Whether registering again while one is already parked replaces it.
118 ///
119 /// `false` for a thread handle: the blocking side registers once per wait, so
120 /// a second registration is a no-op rather than an overwrite. `true` for a
121 /// waker: a re-poll may carry a different waker and the newest must win.
122 const REPLACE_ON_REPEAT: bool;
123
124 /// Wake the parked waiter.
125 fn wake(self);
126}
127
128/// A thread parks once per wait, so a repeat registration is a no-op.
129#[cfg(feature = "std")]
130impl Waiter for std::thread::Thread {
131 const REPLACE_ON_REPEAT: bool = false;
132
133 fn wake(self) {
134 std::thread::Thread::unpark(&self);
135 }
136}
137
138/// A waker may be replaced between polls, so the newest registration wins.
139impl Waiter for core::task::Waker {
140 const REPLACE_ON_REPEAT: bool = true;
141
142 fn wake(self) {
143 core::task::Waker::wake(self);
144 }
145}
146
147/// Single-producer, single-consumer completion cell for one unit of work.
148///
149/// The state word's storage is a parameter ([`StateWord`]), so each caller keeps
150/// the layout it needs; see [the layout note](self#layout).
151pub struct ResultCell<T, W: Waiter, S: StateWord = AtomicU8> {
152 result: UnsafeCell<MaybeUninit<T>>,
153 state: S,
154 waiter: UnsafeCell<MaybeUninit<W>>,
155}
156
157// Safety: the cell has one producer and one consumer. Atomic states serialize
158// result publication, result consumption, and waiter updates, so concurrent
159// `&self` use is sound exactly when the payload and the waiter may each move
160// between threads.
161unsafe impl<T: Send, W: Waiter, S: StateWord> Send for ResultCell<T, W, S> {}
162
163// Safety: shared access is mediated by the state machine; the result and waiter
164// cells are touched only after the corresponding atomic transition succeeds.
165unsafe impl<T: Send, W: Waiter, S: StateWord> Sync for ResultCell<T, W, S> {}
166
167impl<T, W: Waiter, S: StateWord> ResultCell<T, W, S> {
168 /// Create an empty cell.
169 #[must_use]
170 pub fn new() -> Self {
171 Self {
172 result: UnsafeCell::new(MaybeUninit::uninit()),
173 state: S::pending(),
174 waiter: UnsafeCell::new(MaybeUninit::uninit()),
175 }
176 }
177
178 /// Publish the result, waking the parked waiter if one is registered.
179 ///
180 /// Called exactly once by the producer; a second call is a no-op.
181 pub fn complete(&self, result: T) {
182 let Some(waiting) = self.begin_completion() else {
183 return;
184 };
185
186 // Safety: WRITING is reachable only after `begin_completion` wins the
187 // producer transition, so no consumer can read this cell yet.
188 unsafe {
189 (*self.result.get()).write(result);
190 }
191
192 self.state.store(RESULT_READY, Ordering::Release);
193
194 if waiting {
195 // Safety: WAITING is reachable only after `register` writes the
196 // waiter and publishes it with a release transition.
197 let waiter = unsafe { (*self.waiter.get()).assume_init_read() };
198 waiter.wake();
199 }
200 }
201
202 /// Take the result if the producer has published it.
203 pub fn try_take_ready(&self) -> Option<T> {
204 if self
205 .state
206 .compare_exchange(
207 RESULT_READY,
208 RESULT_TAKEN,
209 Ordering::Acquire,
210 Ordering::Relaxed,
211 )
212 .is_ok()
213 {
214 // Safety: READY is published only after the producer initializes the
215 // result cell; READY -> TAKEN is a unique consumer transition.
216 Some(unsafe { (*self.result.get()).assume_init_read() })
217 } else {
218 None
219 }
220 }
221
222 /// Take the result only if a relaxed load already observed it, for a spin
223 /// loop that re-checks: the load is the hint, [`try_take_ready`](ResultCell::try_take_ready)
224 /// is the decision.
225 pub fn try_take_observed_ready(&self) -> Option<T> {
226 if self.state.load(Ordering::Relaxed) == RESULT_READY {
227 self.try_take_ready()
228 } else {
229 None
230 }
231 }
232
233 /// Whether the producer has published a result.
234 #[must_use]
235 pub fn is_completed(&self) -> bool {
236 self.state.load(Ordering::Acquire) == RESULT_READY
237 }
238
239 /// Whether a waiter is parked and no result has been published yet.
240 ///
241 /// The complement of [`is_completed`](ResultCell::is_completed) would also
242 /// answer true before any registration, so this reports the waiting state
243 /// itself: it is what distinguishes a registration that took effect from
244 /// one that never ran.
245 #[must_use]
246 pub fn has_registered_waiter(&self) -> bool {
247 self.state.load(Ordering::Acquire) == RESULT_WAITING
248 }
249
250 /// Park a clone of `waiter`, or replace the parked one when
251 /// [`REPLACE_ON_REPEAT`](Waiter::REPLACE_ON_REPEAT) is set.
252 ///
253 /// Registration retries until the state is settled, so the waiter is stored
254 /// from a clone and a failed publish transition costs only that clone. The
255 /// caller must re-check [`try_take_ready`](ResultCell::try_take_ready) after this
256 /// returns: a `complete` that raced in before the registration does not wake,
257 /// because it saw no waiter. See the
258 /// [module docs](self#cell-access-invariants).
259 ///
260 /// # Safety
261 ///
262 /// The waiter cell has one writer at a time, so no other call to `register`
263 /// may run on this cell concurrently, and `waiter.clone()` must not reach
264 /// back into this cell. The cell has one consumer by construction (a blocking
265 /// join takes the handle by value; an async poll takes `Pin<&mut Self>`), and
266 /// the caller is that consumer.
267 pub unsafe fn register(&self, waiter: &W) {
268 loop {
269 match self.state.load(Ordering::Acquire) {
270 RESULT_PENDING => {
271 // Safety: the caller of `register` is the only consumer
272 // registering. If the publish transition fails, this clone is
273 // dropped before retry.
274 unsafe {
275 (*self.waiter.get()).write(waiter.clone());
276 }
277
278 if self
279 .state
280 .compare_exchange(
281 RESULT_PENDING,
282 RESULT_WAITING,
283 Ordering::Release,
284 Ordering::Acquire,
285 )
286 .is_ok()
287 {
288 return;
289 }
290
291 // Safety: the transition failed, so no producer can observe
292 // this waiter cell as initialized through the WAITING state.
293 unsafe {
294 (*self.waiter.get()).assume_init_drop();
295 }
296 }
297 RESULT_WAITING => {
298 if !W::REPLACE_ON_REPEAT {
299 return;
300 }
301
302 if self
303 .state
304 .compare_exchange(
305 RESULT_WAITING,
306 RESULT_UPDATING_WAITER,
307 Ordering::Acquire,
308 Ordering::Acquire,
309 )
310 .is_ok()
311 {
312 // Safety: UPDATING_WAITER excludes the producer from
313 // reading the waiter cell while the consumer replaces it.
314 unsafe {
315 (*self.waiter.get()).assume_init_drop();
316 (*self.waiter.get()).write(waiter.clone());
317 }
318 self.state.store(RESULT_WAITING, Ordering::Release);
319 return;
320 }
321 }
322 RESULT_UPDATING_WAITER | RESULT_WRITING => core::hint::spin_loop(),
323 _ => return,
324 }
325 }
326 }
327
328 /// Claim the producer side of the hand-off.
329 ///
330 /// Returns `Some(waiting)` with the result cell claimed for writing, or
331 /// `None` when the result was already published (or taken) and the caller
332 /// must abandon its value.
333 fn begin_completion(&self) -> Option<bool> {
334 loop {
335 match self.state.load(Ordering::Acquire) {
336 RESULT_PENDING => {
337 if self
338 .state
339 .compare_exchange(
340 RESULT_PENDING,
341 RESULT_WRITING,
342 Ordering::Relaxed,
343 Ordering::Acquire,
344 )
345 .is_ok()
346 {
347 return Some(false);
348 }
349 }
350 RESULT_WAITING => {
351 if self
352 .state
353 .compare_exchange(
354 RESULT_WAITING,
355 RESULT_WRITING,
356 Ordering::Acquire,
357 Ordering::Acquire,
358 )
359 .is_ok()
360 {
361 return Some(true);
362 }
363 }
364 RESULT_UPDATING_WAITER => core::hint::spin_loop(),
365 _ => return None,
366 }
367 }
368 }
369}
370
371impl<T, W: Waiter, S: StateWord> Default for ResultCell<T, W, S> {
372 fn default() -> Self {
373 Self::new()
374 }
375}
376
377impl<T, W: Waiter, S: StateWord> Drop for ResultCell<T, W, S> {
378 fn drop(&mut self) {
379 match *self.state.get_mut() {
380 RESULT_READY => {
381 // Safety: READY means the result cell is initialized and no
382 // consumer took it, because drop has exclusive access.
383 unsafe {
384 self.result.get_mut().assume_init_drop();
385 }
386 }
387 RESULT_WAITING => {
388 // Safety: WAITING means the waiter cell is initialized and no
389 // producer read it, because drop has exclusive access.
390 unsafe {
391 self.waiter.get_mut().assume_init_drop();
392 }
393 }
394 _ => {}
395 }
396 }
397}
398
399#[cfg(all(test, feature = "std"))]
400mod tests;