diff --git a/Cargo.toml b/Cargo.toml index 3a875beb8..9c0250f19 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -2,7 +2,6 @@ members = [ "tokio", - "tokio-compat", "tokio-macros", "tokio-test", "tokio-tls", diff --git a/tokio-compat/Cargo.toml b/tokio-compat/Cargo.toml deleted file mode 100644 index 8ac9ac399..000000000 --- a/tokio-compat/Cargo.toml +++ /dev/null @@ -1,60 +0,0 @@ -[package] -name = "tokio-compat" -# When releasing to crates.io: -# - Remove path dependencies -# - Update html_root_url. -# - Update doc url -# - Cargo.toml -# - README.md -# - Update CHANGELOG.md. -# - Create "v0.1.x" git tag. -version = "0.1.0-alpha.1" -edition = "2018" -authors = ["Tokio Contributors "] -license = "MIT" -readme = "README.md" -documentation = "https://docs.rs/tokio-compat/0.1.0-alpha.1/tokio-compat/" -repository = "https://github.com/tokio-rs/tokio" -homepage = "https://tokio.rs" -description = """ -Compatibility between `tokio` 0.2 and legacy versions. -""" -categories = ["asynchronous", "network-programming"] -keywords = ["io", "async", "non-blocking", "futures"] - - -[features] -default = ["rt-full", "sink"] -# enables the compat runtimes. -rt-current-thread = [ - "tokio-timer-02", - "tokio-reactor-01", - "tokio-executor-01", - "tokio-02/rt-current-thread", - "tokio-02/timer", - "tokio-02/sync", - "tokio-02/net-driver", -] -rt-full = [ - "tracing-core", - "num_cpus", - "tokio-02/rt-full", - "rt-current-thread", -] -sink = ["futures-util/sink"] - -[dependencies] -futures-01 = { package = "futures", version = "0.1" } -futures-03-core = { package = "futures-core-preview", version = "0.3.0-alpha.19" } -futures-util = { package = "futures-util-preview", version = "0.3.0-alpha.19", default-features = false, features = ["compat"] } -tokio-02 = { package = "tokio", path = "../tokio", default_features = false } - -# runtime-only -tokio-timer-02 = { package = "tokio-timer", version = "0.2", optional = true} -tokio-reactor-01 = { package = "tokio-reactor", version = "0.1", optional = true } -tokio-executor-01 = { package = "tokio-executor", version = "0.1", optional = true } -num_cpus = { version = "1", optional = true } -tracing-core = { version = "0.1", optional = true } - -[dev-dependencies] -tokio-01 = { package = "tokio", version = "0.1" } diff --git a/tokio-compat/README.md b/tokio-compat/README.md deleted file mode 100644 index 7d3eea95f..000000000 --- a/tokio-compat/README.md +++ /dev/null @@ -1,88 +0,0 @@ -# Tokio Compat - -Compatibility layers between `tokio` 0.2 and legacy versions. - -[![Crates.io][crates-badge]][crates-url] -[![MIT licensed][mit-badge]][mit-url] -[![Build Status][azure-badge]][azure-url] -[![Gitter chat][gitter-badge]][gitter-url] - -[crates-badge]: https://img.shields.io/crates/v/tokio-compat.svg -[crates-url]: https://crates.io/crates/tokio-compat -[mit-badge]: https://img.shields.io/badge/license-MIT-blue.svg -[mit-url]: LICENSE -[azure-badge]: https://dev.azure.com/tokio-rs/Tokio/_apis/build/status/tokio-rs.tokio?branchName=master -[azure-url]: https://dev.azure.com/tokio-rs/Tokio/_build/latest?definitionId=1&branchName=master -[gitter-badge]: https://img.shields.io/gitter/room/tokio-rs/tokio.svg -[gitter-url]: https://gitter.im/tokio-rs/tokio - -[Website](https://tokio.rs) | -[Guides](https://tokio.rs/docs/) | -[API Docs](https://docs.rs/tokio-compat/0.1.0-alpha.1/tokio-compat) | -[Chat](https://gitter.im/tokio-rs/tokio) - -## Overview - -This crate provides compatibility runtimes that allow running both `futures` 0.1 -futures that use `tokio` 0.1 runtime services _and_ `std::future` futures that -use `tokio` 0.2 runtime services. - -### Examples - -Spawning both `tokio` 0.1 and `tokio` 0.2 futures: - -```rust -use futures_01::future::lazy; - -tokio_compat::run(lazy(|| { - // spawn a `futures` 0.1 future using the `spawn` function from the - // `tokio` 0.1 crate: - tokio_01::spawn(lazy(|| { - println!("hello from tokio 0.1!"); - Ok(()) - })); - - // spawn an `async` block future on the same runtime using `tokio` - // 0.2's `spawn`: - tokio_02::spawn(async { - println!("hello from tokio 0.2!"); - }); - - Ok(()) -})) -``` - -Futures on the compat runtime can use `timer` APIs from both 0.1 and 0.2 -versions of `tokio`: - -```rust -use std::time::{Duration, Instant}; -use futures_01::future::lazy; -use tokio_compat::prelude::*; - -tokio_compat::run_std(async { - // Wait for a `tokio` 0.1 `Delay`... - let when = Instant::now() + Duration::from_millis(10); - tokio_01::timer::Delay::new(when) - // convert the delay future into a `std::future` that we can `await`. - .compat() - .await - .expect("tokio 0.1 timer should work!"); - println!("10 ms have elapsed"); - - // Wait for a `tokio` 0.2 `Delay`... - let when = Instant::now() + Duration::from_millis(20); - tokio_02::timer::delay(when).await; - println!("20 ms have elapsed"); -}); -``` - -## License - -This project is licensed under the [MIT license](LICENSE). - -### Contribution - -Unless you explicitly state otherwise, any contribution intentionally submitted -for inclusion in Tokio by you, shall be licensed as MIT, without any additional -terms or conditions. diff --git a/tokio-compat/src/lib.rs b/tokio-compat/src/lib.rs deleted file mode 100644 index 26dd45f09..000000000 --- a/tokio-compat/src/lib.rs +++ /dev/null @@ -1,78 +0,0 @@ -//! Compatibility between `tokio` 0.2 and legacy versions. -//! -//! ## Overview -//! -//! This crate provides compatibility runtimes that allow running both `futures` 0.1 -//! futures that use `tokio` 0.1 runtime services _and_ `std::future` futures that -//! use `tokio` 0.2 runtime services. -//! -//! ### Examples -//! -//! Spawning both `tokio` 0.1 and `tokio` 0.2 futures: -//! -//! ```rust -//! use futures_01::future::lazy; -//! -//! tokio_compat::run(lazy(|| { -//! // spawn a `futures` 0.1 future using the `spawn` function from the -//! // `tokio` 0.1 crate: -//! tokio_01::spawn(lazy(|| { -//! println!("hello from tokio 0.1!"); -//! Ok(()) -//! })); -//! -//! // spawn an `async` block future on the same runtime using `tokio` -//! // 0.2's `spawn`: -//! tokio_02::spawn(async { -//! println!("hello from tokio 0.2!"); -//! }); -//! -//! Ok(()) -//! })) -//! ``` -//! -//! Futures on the compat runtime can use `timer` APIs from both 0.1 and 0.2 -//! versions of `tokio`: -//! -//! ```rust -//! use std::time::{Duration, Instant}; -//! use tokio_compat::prelude::*; -//! -//! tokio_compat::run_std(async { -//! // Wait for a `tokio` 0.1 `Delay`... -//! let when = Instant::now() + Duration::from_millis(10); -//! tokio_01::timer::Delay::new(when) -//! // convert the delay future into a `std::future` that we can `await`. -//! .compat() -//! .await -//! .expect("tokio 0.1 timer should work!"); -//! println!("10 ms have elapsed"); -//! -//! // Wait for a `tokio` 0.2 `Delay`... -//! let when = Instant::now() + Duration::from_millis(20); -//! tokio_02::timer::delay(when).await; -//! println!("20 ms have elapsed"); -//! }); -//! ``` -//! -//! ## Feature Flags -//! -//! - `rt-current-thread`: enables the `current_thread` compatibilty runtime -//! - `rt-full`: enables the `current_thread` and threadpool compatibility -//! runtimes (enabled by default) -#![warn( - missing_debug_implementations, - missing_docs, - rust_2018_idioms, - unreachable_pub -)] -#![deny(intra_doc_link_resolution_failure)] -#![doc(test( - no_crate_inject, - attr(deny(warnings, rust_2018_idioms), allow(dead_code, unused_variables)) -))] -#[cfg(any(feature = "rt-current-thread", feature = "rt-full"))] -pub mod runtime; -#[cfg(feature = "rt-full")] -pub use self::runtime::{run, run_std}; -pub mod prelude; diff --git a/tokio-compat/src/prelude.rs b/tokio-compat/src/prelude.rs deleted file mode 100644 index 944288d9d..000000000 --- a/tokio-compat/src/prelude.rs +++ /dev/null @@ -1,4 +0,0 @@ -//! A prelude for `tokio` 0.1/0.2 compatibility. -#[cfg(feature = "sink")] -pub use futures_util::compat::Sink01CompatExt as _; -pub use futures_util::compat::{Future01CompatExt as _, Stream01CompatExt as _}; diff --git a/tokio-compat/src/runtime/compat.rs b/tokio-compat/src/runtime/compat.rs deleted file mode 100644 index 4fe0a9cae..000000000 --- a/tokio-compat/src/runtime/compat.rs +++ /dev/null @@ -1,113 +0,0 @@ -use tokio_executor_01::{self as executor_01, park as park_01}; -use tokio_reactor_01 as reactor_01; -use tokio_timer_02::{clock as clock_02, timer as timer_02}; - -use std::{ - io, thread, - time::{Duration, Instant}, -}; -use tokio_02::executor::{current_thread::CurrentThread, park}; -use tokio_02::sync::oneshot; - -#[derive(Debug)] -pub(super) struct Background { - reactor_handle: reactor_01::Handle, - timer_handle: timer_02::Handle, - shutdown_tx: Option>, - thread: Option>, -} - -#[derive(Debug)] -pub(super) struct Now(N); - -#[derive(Debug)] -struct CompatPark

(P); - -impl Background { - pub(super) fn spawn(clock: &tokio_02::timer::clock::Clock) -> io::Result { - let clock = clock_02::Clock::new_with_now(Now(clock.clone())); - - let reactor = reactor_01::Reactor::new()?; - let reactor_handle = reactor.handle(); - - let timer = timer_02::Timer::new_with_now(reactor, clock); - let timer_handle = timer.handle(); - - let (shutdown_tx, shutdown_rx) = oneshot::channel(); - let shutdown_tx = Some(shutdown_tx); - - let thread = thread::spawn(move || { - let mut rt = CurrentThread::new_with_park(CompatPark(timer)); - let _ = rt.block_on(shutdown_rx); - }); - let thread = Some(thread); - - Ok(Self { - reactor_handle, - timer_handle, - thread, - shutdown_tx, - }) - } - - pub(super) fn reactor(&self) -> &reactor_01::Handle { - &self.reactor_handle - } - - pub(super) fn timer(&self) -> &timer_02::Handle { - &self.timer_handle - } -} - -impl Drop for Background { - fn drop(&mut self) { - let _ = self.shutdown_tx.take().unwrap().send(()); - let _ = self.thread.take().unwrap().join(); - } -} - -pub(super) fn spawn_err(new: tokio_02::executor::SpawnError) -> executor_01::SpawnError { - match new { - _ if new.is_shutdown() => executor_01::SpawnError::shutdown(), - _ if new.is_at_capacity() => executor_01::SpawnError::at_capacity(), - e => unreachable!("weird spawn error {:?}", e), - } -} - -impl

park::Park for CompatPark

-where - P: park_01::Park, -{ - type Unpark = CompatPark; - type Error = P::Error; - - fn unpark(&self) -> Self::Unpark { - CompatPark(self.0.unpark()) - } - - #[inline] - fn park(&mut self) -> Result<(), Self::Error> { - self.0.park() - } - - #[inline] - fn park_timeout(&mut self, duration: Duration) -> Result<(), Self::Error> { - self.0.park_timeout(duration) - } -} - -impl park::Unpark for CompatPark -where - U: park_01::Unpark, -{ - #[inline] - fn unpark(&self) { - self.0.unpark() - } -} - -impl clock_02::Now for Now { - fn now(&self) -> Instant { - self.0.now() - } -} diff --git a/tokio-compat/src/runtime/current_thread/builder.rs b/tokio-compat/src/runtime/current_thread/builder.rs deleted file mode 100644 index b8f033448..000000000 --- a/tokio-compat/src/runtime/current_thread/builder.rs +++ /dev/null @@ -1,86 +0,0 @@ -use crate::runtime::current_thread::Runtime; - -use tokio_02::executor::current_thread::CurrentThread; -use tokio_02::net::driver::Reactor; -use tokio_02::timer::clock::Clock; -use tokio_02::timer::timer::Timer; - -use std::io; - -/// Builds a single-threaded compatibility runtime with custom configuration -/// values. -/// -/// Methods can be chained in order to set the configuration values. The -/// Runtime is constructed by calling [`build`]. -/// -/// New instances of `Builder` are obtained via [`Builder::new`]. -/// -/// See function level documentation for details on the various configuration -/// settings. -/// -/// [`build`]: #method.build -/// [`Builder::new`]: #method.new -/// -/// # Examples -/// -/// ``` -/// use tokio_compat::runtime::current_thread::Builder; -/// use tokio_02::timer::clock::Clock; -/// -/// # pub fn main() { -/// // build Runtime -/// let runtime = Builder::new() -/// .clock(Clock::new()) -/// .build(); -/// // ... call runtime.run(...) -/// # let _ = runtime; -/// # } -/// ``` -#[derive(Debug)] -pub struct Builder { - /// The clock to use - clock: Clock, -} - -impl Builder { - /// Returns a new compatibility runtime builder initialized with default - /// configuration values. - /// - /// Configuration methods can be chained on the return value. - pub fn new() -> Builder { - Builder { - clock: Clock::new(), - } - } - - /// Set the `Clock` instance that will be used by the runtime. - pub fn clock(&mut self, clock: Clock) -> &mut Self { - self.clock = clock; - self - } - - /// Create the configured `Runtime`. - pub fn build(&mut self) -> io::Result { - // We need a reactor to receive events about IO objects from kernel - let reactor = Reactor::new()?; - let reactor_handle = reactor.handle(); - - // Place a timer wheel on top of the reactor. If there are no timeouts to fire, it'll let the - // reactor pick up some new external events. - let timer = Timer::new_with_clock(reactor, self.clock.clone()); - let timer_handle = timer.handle(); - - // And now put a single-threaded executor on top of the timer. When there are no futures ready - // to do something, it'll let the timer or the reactor to generate some new stimuli for the - // futures to continue in their life. - let executor = CurrentThread::new_with_park(timer); - - Runtime::new2(reactor_handle, timer_handle, self.clock.clone(), executor) - } -} - -impl Default for Builder { - fn default() -> Self { - Self::new() - } -} diff --git a/tokio-compat/src/runtime/current_thread/mod.rs b/tokio-compat/src/runtime/current_thread/mod.rs deleted file mode 100644 index 4f6d193bb..000000000 --- a/tokio-compat/src/runtime/current_thread/mod.rs +++ /dev/null @@ -1,145 +0,0 @@ -//! A compatibility implementation that runs everything on the current thread. -//! -//! [`current_thread::Runtime`][rt] is similar to the primary -//! [compatibility `Runtime`][concurrent-rt] except that it runs all components -//! on the current thread instead of using a thread pool. This means that it is -//! able to spawn futures that do not implement `Send`. -//! -//! Same as the default [`current_thread::Runtime`][default-rt] in the main -//! `tokio` crate, the [`tokio_compat::current_thread::Runtime`][rt] includes: -//! -//! * A [reactor] to drive I/O resources. -//! * An [executor] to execute tasks that use these I/O resources. -//! * A [timer] for scheduling work to run after a set period of time. -//! -//! Unlike the default `current_thread::Runtime`, however, the `tokio_compat` -//! version must spawn an additional background thread to run the `tokio` 0.1 -//! [`Reactor`][reactor-01] and [`Timer`][timer-01]. This is necessary to -//! support legacy tasks, as the main thread is already running a `tokio` 0.2 -//! `Reactor` and `Timer`. -//! -//! Note that [`current_thread::Runtime`][rt] does not implement `Send` itself -//! and cannot be safely moved to other threads -//! -//! # Spawning from other threads -//! -//! While [`current_thread::Runtime`][rt] does not implement `Send` and cannot -//! safely be moved to other threads, it provides a `Handle` that can be sent -//! to other threads and allows to spawn new tasks from there. -//! -//! For example: -//! -//! ``` -//! use tokio_compat::runtime::current_thread::Runtime; -//! use std::thread; -//! -//! let runtime = Runtime::new().unwrap(); -//! let handle = runtime.handle(); -//! -//! thread::spawn(move || { -//! // Spawn a `futures` 0.1 task on the other thread's runtime. -//! let _ = handle.spawn(futures_01::future::lazy(|| { -//! println!("hello from futures 0.1!"); -//! Ok(()) -//! })); -//! -//! // Spawn a `std::future` task on the other thread's runtime. -//! let _ = handle.spawn_std(async { -//! println!("hello from std::future!"); -//! }); -//! }).join().unwrap(); -//! ``` -//! -//! # Examples -//! -//! Creating a new `Runtime` and running a future `f` until its completion and -//! returning its result. -//! -//! ``` -//! use tokio_compat::runtime::current_thread::Runtime; -//! -//! let runtime = Runtime::new().unwrap(); -//! -//! // Use the runtime... -//! // runtime.block_on(f); // where f is a future -//! ``` -//! -//! [rt]: struct.Runtime.html -//! [concurrent-rt]: ../struct.Runtime.html -//! [default-rt]: -//! https://docs.rs/tokio/0.2.0-alpha.6/tokio/runtime/current_thread/struct.Runtime.html -//! [chan]: https://docs.rs/futures/0.1/futures/sync/mpsc/fn.channel.html -//! [reactor]: ../../reactor/struct.Reactor.html -//! [executor]: https://tokio.rs/docs/internals/runtime-model/#executors -//! [timer]: ../../timer/index.html -//! [timer-01]: https://docs.rs/tokio/0.1.22/tokio/timer/index.html -//! [reactor-01]: https://docs.rs/tokio/0.1.22/tokio/reactor/struct.Reactor.html -use super::compat; - -mod builder; -mod runtime; -mod task_executor; - -pub use self::builder::Builder; -pub use self::runtime::{Handle, RunError, Runtime}; -pub use self::task_executor::TaskExecutor; - -use futures_01::future::Future as Future01; -use futures_util::{compat::Future01CompatExt, FutureExt}; -use std::future::Future; - -/// Run the provided `futures` 0.1 future to completion using a runtime running on the current thread. -/// -/// This first creates a new [`Runtime`], and calls [`Runtime::block_on`] with the provided future, -/// which blocks the current thread until the provided future completes. It then calls -/// [`Runtime::run`] to wait for any other spawned futures to resolve. -pub fn block_on_all(future: F) -> Result -where - F: Future01, -{ - block_on_all_std(future.compat()) -} - -/// Run the provided `std::future` future to completion using a runtime running on the current thread. -/// -/// This first creates a new [`Runtime`], and calls [`Runtime::block_on`] with the provided future, -/// which blocks the current thread until the provided future completes. It then calls -/// [`Runtime::run`] to wait for any other spawned futures to resolve. -pub fn block_on_all_std(future: F) -> F::Output -where - F: Future, -{ - let mut r = Runtime::new().expect("failed to start runtime on current thread"); - let v = r.block_on_std(future); - r.run().expect("failed to resolve remaining futures"); - v -} - -/// Start a current-thread runtime using the supplied `futures` 0.1 future to bootstrap execution. -/// -/// # Panics -/// -/// This function panics if called from the context of an executor. -pub fn run(future: F) -where - F: Future01 + 'static, -{ - run_std(future.compat().map(|_| ())) -} - -/// Start a current-thread runtime using the supplied `std::future` ture to bootstrap execution. -/// -/// # Panics -/// -/// This function panics if called from the context of an executor. -pub fn run_std(future: F) -where - F: Future + 'static, -{ - let mut r = Runtime::new().expect("failed to start runtime on current thread"); - r.spawn_std(future); - r.run().expect("failed to resolve remaining futures"); -} - -#[cfg(test)] -mod tests; diff --git a/tokio-compat/src/runtime/current_thread/runtime.rs b/tokio-compat/src/runtime/current_thread/runtime.rs deleted file mode 100644 index e9ce007c6..000000000 --- a/tokio-compat/src/runtime/current_thread/runtime.rs +++ /dev/null @@ -1,341 +0,0 @@ -use super::{compat, Builder}; - -use tokio_02::executor::current_thread::Handle as ExecutorHandle; -use tokio_02::executor::current_thread::{self, CurrentThread}; -use tokio_02::net::driver::{self, Reactor}; -use tokio_02::timer::clock::{self, Clock}; -use tokio_02::timer::timer::{self, Timer}; -use tokio_executor_01 as executor_01; -use tokio_reactor_01 as reactor_01; -use tokio_timer_02 as timer_02; - -use futures_01::future::Future as Future01; -use futures_util::{compat::Future01CompatExt, future::FutureExt}; -use std::error::Error; -use std::fmt; -use std::future::Future; -use std::io; - -/// Single-threaded runtime provides a way to start reactor -/// and executor on the current thread. -/// -/// See [module level][mod] documentation for more details. -/// -/// [mod]: index.html -#[derive(Debug)] -pub struct Runtime { - reactor_handle: driver::Handle, - timer_handle: timer::Handle, - clock: Clock, - executor: CurrentThread, - - /// Compatibility background thread. - /// - /// This maintains a `tokio` 0.1 timer and reactor to support running - /// futures that use older tokio APIs. - compat: compat::Background, -} - -pub(super) type Parker = Timer; - -/// Handle to spawn a future on the corresponding `CurrentThread` runtime instance -#[derive(Debug, Clone)] -pub struct Handle(ExecutorHandle); - -impl Handle { - /// Spawn a `futures` 0.1 future onto the `CurrentThread` runtime instance - /// corresponding to this handle. - /// - /// # Panics - /// - /// This function panics if the spawn fails. Failure occurs if the `CurrentThread` - /// instance of the `Handle` does not exist anymore. - pub fn spawn(&self, future: F) -> Result<(), executor_01::SpawnError> - where - F: Future01 + Send + 'static, - { - self.0 - .spawn(future.compat().map(|_| ())) - .map_err(compat::spawn_err) - } - - /// Spawn a `std::future` future onto the `CurrentThread` runtime instance - /// corresponding to this handle. - /// - /// # Panics - /// - /// This function panics if the spawn fails. Failure occurs if the `CurrentThread` - /// instance of the `Handle` does not exist anymore. - pub fn spawn_std(&self, future: F) -> Result<(), tokio_02::executor::SpawnError> - where - F: Future + Send + 'static, - { - self.0.spawn(future) - } - /// Provides a best effort **hint** to whether or not `spawn` will succeed. - /// - /// This function may return both false positives **and** false negatives. - /// If `status` returns `Ok`, then a call to `spawn` will *probably* - /// succeed, but may fail. If `status` returns `Err`, a call to `spawn` will - /// *probably* fail, but may succeed. - /// - /// This allows a caller to avoid creating the task if the call to `spawn` - /// has a high likelihood of failing. - pub fn status(&self) -> Result<(), executor_01::SpawnError> { - self.0.status().map_err(compat::spawn_err) - } -} - -impl tokio_02::executor::TypedExecutor for Handle -where - T: Future + Send + 'static, -{ - fn spawn(&mut self, future: T) -> Result<(), tokio_02::executor::SpawnError> { - Handle::spawn_std(self, future) - } -} - -impl executor_01::TypedExecutor for Handle -where - T: Future01 + Send + 'static, -{ - fn spawn(&mut self, future: T) -> Result<(), executor_01::SpawnError> { - Handle::spawn(self, future) - } -} - -/// Error returned by the `run` function. -#[derive(Debug)] -pub struct RunError { - inner: current_thread::RunError, -} - -impl fmt::Display for RunError { - fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result { - write!(fmt, "{}", self.inner) - } -} - -impl Error for RunError { - fn source(&self) -> Option<&(dyn Error + 'static)> { - self.inner.source() - } -} - -struct CompatExec { - inner: current_thread::TaskExecutor, -} - -impl Runtime { - /// Returns a new runtime initialized with default configuration values. - pub fn new() -> io::Result { - Builder::new().build() - } - - pub(super) fn new2( - reactor_handle: driver::Handle, - timer_handle: timer::Handle, - clock: Clock, - executor: CurrentThread, - ) -> io::Result { - let compat = compat::Background::spawn(&clock)?; - Ok(Runtime { - reactor_handle, - timer_handle, - clock, - executor, - compat, - }) - } - - /// Get a new handle to spawn futures on the single-threaded Tokio runtime - /// - /// Different to the runtime itself, the handle can be sent to different - /// threads. - pub fn handle(&self) -> Handle { - Handle(self.executor.handle().clone()) - } - - /// Spawn a `futures` 0.1 future onto the single-threaded Tokio runtime. - /// - /// See [module level][mod] documentation for more details. - /// - /// [mod]: index.html - /// - /// # Examples - /// - /// ``` - /// use tokio_compat::runtime::current_thread::Runtime; - /// - /// # fn dox() { - /// // Create the runtime - /// let mut rt = Runtime::new().unwrap(); - /// - /// // Spawn a future onto the runtime - /// rt.spawn(futures_01::future::lazy(|| { - /// println!("now running on a worker thread"); - /// Ok(()) - /// })); - /// # } - /// ``` - /// - /// # Panics - /// - /// This function panics if the spawn fails. Failure occurs if the executor - /// is currently at capacity and is unable to spawn a new future. - pub fn spawn(&mut self, future: F) -> &mut Self - where - F: Future01 + 'static, - { - self.executor.spawn(future.compat().map(|_| ())); - self - } - - /// Spawn a `std::future` future onto the single-threaded Tokio runtime. - /// - /// See [module level][mod] documentation for more details. - /// - /// [mod]: index.html - /// - /// # Examples - /// - /// ``` - /// use tokio_compat::runtime::current_thread::Runtime; - /// - /// # fn dox() { - /// // Create the runtime - /// let mut rt = Runtime::new().unwrap(); - /// - /// // Spawn a future onto the runtime - /// rt.spawn_std(async { - /// println!("now running on a worker thread"); - /// }); - /// # } - /// ``` - /// - /// # Panics - /// - /// This function panics if the spawn fails. Failure occurs if the executor - /// is currently at capacity and is unable to spawn a new future. - pub fn spawn_std(&mut self, future: F) -> &mut Self - where - F: Future + 'static, - { - self.executor.spawn(future); - self - } - - /// Runs the provided `futures` 0.1 future, blocking the current thread - /// until the future completes. - /// - /// This function can be used to synchronously block the current thread - /// until the provided `future` has resolved either successfully or with an - /// error. The result of the future is then returned from this function - /// call. - /// - /// Note that this function will **also** execute any spawned futures on the - /// current thread, but will **not** block until these other spawned futures - /// have completed. Once the function returns, any uncompleted futures - /// remain pending in the `Runtime` instance. These futures will not run - /// until `block_on` or `run` is called again. - /// - /// The caller is responsible for ensuring that other spawned futures - /// complete execution by calling `block_on` or `run`. - pub fn block_on(&mut self, f: F) -> Result - where - F: Future01, - { - self.enter(|executor| { - // Run the provided future - executor.block_on(f.compat()) - }) - } - - /// Runs the provided `std::future` future, blocking the current thread - /// until the future completes. - /// - /// This function can be used to synchronously block the current thread - /// until the provided `future` has resolved either successfully or with an - /// error. The result of the future is then returned from this function - /// call. - /// - /// Note that this function will **also** execute any spawned futures on the - /// current thread, but will **not** block until these other spawned futures - /// have completed. Once the function returns, any uncompleted futures - /// remain pending in the `Runtime` instance. These futures will not run - /// until `block_on` or `run` is called again. - /// - /// The caller is responsible for ensuring that other spawned futures - /// complete execution by calling `block_on` or `run`. - pub fn block_on_std(&mut self, f: F) -> F::Output - where - F: Future, - { - self.enter(|executor| { - // Run the provided future - executor.block_on(f) - }) - } - - /// Run the executor to completion, blocking the thread until **all** - /// spawned futures have completed. - pub fn run(&mut self) -> Result<(), RunError> { - self.enter(|executor| executor.run()) - .map_err(|e| RunError { inner: e }) - } - - fn enter(&mut self, f: F) -> R - where - F: FnOnce(&mut current_thread::CurrentThread) -> R, - { - let Runtime { - ref reactor_handle, - ref timer_handle, - ref clock, - ref mut executor, - ref compat, - } = *self; - - let mut enter = executor_01::enter().unwrap(); - // Set the default tokio 0.1 reactor to the background compat reactor. - reactor_01::with_default(compat.reactor(), &mut enter, |enter| { - // This will set the default handle and timer to use inside the closure - // and run the future. - let _reactor = driver::set_default(&reactor_handle); - clock::with_default(clock, || { - // Set up a default timer for tokio 0.1 compat. - timer_02::with_default(compat.timer(), enter, |enter| { - let _timer = timer::set_default(&timer_handle); - // Set default executor for tokio 0.1 futures. - let mut compat_exec = CompatExec { - inner: current_thread::TaskExecutor::current(), - }; - executor_01::with_default(&mut compat_exec, enter, |_enter| { - // The TaskExecutor is a fake executor that looks into the - // current single-threaded executor when used. This is a trick, - // because we need two mutable references to the executor (one - // to run the provided future, another to install as the default - // one). We use the fake one here as the default one. - let mut default_executor = current_thread::TaskExecutor::current(); - tokio_02::executor::with_default(&mut default_executor, || f(executor)) - }) - }) - }) - }) - } -} - -impl executor_01::Executor for CompatExec { - fn spawn( - &mut self, - future: Box + Send>, - ) -> Result<(), executor_01::SpawnError> { - let future = future.compat().map(|_| ()); - tokio_02::executor::Executor::spawn(&mut self.inner, Box::pin(future)) - .map_err(compat::spawn_err) - } - - fn status(&self) -> Result<(), executor_01::SpawnError> { - tokio_02::executor::Executor::status(&self.inner).map_err(compat::spawn_err) - } -} diff --git a/tokio-compat/src/runtime/current_thread/task_executor.rs b/tokio-compat/src/runtime/current_thread/task_executor.rs deleted file mode 100644 index 0f98787ed..000000000 --- a/tokio-compat/src/runtime/current_thread/task_executor.rs +++ /dev/null @@ -1,135 +0,0 @@ -use futures_01::future::{self as future_01, Future as Future01}; -use futures_util::{compat::Future01CompatExt, FutureExt}; -use std::{future::Future, pin::Pin}; -use tokio_02::executor::{self as executor_02, current_thread::TaskExecutor as TaskExecutor02}; -use tokio_executor_01 as executor_01; - -/// Executes futures on the current thread. -/// -/// All futures executed using this executor will be executed on the current -/// thread. As such, `run` will wait for these futures to complete before -/// returning. -/// -/// For more details, see the [module level](../index.html) documentation. -#[derive(Clone, Debug)] -pub struct TaskExecutor { - inner: TaskExecutor02, -} - -impl TaskExecutor { - /// Returns an executor that executes futures on the current thread. - /// - /// The user of `TaskExecutor` must ensure that when a future is submitted, - /// that it is done within the context of a call to `run`. - /// - /// For more details, see the [module level](index.html) documentation. - pub fn current() -> TaskExecutor { - TaskExecutor { - inner: TaskExecutor02::current(), - } - } - - /// Spawn a `futures` 0.1 future onto the current `CurrentThread` instance. - pub fn spawn_local( - &mut self, - future: impl Future01 + 'static, - ) -> Result<(), executor_01::SpawnError> { - let future = Box::pin(future.compat().map(|_| ())); - self.spawn_local_std(future).map_err(map_spawn_err) - } - - /// Spawn a `std::future` future onto the current `CurrentThread` instance. - pub fn spawn_local_std( - &mut self, - future: Pin>>, - ) -> Result<(), executor_02::SpawnError> { - self.inner.spawn_local(future) - } -} - -fn map_spawn_err(new: executor_02::SpawnError) -> executor_01::SpawnError { - match new { - _ if new.is_shutdown() => executor_01::SpawnError::shutdown(), - _ if new.is_at_capacity() => executor_01::SpawnError::at_capacity(), - e => unreachable!("weird spawn error {:?}", e), - } -} - -impl executor_01::Executor for TaskExecutor { - fn spawn( - &mut self, - future: Box + Send>, - ) -> Result<(), executor_01::SpawnError> { - self.spawn_local(future) - } - - fn status(&self) -> Result<(), executor_01::SpawnError> { - executor_02::Executor::status(&self.inner).map_err(map_spawn_err) - } -} - -impl executor_01::TypedExecutor for TaskExecutor -where - F: Future01 + 'static, -{ - fn spawn(&mut self, future: F) -> Result<(), executor_01::SpawnError> { - let future = Box::pin(future.compat().map(|_| ())); - self.spawn_local_std(future).map_err(map_spawn_err) - } - - fn status(&self) -> Result<(), executor_01::SpawnError> { - executor_02::Executor::status(&self.inner).map_err(map_spawn_err) - } -} - -impl executor_02::Executor for TaskExecutor { - fn spawn( - &mut self, - future: Pin + Send>>, - ) -> Result<(), executor_02::SpawnError> { - self.spawn_local_std(future) - } - - fn status(&self) -> Result<(), executor_02::SpawnError> { - executor_02::Executor::status(&self.inner) - } -} - -impl executor_02::TypedExecutor for TaskExecutor -where - F: Future + 'static, -{ - fn spawn(&mut self, future: F) -> Result<(), executor_02::SpawnError> { - self.spawn_local_std(Box::pin(future)) - } - - fn status(&self) -> Result<(), executor_02::SpawnError> { - executor_02::Executor::status(&self.inner) - } -} - -impl future_01::Executor for TaskExecutor -where - F: Future01 + 'static, -{ - fn execute(&self, future: F) -> Result<(), future_01::ExecuteError> { - match executor_02::Executor::status(&self.inner) { - Err(e) if e.is_shutdown() => Err(future_01::ExecuteError::new( - future_01::ExecuteErrorKind::Shutdown, - future, - )), - Err(e) if e.is_at_capacity() => Err(future_01::ExecuteError::new( - future_01::ExecuteErrorKind::NoCapacity, - future, - )), - Err(e) => panic!("unexpected spawn error {:?}", e), - Ok(_) => { - let mut this = self.clone(); - this.spawn_local(future).unwrap_or_else(|e| { - debug_assert!(false, "status succeeded, but spawn failed: {}", e) - }); - Ok(()) - } - } - } -} diff --git a/tokio-compat/src/runtime/current_thread/tests.rs b/tokio-compat/src/runtime/current_thread/tests.rs deleted file mode 100644 index ee0909e51..000000000 --- a/tokio-compat/src/runtime/current_thread/tests.rs +++ /dev/null @@ -1,116 +0,0 @@ -use std::sync::{ - atomic::{AtomicBool, Ordering}, - Arc, -}; -use std::time::{Duration, Instant}; - -use futures_01::future::Future as Future01; -use futures_util::compat::Future01CompatExt; - -#[test] -fn can_run_01_futures() { - let future_ran = Arc::new(AtomicBool::new(false)); - let ran = future_ran.clone(); - - super::run(futures_01::future::lazy(move || { - future_ran.store(true, Ordering::SeqCst); - Ok::<(), ()>(()) - })); - assert!(ran.load(Ordering::SeqCst)); -} - -#[test] -fn can_spawn_01_futures() { - let future_ran = Arc::new(AtomicBool::new(false)); - let ran = future_ran.clone(); - super::run(futures_01::future::lazy(move || { - tokio_01::spawn(futures_01::future::lazy(move || { - future_ran.store(true, Ordering::SeqCst); - Ok(()) - })); - Ok(()) - })); - assert!(ran.load(Ordering::SeqCst)); -} - -#[test] -fn can_spawn_std_futures() { - let future_ran = Arc::new(AtomicBool::new(false)); - let ran = future_ran.clone(); - super::run(futures_01::future::lazy(move || { - tokio_02::spawn(async move { - future_ran.store(true, Ordering::SeqCst); - }); - Ok(()) - })); - assert!(ran.load(Ordering::SeqCst)); -} - -#[test] -fn tokio_01_timers_work() { - let future1_ran = Arc::new(AtomicBool::new(false)); - let ran = future1_ran.clone(); - let future1 = futures_01::future::lazy(|| { - let when = Instant::now() + Duration::from_millis(15); - tokio_01::timer::Delay::new(when).map(move |_| when) - }) - .map(move |when| { - ran.store(true, Ordering::SeqCst); - assert!(Instant::now() >= when); - }) - .map_err(|_| panic!("timer should work")); - - let future2_ran = Arc::new(AtomicBool::new(false)); - let ran = future2_ran.clone(); - let future2 = async move { - let when = Instant::now() + Duration::from_millis(10); - tokio_01::timer::Delay::new(when).compat().await.unwrap(); - ran.store(true, Ordering::SeqCst); - assert!(Instant::now() >= when); - }; - - super::run(futures_01::future::lazy(move || { - tokio_02::spawn(future2); - tokio_01::spawn(future1); - Ok(()) - })); - assert!(future1_ran.load(Ordering::SeqCst)); - assert!(future2_ran.load(Ordering::SeqCst)); -} - -#[test] -fn block_on_01_timer() { - let mut rt = super::Runtime::new().unwrap(); - let when = Instant::now() + Duration::from_millis(10); - rt.block_on(tokio_01::timer::Delay::new(when)).unwrap(); - assert!(Instant::now() >= when); -} - -#[test] -fn block_on_std_01_timer() { - let mut rt = super::Runtime::new().unwrap(); - let when = Instant::now() + Duration::from_millis(10); - rt.block_on_std(async move { - tokio_01::timer::Delay::new(when).compat().await.unwrap(); - }); - assert!(Instant::now() >= when); -} - -#[test] -fn block_on_01_spawn() { - let mut rt = super::Runtime::new().unwrap(); - // other tests assert that spawned 0.1 tasks actually *run*, all we care - // is that we're able to spawn it successfully. - rt.block_on(futures_01::future::lazy(|| { - tokio_01::spawn(futures_01::future::lazy(|| Ok(()))) - })) - .unwrap(); -} - -#[test] -fn block_on_std_01_spawn() { - let mut rt = super::Runtime::new().unwrap(); - // other tests assert that spawned 0.1 tasks actually *run*, all we care - // is that we're able to spawn it successfully. - rt.block_on_std(async { tokio_01::spawn(futures_01::future::lazy(|| Ok(()))) }); -} diff --git a/tokio-compat/src/runtime/mod.rs b/tokio-compat/src/runtime/mod.rs deleted file mode 100644 index 9e6dafa13..000000000 --- a/tokio-compat/src/runtime/mod.rs +++ /dev/null @@ -1,123 +0,0 @@ -//! Runtimes compatible with both `tokio` 0.1 and `tokio` 0.2 futures. -//! -//! This module is similar to the [`tokio::runtime`] module, with one -//! key difference: the runtimes in this crate are capable of executing -//! both `futures` 0.1 futures that use the `tokio` 0.1 runtime services -//! (i.e. `timer`, `reactor`, and `executor`), **and** `std::future` -//! futures that use the `tokio` 0.2 runtime services. -//! -//! The `futures` crate's [`compat` module][futures-compat] provides -//! interoperability between `futures` 0.1 and `std::future` _future types_ -//! (e.g. implementing `std::future::Future` for a type that implements the -//! `futures` 0.1 `Future` trait). However, this on its own is insufficient to -//! run code written against `tokio` 0.1 on a `tokio` 0.2 runtime, if that code -//! also relies on `tokio`'s runtime services. If legacy tasks are executed that -//! rely on `tokio::timer`, perform IO using `tokio`'s reactor, or call -//! `tokio::spawn`, those API calls will fail unless there is also a runtime -//! compatibility layer. -//! -//! `tokio-compat`'s `runtime` module contains modified versions of the `tokio` -//! 0.2 `Runtime` and `current_thread::Runtime` that are capable of providing -//! `tokio` 0.1 and `tokio` 0.2-compatible runtime services. -//! -//! Creating a [`Runtime`] does the following: -//! -//! * Spawn a background thread running a [`Reactor`] instance. -//! * Start a [`ThreadPool`] for executing futures. -//! * Run an instance of [`Timer`] **per** thread pool worker thread. -//! * Run a **single** `tokio` 0.1 [`Reactor`][reactor-01] and -//! [`Timer`][timer-01] on a background thread, for legacy tasks. -//! -//! Legacy `futures` 0.1 tasks will be executed by the `tokio` 0.2 thread pool -//! workers, alongside `std::future` tasks. However, they will use the timer and -//! reactor provided by the compatibility background thread. -//! -//! ## Examples -//! -//! Spawning both `tokio` 0.1 and `tokio` 0.2 futures: -//! -//! ```rust -//! use futures_01::future::lazy; -//! -//! tokio_compat::run(lazy(|| { -//! // spawn a `futures` 0.1 future using the `spawn` function from the -//! // `tokio` 0.1 crate: -//! tokio_01::spawn(lazy(|| { -//! println!("hello from tokio 0.1!"); -//! Ok(()) -//! })); -//! -//! // spawn an `async` block future on the same runtime using `tokio` -//! // 0.2's `spawn`: -//! tokio_02::spawn(async { -//! println!("hello from tokio 0.2!"); -//! }); -//! -//! Ok(()) -//! })) -//! ``` -//! -//! Futures on the compat runtime can use `timer` APIs from both 0.1 and 0.2 -//! versions of `tokio`: -//! -//! ```rust -//! # use std::time::{Duration, Instant}; -//! use tokio_compat::prelude::*; -//! -//! tokio_compat::run_std(async { -//! // Wait for a `tokio` 0.1 `Delay`... -//! let when = Instant::now() + Duration::from_millis(10); -//! tokio_01::timer::Delay::new(when) -//! // convert the delay future into a `std::future` that we can `await`. -//! .compat() -//! .await -//! .expect("tokio 0.1 timer should work!"); -//! println!("10 ms have elapsed"); -//! -//! // Wait for a `tokio` 0.2 `Delay`... -//! let when = Instant::now() + Duration::from_millis(20); -//! tokio_02::timer::delay(when).await; -//! println!("20 ms have elapsed"); -//! }); -//! ``` -//! -//! ## Notes -//! -//! In order to allow drop-in compatibility for legacy codebases using -//! `tokio` 0.1, the [`run`], [`spawn`], and [`block_on`] methods provided by the -//! compatibility runtimes take `futures` 0.1 futures. This allows the -//! compatibility runtimes to replace the `tokio` 0.1 runtimes in those codebases -//! without requiring changes to unrelated code. The compatibility runtimes -//! _also_ provide `std::future`-compatible versions of these methods, named -//! [`run_std`], [`spawn_std`], and [`block_on_std`]. -//! -//! Also, please note that the compatibility thread pool runtime does **not** -//! currently support the `tokio` 0.1 [`tokio_threadpool::blocking][blocking] -//! API. Calls to the legacy version of `blocking` made on the compatibility -//! runtime will currently fail. In the future, `tokio-compat` will allow -//! transparently replacing legacy `blocking` with the `tokio` 0.2 blocking -//! APIs, but in the meantime, it will be necessary to convert this code to call -//! into the `tokio` 0.2 version of `blocking`. -//! -//! [`tokio::runtime`]: https://docs.rs/tokio/0.2.0-alpha.6/tokio/runtime/index.html -//! [futures-compat]: https://rust-lang-nursery.github.io/futures-api-docs/0.3.0-alpha.19/futures/compat/index.html -//! [`Timer`]: https://docs.rs/tokio/0.2.0-alpha.6/tokio/timer/index.html -//! [`Runtime`]: struct.Runtime.html -//! [`Reactor`]:https://docs.rs/tokio/0.2.0-alpha.6/tokio/reactor/struct.Reactor.html -//! [timer-01]: https://docs.rs/tokio/0.1.22/tokio/timer/index.html -//! [reactor-01]: https://docs.rs/tokio/0.1.22/tokio/reactor/struct.Reactor.html -//! [`ThreadPool`]: https://docs.rs/tokio-executor/0.2.0-alpha.2/tokio_executor/threadpool/struct.ThreadPool.html -//! [`run`]: struct.Runtime.html#method.run -//! [`spawn`]: struct.Runtime.html#method.spawn -//! [`block_on`]: struct.Runtime.html#method.block_on -//! [`run_std`]: struct.Runtime.html#method.run_std -//! [`spawn_std`]: struct.Runtime.html#method.spawn_std -//! [`block_on_std`]: struct.Runtime.html#method.spawn_std -//! [blocking]: https://docs.rs/tokio-threadpool/0.1.16/tokio_threadpool/fn.blocking.html -mod compat; -pub mod current_thread; -#[cfg(feature = "rt-full")] -mod threadpool; - -#[cfg(feature = "rt-full")] -pub use threadpool::{run, run_std, Builder, Runtime, TaskExecutor}; diff --git a/tokio-compat/src/runtime/threadpool/background.rs b/tokio-compat/src/runtime/threadpool/background.rs deleted file mode 100644 index d3b41ffcc..000000000 --- a/tokio-compat/src/runtime/threadpool/background.rs +++ /dev/null @@ -1,61 +0,0 @@ -//! Temporary reactor + timer that runs on a background thread. This it to make -//! `block_on` work. - -use tokio_02::executor::current_thread::CurrentThread; -use tokio_02::net::driver::{self, Reactor}; -use tokio_02::sync::oneshot; -use tokio_02::timer::clock::Clock; -use tokio_02::timer::timer::{self, Timer}; - -use std::{io, thread}; - -#[derive(Debug)] -pub(crate) struct Background { - reactor_handle: driver::Handle, - timer_handle: timer::Handle, - shutdown_tx: Option>, - thread: Option>, -} - -pub(crate) fn spawn(clock: &Clock) -> io::Result { - let clock = clock.clone(); - - let reactor = Reactor::new()?; - let reactor_handle = reactor.handle(); - - let timer = Timer::new_with_clock(reactor, clock); - let timer_handle = timer.handle(); - - let (shutdown_tx, shutdown_rx) = oneshot::channel(); - let shutdown_tx = Some(shutdown_tx); - - let thread = thread::spawn(move || { - let mut rt = CurrentThread::new_with_park(timer); - let _ = rt.block_on(shutdown_rx); - }); - let thread = Some(thread); - - Ok(Background { - reactor_handle, - timer_handle, - shutdown_tx, - thread, - }) -} - -impl Background { - pub(super) fn reactor(&self) -> &driver::Handle { - &self.reactor_handle - } - - pub(super) fn timer(&self) -> &timer::Handle { - &self.timer_handle - } -} - -impl Drop for Background { - fn drop(&mut self) { - let _ = self.shutdown_tx.take().unwrap().send(()); - let _ = self.thread.take().unwrap().join(); - } -} diff --git a/tokio-compat/src/runtime/threadpool/builder.rs b/tokio-compat/src/runtime/threadpool/builder.rs deleted file mode 100644 index 0f6d58658..000000000 --- a/tokio-compat/src/runtime/threadpool/builder.rs +++ /dev/null @@ -1,285 +0,0 @@ -use super::{background, compat, Inner, Runtime}; - -use tokio_02::executor::thread_pool; -use tokio_02::net::driver::{self, Reactor}; -use tokio_02::timer::clock::{self, Clock}; -use tokio_02::timer::timer::{self, Timer}; -use tokio_executor_01 as executor_01; -use tokio_reactor_01 as reactor_01; -use tokio_timer_02 as timer_02; - -use num_cpus; -use std::io; -use std::sync::{Arc, Barrier, Mutex, RwLock}; - -/// Builds a compatibility runtime with custom configuration values. -/// -/// This runtime is compatible with code using both the current release version -/// of `tokio` (0.1) and with legacy code using `tokio` 0.1. -/// -/// Methods can be chained in order to set the configuration values. The -/// Runtime is constructed by calling [`build`]. -/// -/// New instances of `Builder` are obtained via [`Builder::new`]. -/// -/// See function level documentation for details on the various configuration -/// settings. -/// -/// [`build`]: #method.build -/// [`Builder::new`]: #method.new -/// -/// # Examples -/// -/// ``` -/// -/// use tokio_compat::runtime::Builder; -/// use tokio_02::timer::clock::Clock; -/// -/// fn main() { -/// // build Runtime -/// let runtime = Builder::new() -/// .clock(Clock::system()) -/// .core_threads(4) -/// .name_prefix("my-custom-name-") -/// .stack_size(3 * 1024 * 1024) -/// .build() -/// .unwrap(); -/// -/// // use runtime ... -/// } -/// ``` -#[derive(Debug)] -pub struct Builder { - /// Thread pool specific builder - threadpool_builder: thread_pool::Builder, - - /// The number of worker threads - core_threads: usize, - - /// The clock to use - clock: Clock, -} - -impl Builder { - /// Returns a new runtime builder initialized with default configuration - /// values. - /// - /// Configuration methods can be chained on the return value. - pub fn new() -> Builder { - let core_threads = num_cpus::get().max(1); - - let mut threadpool_builder = thread_pool::Builder::new(); - threadpool_builder.name("tokio-runtime-worker"); - threadpool_builder.num_threads(core_threads); - - Builder { - threadpool_builder, - core_threads, - clock: Clock::new(), - } - } - - /// Set the `Clock` instance that will be used by the runtime. - pub fn clock(&mut self, clock: Clock) -> &mut Self { - self.clock = clock; - self - } - - /// Set the maximum number of worker threads for the `Runtime`'s thread pool. - /// - /// This must be a number between 1 and 32,768 though it is advised to keep - /// this value on the smaller side. - /// - /// The default value is the number of cores available to the system. - /// - /// # Examples - /// - /// ``` - /// # use tokio_compat::runtime; - /// - /// # pub fn main() { - /// let rt = runtime::Builder::new() - /// .core_threads(4) - /// .build() - /// .unwrap(); - /// # } - /// ``` - pub fn core_threads(&mut self, val: usize) -> &mut Self { - self.core_threads = val; - self.threadpool_builder.num_threads(val); - self - } - - /// Set name prefix of threads spawned by the `Runtime`'s thread pool. - /// - /// Thread name prefix is used for generating thread names. For example, if - /// prefix is `my-pool-`, then threads in the pool will get names like - /// `my-pool-1` etc. - /// - /// The default prefix is "tokio-runtime-worker-". - /// - /// # Examples - /// - /// ``` - /// # use tokio_compat::runtime; - /// - /// # pub fn main() { - /// let rt = runtime::Builder::new() - /// .name_prefix("my-pool-") - /// .build(); - /// # } - /// ``` - pub fn name_prefix>(&mut self, val: S) -> &mut Self { - self.threadpool_builder.name(val); - self - } - - /// Set the stack size (in bytes) for worker threads. - /// - /// The actual stack size may be greater than this value if the platform - /// specifies minimal stack size. - /// - /// The default stack size for spawned threads is 2 MiB, though this - /// particular stack size is subject to change in the future. - /// - /// # Examples - /// - /// ``` - /// # use tokio_compat::runtime; - /// - /// # pub fn main() { - /// let rt = runtime::Builder::new() - /// .stack_size(32 * 1024) - /// .build(); - /// # } - /// ``` - pub fn stack_size(&mut self, val: usize) -> &mut Self { - self.threadpool_builder.stack_size(val); - self - } - - /// Create the configured `Runtime`. - /// - /// The returned `ThreadPool` instance is ready to spawn tasks. - /// - /// # Examples - /// - /// ``` - /// # use tokio_compat::runtime::Builder; - /// # pub fn main() { - /// let runtime = Builder::new().build().unwrap(); - /// // ... call runtime.run(...) - /// # let _ = runtime; - /// # } - /// ``` - pub fn build(&mut self) -> io::Result { - self.threadpool_builder.num_threads(self.core_threads); - - let mut reactor_handles = Vec::new(); - let mut timer_handles = Vec::new(); - let mut timers = Vec::new(); - - for _ in 0..self.core_threads { - // Create a new reactor. - let reactor = Reactor::new()?; - reactor_handles.push(reactor.handle()); - - // Create a new timer. - let timer = Timer::new_with_clock(reactor, self.clock.clone()); - timer_handles.push(timer.handle()); - timers.push(Mutex::new(Some(timer))); - } - - // Get a handle to the clock for the runtime. - let clock = self.clock.clone(); - - // Get the current trace dispatcher. - let dispatch = tracing_core::dispatcher::get_default(tracing_core::Dispatch::clone); - let trace = dispatch.clone(); - - let background = background::spawn(&clock)?; - let compat_bg = compat::Background::spawn(&clock)?; - let compat_reactor = compat_bg.reactor().clone(); - let compat_timer = compat_bg.timer().clone(); - - // The `tokio` 0.2 default executor for the worker threads will be set - // by the threadpool itself, but in order to set a default executor for - // `tokio-executor` 0.1 compatibility, we need a `Sender` for the pool - // in the `around_worker` closure. - // - // Unfortunately, we can't get a sender until the pool is constructed, - // which requires the `around_worker` closure. As a workaround, we can - // an `Arc`, which can be moved into the closure, and then be - // set once the pool is constructed. Since the closures won't _run_ - // until we actually try to run futures on the pool, it's okay to set - // the sender after constructing the pool. - // - // The lock is only acquired in `around_worker`; once it is acquired the - // sender is cloned, and the lock doesn't need to be acquired to spawn a - // future. Since we only use it when creating the pool, there shouldn't - // be much of a performance impact. - let compat_sender = Arc::new((RwLock::new(None), Barrier::new(self.core_threads + 1))); - let compat_sender2 = compat_sender.clone(); - - let pool = self - .threadpool_builder - .around_worker(move |index, work| { - let mut enter = executor_01::enter().unwrap(); - // We need the threadpool's sender to set up the default tokio - // 0.1 executor. - let (compat_sender, sender_ready) = &*compat_sender2; - // Wait for the sender to be set. - sender_ready.wait(); - let mut compat_sender = compat_sender - .read() - .unwrap() - .clone() - .expect("compat executor needs to be set before the pool is run!"); - - // Set the default tokio 0.1 reactor to the background compat reactor. - reactor_01::with_default(&compat_reactor, &mut enter, |enter| { - // Set the default tokio 0.2 reactor to this worker thread's - // reactor. - let _reactor = driver::set_default(&reactor_handles[index]); - clock::with_default(&clock, || { - // Set up a default timer for tokio 0.1 compat. - timer_02::with_default(&compat_timer, enter, |enter| { - let _timer = timer::set_default(&timer_handles[index]); - tracing_core::dispatcher::with_default(&dispatch, || { - // Set the default executor for tokio 0.1 compat. - executor_01::with_default(&mut compat_sender, enter, |_enter| { - work(); - }) - }) - }) - }) - }); - }) - .build_with_park(move |index| timers[index].lock().unwrap().take().unwrap()); - - let (idle, idle_rx) = super::idle::Idle::new(); - let runtime = Runtime { - inner: Some(Inner { - pool, - background, - compat_bg, - trace, - }), - idle_rx, - idle, - }; - - // Set the tokio 0.1 executor to be used by the worker threads. - let (compat_sender, sender_ready) = &*compat_sender; - *compat_sender.write().unwrap() = Some(runtime.spawner()); - sender_ready.wait(); - - Ok(runtime) - } -} - -impl Default for Builder { - fn default() -> Self { - Self::new() - } -} diff --git a/tokio-compat/src/runtime/threadpool/idle.rs b/tokio-compat/src/runtime/threadpool/idle.rs deleted file mode 100644 index 266c6de27..000000000 --- a/tokio-compat/src/runtime/threadpool/idle.rs +++ /dev/null @@ -1,54 +0,0 @@ -use std::future::Future; -use std::sync::{ - atomic::{fence, AtomicUsize, Ordering}, - Arc, -}; -use tokio_02::sync::mpsc; - -/// Tracks the number of tasks spawned on a runtime. -/// -/// This is required to implement `shutdown_on_idle` and `tokio::run` APIs that -/// exist in `tokio` 0.1, as the `tokio` 0.2 threadpool does not expose a -/// `shutdown_on_idle` API. -#[derive(Clone, Debug)] -pub(super) struct Idle { - tx: mpsc::Sender<()>, - spawned: Arc, -} - -/// Wraps a future to decrement the spawned count when it completes. -/// -/// This is obtained from `Idle::reserve`. -pub(super) struct Track(Idle); - -impl Idle { - pub(super) fn new() -> (Self, mpsc::Receiver<()>) { - let (tx, rx) = mpsc::channel(1); - let this = Self { - tx, - spawned: Arc::new(AtomicUsize::new(0)), - }; - (this, rx) - } - - /// Prepare to spawn a task on the runtime, incrementing the spawned count. - pub(super) fn reserve(&self) -> Track { - self.spawned.fetch_add(1, Ordering::Relaxed); - Track(self.clone()) - } -} - -impl Track { - /// Run a task, decrementing the spawn count when it completes. - /// - /// If the spawned count is now 0, this sends a notification on the idle channel. - pub(super) async fn with(mut self, f: impl Future) -> T { - let result = f.await; - let spawned = self.0.spawned.fetch_sub(1, Ordering::Release); - if spawned == 1 { - fence(Ordering::Acquire); - let _ = self.0.tx.send(()).await; - } - result - } -} diff --git a/tokio-compat/src/runtime/threadpool/mod.rs b/tokio-compat/src/runtime/threadpool/mod.rs deleted file mode 100644 index c4e6df2d3..000000000 --- a/tokio-compat/src/runtime/threadpool/mod.rs +++ /dev/null @@ -1,545 +0,0 @@ -mod background; -mod builder; -mod idle; -mod task_executor; - -#[allow(unreachable_pub)] // https://github.com/rust-lang/rust/issues/57411 -pub use builder::Builder; -#[allow(unreachable_pub)] // https://github.com/rust-lang/rust/issues/57411 -pub use task_executor::TaskExecutor; - -use super::compat; -use background::Background; - -use futures_01::future::Future as Future01; -use futures_util::{compat::Future01CompatExt, future::FutureExt}; -use std::{future::Future, io, pin::Pin}; -use tokio_02::executor::enter; -use tokio_02::executor::thread_pool::{Spawner, ThreadPool}; -use tokio_02::net::driver; -use tokio_02::timer::timer; -use tokio_executor_01 as executor_01; -use tokio_reactor_01 as reactor_01; -use tokio_timer_02 as timer_02; - -/// A thread pool runtime that can run tasks that use both `tokio` 0.1 and -/// `tokio` 0.2 APIs. -/// -/// This functions similarly to the [`tokio::runtime::Runtime`][rt] struct in the -/// `tokio` crate. However, unlike that runtime, the `tokio-compat` runtime is -/// capable of running both `std::future::Future` tasks that use `tokio` 0.2 -/// runtime services. and `futures` 0.1 tasks that use `tokio` 0.1 runtime -/// services. -/// -/// [rt]: https://docs.rs/tokio/0.2.0-alpha.6/tokio/runtime/struct.Runtime.html -#[derive(Debug)] -pub struct Runtime { - inner: Option, - idle: idle::Idle, - idle_rx: tokio_02::sync::mpsc::Receiver<()>, -} - -#[derive(Debug)] -struct Inner { - /// Task execution pool. - pool: ThreadPool, - - /// Tracing dispatcher - trace: tracing_core::Dispatch, - - /// Maintains a reactor and timer that are always running on a background - /// thread. This is to support `runtime.block_on` w/o requiring the future - /// to be `Send`. - /// - /// A dedicated background thread is required as the threadpool threads - /// might not be running. However, this is a temporary work around. - background: Background, - - /// Compatibility background thread. - /// - /// This maintains a `tokio` 0.1 timer and reactor to support running - /// futures that use older tokio APIs. - compat_bg: compat::Background, -} - -#[derive(Clone, Debug)] -struct CompatSpawner { - inner: S, - idle: idle::Idle, -} - -// ===== impl Runtime ===== - -/// Start the Tokio runtime using the supplied `futures` 0.1 future to bootstrap -/// execution. -/// -/// This function is used to bootstrap the execution of a Tokio application. It -/// does the following: -/// -/// * Start the Tokio runtime using a default configuration. -/// * Spawn the given future onto the thread pool. -/// * Block the current thread until the runtime shuts down. -/// -/// Note that the function will not return immediately once `future` has -/// completed. Instead it waits for the entire runtime to become idle. -/// -/// See the [module level][mod] documentation for more details. -/// -/// # Examples -/// -/// ```rust -/// use futures_01::{Future as Future01, Stream as Stream01}; -/// use tokio_01::net::TcpListener; -/// -/// # fn process(_: T) -> Box + Send> { -/// # unimplemented!(); -/// # } -/// # fn dox() { -/// # let addr = "127.0.0.1:8080".parse().unwrap(); -/// let listener = TcpListener::bind(&addr).unwrap(); -/// -/// let server = listener.incoming() -/// .map_err(|e| println!("error = {:?}", e)) -/// .for_each(|socket| { -/// tokio_01::spawn(process(socket)) -/// }); -/// -/// tokio_compat::run(server); -/// # } -/// # pub fn main() {} -/// ``` -/// -/// # Panics -/// -/// This function panics if called from the context of an executor. -/// -/// [mod]: ../index.html -pub fn run(future: F) -where - F: Future01 + Send + 'static, -{ - run_std(future.compat().map(|_| ())) -} - -/// Start the Tokio runtime using the supplied `std::future` future to bootstrap -/// execution. -/// -/// This function is used to bootstrap the execution of a Tokio application. It -/// does the following: -/// -/// * Start the Tokio runtime using a default configuration. -/// * Spawn the given future onto the thread pool. -/// * Block the current thread until the runtime shuts down. -/// -/// Note that the function will not return immediately once `future` has -/// completed. Instead it waits for the entire runtime to become idle. -/// -/// See the [module level][mod] documentation for more details. -/// -/// # Examples -/// -/// ```rust -/// use futures_01::{Future as Future01, Stream as Stream01}; -/// use tokio_01::net::TcpListener; -/// -/// # fn process(_: T) -> Box + Send> { -/// # unimplemented!(); -/// # } -/// # fn dox() { -/// # let addr = "127.0.0.1:8080".parse().unwrap(); -/// let listener = TcpListener::bind(&addr).unwrap(); -/// -/// let server = listener.incoming() -/// .map_err(|e| println!("error = {:?}", e)) -/// .for_each(|socket| { -/// tokio_01::spawn(process(socket)) -/// }); -/// -/// tokio_compat::run(server); -/// # } -/// # pub fn main() {} -/// ``` -/// -/// # Panics -/// -/// This function panics if called from the context of an executor. -/// -/// [mod]: ../index.html -pub fn run_std(future: F) -where - F: Future + Send + 'static, -{ - // Check enter before creating a new Runtime... - let runtime = Runtime::new().expect("failed to start new Runtime"); - runtime.spawn_std(future); - runtime.shutdown_on_idle(); -} - -impl Runtime { - /// Create a new runtime instance with default configuration values. - /// - /// This results in a reactor, thread pool, and timer being initialized. The - /// thread pool will not spawn any worker threads until it needs to, i.e. - /// tasks are scheduled to run. - /// - /// Most users will not need to call this function directly, instead they - /// will use [`tokio_compat::run`](fn.run.html). - /// - /// See [module level][mod] documentation for more details. - /// - /// # Examples - /// - /// Creating a new `Runtime` with default configuration values. - /// - /// ``` - /// use tokio_compat::runtime::Runtime; - /// - /// let rt = Runtime::new() - /// .unwrap(); - /// - /// // Use the runtime... - /// ``` - /// - /// [mod]: index.html - pub fn new() -> io::Result { - Builder::new().build() - } - - /// Return a handle to the runtime's executor. - /// - /// The returned handle can be used to spawn both `futures` 0.1 and - /// `std::future` tasks that run on this runtime. - /// - /// # Examples - /// - /// ``` - /// use tokio_compat::runtime::Runtime; - /// - /// let rt = Runtime::new() - /// .unwrap(); - /// - /// let executor_handle = rt.executor(); - /// - /// // use `executor_handle` - /// ``` - pub fn executor(&self) -> TaskExecutor { - let inner = self.spawner(); - TaskExecutor { inner } - } - - /// Spawn a `futures` 0.1 future onto the Tokio runtime. - /// - /// This spawns the given future onto the runtime's executor, usually a - /// thread pool. The thread pool is then responsible for polling the future - /// until it completes. - /// - /// See [module level][mod] documentation for more details. - /// - /// [mod]: index.html - /// - /// # Examples - /// - /// ``` - /// use tokio_compat::runtime::Runtime; - /// - /// fn main() { - /// // Create the runtime - /// let rt = Runtime::new().unwrap(); - /// - /// // Spawn a future onto the runtime - /// rt.spawn(futures_01::future::lazy(|| { - /// println!("now running on a worker thread"); - /// Ok(()) - /// })); - /// - /// rt.shutdown_on_idle(); - /// } - /// ``` - /// - /// # Panics - /// - /// This function panics if the spawn fails. Failure occurs if the executor - /// is currently at capacity and is unable to spawn a new future. - pub fn spawn(&self, future: F) -> &Self - where - F: Future01 + Send + 'static, - { - self.spawn_std(future.compat().map(|_| ())) - } - - /// Spawn a `std::future` future onto the Tokio runtime. - /// - /// This spawns the given future onto the runtime's executor, usually a - /// thread pool. The thread pool is then responsible for polling the future - /// until it completes. - /// - /// See [module level][mod] documentation for more details. - /// - /// [mod]: index.html - /// - /// # Examples - /// - /// ``` - /// use tokio_compat::runtime::Runtime; - /// - /// fn main() { - /// // Create the runtime - /// let rt = Runtime::new().unwrap(); - /// - /// // Spawn a future onto the runtime - /// rt.spawn_std(async { - /// println!("now running on a worker thread"); - /// }); - /// - /// rt.shutdown_on_idle(); - /// } - /// ``` - /// - /// # Panics - /// - /// This function panics if the spawn fails. Failure occurs if the executor - /// is currently at capacity and is unable to spawn a new future. - pub fn spawn_std(&self, future: F) -> &Self - where - F: Future + Send + 'static, - { - let idle = self.idle.reserve(); - self.inner().pool.spawn(idle.with(future)); - self - } - - /// Run a `futures` 0.1 future to completion on the Tokio runtime. - /// - /// This runs the given future on the runtime, blocking until it is - /// complete, and yielding its resolved result. Any tasks or timers which - /// the future spawns internally will be executed on the runtime. - /// - /// This method should not be called from an asynchronous context. - /// - /// # Panics - /// - /// This function panics if the executor is at capacity, if the provided - /// future panics, or if called within an asynchronous execution context. - pub fn block_on(&self, future: F) -> Result - where - F: Future01, - { - self.block_on_std(future.compat()) - } - - /// Run a `std::future` future to completion on the Tokio runtime. - /// - /// This runs the given future on the runtime, blocking until it is - /// complete, and yielding its resolved result. Any tasks or timers which - /// the future spawns internally will be executed on the runtime. - /// - /// This method should not be called from an asynchronous context. - /// - /// # Panics - /// - /// This function panics if the executor is at capacity, if the provided - /// future panics, or if called within an asynchronous execution context. - pub fn block_on_std(&self, future: F) -> F::Output - where - F: Future, - { - let mut entered = enter().expect("nested block_on"); - let mut old_entered = executor_01::enter().expect("nested block_on"); - let bg = &self.inner().background; - let trace = &self.inner().trace; - let compat = &self.inner().compat_bg; - - tokio_02::executor::with_default(&mut self.spawner_ref(), || { - executor_01::with_default( - &mut self.spawner_ref(), - &mut old_entered, - |mut old_entered| { - let _reactor = driver::set_default(bg.reactor()); - let _timer = timer::set_default(bg.timer()); - // Set up a default timer for tokio 0.1 compat. - reactor_01::with_default( - compat.reactor(), - &mut old_entered, - |mut old_entered| { - // Set the default tokio 0.2 reactor to this worker thread's - // reactor. - timer_02::with_default( - compat.timer(), - &mut old_entered, - |_old_entered| { - tracing_core::dispatcher::with_default(trace, || { - entered.block_on(future) - }) - }, - ) - }, - ) - }, - ) - }) - } - - /// Signals the runtime to shutdown once it becomes idle. - /// - /// Blocks the current thread until the shutdown operation has completed. - /// This function can be used to perform a graceful shutdown of the runtime. - /// - /// The runtime enters an idle state once **all** of the following occur. - /// - /// * The thread pool has no tasks to execute, i.e., all tasks that were - /// spawned have completed. - /// * The reactor is not managing any I/O resources. - /// - /// See [module level][mod] documentation for more details. - /// - /// # Examples - /// - /// ``` - /// use tokio_compat::runtime::Runtime; - /// - /// let rt = Runtime::new() - /// .unwrap(); - /// - /// // Use the runtime... - /// # rt.spawn_std(async {}); - /// - /// // Shutdown the runtime - /// rt.shutdown_on_idle(); - /// ``` - /// - /// [mod]: index.html - pub fn shutdown_on_idle(mut self) { - use futures_util::stream::StreamExt; - let mut e = tokio_02::executor::enter().unwrap(); - - e.block_on(self.idle_rx.next()); - } - - /// Signals the runtime to shutdown immediately. - /// - /// Blocks the current thread until the shutdown operation has completed. - /// This function will forcibly shutdown the runtime, causing any - /// in-progress work to become canceled. - /// - /// The shutdown steps are: - /// - /// * Drain any scheduled work queues. - /// * Drop any futures that have not yet completed. - /// * Drop the reactor. - /// - /// Once the reactor has dropped, any outstanding I/O resources bound to - /// that reactor will no longer function. Calling any method on them will - /// result in an error. - /// - /// See [module level][mod] documentation for more details. - /// - /// # Examples - /// - /// ``` - /// use tokio_compat::runtime::Runtime; - /// - /// let rt = Runtime::new() - /// .unwrap(); - /// - /// // Use the runtime... - /// - /// // Shutdown the runtime - /// rt.shutdown_now(); - /// ``` - /// - /// [mod]: index.html - #[allow(warnings)] - pub fn shutdown_now(mut self) { - self.inner.take().unwrap().pool.shutdown_now(); - } - - fn spawner(&self) -> CompatSpawner { - CompatSpawner { - inner: self.inner().pool.spawner().clone(), - idle: self.idle.clone(), - } - } - - fn spawner_ref(&self) -> CompatSpawner<&'_ Spawner> { - CompatSpawner { - inner: self.inner().pool.spawner(), - idle: self.idle.clone(), - } - } - - fn inner(&self) -> &Inner { - self.inner.as_ref().unwrap() - } -} - -impl tokio_02::executor::Executor for CompatSpawner<&'_ Spawner> { - fn spawn( - &mut self, - future: Pin + Send>>, - ) -> Result<(), tokio_02::executor::SpawnError> { - let idle = self.idle.reserve(); - self.inner.spawn(idle.with(future)); - Ok(()) - } -} - -impl executor_01::Executor for CompatSpawner<&'_ Spawner> { - fn spawn( - &mut self, - future: Box + Send>, - ) -> Result<(), executor_01::SpawnError> { - let future = future.compat().map(|_| ()); - let idle = self.idle.reserve(); - self.inner.spawn(idle.with(future)); - Ok(()) - } -} - -impl tokio_02::executor::Executor for CompatSpawner { - fn spawn( - &mut self, - future: Pin + Send>>, - ) -> Result<(), tokio_02::executor::SpawnError> { - let idle = self.idle.reserve(); - self.inner.spawn(idle.with(future)); - Ok(()) - } -} - -impl executor_01::Executor for CompatSpawner { - fn spawn( - &mut self, - future: Box + Send>, - ) -> Result<(), executor_01::SpawnError> { - let future = future.compat().map(|_| ()); - let idle = self.idle.reserve(); - self.inner.spawn(idle.with(future)); - Ok(()) - } -} - -impl executor_01::TypedExecutor for CompatSpawner -where - T: Future01 + Send + 'static, -{ - fn spawn(&mut self, future: T) -> Result<(), executor_01::SpawnError> { - let idle = self.idle.reserve(); - let future = Box::pin(idle.with(future.compat().map(|_| ()))); - // Use the `tokio` 0.2 `TypedExecutor` impl so we don't have to box the - // future twice (once to spawn it using `Executor01::spawn` and a second - // time to pin the compat future). - self.inner.spawn(future); - Ok(()) - } -} - -impl Drop for Runtime { - fn drop(&mut self) { - if let Some(inner) = self.inner.take() { - drop(inner); - } - } -} - -#[cfg(test)] -mod tests; diff --git a/tokio-compat/src/runtime/threadpool/task_executor.rs b/tokio-compat/src/runtime/threadpool/task_executor.rs deleted file mode 100644 index 710928be2..000000000 --- a/tokio-compat/src/runtime/threadpool/task_executor.rs +++ /dev/null @@ -1,133 +0,0 @@ -use tokio_02::executor::{thread_pool::Spawner, Executor}; -use tokio_executor_01::{self as executor_01, Executor as Executor01}; - -use futures_01::future::Future as Future01; -use futures_util::{compat::Future01CompatExt, future::FutureExt}; -use std::future::Future; -use std::pin::Pin; - -/// Executes futures on the runtime -/// -/// All futures spawned using this executor will be submitted to the associated -/// Runtime's executor. This executor is usually a thread pool. -/// -/// For more details, see the [module level](index.html) documentation. -#[derive(Debug, Clone)] -pub struct TaskExecutor { - pub(super) inner: super::CompatSpawner, -} - -impl TaskExecutor { - /// Spawn a `futures` 0.1 future onto the Tokio runtime. - /// - /// This spawns the given future onto the runtime's executor, usually a - /// thread pool. The thread pool is then responsible for polling the future - /// until it completes. - /// - /// See [module level][mod] documentation for more details. - /// - /// [mod]: index.html - /// - /// # Examples - /// - /// ``` - /// use tokio_compat::runtime::Runtime; - /// # fn dox() { - /// // Create the runtime - /// let rt = Runtime::new().unwrap(); - /// let executor = rt.executor(); - /// - /// // Spawn a `futures` 0.1 future onto the runtime - /// executor.spawn(futures_01::future::lazy(|| { - /// println!("now running on a worker thread"); - /// Ok(()) - /// })); - /// # } - /// ``` - /// - /// # Panics - /// - /// This function panics if the spawn fails. Failure occurs if the executor - /// is currently at capacity and is unable to spawn a new future. - pub fn spawn(&self, future: F) - where - F: Future01 + Send + 'static, - { - self.spawn_std(Box::pin(future.compat().map(|_| ()))); - } - - /// Spawn a `std::future` future onto the Tokio runtime. - /// - /// This spawns the given future onto the runtime's executor, usually a - /// thread pool. The thread pool is then responsible for polling the future - /// until it completes. - /// - /// See [module level][mod] documentation for more details. - /// - /// [mod]: index.html - /// - /// # Examples - /// - /// ``` - /// use tokio_compat::runtime::Runtime; - /// - /// # fn dox() { - /// // Create the runtime - /// let rt = Runtime::new().unwrap(); - /// let executor = rt.executor(); - /// - /// // Spawn a `std::future` future onto the runtime - /// executor.spawn_std(async { - /// println!("now running on a worker thread"); - /// }); - /// # } - /// ``` - /// - /// # Panics - /// - /// This function panics if the spawn fails. Failure occurs if the executor - /// is currently at capacity and is unable to spawn a new future. - pub fn spawn_std(&self, future: F) - where - F: Future + Send + 'static, - { - let idle = self.inner.idle.reserve(); - self.inner.inner.spawn(idle.with(future)); - } -} - -impl Executor for TaskExecutor { - fn spawn( - &mut self, - future: Pin + Send>>, - ) -> Result<(), tokio_02::executor::SpawnError> { - Executor::spawn(&mut self.inner, future) - } -} - -impl tokio_02::executor::TypedExecutor for TaskExecutor -where - T: Future + Send + 'static, -{ - fn spawn(&mut self, future: T) -> Result<(), tokio_02::executor::SpawnError> { - Executor::spawn(&mut self.inner, Box::pin(future)) - } -} - -impl Executor01 for TaskExecutor { - fn spawn( - &mut self, - future: Box + Send>, - ) -> Result<(), executor_01::SpawnError> { - Executor01::spawn(&mut self.inner, future) - } -} - -impl executor_01::TypedExecutor for TaskExecutor -where - T: Future01 + Send + 'static, -{ - fn spawn(&mut self, future: T) -> Result<(), executor_01::SpawnError> { - executor_01::TypedExecutor::spawn(&mut self.inner, future) - } -} diff --git a/tokio-compat/src/runtime/threadpool/tests.rs b/tokio-compat/src/runtime/threadpool/tests.rs deleted file mode 100644 index bc87b5b9c..000000000 --- a/tokio-compat/src/runtime/threadpool/tests.rs +++ /dev/null @@ -1,115 +0,0 @@ -use std::sync::{ - atomic::{AtomicBool, Ordering}, - Arc, -}; -use std::time::{Duration, Instant}; - -use futures_01::future::Future as Future01; -use futures_util::compat::Future01CompatExt; - -#[test] -fn can_run_01_futures() { - let future_ran = Arc::new(AtomicBool::new(false)); - let ran = future_ran.clone(); - super::run(futures_01::future::lazy(move || { - future_ran.store(true, Ordering::SeqCst); - Ok(()) - })); - assert!(ran.load(Ordering::SeqCst)); -} - -#[test] -fn can_spawn_01_futures() { - let future_ran = Arc::new(AtomicBool::new(false)); - let ran = future_ran.clone(); - super::run(futures_01::future::lazy(move || { - tokio_01::spawn(futures_01::future::lazy(move || { - future_ran.store(true, Ordering::SeqCst); - Ok(()) - })); - Ok(()) - })); - assert!(ran.load(Ordering::SeqCst)); -} - -#[test] -fn can_spawn_std_futures() { - let future_ran = Arc::new(AtomicBool::new(false)); - let ran = future_ran.clone(); - super::run(futures_01::future::lazy(move || { - tokio_02::spawn(async move { - future_ran.store(true, Ordering::SeqCst); - }); - Ok(()) - })); - assert!(ran.load(Ordering::SeqCst)); -} - -#[test] -fn tokio_01_timers_work() { - let future1_ran = Arc::new(AtomicBool::new(false)); - let ran = future1_ran.clone(); - let future1 = futures_01::future::lazy(|| { - let when = Instant::now() + Duration::from_millis(15); - tokio_01::timer::Delay::new(when).map(move |_| when) - }) - .map(move |when| { - ran.store(true, Ordering::SeqCst); - assert!(Instant::now() >= when); - }) - .map_err(|_| panic!("timer should work")); - - let future2_ran = Arc::new(AtomicBool::new(false)); - let ran = future2_ran.clone(); - let future2 = async move { - let when = Instant::now() + Duration::from_millis(10); - tokio_01::timer::Delay::new(when).compat().await.unwrap(); - ran.store(true, Ordering::SeqCst); - assert!(Instant::now() >= when); - }; - - super::run(futures_01::future::lazy(move || { - tokio_02::spawn(future2); - tokio_01::spawn(future1); - Ok(()) - })); - assert!(future1_ran.load(Ordering::SeqCst)); - assert!(future2_ran.load(Ordering::SeqCst)); -} - -#[test] -fn block_on_01_timer() { - let rt = super::Runtime::new().unwrap(); - let when = Instant::now() + Duration::from_millis(10); - rt.block_on(tokio_01::timer::Delay::new(when)).unwrap(); - assert!(Instant::now() >= when); -} - -#[test] -fn block_on_std_01_timer() { - let rt = super::Runtime::new().unwrap(); - let when = Instant::now() + Duration::from_millis(10); - rt.block_on_std(async move { - tokio_01::timer::Delay::new(when).compat().await.unwrap(); - }); - assert!(Instant::now() >= when); -} - -#[test] -fn block_on_01_spawn() { - let rt = super::Runtime::new().unwrap(); - // other tests assert that spawned 0.1 tasks actually *run*, all we care - // is that we're able to spawn it successfully. - rt.block_on(futures_01::future::lazy(|| { - tokio_01::spawn(futures_01::future::lazy(|| Ok(()))) - })) - .unwrap(); -} - -#[test] -fn block_on_std_01_spawn() { - let rt = super::Runtime::new().unwrap(); - // other tests assert that spawned 0.1 tasks actually *run*, all we care - // is that we're able to spawn it successfully. - rt.block_on_std(async { tokio_01::spawn(futures_01::future::lazy(|| Ok(()))) }); -}