mod builder;
mod task_executor;
#[allow(unreachable_pub)] pub use self::task_executor::TaskExecutor;
#[allow(unreachable_pub)] pub use builder::Builder;
use super::{
compat::{self, CompatSpawner},
idle,
};
use futures_01::future::Future as Future01;
use futures_util::{compat::Future01CompatExt, future::FutureExt};
use std::{
fmt,
future::Future,
io,
sync::{Arc, RwLock},
};
use tokio_02::{
runtime::{self, Handle},
task::JoinHandle,
};
use tokio_executor_01 as executor_01;
use tokio_reactor_01 as reactor_01;
use tokio_timer_02 as timer_02;
#[derive(Debug)]
#[cfg_attr(docsrs, doc(cfg(feature = "rt-full")))]
pub struct Runtime {
inner: Option<Inner>,
idle: idle::Idle,
idle_rx: idle::Rx,
compat_sender: Arc<RwLock<Option<CompatSpawner<tokio_02::runtime::Handle>>>>,
}
#[cfg_attr(docsrs, doc(cfg(feature = "rt-full")))]
pub struct Shutdown {
inner: Box<dyn Future01<Item = (), Error = ()> + Send + Sync>,
}
#[derive(Debug)]
#[cfg_attr(docsrs, doc(cfg(feature = "rt-full")))]
struct Inner {
runtime: runtime::Runtime,
compat_bg: compat::Background,
}
#[cfg_attr(docsrs, doc(cfg(feature = "rt-full")))]
pub fn run<F>(future: F)
where
F: Future01<Item = (), Error = ()> + Send + 'static,
{
run_std(future.compat().map(|_| ()))
}
#[cfg_attr(docsrs, doc(cfg(feature = "rt-full")))]
pub fn run_std<F>(future: F)
where
F: Future<Output = ()> + Send + 'static,
{
let runtime = Runtime::new().expect("failed to start new Runtime");
runtime.spawn_std(future);
runtime.shutdown_on_idle().wait().unwrap();
}
impl Runtime {
pub fn new() -> io::Result<Self> {
Builder::new().build()
}
pub fn executor(&self) -> TaskExecutor {
let inner = self.spawner();
TaskExecutor { inner }
}
pub fn spawn<F>(&self, future: F) -> &Self
where
F: Future01<Item = (), Error = ()> + Send + 'static,
{
self.spawn_std(future.compat().map(|_| ()))
}
pub fn spawn_std<F>(&self, future: F) -> &Self
where
F: Future<Output = ()> + Send + 'static,
{
let idle = self.idle.reserve();
self.inner().runtime.spawn(idle.with(future));
self
}
pub fn spawn_handle<F>(&self, future: F) -> JoinHandle<Result<F::Item, F::Error>>
where
F: Future01 + Send + 'static,
F::Item: Send + 'static,
F::Error: Send + 'static,
{
let future = Box::pin(future.compat());
self.spawn_handle_std(future)
}
pub fn spawn_handle_std<F>(&self, future: F) -> JoinHandle<F::Output>
where
F: Future + Send + 'static,
F::Output: Send + 'static,
{
self.inner().runtime.spawn(future)
}
pub fn block_on<F>(&mut self, future: F) -> Result<F::Item, F::Error>
where
F: Future01,
{
self.block_on_std(future.compat())
}
pub fn block_on_std<F>(&mut self, future: F) -> F::Output
where
F: Future,
{
let idle = self.idle.reserve();
let spawner = self.spawner();
let inner = self.inner_mut();
let compat = &inner.compat_bg;
let _timer = timer_02::timer::set_default(compat.timer());
let _reactor = reactor_01::set_default(compat.reactor());
let _executor = executor_01::set_default(spawner);
inner.runtime.block_on(idle.with(future))
}
pub fn shutdown_on_idle(mut self) -> Shutdown {
let spawner = self.spawner();
let Inner {
compat_bg,
mut runtime,
} = self.inner.take().expect("runtime is only shut down once");
let _timer = timer_02::timer::set_default(compat_bg.timer());
let _reactor = reactor_01::set_default(compat_bg.reactor());
let _executor = executor_01::set_default(spawner);
runtime.block_on(self.idle_rx.idle());
let f = futures_01::future::lazy(move || Ok(()));
Shutdown { inner: Box::new(f) }
}
pub fn shutdown_now(mut self) -> Shutdown {
drop(self.inner.take().unwrap());
let f = futures_01::future::lazy(move || Ok(()));
Shutdown { inner: Box::new(f) }
}
fn spawner(&self) -> CompatSpawner<Handle> {
CompatSpawner {
inner: self.inner().runtime.handle().clone(),
idle: self.idle.clone(),
}
}
fn inner(&self) -> &Inner {
self.inner.as_ref().unwrap()
}
fn inner_mut(&mut self) -> &mut Inner {
self.inner.as_mut().unwrap()
}
pub fn enter<F, R>(&self, f: F) -> R
where
F: FnOnce() -> R,
{
let spawner = self.spawner();
let inner = self.inner();
let compat = &inner.compat_bg;
let _timer = timer_02::timer::set_default(compat.timer());
let _reactor = reactor_01::set_default(compat.reactor());
let _executor = executor_01::set_default(spawner);
inner.runtime.enter(f)
}
}
impl Drop for Runtime {
fn drop(&mut self) {
if let Some(inner) = self.inner.take() {
drop(inner);
}
}
}
impl Future01 for Shutdown {
type Item = ();
type Error = ();
fn poll(&mut self) -> futures_01::Poll<Self::Item, Self::Error> {
self.inner.poll()
}
}
impl fmt::Debug for Shutdown {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.pad("Shutdown { .. }")
}
}