Skip to main content

frust_reactive/
task.rs

1//! The blessed heavy-work idiom: [`AsyncValue<T>`] + [`use_task`].
2//!
3//! This is Frust's direct counterpart to Flutter's
4//! `compute()`/`FutureBuilder` — a one-call way to run heavy work off the UI
5//! thread and get exhaustive load/error/ready states back, with cancellation
6//! semantics Flutter's model doesn't offer.
7//!
8//! # The idiom
9//!
10//! ```ignore
11//! // inside `Component::init`/`build`, under the component's `Owner`:
12//! let data = use_task(|| async { spawn_blocking(parse).await });
13//! // `data.signal()` is an `RwSignal<AsyncValue<T>>` read in `build`;
14//! // `data.restart()` re-runs the fetch.
15//! ```
16//!
17//! # Threading and cancellation contract
18//!
19//! `use_task` splits the work into two halves, wired explicitly (never
20//! assuming any implicit cancellation — the leptos precedent shows
21//! "implicit cancellation" claims are usually wrong; see the
22//! research ledger §9):
23//!
24//! 1. **A background half** — the fetcher future is handed to the process-wide
25//!    tokio runtime via [`ReactiveRuntime`]'s handle. Heavy work inside it
26//!    hops threads through [`crate::spawn_blocking`] (one-off CPU work) or
27//!    `frust::spawn` (async IO). The tokio [`JoinHandle`](tokio::task::JoinHandle)'s
28//!    [`AbortHandle`](tokio::task::AbortHandle) is registered in `on_cleanup`,
29//!    so owner teardown aborts the background task (best-effort: a
30//!    `spawn_blocking` closure *already running* cannot be interrupted — a
31//!    documented limitation shared by every runtime).
32//! 2. **A UI-side coordinator** — a `!Send` future spawned via
33//!    reactive_graph's [`spawn_local_scoped_with_cancellation`] so it aborts
34//!    on owner cleanup. It `await`s the background [`JoinHandle`](tokio::task::JoinHandle)
35//!    and only then writes the result signal.
36//!
37//! Because **every signal write happens on the UI thread** (the coordinator
38//! awaits the background result, then sets), the same-frame cross-thread
39//! write/read race the huddle review deferred as A10 is *structurally
40//! impossible* for idiom users: there is no background-thread `signal.set`,
41//! and a coordinator aborted by owner cleanup never writes to a disposed
42//! signal (closing A9). The `use_task` stress tests below are the audit A10
43//! asked for, executed against the idiom rather than by inspection.
44
45use std::cell::Cell;
46use std::future::Future;
47use std::rc::Rc;
48use std::sync::{Arc, Mutex};
49
50use reactive_graph::owner::{Owner, on_cleanup};
51use reactive_graph::signal::RwSignal;
52use reactive_graph::spawn_local_scoped_with_cancellation;
53use reactive_graph::traits::{Set, Update};
54use tokio::task::AbortHandle;
55
56use crate::runtime::ReactiveRuntime;
57
58/// The boxed error an [`AsyncValue::Error`] carries. `Arc`-wrapped so a
59/// clone of the state is cheap and the error is shareable across the tree.
60pub type TaskError = Arc<dyn std::error::Error + Send + Sync>;
61
62/// The exhaustive state of an asynchronously-loaded value.
63///
64/// Deliberately named for parity with Riverpod's `AsyncValue<T>` (Flutter) —
65/// the genuine prior art for a load/data/error sum type with exhaustive
66/// matching (research §9). The four states are:
67///
68/// - [`Idle`](Self::Idle) — nothing requested yet.
69/// - [`Loading`](Self::Loading) — a fetch is in flight. It carries the
70///   *previous* value (`Some` on a refresh, `None` on a first load), so a UI
71///   can keep showing stale data instead of flickering to a spinner —
72///   mirroring Riverpod's `copyWithPrevious`.
73/// - [`Ready`](Self::Ready) — the fetch resolved to a value.
74/// - [`Error`](Self::Error) — the fetch failed.
75#[derive(Default)]
76pub enum AsyncValue<T> {
77    /// No fetch requested yet.
78    #[default]
79    Idle,
80    /// A fetch is in flight; carries the previous value for
81    /// refresh-without-flicker (`None` on a first load).
82    Loading(Option<T>),
83    /// The fetch resolved.
84    Ready(T),
85    /// The fetch failed.
86    Error(TaskError),
87}
88
89impl<T> AsyncValue<T> {
90    /// Whether this is [`Idle`](Self::Idle).
91    pub fn is_idle(&self) -> bool {
92        matches!(self, AsyncValue::Idle)
93    }
94
95    /// Whether a fetch is in flight.
96    pub fn is_loading(&self) -> bool {
97        matches!(self, AsyncValue::Loading(_))
98    }
99
100    /// Whether the fetch resolved to a value.
101    pub fn is_ready(&self) -> bool {
102        matches!(self, AsyncValue::Ready(_))
103    }
104
105    /// Whether the fetch failed.
106    pub fn is_error(&self) -> bool {
107        matches!(self, AsyncValue::Error(_))
108    }
109
110    /// The resolved value, if [`Ready`](Self::Ready).
111    pub fn ready(&self) -> Option<&T> {
112        match self {
113            AsyncValue::Ready(t) => Some(t),
114            _ => None,
115        }
116    }
117
118    /// The best-available value: the resolved one when [`Ready`](Self::Ready),
119    /// or the carried-over previous one while [`Loading`](Self::Loading) a
120    /// refresh. This is what a flicker-free UI reads.
121    pub fn value(&self) -> Option<&T> {
122        match self {
123            AsyncValue::Ready(t) => Some(t),
124            AsyncValue::Loading(prev) => prev.as_ref(),
125            _ => None,
126        }
127    }
128
129    /// The error, if [`Error`](Self::Error).
130    pub fn error(&self) -> Option<&TaskError> {
131        match self {
132            AsyncValue::Error(e) => Some(e),
133            _ => None,
134        }
135    }
136}
137
138impl<T: Clone> Clone for AsyncValue<T> {
139    fn clone(&self) -> Self {
140        match self {
141            AsyncValue::Idle => AsyncValue::Idle,
142            AsyncValue::Loading(prev) => AsyncValue::Loading(prev.clone()),
143            AsyncValue::Ready(t) => AsyncValue::Ready(t.clone()),
144            AsyncValue::Error(e) => AsyncValue::Error(e.clone()),
145        }
146    }
147}
148
149impl<T: std::fmt::Debug> std::fmt::Debug for AsyncValue<T> {
150    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
151        match self {
152            AsyncValue::Idle => f.write_str("Idle"),
153            AsyncValue::Loading(prev) => f.debug_tuple("Loading").field(prev).finish(),
154            AsyncValue::Ready(t) => f.debug_tuple("Ready").field(t).finish(),
155            AsyncValue::Error(e) => f.debug_tuple("Error").field(e).finish(),
156        }
157    }
158}
159
160/// The handle [`use_task`] returns: a read handle to the task's
161/// [`AsyncValue<T>`] state plus a `restart`/refresh trigger.
162///
163/// Read the state reactively via [`signal`](Self::signal) (or the [`Get`]
164/// trait on it) from `Component::build`; call [`restart`](Self::restart) to
165/// re-run the fetch (e.g. a pull-to-refresh).
166///
167/// [`Get`]: reactive_graph::traits::Get
168pub struct UseTask<T: Send + Sync + 'static> {
169    signal: RwSignal<AsyncValue<T>>,
170    // `Rc<dyn Fn()>` — UI-thread-only; re-runs the fetch under the captured
171    // owner so cancellation stays wired even when `restart` fires from an
172    // arbitrary (owner-less) event-handler context.
173    run: Rc<dyn Fn()>,
174}
175
176impl<T: Send + Sync + 'static> UseTask<T> {
177    /// The reactive state signal. Read it with the [`Get`] trait
178    /// (`task.signal().get()`) inside a tracked `Component::build` so a state
179    /// transition wakes the shell.
180    ///
181    /// [`Get`]: reactive_graph::traits::Get
182    pub fn signal(&self) -> RwSignal<AsyncValue<T>> {
183        self.signal
184    }
185
186    /// Re-runs the fetch: transitions the state to
187    /// [`Loading`](AsyncValue::Loading) (carrying the current value for
188    /// refresh-without-flicker), aborts any in-flight background task, and
189    /// starts a fresh one. Last-write-wins by generation, so a restart storm
190    /// settles on the newest fetch's result regardless of completion order.
191    pub fn restart(&self) {
192        (self.run)();
193    }
194}
195
196impl<T: Send + Sync + 'static> Clone for UseTask<T> {
197    fn clone(&self) -> Self {
198        UseTask {
199            signal: self.signal,
200            run: self.run.clone(),
201        }
202    }
203}
204
205/// Shared, UI-thread-owned coordination state for one [`use_task`] instance.
206struct Coordinator<T: Send + Sync + 'static> {
207    signal: RwSignal<AsyncValue<T>>,
208    /// Bumped on every run; the coordinator only writes its result if its
209    /// captured generation still matches (last-write-wins under a restart
210    /// storm). UI-thread-only, so a plain [`Cell`] suffices.
211    generation: Cell<u64>,
212    /// The current background task's abort handle, shared with the single
213    /// `on_cleanup` registration so owner teardown aborts it. `Arc<Mutex<_>>`
214    /// because `on_cleanup` requires a `Send + Sync` closure — the tokio
215    /// [`AbortHandle`] is itself `Send + Sync`.
216    bg_abort: Arc<Mutex<Option<AbortHandle>>>,
217}
218
219/// Runs a fetch, wiring a heavy-work idiom around it.
220///
221/// Called from `Component::init`/`build` under the component's [`Owner`]. It
222/// immediately starts a first fetch and returns a [`UseTask<T>`] to read the
223/// [`AsyncValue<T>`] state and to `restart` it.
224///
225/// - `T` is the loaded value type (`Send + Sync` — it moves from a background
226///   thread to the UI thread, and lives in a thread-safe signal).
227/// - `E` is any [`std::error::Error`] the fetch may fail with (tokio's
228///   `JoinError` qualifies, so `|| async { spawn_blocking(f).await }` works
229///   directly).
230///
231/// See the [module docs](self) for the full threading/cancellation contract.
232///
233/// # Decision: hand-rolled vs `AsyncDerived`
234///
235/// reactive_graph 0.2 ships `AsyncDerived` (research §9), which this could
236/// wrap for a *signal-driven* restart. It is deliberately **not** used here:
237/// `AsyncDerived` re-runs when a tracked signal it reads changes, whereas
238/// `use_task`'s contract is an *imperative* first-load + explicit `restart`
239/// (the pull-to-refresh / retry shape), and its cancellation story is the
240/// leptos one the research refuted as non-explicit. Hand-rolling keeps all
241/// three guarantees visible in one place — background `AbortHandle` in
242/// `on_cleanup`, coordinator abort via `spawn_local_scoped_with_cancellation`,
243/// and last-write-wins by generation. A future `AsyncDerived`-backed
244/// signal-driven variant can live alongside this without changing it (a
245/// documented Future Enhancement in the plan).
246pub fn use_task<T, E, Fut, F>(fetch: F) -> UseTask<T>
247where
248    T: Send + Sync + 'static,
249    E: std::error::Error + Send + Sync + 'static,
250    Fut: Future<Output = Result<T, E>> + Send + 'static,
251    F: Fn() -> Fut + 'static,
252{
253    let coord = Rc::new(Coordinator {
254        signal: RwSignal::new(AsyncValue::Idle),
255        generation: Cell::new(0),
256        bg_abort: Arc::new(Mutex::new(None)),
257    });
258
259    // Register a single owner-cleanup that aborts whatever background task is
260    // current at teardown time. Capturing only the `Send + Sync` abort slot
261    // (not the `!Send` `Rc<Coordinator>`) keeps the closure within
262    // `on_cleanup`'s bound.
263    {
264        let bg_abort = coord.bg_abort.clone();
265        on_cleanup(move || {
266            if let Some(handle) = bg_abort.lock().expect("bg_abort poisoned").take() {
267                handle.abort();
268            }
269        });
270    }
271
272    let signal = coord.signal;
273    let fetch = Rc::new(fetch);
274
275    // Capture the owner so every run (initial + restart) re-enters it: the
276    // coordinator's `spawn_local_scoped_with_cancellation` and the background
277    // `on_cleanup` both bind to *this* owner, even if `restart` is called from
278    // an owner-less context (an event handler).
279    let owner = Owner::current();
280    let run: Rc<dyn Fn()> = {
281        let coord = coord.clone();
282        let fetch = fetch.clone();
283        Rc::new(move || {
284            let go = || run_once(&coord, &fetch);
285            match &owner {
286                Some(owner) => owner.with(go),
287                None => go(),
288            }
289        })
290    };
291
292    // Kick off the first load.
293    run();
294
295    UseTask { signal, run }
296}
297
298/// One fetch cycle: supersede any prior run, flip to `Loading`, spawn the
299/// background work, and spawn the UI-side coordinator that writes the result.
300fn run_once<T, E, Fut, F>(coord: &Rc<Coordinator<T>>, fetch: &Rc<F>)
301where
302    T: Send + Sync + 'static,
303    E: std::error::Error + Send + Sync + 'static,
304    Fut: Future<Output = Result<T, E>> + Send + 'static,
305    F: Fn() -> Fut + 'static,
306{
307    let rt = ReactiveRuntime::get().expect(
308        "frust-reactive: use_task called before ReactiveRuntime::init — \
309         this is a wiring bug: initialize the reactive runtime (the shell does \
310         this on startup) before mounting components that use use_task",
311    );
312
313    // Abort the previous in-flight background task (restart supersedes it).
314    if let Some(handle) = coord.bg_abort.lock().expect("bg_abort poisoned").take() {
315        handle.abort();
316    }
317
318    // Claim a fresh generation; only this run may write its result.
319    let generation = coord.generation.get().wrapping_add(1);
320    coord.generation.set(generation);
321
322    // Transition to Loading, carrying the current value for a flicker-free
323    // refresh.
324    coord.signal.update(|state| {
325        let prev = match std::mem::take(state) {
326            AsyncValue::Ready(t) => Some(t),
327            AsyncValue::Loading(prev) => prev,
328            _ => None,
329        };
330        *state = AsyncValue::Loading(prev);
331    });
332
333    // Background half: hand the fetcher future to the tokio runtime and
334    // register its abort handle for owner teardown / the next restart.
335    let join = rt.handle().spawn((fetch)());
336    *coord.bg_abort.lock().expect("bg_abort poisoned") = Some(join.abort_handle());
337
338    // UI-side coordinator: await the background result, then write the signal
339    // on the UI thread. Scoped-with-cancellation so owner cleanup aborts it —
340    // a disposed signal is never written.
341    let coord = coord.clone();
342    spawn_local_scoped_with_cancellation(async move {
343        let outcome = join.await;
344
345        // Last-write-wins: a newer run has already claimed the signal.
346        if coord.generation.get() != generation {
347            return;
348        }
349
350        match outcome {
351            Ok(Ok(value)) => coord.signal.set(AsyncValue::Ready(value)),
352            Ok(Err(err)) => {
353                let err: TaskError = Arc::new(err);
354                coord.signal.set(AsyncValue::Error(err));
355            }
356            // The background task was aborted (owner teardown or a restart) or
357            // panicked. On abort the coordinator is normally torn down too, so
358            // this arm is a race-safe fallback: leave the state as the caller
359            // (or the newer run) set it rather than clobbering it.
360            Err(_join_err) => {}
361        }
362    });
363}
364
365#[cfg(test)]
366mod tests {
367    use super::*;
368    use reactive_graph::owner::Owner;
369    use reactive_graph::traits::GetUntracked;
370    use std::sync::atomic::{AtomicUsize, Ordering};
371    use std::sync::mpsc;
372    use std::time::{Duration, Instant};
373
374    use crate::ReactiveRuntime;
375
376    fn noop_waker() -> crate::FrameWaker {
377        Arc::new(|| {})
378    }
379
380    /// Ensures a shared, initialized runtime on the calling (UI) thread.
381    fn init_rt() -> &'static ReactiveRuntime {
382        ReactiveRuntime::init(noop_waker())
383    }
384
385    /// Pump the UI-thread local queue until `cond` holds or the deadline
386    /// passes; returns whether `cond` became true. The 1ms yield between pumps
387    /// is not a *synchronization* device — it only lets the background tokio
388    /// task make progress; correctness is asserted on `cond`, not on elapsed
389    /// time.
390    fn pump_until(rt: &ReactiveRuntime, timeout: Duration, mut cond: impl FnMut() -> bool) -> bool {
391        let start = Instant::now();
392        loop {
393            rt.pump_local();
394            if cond() {
395                return true;
396            }
397            if start.elapsed() >= timeout {
398                return false;
399            }
400            std::thread::sleep(Duration::from_millis(1));
401        }
402    }
403
404    /// Pump a fixed handful of turns to drain aborted coordinators; used where
405    /// the assertion is "no panic" rather than a state condition.
406    fn pump_a_few(rt: &ReactiveRuntime) {
407        for _ in 0..4 {
408            rt.pump_local();
409        }
410    }
411
412    #[derive(Debug)]
413    struct TestError(&'static str);
414    impl std::fmt::Display for TestError {
415        fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
416            f.write_str(self.0)
417        }
418    }
419    impl std::error::Error for TestError {}
420
421    /// A `use_task` fetch resolves to `Ready` on the UI thread after the
422    /// background work completes. This is the milestone shape:
423    /// `use_task(|| async { spawn_blocking(f).await })` (fetch error type is
424    /// tokio's `JoinError`).
425    #[test]
426    fn use_task_resolves_to_ready() {
427        let _guard = crate::WAKER_TEST_LOCK
428            .lock()
429            .unwrap_or_else(|e| e.into_inner());
430        let rt = init_rt();
431
432        let owner = Owner::new();
433        let task = owner.with(|| use_task(|| async { crate::spawn_blocking(|| 6 * 7).await }));
434
435        assert!(
436            task.signal().get_untracked().is_loading(),
437            "state must be Loading immediately after use_task"
438        );
439        assert!(
440            pump_until(rt, Duration::from_secs(5), || task
441                .signal()
442                .get_untracked()
443                .is_ready()),
444            "task should reach Ready after pumping"
445        );
446        assert_eq!(task.signal().get_untracked().ready().copied(), Some(42));
447
448        owner.cleanup();
449    }
450
451    /// A failing fetch lands in `Error`, carrying the fetch's own error type.
452    #[test]
453    fn use_task_reports_error() {
454        let _guard = crate::WAKER_TEST_LOCK
455            .lock()
456            .unwrap_or_else(|e| e.into_inner());
457        let rt = init_rt();
458
459        let owner = Owner::new();
460        let task = owner.with(|| {
461            use_task(|| async {
462                crate::spawn_blocking(|| ()).await.expect("join");
463                Result::<i32, TestError>::Err(TestError("boom"))
464            })
465        });
466
467        assert!(
468            pump_until(rt, Duration::from_secs(5), || task
469                .signal()
470                .get_untracked()
471                .is_error()),
472            "task should reach Error after pumping"
473        );
474        assert_eq!(
475            task.signal().get_untracked().error().map(|e| e.to_string()),
476            Some("boom".to_string())
477        );
478
479        owner.cleanup();
480    }
481
482    /// `Loading` carries the previous `Ready` value across a `restart`, so a
483    /// refresh can render stale data instead of flickering to a spinner.
484    #[test]
485    fn restart_carries_previous_value_in_loading() {
486        let _guard = crate::WAKER_TEST_LOCK
487            .lock()
488            .unwrap_or_else(|e| e.into_inner());
489        let rt = init_rt();
490
491        let seq = Arc::new(AtomicUsize::new(0));
492        let owner = Owner::new();
493        let task = {
494            let seq = seq.clone();
495            owner.with(|| {
496                use_task(move || {
497                    let seq = seq.clone();
498                    async move {
499                        crate::spawn_blocking(move || seq.fetch_add(1, Ordering::SeqCst)).await
500                    }
501                })
502            })
503        };
504
505        assert!(pump_until(rt, Duration::from_secs(5), || task
506            .signal()
507            .get_untracked()
508            .is_ready()));
509        assert_eq!(task.signal().get_untracked().ready().copied(), Some(0));
510
511        task.restart();
512        // Immediately after restart: Loading, but carrying the previous value.
513        assert_eq!(
514            task.signal().get_untracked().value().copied(),
515            Some(0),
516            "Loading must carry the previous Ready value for flicker-free refresh"
517        );
518        assert!(task.signal().get_untracked().is_loading());
519
520        owner.cleanup();
521    }
522
523    /// A9/A10 core: mount, start a load, unmount *while loading*, then let the
524    /// background task complete — no panic, and no write to the disposed
525    /// signal (the coordinator is aborted on cleanup). ≥1000 iterations,
526    /// headless.
527    #[test]
528    fn unmount_while_loading_never_writes_disposed_signal() {
529        let _guard = crate::WAKER_TEST_LOCK
530            .lock()
531            .unwrap_or_else(|e| e.into_inner());
532        let rt = init_rt();
533
534        for i in 0..1_000u32 {
535            let ran = Arc::new(AtomicUsize::new(0));
536            let owner = Owner::new();
537            let task = {
538                let ran = ran.clone();
539                owner.with(|| {
540                    use_task(move || {
541                        let ran = ran.clone();
542                        async move {
543                            crate::spawn_blocking(move || {
544                                ran.fetch_add(1, Ordering::SeqCst);
545                                i
546                            })
547                            .await
548                        }
549                    })
550                })
551            };
552
553            // Tear the owner down while the fetch is (almost certainly) still
554            // in flight: runs the on_cleanup that aborts the coordinator and
555            // the background task, and disposes the signal.
556            owner.cleanup();
557            drop(task);
558
559            // Pump a few turns; the background task may still complete, but the
560            // aborted coordinator never writes the disposed signal. The only
561            // guarantee under test is "no panic".
562            pump_a_few(rt);
563        }
564    }
565
566    /// A task that only completes *after* teardown must not panic when its
567    /// background work finishes. ≥1000 iterations.
568    #[test]
569    fn task_completing_after_teardown_is_safe() {
570        let _guard = crate::WAKER_TEST_LOCK
571            .lock()
572            .unwrap_or_else(|e| e.into_inner());
573        let rt = init_rt();
574
575        for _ in 0..1_000u32 {
576            let (tx, rx) = mpsc::channel::<()>();
577            let rx = Arc::new(Mutex::new(rx));
578            let owner = Owner::new();
579            let task = {
580                let rx = rx.clone();
581                owner.with(|| {
582                    use_task(move || {
583                        let rx = rx.clone();
584                        async move {
585                            // Block the background task until *after* teardown,
586                            // guaranteeing "completes after teardown".
587                            crate::spawn_blocking(move || {
588                                let _ = rx.lock().expect("rx").recv();
589                                1u32
590                            })
591                            .await
592                        }
593                    })
594                })
595            };
596
597            owner.cleanup();
598            drop(task);
599            // Release the background task only now: it completes post-teardown.
600            let _ = tx.send(());
601            pump_a_few(rt);
602        }
603    }
604
605    /// Restart storm: hammer `restart` many times; the final state settles on
606    /// the newest fetch's value regardless of completion order (the generation
607    /// guard is last-write-wins). ≥1000 restarts.
608    #[test]
609    fn restart_storm_settles_on_latest() {
610        let _guard = crate::WAKER_TEST_LOCK
611            .lock()
612            .unwrap_or_else(|e| e.into_inner());
613        let rt = init_rt();
614
615        let counter = Arc::new(AtomicUsize::new(0));
616        let owner = Owner::new();
617        let task = {
618            let counter = counter.clone();
619            owner.with(|| {
620                use_task(move || {
621                    // Assign the sequence number at fetch-*call* time (run_once
622                    // runs synchronously in generation order), not at poll time
623                    // (background scheduling order is nondeterministic) — so the
624                    // newest generation deterministically owns the highest seq.
625                    let seq = counter.fetch_add(1, Ordering::SeqCst);
626                    async move { crate::spawn_blocking(move || seq).await }
627                })
628            })
629        };
630
631        for _ in 0..1_000u32 {
632            task.restart();
633        }
634        let last_seq = counter.load(Ordering::SeqCst) - 1;
635
636        assert!(
637            pump_until(rt, Duration::from_secs(10), || {
638                matches!(
639                    task.signal().get_untracked().ready().copied(),
640                    Some(seq) if seq == last_seq
641                )
642            }),
643            "restart storm must settle on the newest fetch's value (last-write-wins)"
644        );
645
646        owner.cleanup();
647    }
648
649    /// Cross-thread completion ordering: an *earlier* fetch that finishes
650    /// *later* must not clobber a *newer* fetch that already resolved. The
651    /// ordering is enforced explicitly with a channel, not wall-clock timing.
652    /// ≥1000 iterations.
653    #[test]
654    fn out_of_order_completion_respects_generation() {
655        let _guard = crate::WAKER_TEST_LOCK
656            .lock()
657            .unwrap_or_else(|e| e.into_inner());
658        let rt = init_rt();
659
660        for _ in 0..1_000u32 {
661            let seq = Arc::new(AtomicUsize::new(0));
662            // The first fetch blocks on this receiver until we release it.
663            let (release_tx, release_rx) = mpsc::channel::<()>();
664            let release_rx = Arc::new(Mutex::new(release_rx));
665
666            let owner = Owner::new();
667            let task = {
668                let seq = seq.clone();
669                let release_rx = release_rx.clone();
670                owner.with(|| {
671                    use_task(move || {
672                        let n = seq.fetch_add(1, Ordering::SeqCst);
673                        let release_rx = release_rx.clone();
674                        async move {
675                            crate::spawn_blocking(move || {
676                                if n == 0 {
677                                    // First fetch: block until released, so it
678                                    // completes AFTER the newer one.
679                                    let _ = release_rx.lock().expect("rx").recv();
680                                }
681                                n
682                            })
683                            .await
684                        }
685                    })
686                })
687            };
688
689            // Second fetch supersedes the (blocked) first.
690            task.restart();
691
692            // The newer fetch (seq == 1) resolves first.
693            assert!(
694                pump_until(rt, Duration::from_secs(5), || matches!(
695                    task.signal().get_untracked().ready().copied(),
696                    Some(1)
697                )),
698                "the newer fetch must resolve to 1"
699            );
700
701            // Release the stale first fetch; it completes now but must NOT
702            // overwrite 1 (its generation is stale).
703            let _ = release_tx.send(());
704            pump_a_few(rt);
705            assert_eq!(
706                task.signal().get_untracked().ready().copied(),
707                Some(1),
708                "a stale, later-completing fetch must not clobber the newer result"
709            );
710
711            owner.cleanup();
712        }
713    }
714}