From e538329646623522f9a697877413a1b24d812a1f Mon Sep 17 00:00:00 2001 From: tison Date: Mon, 14 Sep 2026 23:54:08 +0800 Subject: [PATCH 1/3] perf(event): avoid locking state queries and immediate waits --- CHANGELOG.md | 1 + asyncband/src/event/auto_reset.rs | 41 +++-- asyncband/src/event/manual_reset.rs | 29 ++-- benchmarks/asyncband/event/auto_reset.rs | 150 ++++++++++++++++++ benchmarks/asyncband/event/mod.rs | 1 + benchmarks/asyncband/event/wait.rs | 7 + .../tests/auto_reset_event_test.rs | 53 ++++++- tests-integration/tests/event_test.rs | 41 ++++- 8 files changed, 295 insertions(+), 28 deletions(-) create mode 100644 benchmarks/asyncband/event/auto_reset.rs 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/auto_reset.rs b/asyncband/src/event/auto_reset.rs index af72eaec..5306b872 100644 --- a/asyncband/src/event/auto_reset.rs +++ b/asyncband/src/event/auto_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; @@ -76,6 +78,7 @@ use crate::internal::waitlist::WaiterId; /// # } /// ``` pub struct AutoResetEvent { + is_set: AtomicBool, state: Mutex, } @@ -90,8 +93,8 @@ impl AutoResetEvent { /// If `is_set` is `true`, the event stores one signal for a future wait. pub const fn with_state(is_set: bool) -> Self { Self { + is_set: AtomicBool::new(is_set), state: Mutex::new(State { - is_set, waiters: WaitList::new(), }), } @@ -108,7 +111,7 @@ impl AutoResetEvent { /// Panics if waking a selected task panics. Its signal remains assigned and can still be /// consumed by polling that wait or passed on by dropping it. pub fn set(&self) { - let waker = self.state.lock().signal(); + let waker = self.state.lock().signal(&self.is_set); if let Some(waker) = waker { waker.wake(); } @@ -119,7 +122,8 @@ impl AutoResetEvent { /// Signals already assigned to waits remain theirs. Cancelling such a wait can still transfer /// or restore its signal after this call. If no signal is stored, this has no effect. pub fn reset(&self) { - self.state.lock().is_set = false; + let _state = self.state.lock(); + self.is_set.store(false, Ordering::Relaxed); } /// Returns whether the event is currently set. @@ -140,7 +144,7 @@ impl AutoResetEvent { /// 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. @@ -159,7 +163,13 @@ impl AutoResetEvent { /// assert!(!event.try_wait()); // A successful wait consumes the signal. /// ``` pub fn try_wait(&self) -> bool { - mem::take(&mut self.state.lock().is_set) + // Only consumption bypasses the waiter lock. Signals are stored under that lock only + // when no wait is queued, so this cannot take a signal assigned to a registered wait. + self.is_set.load(Ordering::Relaxed) + && self + .is_set + .compare_exchange(true, false, Ordering::Acquire, Ordering::Relaxed) + .is_ok() } /// Waits for and consumes one signal. @@ -207,6 +217,11 @@ impl AutoResetEvent { } fn poll_wait(&self, waiter_id: &mut Option, cx: &mut Context<'_>) -> Poll<()> { + // Registered waits own separate signals and must remove their nodes under the lock. + if waiter_id.is_none() && self.try_wait() { + return Poll::Ready(()); + } + let (poll, retired_waker) = { let mut state = self.state.lock(); match *waiter_id { @@ -222,10 +237,8 @@ impl AutoResetEvent { (Poll::Pending, retired) } }, - None if state.is_set => { - state.is_set = false; - (Poll::Ready(()), None) - } + // Recheck before enqueueing: set and cancellation handoff use this same lock. + None if self.try_wait() => (Poll::Ready(()), None), None => { *waiter_id = Some(state.waiters.push_back(Waiter::Waiting(cx.waker().clone()))); (Poll::Pending, None) @@ -243,7 +256,7 @@ impl AutoResetEvent { state.waiters.unlink_waiter(id, |_| true); let waiter = state.waiters.remove_unlinked_waiter(id); let waker = match &waiter { - Waiter::Notified => state.signal(), + Waiter::Notified => state.signal(&self.is_set), Waiter::Waiting(_) => None, }; (waiter, waker) @@ -263,7 +276,7 @@ impl Default for AutoResetEvent { impl fmt::Debug for AutoResetEvent { 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("AutoResetEvent") .field("is_set", &is_set) .finish_non_exhaustive() @@ -273,19 +286,19 @@ impl fmt::Debug for AutoResetEvent { struct State { // A stored signal and queued (unselected) waits never coexist. Detached, selected waits can // coexist with either: their signals are reserved until consumption or cancellation. - is_set: bool, waiters: WaitList, } impl State { - fn signal(&mut self) -> Option { + fn signal(&mut self, is_set: &AtomicBool) -> Option { if let Some((_, waiter)) = self.waiters.unlink_first_waiter(|_| true) { let Waiter::Waiting(waker) = mem::replace(waiter, Waiter::Notified) else { unreachable!("only unselected waits remain queued") }; Some(waker) } else { - self.is_set = true; + // Publish every set, including coalesced sets and returned assigned signals. + is_set.store(true, Ordering::Release); None } } diff --git a/asyncband/src/event/manual_reset.rs b/asyncband/src/event/manual_reset.rs index d51b16a4..32689200 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,6 +74,7 @@ use crate::internal::waker_batch::WakerBatch; /// # } /// ``` pub struct ManualResetEvent { + is_set: AtomicBool, state: Mutex, } @@ -86,8 +89,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 { + is_set: AtomicBool::new(is_set), state: Mutex::new(State { - is_set, waiters: WaitList::new(), }), } @@ -105,11 +108,12 @@ impl ManualResetEvent { pub fn set(&self) { let wakers = { let mut state = self.state.lock(); - if state.is_set { + if self.is_set.load(Ordering::Relaxed) { return; } - state.is_set = true; + // Serialize publication, cohort detach, reset, and registration with the same 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(); @@ -132,7 +136,8 @@ 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 _state = self.state.lock(); + 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,6 +224,11 @@ 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(); match *waiter_id { @@ -229,7 +239,7 @@ impl ManualResetEvent { } 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); @@ -237,7 +247,9 @@ impl ManualResetEvent { .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 { notified: false, @@ -269,7 +281,7 @@ 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() @@ -278,7 +290,6 @@ impl fmt::Debug for ManualResetEvent { #[derive(Debug)] struct State { - is_set: bool, waiters: WaitList, } diff --git a/benchmarks/asyncband/event/auto_reset.rs b/benchmarks/asyncband/event/auto_reset.rs new file mode 100644 index 00000000..f4fc1fb7 --- /dev/null +++ b/benchmarks/asyncband/event/auto_reset.rs @@ -0,0 +1,150 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use std::pin::pin; + +use asyncband::event::AutoResetEvent; +use divan::Bencher; +use divan::black_box; +use divan::counter::ItemsCount; + +use crate::support::bench_context; +use crate::support::poll_pending; +use crate::support::poll_pinned_ready; + +const WAITER_COUNTS: &[usize] = &[1, 8, 32]; +const THREAD_COUNTS: &[usize] = &[1, 2, 8, 32]; +const CONTENDED_SAMPLE_SIZE: u32 = 256; + +#[divan::bench] +fn set_then_wait(bencher: Bencher) { + let event = AutoResetEvent::new(); + + bencher.bench_local(|| { + event.set(); + let mut context = bench_context(); + let mut wait = pin!(event.wait()); + poll_pinned_ready(wait.as_mut(), &mut context); + black_box(&event) + }); +} + +#[divan::bench(threads = THREAD_COUNTS, sample_size = CONTENDED_SAMPLE_SIZE)] +fn is_set_contended(bencher: Bencher) { + let event = AutoResetEvent::with_state(true); + + bencher.bench(|| black_box(event.is_set())); +} + +#[divan::bench] +fn set_reset_cycle(bencher: Bencher) { + let event = AutoResetEvent::new(); + + bencher.bench_local(|| { + event.set(); + event.reset(); + black_box(&event) + }); +} + +#[divan::bench] +fn cancel_pending(bencher: Bencher) { + let mut context = bench_context(); + + bencher.bench_local(|| { + let event = AutoResetEvent::new(); + { + let mut wait = pin!(event.wait()); + poll_pending(wait.as_mut(), &mut context); + } + black_box(event) + }); +} + +#[divan::bench] +fn wake_waiter(bencher: Bencher) { + let mut context = bench_context(); + + bencher.bench_local(|| { + let event = AutoResetEvent::new(); + { + let mut wait = pin!(event.wait()); + poll_pending(wait.as_mut(), &mut context); + + event.set(); + poll_pinned_ready(wait.as_mut(), &mut context); + } + black_box(event) + }); +} + +#[divan::bench] +fn wake_waiter_reused(bencher: Bencher) { + let mut context = bench_context(); + let event = AutoResetEvent::new(); + + bencher.bench_local(|| { + let mut wait = pin!(event.wait()); + poll_pending(wait.as_mut(), &mut context); + + event.set(); + poll_pinned_ready(wait.as_mut(), &mut context); + black_box(&event) + }); +} + +#[divan::bench(args = WAITER_COUNTS)] +fn wake_waiters(bencher: Bencher, waiter_count: usize) { + let mut context = bench_context(); + + bencher + .counter(ItemsCount::new(waiter_count)) + .bench_local(|| { + let event = AutoResetEvent::new(); + let mut waiters = (0..waiter_count) + .map(|_| Box::pin(event.wait())) + .collect::>(); + for waiter in &mut waiters { + poll_pending(waiter.as_mut(), &mut context); + } + + for _ in 0..waiter_count { + event.set(); + } + for mut waiter in waiters { + poll_pinned_ready(waiter.as_mut(), &mut context); + } + black_box(event) + }); +} + +#[divan::bench(threads = THREAD_COUNTS, sample_size = CONTENDED_SAMPLE_SIZE)] +fn try_wait_unset(bencher: Bencher) { + let event = AutoResetEvent::new(); + + bencher.bench(|| black_box(event.try_wait())); +} + +#[divan::bench] +fn set_then_try_wait(bencher: Bencher) { + let event = AutoResetEvent::new(); + + bencher.bench_local(|| { + event.set(); + assert!(event.try_wait()); + }); +} diff --git a/benchmarks/asyncband/event/mod.rs b/benchmarks/asyncband/event/mod.rs index 35540691..f4f10365 100644 --- a/benchmarks/asyncband/event/mod.rs +++ b/benchmarks/asyncband/event/mod.rs @@ -15,4 +15,5 @@ // specific language governing permissions and limitations // under the License. +mod auto_reset; mod wait; 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/auto_reset_event_test.rs b/tests-integration/tests/auto_reset_event_test.rs index dc0e66b8..a7c3ef38 100644 --- a/tests-integration/tests/auto_reset_event_test.rs +++ b/tests-integration/tests/auto_reset_event_test.rs @@ -257,13 +257,13 @@ struct ReentrantWaker(Arc); impl Wake for ReentrantWaker { fn wake(self: Arc) { - self.0.try_wait(); + self.0.reset(); } } impl Drop for ReentrantWaker { fn drop(&mut self) { - self.0.try_wait(); + self.0.reset(); } } @@ -352,10 +352,57 @@ fn concurrent_sets_and_waits_publish_state_without_losing_signals() { // Registration races with set; the second barrier prevents the next set from // coalescing before this round's signal has been consumed. round.wait(); - event.wait().block_on(); + 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(); } }); }); } + +#[test] +fn stored_signals_are_consumed_once_when_immediate_and_async_waits_race() { + assert_completes_without_deadlock(|| { + let event = AutoResetEvent::new(); + let round = Barrier::new(3); + let completed = AtomicUsize::new(0); + thread::scope(|scope| { + scope.spawn(|| { + for _ in 0..100 { + round.wait(); + if event.try_wait() { + completed.fetch_add(1, Ordering::Relaxed); + } + round.wait(); + } + }); + scope.spawn(|| { + for _ in 0..100 { + round.wait(); + { + let mut wait = pin!(event.wait()); + if poll_once(wait.as_mut()).is_ready() { + completed.fetch_add(1, Ordering::Relaxed); + } + } + // Remove any pending registration before the next signal is stored. + round.wait(); + } + }); + for _ in 0..100 { + completed.store(0, Ordering::Relaxed); + event.set(); + round.wait(); + round.wait(); + assert_eq!(completed.load(Ordering::Relaxed), 1); + assert!(!event.is_set()); + } + }); + }); +} 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(); + } + }); + }); +} From b25e6df9e6c1951f1e30aaaba4867cbfb94f4370 Mon Sep 17 00:00:00 2001 From: tison Date: Tue, 15 Sep 2026 06:37:41 +0800 Subject: [PATCH 2/3] refactor(event): inline waiter state and cover coalesced signals --- asyncband/src/event/auto_reset.rs | 62 +++++++++++------------- asyncband/src/event/manual_reset.rs | 43 ++++++---------- benchmarks/asyncband/event/auto_reset.rs | 26 ++++++++++ 3 files changed, 67 insertions(+), 64 deletions(-) diff --git a/asyncband/src/event/auto_reset.rs b/asyncband/src/event/auto_reset.rs index 5306b872..2a38b5d8 100644 --- a/asyncband/src/event/auto_reset.rs +++ b/asyncband/src/event/auto_reset.rs @@ -79,7 +79,9 @@ use crate::internal::waitlist::WaiterId; /// ``` pub struct AutoResetEvent { is_set: AtomicBool, - state: Mutex, + // A stored signal and queued (unselected) waits never coexist. Detached, selected waits can + // coexist with either: their signals are reserved until consumption or cancellation. + waiters: Mutex>, } impl AutoResetEvent { @@ -94,9 +96,7 @@ impl AutoResetEvent { pub const fn with_state(is_set: bool) -> Self { Self { is_set: AtomicBool::new(is_set), - state: Mutex::new(State { - waiters: WaitList::new(), - }), + waiters: Mutex::new(WaitList::new()), } } @@ -111,7 +111,7 @@ impl AutoResetEvent { /// Panics if waking a selected task panics. Its signal remains assigned and can still be /// consumed by polling that wait or passed on by dropping it. pub fn set(&self) { - let waker = self.state.lock().signal(&self.is_set); + let waker = self.signal(&mut self.waiters.lock()); if let Some(waker) = waker { waker.wake(); } @@ -122,7 +122,7 @@ impl AutoResetEvent { /// Signals already assigned to waits remain theirs. Cancelling such a wait can still transfer /// or restore its signal after this call. If no signal is stored, this has no effect. pub fn reset(&self) { - let _state = self.state.lock(); + let _waiters = self.waiters.lock(); self.is_set.store(false, Ordering::Relaxed); } @@ -223,11 +223,11 @@ impl AutoResetEvent { } let (poll, retired_waker) = { - let mut state = self.state.lock(); + let mut waiters = self.waiters.lock(); match *waiter_id { - Some(id) => match state.waiters.waiter_mut(id) { + Some(id) => match waiters.waiter_mut(id) { Waiter::Notified => { - state.waiters.remove_unlinked_waiter(id); + waiters.remove_unlinked_waiter(id); *waiter_id = None; (Poll::Ready(()), None) } @@ -240,7 +240,7 @@ impl AutoResetEvent { // Recheck before enqueueing: set and cancellation handoff use this same lock. None if self.try_wait() => (Poll::Ready(()), None), None => { - *waiter_id = Some(state.waiters.push_back(Waiter::Waiting(cx.waker().clone()))); + *waiter_id = Some(waiters.push_back(Waiter::Waiting(cx.waker().clone()))); (Poll::Pending, None) } } @@ -251,12 +251,12 @@ impl AutoResetEvent { fn unregister_waiter(&self, id: WaiterId) { let (waiter, waker) = { - let mut state = self.state.lock(); + let mut waiters = self.waiters.lock(); // A selected waiter is already detached, but still owns its signal until removal. - state.waiters.unlink_waiter(id, |_| true); - let waiter = state.waiters.remove_unlinked_waiter(id); + waiters.unlink_waiter(id, |_| true); + let waiter = waiters.remove_unlinked_waiter(id); let waker = match &waiter { - Waiter::Notified => state.signal(&self.is_set), + Waiter::Notified => self.signal(&mut waiters), Waiter::Waiting(_) => None, }; (waiter, waker) @@ -266,6 +266,19 @@ impl AutoResetEvent { waker.wake(); } } + + fn signal(&self, waiters: &mut WaitList) -> Option { + if let Some((_, waiter)) = waiters.unlink_first_waiter(|_| true) { + let Waiter::Waiting(waker) = mem::replace(waiter, Waiter::Notified) else { + unreachable!("only unselected waits remain queued") + }; + Some(waker) + } else { + // Publish every set, including coalesced sets and returned assigned signals. + self.is_set.store(true, Ordering::Release); + None + } + } } impl Default for AutoResetEvent { @@ -283,27 +296,6 @@ impl fmt::Debug for AutoResetEvent { } } -struct State { - // A stored signal and queued (unselected) waits never coexist. Detached, selected waits can - // coexist with either: their signals are reserved until consumption or cancellation. - waiters: WaitList, -} - -impl State { - fn signal(&mut self, is_set: &AtomicBool) -> Option { - if let Some((_, waiter)) = self.waiters.unlink_first_waiter(|_| true) { - let Waiter::Waiting(waker) = mem::replace(waiter, Waiter::Notified) else { - unreachable!("only unselected waits remain queued") - }; - Some(waker) - } else { - // Publish every set, including coalesced sets and returned assigned signals. - is_set.store(true, Ordering::Release); - None - } - } -} - enum Waiter { Waiting(Waker), Notified, diff --git a/asyncband/src/event/manual_reset.rs b/asyncband/src/event/manual_reset.rs index 32689200..b365ebec 100644 --- a/asyncband/src/event/manual_reset.rs +++ b/asyncband/src/event/manual_reset.rs @@ -75,7 +75,7 @@ use crate::internal::waker_batch::WakerBatch; /// ``` pub struct ManualResetEvent { is_set: AtomicBool, - state: Mutex, + waiters: Mutex>, } impl ManualResetEvent { @@ -90,9 +90,7 @@ impl ManualResetEvent { pub const fn with_state(is_set: bool) -> Self { Self { is_set: AtomicBool::new(is_set), - state: Mutex::new(State { - waiters: WaitList::new(), - }), + waiters: Mutex::new(WaitList::new()), } } @@ -107,7 +105,7 @@ impl ManualResetEvent { /// attempted for every other selected task before the panic resumes. pub fn set(&self) { let wakers = { - let mut state = self.state.lock(); + let mut waiters = self.waiters.lock(); if self.is_set.load(Ordering::Relaxed) { return; } @@ -117,7 +115,7 @@ impl ManualResetEvent { // 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 }) { @@ -136,7 +134,7 @@ 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) { - let _state = self.state.lock(); + let _waiters = self.waiters.lock(); self.is_set.store(false, Ordering::Relaxed); } @@ -230,10 +228,10 @@ impl ManualResetEvent { } 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) } @@ -242,7 +240,7 @@ impl ManualResetEvent { !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) @@ -251,7 +249,7 @@ impl ManualResetEvent { // 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()), })); @@ -266,8 +264,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); } @@ -288,21 +288,6 @@ impl fmt::Debug for ManualResetEvent { } } -#[derive(Debug)] -struct State { - 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/auto_reset.rs b/benchmarks/asyncband/event/auto_reset.rs index f4fc1fb7..e4d08037 100644 --- a/benchmarks/asyncband/event/auto_reset.rs +++ b/benchmarks/asyncband/event/auto_reset.rs @@ -30,6 +30,13 @@ const WAITER_COUNTS: &[usize] = &[1, 8, 32]; const THREAD_COUNTS: &[usize] = &[1, 2, 8, 32]; const CONTENDED_SAMPLE_SIZE: u32 = 256; +#[divan::bench] +fn set_coalesced(bencher: Bencher) { + let event = AutoResetEvent::with_state(true); + + bencher.bench_local(|| black_box(&event).set()); +} + #[divan::bench] fn set_then_wait(bencher: Bencher) { let event = AutoResetEvent::new(); @@ -107,6 +114,25 @@ fn wake_waiter_reused(bencher: Bencher) { }); } +#[divan::bench] +fn consume_stale_signal_then_wait(bencher: Bencher) { + let mut context = bench_context(); + let event = AutoResetEvent::new(); + + bencher.bench_local(|| { + // A predicate loop consumes a leftover signal, rechecks, and waits for new work. + event.set(); + let mut stale = pin!(event.wait()); + poll_pinned_ready(stale.as_mut(), &mut context); + + let mut wait = pin!(event.wait()); + poll_pending(wait.as_mut(), &mut context); + event.set(); + poll_pinned_ready(wait.as_mut(), &mut context); + black_box(&event) + }); +} + #[divan::bench(args = WAITER_COUNTS)] fn wake_waiters(bencher: Bencher, waiter_count: usize) { let mut context = bench_context(); From 50e50ad5b9a50feb68a659eb386708e1b20971a2 Mon Sep 17 00:00:00 2001 From: tison Date: Tue, 15 Sep 2026 07:02:43 +0800 Subject: [PATCH 3/3] refactor(event): limit state fast paths to manual-reset events --- asyncband/src/event/auto_reset.rs | 91 +++++---- asyncband/src/event/manual_reset.rs | 4 +- benchmarks/asyncband/event/auto_reset.rs | 176 ------------------ benchmarks/asyncband/event/mod.rs | 1 - .../tests/auto_reset_event_test.rs | 53 +----- 5 files changed, 49 insertions(+), 276 deletions(-) delete mode 100644 benchmarks/asyncband/event/auto_reset.rs diff --git a/asyncband/src/event/auto_reset.rs b/asyncband/src/event/auto_reset.rs index 2a38b5d8..af72eaec 100644 --- a/asyncband/src/event/auto_reset.rs +++ b/asyncband/src/event/auto_reset.rs @@ -20,8 +20,6 @@ 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; @@ -78,10 +76,7 @@ use crate::internal::waitlist::WaiterId; /// # } /// ``` pub struct AutoResetEvent { - is_set: AtomicBool, - // A stored signal and queued (unselected) waits never coexist. Detached, selected waits can - // coexist with either: their signals are reserved until consumption or cancellation. - waiters: Mutex>, + state: Mutex, } impl AutoResetEvent { @@ -95,8 +90,10 @@ impl AutoResetEvent { /// If `is_set` is `true`, the event stores one signal for a future wait. pub const fn with_state(is_set: bool) -> Self { Self { - is_set: AtomicBool::new(is_set), - waiters: Mutex::new(WaitList::new()), + state: Mutex::new(State { + is_set, + waiters: WaitList::new(), + }), } } @@ -111,7 +108,7 @@ impl AutoResetEvent { /// Panics if waking a selected task panics. Its signal remains assigned and can still be /// consumed by polling that wait or passed on by dropping it. pub fn set(&self) { - let waker = self.signal(&mut self.waiters.lock()); + let waker = self.state.lock().signal(); if let Some(waker) = waker { waker.wake(); } @@ -122,8 +119,7 @@ impl AutoResetEvent { /// Signals already assigned to waits remain theirs. Cancelling such a wait can still transfer /// or restore its signal after this call. If no signal is stored, this has no effect. pub fn reset(&self) { - let _waiters = self.waiters.lock(); - self.is_set.store(false, Ordering::Relaxed); + self.state.lock().is_set = false; } /// Returns whether the event is currently set. @@ -144,7 +140,7 @@ impl AutoResetEvent { /// assert!(event.is_set()); /// ``` pub fn is_set(&self) -> bool { - self.is_set.load(Ordering::Acquire) + self.state.lock().is_set } /// Attempts to wait without registering a waiter. @@ -163,13 +159,7 @@ impl AutoResetEvent { /// assert!(!event.try_wait()); // A successful wait consumes the signal. /// ``` pub fn try_wait(&self) -> bool { - // Only consumption bypasses the waiter lock. Signals are stored under that lock only - // when no wait is queued, so this cannot take a signal assigned to a registered wait. - self.is_set.load(Ordering::Relaxed) - && self - .is_set - .compare_exchange(true, false, Ordering::Acquire, Ordering::Relaxed) - .is_ok() + mem::take(&mut self.state.lock().is_set) } /// Waits for and consumes one signal. @@ -217,17 +207,12 @@ impl AutoResetEvent { } fn poll_wait(&self, waiter_id: &mut Option, cx: &mut Context<'_>) -> Poll<()> { - // Registered waits own separate signals and must remove their nodes under the lock. - if waiter_id.is_none() && self.try_wait() { - return Poll::Ready(()); - } - let (poll, retired_waker) = { - let mut waiters = self.waiters.lock(); + let mut state = self.state.lock(); match *waiter_id { - Some(id) => match waiters.waiter_mut(id) { + Some(id) => match state.waiters.waiter_mut(id) { Waiter::Notified => { - waiters.remove_unlinked_waiter(id); + state.waiters.remove_unlinked_waiter(id); *waiter_id = None; (Poll::Ready(()), None) } @@ -237,10 +222,12 @@ impl AutoResetEvent { (Poll::Pending, retired) } }, - // Recheck before enqueueing: set and cancellation handoff use this same lock. - None if self.try_wait() => (Poll::Ready(()), None), + None if state.is_set => { + state.is_set = false; + (Poll::Ready(()), None) + } None => { - *waiter_id = Some(waiters.push_back(Waiter::Waiting(cx.waker().clone()))); + *waiter_id = Some(state.waiters.push_back(Waiter::Waiting(cx.waker().clone()))); (Poll::Pending, None) } } @@ -251,12 +238,12 @@ impl AutoResetEvent { fn unregister_waiter(&self, id: WaiterId) { let (waiter, waker) = { - let mut waiters = self.waiters.lock(); + let mut state = self.state.lock(); // A selected waiter is already detached, but still owns its signal until removal. - waiters.unlink_waiter(id, |_| true); - let waiter = waiters.remove_unlinked_waiter(id); + state.waiters.unlink_waiter(id, |_| true); + let waiter = state.waiters.remove_unlinked_waiter(id); let waker = match &waiter { - Waiter::Notified => self.signal(&mut waiters), + Waiter::Notified => state.signal(), Waiter::Waiting(_) => None, }; (waiter, waker) @@ -266,19 +253,6 @@ impl AutoResetEvent { waker.wake(); } } - - fn signal(&self, waiters: &mut WaitList) -> Option { - if let Some((_, waiter)) = waiters.unlink_first_waiter(|_| true) { - let Waiter::Waiting(waker) = mem::replace(waiter, Waiter::Notified) else { - unreachable!("only unselected waits remain queued") - }; - Some(waker) - } else { - // Publish every set, including coalesced sets and returned assigned signals. - self.is_set.store(true, Ordering::Release); - None - } - } } impl Default for AutoResetEvent { @@ -289,13 +263,34 @@ impl Default for AutoResetEvent { impl fmt::Debug for AutoResetEvent { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - let is_set = self.is_set(); + let is_set = self.state.lock().is_set; f.debug_struct("AutoResetEvent") .field("is_set", &is_set) .finish_non_exhaustive() } } +struct State { + // A stored signal and queued (unselected) waits never coexist. Detached, selected waits can + // coexist with either: their signals are reserved until consumption or cancellation. + is_set: bool, + waiters: WaitList, +} + +impl State { + fn signal(&mut self) -> Option { + if let Some((_, waiter)) = self.waiters.unlink_first_waiter(|_| true) { + let Waiter::Waiting(waker) = mem::replace(waiter, Waiter::Notified) else { + unreachable!("only unselected waits remain queued") + }; + Some(waker) + } else { + self.is_set = true; + None + } + } +} + enum Waiter { Waiting(Waker), Notified, diff --git a/asyncband/src/event/manual_reset.rs b/asyncband/src/event/manual_reset.rs index b365ebec..b7f3653e 100644 --- a/asyncband/src/event/manual_reset.rs +++ b/asyncband/src/event/manual_reset.rs @@ -74,6 +74,7 @@ use crate::internal::waker_batch::WakerBatch; /// # } /// ``` pub struct ManualResetEvent { + // Flag writes and waiter-list changes hold `waiters`. Lock-free paths only read the flag. is_set: AtomicBool, waiters: Mutex>, } @@ -110,7 +111,7 @@ impl ManualResetEvent { return; } - // Serialize publication, cohort detach, reset, and registration with the same lock. + // 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. @@ -135,6 +136,7 @@ impl ManualResetEvent { /// already unset, this has no effect. pub fn reset(&self) { let _waiters = self.waiters.lock(); + // Clearing the flag does not publish data to successful waits. self.is_set.store(false, Ordering::Relaxed); } diff --git a/benchmarks/asyncband/event/auto_reset.rs b/benchmarks/asyncband/event/auto_reset.rs deleted file mode 100644 index e4d08037..00000000 --- a/benchmarks/asyncband/event/auto_reset.rs +++ /dev/null @@ -1,176 +0,0 @@ -// Licensed to the Apache Software Foundation (ASF) under one -// or more contributor license agreements. See the NOTICE file -// distributed with this work for additional information -// regarding copyright ownership. The ASF licenses this file -// to you under the Apache License, Version 2.0 (the -// "License"); you may not use this file except in compliance -// with the License. You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, -// software distributed under the License is distributed on an -// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY -// KIND, either express or implied. See the License for the -// specific language governing permissions and limitations -// under the License. - -use std::pin::pin; - -use asyncband::event::AutoResetEvent; -use divan::Bencher; -use divan::black_box; -use divan::counter::ItemsCount; - -use crate::support::bench_context; -use crate::support::poll_pending; -use crate::support::poll_pinned_ready; - -const WAITER_COUNTS: &[usize] = &[1, 8, 32]; -const THREAD_COUNTS: &[usize] = &[1, 2, 8, 32]; -const CONTENDED_SAMPLE_SIZE: u32 = 256; - -#[divan::bench] -fn set_coalesced(bencher: Bencher) { - let event = AutoResetEvent::with_state(true); - - bencher.bench_local(|| black_box(&event).set()); -} - -#[divan::bench] -fn set_then_wait(bencher: Bencher) { - let event = AutoResetEvent::new(); - - bencher.bench_local(|| { - event.set(); - let mut context = bench_context(); - let mut wait = pin!(event.wait()); - poll_pinned_ready(wait.as_mut(), &mut context); - black_box(&event) - }); -} - -#[divan::bench(threads = THREAD_COUNTS, sample_size = CONTENDED_SAMPLE_SIZE)] -fn is_set_contended(bencher: Bencher) { - let event = AutoResetEvent::with_state(true); - - bencher.bench(|| black_box(event.is_set())); -} - -#[divan::bench] -fn set_reset_cycle(bencher: Bencher) { - let event = AutoResetEvent::new(); - - bencher.bench_local(|| { - event.set(); - event.reset(); - black_box(&event) - }); -} - -#[divan::bench] -fn cancel_pending(bencher: Bencher) { - let mut context = bench_context(); - - bencher.bench_local(|| { - let event = AutoResetEvent::new(); - { - let mut wait = pin!(event.wait()); - poll_pending(wait.as_mut(), &mut context); - } - black_box(event) - }); -} - -#[divan::bench] -fn wake_waiter(bencher: Bencher) { - let mut context = bench_context(); - - bencher.bench_local(|| { - let event = AutoResetEvent::new(); - { - let mut wait = pin!(event.wait()); - poll_pending(wait.as_mut(), &mut context); - - event.set(); - poll_pinned_ready(wait.as_mut(), &mut context); - } - black_box(event) - }); -} - -#[divan::bench] -fn wake_waiter_reused(bencher: Bencher) { - let mut context = bench_context(); - let event = AutoResetEvent::new(); - - bencher.bench_local(|| { - let mut wait = pin!(event.wait()); - poll_pending(wait.as_mut(), &mut context); - - event.set(); - poll_pinned_ready(wait.as_mut(), &mut context); - black_box(&event) - }); -} - -#[divan::bench] -fn consume_stale_signal_then_wait(bencher: Bencher) { - let mut context = bench_context(); - let event = AutoResetEvent::new(); - - bencher.bench_local(|| { - // A predicate loop consumes a leftover signal, rechecks, and waits for new work. - event.set(); - let mut stale = pin!(event.wait()); - poll_pinned_ready(stale.as_mut(), &mut context); - - let mut wait = pin!(event.wait()); - poll_pending(wait.as_mut(), &mut context); - event.set(); - poll_pinned_ready(wait.as_mut(), &mut context); - black_box(&event) - }); -} - -#[divan::bench(args = WAITER_COUNTS)] -fn wake_waiters(bencher: Bencher, waiter_count: usize) { - let mut context = bench_context(); - - bencher - .counter(ItemsCount::new(waiter_count)) - .bench_local(|| { - let event = AutoResetEvent::new(); - let mut waiters = (0..waiter_count) - .map(|_| Box::pin(event.wait())) - .collect::>(); - for waiter in &mut waiters { - poll_pending(waiter.as_mut(), &mut context); - } - - for _ in 0..waiter_count { - event.set(); - } - for mut waiter in waiters { - poll_pinned_ready(waiter.as_mut(), &mut context); - } - black_box(event) - }); -} - -#[divan::bench(threads = THREAD_COUNTS, sample_size = CONTENDED_SAMPLE_SIZE)] -fn try_wait_unset(bencher: Bencher) { - let event = AutoResetEvent::new(); - - bencher.bench(|| black_box(event.try_wait())); -} - -#[divan::bench] -fn set_then_try_wait(bencher: Bencher) { - let event = AutoResetEvent::new(); - - bencher.bench_local(|| { - event.set(); - assert!(event.try_wait()); - }); -} diff --git a/benchmarks/asyncband/event/mod.rs b/benchmarks/asyncband/event/mod.rs index f4f10365..35540691 100644 --- a/benchmarks/asyncband/event/mod.rs +++ b/benchmarks/asyncband/event/mod.rs @@ -15,5 +15,4 @@ // specific language governing permissions and limitations // under the License. -mod auto_reset; mod wait; diff --git a/tests-integration/tests/auto_reset_event_test.rs b/tests-integration/tests/auto_reset_event_test.rs index a7c3ef38..dc0e66b8 100644 --- a/tests-integration/tests/auto_reset_event_test.rs +++ b/tests-integration/tests/auto_reset_event_test.rs @@ -257,13 +257,13 @@ struct ReentrantWaker(Arc); impl Wake for ReentrantWaker { fn wake(self: Arc) { - self.0.reset(); + self.0.try_wait(); } } impl Drop for ReentrantWaker { fn drop(&mut self) { - self.0.reset(); + self.0.try_wait(); } } @@ -352,57 +352,10 @@ fn concurrent_sets_and_waits_publish_state_without_losing_signals() { // Registration races with set; the second barrier prevents the next set from // coalescing before this round's signal has been consumed. round.wait(); - if expected % 2 == 0 { - event.wait().block_on(); - } else { - while !event.try_wait() { - thread::yield_now(); - } - } + event.wait().block_on(); assert_eq!(value.load(Ordering::Relaxed), expected); round.wait(); } }); }); } - -#[test] -fn stored_signals_are_consumed_once_when_immediate_and_async_waits_race() { - assert_completes_without_deadlock(|| { - let event = AutoResetEvent::new(); - let round = Barrier::new(3); - let completed = AtomicUsize::new(0); - thread::scope(|scope| { - scope.spawn(|| { - for _ in 0..100 { - round.wait(); - if event.try_wait() { - completed.fetch_add(1, Ordering::Relaxed); - } - round.wait(); - } - }); - scope.spawn(|| { - for _ in 0..100 { - round.wait(); - { - let mut wait = pin!(event.wait()); - if poll_once(wait.as_mut()).is_ready() { - completed.fetch_add(1, Ordering::Relaxed); - } - } - // Remove any pending registration before the next signal is stored. - round.wait(); - } - }); - for _ in 0..100 { - completed.store(0, Ordering::Relaxed); - event.set(); - round.wait(); - round.wait(); - assert_eq!(completed.load(Ordering::Relaxed), 1); - assert!(!event.is_set()); - } - }); - }); -}