From 2a416cba3bd086f8339305c4e7718ca21f102ccb Mon Sep 17 00:00:00 2001 From: Carl Lerche Date: Wed, 21 Jun 2023 20:31:34 +0000 Subject: [PATCH] wip --- .../runtime/scheduler/multi_thread/idle.rs | 67 ++++++---- .../runtime/scheduler/multi_thread/stats.rs | 21 ++++ .../runtime/scheduler/multi_thread/worker.rs | 119 +++++++++--------- 3 files changed, 124 insertions(+), 83 deletions(-) diff --git a/tokio/src/runtime/scheduler/multi_thread/idle.rs b/tokio/src/runtime/scheduler/multi_thread/idle.rs index 53ee01106..ee38761f8 100644 --- a/tokio/src/runtime/scheduler/multi_thread/idle.rs +++ b/tokio/src/runtime/scheduler/multi_thread/idle.rs @@ -99,6 +99,7 @@ impl Idle { } if self.num_idle.load(Acquire) == 0 { + self.needs_searching.store(true, Release); return; } @@ -117,27 +118,36 @@ impl Idle { // Acquire the lock let synced = shared.synced.lock(); - self.notify_synced(synced, shared, true); + self.notify_synced(synced, shared); } /// Notifies a single worker pub(super) fn notify_remote(&self, synced: MutexGuard<'_, worker::Synced>, shared: &Shared) { - self.notify_synced(synced, shared, false); + if synced.idle.sleepers.is_empty() { + self.needs_searching.store(true, Release); + return; + } + + // We need to establish a stronger barrier than with `notify_local` + if self + .num_searching + .compare_exchange(0, 1, AcqRel, Acquire) + .is_err() + { + return; + } + + self.notify_synced(synced, shared); } /// Notify a worker while synced - fn notify_synced( - &self, - mut synced: MutexGuard<'_, worker::Synced>, - shared: &Shared, - is_searching: bool, - ) { + fn notify_synced(&self, mut synced: MutexGuard<'_, worker::Synced>, shared: &Shared) { // Find a sleeping worker if let Some(worker) = synced.idle.sleepers.pop() { // Find an available core if let Some(mut core) = synced.idle.available_cores.pop() { debug_assert!(!core.is_searching); - core.is_searching = is_searching; + core.is_searching = true; self.idle_map.unset(core.index); debug_assert!(self.idle_map.matches(&synced.idle.available_cores)); @@ -164,10 +174,7 @@ impl Idle { // Set the `needs_searching` flag, this happens *while* the lock is held. self.needs_searching.store(true, Release); - - if is_searching { - self.num_searching.fetch_sub(1, Release); - } + self.num_searching.fetch_sub(1, Release); // Explicit mutex guard drop to show that holding the guard to this // point is significant. `needs_searching` and `num_searching` must be @@ -327,18 +334,17 @@ const BIT_MASK: usize = (usize::BITS - 1) as usize; impl IdleMap { fn new(cores: &[Box]) -> IdleMap { - let chunks = (0..num_chunks(cores.len())) - .map(|_| AtomicUsize::new(0)) - .collect(); - let ret = IdleMap { chunks }; - - for core in cores { - ret.set(core.index); - } + let ret = IdleMap::new_n(num_chunks(cores.len())); + ret.set_all(cores); ret } + fn new_n(n: usize) -> IdleMap { + let chunks = (0..n).map(|_| AtomicUsize::new(0)).collect(); + IdleMap { chunks } + } + fn get(&self, index: usize) -> bool { let (chunk, mask) = index_to_mask(index); self.chunks[chunk].load(Acquire) & mask == mask @@ -351,6 +357,12 @@ impl IdleMap { self.chunks[chunk].store(next, Release); } + fn set_all(&self, cores: &[Box]) { + for core in cores { + self.set(core.index); + } + } + fn unset(&self, index: usize) { let (chunk, mask) = index_to_mask(index); let prev = self.chunks[chunk].load(Acquire); @@ -359,19 +371,22 @@ impl IdleMap { } fn matches(&self, idle_cores: &[Box]) -> bool { - let expect = IdleMap::new(idle_cores); + let expect = IdleMap::new_n(self.chunks.len()); + expect.set_all(idle_cores); + for (i, chunk) in expect.chunks.iter().enumerate() { if chunk.load(Acquire) != self.chunks[i].load(Acquire) { return false; } } + true } } impl Snapshot { pub(crate) fn new(idle: &Idle) -> Snapshot { - let chunks = vec![0; num_chunks(idle.idle_map.chunks.len())]; + let chunks = vec![0; idle.idle_map.chunks.len()]; let mut ret = Snapshot { chunks }; ret.update(&idle.idle_map); ret @@ -385,6 +400,12 @@ impl Snapshot { pub(super) fn is_idle(&self, index: usize) -> bool { let (chunk, mask) = index_to_mask(index); + debug_assert!( + chunk < self.chunks.len(), + "index={}; chunks={}", + index, + self.chunks.len() + ); self.chunks[chunk] & mask == mask } } diff --git a/tokio/src/runtime/scheduler/multi_thread/stats.rs b/tokio/src/runtime/scheduler/multi_thread/stats.rs index ce215e3fd..57657bb03 100644 --- a/tokio/src/runtime/scheduler/multi_thread/stats.rs +++ b/tokio/src/runtime/scheduler/multi_thread/stats.rs @@ -28,6 +28,10 @@ pub(crate) struct Ephemeral { /// Number of tasks polled in the batch of scheduled tasks tasks_polled_in_batch: usize, + + /// Used to ensure calls to start / stop batch are paired + #[cfg(debug_assertions)] + batch_started: bool, } impl Ephemeral { @@ -35,6 +39,8 @@ impl Ephemeral { Ephemeral { processing_scheduled_tasks_started_at: Instant::now(), tasks_polled_in_batch: 0, + #[cfg(debug_assertions)] + batch_started: false, } } } @@ -52,6 +58,9 @@ const MAX_TASKS_POLLED_PER_GLOBAL_QUEUE_INTERVAL: u32 = 127; const TARGET_TASKS_POLLED_PER_GLOBAL_QUEUE_INTERVAL: u32 = 61; impl Stats { + pub(crate) const DEFAULT_GLOBAL_QUEUE_INTERVAL: u32 = + TARGET_TASKS_POLLED_PER_GLOBAL_QUEUE_INTERVAL; + pub(crate) fn new(worker_metrics: &WorkerMetrics) -> Stats { // Seed the value with what we hope to see. let task_poll_time_ewma = @@ -98,6 +107,12 @@ impl Stats { pub(crate) fn start_processing_scheduled_tasks(&mut self, ephemeral: &mut Ephemeral) { self.batch.start_processing_scheduled_tasks(); + #[cfg(debug_assertions)] + { + debug_assert!(!ephemeral.batch_started); + ephemeral.batch_started = true; + } + ephemeral.processing_scheduled_tasks_started_at = Instant::now(); ephemeral.tasks_polled_in_batch = 0; } @@ -105,6 +120,12 @@ impl Stats { pub(crate) fn end_processing_scheduled_tasks(&mut self, ephemeral: &mut Ephemeral) { self.batch.end_processing_scheduled_tasks(); + #[cfg(debug_assertions)] + { + debug_assert!(ephemeral.batch_started); + ephemeral.batch_started = false; + } + // Update the EWMA task poll time if ephemeral.tasks_polled_in_batch > 0 { let now = Instant::now(); diff --git a/tokio/src/runtime/scheduler/multi_thread/worker.rs b/tokio/src/runtime/scheduler/multi_thread/worker.rs index efac120fb..58211d0a4 100644 --- a/tokio/src/runtime/scheduler/multi_thread/worker.rs +++ b/tokio/src/runtime/scheduler/multi_thread/worker.rs @@ -490,7 +490,7 @@ fn run( let mut worker = Worker { tick: 0, num_seq_local_queue_polls: 0, - global_queue_interval: 0, + global_queue_interval: Stats::DEFAULT_GLOBAL_QUEUE_INTERVAL, is_shutdown: false, is_traced: false, workers_to_notify: Vec::with_capacity(num_workers - 1), @@ -530,7 +530,7 @@ fn run( }); } -macro_rules! n { +macro_rules! try_task { ($e:expr) => {{ let (task, core) = $e?; if task.is_some() { @@ -540,6 +540,17 @@ macro_rules! n { }}; } +macro_rules! try_task_new_batch { + ($w:expr, $e:expr) => {{ + let (task, mut core) = $e?; + if task.is_some() { + core.stats.start_processing_scheduled_tasks(&mut $w.stats); + return Ok((task, core)); + } + core + }}; +} + impl Worker { fn run(&mut self, cx: &Context, blocking_in_place: bool) -> RunResult { let (maybe_task, mut core) = { @@ -571,7 +582,7 @@ impl Worker { core = self.run_task(cx, core, task)?; } - loop { + while !self.is_shutdown { let (maybe_task, c) = self.next_task(cx, core)?; core = c; @@ -585,8 +596,6 @@ impl Worker { } } - debug_assert!(cx.defer.borrow().is_empty()); - self.pre_shutdown(cx, &mut core); // Signal shutdown @@ -649,11 +658,6 @@ impl Worker { return Ok((None, core)); } - // The core was notified to search for work, don't try to take tasks from the injection queue - if core.is_searching { - return Ok((None, core)); - } - let n = core.run_queue.max_capacity() / 2; let maybe_task = self.next_remote_task_batch(cx, &mut synced, &mut core, n); @@ -663,6 +667,7 @@ impl Worker { /// Ensure core's state is set correctly for the worker to start using. fn reset_acquired_core(&mut self, cx: &Context, synced: &mut Synced, core: &mut Core) { self.global_queue_interval = core.stats.tuned_global_queue_interval(&cx.shared().config); + debug_assert!(self.global_queue_interval > 1); // Reset `lifo_enabled` here in case the core was previously stolen from // a task that had the LIFO slot disabled. @@ -678,42 +683,37 @@ impl Worker { /// 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); + + if self.is_traced { + core = cx.handle.trace_core(core); + } + + // Increment the tick + self.tick = self.tick.wrapping_add(1); + + // Runs maintenance every so often. When maintenance is run, the + // driver is checked, which may result in a task being found. + core = try_task!(self.maybe_maintenance(&cx, core)); + + // Check the LIFO slot, local run queue, and the injection queue for + // a notified task. + core = try_task!(self.next_notified_task(cx, core)); + + // We consumed all work in the queues and will start searching for work. + core.stats.end_processing_scheduled_tasks(&mut self.stats); + + core = try_task_new_batch!(self, self.poll_driver(cx, core)); + while !self.is_shutdown { - self.assert_lifo_enabled_is_correct(cx, &core); - - if self.is_traced { - core = cx.handle.trace_core(core); - } - - // Increment the tick - self.tick = self.tick.wrapping_add(1); - - // Runs maintenance every so often. When maintenance is run, the - // driver is checked, which may result in a task being found. - core = n!(self.maybe_maintenance(&cx, core)); - - // Check the LIFO slot, local run queue, and the injection queue for - // a notified task. - if let Some(task) = self.next_notified_task(cx, &mut core) { - return Ok((Some(task), core)); - } - - // We consumed all work in the queues and will start searching for work. - core.stats.end_processing_scheduled_tasks(&mut self.stats); - - core = n!(self.poll_driver(cx, core)); - // Try to steal a task from other workers - if let Some(task) = self.steal_work(cx, &mut core) { - core.stats.start_processing_scheduled_tasks(&mut self.stats); - return Ok((Some(task), core)); - } + core = try_task_new_batch!(self, self.steal_work(cx, core)); if !cx.defer.borrow().is_empty() { - core = n!(self.park_yield(cx, core)); + core = try_task_new_batch!(self, self.park_yield(cx, core)); } else { super::counters::inc_num_parks(); - core = n!(self.park(cx, core)); + core = try_task_new_batch!(self, self.park(cx, core)); } } @@ -723,26 +723,26 @@ impl Worker { Ok((None, core)) } - fn next_notified_task(&mut self, cx: &Context, core: &mut Core) -> Option { + fn next_notified_task(&mut self, cx: &Context, mut core: Box) -> NextTaskResult { self.num_seq_local_queue_polls += 1; if self.num_seq_local_queue_polls % self.global_queue_interval == 0 { self.num_seq_local_queue_polls = 0; // Update the global queue interval, if needed - self.tune_global_queue_interval(cx, core); + self.tune_global_queue_interval(cx, &mut core); if let Some(task) = self.next_remote_task(cx) { - return Some(task); + return Ok((Some(task), core)); } } - if let Some(task) = self.next_local_task(cx, core) { - return Some(task); + if let Some(task) = self.next_local_task(cx, &mut core) { + return Ok((Some(task), core)); } if cx.shared().inject.is_empty() { - return None; + return Ok((None, core)); } // Other threads can only **remove** tasks from the current worker's @@ -755,7 +755,8 @@ impl Worker { ); let mut synced = cx.shared().synced.lock(); - self.next_remote_task_batch(cx, &mut synced, core, cap) + let maybe_task = self.next_remote_task_batch(cx, &mut synced, &mut core, cap); + Ok((maybe_task, core)) } fn next_remote_task(&self, cx: &Context) -> Option { @@ -819,7 +820,7 @@ impl Worker { /// Note: Only if less than half the workers are searching for tasks to steal /// a new worker will actually try to steal. The idea is to make sure not all /// workers will be trying to steal at the same time. - fn steal_work(&mut self, cx: &Context, core: &mut Core) -> Option { + fn steal_work(&mut self, cx: &Context, mut core: Box) -> NextTaskResult { #[cfg(not(loom))] const ROUNDS: usize = 1; @@ -829,8 +830,8 @@ impl Worker { debug_assert!(core.lifo_slot.is_none()); debug_assert!(core.run_queue.is_empty()); - if !self.transition_to_searching(cx, core) { - return None; + if !self.transition_to_searching(cx, &mut core) { + return Ok((None, core)); } // Get a snapshot of which workers are idle @@ -844,12 +845,12 @@ impl Worker { // 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, core, start, last) { - return Some(task); + if let Some(task) = self.steal_one_round(cx, &mut core, start, last) { + return Ok((Some(task), core)); } } - None + Ok((None, core)) } fn steal_one_round( @@ -1066,7 +1067,7 @@ impl Worker { core.stats.end_processing_scheduled_tasks(&mut self.stats); // Run regularly scheduled maintenance - core = n!(self.park_yield(cx, core)); + core = try_task_new_batch!(self, self.park_yield(cx, core)); core.stats.start_processing_scheduled_tasks(&mut self.stats); } @@ -1088,7 +1089,7 @@ impl Worker { } } - fn park_yield(&mut self, cx: &Context, mut core: Box) -> NextTaskResult { + fn park_yield(&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. if let Some(mut driver) = cx.shared().driver.take() { @@ -1098,10 +1099,8 @@ impl Worker { } // If there are more I/O events, schedule them. - let res = self.schedule_deferred_with_core(cx, core, || cx.shared().synced.lock())?; - - let maybe_task = res.0; - core = res.1; + let (maybe_task, mut core) = + 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); @@ -1133,7 +1132,7 @@ impl Worker { debug_assert!(!self.is_shutdown); debug_assert!(!self.is_traced); - core = n!(self.do_park(cx, core)); + core = try_task!(self.do_park(cx, core)); } if let Some(f) = &cx.shared().config.after_unpark {