diff --git a/tokio/src/runtime/mod.rs b/tokio/src/runtime/mod.rs index cb198f51f..129a06e39 100644 --- a/tokio/src/runtime/mod.rs +++ b/tokio/src/runtime/mod.rs @@ -186,6 +186,7 @@ pub(crate) mod coop; pub(crate) mod park; mod driver; +use driver::Driver; pub(crate) mod scheduler; diff --git a/tokio/src/runtime/scheduler/multi_thread/worker.rs b/tokio/src/runtime/scheduler/multi_thread/worker.rs index ef6961750..b265e0eae 100644 --- a/tokio/src/runtime/scheduler/multi_thread/worker.rs +++ b/tokio/src/runtime/scheduler/multi_thread/worker.rs @@ -65,7 +65,7 @@ use crate::runtime::scheduler::multi_thread::{ use crate::runtime::scheduler::{inject, Defer, Lock}; use crate::runtime::task::OwnedTasks; use crate::runtime::{ - blocking, coop, driver, scheduler, task, Config, SchedulerMetrics, WorkerMetrics, + blocking, coop, driver, scheduler, task, Config, Driver, SchedulerMetrics, WorkerMetrics, }; use crate::util::atomic_cell::AtomicCell; use crate::util::rand::{FastRand, RngSeedGenerator}; @@ -157,6 +157,10 @@ pub(crate) struct Shared { /// Condition variable used to unblock waiting workers condvar: Condvar, + /// Power's Tokio's I/O, timers, etc... the responsibility of polling the + /// driver is shared across workers. + driver: AtomicCell, + /// Cores that have observed the shutdown signal /// /// The core is **not** placed back in the worker to avoid it from being @@ -207,7 +211,6 @@ struct Remote { pub(crate) struct Context { // /// Worker // worker: Arc, - /// Core data core: RefCell>>, @@ -520,11 +523,11 @@ impl Worker { core.tick(); // Run maintenance, if needed - core = self.maintenance(core); + core = self.maybe_maintenance(core); // First, check work available to the current worker. if let Some(task) = self.next_task(&mut core) { - core = match self.run_task(task, core) { + core = match self.run_task(cx, core, task) { Ok(core) => core, Err(_) => return, }; @@ -541,7 +544,7 @@ impl Worker { // Found work, switch back to processing core.stats.start_processing_scheduled_tasks(); - core = match self.run_task(task, core) { + core = match self.run_task(cx, core, task) { Ok(core) => core, Err(_) => return, }; @@ -563,7 +566,6 @@ impl Worker { // Signal shutdown self.shutdown_core(core); - todo!() } @@ -577,23 +579,20 @@ impl Worker { } fn next_task(&self, core: &mut Core) -> Option { - /* - if self.tick % self.global_queue_interval == 0 { + if core.tick % core.global_queue_interval == 0 { // Update the global queue interval, if needed - self.tune_global_queue_interval(worker); + self.tune_global_queue_interval(core); - worker - .handle - .next_remote_task() - .or_else(|| self.next_local_task()) + self.next_remote_task() + .or_else(|| self.next_local_task(core)) } else { - let maybe_task = self.next_local_task(); + let maybe_task = self.next_local_task(core); if maybe_task.is_some() { return maybe_task; } - if worker.inject().is_empty() { + if self.shared().inject.is_empty() { return None; } @@ -602,32 +601,44 @@ impl Worker { // `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, + core.run_queue.remaining_slots(), + core.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, + self.shared().inject.len() / self.shared().remotes.len() + 1, cap, ); - let mut synced = worker.handle.shared.synced.lock(); + let mut synced = self.shared().synced.lock(); // safety: passing in the correct `inject::Synced`. - let mut tasks = unsafe { worker.inject().pop_n(&mut synced.inject, n) }; + let mut tasks = unsafe { self.shared().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); + core.run_queue.push_back(tasks); ret } - */ - todo!(); + } + + 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 next_local_task(&self, core: &mut Core) -> Option { + core.lifo_slot.take().or_else(|| core.run_queue.pop()) } /// Function responsible for stealing tasks from another worker @@ -636,45 +647,42 @@ impl Worker { /// 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(&self, core: &mut Core) -> Option { - /* - if !self.transition_to_searching(worker) { + if !self.transition_to_searching(core) { return None; } - let num = worker.handle.shared.remotes.len(); + let num = self.shared().remotes.len(); // Start from a random worker - let start = self.rand.fastrand_n(num as u32) as usize; + let start = core.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 { + if i == core.index { continue; } - let target = &worker.handle.shared.remotes[i]; + let target = &self.shared().remotes[i]; + if let Some(task) = target .steal - .steal_into(&mut self.run_queue, &mut self.stats) + .steal_into(&mut core.run_queue, &mut core.stats) { return Some(task); } } // Fallback on checking the global queue - worker.handle.next_remote_task() - */ - todo!() + self.next_remote_task() } - fn run_task(&self, task: Notified, mut core: Box) -> RunResult { - /* - let task = self.worker.handle.shared.owned.assert_owner(task); + fn run_task(&self, cx: &Context, mut core: Box, task: Notified) -> RunResult { + let task = self.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.transition_from_searching(&mut core); self.assert_lifo_enabled_is_correct(&core); @@ -685,7 +693,7 @@ impl Worker { core.stats.start_poll(); // Make the core available to the runtime context - *self.core.borrow_mut() = Some(core); + *cx.core.borrow_mut() = Some(core); // Run the task coop::budget(|| { @@ -697,7 +705,7 @@ impl Worker { 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() { + let mut core = match cx.core.borrow_mut().take() { Some(core) => core, None => { // In this case, we cannot call `reset_lifo_enabled()` @@ -722,11 +730,8 @@ impl Worker { // 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, - ); + core.run_queue + .push_back_or_overflow(task, self.shared(), &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); @@ -750,48 +755,80 @@ impl Worker { } // Run the LIFO task, then loop - *self.core.borrow_mut() = Some(core); - let task = self.worker.handle.shared.owned.assert_owner(task); + *cx.core.borrow_mut() = Some(core); + let task = self.shared().owned.assert_owner(task); task.run(); } }) - */ - todo!() } - fn maintenance(&self, mut core: Box) -> Box { - /* - if core.tick % self.worker.handle.shared.config.event_interval == 0 { + fn maybe_maintenance(&self, mut core: Box) -> Box { + if core.tick % self.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))); + */ + todo!(); // Run regularly scheduled maintenance - core.maintenance(&self.worker); + self.maintenance(&mut core); core.stats.start_processing_scheduled_tasks(); } core - */ - todo!() + } + + /// Runs maintenance work such as checking the pool's state. + fn maintenance(&self, core: &mut Core) { + core.stats.submit(&self.shared().worker_metrics[core.index]); + + if !core.is_shutdown { + // Check if the scheduler has been shutdown + let synced = self.shared().synced.lock(); + core.is_shutdown = self.shared().inject.is_closed(&synced.inject); + } + + if !core.is_traced { + // Check if the worker should be tracing. + core.is_traced = self.shared().trace_status.trace_requested(); + } + } + + fn transition_to_searching(&self, core: &mut Core) -> bool { + if !core.is_searching { + core.is_searching = self.shared().idle.transition_worker_to_searching(); + } + + core.is_searching + } + + fn transition_from_searching(&self, core: &mut Core) { + 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(); + } } /// 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(&self, core: &mut Core) { - /* // Signal to all tasks to shut down. - worker.handle.shared.owned.close_and_shutdown_all(); + self.shared().owned.close_and_shutdown_all(); - self.stats - .submit(&worker.handle.shared.worker_metrics[self.index]); - */ - todo!() + core.stats.submit(&self.shared().worker_metrics[core.index]); } /// Signals that a worker has observed the shutdown signal and has replaced @@ -799,28 +836,30 @@ impl Worker { /// /// 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(); + let mut cores = self.shared().shutdown_cores.lock(); cores.push(core); - if cores.len() != self.shared.remotes.len() { + if cores.len() != self.shared().remotes.len() { return; } - debug_assert!(self.shared.owned.is_empty()); + debug_assert!(self.shared().owned.is_empty()); for mut core in cores.drain(..) { - core.shutdown(self); + // Drain tasks from the local queue + while self.next_local_task(&mut core).is_some() {} } + // Shutdown the driver + let mut driver = self.shared().driver.take().expect("driver missing"); + driver.shutdown(&self.handle.driver); + // 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); } - */ - todo!() } fn reset_lifo_enabled(&self, core: &mut Core) { @@ -833,12 +872,86 @@ impl Worker { !self.handle.shared.config.disable_lifo_slot ); } + + fn tune_global_queue_interval(&self, core: &mut Core) { + let next = core + .stats + .tuned_global_queue_interval(&self.shared().config); + + debug_assert!(next > 1); + + // Smooth out jitter + if abs_diff(core.global_queue_interval, next) > 2 { + core.global_queue_interval = next; + } + } + + fn shared(&self) -> &Shared { + &self.handle.shared + } } impl Context { pub(crate) fn defer(&self, waker: &Waker) { - // self.defer.defer(waker); - todo!(); + self.defer.defer(waker); + } +} + +impl Shared { + 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!() + } + + 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) { + 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); + } + } +} + +impl Overflow> for Shared { + fn push(&self, task: task::Notified>) { + self.push_remote_task(task); + } + + fn push_batch(&self, iter: I) + where + I: Iterator>>, + { + unsafe { + self.inject.push_batch(self, iter); + } + } +} + +impl<'a> Lock for &'a Shared { + type Handle = InjectGuard<'a>; + + fn lock(self) -> Self::Handle { + InjectGuard { + lock: self.synced.lock(), + } } } @@ -866,6 +979,16 @@ impl Core { } } +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 {