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}