mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-27 00:00:12 +02:00
rt: cleanup runtime::context (#2063)
Tweak context to remove more fns and usage of `Option`. Remove `ThreadContext` struct as it is reduced to just `Handle`. Avoid passing around individual driver handles and instead limit to the `runtime::Handle` struct.
This commit is contained in:
@@ -5,7 +5,7 @@ use crate::loom::thread;
|
||||
use crate::runtime::blocking::schedule::NoopSchedule;
|
||||
use crate::runtime::blocking::shutdown;
|
||||
use crate::runtime::blocking::task::BlockingTask;
|
||||
use crate::runtime::{self, context::ThreadContext, io, time, Builder, Callback};
|
||||
use crate::runtime::{Builder, Callback, Handle};
|
||||
use crate::task::{self, JoinHandle};
|
||||
|
||||
use std::collections::VecDeque;
|
||||
@@ -41,18 +41,6 @@ struct Inner {
|
||||
/// Call before a thread stops
|
||||
before_stop: Option<Callback>,
|
||||
|
||||
/// Spawns async tasks
|
||||
spawner: runtime::Spawner,
|
||||
|
||||
/// Runtime I/O driver handle
|
||||
io_handle: io::Handle,
|
||||
|
||||
/// Runtime time driver handle
|
||||
time_handle: time::Handle,
|
||||
|
||||
/// Source of `Instant::now()`
|
||||
clock: time::Clock,
|
||||
|
||||
thread_cap: usize,
|
||||
}
|
||||
|
||||
@@ -74,27 +62,17 @@ pub(crate) fn spawn_blocking<F, R>(func: F) -> JoinHandle<R>
|
||||
where
|
||||
F: FnOnce() -> R + Send + 'static,
|
||||
{
|
||||
use crate::runtime::context::ThreadContext;
|
||||
|
||||
let schedule =
|
||||
ThreadContext::blocking_spawner().expect("not currently running on the Tokio runtime.");
|
||||
let rt = Handle::current();
|
||||
|
||||
let (task, handle) = task::joinable(BlockingTask::new(func));
|
||||
schedule.schedule(task);
|
||||
rt.blocking_spawner.spawn(task, &rt);
|
||||
handle
|
||||
}
|
||||
|
||||
// ===== impl BlockingPool =====
|
||||
|
||||
impl BlockingPool {
|
||||
pub(crate) fn new(
|
||||
builder: &Builder,
|
||||
spawner: &runtime::Spawner,
|
||||
io: &io::Handle,
|
||||
time: &time::Handle,
|
||||
clock: &time::Clock,
|
||||
thread_cap: usize,
|
||||
) -> BlockingPool {
|
||||
pub(crate) fn new(builder: &Builder, thread_cap: usize) -> BlockingPool {
|
||||
let (shutdown_tx, shutdown_rx) = shutdown::channel();
|
||||
|
||||
BlockingPool {
|
||||
@@ -113,10 +91,6 @@ impl BlockingPool {
|
||||
stack_size: builder.thread_stack_size,
|
||||
after_start: builder.after_start.clone(),
|
||||
before_stop: builder.before_stop.clone(),
|
||||
spawner: spawner.clone(),
|
||||
io_handle: io.clone(),
|
||||
time_handle: time.clone(),
|
||||
clock: clock.clone(),
|
||||
thread_cap,
|
||||
}),
|
||||
},
|
||||
@@ -152,21 +126,7 @@ impl fmt::Debug for BlockingPool {
|
||||
// ===== impl Spawner =====
|
||||
|
||||
impl Spawner {
|
||||
/// Set the blocking pool for the duration of the closure
|
||||
///
|
||||
/// If a blocking pool is already set, it will be restored when the closure
|
||||
/// returns or if it panics.
|
||||
pub(crate) fn enter<F, R>(&self, f: F) -> R
|
||||
where
|
||||
F: FnOnce() -> R,
|
||||
{
|
||||
let ctx = crate::runtime::context::ThreadContext::clone_current();
|
||||
let _e = ctx.with_blocking_spawner(self.clone()).enter();
|
||||
|
||||
f()
|
||||
}
|
||||
|
||||
fn schedule(&self, task: Task) {
|
||||
fn spawn(&self, task: Task, rt: &Handle) {
|
||||
let shutdown_tx = {
|
||||
let mut shared = self.inner.shared.lock().unwrap();
|
||||
|
||||
@@ -205,41 +165,32 @@ impl Spawner {
|
||||
};
|
||||
|
||||
if let Some(shutdown_tx) = shutdown_tx {
|
||||
self.spawn_thread(shutdown_tx);
|
||||
self.spawn_thread(shutdown_tx, rt);
|
||||
}
|
||||
}
|
||||
|
||||
fn spawn_thread(&self, shutdown_tx: shutdown::Sender) {
|
||||
fn spawn_thread(&self, shutdown_tx: shutdown::Sender, rt: &Handle) {
|
||||
let mut builder = thread::Builder::new().name(self.inner.thread_name.clone());
|
||||
|
||||
if let Some(stack_size) = self.inner.stack_size {
|
||||
builder = builder.stack_size(stack_size);
|
||||
}
|
||||
let thread_context = ThreadContext::new(
|
||||
self.inner.spawner.clone(),
|
||||
self.inner.io_handle.clone(),
|
||||
self.inner.time_handle.clone(),
|
||||
Some(self.inner.clock.clone()),
|
||||
Some(self.clone()),
|
||||
);
|
||||
let spawner = self.clone();
|
||||
|
||||
let rt = rt.clone();
|
||||
|
||||
builder
|
||||
.spawn(move || {
|
||||
let _e = thread_context.enter();
|
||||
run_thread(spawner);
|
||||
drop(shutdown_tx);
|
||||
// Only the reference should be moved into the closure
|
||||
let rt = &rt;
|
||||
rt.enter(move || {
|
||||
rt.blocking_spawner.inner.run();
|
||||
drop(shutdown_tx);
|
||||
})
|
||||
})
|
||||
.unwrap();
|
||||
}
|
||||
}
|
||||
|
||||
fn run_thread(spawner: Spawner) {
|
||||
spawner.enter(|| {
|
||||
let inner = &*spawner.inner;
|
||||
inner.run()
|
||||
});
|
||||
}
|
||||
|
||||
impl Inner {
|
||||
fn run(&self) {
|
||||
if let Some(f) = &self.after_start {
|
||||
|
||||
Reference in New Issue
Block a user