From: Thomas Lamprecht <t.lamprecht@proxmox.com>
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 [thread overview]
Message-ID: <20260909010400.684482-3-t.lamprecht@proxmox.com> (raw)
In-Reply-To: <20260909010400.684482-1-t.lamprecht@proxmox.com>
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 <t.lamprecht@proxmox.com>
---
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<F>(future: F) -> Self
where
F: Future<Output = ()> + '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<F>(future: F) -> (Self, impl Future<Output = ()>)
+ where
+ F: Future<Output = ()>,
{
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<T> {
- loading: u64,
storage_location: Option<StorageLocation>,
- async_abort_guard: Option<AsyncAbortGuard>,
+ active_load: Option<ActiveLoad>,
pub loader: Option<LoadCallback<T>>,
pub data: Option<Result<Rc<T>, 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<T> {
+ state: SharedState<LoaderState<T>>,
+}
+
+impl<T> Drop for LoaderInner<T> {
+ 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<T: 'static + DeserializeOwned + Serialize> LoaderState<T> {
fn load_from_cache(&mut self) {
let storage_location = match &self.storage_location {
@@ -61,9 +84,18 @@ impl<T: 'static + DeserializeOwned + Serialize> LoaderState<T> {
/// - 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<T>(SharedState<LoaderState<T>>);
+pub struct Loader<T> {
+ #[derivative(PartialEq(compare_with = "optional_rc_ptr_eq"))]
+ on_change: Option<Rc<SharedStateObserver<LoaderState<T>>>>,
+ #[derivative(PartialEq(compare_with = "Rc::ptr_eq"))]
+ inner: Rc<LoaderInner<T>>,
+}
impl<T: 'static + DeserializeOwned + Serialize> Default for Loader<T> {
fn default() -> Self {
@@ -75,13 +107,17 @@ impl<T: 'static + DeserializeOwned + Serialize> Loader<T> {
/// 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<T: 'static + DeserializeOwned + Serialize> Loader<T> {
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<Loader<T>>) -> 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::<Callback<SharedState<LoaderState<T>>>>),
- };
+ self.on_change = cb
+ .into_event_callback()
+ .map(|cb| Rc::new(self.add_listener(cb)));
self
}
@@ -121,24 +157,33 @@ impl<T: 'static + DeserializeOwned + Serialize> Loader<T> {
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<Callback<Loader<T>>>,
) -> SharedStateObserver<LoaderState<T>> {
- 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<T>> {
- self.0.read()
+ self.inner.state.read()
}
pub fn write(&self) -> SharedStateWriteGuard<'_, LoaderState<T>> {
- 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<T: 'static + DeserializeOwned + Serialize> Loader<T> {
}
}
+ /// 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<impl Future<Output = ()> + use<T>> {
+ 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<T: 'static + DeserializeOwned + Serialize> Loader<T> {
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<Box<dyn Future<Output = ()>>>;
+type Snapshot = (bool, Option<Result<u32, String>>);
+type Sender = oneshot::Sender<Result<u32, Error>>;
+type Requests = Rc<RefCell<VecDeque<Sender>>>;
+
+fn poll(task: &mut Task) -> Poll<()> {
+ let waker = noop_waker();
+ task.as_mut().poll(&mut Context::from_waker(&waker))
+}
+
+fn configure(loader: &mut Loader<u32>) -> 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<u32>) -> Task {
+ Box::pin(loader.start_load().unwrap())
+}
+
+fn pending_request(loader: &Loader<u32>, 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<u32>) -> 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<RefCell<Vec<Snapshot>>>) -> Callback<Loader<u32>> {
+ 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::<u32>::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::<LoadCallback<u32>>);
+ 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::<Callback<Loader<u32>>>);
+ 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::<u32>::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::<u32>::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::<Task>::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::<u32>::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
prev parent reply other threads:[~2026-09-09 1:04 UTC|newest]
Thread overview: 3+ messages / expand[flat|nested] mbox.gz Atom feed top
2026-09-09 1:01 [PATCH 0/2] fix listener teardown and loader cancellation Thomas Lamprecht
2026-09-09 1:01 ` [PATCH 1/2] shared state: avoid borrow panics when removing listeners Thomas Lamprecht
2026-09-09 1:01 ` Thomas Lamprecht [this message]
Reply instructions:
You may reply publicly to this message via plain-text email
using any one of the following methods:
* Save the following mbox file, import it into your mail client,
and reply-to-all from there: mbox
Avoid top-posting and favor interleaved quoting:
https://en.wikipedia.org/wiki/Posting_style#Interleaved_style
* Reply using the --to, --cc, and --in-reply-to
switches of git-send-email(1):
git send-email \
--in-reply-to=20260909010400.684482-3-t.lamprecht@proxmox.com \
--to=t.lamprecht@proxmox.com \
--cc=yew-devel@lists.proxmox.com \
/path/to/YOUR_REPLY
https://kernel.org/pub/software/scm/git/docs/git-send-email.html
* If your mail client supports setting the In-Reply-To header
via mailto: links, try the mailto: link
Be sure your reply has a Subject: header at the top and a blank line
before the message body.
This is a public inbox, see mirroring instructions
for how to clone and mirror all data and code used for this inbox