Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ percent-encoding = { version = "2.3", optional = true }
pin-project-lite = "0.2.4"
socket2 = { version = "0.6", optional = true, features = ["all"] }
tracing = { version = "0.1.36", default-features = false, features = ["std"], optional = true }
tokio = { version = "1.13", optional = true, default-features = false }
tokio = { version = "1.49", optional = true, default-features = false }
tower-layer = { version = "0.3", optional = true }
tower-service = { version = "0.3", optional = true }
sync_wrapper = { version = "1", optional = true }
Expand Down
2 changes: 1 addition & 1 deletion src/rt/tokio.rs
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,7 @@ use hyper::rt::{Sleep, Timer};
use pin_project_lite::pin_project;

pub use self::{
executor::{TokioExecutor, TokioLocalExecutor},
executor::{TokioExecutor, TokioHandleExecutor, TokioLocalExecutor},
with_hyper_io::WithHyperIo,
with_tokio_io::WithTokioIo,
};
Expand Down
129 changes: 128 additions & 1 deletion src/rt/tokio/executor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,45 @@ pub struct TokioExecutor {}
#[derive(Default, Debug, Clone)]
pub struct TokioLocalExecutor {}

/// Future executor backed by a runtime [`Handle`].
///
/// This executor, like [`TokioExecutor`], utilises [`tokio`] threads. This
/// executor spawns tasks using [`Handle::spawn()`] rather than
/// [`tokio::spawn()`], however.
///
/// A runtime handle may be obtained by calling [`Runtime::handle()`].
///
/// This is useful for situations in which you wish to run tasks on a
/// *separate* runtime. If your application only runs using a single tokio
/// runtime, [`TokioExecutor`] should be used instead.
///
/// This may be applicable to those configuring hyper clients and servers in
/// applications running on NUMA (Non-Uniform Memory Awareneses) systems, or
/// if you wish to manage provision separate resources for background tasks
/// associated with a client or server.
///
/// See the [`tokio::runtime`] documentation for more information about
/// choosing the correct runtime for your application.
///
/// # Examples
///
/// ```
/// use hyper_util::rt::tokio::TokioHandleExecutor;
///
/// let runtime = tokio::runtime::Builder::new_current_thread()
/// .build()
/// .unwrap();
/// let handle = runtime.handle().clone();
/// let executor = TokioHandleExecutor::new(handle);
/// ```
///
/// [`Handle`]: tokio::runtime::Handle
/// [`Handle::spawn()`]: tokio::runtime::Handle::spawn
/// [`Runtime::handle()`]: tokio::runtime::Runtime::handle
pub struct TokioHandleExecutor {
handle: tokio::runtime::Handle,
}

// ===== impl TokioExecutor =====

impl<Fut> Executor<Fut> for TokioExecutor
Expand Down Expand Up @@ -105,9 +144,31 @@ where
}
}

// ===== impl TokioHandleExecutor =====

impl TokioHandleExecutor {
/// TK
pub fn new(handle: tokio::runtime::Handle) -> Self {
Self { handle }
}
}

impl<Fut> Executor<Fut> for TokioHandleExecutor
where
Fut: Future + Send + 'static,
Fut::Output: Send + 'static,
{
fn execute(&self, fut: Fut) {
self.handle.spawn(fut);
}
}

#[cfg(test)]
mod tests {
use crate::rt::{TokioExecutor, tokio::executor::TokioLocalExecutor};
use crate::rt::{
TokioExecutor,
tokio::{TokioHandleExecutor, TokioLocalExecutor},
};
use hyper::rt::Executor;
use tokio::sync::oneshot;

Expand Down Expand Up @@ -220,4 +281,70 @@ mod tests {
let runtime = tokio::runtime::LocalRuntime::new().unwrap();
runtime.block_on(fut);
}

#[test]
fn handle_executor_can_execute_task_on_separate_runtime() {
// Create a "foreground" runtime we will run our top-level on.
let rt = tokio::runtime::Builder::new_current_thread()
.worker_threads(1)
.name("foreground")
.build()
.unwrap();

// Create a "background" runtime, whose handle will be used to spawn
// background tasks by our executor.
let background = tokio::runtime::Builder::new_current_thread()
.worker_threads(1)
.name("background")
.build()
.unwrap();
let handle = background.handle().clone();
let executor = TokioHandleExecutor::new(handle);

// Begin running the background runtime on a separate worker thread.
let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
let worker = std::thread::Builder::new()
.name("execute-task-on-separate-runtime-worker".into())
.spawn(move || {
use futures_util::FutureExt;
let fut = shutdown_rx.map(drop);
background.block_on(fut);
})
.expect("should spawn thread");

// Run a future that, when polled, spawns a background task onto the
// handle executor. This background task retrieves the name of the
// runtime that it is running on, and sends the name back to its
// caller. The parent then asserts that the child was run on the
// "background" runtime.
rt.block_on(async move {
let handle = tokio::runtime::Handle::current();
let name = handle.name().unwrap().to_string();
assert_eq!(
name, "foreground",
"future should be spawned onto foreground runtime"
);

let (tx, rx) = oneshot::channel();
let fut = async move {
let handle = tokio::runtime::Handle::current();
let name = handle.name().unwrap().to_string();
tx.send(name).unwrap();
};

executor.execute(fut);
let name = rx.await.unwrap();
assert_eq!(
name, "background",
"worker should be spawned onto background runtime"
);
});

// Signal to the background runtime that it should shutdown now, and
// then wait for the thread running it to finish.
shutdown_tx
.send(())
.expect("shutdown signal should be sent");
worker.join().expect("worker thread should finish");
}
}
Loading