From 97123db204aaa569493577a4a9f3a0f38fea044e Mon Sep 17 00:00:00 2001 From: Carl Lerche Date: Fri, 9 Jun 2023 11:00:31 -0700 Subject: [PATCH] wip --- .../runtime/scheduler/multi_thread/idle.rs | 253 +++-- .../src/runtime/scheduler/multi_thread/mod.rs | 4 +- .../runtime/scheduler/multi_thread/worker.rs | 980 +++--------------- 3 files changed, 296 insertions(+), 941 deletions(-) diff --git a/tokio/src/runtime/scheduler/multi_thread/idle.rs b/tokio/src/runtime/scheduler/multi_thread/idle.rs index eaa3aae8e..61bb0c338 100644 --- a/tokio/src/runtime/scheduler/multi_thread/idle.rs +++ b/tokio/src/runtime/scheduler/multi_thread/idle.rs @@ -1,10 +1,11 @@ //! Coordinates idling workers use crate::loom::sync::atomic::{AtomicBool, AtomicUsize}; -use crate::runtime::scheduler::multi_thread::Shared; +use crate::loom::sync::MutexGuard; +use crate::runtime::scheduler::multi_thread::{worker, Core, Shared}; use std::fmt; -use std::sync::atomic::Ordering::{self, SeqCst}; +use std::sync::atomic::Ordering::{self, AcqRel, Acquire, Release}; pub(super) struct Idle { /// Number of searching workers @@ -15,21 +16,23 @@ pub(super) struct Idle { /// Used to catch false-negatives when waking workers needs_searching: AtomicBool, + + /// Total number of workers + num_workers: usize, } /// Data synchronized by the scheduler mutex pub(super) struct Synced { - /// Sleeping workers + /// Worker IDs that are currently sleeping sleepers: Vec, } impl Idle { pub(super) fn new(num_workers: usize) -> (Idle, Synced) { - /* - let init = State::new(num_workers); - let idle = Idle { - state: AtomicUsize::new(init.into()), + num_searching: AtomicUsize::new(0), + num_sleeping: AtomicUsize::new(0), + needs_searching: AtomicBool::new(false), num_workers, }; @@ -38,131 +41,167 @@ impl Idle { }; (idle, synced) - */ - todo!() } - /// If there are no workers actively searching, returns the index of a - /// worker currently sleeping. - pub(super) fn worker_to_notify(&self, shared: &Shared) -> Option { - /* - // If at least one worker is spinning, work being notified will - // eventually be found. A searching thread will find **some** work and - // notify another worker, eventually leading to our work being found. - // - // For this to happen, this load must happen before the thread - // transitioning `num_searching` to zero. Acquire / Release does not - // provide sufficient guarantees, so this load is done with `SeqCst` and - // will pair with the `fetch_sub(1)` when transitioning out of - // searching. - if !self.notify_should_wakeup() { - return None; + /// We need at least one searching worker + pub(super) fn notify_local(&self, shared: &Shared) { + if self.num_searching.load(Acquire) != 0 { + // There already is a searching worker. Note, that this could be a + // false positive. However, because this method is called **from** a + // worker, we know that there is at least one worker currently + // awake, so the scheduler won't deadlock. + return; + } + + // There aren't any searching workers. Try to initialize one + if self + .num_searching + .compare_exchange(0, 1, AcqRel, Acquire) + .is_err() + { + // Failing the compare_exchange means another thread concurrently + // launched a searching worker. + return; } // Acquire the lock - let mut lock = shared.synced.lock(); - - // Check again, now that the lock is acquired - if !self.notify_should_wakeup() { - return None; - } - - // A worker should be woken up, atomically increment the number of - // searching workers as well as the number of unparked workers. - State::unpark_one(&self.state, 1); - - // Get the worker to unpark - let ret = lock.idle.sleepers.pop(); - debug_assert!(ret.is_some()); - - ret - */ - todo!() + let synced = shared.synced.lock(); + self.notify_synced(synced, shared, true); + } + + /// Notifies a single worker + pub(super) fn notify_remote(&self, synced: MutexGuard<'_, worker::Synced>, shared: &Shared) { + self.notify_synced(synced, shared, false); + } + + /// Notify a worker while synced + fn notify_synced( + &self, + mut synced: MutexGuard<'_, worker::Synced>, + shared: &Shared, + is_searching: bool, + ) { + // Find a sleeping worker + if let Some(worker) = synced.idle.sleepers.pop() { + // Find an available core + if let Some(mut core) = synced.available_cores.pop() { + debug_assert!(!core.is_searching); + core.is_searching = is_searching; + + // Assign the core to the worker + synced.assigned_cores[worker] = Some(core); + + let num_sleeping = self.num_sleeping.load(Acquire) - 1; + debug_assert_eq!(num_sleeping, synced.idle.sleepers.len()); + + // Update the number of sleeping workers + self.num_sleeping.store(num_sleeping, Release); + + // Drop the lock before notifying the condvar. + drop(synced); + + // Notify the worker + shared.condvars[worker].notify_one(); + return; + } else { + synced.idle.sleepers.push(worker); + } + } + + // 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); + } + + // Explicit mutex guard drop to show that holding the guard to this + // point is significant. `needs_searching` and `num_searching` must be + // updated in the critical section. + drop(synced); } - /// Returns `true` if the worker needs to do a final check for submitted - /// work. pub(super) fn transition_worker_to_parked( &self, - synced: &mut Synced, - is_searching: bool, - ) -> bool { - /* - // Acquire the lock - let mut lock = shared.synced.lock(); + synced: &mut worker::Synced, + core: Box, + index: usize, + ) { + // The core should not be searching at this point + debug_assert!(!core.is_searching); - // Decrement the number of unparked threads - let ret = State::dec_num_unparked(&self.state, is_searching); + // Check that this isn't the final worker to go idle *and* + // `needs_searching` is set. + debug_assert!(!self.needs_searching.load(Acquire) || synced.idle.num_active_workers() > 1); - // Track the sleeping worker - lock.idle.sleepers.push(worker); + let num_sleeping = synced.idle.sleepers.len(); + debug_assert_eq!(num_sleeping, self.num_sleeping.load(Acquire)); - ret - */ - todo!() + // Store the worker index in the list of sleepers + synced.idle.sleepers.push(index); + + // Store the core in the list of available cores + synced.available_cores.push(core); + + // The worker's assigned core slot should be empty + debug_assert!(synced.assigned_cores[index].is_none()); } - pub(super) fn transition_worker_to_searching(&self) -> bool { - /* - let state = State::load(&self.state, SeqCst); - if 2 * state.num_searching() >= self.num_workers { - return false; + pub(super) fn try_transition_worker_to_searching(&self, core: &mut Core) { + debug_assert!(!core.is_searching); + + let num_searching = self.num_searching.load(Acquire); + let num_sleeping = self.num_sleeping.load(Acquire); + + if 2 * num_searching >= self.num_workers - num_sleeping { + return; } - // It is possible for this routine to allow more than 50% of the workers - // to search. That is OK. Limiting searchers is only an optimization to - // prevent too much contention. - State::inc_num_searching(&self.state, SeqCst); - true - */ - todo!() + self.transition_worker_to_searching(core); + } + + /// Needs to happen while synchronized in order to avoid races + pub(super) fn transition_worker_to_searching_if_needed( + &self, + _synced: &mut Synced, + core: &mut Core, + ) -> bool { + if self.needs_searching.load(Acquire) { + // Needs to be called while holding the lock + self.transition_worker_to_searching(core); + true + } else { + false + } + } + + fn transition_worker_to_searching(&self, core: &mut Core) { + core.is_searching = true; + self.num_searching.fetch_add(1, AcqRel); + self.needs_searching.store(false, Release); } /// A lightweight transition from searching -> running. /// /// Returns `true` if this is the final searching worker. The caller /// **must** notify a new worker. - pub(super) fn transition_worker_from_searching(&self) -> bool { - // State::dec_num_searching(&self.state) - todo!() - } + pub(super) fn transition_worker_from_searching(&self, core: &mut Core) -> bool { + debug_assert!(core.is_searching); - /// Unpark a specific worker. This happens if tasks are submitted from - /// within the worker's park routine. - /// - /// Returns `true` if the worker was parked before calling the method. - pub(super) fn unpark_worker_by_id(&self, shared: &Shared, worker_id: usize) -> bool { - /* - let mut lock = shared.synced.lock(); - let sleepers = &mut lock.idle.sleepers; + let prev = self.num_searching.fetch_sub(1, AcqRel); + debug_assert!(prev > 0); - for index in 0..sleepers.len() { - if sleepers[index] == worker_id { - sleepers.swap_remove(index); - - // Update the state accordingly while the lock is held. - State::unpark_one(&self.state, 0); - - return true; - } + if prev == 1 { + false + } else { + core.is_searching = false; + false } - - false - */ - todo!() } - - /// Returns `true` if `worker_id` is contained in the sleep set. - pub(super) fn is_parked(&self, shared: &Shared, worker_id: usize) -> bool { - let lock = shared.synced.lock(); - lock.idle.sleepers.contains(&worker_id) - } - -/* - fn notify_should_wakeup(&self) -> bool { - let state = State(self.state.fetch_add(0, SeqCst)); - state.num_searching() == 0 && state.num_unparked() < self.num_workers - } - */ } +impl Synced { + fn num_active_workers(&self) -> usize { + self.sleepers.capacity() - self.sleepers.len() + } +} diff --git a/tokio/src/runtime/scheduler/multi_thread/mod.rs b/tokio/src/runtime/scheduler/multi_thread/mod.rs index 2599a80aa..11be279ca 100644 --- a/tokio/src/runtime/scheduler/multi_thread/mod.rs +++ b/tokio/src/runtime/scheduler/multi_thread/mod.rs @@ -18,6 +18,7 @@ pub(crate) use stats::Stats; pub(crate) mod queue; mod worker; +use worker::Core; pub(crate) use worker::{Context, Shared}; cfg_taskdump! { @@ -35,8 +36,7 @@ cfg_not_taskdump! { pub(crate) use worker::block_in_place; use crate::runtime::{ - self, - blocking, + self, blocking, driver::{self, Driver}, scheduler, Config, }; diff --git a/tokio/src/runtime/scheduler/multi_thread/worker.rs b/tokio/src/runtime/scheduler/multi_thread/worker.rs index 2949c72fa..8a1a5a0b9 100644 --- a/tokio/src/runtime/scheduler/multi_thread/worker.rs +++ b/tokio/src/runtime/scheduler/multi_thread/worker.rs @@ -67,12 +67,10 @@ use crate::runtime::task::OwnedTasks; use crate::runtime::{ blocking, coop, driver, task, Config, Driver, SchedulerMetrics, WorkerMetrics, }; -use crate::util::atomic_cell::AtomicCell; use crate::util::rand::{FastRand, RngSeedGenerator}; use std::cell::RefCell; use std::task::Waker; -use std::time::Duration; cfg_metrics! { mod metrics; @@ -90,10 +88,13 @@ cfg_not_taskdump! { pub(super) struct Worker { /// Reference to scheduler's handle handle: Arc, + + /// This worker's index in `available_cores` and `condvars`. + index: usize, } /// Core data -struct Core { +pub(super) struct Core { /// Index holding this core's remote/shared state. index: usize, @@ -116,7 +117,7 @@ struct Core { /// True if the worker is currently searching for more work. Searching /// involves attempting to steal from other workers. - is_searching: bool, + pub(super) is_searching: bool, /// True if the scheduler is being shutdown is_shutdown: bool, @@ -136,8 +137,7 @@ struct Core { /// State shared across all workers pub(crate) struct Shared { - /// Per-worker remote state. All other workers have access to this and is - /// how they communicate between each other. + /// Per-core remote state. remotes: Box<[Remote]>, /// Global task queue used for: @@ -154,8 +154,9 @@ pub(crate) struct Shared { /// Data synchronized by the scheduler mutex pub(super) synced: Mutex, - /// Condition variable used to unblock waiting workers - condvar: Condvar, + /// Condition variables used to unblock worker threads. Each worker thread + /// has its own condvar it waits on. + pub(super) condvars: Vec, /// The number of cores that have observed the trace signal. pub(super) trace_status: TraceStatus, @@ -178,7 +179,11 @@ pub(crate) struct Shared { /// Data synchronized by the scheduler mutex pub(crate) struct Synced { /// Cores not currently assigned to workers - cores: Vec>, + pub(super) available_cores: Vec>, + + /// When worker is notified, it is assigned a core. The core is placed here + /// until the worker wakes up to take it. + pub(super) assigned_cores: Vec>>, /// Cores that have observed the shutdown signal /// @@ -205,8 +210,9 @@ struct Remote { /// Thread-local context pub(crate) struct Context { - // /// Worker - // worker: Arc, + // Current scheduler's handle + handle: Arc, + /// Core data core: RefCell>>, @@ -251,7 +257,7 @@ pub(super) fn create( let metrics = WorkerMetrics::from_config(&config); let stats = Stats::new(&metrics); - cores.push(Box::new(Core { + cores.push(Some(Box::new(Core { index: i, tick: 0, lifo_slot: None, @@ -263,7 +269,7 @@ pub(super) fn create( global_queue_interval: stats.tuned_global_queue_interval(&config), stats, rand: FastRand::from_seed(config.seed_generator.next_seed()), - })); + }))); remotes.push(Remote { steal }); worker_metrics.push(metrics); @@ -280,13 +286,14 @@ pub(super) fn create( idle, owned: OwnedTasks::new(), synced: Mutex::new(Synced { - cores, + available_cores: Vec::with_capacity(size), + assigned_cores: cores, shutdown_cores: Vec::with_capacity(size), idle: idle_synced, inject: inject_synced, driver: Some(Box::new(driver)), }), - condvar: Condvar::new(), + condvars: (0..size).map(|_| Condvar::new()).collect(), trace_status: TraceStatus::new(remotes_len), config, scheduler_metrics: SchedulerMetrics::new(), @@ -303,10 +310,11 @@ pub(super) fn create( }; // Eagerly start worker threads - for _ in 0..size { + for index in 0..size { let handle = rt_handle.inner.expect_multi_thread(); let worker = Worker { handle: handle.clone(), + index, }; handle @@ -318,7 +326,7 @@ pub(super) fn create( } #[track_caller] -pub(crate) fn block_in_place(f: F) -> R +pub(crate) fn block_in_place(_f: F) -> R where F: FnOnce() -> R, { @@ -462,6 +470,7 @@ fn run(mut worker: Worker) { crate::runtime::context::enter_runtime(&handle, true, |_| { // Set the worker context. let cx = scheduler::Context::MultiThread(Context { + handle: worker.handle.clone(), core: RefCell::new(None), defer: Defer::new(), }); @@ -550,17 +559,17 @@ impl Worker { todo!() } - fn acquire_core(&self, cx: &Context, mut synced: MutexGuard) -> RunResult { + fn acquire_core(&self, cx: &Context, mut synced: MutexGuard<'_, Synced>) -> RunResult { // Wait until a core is available, then exit the loop. let mut core = loop { - if let Some(core) = synced.cores.pop() { + if let Some(core) = synced.assigned_cores[self.index].take() { break core; } // TODO: not always the case assert!(cx.defer.is_empty()); - synced = self.shared().condvar.wait(synced).unwrap(); + synced = self.shared().condvars[self.index].wait(synced).unwrap(); }; // Reset `lifo_enabled` here in case the core was previously stolen from @@ -578,6 +587,11 @@ impl Worker { return Ok(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(core); + } + // TODO: don't hardcode 128 let n = core.run_queue.max_capacity() / 2; let maybe_task = self.next_remote_task_batch(&mut synced, &mut core, n); @@ -685,6 +699,12 @@ impl Worker { // Start from a random worker let start = core.rand.fastrand_n(num as u32) as usize; + self.steal_one_round(core, start) + } + + fn steal_one_round(&self, core: &mut Core, start: usize) -> Option { + let num = self.shared().remotes.len(); + for i in 0..num { let i = (start + i) % num; @@ -703,8 +723,7 @@ impl Worker { } } - // Fallback on checking the global queue - self.next_remote_task() + None } fn run_task(&self, cx: &Context, mut core: Box, task: Notified) -> RunResult { @@ -712,7 +731,9 @@ impl Worker { // Make sure the worker is not in the **searching** state. This enables // another idle worker to try to steal work. - self.transition_from_searching(&mut core); + if self.transition_from_searching(&mut core) { + self.shared().notify_parked_local(); + } self.assert_lifo_enabled_is_correct(&core); @@ -851,15 +872,11 @@ impl Worker { f(); } - if self.transition_to_parked(&mut core) { + if self.can_transition_to_parked(&mut core) { debug_assert!(!core.is_shutdown); debug_assert!(!core.is_traced); - core.stats.about_to_park(); core = self.do_park(cx, core)?; - } else { - // Just run maintenance and carry on - core = self.park_yield(core); } if let Some(f) = &self.shared().config.after_unpark { @@ -870,68 +887,80 @@ impl Worker { } fn do_park(&self, cx: &Context, mut core: Box) -> RunResult { + core.stats.about_to_park(); + + let was_searching = core.is_searching; + + // Before we park, if we are searching, we need to transition away from searching + if self.transition_from_searching(&mut core) { + // We were the last searching worker, we need to do one last check + if let Some(task) = self.steal_one_round(&mut core, 0) { + self.shared().notify_parked_local(); + + return self.run_task(cx, core, task); + } + } + + // Acquire the lock let mut synced = self.shared().synced.lock(); - // Return `core` to shared - synced.cores.push(core); + // Try one last time to get tasks + let n = core.run_queue.max_capacity() / 2; + if let Some(task) = self.next_remote_task_batch(&mut synced, &mut core, n) { + drop(synced); - if let Some(mut driver) = synced.driver.take() { + return self.run_task(cx, core, task); + } + + if !was_searching { + if self + .shared() + .idle + .transition_worker_to_searching_if_needed(&mut synced.idle, &mut core) + { + // Skip parking, go back to searching + return Ok(core); + } + } + + // Core being returned must not be in the searching state + debug_assert!(!core.is_searching); + + self.shared() + .idle + .transition_worker_to_parked(&mut synced, core, self.index); + + /* + if let Some(_driver) = synced.driver.take() { todo!() } else { - synced = self.shared().condvar.wait(synced).unwrap(); + // Wait for a core to be assigned to us self.acquire_core(cx, synced) } + */ + // TODO: poll driver if needed + self.acquire_core(cx, synced) } fn transition_to_searching(&self, core: &mut Core) -> bool { if !core.is_searching { - core.is_searching = self.shared().idle.transition_worker_to_searching(); + self.shared().idle.try_transition_worker_to_searching(core); } core.is_searching } - fn transition_from_searching(&self, core: &mut Core) { + /// Returns `true` if another worker must be notified + fn transition_from_searching(&self, core: &mut Core) -> bool { if !core.is_searching { - return; - } - - core.is_searching = false; - - if self.shared().idle.transition_worker_from_searching() { - // We are the final searching worker. Because work was found, we - // need to notify another worker. - self.shared().notify_parked_local(); - } - } - - /// Prepares the worker state for parking. - /// - /// Returns true if the transition happened, false if there is work to do first. - fn transition_to_parked(&self, core: &mut Core) -> bool { - // Workers should not park if they have work to do - if core.lifo_slot.is_some() || core.run_queue.has_tasks() || core.is_traced { return false; } - // When the final worker transitions **out** of searching to parked, it - // must check all the queues one last time in case work materialized - // between the last work scan and transitioning out of searching. - let is_last_searcher = self - .shared() - .idle - .transition_worker_to_parked(todo!(), core.is_searching); + self.shared().idle.transition_worker_from_searching(core) + } - // The worker is no longer searching. Setting this is the local cache - // only. - core.is_searching = false; - - if is_last_searcher { - // worker.handle.notify_if_work_pending(); - todo!() - } - - true + fn can_transition_to_parked(&self, core: &mut Core) -> bool { + core.lifo_slot.is_none() && core.run_queue.is_empty() && !core.is_traced } fn transition_from_parked(&self, core: &mut Core) -> bool { @@ -953,15 +982,15 @@ impl Worker { /// If all workers have reached this point, the final cleanup is performed. fn shutdown_core(&self, core: Box) { let mut synced = self.shared().synced.lock(); - synced.cores.push(core); + synced.available_cores.push(core); - if synced.cores.len() != self.shared().remotes.len() { + if synced.available_cores.len() != self.shared().remotes.len() { return; } debug_assert!(self.shared().owned.is_empty()); - for mut core in synced.cores.drain(..) { + for mut core in synced.available_cores.drain(..) { // Drain tasks from the local queue while self.next_local_task(&mut core).is_some() {} } @@ -1018,93 +1047,65 @@ impl Context { impl Shared { pub(super) fn schedule_task(&self, task: Notified, is_yield: bool) { - /* + use std::ptr; + with_current(|maybe_cx| { if let Some(cx) = maybe_cx { // Make sure the task is part of the **current** scheduler. - if self.ptr_eq(&cx.worker.handle) { + if ptr::eq(self, &cx.handle.shared) { // And the current thread still holds a core if let Some(core) = cx.core.borrow_mut().as_mut() { - self.schedule_local(core, task, is_yield); + if is_yield { + // In this case, we want to defer the wake + todo!(); + } else { + self.schedule_local(core, task); + } + return; + } else { + // This can happen if either the core was stolen + // (`block_in_place`) or the notification happens from + // the driver. + todo!(); } } } // Otherwise, use the inject queue. - self.push_remote_task(task); - self.notify_parked_remote(); + self.schedule_remote(task); }) - */ - todo!() } - fn schedule_local(&self, core: &mut Core, task: Notified, is_yield: bool) { - /* - core.stats.inc_local_schedule_count(); + fn schedule_local(&self, core: &mut Core, task: Notified) { + // Push to the LIFO slot + let prev = core.lifo_slot.take(); - // Spawning from the worker thread. If scheduling a "yield" then the - // task must always be pushed to the back of the queue, enabling other - // tasks to be executed. If **not** a yield, then there is more - // flexibility and the task may go to the front of the queue. - let should_notify = if is_yield || !core.lifo_enabled { + if let Some(prev) = prev { core.run_queue - .push_back_or_overflow(task, self, &mut core.stats); - true - } else { - // Push to the LIFO slot - let prev = core.lifo_slot.take(); - let ret = prev.is_some(); - - if let Some(prev) = prev { - core.run_queue - .push_back_or_overflow(prev, self, &mut core.stats); - } - - core.lifo_slot = Some(task); - - ret - }; - - // Only notify if not currently parked. If `park` is `None`, then the - // scheduling is from a resource driver. As notifications often come in - // batches, the notification is delayed until the park is complete. - if should_notify && core.park.is_some() { - self.notify_parked_local(); + .push_back_or_overflow(prev, self, &mut core.stats); } - */ - todo!() + + core.lifo_slot = Some(task); + + self.notify_parked_local(); } fn notify_parked_local(&self) { - /* super::counters::inc_num_inc_notify_local(); - - if let Some(index) = self.idle.worker_to_notify(self) { - super::counters::inc_num_unparks_local(); - self.remotes[index].unpark.unpark(&self.driver); - } - */ - todo!() + self.idle.notify_local(self); } - fn notify_parked_remote(&self) { - /* - if let Some(index) = self.shared.idle.worker_to_notify(&self.shared) { - self.shared.remotes[index].unpark.unpark(&self.driver); - } - */ - todo!() - } - - fn push_remote_task(&self, task: Notified) { + fn schedule_remote(&self, task: Notified) { self.scheduler_metrics.inc_remote_schedule_count(); let mut synced = self.synced.lock(); - // safety: passing in correct `idle::Synced` - unsafe { - self.inject.push(&mut synced.inject, task); - } + // Push the task in the + self.push_remote_task(&mut synced, task); + + // Notify a worker. The mutex is passed in and will be released as part + // of the method call. + self.idle.notify_remote(synced, self); } pub(super) fn close(&self) { @@ -1113,11 +1114,18 @@ impl Shared { todo!() } } + + fn push_remote_task(&self, synced: &mut Synced, task: Notified) { + // safety: passing in correct `idle::Synced` + unsafe { + self.inject.push(&mut synced.inject, task); + } + } } impl Overflow> for Shared { fn push(&self, task: task::Notified>) { - self.push_remote_task(task); + self.push_remote_task(&mut self.synced.lock(), task); } fn push_batch(&self, iter: I) @@ -1140,705 +1148,24 @@ impl<'a> Lock for &'a Shared { } } -impl task::Schedule for Arc { - fn release(&self, task: &Task) -> Option { - // self.shared.owned.remove(task) - todo!() - } - - fn schedule(&self, task: Notified) { - // self.schedule_task(task, false); - todo!() - } - - fn yield_now(&self, task: Notified) { - // self.schedule_task(task, true); - todo!() - } -} - -impl Core { - /// Increment the tick - fn tick(&mut self) { - self.tick = self.tick.wrapping_add(1); - } -} - -pub(crate) struct InjectGuard<'a> { - lock: crate::loom::sync::MutexGuard<'a, Synced>, -} - -impl<'a> AsMut for InjectGuard<'a> { - fn as_mut(&mut self) -> &mut inject::Synced { - &mut self.lock.inject - } -} - -/* -impl Context { - fn run(&self, mut core: Box) -> RunResult { - // Reset `lifo_enabled` here in case the core was previously stolen from - // a task that had the LIFO slot disabled. - self.reset_lifo_enabled(&mut core); - - // Start as "processing" tasks as polling tasks from the local queue - // will be one of the first things we do. - core.stats.start_processing_scheduled_tasks(); - - while !core.is_shutdown { - self.assert_lifo_enabled_is_correct(&core); - - if core.is_traced { - core = self.worker.handle.trace_core(core); - } - - // Increment the tick - core.tick(); - - // Run maintenance, if needed - core = self.maintenance(core); - - // First, check work available to the current worker. - if let Some(task) = core.next_task(&self.worker) { - core = self.run_task(task, core)?; - continue; - } - - // We consumed all work in the queues and will start searching for work. - core.stats.end_processing_scheduled_tasks(); - - // There is no more **local** work to process, try to steal work - // from other workers. - if let Some(task) = core.steal_work(&self.worker) { - // Found work, switch back to processing - core.stats.start_processing_scheduled_tasks(); - core = self.run_task(task, core)?; - } else { - // Wait for work - core = if !self.defer.is_empty() { - self.park_timeout(core, Some(Duration::from_millis(0))) - } else { - self.park(core) - }; - } - } - - core.pre_shutdown(&self.worker); - - // Signal shutdown - self.worker.handle.shutdown_core(core); - Err(()) - } - - fn run_task(&self, task: Notified, mut core: Box) -> RunResult { - let task = self.worker.handle.shared.owned.assert_owner(task); - - // Make sure the worker is not in the **searching** state. This enables - // another idle worker to try to steal work. - core.transition_from_searching(&self.worker); - - self.assert_lifo_enabled_is_correct(&core); - - // 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 - // LIFO slot and are considered part of the current task for scheduling - // purposes. These tasks inherent the "parent"'s limits. - core.stats.start_poll(); - - // Make the core available to the runtime context - *self.core.borrow_mut() = Some(core); - - // Run the task - coop::budget(|| { - task.run(); - let mut lifo_polls = 0; - - // As long as there is budget remaining and a task exists in the - // `lifo_slot`, then keep running. - loop { - // Check if we still have the core. If not, the core was stolen - // by another worker. - let mut core = match self.core.borrow_mut().take() { - Some(core) => core, - None => { - // In this case, we cannot call `reset_lifo_enabled()` - // because the core was stolen. The stealer will handle - // that at the top of `Context::run` - return Err(()); - } - }; - - // Check for a task in the LIFO slot - let task = match core.lifo_slot.take() { - Some(task) => task, - None => { - self.reset_lifo_enabled(&mut core); - core.stats.end_poll(); - return Ok(core); - } - }; - - if !coop::has_budget_remaining() { - core.stats.end_poll(); - - // Not enough budget left to run the LIFO task, push it to - // the back of the queue and return. - core.run_queue.push_back_or_overflow( - task, - &*self.worker.handle, - &mut core.stats, - ); - // If we hit this point, the LIFO slot should be enabled. - // There is no need to reset it. - debug_assert!(core.lifo_enabled); - return Ok(core); - } - - // Track that we are about to run a task from the LIFO slot. - lifo_polls += 1; - super::counters::inc_lifo_schedules(); - - // Disable the LIFO slot if we reach our limit - // - // In ping-ping style workloads where task A notifies task B, - // which notifies task A again, continuously prioritizing the - // LIFO slot can cause starvation as these two tasks will - // repeatedly schedule the other. To mitigate this, we limit the - // number of times the LIFO slot is prioritized. - if lifo_polls >= MAX_LIFO_POLLS_PER_TICK { - core.lifo_enabled = false; - super::counters::inc_lifo_capped(); - } - - // Run the LIFO task, then loop - *self.core.borrow_mut() = Some(core); - let task = self.worker.handle.shared.owned.assert_owner(task); - task.run(); - } - }) - } - - fn reset_lifo_enabled(&self, core: &mut Core) { - core.lifo_enabled = !self.worker.handle.shared.config.disable_lifo_slot; - } - - fn assert_lifo_enabled_is_correct(&self, core: &Core) { - debug_assert_eq!( - core.lifo_enabled, - !self.worker.handle.shared.config.disable_lifo_slot - ); - } - - fn maintenance(&self, mut core: Box) -> Box { - if core.tick % self.worker.handle.shared.config.event_interval == 0 { - super::counters::inc_num_maintenance(); - - core.stats.end_processing_scheduled_tasks(); - - // Call `park` with a 0 timeout. This enables the I/O driver, timer, ... - // to run without actually putting the thread to sleep. - core = self.park_timeout(core, Some(Duration::from_millis(0))); - - // Run regularly scheduled maintenance - core.maintenance(&self.worker); - - core.stats.start_processing_scheduled_tasks(); - } - - core - } - - /// Parks the worker thread while waiting for tasks to execute. - /// - /// This function checks if indeed there's no more work left to be done before parking. - /// Also important to notice that, before parking, the worker thread will try to take - /// ownership of the Driver (IO/Time) and dispatch any events that might have fired. - /// Whenever a worker thread executes the Driver loop, all waken tasks are scheduled - /// in its own local queue until the queue saturates (ntasks > LOCAL_QUEUE_CAPACITY). - /// When the local queue is saturated, the overflow tasks are added to the injection queue - /// from where other workers can pick them up. - /// Also, we rely on the workstealing algorithm to spread the tasks amongst workers - /// after all the IOs get dispatched - fn park(&self, mut core: Box) -> Box { - if let Some(f) = &self.worker.handle.shared.config.before_park { - f(); - } - - if core.transition_to_parked(&self.worker) { - while !core.is_shutdown && !core.is_traced { - core.stats.about_to_park(); - core = self.park_timeout(core, None); - - // Run regularly scheduled maintenance - core.maintenance(&self.worker); - - if core.transition_from_parked(&self.worker) { - break; - } - } - } - - if let Some(f) = &self.worker.handle.shared.config.after_unpark { - f(); - } - core - } - - fn park_timeout(&self, mut core: Box, duration: Option) -> Box { - self.assert_lifo_enabled_is_correct(&core); - - // Take the parker out of core - let mut park = core.park.take().expect("park missing"); - - // Store `core` in context - *self.core.borrow_mut() = Some(core); - - // Park thread - if let Some(timeout) = duration { - park.park_timeout(&self.worker.handle.driver, timeout); - } else { - park.park(&self.worker.handle.driver); - } - - self.defer.wake(); - - // Remove `core` from context - core = self.core.borrow_mut().take().expect("core missing"); - - // Place `park` back in `core` - core.park = Some(park); - - // If there are tasks available to steal, but this worker is not - // looking for tasks to steal, notify another worker. - if !core.is_searching && core.run_queue.is_stealable() { - self.worker.handle.notify_parked_local(); - } - - core - } - - pub(crate) fn defer(&self, waker: &Waker) { - self.defer.defer(waker); - } -} - -impl Core { - /// Increment the tick - fn tick(&mut self) { - self.tick = self.tick.wrapping_add(1); - } - - /// Return the next notified task available to this worker. - fn next_task(&mut self, worker: &Worker) -> Option { - if self.tick % self.global_queue_interval == 0 { - // Update the global queue interval, if needed - self.tune_global_queue_interval(worker); - - worker - .handle - .next_remote_task() - .or_else(|| self.next_local_task()) - } else { - let maybe_task = self.next_local_task(); - - if maybe_task.is_some() { - return maybe_task; - } - - if worker.inject().is_empty() { - return None; - } - - // Other threads can only **remove** tasks from the current worker's - // `run_queue`. So, we can be confident that by the time we call - // `run_queue.push_back` below, there will be *at least* `cap` - // available slots in the queue. - let cap = usize::min( - self.run_queue.remaining_slots(), - self.run_queue.max_capacity() / 2, - ); - - // The worker is currently idle, pull a batch of work from the - // injection queue. We don't want to pull *all* the work so other - // workers can also get some. - let n = usize::min( - worker.inject().len() / worker.handle.shared.remotes.len() + 1, - cap, - ); - - let mut synced = worker.handle.shared.synced.lock(); - // safety: passing in the correct `inject::Synced`. - let mut tasks = unsafe { worker.inject().pop_n(&mut synced.inject, n) }; - - // Pop the first task to return immedietly - let ret = tasks.next(); - - // Push the rest of the on the run queue - self.run_queue.push_back(tasks); - - ret - } - } - - fn next_local_task(&mut self) -> Option { - self.lifo_slot.take().or_else(|| self.run_queue.pop()) - } - - /// Function responsible for stealing tasks from another 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, worker: &Worker) -> Option { - if !self.transition_to_searching(worker) { - return None; - } - - let num = worker.handle.shared.remotes.len(); - // Start from a random worker - let start = self.rand.fastrand_n(num as u32) as usize; - - for i in 0..num { - let i = (start + i) % num; - - // Don't steal from ourself! We know we don't have work. - if i == self.index { - continue; - } - - let target = &worker.handle.shared.remotes[i]; - if let Some(task) = target - .steal - .steal_into(&mut self.run_queue, &mut self.stats) - { - return Some(task); - } - } - - // Fallback on checking the global queue - worker.handle.next_remote_task() - } - - fn transition_to_searching(&mut self, worker: &Worker) -> bool { - if !self.is_searching { - self.is_searching = worker.handle.shared.idle.transition_worker_to_searching(); - } - - self.is_searching - } - - fn transition_from_searching(&mut self, worker: &Worker) { - if !self.is_searching { - return; - } - - self.is_searching = false; - worker.handle.transition_worker_from_searching(); - } - - /// Prepares the worker state for parking. - /// - /// Returns true if the transition happened, false if there is work to do first. - fn transition_to_parked(&mut self, worker: &Worker) -> bool { - // Workers should not park if they have work to do - if self.lifo_slot.is_some() || self.run_queue.has_tasks() || self.is_traced { - return false; - } - - // When the final worker transitions **out** of searching to parked, it - // must check all the queues one last time in case work materialized - // between the last work scan and transitioning out of searching. - let is_last_searcher = worker.handle.shared.idle.transition_worker_to_parked( - &worker.handle.shared, - self.index, - self.is_searching, - ); - - // The worker is no longer searching. Setting this is the local cache - // only. - self.is_searching = false; - - if is_last_searcher { - worker.handle.notify_if_work_pending(); - } - - true - } - - /// Returns `true` if the transition happened. - fn transition_from_parked(&mut self, worker: &Worker) -> bool { - // If a task is in the lifo slot, then we must unpark regardless of - // being notified - if self.lifo_slot.is_some() { - // When a worker wakes, it should only transition to the "searching" - // state when the wake originates from another worker *or* a new task - // is pushed. We do *not* want the worker to transition to "searching" - // when it wakes when the I/O driver receives new events. - self.is_searching = !worker - .handle - .shared - .idle - .unpark_worker_by_id(&worker.handle.shared, self.index); - return true; - } - - if worker - .handle - .shared - .idle - .is_parked(&worker.handle.shared, self.index) - { - return false; - } - - // When unparked, the worker is in the searching state. - self.is_searching = true; - true - } - - /// Runs maintenance work such as checking the pool's state. - fn maintenance(&mut self, worker: &Worker) { - self.stats - .submit(&worker.handle.shared.worker_metrics[self.index]); - - if !self.is_shutdown { - // Check if the scheduler has been shutdown - let synced = worker.handle.shared.synced.lock(); - self.is_shutdown = worker.inject().is_closed(&synced.inject); - } - - if !self.is_traced { - // Check if the worker should be tracing. - self.is_traced = worker.handle.shared.trace_status.trace_requested(); - } - } - - /// Signals all tasks to shut down, and waits for them to complete. Must run - /// before we enter the single-threaded phase of shutdown processing. - fn pre_shutdown(&mut self, worker: &Worker) { - // Signal to all tasks to shut down. - worker.handle.shared.owned.close_and_shutdown_all(); - - self.stats - .submit(&worker.handle.shared.worker_metrics[self.index]); - } - - /// Shuts down the core. - fn shutdown(&mut self, handle: &Handle) { - // Take the core - let mut park = self.park.take().expect("park missing"); - - // Drain the queue - while self.next_local_task().is_some() {} - - park.shutdown(&handle.driver); - } - - fn tune_global_queue_interval(&mut self, worker: &Worker) { - let next = self - .stats - .tuned_global_queue_interval(&worker.handle.shared.config); - - debug_assert!(next > 1); - - // Smooth out jitter - if abs_diff(self.global_queue_interval, next) > 2 { - self.global_queue_interval = next; - } - } -} - -impl Worker { - /// Returns a reference to the scheduler's injection queue. - fn inject(&self) -> &inject::Shared> { - &self.handle.shared.inject - } -} - -// TODO: Move `Handle` impls into handle.rs impl task::Schedule for Arc { fn release(&self, task: &Task) -> Option { self.shared.owned.remove(task) } fn schedule(&self, task: Notified) { - self.schedule_task(task, false); + self.shared.schedule_task(task, false); } fn yield_now(&self, task: Notified) { - self.schedule_task(task, true); + self.shared.schedule_task(task, true); } } -impl Handle { - pub(super) fn schedule_task(&self, task: Notified, is_yield: bool) { - with_current(|maybe_cx| { - if let Some(cx) = maybe_cx { - // Make sure the task is part of the **current** scheduler. - if self.ptr_eq(&cx.worker.handle) { - // And the current thread still holds a core - if let Some(core) = cx.core.borrow_mut().as_mut() { - self.schedule_local(core, task, is_yield); - return; - } - } - } - - // Otherwise, use the inject queue. - self.push_remote_task(task); - self.notify_parked_remote(); - }) - } - - fn schedule_local(&self, core: &mut Core, task: Notified, is_yield: bool) { - core.stats.inc_local_schedule_count(); - - // Spawning from the worker thread. If scheduling a "yield" then the - // task must always be pushed to the back of the queue, enabling other - // tasks to be executed. If **not** a yield, then there is more - // flexibility and the task may go to the front of the queue. - let should_notify = if is_yield || !core.lifo_enabled { - core.run_queue - .push_back_or_overflow(task, self, &mut core.stats); - true - } else { - // Push to the LIFO slot - let prev = core.lifo_slot.take(); - let ret = prev.is_some(); - - if let Some(prev) = prev { - core.run_queue - .push_back_or_overflow(prev, self, &mut core.stats); - } - - core.lifo_slot = Some(task); - - ret - }; - - // Only notify if not currently parked. If `park` is `None`, then the - // scheduling is from a resource driver. As notifications often come in - // batches, the notification is delayed until the park is complete. - if should_notify && core.park.is_some() { - self.notify_parked_local(); - } - } - - fn next_remote_task(&self) -> Option { - if self.shared.inject.is_empty() { - return None; - } - - let mut synced = self.shared.synced.lock(); - // safety: passing in correct `idle::Synced` - unsafe { self.shared.inject.pop(&mut synced.inject) } - } - - fn push_remote_task(&self, task: Notified) { - self.shared.scheduler_metrics.inc_remote_schedule_count(); - - let mut synced = self.shared.synced.lock(); - // safety: passing in correct `idle::Synced` - unsafe { - self.shared.inject.push(&mut synced.inject, task); - } - } - - pub(super) fn close(&self) { - if self - .shared - .inject - .close(&mut self.shared.synced.lock().inject) - { - self.notify_all(); - } - } - - fn notify_parked_local(&self) { - super::counters::inc_num_inc_notify_local(); - - if let Some(index) = self.shared.idle.worker_to_notify(&self.shared) { - super::counters::inc_num_unparks_local(); - self.shared.remotes[index].unpark.unpark(&self.driver); - } - } - - fn notify_parked_remote(&self) { - if let Some(index) = self.shared.idle.worker_to_notify(&self.shared) { - self.shared.remotes[index].unpark.unpark(&self.driver); - } - } - - pub(super) fn notify_all(&self) { - for remote in &self.shared.remotes[..] { - remote.unpark.unpark(&self.driver); - } - } - - fn notify_if_work_pending(&self) { - for remote in &self.shared.remotes[..] { - if !remote.steal.is_empty() { - self.notify_parked_local(); - return; - } - } - - if !self.shared.inject.is_empty() { - self.notify_parked_local(); - } - } - - fn transition_worker_from_searching(&self) { - if self.shared.idle.transition_worker_from_searching() { - // We are the final searching worker. Because work was found, we - // need to notify another worker. - self.notify_parked_local(); - } - } - - /// Signals that a worker has observed the shutdown signal and has replaced - /// its core back into its handle. - /// - /// If all workers have reached this point, the final cleanup is performed. - fn shutdown_core(&self, core: Box) { - let mut cores = self.shared.shutdown_cores.lock(); - cores.push(core); - - if cores.len() != self.shared.remotes.len() { - return; - } - - debug_assert!(self.shared.owned.is_empty()); - - for mut core in cores.drain(..) { - core.shutdown(self); - } - - // Drain the injection queue - // - // We already shut down every task, so we can simply drop the tasks. - while let Some(task) = self.next_remote_task() { - drop(task); - } - } - - fn ptr_eq(&self, other: &Handle) -> bool { - std::ptr::eq(self, other) - } -} - -impl Overflow> for Handle { - fn push(&self, task: task::Notified>) { - self.push_remote_task(task); - } - - fn push_batch(&self, iter: I) - where - I: Iterator>>, - { - unsafe { - self.shared.inject.push_batch(self, iter); - } +impl Core { + /// Increment the tick + fn tick(&mut self) { + self.tick = self.tick.wrapping_add(1); } } @@ -1852,16 +1179,6 @@ impl<'a> AsMut for InjectGuard<'a> { } } -impl<'a> Lock for &'a Handle { - type Handle = InjectGuard<'a>; - - fn lock(self) -> Self::Handle { - InjectGuard { - lock: self.shared.synced.lock(), - } - } -} - #[track_caller] fn with_current(f: impl FnOnce(Option<&Context>) -> R) -> R { use scheduler::Context::MultiThread; @@ -1871,7 +1188,6 @@ fn with_current(f: impl FnOnce(Option<&Context>) -> R) -> R { _ => f(None), }) } -*/ // `u32::abs_diff` is not available on Tokio's MSRV. fn abs_diff(a: u32, b: u32) -> u32 {