Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
72 changes: 35 additions & 37 deletions asyncband/src/event/manual_reset.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -72,7 +74,9 @@ use crate::internal::waker_batch::WakerBatch;
/// # }
/// ```
pub struct ManualResetEvent {
state: Mutex<State>,
// Flag writes and waiter-list changes hold `waiters`. Lock-free paths only read the flag.
is_set: AtomicBool,
waiters: Mutex<WaitList<Waiter>>,
}

impl ManualResetEvent {
Expand All @@ -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()),
}
}

Expand All @@ -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
}) {
Expand All @@ -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.
Expand All @@ -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.
Expand Down Expand Up @@ -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<WaiterId>, 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()),
}));
Expand All @@ -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);
}
Expand All @@ -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<Waiter>,
}

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,
Expand Down
7 changes: 7 additions & 0 deletions benchmarks/asyncband/event/wait.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()));
}
41 changes: 39 additions & 2 deletions tests-integration/tests/event_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -173,7 +177,7 @@ impl Wake for ReentrantWaker {

impl Drop for ReentrantWaker {
fn drop(&mut self) {
self.0.is_set();
self.0.reset();
}
}

Expand Down Expand Up @@ -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);
});
Expand Down Expand Up @@ -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();
}
});
});
}