From mboxrd@z Thu Jan 1 00:00:00 1970 Return-Path: Received: from gate001.proxmox.com (gate001.proxmox.com [IPv6:2a0f:8001:1:32::40]) by lore.proxmox.com (Postfix) with ESMTPS id 74DB11FF0C0 for ; Wed, 09 Sep 2026 03:04:13 +0200 (CEST) Received: from gate001.proxmox.com (localhost.localdomain [127.0.0.1]) by gate001.proxmox.com (Proxmox) with ESMTP id 3D1C021492; Wed, 09 Sep 2026 03:04:13 +0200 (CEST) From: Thomas Lamprecht To: yew-devel@lists.proxmox.com Subject: [PATCH 2/2] loader: fix stuck loading state and cancel loads without owners Date: Wed, 9 Sep 2026 03:01:13 +0200 Message-ID: <20260909010400.684482-3-t.lamprecht@proxmox.com> X-Mailer: git-send-email 2.47.3 In-Reply-To: <20260909010400.684482-1-t.lamprecht@proxmox.com> References: <20260909010400.684482-1-t.lamprecht@proxmox.com> MIME-Version: 1.0 Content-Transfer-Encoding: 8bit X-Bm-Milter-Handled: 55990f41-d878-4baa-be0a-ee34c49e34d2 X-Bm-Transport-Timestamp: 1788915837297 X-SPAM-LEVEL: Spam detection results: 0 AWL 0.843 Adjusted score from AWL reputation of From: address DMARC_MISSING 0.1 Missing DMARC policy KAM_DMARC_STATUS 0.01 Test Rule for DKIM or SPF Failure with Strict Alignment (newer systems) RCVD_IN_DNSWL_MED -2.3 Sender listed at https://www.dnswl.org/, medium trust SPF_HELO_NONE 0.001 SPF: HELO does not publish an SPF Record SPF_PASS -0.001 SPF: sender matches SPF record Message-ID-Hash: JTIM5FDUQKBZ7LMPYSZLFTO557YFDTYY X-Message-ID-Hash: JTIM5FDUQKBZ7LMPYSZLFTO557YFDTYY X-MailFrom: t.lamprecht@proxmox.com X-Mailman-Rule-Misses: dmarc-mitigation; no-senders; approved; loop; banned-address; emergency; member-moderation; nonmember-moderation; administrivia; implicit-dest; max-recipients; max-size; news-moderation; no-subject; digests; suspicious-header X-Mailman-Version: 3.3.10 Precedence: list List-Id: Yew framework devel list at Proxmox List-Help: List-Owner: List-Post: List-Subscribe: List-Unsubscribe: The loading counter is decremented only on completion. Canceling a pending load can skip that decrement, leaving the loader busy and its refresh button spinning and disabled even after later loads finish. Track only the current request instead. A callback can still complete after replacing or aborting its own request during its final poll, so cancellation alone cannot prevent stale results or cancellation of a replacement. Pending futures retain Loader handles until completion, preventing cancellation when consumers disappear. Listener adapters also capture Loader handles and can keep an earlier change callback registered after callers replace it and drop their old handles. Use non-owning internal references so the last Loader handle cancels pending work even when observers remain. Keep load starts silent because existing task viewers schedule their next poll from change notifications. Notifying on entry could repeatedly cancel a slow request before it completes. Fixes: f8669c6eecf1 ("Loader: use SharedState") Fixes: 33a8f0f592b8 ("Loader: use AsyncAbortGuard") Fixes: d5fb74e56e1d ("loader: add helper to allow aborting a load") Signed-off-by: Thomas Lamprecht --- src/async_abort_guard.rs | 22 +- src/state/loader.rs | 144 +++++++++---- src/state/loader/tests.rs | 430 ++++++++++++++++++++++++++++++++++++++ 3 files changed, 549 insertions(+), 47 deletions(-) create mode 100644 src/state/loader/tests.rs diff --git a/src/async_abort_guard.rs b/src/async_abort_guard.rs index a5f21d00..4b9bab04 100644 --- a/src/async_abort_guard.rs +++ b/src/async_abort_guard.rs @@ -12,16 +12,22 @@ impl AsyncAbortGuard { pub fn spawn(future: F) -> Self where F: Future + 'static, + { + let (guard, future) = Self::new(future); + wasm_bindgen_futures::spawn_local(future); + guard + } + + /// Wrap a future and return its abort guard without scheduling it. + pub(crate) fn new(future: F) -> (Self, impl Future) + where + F: Future, { let (future, abort_handle) = abortable(future); - - wasm_bindgen_futures::spawn_local(async move { - match future.await { - Ok(()) => { /* do nothing */ } - Err(futures::future::Aborted) => { /* do nothing (maybe we want to log this?) */ } - } - }); - AsyncAbortGuard { abort_handle } + let future = async move { + let _ = future.await; + }; + (Self { abort_handle }, future) } } diff --git a/src/state/loader.rs b/src/state/loader.rs index 4ad461f7..06955c78 100644 --- a/src/state/loader.rs +++ b/src/state/loader.rs @@ -1,3 +1,4 @@ +use std::future::Future; use std::rc::Rc; use anyhow::Error; @@ -10,20 +11,42 @@ use yew::prelude::*; use crate::AsyncAbortGuard; use crate::prelude::*; use crate::props::{IntoLoadCallback, IntoStorageLocation, LoadCallback, StorageLocation}; -use crate::state::{SharedState, SharedStateObserver, SharedStateReadGuard, SharedStateWriteGuard}; +use crate::state::{ + SharedState, SharedStateObserver, SharedStateReadGuard, SharedStateWriteGuard, + optional_rc_ptr_eq, +}; use crate::widget::{Button, Container, Fa, error_message}; /// Shared HTTP load state /// -/// This struct stores the state (loading) and the result of the load. +/// Stores the active request and the load result. pub struct LoaderState { - loading: u64, storage_location: Option, - async_abort_guard: Option, + active_load: Option, pub loader: Option>, pub data: Option, Error>>, } +struct ActiveLoad { + // The future holds a clone, preventing address reuse while it can still compare the request ID. + id: Rc<()>, + _abort_guard: AsyncAbortGuard, +} + +struct LoaderInner { + state: SharedState>, +} + +impl Drop for LoaderInner { + fn drop(&mut self) { + // Observer handles can outlive the Loader and keep its shared state allocated, but must not + // keep a request running after the last Loader clone is dropped. + let mut state = self.state.write(); + state.notify = false; + state.active_load = None; + } +} + impl LoaderState { fn load_from_cache(&mut self) { let storage_location = match &self.storage_location { @@ -61,9 +84,18 @@ impl LoaderState { /// - ability to cache result in local (default) or session storage by setting `state_id`. /// - helper to simplify renderering `self.render`. /// +/// Starting a load replaces the previous request. Dropping the last Loader clone cancels the active +/// request; pending futures and listener registrations do not themselves keep the Loader alive. +/// The canceled future is dropped when the executor observes the abort. Aborting an underlying HTTP +/// request depends on the load callback's cancellation behavior. #[derive(Derivative)] #[derivative(Clone(bound = ""), PartialEq(bound = ""))] -pub struct Loader(SharedState>); +pub struct Loader { + #[derivative(PartialEq(compare_with = "optional_rc_ptr_eq"))] + on_change: Option>>>, + #[derivative(PartialEq(compare_with = "Rc::ptr_eq"))] + inner: Rc>, +} impl Default for Loader { fn default() -> Self { @@ -75,13 +107,17 @@ impl Loader { /// Create a new instance. pub fn new() -> Self { let state = LoaderState { - loading: 0, data: None, loader: None, storage_location: None, - async_abort_guard: None, + active_load: None, }; - Self(SharedState::new(state)) + Self { + on_change: None, + inner: Rc::new(LoaderInner { + state: SharedState::new(state), + }), + } } /// Builder style method to set the persistent state ID. @@ -97,14 +133,14 @@ impl Loader { me.load_from_cache(); } + /// Register a listener kept alive by this handle and its clones. + /// + /// Starting a load does not notify listeners. Accepted completion, explicit abort, and writes + /// through [Self::write] do. The Loader passed to the callback does not retain registrations. pub fn on_change(mut self, cb: impl IntoEventCallback>) -> Self { - let me = self.clone(); - match cb.into_event_callback() { - Some(cb) => self.0.set_on_change(move |_| cb.emit(me.clone())), - _ => self - .0 - .set_on_change(None::>>>), - }; + self.on_change = cb + .into_event_callback() + .map(|cb| Rc::new(self.add_listener(cb))); self } @@ -121,24 +157,33 @@ impl Loader { me.loader = callback.into_load_callback(); } + /// Register a listener until the returned observer is dropped, without retaining the Loader. pub fn add_listener( &self, cb: impl Into>>, ) -> SharedStateObserver> { - let me = self.clone(); + let inner = Rc::downgrade(&self.inner); let cb = cb.into(); - self.0.add_listener(move |_| cb.emit(me.clone())) + self.inner.state.add_listener(move |_| { + if let Some(inner) = inner.upgrade() { + cb.emit(Self { + on_change: None, + inner, + }); + } + }) } pub fn read(&self) -> SharedStateReadGuard<'_, LoaderState> { - self.0.read() + self.inner.state.read() } pub fn write(&self) -> SharedStateWriteGuard<'_, LoaderState> { - self.0.write() + self.inner.state.write() } + /// Whether a current request is pending, excluding canceled and superseded requests. pub fn loading(&self) -> bool { - self.read().loading > 0 + self.read().active_load.is_some() } pub fn has_valid_data(&self) -> bool { @@ -158,33 +203,51 @@ impl Loader { } } + /// Start a request, replacing any pending request without notifying listeners on entry. pub fn load(&self) { - let loader = match &self.read().loader { - Some(loader) => loader.clone(), - None => return, // do nothing - }; + if let Some(future) = self.start_load() { + wasm_bindgen_futures::spawn_local(future); + } + } - let me = self.clone(); - - let mut state = self.write(); - state.loading += 1; - state.async_abort_guard = Some(AsyncAbortGuard::spawn(async move { + fn start_load(&self) -> Option + use> { + let loader = self.read().loader.clone()?; + let inner = Rc::downgrade(&self.inner); + let id = Rc::new(()); + let request_id = Rc::clone(&id); + let (abort_guard, future) = AsyncAbortGuard::new(async move { let res = loader.apply().await; - let mut me = me.write(); - me.async_abort_guard = None; - me.loading -= 1; - me.data = Some(res.map(|data| Rc::new(data))); - me.store_to_cache(); - })); + let Some(inner) = inner.upgrade() else { + return; + }; + let mut state = inner.state.write(); + // A callback can replace or abort its own request during its final poll. Aborting the + // future alone does not prevent that poll from completing. + if !state + .active_load + .as_ref() + .is_some_and(|active| Rc::ptr_eq(&active.id, &request_id)) + { + state.notify = false; + return; + } + state.active_load = None; + state.data = Some(res.map(Rc::new)); + state.store_to_cache(); + }); - // we don't want to notify listeners here, because we just set up the loading + let mut state = self.write(); state.notify = false; - drop(state); + state.active_load = Some(ActiveLoad { + id, + _abort_guard: abort_guard, + }); + Some(future) } - /// Abort any currently running load. + /// Abort the current request and notify listeners, without changing data or persistent cache. pub fn abort(&mut self) { - self.write().async_abort_guard = None; + self.write().active_load = None; } pub fn reload_button(&self) -> Button { @@ -192,3 +255,6 @@ impl Loader { Button::refresh(self.loading()).onclick(move |_| loader.load()) } } + +#[cfg(test)] +mod tests; diff --git a/src/state/loader/tests.rs b/src/state/loader/tests.rs new file mode 100644 index 00000000..191d63d4 --- /dev/null +++ b/src/state/loader/tests.rs @@ -0,0 +1,430 @@ +use std::cell::RefCell; +use std::collections::VecDeque; +use std::pin::Pin; +use std::task::{Context, Poll}; + +use futures::channel::oneshot; +use futures::task::noop_waker; + +use super::*; + +type Task = Pin>>; +type Snapshot = (bool, Option>); +type Sender = oneshot::Sender>; +type Requests = Rc>>; + +fn poll(task: &mut Task) -> Poll<()> { + let waker = noop_waker(); + task.as_mut().poll(&mut Context::from_waker(&waker)) +} + +fn configure(loader: &mut Loader) -> Requests { + let requests = Rc::new(RefCell::new(VecDeque::new())); + loader.set_loader({ + let requests = Rc::clone(&requests); + move || { + let (sender, receiver) = oneshot::channel(); + requests.borrow_mut().push_back(sender); + async move { receiver.await.unwrap() } + } + }); + requests +} + +fn start(loader: &Loader) -> Task { + Box::pin(loader.start_load().unwrap()) +} + +fn pending_request(loader: &Loader, requests: &Requests) -> (Sender, Task) { + let mut task = start(loader); + assert!(poll(&mut task).is_pending()); + let sender = requests.borrow_mut().pop_front().unwrap(); + (sender, task) +} + +fn snapshot(loader: &Loader) -> Snapshot { + let data = loader.read().data.as_ref().map(|result| match result { + Ok(data) => Ok(**data), + Err(err) => Err(err.to_string()), + }); + (loader.loading(), data) +} + +fn record(events: &Rc>>) -> Callback> { + let events = Rc::clone(events); + Callback::from(move |loader| events.borrow_mut().push(snapshot(&loader))) +} + +#[test] +fn missing_or_cleared_callback_does_not_replace_request() { + let events = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::::new().on_change(record(&events)); + assert!(loader.start_load().is_none()); + assert_eq!(snapshot(&loader), (false, None)); + assert!(events.borrow().is_empty()); + + let requests = configure(&mut loader); + let mut task = start(&loader); + loader.set_loader(None::>); + assert!(loader.start_load().is_none()); + assert_eq!(snapshot(&loader), (true, None)); + assert!(events.borrow().is_empty()); + assert!(poll(&mut task).is_pending()); + let sender = requests.borrow_mut().pop_front().unwrap(); + sender.send(Ok(1)).unwrap(); + assert!(poll(&mut task).is_ready()); + assert_eq!(snapshot(&loader), (false, Some(Ok(1)))); + assert_eq!(*events.borrow(), vec![(false, Some(Ok(1)))]); +} + +#[test] +fn sequential_loads_reuse_callback_and_notify_only_on_completion() { + let events = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::new().on_change(record(&events)); + let requests = configure(&mut loader); + let mut expected = Vec::new(); + for value in 1..=3 { + let (sender, mut task) = pending_request(&loader, &requests); + assert_eq!( + snapshot(&loader), + (true, (value > 1).then_some(Ok(value - 1))) + ); + assert_eq!(*events.borrow(), expected); + sender.send(Ok(value)).unwrap(); + assert!(poll(&mut task).is_ready()); + assert_eq!(snapshot(&loader), (false, Some(Ok(value)))); + assert!(loader.has_valid_data()); + expected.push((false, Some(Ok(value)))); + assert_eq!(*events.borrow(), expected); + } +} + +#[test] +fn replacement_retires_old_request_in_either_poll_order() { + for retire_old_first in [true, false] { + let events = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::new().on_change(record(&events)); + let requests = configure(&mut loader); + let (old_sender, mut old_task) = pending_request(&loader, &requests); + let (sender, mut task) = pending_request(&loader, &requests); + assert_eq!(snapshot(&loader), (true, None)); + assert!(events.borrow().is_empty()); + + let old_sender = if retire_old_first { + // A result can become ready after cancellation but before the executor observes it. + old_sender.send(Ok(1)).unwrap(); + assert!(poll(&mut old_task).is_ready()); + assert_eq!(snapshot(&loader), (true, None)); + assert!(events.borrow().is_empty()); + None + } else { + Some(old_sender) + }; + + sender.send(Ok(2)).unwrap(); + assert!(poll(&mut task).is_ready()); + assert_eq!(snapshot(&loader), (false, Some(Ok(2)))); + assert_eq!(*events.borrow(), vec![(false, Some(Ok(2)))]); + if let Some(old_sender) = old_sender { + old_sender.send(Err(Error::msg("superseded"))).unwrap(); + assert!(poll(&mut old_task).is_ready()); + assert_eq!(snapshot(&loader), (false, Some(Ok(2)))); + assert_eq!(*events.borrow(), vec![(false, Some(Ok(2)))]); + } + } +} + +#[test] +fn queued_replacements_invoke_only_current_callback() { + let events = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::new().on_change(record(&events)); + let requests = configure(&mut loader); + let mut tasks: Vec<_> = (0..3).map(|_| start(&loader)).collect(); + assert_eq!(snapshot(&loader), (true, None)); + assert!(requests.borrow().is_empty()); + assert!(events.borrow().is_empty()); + + let mut task = tasks.pop().unwrap(); + assert!(poll(&mut task).is_pending()); + assert_eq!(requests.borrow().len(), 1); + let sender = requests.borrow_mut().pop_front().unwrap(); + sender.send(Ok(3)).unwrap(); + assert!(poll(&mut task).is_ready()); + for mut task in tasks { + assert!(poll(&mut task).is_ready()); + } + assert!(requests.borrow().is_empty()); + assert_eq!(snapshot(&loader), (false, Some(Ok(3)))); + assert_eq!(*events.borrow(), vec![(false, Some(Ok(3)))]); +} + +#[test] +fn abort_retires_queued_and_pending_requests_without_changing_data() { + for poll_first in [false, true] { + let events = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::new(); + loader.write().data = Some(Ok(Rc::new(7))); + let requests = configure(&mut loader); + let _observer = loader.add_listener(record(&events)); + let mut task = start(&loader); + assert_eq!(snapshot(&loader), (true, Some(Ok(7)))); + assert!(requests.borrow().is_empty()); + let sender = if poll_first { + assert!(poll(&mut task).is_pending()); + Some(requests.borrow_mut().pop_front().unwrap()) + } else { + None + }; + loader.abort(); + assert_eq!(snapshot(&loader), (false, Some(Ok(7)))); + assert_eq!(*events.borrow(), vec![(false, Some(Ok(7)))]); + assert!(poll(&mut task).is_ready()); + if let Some(sender) = sender { + assert!(sender.is_canceled()); + } + assert!(requests.borrow().is_empty()); + + loader.abort(); + assert_eq!(*events.borrow(), vec![(false, Some(Ok(7))); 2]); + } +} + +#[test] +fn load_after_abort_and_failure_recovers() { + let events = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::new().on_change(record(&events)); + let requests = configure(&mut loader); + let mut canceled = start(&loader); + loader.abort(); + assert!(poll(&mut canceled).is_ready()); + assert!(requests.borrow().is_empty()); + assert_eq!(*events.borrow(), vec![(false, None)]); + + let (sender, mut task) = pending_request(&loader, &requests); + sender.send(Err(Error::msg("failure"))).unwrap(); + assert!(poll(&mut task).is_ready()); + assert_eq!(snapshot(&loader), (false, Some(Err("failure".into())))); + assert!(!loader.has_valid_data()); + let mut expected = vec![(false, None), (false, Some(Err("failure".into())))]; + assert_eq!(*events.borrow(), expected); + + let (sender, mut task) = pending_request(&loader, &requests); + assert_eq!(snapshot(&loader), (true, Some(Err("failure".into())))); + assert_eq!(*events.borrow(), expected); + sender.send(Ok(3)).unwrap(); + assert!(poll(&mut task).is_ready()); + assert_eq!(snapshot(&loader), (false, Some(Ok(3)))); + assert!(loader.has_valid_data()); + expected.push((false, Some(Ok(3)))); + assert_eq!(*events.borrow(), expected); +} + +#[test] +fn failed_refresh_replaces_previous_data_with_error() { + let mut loader = Loader::new(); + let requests = configure(&mut loader); + let (sender, mut task) = pending_request(&loader, &requests); + sender.send(Ok(1)).unwrap(); + assert!(poll(&mut task).is_ready()); + + let (sender, mut task) = pending_request(&loader, &requests); + assert_eq!(snapshot(&loader), (true, Some(Ok(1)))); + sender.send(Err(Error::msg("failure"))).unwrap(); + assert!(poll(&mut task).is_ready()); + assert_eq!(snapshot(&loader), (false, Some(Err("failure".into())))); + assert!(!loader.has_valid_data()); +} + +#[test] +fn dropping_last_owner_cancels_queued_and_pending_requests() { + for (with_on_change, with_observer) in + [(false, false), (false, true), (true, false), (true, true)] + { + for poll_first in [false, true] { + let events = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::new(); + if with_on_change { + loader = loader.on_change(record(&events)); + } + let observer = with_observer.then(|| loader.add_listener(record(&events))); + let requests = configure(&mut loader); + let mut task = start(&loader); + assert!(requests.borrow().is_empty()); + let sender = if poll_first { + assert!(poll(&mut task).is_pending()); + Some(requests.borrow_mut().pop_front().unwrap()) + } else { + None + }; + drop(loader); + assert!(events.borrow().is_empty()); + assert!( + poll(&mut task).is_ready(), + "on_change={with_on_change}, observer={with_observer}, polled={poll_first}" + ); + if let Some(sender) = sender { + assert!(sender.is_canceled()); + } + assert!(requests.borrow().is_empty()); + assert!(events.borrow().is_empty()); + drop(observer); + } + } +} + +#[test] +fn dropping_one_clone_keeps_request_and_listener_alive() { + let events = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::new().on_change(record(&events)); + let requests = configure(&mut loader); + let clone = loader.clone(); + let (sender, mut task) = pending_request(&loader, &requests); + drop(loader); + assert!(clone.loading()); + assert!(poll(&mut task).is_pending()); + sender.send(Ok(1)).unwrap(); + assert!(poll(&mut task).is_ready()); + assert_eq!(snapshot(&clone), (false, Some(Ok(1)))); + assert_eq!(*events.borrow(), vec![(false, Some(Ok(1)))]); +} + +#[test] +fn replacing_listener_does_not_retain_previous_registration() { + let first = Rc::new(RefCell::new(Vec::new())); + let second = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::new().on_change(record(&first)); + let clone = loader.clone(); + loader = loader.on_change(record(&second)); + loader.abort(); + assert_eq!(*first.borrow(), vec![(false, None)]); + assert_eq!(*second.borrow(), vec![(false, None)]); + drop(clone); + loader.abort(); + assert_eq!(*first.borrow(), vec![(false, None)]); + assert_eq!(*second.borrow(), vec![(false, None); 2]); + loader = loader.on_change(None::>>); + loader.abort(); + assert_eq!(*second.borrow(), vec![(false, None); 2]); +} + +#[test] +fn callback_handle_owns_request_without_retaining_registration() { + let received = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::::new().on_change({ + let received = Rc::clone(&received); + move |loader| received.borrow_mut().push(loader) + }); + let requests = configure(&mut loader); + loader.abort(); + let callback_loader = received.borrow_mut().pop().unwrap(); + let (sender, mut task) = pending_request(&loader, &requests); + drop(loader); + assert!(callback_loader.loading()); + assert!(poll(&mut task).is_pending()); + sender.send(Ok(1)).unwrap(); + assert!(poll(&mut task).is_ready()); + assert_eq!(snapshot(&callback_loader), (false, Some(Ok(1)))); + assert!(received.borrow().is_empty()); +} + +#[test] +fn dropping_observer_can_release_last_captured_owner() { + for keep_owner in [false, true] { + for with_request in [false, true] { + let events = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::::new(); + let requests = configure(&mut loader); + let _observer = loader.add_listener(record(&events)); + let captured = loader.clone(); + let observer = loader.add_listener(move |_| { + let _ = captured.loading(); + }); + let mut pending = with_request.then(|| pending_request(&loader, &requests)); + let owner = keep_owner.then_some(loader); + drop(observer); + if let Some(owner) = &owner { + assert_eq!(snapshot(owner), (with_request, None)); + if let Some((_, task)) = &mut pending { + assert!(poll(task).is_pending()); + } + } + drop(owner); + if let Some((sender, mut task)) = pending { + assert!(poll(&mut task).is_ready()); + assert!(sender.is_canceled()); + } + assert!(events.borrow().is_empty()); + } + } +} + +#[test] +fn reentrant_replacement_cannot_publish_or_abort_new_request() { + let events = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::new().on_change(record(&events)); + let owner = Rc::new(loader.clone()); + let weak_owner = Rc::downgrade(&owner); + let tasks = Rc::new(RefCell::new(Vec::::new())); + loader.set_loader({ + let tasks = Rc::clone(&tasks); + move || { + let mut loader = weak_owner.upgrade().unwrap().as_ref().clone(); + let tasks = Rc::clone(&tasks); + async move { + loader.set_loader(|| async { Ok(2) }); + tasks.borrow_mut().push(start(&loader)); + Ok(1) + } + } + }); + let mut task = start(&loader); + assert!(poll(&mut task).is_ready()); + assert_eq!(snapshot(&loader), (true, None)); + assert!(events.borrow().is_empty()); + assert!(poll(&mut tasks.borrow_mut().pop().unwrap()).is_ready()); + assert_eq!(snapshot(&loader), (false, Some(Ok(2)))); + assert_eq!(*events.borrow(), vec![(false, Some(Ok(2)))]); +} + +#[test] +fn reentrant_abort_cannot_publish_result() { + let events = Rc::new(RefCell::new(Vec::new())); + let mut loader = Loader::new().on_change(record(&events)); + let owner = Rc::new(loader.clone()); + let weak_owner = Rc::downgrade(&owner); + loader.set_loader(move || { + let mut loader = weak_owner.upgrade().unwrap().as_ref().clone(); + async move { + loader.abort(); + Ok(1) + } + }); + let mut task = start(&loader); + assert!(poll(&mut task).is_ready()); + assert_eq!(snapshot(&loader), (false, None)); + assert_eq!(*events.borrow(), vec![(false, None)]); +} + +#[test] +fn dropping_last_owner_during_final_poll_cannot_publish_result() { + let mut loader = Loader::::new(); + let owner = Rc::new(RefCell::new(None)); + loader.set_loader({ + let owner = Rc::clone(&owner); + move || { + let owner = Rc::clone(&owner); + async move { + owner.borrow_mut().take(); + Ok(1) + } + } + }); + // Retain storage, not an owner, to observe whether completion publishes after teardown. + let state = loader.inner.state.clone(); + let inner = Rc::downgrade(&loader.inner); + let mut task = start(&loader); + *owner.borrow_mut() = Some(loader); + assert!(poll(&mut task).is_ready()); + assert!(inner.upgrade().is_none()); + assert!(state.read().data.is_none()); +} -- 2.47.3