diff --git a/CHANGELOG.md b/CHANGELOG.md index 31929603..e659fba1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -41,6 +41,7 @@ All notable changes to this project will be documented in this file. ### Improvements +* Allow `Condvar` waits with unsized mutex contents, including slices and trait objects, for both borrowed and owned guards. * Allow manual unbounded pools to use non-`Send` factory closures with `get_or_create`, including factories that borrow task-local state. * Allow `watch::channel` to store non-`Clone` values for publication and change notification; only owning reads through `Receiver::get` and `Receiver::recv` require `Clone`. * Finish releasing buffered bounded MPSC messages even if one message destructor panics. diff --git a/asyncband/src/condvar/mod.rs b/asyncband/src/condvar/mod.rs index 06a8550a..f2b06524 100644 --- a/asyncband/src/condvar/mod.rs +++ b/asyncband/src/condvar/mod.rs @@ -193,7 +193,7 @@ impl Condvar { /// [`notify_one`](Self::notify_one) but has not yet reacquired the mutex, the notification is /// passed to another task that is waiting at that point, if one exists. It is never buffered /// for a future waiter. - pub async fn wait<'a, T>(&self, guard: MutexGuard<'a, T>) -> MutexGuard<'a, T> { + pub async fn wait<'a, T: ?Sized>(&self, guard: MutexGuard<'a, T>) -> MutexGuard<'a, T> { let mutex = mutex::guard_lock(&guard); let notify_one_baton = Wait { condvar: self, @@ -212,7 +212,7 @@ impl Condvar { /// /// This has the same notification and cancellation semantics as [`wait`](Self::wait), but /// accepts and returns an owned guard. - pub async fn wait_owned(&self, guard: OwnedMutexGuard) -> OwnedMutexGuard { + pub async fn wait_owned(&self, guard: OwnedMutexGuard) -> OwnedMutexGuard { let mutex = mutex::owned_guard_lock(&guard); let notify_one_baton = Wait { condvar: self, @@ -263,7 +263,7 @@ impl Condvar { /// /// Each wait iteration has the same cancellation semantics as [`wait`](Self::wait). Cancelling /// drops the mutex guard; mutations already made by `condition` are not rolled back. - pub async fn wait_while<'a, T, F>( + pub async fn wait_while<'a, T: ?Sized, F>( &self, mut guard: MutexGuard<'a, T>, mut condition: F, @@ -314,7 +314,7 @@ impl Condvar { /// Each wait iteration has the same cancellation semantics as /// [`wait_owned`](Self::wait_owned). Cancelling drops the owned mutex guard; mutations already /// made by `condition` are not rolled back. - pub async fn wait_while_owned( + pub async fn wait_while_owned( &self, mut guard: OwnedMutexGuard, mut condition: F, diff --git a/asyncband/src/once/once/mod.rs b/asyncband/src/once/once/mod.rs index 4870435a..5389d858 100644 --- a/asyncband/src/once/once/mod.rs +++ b/asyncband/src/once/once/mod.rs @@ -32,8 +32,11 @@ use crate::semaphore::Semaphore; /// /// This type also intentionally omits "poisoning" semantics. If an initialization future is /// cancelled or panics, the attempt is abandoned and other tasks may retry the operation. -/// Encode partial-initialization detection in the future itself (e.g. return a `Result`) -/// when needed. +/// Retrying does not undo side effects from the abandoned attempt. +/// +/// [`call_once`](Self::call_once) accepts only initializers returning `()`. For fallible +/// initialization, enable the `once-cell` feature and use `OnceCell::get_or_try_init`, which leaves +/// the cell empty on error. Use `OnceCell<()>` when no initialized value needs to be stored. /// /// See the [module level documentation](super) for additional context. /// diff --git a/asyncband/src/pool/bounded.rs b/asyncband/src/pool/bounded.rs index 4ac17c65..c6309911 100644 --- a/asyncband/src/pool/bounded.rs +++ b/asyncband/src/pool/bounded.rs @@ -264,6 +264,10 @@ impl Pool { /// If the pool has reached its maximum size and has no idle object, this method waits until an /// object is returned to or detached from the pool. /// + /// Idle objects are checked with [`ManageObject::is_recyclable`]. A failed check detaches the + /// object and retries checkout; its error is discarded. Only errors from + /// [`ManageObject::create`] are returned to the caller. + /// /// # Cancel safety /// /// Cancelling while waiting for capacity or creating a new object restores the reserved pool diff --git a/asyncband/src/pool/common.rs b/asyncband/src/pool/common.rs index 31019612..894b777b 100644 --- a/asyncband/src/pool/common.rs +++ b/asyncband/src/pool/common.rs @@ -89,7 +89,9 @@ pub trait ManageObject: Send + Sync { /// Whether the object `o` is recyclable. /// - /// Returns `Ok(())` if the object is recyclable; otherwise, returns an error. + /// Returns `Ok(())` if the object is recyclable. On error, the pool detaches the object and + /// retries checkout with another idle object or creates a replacement. The error is discarded + /// rather than returned by `Pool::get`; record any needed diagnostics in this implementation. fn is_recyclable( &self, o: &mut Self::Object, diff --git a/asyncband/src/pool/unbounded.rs b/asyncband/src/pool/unbounded.rs index ebc5dd3d..927530fd 100644 --- a/asyncband/src/pool/unbounded.rs +++ b/asyncband/src/pool/unbounded.rs @@ -316,6 +316,10 @@ impl> Pool { /// /// If no idle object is available, this method calls [`ManageObject::create`]. /// + /// Idle objects are checked with [`ManageObject::is_recyclable`]. A failed check detaches the + /// object and retries checkout; its error is discarded. Only errors from + /// [`ManageObject::create`] are returned to the caller. + /// /// # Cancel safety /// /// Cancelling while creating a new object leaves the pool unchanged. Cancelling while diff --git a/tests-integration/tests/condvar_test.rs b/tests-integration/tests/condvar_test.rs index 2f95bf25..fb7128fe 100644 --- a/tests-integration/tests/condvar_test.rs +++ b/tests-integration/tests/condvar_test.rs @@ -256,3 +256,28 @@ fn wait_owned_reacquires_the_mutex() { assert_eq!(*mutex.lock().await, 1); }); } + +#[test] +fn predicate_waits_support_unsized_state() { + let mutex: Arc> = Arc::new(Mutex::new([0, 0])); + let condvar = Condvar::new(); + + let guard = mutex.try_lock().unwrap(); + let mut borrowed = Box::pin(condvar.wait_while(guard, |bytes| bytes[0] == 0)); + assert!(poll_once(borrowed.as_mut()).is_pending()); + mutex.try_lock().unwrap()[0] = 1; + condvar.notify_one(); + let guard = expect_ready(poll_once(borrowed.as_mut())); + assert_eq!(&*guard, &[1, 0]); + assert!(mutex.try_lock().is_none()); + drop(guard); + + let guard = mutex.clone().try_lock_owned().unwrap(); + let mut owned = Box::pin(condvar.wait_while_owned(guard, |bytes| bytes[1] == 0)); + assert!(poll_once(owned.as_mut()).is_pending()); + mutex.try_lock().unwrap()[1] = 2; + condvar.notify_one(); + let guard = expect_ready(poll_once(owned.as_mut())); + assert_eq!(&*guard, &[1, 2]); + assert!(mutex.try_lock().is_none()); +}