From 1ae9434e8e4a419ce25644e6c8d2b2e2e8c34750 Mon Sep 17 00:00:00 2001 From: Carl Lerche Date: Mon, 5 May 2025 10:48:02 -0700 Subject: [PATCH] time: revert "use sharding for timer implementation" related changes (#7226) The work on sharding the timer implementation has caused a measurable performance regression due to increased contention. This patch reverts the current work on sharding. The next step will be to work on a per-worker timer wheel. --- tokio/src/loom/mocked.rs | 5 - tokio/src/loom/std/mutex.rs | 8 - tokio/src/runtime/builder.rs | 7 +- tokio/src/runtime/context.rs | 12 +- tokio/src/runtime/driver.rs | 8 +- .../runtime/scheduler/multi_thread/worker.rs | 5 - tokio/src/runtime/time/entry.rs | 35 +--- tokio/src/runtime/time/mod.rs | 183 +++++------------- tokio/src/runtime/time/tests/mod.rs | 16 +- tokio/src/util/mod.rs | 2 +- tokio/src/util/rand.rs | 1 - 11 files changed, 70 insertions(+), 212 deletions(-) diff --git a/tokio/src/loom/mocked.rs b/tokio/src/loom/mocked.rs index cfcbb2967..cd996bc97 100644 --- a/tokio/src/loom/mocked.rs +++ b/tokio/src/loom/mocked.rs @@ -24,11 +24,6 @@ pub(crate) mod sync { pub(crate) fn try_lock(&self) -> Option> { self.0.try_lock().ok() } - - #[inline] - pub(crate) fn get_mut(&mut self) -> &mut T { - self.0.get_mut().unwrap() - } } #[derive(Debug)] diff --git a/tokio/src/loom/std/mutex.rs b/tokio/src/loom/std/mutex.rs index 9593ec487..95f6d73ba 100644 --- a/tokio/src/loom/std/mutex.rs +++ b/tokio/src/loom/std/mutex.rs @@ -33,12 +33,4 @@ impl Mutex { Err(TryLockError::WouldBlock) => None, } } - - #[inline] - pub(crate) fn get_mut(&mut self) -> &mut T { - match self.0.get_mut() { - Ok(val) => val, - Err(p_err) => p_err.into_inner(), - } - } } diff --git a/tokio/src/runtime/builder.rs b/tokio/src/runtime/builder.rs index 994fcfa5c..47ba18c96 100644 --- a/tokio/src/runtime/builder.rs +++ b/tokio/src/runtime/builder.rs @@ -924,7 +924,7 @@ impl Builder { } } - fn get_cfg(&self, workers: usize) -> driver::Cfg { + fn get_cfg(&self) -> driver::Cfg { driver::Cfg { enable_pause_time: match self.kind { Kind::CurrentThread => true, @@ -935,7 +935,6 @@ impl Builder { enable_time: self.enable_time, start_paused: self.start_paused, nevents: self.nevents, - workers, } } @@ -1453,7 +1452,7 @@ impl Builder { use crate::runtime::scheduler; use crate::runtime::Config; - let (driver, driver_handle) = driver::Driver::new(self.get_cfg(1))?; + let (driver, driver_handle) = driver::Driver::new(self.get_cfg())?; // Blocking pool let blocking_pool = blocking::create_blocking_pool(self, self.max_blocking_threads); @@ -1608,7 +1607,7 @@ cfg_rt_multi_thread! { let worker_threads = self.worker_threads.unwrap_or_else(num_cpus); - let (driver, driver_handle) = driver::Driver::new(self.get_cfg(worker_threads))?; + let (driver, driver_handle) = driver::Driver::new(self.get_cfg())?; // Create the blocking pool let blocking_pool = diff --git a/tokio/src/runtime/context.rs b/tokio/src/runtime/context.rs index e8f17bb37..0d54f6ca5 100644 --- a/tokio/src/runtime/context.rs +++ b/tokio/src/runtime/context.rs @@ -3,7 +3,7 @@ use crate::task::coop; use std::cell::Cell; -#[cfg(any(feature = "rt", feature = "macros", feature = "time"))] +#[cfg(any(feature = "rt", feature = "macros"))] use crate::util::rand::FastRand; cfg_rt! { @@ -57,7 +57,7 @@ struct Context { #[cfg(feature = "rt")] runtime: Cell, - #[cfg(any(feature = "rt", feature = "macros", feature = "time"))] + #[cfg(any(feature = "rt", feature = "macros"))] rng: Cell>, /// Tracks the amount of "work" a task may still do before yielding back to @@ -100,7 +100,7 @@ tokio_thread_local! { #[cfg(feature = "rt")] runtime: Cell::new(EnterRuntime::NotEntered), - #[cfg(any(feature = "rt", feature = "macros", feature = "time"))] + #[cfg(any(feature = "rt", feature = "macros"))] rng: Cell::new(None), budget: Cell::new(coop::Budget::unconstrained()), @@ -121,11 +121,7 @@ tokio_thread_local! { } } -#[cfg(any( - feature = "time", - feature = "macros", - all(feature = "sync", feature = "rt") -))] +#[cfg(any(feature = "macros", all(feature = "sync", feature = "rt")))] pub(crate) fn thread_rng_n(n: u32) -> u32 { CONTEXT.with(|ctx| { let mut rng = ctx.rng.get().unwrap_or_else(FastRand::new); diff --git a/tokio/src/runtime/driver.rs b/tokio/src/runtime/driver.rs index 4ff2bec29..3b84a8669 100644 --- a/tokio/src/runtime/driver.rs +++ b/tokio/src/runtime/driver.rs @@ -40,7 +40,6 @@ pub(crate) struct Cfg { pub(crate) enable_pause_time: bool, pub(crate) start_paused: bool, pub(crate) nevents: usize, - pub(crate) workers: usize, } impl Driver { @@ -49,8 +48,7 @@ impl Driver { let clock = create_clock(cfg.enable_pause_time, cfg.start_paused); - let (time_driver, time_handle) = - create_time_driver(cfg.enable_time, io_stack, &clock, cfg.workers); + let (time_driver, time_handle) = create_time_driver(cfg.enable_time, io_stack, &clock); Ok(( Self { inner: time_driver }, @@ -297,10 +295,9 @@ cfg_time! { enable: bool, io_stack: IoStack, clock: &Clock, - workers: usize, ) -> (TimeDriver, TimeHandle) { if enable { - let (driver, handle) = crate::runtime::time::Driver::new(io_stack, clock, workers as u32); + let (driver, handle) = crate::runtime::time::Driver::new(io_stack, clock); (TimeDriver::Enabled { driver }, Some(handle)) } else { @@ -346,7 +343,6 @@ cfg_not_time! { _enable: bool, io_stack: IoStack, _clock: &Clock, - _workers: usize, ) -> (TimeDriver, TimeHandle) { (io_stack, ()) } diff --git a/tokio/src/runtime/scheduler/multi_thread/worker.rs b/tokio/src/runtime/scheduler/multi_thread/worker.rs index e33b9baea..73f033720 100644 --- a/tokio/src/runtime/scheduler/multi_thread/worker.rs +++ b/tokio/src/runtime/scheduler/multi_thread/worker.rs @@ -790,11 +790,6 @@ impl Context { self.defer.defer(waker); } } - - #[allow(dead_code)] - pub(crate) fn get_worker_index(&self) -> usize { - self.worker.index - } } impl Core { diff --git a/tokio/src/runtime/time/entry.rs b/tokio/src/runtime/time/entry.rs index 6f29f2901..7991ee0dc 100644 --- a/tokio/src/runtime/time/entry.rs +++ b/tokio/src/runtime/time/entry.rs @@ -58,7 +58,6 @@ use crate::loom::cell::UnsafeCell; use crate::loom::sync::atomic::AtomicU64; use crate::loom::sync::atomic::Ordering; -use crate::runtime::context; use crate::runtime::scheduler; use crate::sync::AtomicWaker; use crate::time::Instant; @@ -329,8 +328,6 @@ pub(super) type EntryList = crate::util::linked_list::LinkedList Self { + pub(super) fn new() -> Self { Self { - shard_id, cached_when: AtomicU64::new(0), pointers: linked_list::Pointers::new(), state: StateCell::default(), @@ -442,11 +438,6 @@ impl TimerShared { pub(super) fn might_be_registered(&self) -> bool { self.state.might_be_registered() } - - /// Gets the shard id. - pub(super) fn shard_id(&self) -> u32 { - self.shard_id - } } unsafe impl linked_list::Link for TimerShared { @@ -494,10 +485,8 @@ impl TimerEntry { fn inner(&self) -> &TimerShared { let inner = unsafe { &*self.inner.get() }; if inner.is_none() { - let shard_size = self.driver.driver().time().inner.get_shard_size(); - let shard_id = generate_shard_id(shard_size); unsafe { - *self.inner.get() = Some(TimerShared::new(shard_id)); + *self.inner.get() = Some(TimerShared::new()); } } return inner.as_ref().unwrap(); @@ -654,23 +643,3 @@ impl Drop for TimerEntry { unsafe { Pin::new_unchecked(self) }.as_mut().cancel(); } } - -// Generates a shard id. If current thread is a worker thread, we use its worker index as a shard id. -// Otherwise, we use a random number generator to obtain the shard id. -cfg_rt! { - fn generate_shard_id(shard_size: u32) -> u32 { - let id = context::with_scheduler(|ctx| match ctx { - Some(scheduler::Context::CurrentThread(_ctx)) => 0, - #[cfg(feature = "rt-multi-thread")] - Some(scheduler::Context::MultiThread(ctx)) => ctx.get_worker_index() as u32, - None => context::thread_rng_n(shard_size), - }); - id % shard_size - } -} - -cfg_not_rt! { - fn generate_shard_id(shard_size: u32) -> u32 { - context::thread_rng_n(shard_size) - } -} diff --git a/tokio/src/runtime/time/mod.rs b/tokio/src/runtime/time/mod.rs index 56e0ba64d..8cd51c5cb 100644 --- a/tokio/src/runtime/time/mod.rs +++ b/tokio/src/runtime/time/mod.rs @@ -12,7 +12,6 @@ use entry::{EntryList, TimerHandle, TimerShared, MAX_SAFE_MILLIS_DURATION}; mod handle; pub(crate) use self::handle::Handle; -use self::wheel::Wheel; mod source; pub(crate) use source::TimeSource; @@ -20,34 +19,15 @@ pub(crate) use source::TimeSource; mod wheel; use crate::loom::sync::atomic::{AtomicBool, Ordering}; -use crate::loom::sync::{Mutex, RwLock}; +use crate::loom::sync::Mutex; use crate::runtime::driver::{self, IoHandle, IoStack}; use crate::time::error::Error; use crate::time::{Clock, Duration}; use crate::util::WakeList; -use crate::loom::sync::atomic::AtomicU64; use std::fmt; use std::{num::NonZeroU64, ptr::NonNull}; -struct AtomicOptionNonZeroU64(AtomicU64); - -// A helper type to store the `next_wake`. -impl AtomicOptionNonZeroU64 { - fn new(val: Option) -> Self { - Self(AtomicU64::new(val.map_or(0, NonZeroU64::get))) - } - - fn store(&self, val: Option) { - self.0 - .store(val.map_or(0, NonZeroU64::get), Ordering::Relaxed); - } - - fn load(&self) -> Option { - NonZeroU64::new(self.0.load(Ordering::Relaxed)) - } -} - /// Time implementation that drives [`Sleep`][sleep], [`Interval`][interval], and [`Timeout`][timeout]. /// /// A `Driver` instance tracks the state necessary for managing time and @@ -111,14 +91,8 @@ pub(crate) struct Driver { /// Timer state shared between `Driver`, `Handle`, and `Registration`. struct Inner { - /// The earliest time at which we promise to wake up without unparking. - next_wake: AtomicOptionNonZeroU64, - - /// Sharded Timer wheels. - wheels: RwLock, - - /// Number of entries in the sharded timer wheels. - wheels_len: u32, + // The state is split like this so `Handle` can access `is_shutdown` without locking the mutex + pub(super) state: Mutex, /// True if the driver is being shutdown. pub(super) is_shutdown: AtomicBool, @@ -133,8 +107,14 @@ struct Inner { did_wake: AtomicBool, } -/// Wrapper around the sharded timer wheels. -struct ShardedWheel(Box<[Mutex]>); +/// Time state shared which must be protected by a `Mutex` +struct InnerState { + /// The earliest time at which we promise to wake up without unparking. + next_wake: Option, + + /// Timer wheel. + wheel: wheel::Wheel, +} // ===== impl Driver ===== @@ -143,21 +123,18 @@ impl Driver { /// thread and `time_source` to get the current time and convert to ticks. /// /// Specifying the source of time is useful when testing. - pub(crate) fn new(park: IoStack, clock: &Clock, shards: u32) -> (Driver, Handle) { - assert!(shards > 0); - + pub(crate) fn new(park: IoStack, clock: &Clock) -> (Driver, Handle) { let time_source = TimeSource::new(clock); - let wheels: Vec<_> = (0..shards) - .map(|_| Mutex::new(wheel::Wheel::new())) - .collect(); let handle = Handle { time_source, inner: Inner { - next_wake: AtomicOptionNonZeroU64::new(None), - wheels: RwLock::new(ShardedWheel(wheels.into_boxed_slice())), - wheels_len: shards, + state: Mutex::new(InnerState { + next_wake: None, + wheel: wheel::Wheel::new(), + }), is_shutdown: AtomicBool::new(false), + #[cfg(feature = "test-util")] did_wake: AtomicBool::new(false), }, @@ -187,34 +164,24 @@ impl Driver { // Advance time forward to the end of time. - handle.process_at_time(0, u64::MAX); + handle.process_at_time(u64::MAX); self.park.shutdown(rt_handle); } fn park_internal(&mut self, rt_handle: &driver::Handle, limit: Option) { let handle = rt_handle.time(); + let mut lock = handle.inner.state.lock(); + assert!(!handle.is_shutdown()); - // Finds out the min expiration time to park. - let expiration_time = { - let mut wheels_lock = rt_handle.time().inner.wheels.write(); - let expiration_time = wheels_lock - .0 - .iter_mut() - .filter_map(|wheel| wheel.get_mut().next_expiration_time()) - .min(); + let next_wake = lock.wheel.next_expiration_time(); + lock.next_wake = + next_wake.map(|t| NonZeroU64::new(t).unwrap_or_else(|| NonZeroU64::new(1).unwrap())); - rt_handle - .time() - .inner - .next_wake - .store(next_wake_time(expiration_time)); + drop(lock); - expiration_time - }; - - match expiration_time { + match next_wake { Some(when) => { let now = handle.time_source.now(rt_handle.clock()); // Note that we effectively round up to 1ms here - this avoids @@ -278,60 +245,30 @@ impl Driver { } } -// Helper function to turn expiration_time into next_wake_time. -// Since the `park_timeout` will round up to 1ms for avoiding very -// short-duration microsecond-resolution sleeps, we do the same here. -// The conversion is as follows -// None => None -// Some(0) => Some(1) -// Some(i) => Some(i) -fn next_wake_time(expiration_time: Option) -> Option { - expiration_time.and_then(|v| { - if v == 0 { - NonZeroU64::new(1) - } else { - NonZeroU64::new(v) - } - }) -} - impl Handle { /// Runs timer related logic, and returns the next wakeup time pub(self) fn process(&self, clock: &Clock) { let now = self.time_source().now(clock); - // For fairness, randomly select one to start. - let shards = self.inner.get_shard_size(); - let start = crate::runtime::context::thread_rng_n(shards); - self.process_at_time(start, now); + + self.process_at_time(now); } - pub(self) fn process_at_time(&self, start: u32, now: u64) { - let shards = self.inner.get_shard_size(); - - let expiration_time = (start..shards + start) - .filter_map(|i| self.process_at_sharded_time(i, now)) - .min(); - - self.inner.next_wake.store(next_wake_time(expiration_time)); - } - - // Returns the next wakeup time of this shard. - pub(self) fn process_at_sharded_time(&self, id: u32, mut now: u64) -> Option { + pub(self) fn process_at_time(&self, mut now: u64) { let mut waker_list = WakeList::new(); - let mut wheels_lock = self.inner.wheels.read(); - let mut lock = wheels_lock.lock_sharded_wheel(id); - if now < lock.elapsed() { + let mut lock = self.inner.lock(); + + if now < lock.wheel.elapsed() { // Time went backwards! This normally shouldn't happen as the Rust language // guarantees that an Instant is monotonic, but can happen when running // Linux in a VM on a Windows host due to std incorrectly trusting the // hardware clock to be monotonic. // // See for more information. - now = lock.elapsed(); + now = lock.wheel.elapsed(); } - while let Some(entry) = lock.poll(now) { + while let Some(entry) = lock.wheel.poll(now) { debug_assert!(unsafe { entry.is_pending() }); // SAFETY: We hold the driver lock, and just removed the entry from any linked lists. @@ -341,21 +278,22 @@ impl Handle { if !waker_list.can_push() { // Wake a batch of wakers. To avoid deadlock, we must do this with the lock temporarily dropped. drop(lock); - drop(wheels_lock); waker_list.wake_all(); - wheels_lock = self.inner.wheels.read(); - lock = wheels_lock.lock_sharded_wheel(id); + lock = self.inner.lock(); } } } - let next_wake_up = lock.poll_at(); + + lock.next_wake = lock + .wheel + .poll_at() + .map(|t| NonZeroU64::new(t).unwrap_or_else(|| NonZeroU64::new(1).unwrap())); + drop(lock); - drop(wheels_lock); waker_list.wake_all(); - next_wake_up } /// Removes a registered timer from the driver. @@ -370,11 +308,10 @@ impl Handle { /// `add_entry` must not be called concurrently. pub(self) unsafe fn clear_entry(&self, entry: NonNull) { unsafe { - let wheels_lock = self.inner.wheels.read(); - let mut lock = wheels_lock.lock_sharded_wheel(entry.as_ref().shard_id()); + let mut lock = self.inner.lock(); if entry.as_ref().might_be_registered() { - lock.remove(entry); + lock.wheel.remove(entry); } entry.as_ref().handle().fire(Ok(())); @@ -394,14 +331,12 @@ impl Handle { entry: NonNull, ) { let waker = unsafe { - let wheels_lock = self.inner.wheels.read(); - - let mut lock = wheels_lock.lock_sharded_wheel(entry.as_ref().shard_id()); + let mut lock = self.inner.lock(); // We may have raced with a firing/deregistration, so check before // deregistering. if unsafe { entry.as_ref().might_be_registered() } { - lock.remove(entry); + lock.wheel.remove(entry); } // Now that we have exclusive control of this entry, mint a handle to reinsert it. @@ -415,12 +350,10 @@ impl Handle { // Note: We don't have to worry about racing with some other resetting // thread, because add_entry and reregister require exclusive control of // the timer entry. - match unsafe { lock.insert(entry) } { + match unsafe { lock.wheel.insert(entry) } { Ok(when) => { - if self - .inner + if lock .next_wake - .load() .map(|next_wake| when < next_wake.get()) .unwrap_or(true) { @@ -456,15 +389,15 @@ impl Handle { // ===== impl Inner ===== impl Inner { + /// Locks the driver's inner structure + pub(super) fn lock(&self) -> crate::loom::sync::MutexGuard<'_, InnerState> { + self.state.lock() + } + // Check whether the driver has been shutdown pub(super) fn is_shutdown(&self) -> bool { self.is_shutdown.load(Ordering::SeqCst) } - - // Gets the number of shards. - fn get_shard_size(&self) -> u32 { - self.wheels_len - } } impl fmt::Debug for Inner { @@ -473,19 +406,5 @@ impl fmt::Debug for Inner { } } -// ===== impl ShardedWheel ===== - -impl ShardedWheel { - /// Locks the driver's sharded wheel structure. - pub(super) fn lock_sharded_wheel( - &self, - shard_id: u32, - ) -> crate::loom::sync::MutexGuard<'_, Wheel> { - let index = shard_id % (self.0.len() as u32); - // Safety: This modulo operation ensures that the index is not out of bounds. - unsafe { self.0.get_unchecked(index as usize) }.lock() - } -} - #[cfg(test)] mod tests; diff --git a/tokio/src/runtime/time/tests/mod.rs b/tokio/src/runtime/time/tests/mod.rs index a2271b6fb..c2138a82d 100644 --- a/tokio/src/runtime/time/tests/mod.rs +++ b/tokio/src/runtime/time/tests/mod.rs @@ -65,7 +65,7 @@ fn single_timer() { // This may or may not return Some (depending on how it races with the // thread). If it does return None, however, the timer should complete // synchronously. - time.process_at_time(0, time.time_source().now(clock) + 2_000_000_000); + time.process_at_time(time.time_source().now(clock) + 2_000_000_000); jh.join().unwrap(); }) @@ -99,7 +99,7 @@ fn drop_timer() { let clock = handle.inner.driver().clock(); // advance 2s in the future. - time.process_at_time(0, time.time_source().now(clock) + 2_000_000_000); + time.process_at_time(time.time_source().now(clock) + 2_000_000_000); jh.join().unwrap(); }) @@ -132,7 +132,7 @@ fn change_waker() { let clock = handle.inner.driver().clock(); // advance 2s - time.process_at_time(0, time.time_source().now(clock) + 2_000_000_000); + time.process_at_time(time.time_source().now(clock) + 2_000_000_000); jh.join().unwrap(); }) @@ -172,7 +172,6 @@ fn reset_future() { // This may or may not return a wakeup time. handle.process_at_time( - 0, handle .time_source() .instant_to_tick(start + Duration::from_millis(1500)), @@ -181,7 +180,6 @@ fn reset_future() { assert!(!finished_early.load(Ordering::Relaxed)); handle.process_at_time( - 0, handle .time_source() .instant_to_tick(start + Duration::from_millis(2500)), @@ -224,7 +222,7 @@ fn poll_process_levels() { } for t in 1..normal_or_miri(1024, 64) { - handle.inner.driver().time().process_at_time(0, t as u64); + handle.inner.driver().time().process_at_time(t as u64); for (deadline, future) in entries.iter_mut().enumerate() { let mut context = Context::from_waker(noop_waker_ref()); @@ -253,10 +251,10 @@ fn poll_process_levels_targeted() { let handle = handle.inner.driver().time(); - handle.process_at_time(0, 62); + handle.process_at_time(62); assert!(e1.as_mut().poll_elapsed(&mut context).is_pending()); - handle.process_at_time(0, 192); - handle.process_at_time(0, 192); + handle.process_at_time(192); + handle.process_at_time(192); } #[test] diff --git a/tokio/src/util/mod.rs b/tokio/src/util/mod.rs index 328c4094a..b57c6acfe 100644 --- a/tokio/src/util/mod.rs +++ b/tokio/src/util/mod.rs @@ -57,7 +57,7 @@ cfg_rt! { pub(crate) mod sharded_list; } -#[cfg(any(feature = "rt", feature = "macros", feature = "time"))] +#[cfg(any(feature = "rt", feature = "macros"))] pub(crate) mod rand; cfg_rt! { diff --git a/tokio/src/util/rand.rs b/tokio/src/util/rand.rs index aad85b973..67c45693c 100644 --- a/tokio/src/util/rand.rs +++ b/tokio/src/util/rand.rs @@ -71,7 +71,6 @@ impl FastRand { #[cfg(any( feature = "macros", feature = "rt-multi-thread", - feature = "time", all(feature = "sync", feature = "rt") ))] pub(crate) fn fastrand_n(&mut self, n: u32) -> u32 {