util: Add TokioContext future (#2791) (#2958)

Co-authored-by: Lucio Franco <[email protected]>
Co-authored-by: Blas Rodriguez Irizar <[email protected]>
This commit is contained in:
Lucio Franco
2020-10-14 19:36:24 -04:00
committed by GitHub
co-authored by Blas Rodriguez Irizar
parent b69d4a6108
commit 85d8029e9b
12 changed files with 138 additions and 9 deletions
+2
View File
@@ -1,3 +1,5 @@
#![allow(clippy::unnecessary_lazy_evaluations)]
use proc_macro::TokenStream;
use quote::quote;
use std::num::NonZeroUsize;
+1
View File
@@ -29,6 +29,7 @@ full = ["codec", "udp", "compat"]
compat = ["futures-io",]
codec = ["tokio/stream"]
udp = ["tokio/udp"]
rt = ["tokio/rt-core"]
[dependencies]
tokio = { version = "0.2.5", path = "../tokio" }
+10
View File
@@ -18,6 +18,16 @@ macro_rules! cfg_compat {
}
}
macro_rules! cfg_rt {
($($item:item)*) => {
$(
#[cfg(feature = "rt")]
#[cfg_attr(docsrs, doc(cfg(feature = "rt")))]
$item
)*
}
}
macro_rules! cfg_udp {
($($item:item)*) => {
$(
+78
View File
@@ -0,0 +1,78 @@
//! Tokio context aware futures utilities.
//!
//! This module includes utilities around integrating tokio with other runtimes
//! by allowing the context to be attached to futures. This allows spawning
//! futures on other executors while still using tokio to drive them. This
//! can be useful if you need to use a tokio based library in an executor/runtime
//! that does not provide a tokio context.
use pin_project_lite::pin_project;
use std::{
future::Future,
pin::Pin,
task::{Context, Poll},
};
use tokio::runtime::Handle;
pin_project! {
/// `TokioContext` allows connecting a custom executor with the tokio runtime.
///
/// It contains a `Handle` to the runtime. A handle to the runtime can be
/// obtain by calling the `Runtime::handle()` method.
pub struct TokioContext<F> {
#[pin]
inner: F,
handle: Handle,
}
}
impl<F: Future> Future for TokioContext<F> {
type Output = F::Output;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
let me = self.project();
let handle = me.handle;
let fut = me.inner;
handle.enter(|| fut.poll(cx))
}
}
/// Trait extension that simplifies bundling a `Handle` with a `Future`.
pub trait HandleExt {
/// Convenience method that takes a Future and returns a `TokioContext`.
///
/// # Example: calling Tokio Runtime from a custom ThreadPool
///
/// ```no_run
/// use tokio_util::context::HandleExt;
/// use tokio::time::{delay_for, Duration};
///
/// let mut rt = tokio::runtime::Builder::new()
/// .threaded_scheduler()
/// .enable_all()
/// .build().unwrap();
///
/// let rt2 = tokio::runtime::Builder::new()
/// .threaded_scheduler()
/// .build().unwrap();
///
/// let fut = delay_for(Duration::from_millis(2));
///
/// rt.block_on(
/// rt2
/// .handle()
/// .wrap(async { delay_for(Duration::from_millis(2)).await }),
/// );
///```
fn wrap<F: Future>(&self, fut: F) -> TokioContext<F>;
}
impl HandleExt for Handle {
fn wrap<F: Future>(&self, fut: F) -> TokioContext<F> {
TokioContext {
inner: fut,
handle: self.clone(),
}
}
}
+4
View File
@@ -35,3 +35,7 @@ cfg_udp! {
cfg_compat! {
pub mod compat;
}
cfg_rt! {
pub mod context;
}
+29
View File
@@ -0,0 +1,29 @@
#![warn(rust_2018_idioms)]
#![cfg(feature = "rt")]
use tokio::runtime::Builder;
use tokio::time::*;
use tokio_util::context::HandleExt;
#[test]
fn tokio_context_with_another_runtime() {
let mut rt1 = Builder::new()
.threaded_scheduler()
.core_threads(1)
// no timer!
.build()
.unwrap();
let rt2 = Builder::new()
.threaded_scheduler()
.core_threads(1)
.enable_all()
.build()
.unwrap();
// Without the `HandleExt.wrap()` there would be a panic because there is
// no timer running, since it would be referencing runtime r1.
let _ = rt1.block_on(
rt2.handle()
.wrap(async move { delay_for(Duration::from_millis(2)).await }),
);
}
+3 -1
View File
@@ -2,7 +2,9 @@
#![allow(
clippy::cognitive_complexity,
clippy::large_enum_variant,
clippy::needless_doctest_main
clippy::needless_doctest_main,
clippy::match_like_matches_macro,
clippy::stable_sort_primitive
)]
#![warn(
missing_debug_implementations,
+2 -1
View File
@@ -1,4 +1,5 @@
#![allow(clippy::blacklisted_name)]
#![allow(clippy::blacklisted_name, clippy::stable_sort_primitive)]
use tokio::sync::{mpsc, oneshot};
use tokio::task;
use tokio_test::{assert_ok, assert_pending, assert_ready};
+1 -1
View File
@@ -1,4 +1,4 @@
#![allow(clippy::needless_range_loop)]
#![allow(clippy::needless_range_loop, clippy::stable_sort_primitive)]
#![warn(rust_2018_idioms)]
#![cfg(feature = "full")]
+6 -4
View File
@@ -1,3 +1,5 @@
#![allow(clippy::stable_sort_primitive)]
use tokio::stream::{self, pending, Stream, StreamExt, StreamMap};
use tokio::sync::mpsc;
use tokio_test::{assert_ok, assert_pending, assert_ready, task};
@@ -213,8 +215,8 @@ fn new_capacity_zero() {
let map = StreamMap::<&str, stream::Pending<()>>::new();
assert_eq!(0, map.capacity());
let keys = map.keys().collect::<Vec<_>>();
assert!(keys.is_empty());
let mut keys = map.keys();
assert!(keys.next().is_none());
}
#[test]
@@ -222,8 +224,8 @@ fn with_capacity() {
let map = StreamMap::<&str, stream::Pending<()>>::with_capacity(10);
assert!(10 <= map.capacity());
let keys = map.keys().collect::<Vec<_>>();
assert!(keys.is_empty());
let mut keys = map.keys();
assert!(keys.next().is_none());
}
#[test]
+1 -1
View File
@@ -1,4 +1,4 @@
#![allow(clippy::cognitive_complexity)]
#![allow(clippy::cognitive_complexity, clippy::match_like_matches_macro)]
#![warn(rust_2018_idioms)]
#![cfg(feature = "sync")]
+1 -1
View File
@@ -1,4 +1,4 @@
#![allow(clippy::blacklisted_name)]
#![allow(clippy::blacklisted_name, clippy::stable_sort_primitive)]
#![warn(rust_2018_idioms)]
#![cfg(feature = "full")]