mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-22 00:00:11 +02:00
fs: add support for non-threadpool executors (#1495)
Provides a thread pool dedicated to running blocking operations (#588) and update `tokio-fs` to use this pool. In an effort to make incremental progress, this is an initial step towards a final solution. First, it provides a very basic pool implementation with the intend that the pool will be replaced before the final release. Second, it updates `tokio-fs` to always use this blocking pool instead of conditionally using `threadpool::blocking`. Issue #588 contains additional discussion around potential improvements to the "blocking for all" strategy. The implementation provided here builds on work started in #954 and continued in #1045. The general idea is th same as #1045, but the PR improves on some of the details: * The number of explicit operations tracked by `File` is reduced only to the ones that could interact. All other ops are spawned on the blocking pool without being tracked by the `File` instance. * The `seek` implementation is not backed by a trait and `poll_seek` function. This avoids the question of how to model non-blocking seeks on top of a blocking file. In this patch, `seek` is represented as an `async fn`. If the associated future is dropped before the caller observes the return value, we make no effort to define the state in which the file ends up.
This commit is contained in:
@@ -0,0 +1,142 @@
|
||||
//! Thread pool for blocking operations
|
||||
|
||||
use tokio_sync::oneshot;
|
||||
|
||||
use lazy_static::lazy_static;
|
||||
use std::collections::VecDeque;
|
||||
use std::future::Future;
|
||||
use std::pin::Pin;
|
||||
use std::sync::{Condvar, Mutex};
|
||||
use std::task::{Context, Poll};
|
||||
use std::thread;
|
||||
use std::time::Duration;
|
||||
|
||||
struct Pool {
|
||||
shared: Mutex<Shared>,
|
||||
condvar: Condvar,
|
||||
}
|
||||
|
||||
struct Shared {
|
||||
queue: VecDeque<Box<dyn FnOnce() + Send>>,
|
||||
num_th: u32,
|
||||
num_idle: u32,
|
||||
}
|
||||
|
||||
lazy_static! {
|
||||
static ref POOL: Pool = Pool::new();
|
||||
}
|
||||
|
||||
const MAX_THREADS: u32 = 1_000;
|
||||
const KEEP_ALIVE: Duration = Duration::from_secs(10);
|
||||
|
||||
/// Result of a blocking operation running on the blocking thread pool.
|
||||
#[derive(Debug)]
|
||||
pub struct Blocking<T> {
|
||||
rx: oneshot::Receiver<T>,
|
||||
}
|
||||
|
||||
/// Run the provided function on a threadpool dedicated to blocking operations.
|
||||
pub fn run<F, R>(f: F) -> Blocking<R>
|
||||
where
|
||||
F: FnOnce() -> R + Send + 'static,
|
||||
R: Send + 'static,
|
||||
{
|
||||
let (tx, rx) = oneshot::channel();
|
||||
|
||||
let should_spawn = {
|
||||
let mut shared = POOL.shared.lock().unwrap();
|
||||
|
||||
shared.queue.push_back(Box::new(move || {
|
||||
// The receiver may have dropped
|
||||
let _ = tx.send(f());
|
||||
}));
|
||||
|
||||
if shared.num_idle == 0 {
|
||||
// No threads are able to process the task
|
||||
|
||||
if shared.num_th == MAX_THREADS {
|
||||
// At max number of threads
|
||||
false
|
||||
} else {
|
||||
shared.num_th += 1;
|
||||
true
|
||||
}
|
||||
} else {
|
||||
shared.num_idle -= 1;
|
||||
POOL.condvar.notify_one();
|
||||
false
|
||||
}
|
||||
};
|
||||
|
||||
if should_spawn {
|
||||
spawn_thread();
|
||||
}
|
||||
|
||||
Blocking { rx }
|
||||
}
|
||||
|
||||
impl<T> Future for Blocking<T> {
|
||||
type Output = T;
|
||||
|
||||
fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
|
||||
use std::task::Poll::*;
|
||||
|
||||
match Pin::new(&mut self.rx).poll(cx) {
|
||||
Ready(Ok(v)) => Ready(v),
|
||||
Ready(Err(_)) => panic!(
|
||||
"the blocking operation has been dropped before completing. \
|
||||
This should not happen and is a bug."
|
||||
),
|
||||
Pending => Pending,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn spawn_thread() {
|
||||
thread::Builder::new()
|
||||
.name("tokio-blocking-driver".to_string())
|
||||
.spawn(|| {
|
||||
'outer: loop {
|
||||
let mut shared = POOL.shared.lock().unwrap();
|
||||
|
||||
if let Some(task) = shared.queue.pop_front() {
|
||||
drop(shared);
|
||||
run_task(task);
|
||||
continue;
|
||||
}
|
||||
|
||||
// IDLE
|
||||
shared.num_idle += 1;
|
||||
|
||||
loop {
|
||||
shared = POOL.condvar.wait_timeout(shared, KEEP_ALIVE).unwrap().0;
|
||||
|
||||
if let Some(task) = shared.queue.pop_front() {
|
||||
drop(shared);
|
||||
run_task(task);
|
||||
continue 'outer;
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
fn run_task(f: Box<dyn FnOnce() + Send>) {
|
||||
use std::panic::{catch_unwind, AssertUnwindSafe};
|
||||
|
||||
let _ = catch_unwind(AssertUnwindSafe(|| f()));
|
||||
}
|
||||
|
||||
impl Pool {
|
||||
fn new() -> Pool {
|
||||
Pool {
|
||||
shared: Mutex::new(Shared {
|
||||
queue: VecDeque::new(),
|
||||
num_th: 0,
|
||||
num_idle: 0,
|
||||
}),
|
||||
condvar: Condvar::new(),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -67,6 +67,9 @@ mod global;
|
||||
pub mod park;
|
||||
mod typed;
|
||||
|
||||
#[cfg(feature = "blocking")]
|
||||
pub mod blocking;
|
||||
|
||||
#[cfg(feature = "current-thread")]
|
||||
pub mod current_thread;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user