use super::{compat, Inner, Runtime};
use std::{
io,
sync::{Arc, RwLock},
};
use tokio_02::runtime;
use tokio_timer_02::clock as clock_02;
#[derive(Debug)]
#[cfg_attr(docsrs, doc(cfg(feature = "rt-full")))]
pub struct Builder {
inner: runtime::Builder,
clock: clock_02::Clock,
}
impl Builder {
pub fn new() -> Builder {
Builder {
clock: clock_02::Clock::system(),
inner: runtime::Builder::new(),
}
}
pub fn clock(&mut self, clock: clock_02::Clock) -> &mut Self {
self.clock = clock;
self
}
pub fn core_threads(&mut self, val: usize) -> &mut Self {
#[allow(deprecated)]
self.inner.num_threads(val);
self
}
pub fn name_prefix<S: Into<String>>(&mut self, val: S) -> &mut Self {
self.inner.thread_name(val);
self
}
pub fn stack_size(&mut self, val: usize) -> &mut Self {
self.inner.thread_stack_size(val);
self
}
pub fn build(&mut self) -> io::Result<Runtime> {
let compat_bg = compat::Background::spawn(&self.clock)?;
let compat_timer = compat_bg.timer().clone();
let compat_reactor = compat_bg.reactor().clone();
let compat_sender: Arc<RwLock<Option<super::CompatSpawner<tokio_02::runtime::Handle>>>> =
Arc::new(RwLock::new(None));
let compat_sender2 = Arc::downgrade(&compat_sender);
let mut lock = compat_sender.write().unwrap();
let runtime = self
.inner
.threaded_scheduler()
.enable_all()
.on_thread_start(move || {
let sender = compat_sender2
.upgrade()
.expect("Runtime dropped but thread started; this is a bug!");
let lock = sender.read().unwrap();
let compat_sender = lock
.as_ref()
.expect("compat executor needs to be set before the pool is run!")
.clone();
compat::set_guards(compat_sender, &compat_timer, &compat_reactor);
})
.on_thread_stop(|| {
compat::unset_guards();
})
.build()?;
let (idle, idle_rx) = super::idle::Idle::new();
*lock = Some(super::CompatSpawner::new(runtime.handle().clone(), &idle));
drop(lock);
let runtime = Runtime {
inner: Some(Inner { runtime, compat_bg }),
idle_rx,
idle,
compat_sender,
};
Ok(runtime)
}
}
impl Default for Builder {
fn default() -> Self {
Self::new()
}
}