mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-26 00:00:16 +02:00
sync: reduce contention in Notify (#5503)
This commit is contained in:
@@ -45,6 +45,10 @@ name = "rt_multi_threaded"
|
||||
path = "rt_multi_threaded.rs"
|
||||
harness = false
|
||||
|
||||
[[bench]]
|
||||
name = "sync_notify"
|
||||
path = "sync_notify.rs"
|
||||
harness = false
|
||||
|
||||
[[bench]]
|
||||
name = "sync_rwlock"
|
||||
|
||||
@@ -0,0 +1,90 @@
|
||||
use bencher::Bencher;
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
use std::sync::Arc;
|
||||
|
||||
use tokio::sync::Notify;
|
||||
|
||||
fn rt() -> tokio::runtime::Runtime {
|
||||
tokio::runtime::Builder::new_multi_thread()
|
||||
.worker_threads(6)
|
||||
.build()
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
fn notify_waiters<const N_WAITERS: usize>(b: &mut Bencher) {
|
||||
let rt = rt();
|
||||
let notify = Arc::new(Notify::new());
|
||||
let counter = Arc::new(AtomicUsize::new(0));
|
||||
for _ in 0..N_WAITERS {
|
||||
rt.spawn({
|
||||
let notify = notify.clone();
|
||||
let counter = counter.clone();
|
||||
async move {
|
||||
loop {
|
||||
notify.notified().await;
|
||||
counter.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
const N_ITERS: usize = 500;
|
||||
b.iter(|| {
|
||||
counter.store(0, Ordering::Relaxed);
|
||||
loop {
|
||||
notify.notify_waiters();
|
||||
if counter.load(Ordering::Relaxed) >= N_ITERS {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
fn notify_one<const N_WAITERS: usize>(b: &mut Bencher) {
|
||||
let rt = rt();
|
||||
let notify = Arc::new(Notify::new());
|
||||
let counter = Arc::new(AtomicUsize::new(0));
|
||||
for _ in 0..N_WAITERS {
|
||||
rt.spawn({
|
||||
let notify = notify.clone();
|
||||
let counter = counter.clone();
|
||||
async move {
|
||||
loop {
|
||||
notify.notified().await;
|
||||
counter.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
const N_ITERS: usize = 500;
|
||||
b.iter(|| {
|
||||
counter.store(0, Ordering::Relaxed);
|
||||
loop {
|
||||
notify.notify_one();
|
||||
if counter.load(Ordering::Relaxed) >= N_ITERS {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
bencher::benchmark_group!(
|
||||
notify_waiters_simple,
|
||||
notify_waiters::<10>,
|
||||
notify_waiters::<50>,
|
||||
notify_waiters::<100>,
|
||||
notify_waiters::<200>,
|
||||
notify_waiters::<500>
|
||||
);
|
||||
|
||||
bencher::benchmark_group!(
|
||||
notify_one_simple,
|
||||
notify_one::<10>,
|
||||
notify_one::<50>,
|
||||
notify_one::<100>,
|
||||
notify_one::<200>,
|
||||
notify_one::<500>
|
||||
);
|
||||
|
||||
bencher::benchmark_main!(notify_waiters_simple, notify_one_simple);
|
||||
+151
-81
@@ -5,18 +5,18 @@
|
||||
// triggers this warning but it is safe to ignore in this case.
|
||||
#![cfg_attr(not(feature = "sync"), allow(unreachable_pub, dead_code))]
|
||||
|
||||
use crate::loom::cell::UnsafeCell;
|
||||
use crate::loom::sync::atomic::AtomicUsize;
|
||||
use crate::loom::sync::Mutex;
|
||||
use crate::util::linked_list::{self, GuardedLinkedList, LinkedList};
|
||||
use crate::util::WakeList;
|
||||
|
||||
use std::cell::UnsafeCell;
|
||||
use std::future::Future;
|
||||
use std::marker::PhantomPinned;
|
||||
use std::panic::{RefUnwindSafe, UnwindSafe};
|
||||
use std::pin::Pin;
|
||||
use std::ptr::NonNull;
|
||||
use std::sync::atomic::Ordering::SeqCst;
|
||||
use std::sync::atomic::Ordering::{self, Acquire, Relaxed, Release, SeqCst};
|
||||
use std::task::{Context, Poll, Waker};
|
||||
|
||||
type WaitList = LinkedList<Waiter, <Waiter as linked_list::Link>::Target>;
|
||||
@@ -213,24 +213,21 @@ pub struct Notify {
|
||||
waiters: Mutex<WaitList>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
enum NotificationType {
|
||||
// Notification triggered by calling `notify_waiters`
|
||||
AllWaiters,
|
||||
// Notification triggered by calling `notify_one`
|
||||
OneWaiter,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct Waiter {
|
||||
/// Intrusive linked-list pointers.
|
||||
pointers: linked_list::Pointers<Waiter>,
|
||||
|
||||
/// Waiting task's waker.
|
||||
waker: Option<Waker>,
|
||||
/// Waiting task's waker. Depending on the value of `notification`,
|
||||
/// this field is either protected by the `waiters` lock in
|
||||
/// `Notify`, or it is exclusively owned by the enclosing `Waiter`.
|
||||
waker: UnsafeCell<Option<Waker>>,
|
||||
|
||||
/// `true` if the notification has been assigned to this waiter.
|
||||
notified: Option<NotificationType>,
|
||||
/// Notification for this waiter.
|
||||
/// * if it's `None`, then `waker` is protected by the `waiters` lock.
|
||||
/// * if it's `Some`, then `waker` is exclusively owned by the
|
||||
/// enclosing `Waiter` and can be accessed without locking.
|
||||
notification: AtomicNotification,
|
||||
|
||||
/// Should not be `Unpin`.
|
||||
_p: PhantomPinned,
|
||||
@@ -240,8 +237,8 @@ impl Waiter {
|
||||
fn new() -> Waiter {
|
||||
Waiter {
|
||||
pointers: linked_list::Pointers::new(),
|
||||
waker: None,
|
||||
notified: None,
|
||||
waker: UnsafeCell::new(None),
|
||||
notification: AtomicNotification::none(),
|
||||
_p: PhantomPinned,
|
||||
}
|
||||
}
|
||||
@@ -255,6 +252,57 @@ generate_addr_of_methods! {
|
||||
}
|
||||
}
|
||||
|
||||
// No notification.
|
||||
const NOTIFICATION_NONE: usize = 0;
|
||||
|
||||
// Notification type used by `notify_one`.
|
||||
const NOTIFICATION_ONE: usize = 1;
|
||||
|
||||
// Notification type used by `notify_waiters`.
|
||||
const NOTIFICATION_ALL: usize = 2;
|
||||
|
||||
/// Notification for a `Waiter`.
|
||||
/// This struct is equivalent to `Option<Notification>`, but uses
|
||||
/// `AtomicUsize` inside for atomic operations.
|
||||
#[derive(Debug)]
|
||||
struct AtomicNotification(AtomicUsize);
|
||||
|
||||
impl AtomicNotification {
|
||||
fn none() -> Self {
|
||||
AtomicNotification(AtomicUsize::new(NOTIFICATION_NONE))
|
||||
}
|
||||
|
||||
/// Store-release a notification.
|
||||
/// This method should be called exactly once.
|
||||
fn store_release(&self, notification: Notification) {
|
||||
self.0.store(notification as usize, Release);
|
||||
}
|
||||
|
||||
fn load(&self, ordering: Ordering) -> Option<Notification> {
|
||||
match self.0.load(ordering) {
|
||||
NOTIFICATION_NONE => None,
|
||||
NOTIFICATION_ONE => Some(Notification::One),
|
||||
NOTIFICATION_ALL => Some(Notification::All),
|
||||
_ => unreachable!(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Clears the notification.
|
||||
/// This method is used by a `Notified` future to consume the
|
||||
/// notification. It uses relaxed ordering and should be only
|
||||
/// used once the atomic notification is no longer shared.
|
||||
fn clear(&self) {
|
||||
self.0.store(NOTIFICATION_NONE, Relaxed);
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
#[repr(usize)]
|
||||
enum Notification {
|
||||
One = NOTIFICATION_ONE,
|
||||
All = NOTIFICATION_ALL,
|
||||
}
|
||||
|
||||
/// List used in `Notify::notify_waiters`. It wraps a guarded linked list
|
||||
/// and gates the access to it on `notify.waiters` mutex. It also empties
|
||||
/// the list on drop.
|
||||
@@ -267,11 +315,10 @@ struct NotifyWaitersList<'a> {
|
||||
impl<'a> NotifyWaitersList<'a> {
|
||||
fn new(
|
||||
unguarded_list: WaitList,
|
||||
guard: Pin<&'a mut UnsafeCell<Waiter>>,
|
||||
guard: Pin<&'a Waiter>,
|
||||
notify: &'a Notify,
|
||||
) -> NotifyWaitersList<'a> {
|
||||
// Safety: pointer to the guarding waiter is not null.
|
||||
let guard_ptr = unsafe { NonNull::new_unchecked(guard.get()) };
|
||||
let guard_ptr = NonNull::from(guard.get_ref());
|
||||
let list = unguarded_list.into_guarded(guard_ptr);
|
||||
NotifyWaitersList {
|
||||
list,
|
||||
@@ -299,10 +346,10 @@ impl Drop for NotifyWaitersList<'_> {
|
||||
// We do not wake the waiters to avoid double panics.
|
||||
if !self.is_empty {
|
||||
let _lock_guard = self.notify.waiters.lock();
|
||||
while let Some(mut waiter) = self.list.pop_back() {
|
||||
// Safety: we hold the lock.
|
||||
let waiter = unsafe { waiter.as_mut() };
|
||||
waiter.notified = Some(NotificationType::AllWaiters);
|
||||
while let Some(waiter) = self.list.pop_back() {
|
||||
// Safety: we never make mutable references to waiters.
|
||||
let waiter = unsafe { waiter.as_ref() };
|
||||
waiter.notification.store_release(Notification::All);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -324,7 +371,7 @@ pub struct Notified<'a> {
|
||||
notify_waiters_calls: usize,
|
||||
|
||||
/// Entry in the waiter `LinkedList`.
|
||||
waiter: UnsafeCell<Waiter>,
|
||||
waiter: Waiter,
|
||||
}
|
||||
|
||||
unsafe impl<'a> Send for Notified<'a> {}
|
||||
@@ -463,7 +510,7 @@ impl Notify {
|
||||
notify: self,
|
||||
state: State::Init,
|
||||
notify_waiters_calls: get_num_notify_waiters_calls(state),
|
||||
waiter: UnsafeCell::new(Waiter::new()),
|
||||
waiter: Waiter::new(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -590,7 +637,7 @@ impl Notify {
|
||||
|
||||
// It is critical for `GuardedLinkedList` safety that the guard node is
|
||||
// pinned in memory and is not dropped until the guarded list is dropped.
|
||||
let guard = UnsafeCell::new(Waiter::new());
|
||||
let guard = Waiter::new();
|
||||
pin!(guard);
|
||||
|
||||
// We move all waiters to a secondary list. It uses a `GuardedLinkedList`
|
||||
@@ -601,23 +648,25 @@ impl Notify {
|
||||
// * This wrapper will empty the list on drop. It is critical for safety
|
||||
// that we will not leave any list entry with a pointer to the local
|
||||
// guard node after this function returns / panics.
|
||||
let mut list = NotifyWaitersList::new(std::mem::take(&mut *waiters), guard, self);
|
||||
let mut list = NotifyWaitersList::new(std::mem::take(&mut *waiters), guard.as_ref(), self);
|
||||
|
||||
let mut wakers = WakeList::new();
|
||||
'outer: loop {
|
||||
while wakers.can_push() {
|
||||
match list.pop_back_locked(&mut waiters) {
|
||||
Some(mut waiter) => {
|
||||
// Safety: `waiters` lock is still held.
|
||||
let waiter = unsafe { waiter.as_mut() };
|
||||
Some(waiter) => {
|
||||
// Safety: we never make mutable references to waiters.
|
||||
let waiter = unsafe { waiter.as_ref() };
|
||||
|
||||
assert!(waiter.notified.is_none());
|
||||
|
||||
waiter.notified = Some(NotificationType::AllWaiters);
|
||||
|
||||
if let Some(waker) = waiter.waker.take() {
|
||||
// Safety: we hold the lock, so we can access the waker.
|
||||
if let Some(waker) =
|
||||
unsafe { waiter.waker.with_mut(|waker| (*waker).take()) }
|
||||
{
|
||||
wakers.push(waker);
|
||||
}
|
||||
|
||||
// This waiter is unlinked and will not be shared ever again, release it.
|
||||
waiter.notification.store_release(Notification::All);
|
||||
}
|
||||
None => {
|
||||
break 'outer;
|
||||
@@ -674,15 +723,16 @@ fn notify_locked(waiters: &mut WaitList, state: &AtomicUsize, curr: usize) -> Op
|
||||
// transition **out** of `WAITING`.
|
||||
//
|
||||
// Get a pending waiter
|
||||
let mut waiter = waiters.pop_back().unwrap();
|
||||
let waiter = waiters.pop_back().unwrap();
|
||||
|
||||
// Safety: `waiters` lock is still held.
|
||||
let waiter = unsafe { waiter.as_mut() };
|
||||
// Safety: we never make mutable references to waiters.
|
||||
let waiter = unsafe { waiter.as_ref() };
|
||||
|
||||
assert!(waiter.notified.is_none());
|
||||
// Safety: we hold the lock, so we can access the waker.
|
||||
let waker = unsafe { waiter.waker.with_mut(|waker| (*waker).take()) };
|
||||
|
||||
waiter.notified = Some(NotificationType::OneWaiter);
|
||||
let waker = waiter.waker.take();
|
||||
// This waiter is unlinked and will not be shared ever again, release it.
|
||||
waiter.notification.store_release(Notification::One);
|
||||
|
||||
if waiters.is_empty() {
|
||||
// As this the **final** waiter in the list, the state
|
||||
@@ -812,12 +862,12 @@ impl Notified<'_> {
|
||||
|
||||
/// A custom `project` implementation is used in place of `pin-project-lite`
|
||||
/// as a custom drop implementation is needed.
|
||||
fn project(self: Pin<&mut Self>) -> (&Notify, &mut State, &usize, &UnsafeCell<Waiter>) {
|
||||
fn project(self: Pin<&mut Self>) -> (&Notify, &mut State, &usize, &Waiter) {
|
||||
unsafe {
|
||||
// Safety: `notify`, `state` and `notify_waiters_calls` are `Unpin`.
|
||||
|
||||
is_unpin::<&Notify>();
|
||||
is_unpin::<AtomicUsize>();
|
||||
is_unpin::<State>();
|
||||
is_unpin::<usize>();
|
||||
|
||||
let me = self.get_unchecked_mut();
|
||||
@@ -924,14 +974,13 @@ impl Notified<'_> {
|
||||
// The use of `old_waiter` here is not necessary, as the field is always
|
||||
// None when we reach this line.
|
||||
unsafe {
|
||||
old_waker = std::mem::replace(&mut (*waiter.get()).waker, waker);
|
||||
old_waker =
|
||||
waiter.waker.with_mut(|v| std::mem::replace(&mut *v, waker));
|
||||
}
|
||||
}
|
||||
|
||||
// Insert the waiter into the linked list
|
||||
//
|
||||
// safety: pointers from `UnsafeCell` are never null.
|
||||
waiters.push_front(unsafe { NonNull::new_unchecked(waiter.get()) });
|
||||
waiters.push_front(NonNull::from(waiter));
|
||||
|
||||
*state = Waiting;
|
||||
|
||||
@@ -941,26 +990,45 @@ impl Notified<'_> {
|
||||
return Poll::Pending;
|
||||
}
|
||||
Waiting => {
|
||||
// Currently in the "Waiting" state, implying the caller has a waiter stored in
|
||||
// a waiter list (guarded by `notify.waiters`). In order to access the waker
|
||||
if waiter.notification.load(Acquire).is_some() {
|
||||
// Safety: waiter is already unlinked and will not be shared again,
|
||||
// so we have an exclusive access to `waker`.
|
||||
drop(unsafe { waiter.waker.with_mut(|waker| (*waker).take()) });
|
||||
|
||||
waiter.notification.clear();
|
||||
*state = Done;
|
||||
return Poll::Ready(());
|
||||
}
|
||||
|
||||
// Our waiter was not notified, implying it is still stored in a waiter
|
||||
// list (guarded by `notify.waiters`). In order to access the waker
|
||||
// fields, we must acquire the lock.
|
||||
|
||||
let mut old_waker = None;
|
||||
let mut waiters = notify.waiters.lock();
|
||||
|
||||
// We hold the lock and notifications are set only with the lock held,
|
||||
// so this can be relaxed, because the happens-before relationship is
|
||||
// established through the mutex.
|
||||
if waiter.notification.load(Relaxed).is_some() {
|
||||
// Safety: waiter is already unlinked and will not be shared again,
|
||||
// so we have an exclusive access to `waker`.
|
||||
old_waker = unsafe { waiter.waker.with_mut(|waker| (*waker).take()) };
|
||||
|
||||
waiter.notification.clear();
|
||||
|
||||
// Drop the old waker after releasing the lock.
|
||||
drop(waiters);
|
||||
drop(old_waker);
|
||||
|
||||
*state = Done;
|
||||
return Poll::Ready(());
|
||||
}
|
||||
|
||||
// Load the state with the lock held.
|
||||
let curr = notify.state.load(SeqCst);
|
||||
|
||||
// Safety: called while locked
|
||||
let w = unsafe { &mut *waiter.get() };
|
||||
let mut old_waker = None;
|
||||
|
||||
if w.notified.is_some() {
|
||||
// Our waker has been notified and our waiter is already removed from
|
||||
// the list. Reset the notification and convert to `Done`.
|
||||
old_waker = std::mem::take(&mut w.waker);
|
||||
w.notified = None;
|
||||
*state = Done;
|
||||
} else if get_num_notify_waiters_calls(curr) != *notify_waiters_calls {
|
||||
if get_num_notify_waiters_calls(curr) != *notify_waiters_calls {
|
||||
// Before we add a waiter to the list we check if these numbers are
|
||||
// different while holding the lock. If these numbers are different now,
|
||||
// it means that there is a call to `notify_waiters` in progress and this
|
||||
@@ -968,23 +1036,28 @@ impl Notified<'_> {
|
||||
// We can treat the waiter as notified and remove it from the list, as
|
||||
// it would have been notified in the `notify_waiters` call anyways.
|
||||
|
||||
old_waker = std::mem::take(&mut w.waker);
|
||||
// Safety: we hold the lock, so we can modify the waker.
|
||||
old_waker = unsafe { waiter.waker.with_mut(|waker| (*waker).take()) };
|
||||
|
||||
// Safety: we hold the lock, so we have an exclusive access to the list.
|
||||
// The list is used in `notify_waiters`, so it must be guarded.
|
||||
unsafe { waiters.remove(NonNull::new_unchecked(w)) };
|
||||
unsafe { waiters.remove(NonNull::from(waiter)) };
|
||||
|
||||
*state = Done;
|
||||
} else {
|
||||
// Update the waker, if necessary.
|
||||
if let Some(waker) = waker {
|
||||
let should_update = match w.waker.as_ref() {
|
||||
Some(current_waker) => !current_waker.will_wake(waker),
|
||||
None => true,
|
||||
};
|
||||
if should_update {
|
||||
old_waker = std::mem::replace(&mut w.waker, Some(waker.clone()));
|
||||
}
|
||||
// Safety: we hold the lock, so we can modify the waker.
|
||||
unsafe {
|
||||
waiter.waker.with_mut(|v| {
|
||||
if let Some(waker) = waker {
|
||||
let should_update = match &*v {
|
||||
Some(current_waker) => !current_waker.will_wake(waker),
|
||||
None => true,
|
||||
};
|
||||
if should_update {
|
||||
old_waker = std::mem::replace(&mut *v, Some(waker.clone()));
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
// Drop the old waker after releasing the lock.
|
||||
@@ -1034,13 +1107,16 @@ impl Drop for Notified<'_> {
|
||||
let mut waiters = notify.waiters.lock();
|
||||
let mut notify_state = notify.state.load(SeqCst);
|
||||
|
||||
// We hold the lock, so this field is not concurrently accessed by
|
||||
// `notify_*` functions and we can use the relaxed ordering.
|
||||
let notification = waiter.notification.load(Relaxed);
|
||||
|
||||
// remove the entry from the list (if not already removed)
|
||||
//
|
||||
// Safety: we hold the lock, so we have an exclusive access to every list the
|
||||
// waiter may be contained in. If the node is not contained in the `waiters`
|
||||
// list, then it is contained by a guarded list used by `notify_waiters` and
|
||||
// in such case it must be a middle node.
|
||||
unsafe { waiters.remove(NonNull::new_unchecked(waiter.get())) };
|
||||
// list, then it is contained by a guarded list used by `notify_waiters`.
|
||||
unsafe { waiters.remove(NonNull::from(waiter)) };
|
||||
|
||||
if waiters.is_empty() && get_state(notify_state) == WAITING {
|
||||
notify_state = set_state(notify_state, EMPTY);
|
||||
@@ -1050,13 +1126,7 @@ impl Drop for Notified<'_> {
|
||||
// See if the node was notified but not received. In this case, if
|
||||
// the notification was triggered via `notify_one`, it must be sent
|
||||
// to the next waiter.
|
||||
//
|
||||
// Safety: with the entry removed from the linked list, there can be
|
||||
// no concurrent access to the entry
|
||||
if matches!(
|
||||
unsafe { (*waiter.get()).notified },
|
||||
Some(NotificationType::OneWaiter)
|
||||
) {
|
||||
if notification == Some(Notification::One) {
|
||||
if let Some(waker) = notify_locked(&mut waiters, ¬ify.state, notify_state) {
|
||||
drop(waiters);
|
||||
waker.wake();
|
||||
|
||||
Reference in New Issue
Block a user