diff --git a/CHANGELOG.md b/CHANGELOG.md index 81e16b5f..73204070 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -42,6 +42,7 @@ All notable changes to this project will be documented in this file. ### Improvements * Allow `watch::channel` to store non-`Clone` values for publication and change notification; only owning reads through `Receiver::get` and `Receiver::recv` require `Clone`. +* Reduce `ManualResetEvent` state-query and already-set wait overhead, especially with concurrent readers. * Finish releasing buffered bounded MPSC messages even if one message destructor panics. * Improve unbounded MPSC throughput with batched receiving and incremental storage reclamation; empty-buffer retention is bounded independently of previous peak occupancy. * Make completed and abandoned `Completion` waits lock-free while preserving cancellable pending registration. diff --git a/asyncband/src/event/manual_reset.rs b/asyncband/src/event/manual_reset.rs index d51b16a4..b7f3653e 100644 --- a/asyncband/src/event/manual_reset.rs +++ b/asyncband/src/event/manual_reset.rs @@ -20,6 +20,8 @@ use std::future::Future; use std::mem; use std::pin::Pin; use std::sync::Arc; +use std::sync::atomic::AtomicBool; +use std::sync::atomic::Ordering; use std::task::Context; use std::task::Poll; use std::task::Waker; @@ -72,7 +74,9 @@ use crate::internal::waker_batch::WakerBatch; /// # } /// ``` pub struct ManualResetEvent { - state: Mutex, + // Flag writes and waiter-list changes hold `waiters`. Lock-free paths only read the flag. + is_set: AtomicBool, + waiters: Mutex>, } impl ManualResetEvent { @@ -86,10 +90,8 @@ impl ManualResetEvent { /// If `is_set` is `true`, waits complete immediately until the event is reset. pub const fn with_state(is_set: bool) -> Self { Self { - state: Mutex::new(State { - is_set, - waiters: WaitList::new(), - }), + is_set: AtomicBool::new(is_set), + waiters: Mutex::new(WaitList::new()), } } @@ -104,16 +106,17 @@ impl ManualResetEvent { /// attempted for every other selected task before the panic resumes. pub fn set(&self) { let wakers = { - let mut state = self.state.lock(); - if state.is_set { + let mut waiters = self.waiters.lock(); + if self.is_set.load(Ordering::Relaxed) { return; } - state.is_set = true; + // Publish to waits that observe the set state without taking the lock. + self.is_set.store(true, Ordering::Release); // Detach the complete cohort before invoking any waker. A wake callback may reset the // event and register a new wait, which must belong to the state current at that point. let mut wakers = WakerBatch::new(); - while let Some((_id, waiter)) = state.waiters.unlink_first_waiter(|waiter| { + while let Some((_id, waiter)) = waiters.unlink_first_waiter(|waiter| { waiter.notified = true; true }) { @@ -132,7 +135,9 @@ impl ManualResetEvent { /// Waits already released by a preceding [`set`](Self::set) remain ready. If the event is /// already unset, this has no effect. pub fn reset(&self) { - self.state.lock().is_set = false; + let _waiters = self.waiters.lock(); + // Clearing the flag does not publish data to successful waits. + self.is_set.store(false, Ordering::Relaxed); } /// Returns whether the event is currently set. @@ -150,7 +155,7 @@ impl ManualResetEvent { /// assert!(event.is_set()); /// ``` pub fn is_set(&self) -> bool { - self.state.lock().is_set + self.is_set.load(Ordering::Acquire) } /// Attempts to wait without registering a waiter. @@ -219,27 +224,34 @@ impl ManualResetEvent { /// event is set never enqueues. A linked waiter therefore always belongs to an unset event, so /// `notified` alone decides whether a registered waiter is already committed. fn poll_wait(&self, waiter_id: &mut Option, cx: &mut Context<'_>) -> Poll<()> { + // A registered wait must still remove its node, even if the event has since been set. + if waiter_id.is_none() && self.is_set() { + return Poll::Ready(()); + } + let (poll, retired_waker) = { - let mut state = self.state.lock(); + let mut waiters = self.waiters.lock(); match *waiter_id { - Some(id) if state.waiters.waiter_mut(id).notified => { - let waiter = state.remove_waiter(id); + Some(id) if waiters.waiter_mut(id).notified => { + let waiter = waiters.remove_unlinked_waiter(id); *waiter_id = None; (Poll::Ready(()), waiter.waker) } Some(id) => { debug_assert!( - !state.is_set, + !self.is_set.load(Ordering::Relaxed), "a linked waiter must belong to an unset event" ); - let waiter = state.waiters.waiter_mut(id); + let waiter = waiters.waiter_mut(id); let retired = (!waiter.will_wake(cx.waker())) .then(|| waiter.replace_waker(cx.waker().clone())); (Poll::Pending, retired) } - None if state.is_set => (Poll::Ready(()), None), + // Recheck under the lock so a set between the fast probe and registration + // either completes this wait here or selects its registered node later. + None if self.is_set.load(Ordering::Relaxed) => (Poll::Ready(()), None), None => { - *waiter_id = Some(state.waiters.push_back(Waiter { + *waiter_id = Some(waiters.push_back(Waiter { notified: false, waker: Some(cx.waker().clone()), })); @@ -254,8 +266,10 @@ impl ManualResetEvent { fn unregister_waiter(&self, id: WaiterId) { let waiter = { - let mut state = self.state.lock(); - state.remove_waiter(id) + let mut waiters = self.waiters.lock(); + // A released waiter is already detached, but retains its node until removal. + waiters.unlink_waiter(id, |_| true); + waiters.remove_unlinked_waiter(id) }; drop(waiter); } @@ -269,29 +283,13 @@ impl Default for ManualResetEvent { impl fmt::Debug for ManualResetEvent { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - let is_set = self.state.lock().is_set; + let is_set = self.is_set(); f.debug_struct("ManualResetEvent") .field("is_set", &is_set) .finish_non_exhaustive() } } -#[derive(Debug)] -struct State { - is_set: bool, - waiters: WaitList, -} - -impl State { - /// Removes a waiter whether or not [`ManualResetEvent::set`] already unlinked it. - fn remove_waiter(&mut self, id: WaiterId) -> Waiter { - // Unlinking is idempotent: a waiter that `set` detached keeps its node until it is removed - // here, and an unconditional predicate never declines. - self.waiters.unlink_waiter(id, |_| true); - self.waiters.remove_unlinked_waiter(id) - } -} - #[derive(Debug)] struct Waiter { notified: bool, diff --git a/benchmarks/asyncband/event/wait.rs b/benchmarks/asyncband/event/wait.rs index 18ada489..cf0b2b85 100644 --- a/benchmarks/asyncband/event/wait.rs +++ b/benchmarks/asyncband/event/wait.rs @@ -129,3 +129,10 @@ fn waiter_fan_out(bencher: Bencher, waiter_count: usize) { black_box(event) }); } + +#[divan::bench(args = [false, true], threads = THREAD_COUNTS, sample_size = CONTENDED_SAMPLE_SIZE)] +fn try_wait_contended(bencher: Bencher, is_set: bool) { + let event = ManualResetEvent::with_state(is_set); + + bencher.bench(|| black_box(event.try_wait())); +} diff --git a/tests-integration/tests/event_test.rs b/tests-integration/tests/event_test.rs index f234a7cf..cb8479bf 100644 --- a/tests-integration/tests/event_test.rs +++ b/tests-integration/tests/event_test.rs @@ -21,13 +21,17 @@ use std::panic::AssertUnwindSafe; use std::pin::Pin; use std::pin::pin; use std::sync::Arc; +use std::sync::Barrier; use std::sync::Mutex; use std::sync::atomic::AtomicBool; +use std::sync::atomic::AtomicUsize; use std::sync::atomic::Ordering; use std::task::Context; use std::task::Wake; use std::task::Waker; +use std::thread; +use asyncband::blocking::FutureExt; use asyncband::event::ManualResetEvent; use tests_integration::PanicWake; use tests_integration::WakeCounter; @@ -173,7 +177,7 @@ impl Wake for ReentrantWaker { impl Drop for ReentrantWaker { fn drop(&mut self) { - self.0.is_set(); + self.0.reset(); } } @@ -225,7 +229,7 @@ fn wakers_are_woken_and_dropped_outside_the_internal_lock() { ); } - // Cancelling a pending wait drops the registered waker, which re-enters `is_set`. The + // Cancelling a pending wait drops the registered waker, which re-enters `reset`. The // registration holds the last reference, so the drop runs here. drop(replaced); }); @@ -384,3 +388,36 @@ fn cancelling_an_owned_waiter_releases_its_waker_and_event_handle() { event.set(); assert_eq!(tracker.count(), 0); } + +#[test] +fn concurrent_sets_and_waits_publish_state_across_reset_cycles() { + assert_completes_without_deadlock(|| { + let event = ManualResetEvent::new(); + let round = Barrier::new(2); + let value = AtomicUsize::new(0); + thread::scope(|scope| { + scope.spawn(|| { + for expected in 1..=100 { + round.wait(); + value.store(expected, Ordering::Relaxed); + event.set(); + round.wait(); + } + }); + for expected in 1..=100 { + // The opening barrier precedes publication; only the event publishes value. + round.wait(); + if expected % 2 == 0 { + event.wait().block_on(); + } else { + while !event.try_wait() { + thread::yield_now(); + } + } + assert_eq!(value.load(Ordering::Relaxed), expected); + round.wait(); + event.reset(); + } + }); + }); +}