Skip to main content

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;