From 880ae378bcb37b0b20e562b0736727a70c5d0495 Mon Sep 17 00:00:00 2001 From: Carl Lerche Date: Mon, 26 Jun 2023 18:18:52 +0000 Subject: [PATCH] wip --- tokio/src/runtime/scheduler/lock.rs | 2 - .../runtime/scheduler/multi_thread/idle.rs | 9 -- .../runtime/scheduler/multi_thread/queue.rs | 17 ---- .../runtime/scheduler/multi_thread/stats.rs | 4 - .../runtime/scheduler/multi_thread/worker.rs | 89 ++++++------------- tokio/src/util/atomic_cell.rs | 10 +-- tokio/src/util/mod.rs | 5 -- 7 files changed, 29 insertions(+), 107 deletions(-) diff --git a/tokio/src/runtime/scheduler/lock.rs b/tokio/src/runtime/scheduler/lock.rs index 0111e301f..0901c2b37 100644 --- a/tokio/src/runtime/scheduler/lock.rs +++ b/tokio/src/runtime/scheduler/lock.rs @@ -1,5 +1,3 @@ -use crate::loom::sync::{Mutex, MutexGuard}; - /// A lock (mutex) yielding generic data. pub(crate) trait Lock { type Handle: AsMut; diff --git a/tokio/src/runtime/scheduler/multi_thread/idle.rs b/tokio/src/runtime/scheduler/multi_thread/idle.rs index aa3bc76ae..bd8ffa399 100644 --- a/tokio/src/runtime/scheduler/multi_thread/idle.rs +++ b/tokio/src/runtime/scheduler/multi_thread/idle.rs @@ -67,10 +67,6 @@ impl Idle { self.num_searching.load(Acquire) } - pub(super) fn is_idle(&self, index: usize) -> bool { - self.idle_map.get(index) - } - pub(super) fn snapshot(&self, snapshot: &mut Snapshot) { snapshot.update(&self.idle_map) } @@ -353,11 +349,6 @@ impl IdleMap { IdleMap { chunks } } - fn get(&self, index: usize) -> bool { - let (chunk, mask) = index_to_mask(index); - self.chunks[chunk].load(Acquire) & mask == mask - } - fn set(&self, index: usize) { let (chunk, mask) = index_to_mask(index); let prev = self.chunks[chunk].load(Acquire); diff --git a/tokio/src/runtime/scheduler/multi_thread/queue.rs b/tokio/src/runtime/scheduler/multi_thread/queue.rs index dc34b7ed4..a07d76c96 100644 --- a/tokio/src/runtime/scheduler/multi_thread/queue.rs +++ b/tokio/src/runtime/scheduler/multi_thread/queue.rs @@ -106,11 +106,6 @@ pub(crate) fn local() -> (Steal, Local) { } impl Local { - /// Returns the number of entries in the queue - pub(crate) fn len(&self) -> usize { - self.inner.len() as usize - } - /// How many tasks can be pushed into the queue pub(crate) fn remaining_slots(&self) -> usize { self.inner.remaining_slots() @@ -125,14 +120,6 @@ impl Local { self.inner.is_empty() } - /// Returns false if there are any entries in the queue - /// - /// Separate to is_stealable so that refactors of is_stealable to "protect" - /// some tasks from stealing won't affect this - pub(crate) fn has_tasks(&self) -> bool { - !self.inner.is_empty() - } - /// Pushes a batch of tasks to the back of the queue. All tasks must fit in /// the local queue. /// @@ -394,10 +381,6 @@ impl Local { } impl Steal { - pub(crate) fn is_empty(&self) -> bool { - self.0.is_empty() - } - /// Steals half the tasks from self and place them into `dst`. pub(crate) fn steal_into( &self, diff --git a/tokio/src/runtime/scheduler/multi_thread/stats.rs b/tokio/src/runtime/scheduler/multi_thread/stats.rs index f9e864b6c..57657bb03 100644 --- a/tokio/src/runtime/scheduler/multi_thread/stats.rs +++ b/tokio/src/runtime/scheduler/multi_thread/stats.rs @@ -72,10 +72,6 @@ impl Stats { } } - pub(crate) fn mean_task_poll_duration(&self) -> f64 { - self.task_poll_time_ewma - } - pub(crate) fn tuned_global_queue_interval(&self, config: &Config) -> u32 { // If an interval is explicitly set, don't tune. if let Some(configured) = config.global_queue_interval { diff --git a/tokio/src/runtime/scheduler/multi_thread/worker.rs b/tokio/src/runtime/scheduler/multi_thread/worker.rs index cca3a19df..a026d42e4 100644 --- a/tokio/src/runtime/scheduler/multi_thread/worker.rs +++ b/tokio/src/runtime/scheduler/multi_thread/worker.rs @@ -140,12 +140,6 @@ pub(super) struct Core { rand: FastRand, } -#[test] -fn test_size_of_core() { - let size_of = std::mem::size_of::(); - assert!(size_of <= 64, "actual={}", size_of); -} - /// State shared across all workers pub(crate) struct Shared { /// Per-core remote state. @@ -668,19 +662,19 @@ impl Worker { // Reset `lifo_enabled` here in case the core was previously stolen from // a task that had the LIFO slot disabled. - self.reset_lifo_enabled(cx, core); + self.reset_lifo_enabled(cx); // At this point, the local queue should be empty debug_assert!(core.run_queue.is_empty()); // Update shutdown state while locked - self.update_global_flags(cx, synced, core); + self.update_global_flags(cx, synced); } /// Finds the next task to run, this could be from a queue or stealing. If /// none are available, the thread sleeps and tries again. fn next_task(&mut self, cx: &Context, mut core: Box) -> NextTaskResult { - self.assert_lifo_enabled_is_correct(cx, &core); + self.assert_lifo_enabled_is_correct(cx); if self.is_traced { core = cx.handle.trace_core(core); @@ -703,10 +697,9 @@ impl Worker { super::counters::inc_num_no_local_work(); if !cx.defer.borrow().is_empty() { - // assert!(!core.is_searching); // We are deferring tasks, so poll the resource driver and schedule // the deferred tasks. - core = try_task_new_batch!(self, self.park_yield(cx, core)); + try_task_new_batch!(self, self.park_yield(cx, core)); panic!("what happened to the deferred tasks? 🤔"); } @@ -743,7 +736,7 @@ impl Worker { } } - if let Some(task) = self.next_local_task(cx, &mut core) { + if let Some(task) = self.next_local_task(&mut core) { return Ok((Some(task), core)); } @@ -815,18 +808,11 @@ impl Worker { ret } - fn next_local_task(&self, cx: &Context, core: &mut Core) -> Option { - /* - cx.shared().remotes[core.index] - .lifo_slot - .take_local() - .or_else(|| core.run_queue.pop()) - */ - self.next_lifo_task(cx, core) - .or_else(|| core.run_queue.pop()) + fn next_local_task(&self, core: &mut Core) -> Option { + self.next_lifo_task(core).or_else(|| core.run_queue.pop()) } - fn next_lifo_task(&self, cx: &Context, core: &mut Core) -> Option { + fn next_lifo_task(&self, core: &mut Core) -> Option { core.lifo_slot.take() } @@ -857,21 +843,15 @@ impl Worker { let num = cx.shared().remotes.len(); for i in 0..ROUNDS { - let last = ROUNDS - 1 == i; - // Start from a random worker let start = core.rand.fastrand_n(num as u32) as usize; - if let Some(task) = self.steal_one_round(cx, &mut core, start, last) { + if let Some(task) = self.steal_one_round(cx, &mut core, start) { return Ok((Some(task), core)); } core = try_task_new_batch!(self, self.next_remote_task_batch(cx, core)); - // if cx.shared().idle.num_searching() > 1 { - // break; - // } - if i > 0 { super::counters::inc_num_spin_stall(); std::thread::sleep(std::time::Duration::from_micros(i as u64)); @@ -881,13 +861,7 @@ impl Worker { Ok((None, core)) } - fn steal_one_round( - &self, - cx: &Context, - core: &mut Core, - start: usize, - lifo: bool, - ) -> Option { + fn steal_one_round(&self, cx: &Context, core: &mut Core, start: usize) -> Option { let num = cx.shared().remotes.len(); for i in 0..num { @@ -905,14 +879,6 @@ impl Worker { let target = &cx.shared().remotes[i]; - /* - if lifo { - if let Some(task) = target.lifo_slot.take_remote() { - return Some(task); - } - } - */ - if let Some(task) = target .steal .steal_into(&mut core.run_queue, &mut core.stats) @@ -934,7 +900,7 @@ impl Worker { cx.shared().notify_parked_local(); } - self.assert_lifo_enabled_is_correct(cx, &core); + self.assert_lifo_enabled_is_correct(cx); // Measure the poll start time. Note that we may end up polling other // tasks under this measurement. In this case, the tasks came from the @@ -967,10 +933,10 @@ impl Worker { }; // Check for a task in the LIFO slot - let task = match self.next_lifo_task(cx, &mut core) { + let task = match self.next_lifo_task(&mut core) { Some(task) => task, None => { - self.reset_lifo_enabled(cx, &mut core); + self.reset_lifo_enabled(cx); core.stats.end_poll(); return Ok(core); } @@ -1108,7 +1074,7 @@ impl Worker { core.stats.submit(&cx.shared().worker_metrics[core.index]); } - fn update_global_flags(&mut self, cx: &Context, synced: &mut Synced, core: &mut Core) { + fn update_global_flags(&mut self, cx: &Context, synced: &mut Synced) { if !self.is_shutdown { self.is_shutdown = cx.shared().inject.is_closed(&synced.inject); } @@ -1132,11 +1098,12 @@ impl Worker { self.schedule_deferred_with_core(cx, core, || cx.shared().synced.lock())?; self.flush_metrics(cx, &mut core); - self.update_global_flags(cx, &mut cx.shared().synced.lock(), &mut core); + self.update_global_flags(cx, &mut cx.shared().synced.lock()); Ok((maybe_task, core)) } + /* fn poll_driver(&mut self, cx: &Context, core: Box) -> NextTaskResult { // Call `park` with a 0 timeout. This enables the I/O driver, timer, ... // to run without actually putting the thread to sleep. @@ -1151,13 +1118,14 @@ impl Worker { Ok((None, core)) } } + */ fn park(&mut self, cx: &Context, mut core: Box) -> NextTaskResult { if let Some(f) = &cx.shared().config.before_park { f(); } - if self.can_transition_to_parked(cx, &mut core) { + if self.can_transition_to_parked(&mut core) { debug_assert!(!self.is_shutdown); debug_assert!(!self.is_traced); @@ -1178,7 +1146,7 @@ impl Worker { if self.transition_from_searching(cx, &mut core) { cx.shared().idle.snapshot(&mut self.idle_snapshot); // We were the last searching worker, we need to do one last check - if let Some(task) = self.steal_one_round(cx, &mut core, 0, true) { + if let Some(task) = self.steal_one_round(cx, &mut core, 0) { cx.shared().notify_parked_local(); return Ok((Some(task), core)); @@ -1211,7 +1179,7 @@ impl Worker { self.flush_metrics(cx, &mut core); // If the runtime is shutdown, skip parking - self.update_global_flags(cx, &mut synced, &mut core); + self.update_global_flags(cx, &mut synced); if self.is_shutdown { return Ok((None, core)); @@ -1275,7 +1243,7 @@ impl Worker { cx.shared().idle.transition_worker_from_searching(core) } - fn can_transition_to_parked(&self, cx: &Context, core: &mut Core) -> bool { + fn can_transition_to_parked(&self, core: &mut Core) -> bool { !self.has_tasks(core) && !self.is_shutdown && !self.is_traced } @@ -1309,20 +1277,19 @@ impl Worker { return; } - // Wait for driver - if cx.shared().driver.is_none() { - return; - } + let mut driver = match cx.shared().driver.take() { + Some(driver) => driver, + None => return, + }; debug_assert!(cx.shared().owned.is_empty()); for mut core in synced.shutdown_cores.drain(..) { // Drain tasks from the local queue - while self.next_local_task(cx, &mut core).is_some() {} + while self.next_local_task(&mut core).is_some() {} } // Shutdown the driver - let mut driver = cx.shared().driver.take().expect("driver missing"); driver.shutdown(&cx.handle.driver); // Drain the injection queue @@ -1336,12 +1303,12 @@ impl Worker { } } - fn reset_lifo_enabled(&self, cx: &Context, core: &mut Core) { + fn reset_lifo_enabled(&self, cx: &Context) { cx.lifo_enabled .set(!cx.handle.shared.config.disable_lifo_slot); } - fn assert_lifo_enabled_is_correct(&self, cx: &Context, core: &Core) { + fn assert_lifo_enabled_is_correct(&self, cx: &Context) { debug_assert_eq!( cx.lifo_enabled.get(), !cx.handle.shared.config.disable_lifo_slot diff --git a/tokio/src/util/atomic_cell.rs b/tokio/src/util/atomic_cell.rs index d889654a7..07e37303a 100644 --- a/tokio/src/util/atomic_cell.rs +++ b/tokio/src/util/atomic_cell.rs @@ -1,7 +1,7 @@ use crate::loom::sync::atomic::AtomicPtr; use std::ptr; -use std::sync::atomic::Ordering::{AcqRel, Acquire}; +use std::sync::atomic::Ordering::AcqRel; pub(crate) struct AtomicCell { data: AtomicPtr, @@ -27,16 +27,8 @@ impl AtomicCell { } pub(crate) fn take(&self) -> Option> { - if self.data.load(Acquire).is_null() { - return None; - } - self.swap(None) } - - pub(crate) fn is_none(&self) -> bool { - self.data.load(Acquire).is_null() - } } fn to_raw(data: Option>) -> *mut T { diff --git a/tokio/src/util/mod.rs b/tokio/src/util/mod.rs index 6a7d4b103..ba0f6dd74 100644 --- a/tokio/src/util/mod.rs +++ b/tokio/src/util/mod.rs @@ -68,11 +68,6 @@ cfg_rt! { pub(crate) use rc_cell::RcCell; } -cfg_rt_multi_thread! { - mod try_lock; - pub(crate) use try_lock::TryLock; -} - pub(crate) mod trace; pub(crate) mod error;