From f1f61a3b15d17767b12bfd8c0c5712db1b089b0b Mon Sep 17 00:00:00 2001 From: Carl Lerche Date: Thu, 15 Aug 2019 15:04:21 -0700 Subject: [PATCH] net: reorganize crate in anticipation of #1264 (#1453) Space is made to add `tcp`, `udp`, `uds`, ... modules. --- tokio-fs/examples/std-echo.rs | 2 +- .../reactor.rs => tokio-net/src/driver/mod.rs | 8 +- tokio-net/src/driver/platform.rs | 28 + tokio-net/src/driver/reactor.rs | 532 ++++++++++++++++ tokio-net/src/{ => driver}/registration.rs | 6 +- tokio-net/src/{ => driver}/sharded_rwlock.rs | 0 tokio-net/src/lib.rs | 569 +----------------- tokio-net/src/util/mod.rs | 5 + tokio-net/src/{ => util}/poll_evented.rs | 11 +- tokio-process/src/lib.rs | 2 +- tokio-process/src/unix/mod.rs | 3 +- tokio-process/src/windows.rs | 16 +- tokio-signal/src/ctrl_c.rs | 2 +- tokio-signal/src/unix.rs | 15 +- tokio-signal/src/windows.rs | 2 +- tokio-tcp/src/listener.rs | 9 +- tokio-tcp/src/stream.rs | 7 +- tokio-udp/src/socket.rs | 5 +- tokio-uds/src/datagram.rs | 5 +- tokio-uds/src/listener.rs | 5 +- tokio-uds/src/stream.rs | 5 +- tokio/src/lib.rs | 2 - tokio/src/runtime/current_thread/builder.rs | 2 +- tokio/src/runtime/current_thread/runtime.rs | 8 +- tokio/src/runtime/threadpool/background.rs | 6 +- tokio/src/runtime/threadpool/builder.rs | 4 +- tokio/src/runtime/threadpool/mod.rs | 3 +- tokio/tests/drop-core.rs | 2 +- tokio/tests/reactor.rs | 4 +- ui-tests/tests/features.rs | 5 +- .../tests/ui/tokio_without_net_missing_net.rs | 3 + .../ui/tokio_without_net_missing_net.stderr | 7 + .../ui/tokio_without_net_missing_reactor.rs | 3 - .../tokio_without_net_missing_reactor.stderr | 7 - 34 files changed, 652 insertions(+), 641 deletions(-) rename tokio/src/reactor.rs => tokio-net/src/driver/mod.rs (97%) create mode 100644 tokio-net/src/driver/platform.rs create mode 100644 tokio-net/src/driver/reactor.rs rename tokio-net/src/{ => driver}/registration.rs (99%) rename tokio-net/src/{ => driver}/sharded_rwlock.rs (100%) create mode 100644 tokio-net/src/util/mod.rs rename tokio-net/src/{ => util}/poll_evented.rs (98%) create mode 100644 ui-tests/tests/ui/tokio_without_net_missing_net.rs create mode 100644 ui-tests/tests/ui/tokio_without_net_missing_net.stderr delete mode 100644 ui-tests/tests/ui/tokio_without_net_missing_reactor.rs delete mode 100644 ui-tests/tests/ui/tokio_without_net_missing_reactor.stderr diff --git a/tokio-fs/examples/std-echo.rs b/tokio-fs/examples/std-echo.rs index 3a5eca676..5a46e3529 100644 --- a/tokio-fs/examples/std-echo.rs +++ b/tokio-fs/examples/std-echo.rs @@ -5,8 +5,8 @@ use futures_util::{FutureExt, SinkExt, StreamExt, TryFutureExt}; use tokio::codec::{FramedRead, FramedWrite, LinesCodec, LinesCodecError}; use tokio::future::ready; +use tokio_executor::threadpool::Builder; use tokio_fs::{stderr, stdin, stdout}; -use tokio_threadpool::Builder; #[tokio::main] async fn main() -> Result<(), Box> { diff --git a/tokio/src/reactor.rs b/tokio-net/src/driver/mod.rs similarity index 97% rename from tokio/src/reactor.rs rename to tokio-net/src/driver/mod.rs index 53351e3ea..2ba54d9ad 100644 --- a/tokio/src/reactor.rs +++ b/tokio-net/src/driver/mod.rs @@ -131,4 +131,10 @@ //! [`std::io::Read`]: https://doc.rust-lang.org/std/io/trait.Read.html //! [`std::io::Write`]: https://doc.rust-lang.org/std/io/trait.Write.html -pub use tokio_net::{Handle, PollEvented, Reactor, Registration, Turn}; +pub(crate) mod platform; +mod reactor; +mod registration; +mod sharded_rwlock; + +pub use self::reactor::{set_default, Handle, Reactor}; +pub use self::registration::Registration; diff --git a/tokio-net/src/driver/platform.rs b/tokio-net/src/driver/platform.rs new file mode 100644 index 000000000..4cfe7345b --- /dev/null +++ b/tokio-net/src/driver/platform.rs @@ -0,0 +1,28 @@ +pub(crate) use self::sys::*; + +#[cfg(unix)] +mod sys { + use mio::unix::UnixReady; + use mio::Ready; + + pub(crate) fn hup() -> Ready { + UnixReady::hup().into() + } + + pub(crate) fn is_hup(ready: Ready) -> bool { + UnixReady::from(ready).is_hup() + } +} + +#[cfg(windows)] +mod sys { + use mio::Ready; + + pub(crate) fn hup() -> Ready { + Ready::empty() + } + + pub(crate) fn is_hup(_: Ready) -> bool { + false + } +} diff --git a/tokio-net/src/driver/reactor.rs b/tokio-net/src/driver/reactor.rs new file mode 100644 index 000000000..2aa60374d --- /dev/null +++ b/tokio-net/src/driver/reactor.rs @@ -0,0 +1,532 @@ +use super::platform; +use super::sharded_rwlock::RwLock; + +use tokio_executor::park::{Park, Unpark}; +use tokio_sync::AtomicWaker; + +use log::{debug, log_enabled, trace, Level}; +use mio::event::Evented; +use slab::Slab; +use std::cell::RefCell; +use std::io; +use std::marker::PhantomData; +#[cfg(all(unix, not(target_os = "fuchsia")))] +use std::os::unix::io::{AsRawFd, RawFd}; +use std::sync::atomic::AtomicUsize; +use std::sync::atomic::Ordering::{Relaxed, SeqCst}; +use std::sync::{Arc, Weak}; +use std::task::Waker; +use std::time::{Duration, Instant}; +use std::{fmt, usize}; + +/// The core reactor, or event loop. +/// +/// The event loop is the main source of blocking in an application which drives +/// all other I/O events and notifications happening. Each event loop can have +/// multiple handles pointing to it, each of which can then be used to create +/// various I/O objects to interact with the event loop in interesting ways. +pub struct Reactor { + /// Reuse the `mio::Events` value across calls to poll. + events: mio::Events, + + /// State shared between the reactor and the handles. + inner: Arc, + + _wakeup_registration: mio::Registration, +} + +/// A reference to a reactor. +/// +/// A `Handle` is used for associating I/O objects with an event loop +/// explicitly. Typically though you won't end up using a `Handle` that often +/// and will instead use the default reactor for the execution context. +/// +/// By default, most components bind lazily to reactors. +/// To get this behavior when manually passing a `Handle`, use `default()`. +#[derive(Clone)] +pub struct Handle { + inner: Option, +} + +/// Like `Handle`, but never `None`. +#[derive(Clone)] +pub(crate) struct HandlePriv { + inner: Weak, +} + +/// Return value from the `turn` method on `Reactor`. +/// +/// Currently this value doesn't actually provide any functionality, but it may +/// in the future give insight into what happened during `turn`. +#[derive(Debug)] +pub struct Turn { + _priv: (), +} + +#[test] +fn test_handle_size() { + use std::mem; + assert_eq!(mem::size_of::(), mem::size_of::()); +} + +pub(super) struct Inner { + /// The underlying system event queue. + io: mio::Poll, + + /// ABA guard counter + next_aba_guard: AtomicUsize, + + /// Dispatch slabs for I/O and futures events + pub(super) io_dispatch: RwLock>, + + /// Used to wake up the reactor from a call to `turn` + wakeup: mio::SetReadiness, +} + +pub(super) struct ScheduledIo { + aba_guard: usize, + pub(super) readiness: AtomicUsize, + pub(super) reader: AtomicWaker, + pub(super) writer: AtomicWaker, +} + +#[derive(Debug, Eq, PartialEq, Clone, Copy)] +pub(super) enum Direction { + Read, + Write, +} + +thread_local! { + /// Tracks the reactor for the current execution context. + static CURRENT_REACTOR: RefCell> = RefCell::new(None) +} + +const TOKEN_SHIFT: usize = 22; + +// Kind of arbitrary, but this reserves some token space for later usage. +const MAX_SOURCES: usize = (1 << TOKEN_SHIFT) - 1; +const TOKEN_WAKEUP: mio::Token = mio::Token(MAX_SOURCES); + +fn _assert_kinds() { + fn _assert() {} + + _assert::(); +} + +// ===== impl Reactor ===== + +#[derive(Debug)] +///Guard that resets current reactor on drop. +pub struct DefaultGuard<'a> { + _lifetime: PhantomData<&'a u8>, +} + +impl Drop for DefaultGuard<'_> { + fn drop(&mut self) { + CURRENT_REACTOR.with(|current| { + let mut current = current.borrow_mut(); + *current = None; + }); + } +} + +///Sets handle for a default reactor, returning guard that unsets it on drop. +pub fn set_default(handle: &Handle) -> DefaultGuard<'_> { + CURRENT_REACTOR.with(|current| { + let mut current = current.borrow_mut(); + + assert!( + current.is_none(), + "default Tokio reactor already set \ + for execution context" + ); + + let handle = match handle.as_priv() { + Some(handle) => handle, + None => { + panic!("`handle` does not reference a reactor"); + } + }; + + *current = Some(handle.clone()); + }); + + DefaultGuard { + _lifetime: PhantomData, + } +} + +impl Reactor { + /// Creates a new event loop, returning any error that happened during the + /// creation. + pub fn new() -> io::Result { + let io = mio::Poll::new()?; + let wakeup_pair = mio::Registration::new2(); + + io.register( + &wakeup_pair.0, + TOKEN_WAKEUP, + mio::Ready::readable(), + mio::PollOpt::level(), + )?; + + Ok(Reactor { + events: mio::Events::with_capacity(1024), + _wakeup_registration: wakeup_pair.0, + inner: Arc::new(Inner { + io, + next_aba_guard: AtomicUsize::new(0), + io_dispatch: RwLock::new(Slab::with_capacity(1)), + wakeup: wakeup_pair.1, + }), + }) + } + + /// Returns a handle to this event loop which can be sent across threads + /// and can be used as a proxy to the event loop itself. + /// + /// Handles are cloneable and clones always refer to the same event loop. + /// This handle is typically passed into functions that create I/O objects + /// to bind them to this event loop. + pub fn handle(&self) -> Handle { + Handle { + inner: Some(HandlePriv { + inner: Arc::downgrade(&self.inner), + }), + } + } + + /// Performs one iteration of the event loop, blocking on waiting for events + /// for at most `max_wait` (forever if `None`). + /// + /// This method is the primary method of running this reactor and processing + /// I/O events that occur. This method executes one iteration of an event + /// loop, blocking at most once waiting for events to happen. + /// + /// If a `max_wait` is specified then the method should block no longer than + /// the duration specified, but this shouldn't be used as a super-precise + /// timer but rather a "ballpark approximation" + /// + /// # Return value + /// + /// This function returns an instance of `Turn` + /// + /// `Turn` as of today has no extra information with it and can be safely + /// discarded. In the future `Turn` may contain information about what + /// happened while this reactor blocked. + /// + /// # Errors + /// + /// This function may also return any I/O error which occurs when polling + /// for readiness of I/O objects with the OS. This is quite unlikely to + /// arise and typically mean that things have gone horribly wrong at that + /// point. Currently this is primarily only known to happen for internal + /// bugs to `tokio` itself. + pub fn turn(&mut self, max_wait: Option) -> io::Result { + self.poll(max_wait)?; + Ok(Turn { _priv: () }) + } + + /// Returns true if the reactor is currently idle. + /// + /// Idle is defined as all tasks that have been spawned have completed, + /// either successfully or with an error. + pub fn is_idle(&self) -> bool { + self.inner.io_dispatch.read().is_empty() + } + + fn poll(&mut self, max_wait: Option) -> io::Result<()> { + // Block waiting for an event to happen, peeling out how many events + // happened. + match self.inner.io.poll(&mut self.events, max_wait) { + Ok(_) => {} + Err(e) => return Err(e), + } + + let start = if log_enabled!(Level::Debug) { + Some(Instant::now()) + } else { + None + }; + + // Process all the events that came in, dispatching appropriately + let mut events = 0; + for event in self.events.iter() { + events += 1; + let token = event.token(); + trace!("event {:?} {:?}", event.readiness(), event.token()); + + if token == TOKEN_WAKEUP { + self.inner + .wakeup + .set_readiness(mio::Ready::empty()) + .unwrap(); + } else { + self.dispatch(token, event.readiness()); + } + } + + if let Some(start) = start { + let dur = start.elapsed(); + trace!( + "loop process - {} events, {}.{:03}s", + events, + dur.as_secs(), + dur.subsec_millis() + ); + } + + Ok(()) + } + + fn dispatch(&self, token: mio::Token, ready: mio::Ready) { + let aba_guard = token.0 & !MAX_SOURCES; + let token = token.0 & MAX_SOURCES; + + let mut rd = None; + let mut wr = None; + + // Create a scope to ensure that notifying the tasks stays out of the + // lock's critical section. + { + let io_dispatch = self.inner.io_dispatch.read(); + + let io = match io_dispatch.get(token) { + Some(io) => io, + None => return, + }; + + if aba_guard != io.aba_guard { + return; + } + + io.readiness.fetch_or(ready.as_usize(), Relaxed); + + if ready.is_writable() || platform::is_hup(ready) { + wr = io.writer.take_waker(); + } + + if !(ready & (!mio::Ready::writable())).is_empty() { + rd = io.reader.take_waker(); + } + } + + if let Some(w) = rd { + w.wake(); + } + + if let Some(w) = wr { + w.wake(); + } + } +} + +#[cfg(all(unix, not(target_os = "fuchsia")))] +impl AsRawFd for Reactor { + fn as_raw_fd(&self) -> RawFd { + self.inner.io.as_raw_fd() + } +} + +impl Park for Reactor { + type Unpark = Handle; + type Error = io::Error; + + fn unpark(&self) -> Self::Unpark { + self.handle() + } + + fn park(&mut self) -> io::Result<()> { + self.turn(None)?; + Ok(()) + } + + fn park_timeout(&mut self, duration: Duration) -> io::Result<()> { + self.turn(Some(duration))?; + Ok(()) + } +} + +impl fmt::Debug for Reactor { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "Reactor") + } +} + +// ===== impl Handle ===== + +impl Handle { + #[doc(hidden)] + #[deprecated(note = "semantics were sometimes surprising, use Handle::default()")] + pub fn current() -> Handle { + // TODO: Should this panic on error? + HandlePriv::try_current() + .map(|handle| Handle { + inner: Some(handle), + }) + .unwrap_or(Handle { + inner: Some(HandlePriv { inner: Weak::new() }), + }) + } + + pub(crate) fn as_priv(&self) -> Option<&HandlePriv> { + self.inner.as_ref() + } +} + +impl Unpark for Handle { + fn unpark(&self) { + if let Some(ref h) = self.inner { + h.wakeup(); + } + } +} + +impl Default for Handle { + /// Returns a "default" handle, i.e., a handle that lazily binds to a reactor. + fn default() -> Handle { + Handle { inner: None } + } +} + +impl fmt::Debug for Handle { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "Handle") + } +} + +// ===== impl HandlePriv ===== + +impl HandlePriv { + /// Try to get a handle to the current reactor. + /// + /// Returns `Err` if no handle is found. + pub(super) fn try_current() -> io::Result { + CURRENT_REACTOR.with(|current| match *current.borrow() { + Some(ref handle) => Ok(handle.clone()), + None => Err(io::Error::new(io::ErrorKind::Other, "no current reactor")), + }) + } + + /// Forces a reactor blocked in a call to `turn` to wakeup, or otherwise + /// makes the next call to `turn` return immediately. + /// + /// This method is intended to be used in situations where a notification + /// needs to otherwise be sent to the main reactor. If the reactor is + /// currently blocked inside of `turn` then it will wake up and soon return + /// after this method has been called. If the reactor is not currently + /// blocked in `turn`, then the next call to `turn` will not block and + /// return immediately. + fn wakeup(&self) { + if let Some(inner) = self.inner() { + inner.wakeup.set_readiness(mio::Ready::readable()).unwrap(); + } + } + + pub(super) fn inner(&self) -> Option> { + self.inner.upgrade() + } +} + +impl fmt::Debug for HandlePriv { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "HandlePriv") + } +} + +// ===== impl Inner ===== + +impl Inner { + /// Register an I/O resource with the reactor. + /// + /// The registration token is returned. + pub(super) fn add_source(&self, source: &dyn Evented) -> io::Result { + // Get an ABA guard value + let aba_guard = self.next_aba_guard.fetch_add(1 << TOKEN_SHIFT, Relaxed); + + let key = { + // Block to contain the write lock + let mut io_dispatch = self.io_dispatch.write(); + + if io_dispatch.len() == MAX_SOURCES { + return Err(io::Error::new( + io::ErrorKind::Other, + "reactor at max \ + registered I/O resources", + )); + } + + io_dispatch.insert(ScheduledIo { + aba_guard, + readiness: AtomicUsize::new(0), + reader: AtomicWaker::new(), + writer: AtomicWaker::new(), + }) + }; + + let token = aba_guard | key; + debug!("adding I/O source: {}", token); + + self.io.register( + source, + mio::Token(token), + mio::Ready::all(), + mio::PollOpt::edge(), + )?; + + Ok(key) + } + + /// Deregisters an I/O resource from the reactor. + pub(super) fn deregister_source(&self, source: &dyn Evented) -> io::Result<()> { + self.io.deregister(source) + } + + pub(super) fn drop_source(&self, token: usize) { + debug!("dropping I/O source: {}", token); + self.io_dispatch.write().remove(token); + } + + /// Registers interest in the I/O resource associated with `token`. + pub(super) fn register(&self, token: usize, dir: Direction, w: Waker) { + debug!("scheduling {:?} for: {}", dir, token); + let io_dispatch = self.io_dispatch.read(); + let sched = io_dispatch.get(token).unwrap(); + + let (waker, ready) = match dir { + Direction::Read => (&sched.reader, !mio::Ready::writable()), + Direction::Write => (&sched.writer, mio::Ready::writable()), + }; + + waker.register(w); + + if sched.readiness.load(SeqCst) & ready.as_usize() != 0 { + waker.wake(); + } + } +} + +impl Drop for Inner { + fn drop(&mut self) { + // When a reactor is dropped it needs to wake up all blocked tasks as + // they'll never receive a notification, and all connected I/O objects + // will start returning errors pretty quickly. + let io = self.io_dispatch.read(); + for (_, io) in io.iter() { + io.writer.wake(); + io.reader.wake(); + } + } +} + +impl Direction { + pub(super) fn mask(self) -> mio::Ready { + match self { + Direction::Read => { + // Everything except writable is signaled through read. + mio::Ready::all() - mio::Ready::writable() + } + Direction::Write => mio::Ready::writable() | platform::hup(), + } + } +} diff --git a/tokio-net/src/registration.rs b/tokio-net/src/driver/registration.rs similarity index 99% rename from tokio-net/src/registration.rs rename to tokio-net/src/driver/registration.rs index d094e2183..eb24d8f78 100644 --- a/tokio-net/src/registration.rs +++ b/tokio-net/src/driver/registration.rs @@ -1,4 +1,6 @@ -use crate::{Direction, Handle, HandlePriv}; +use super::platform; +use super::reactor::{Direction, Handle, HandlePriv}; + use log::debug; use mio::{self, Evented}; use std::cell::UnsafeCell; @@ -504,7 +506,7 @@ impl Inner { }; let mask = direction.mask(); - let mask_no_hup = (mask - crate::platform::hup()).as_usize(); + let mask_no_hup = (mask - platform::hup()).as_usize(); let io_dispatch = inner.io_dispatch.read(); let sched = &io_dispatch[self.token]; diff --git a/tokio-net/src/sharded_rwlock.rs b/tokio-net/src/driver/sharded_rwlock.rs similarity index 100% rename from tokio-net/src/sharded_rwlock.rs rename to tokio-net/src/driver/sharded_rwlock.rs diff --git a/tokio-net/src/lib.rs b/tokio-net/src/lib.rs index 09a9d5e74..c13933003 100644 --- a/tokio-net/src/lib.rs +++ b/tokio-net/src/lib.rs @@ -36,570 +36,5 @@ //! [`PollEvented`]: struct.PollEvented.html //! [reactor module]: https://docs.rs/tokio/0.1/tokio/reactor/index.html -mod poll_evented; -mod registration; -mod sharded_rwlock; - -// ===== Public re-exports ===== - -pub use self::poll_evented::PollEvented; -pub use self::registration::Registration; - -// ===== Private imports ===== - -use crate::sharded_rwlock::RwLock; -use log::{debug, log_enabled, trace, Level}; -use mio::event::Evented; -use slab::Slab; -use std::cell::RefCell; -use std::io; -use std::marker::PhantomData; -#[cfg(all(unix, not(target_os = "fuchsia")))] -use std::os::unix::io::{AsRawFd, RawFd}; -use std::sync::atomic::AtomicUsize; -use std::sync::atomic::Ordering::{Relaxed, SeqCst}; -use std::sync::{Arc, Weak}; -use std::task::Waker; -use std::time::{Duration, Instant}; -use std::{fmt, usize}; -use tokio_executor::park::{Park, Unpark}; -use tokio_sync::AtomicWaker; - -/// The core reactor, or event loop. -/// -/// The event loop is the main source of blocking in an application which drives -/// all other I/O events and notifications happening. Each event loop can have -/// multiple handles pointing to it, each of which can then be used to create -/// various I/O objects to interact with the event loop in interesting ways. -pub struct Reactor { - /// Reuse the `mio::Events` value across calls to poll. - events: mio::Events, - - /// State shared between the reactor and the handles. - inner: Arc, - - _wakeup_registration: mio::Registration, -} - -/// A reference to a reactor. -/// -/// A `Handle` is used for associating I/O objects with an event loop -/// explicitly. Typically though you won't end up using a `Handle` that often -/// and will instead use the default reactor for the execution context. -/// -/// By default, most components bind lazily to reactors. -/// To get this behavior when manually passing a `Handle`, use `default()`. -#[derive(Clone)] -pub struct Handle { - inner: Option, -} - -/// Like `Handle`, but never `None`. -#[derive(Clone)] -struct HandlePriv { - inner: Weak, -} - -/// Return value from the `turn` method on `Reactor`. -/// -/// Currently this value doesn't actually provide any functionality, but it may -/// in the future give insight into what happened during `turn`. -#[derive(Debug)] -pub struct Turn { - _priv: (), -} - -#[test] -fn test_handle_size() { - use std::mem; - assert_eq!(mem::size_of::(), mem::size_of::()); -} - -struct Inner { - /// The underlying system event queue. - io: mio::Poll, - - /// ABA guard counter - next_aba_guard: AtomicUsize, - - /// Dispatch slabs for I/O and futures events - io_dispatch: RwLock>, - - /// Used to wake up the reactor from a call to `turn` - wakeup: mio::SetReadiness, -} - -struct ScheduledIo { - aba_guard: usize, - readiness: AtomicUsize, - reader: AtomicWaker, - writer: AtomicWaker, -} - -#[derive(Debug, Eq, PartialEq, Clone, Copy)] -pub(crate) enum Direction { - Read, - Write, -} - -thread_local! { - /// Tracks the reactor for the current execution context. - static CURRENT_REACTOR: RefCell> = RefCell::new(None) -} - -const TOKEN_SHIFT: usize = 22; - -// Kind of arbitrary, but this reserves some token space for later usage. -const MAX_SOURCES: usize = (1 << TOKEN_SHIFT) - 1; -const TOKEN_WAKEUP: mio::Token = mio::Token(MAX_SOURCES); - -fn _assert_kinds() { - fn _assert() {} - - _assert::(); -} - -// ===== impl Reactor ===== - -#[derive(Debug)] -///Guard that resets current reactor on drop. -pub struct DefaultGuard<'a> { - _lifetime: PhantomData<&'a u8>, -} - -impl Drop for DefaultGuard<'_> { - fn drop(&mut self) { - CURRENT_REACTOR.with(|current| { - let mut current = current.borrow_mut(); - *current = None; - }); - } -} - -///Sets handle for a default reactor, returning guard that unsets it on drop. -pub fn set_default(handle: &Handle) -> DefaultGuard<'_> { - CURRENT_REACTOR.with(|current| { - let mut current = current.borrow_mut(); - - assert!( - current.is_none(), - "default Tokio reactor already set \ - for execution context" - ); - - let handle = match handle.as_priv() { - Some(handle) => handle, - None => { - panic!("`handle` does not reference a reactor"); - } - }; - - *current = Some(handle.clone()); - }); - - DefaultGuard { - _lifetime: PhantomData, - } -} - -impl Reactor { - /// Creates a new event loop, returning any error that happened during the - /// creation. - pub fn new() -> io::Result { - let io = mio::Poll::new()?; - let wakeup_pair = mio::Registration::new2(); - - io.register( - &wakeup_pair.0, - TOKEN_WAKEUP, - mio::Ready::readable(), - mio::PollOpt::level(), - )?; - - Ok(Reactor { - events: mio::Events::with_capacity(1024), - _wakeup_registration: wakeup_pair.0, - inner: Arc::new(Inner { - io, - next_aba_guard: AtomicUsize::new(0), - io_dispatch: RwLock::new(Slab::with_capacity(1)), - wakeup: wakeup_pair.1, - }), - }) - } - - /// Returns a handle to this event loop which can be sent across threads - /// and can be used as a proxy to the event loop itself. - /// - /// Handles are cloneable and clones always refer to the same event loop. - /// This handle is typically passed into functions that create I/O objects - /// to bind them to this event loop. - pub fn handle(&self) -> Handle { - Handle { - inner: Some(HandlePriv { - inner: Arc::downgrade(&self.inner), - }), - } - } - - /// Performs one iteration of the event loop, blocking on waiting for events - /// for at most `max_wait` (forever if `None`). - /// - /// This method is the primary method of running this reactor and processing - /// I/O events that occur. This method executes one iteration of an event - /// loop, blocking at most once waiting for events to happen. - /// - /// If a `max_wait` is specified then the method should block no longer than - /// the duration specified, but this shouldn't be used as a super-precise - /// timer but rather a "ballpark approximation" - /// - /// # Return value - /// - /// This function returns an instance of `Turn` - /// - /// `Turn` as of today has no extra information with it and can be safely - /// discarded. In the future `Turn` may contain information about what - /// happened while this reactor blocked. - /// - /// # Errors - /// - /// This function may also return any I/O error which occurs when polling - /// for readiness of I/O objects with the OS. This is quite unlikely to - /// arise and typically mean that things have gone horribly wrong at that - /// point. Currently this is primarily only known to happen for internal - /// bugs to `tokio` itself. - pub fn turn(&mut self, max_wait: Option) -> io::Result { - self.poll(max_wait)?; - Ok(Turn { _priv: () }) - } - - /// Returns true if the reactor is currently idle. - /// - /// Idle is defined as all tasks that have been spawned have completed, - /// either successfully or with an error. - pub fn is_idle(&self) -> bool { - self.inner.io_dispatch.read().is_empty() - } - - fn poll(&mut self, max_wait: Option) -> io::Result<()> { - // Block waiting for an event to happen, peeling out how many events - // happened. - match self.inner.io.poll(&mut self.events, max_wait) { - Ok(_) => {} - Err(e) => return Err(e), - } - - let start = if log_enabled!(Level::Debug) { - Some(Instant::now()) - } else { - None - }; - - // Process all the events that came in, dispatching appropriately - let mut events = 0; - for event in self.events.iter() { - events += 1; - let token = event.token(); - trace!("event {:?} {:?}", event.readiness(), event.token()); - - if token == TOKEN_WAKEUP { - self.inner - .wakeup - .set_readiness(mio::Ready::empty()) - .unwrap(); - } else { - self.dispatch(token, event.readiness()); - } - } - - if let Some(start) = start { - let dur = start.elapsed(); - trace!( - "loop process - {} events, {}.{:03}s", - events, - dur.as_secs(), - dur.subsec_millis() - ); - } - - Ok(()) - } - - fn dispatch(&self, token: mio::Token, ready: mio::Ready) { - let aba_guard = token.0 & !MAX_SOURCES; - let token = token.0 & MAX_SOURCES; - - let mut rd = None; - let mut wr = None; - - // Create a scope to ensure that notifying the tasks stays out of the - // lock's critical section. - { - let io_dispatch = self.inner.io_dispatch.read(); - - let io = match io_dispatch.get(token) { - Some(io) => io, - None => return, - }; - - if aba_guard != io.aba_guard { - return; - } - - io.readiness.fetch_or(ready.as_usize(), Relaxed); - - if ready.is_writable() || platform::is_hup(ready) { - wr = io.writer.take_waker(); - } - - if !(ready & (!mio::Ready::writable())).is_empty() { - rd = io.reader.take_waker(); - } - } - - if let Some(w) = rd { - w.wake(); - } - - if let Some(w) = wr { - w.wake(); - } - } -} - -#[cfg(all(unix, not(target_os = "fuchsia")))] -impl AsRawFd for Reactor { - fn as_raw_fd(&self) -> RawFd { - self.inner.io.as_raw_fd() - } -} - -impl Park for Reactor { - type Unpark = Handle; - type Error = io::Error; - - fn unpark(&self) -> Self::Unpark { - self.handle() - } - - fn park(&mut self) -> io::Result<()> { - self.turn(None)?; - Ok(()) - } - - fn park_timeout(&mut self, duration: Duration) -> io::Result<()> { - self.turn(Some(duration))?; - Ok(()) - } -} - -impl fmt::Debug for Reactor { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - write!(f, "Reactor") - } -} - -// ===== impl Handle ===== - -impl Handle { - #[doc(hidden)] - #[deprecated(note = "semantics were sometimes surprising, use Handle::default()")] - pub fn current() -> Handle { - // TODO: Should this panic on error? - HandlePriv::try_current() - .map(|handle| Handle { - inner: Some(handle), - }) - .unwrap_or(Handle { - inner: Some(HandlePriv { inner: Weak::new() }), - }) - } - - fn as_priv(&self) -> Option<&HandlePriv> { - self.inner.as_ref() - } -} - -impl Unpark for Handle { - fn unpark(&self) { - if let Some(ref h) = self.inner { - h.wakeup(); - } - } -} - -impl Default for Handle { - /// Returns a "default" handle, i.e., a handle that lazily binds to a reactor. - fn default() -> Handle { - Handle { inner: None } - } -} - -impl fmt::Debug for Handle { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - write!(f, "Handle") - } -} - -// ===== impl HandlePriv ===== - -impl HandlePriv { - /// Try to get a handle to the current reactor. - /// - /// Returns `Err` if no handle is found. - pub(crate) fn try_current() -> io::Result { - CURRENT_REACTOR.with(|current| match *current.borrow() { - Some(ref handle) => Ok(handle.clone()), - None => Err(io::Error::new(io::ErrorKind::Other, "no current reactor")), - }) - } - - /// Forces a reactor blocked in a call to `turn` to wakeup, or otherwise - /// makes the next call to `turn` return immediately. - /// - /// This method is intended to be used in situations where a notification - /// needs to otherwise be sent to the main reactor. If the reactor is - /// currently blocked inside of `turn` then it will wake up and soon return - /// after this method has been called. If the reactor is not currently - /// blocked in `turn`, then the next call to `turn` will not block and - /// return immediately. - fn wakeup(&self) { - if let Some(inner) = self.inner() { - inner.wakeup.set_readiness(mio::Ready::readable()).unwrap(); - } - } - - fn inner(&self) -> Option> { - self.inner.upgrade() - } -} - -impl fmt::Debug for HandlePriv { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - write!(f, "HandlePriv") - } -} - -// ===== impl Inner ===== - -impl Inner { - /// Register an I/O resource with the reactor. - /// - /// The registration token is returned. - fn add_source(&self, source: &dyn Evented) -> io::Result { - // Get an ABA guard value - let aba_guard = self.next_aba_guard.fetch_add(1 << TOKEN_SHIFT, Relaxed); - - let key = { - // Block to contain the write lock - let mut io_dispatch = self.io_dispatch.write(); - - if io_dispatch.len() == MAX_SOURCES { - return Err(io::Error::new( - io::ErrorKind::Other, - "reactor at max \ - registered I/O resources", - )); - } - - io_dispatch.insert(ScheduledIo { - aba_guard, - readiness: AtomicUsize::new(0), - reader: AtomicWaker::new(), - writer: AtomicWaker::new(), - }) - }; - - let token = aba_guard | key; - debug!("adding I/O source: {}", token); - - self.io.register( - source, - mio::Token(token), - mio::Ready::all(), - mio::PollOpt::edge(), - )?; - - Ok(key) - } - - /// Deregisters an I/O resource from the reactor. - fn deregister_source(&self, source: &dyn Evented) -> io::Result<()> { - self.io.deregister(source) - } - - fn drop_source(&self, token: usize) { - debug!("dropping I/O source: {}", token); - self.io_dispatch.write().remove(token); - } - - /// Registers interest in the I/O resource associated with `token`. - fn register(&self, token: usize, dir: Direction, w: Waker) { - debug!("scheduling {:?} for: {}", dir, token); - let io_dispatch = self.io_dispatch.read(); - let sched = io_dispatch.get(token).unwrap(); - - let (waker, ready) = match dir { - Direction::Read => (&sched.reader, !mio::Ready::writable()), - Direction::Write => (&sched.writer, mio::Ready::writable()), - }; - - waker.register(w); - - if sched.readiness.load(SeqCst) & ready.as_usize() != 0 { - waker.wake(); - } - } -} - -impl Drop for Inner { - fn drop(&mut self) { - // When a reactor is dropped it needs to wake up all blocked tasks as - // they'll never receive a notification, and all connected I/O objects - // will start returning errors pretty quickly. - let io = self.io_dispatch.read(); - for (_, io) in io.iter() { - io.writer.wake(); - io.reader.wake(); - } - } -} - -impl Direction { - fn mask(self) -> mio::Ready { - match self { - Direction::Read => { - // Everything except writable is signaled through read. - mio::Ready::all() - mio::Ready::writable() - } - Direction::Write => mio::Ready::writable() | platform::hup(), - } - } -} - -#[cfg(unix)] -mod platform { - use mio::unix::UnixReady; - use mio::Ready; - - pub(crate) fn hup() -> Ready { - UnixReady::hup().into() - } - - pub(crate) fn is_hup(ready: Ready) -> bool { - UnixReady::from(ready).is_hup() - } -} - -#[cfg(windows)] -mod platform { - use mio::Ready; - - pub(crate) fn hup() -> Ready { - Ready::empty() - } - - pub(crate) fn is_hup(_: Ready) -> bool { - false - } -} +pub mod driver; +pub mod util; diff --git a/tokio-net/src/util/mod.rs b/tokio-net/src/util/mod.rs new file mode 100644 index 000000000..a3ec246b0 --- /dev/null +++ b/tokio-net/src/util/mod.rs @@ -0,0 +1,5 @@ +//! Utilities for implementing networking types. + +mod poll_evented; + +pub use self::poll_evented::PollEvented; diff --git a/tokio-net/src/poll_evented.rs b/tokio-net/src/util/poll_evented.rs similarity index 98% rename from tokio-net/src/poll_evented.rs rename to tokio-net/src/util/poll_evented.rs index 145a3bc9a..29dd94fca 100644 --- a/tokio-net/src/poll_evented.rs +++ b/tokio-net/src/util/poll_evented.rs @@ -1,4 +1,4 @@ -use crate::{Handle, Registration}; +use crate::driver::{platform, Handle, Registration}; use tokio_io::{AsyncRead, AsyncWrite}; @@ -55,7 +55,7 @@ use std::task::{Context, Poll}; /// [`clear_read_ready`]. /// /// ```rust -/// use tokio_net::PollEvented; +/// use tokio_net::util::PollEvented; /// /// use futures_core::ready; /// use mio::Ready; @@ -125,7 +125,7 @@ macro_rules! poll_ready { // Load cached & encoded readiness. let mut cached = $me.inner.$cache.load(Relaxed); - let mask = $mask | crate::platform::hup(); + let mask = $mask | platform::hup(); // See if the current readiness matches any bits. let mut ret = mio::Ready::from_usize(cached) & $mask; @@ -272,10 +272,7 @@ where pub fn clear_read_ready(&self, cx: &mut Context<'_>, ready: mio::Ready) -> io::Result<()> { // Cannot clear write readiness assert!(!ready.is_writable(), "cannot clear write readiness"); - assert!( - !crate::platform::is_hup(ready), - "cannot clear HUP readiness" - ); + assert!(!platform::is_hup(ready), "cannot clear HUP readiness"); self.inner .read_readiness diff --git a/tokio-process/src/lib.rs b/tokio-process/src/lib.rs index d986d2a30..42c5f386c 100644 --- a/tokio-process/src/lib.rs +++ b/tokio-process/src/lib.rs @@ -133,7 +133,7 @@ extern crate lazy_static; extern crate log; use tokio_io::{AsyncRead, AsyncReadExt, AsyncWrite}; -use tokio_net::Handle; +use tokio_net::driver::Handle; use futures_core::future::TryFuture; use futures_util::future; diff --git a/tokio-process/src/unix/mod.rs b/tokio-process/src/unix/mod.rs index d4b1c73d6..42aaca18c 100644 --- a/tokio-process/src/unix/mod.rs +++ b/tokio-process/src/unix/mod.rs @@ -29,7 +29,8 @@ use self::reap::Reaper; use super::SpawnedChild; use crate::kill::Kill; -use tokio_net::{Handle, PollEvented}; +use tokio_net::driver::Handle; +use tokio_net::util::PollEvented; use tokio_signal::unix::{Signal, SignalKind}; use mio::event::Evented; diff --git a/tokio-process/src/windows.rs b/tokio-process/src/windows.rs index 28a0df723..68ab15fa1 100644 --- a/tokio-process/src/windows.rs +++ b/tokio-process/src/windows.rs @@ -15,8 +15,16 @@ //! `RegisterWaitForSingleObject` and then wait on the other end of the oneshot //! from then on out. +use super::SpawnedChild; use crate::kill::Kill; +use tokio_net::driver::Handle; +use tokio_net::util::PollEvented; +use tokio_sync::oneshot; + +use futures_util::future::Fuse; +use futures_util::future::FutureExt; +use mio_named_pipes::NamedPipe; use std::fmt; use std::future::Future; use std::io; @@ -27,14 +35,6 @@ use std::process::{self, ExitStatus}; use std::ptr; use std::task::Context; use std::task::Poll; - -use futures_util::future::Fuse; -use futures_util::future::FutureExt; - -use super::SpawnedChild; -use mio_named_pipes::NamedPipe; -use tokio_net::{Handle, PollEvented}; -use tokio_sync::oneshot; use winapi::shared::minwindef::*; use winapi::shared::winerror::*; use winapi::um::handleapi::*; diff --git a/tokio-signal/src/ctrl_c.rs b/tokio-signal/src/ctrl_c.rs index 1f2ddc41b..46ae92583 100644 --- a/tokio-signal/src/ctrl_c.rs +++ b/tokio-signal/src/ctrl_c.rs @@ -3,7 +3,7 @@ use crate::unix::Signal as Inner; #[cfg(windows)] use crate::windows::Event as Inner; -use tokio_net::Handle; +use tokio_net::driver::Handle; use futures_core::stream::Stream; use std::io; diff --git a/tokio-signal/src/unix.rs b/tokio-signal/src/unix.rs index f05d7e186..77c994471 100644 --- a/tokio-signal/src/unix.rs +++ b/tokio-signal/src/unix.rs @@ -7,19 +7,20 @@ pub use libc; -use std::io::{self, Error, ErrorKind, Write}; -use std::pin::Pin; -use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::Once; +use tokio_io::AsyncRead; +use tokio_net::driver::Handle; +use tokio_net::util::PollEvented; +use tokio_sync::mpsc::{channel, Receiver}; use futures_core::stream::Stream; use libc::c_int; use mio_uds::UnixStream; use std::future::Future; +use std::io::{self, Error, ErrorKind, Write}; +use std::pin::Pin; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Once; use std::task::{Context, Poll}; -use tokio_io::AsyncRead; -use tokio_net::{Handle, PollEvented}; -use tokio_sync::mpsc::{channel, Receiver}; use crate::registry::{globals, EventId, EventInfo, Globals, Init, Storage}; diff --git a/tokio-signal/src/windows.rs b/tokio-signal/src/windows.rs index b49361f01..2debadc20 100644 --- a/tokio-signal/src/windows.rs +++ b/tokio-signal/src/windows.rs @@ -9,7 +9,7 @@ use crate::registry::{globals, EventId, EventInfo, Init, Storage}; -use tokio_net::Handle; +use tokio_net::driver::Handle; use tokio_sync::mpsc::{channel, Receiver}; use futures_core::stream::Stream; diff --git a/tokio-tcp/src/listener.rs b/tokio-tcp/src/listener.rs index 8fc7051b0..ac928c12c 100644 --- a/tokio-tcp/src/listener.rs +++ b/tokio-tcp/src/listener.rs @@ -1,7 +1,9 @@ #[cfg(feature = "async-traits")] use super::incoming::Incoming; use super::TcpStream; -use tokio_net::{Handle, PollEvented}; + +use tokio_net::driver::Handle; +use tokio_net::util::PollEvented; use futures_core::ready; use futures_util::future::poll_fn; @@ -153,8 +155,9 @@ impl TcpListener { /// /// ```no_run /// use tokio::net::TcpListener; + /// use tokio_net::driver::Handle; + /// /// use std::net::TcpListener as StdTcpListener; - /// use tokio::reactor::Handle; /// /// let std_listener = StdTcpListener::bind("127.0.0.1:0")?; /// let listener = TcpListener::from_std(std_listener, &Handle::default())?; @@ -260,7 +263,7 @@ impl TryFrom for mio::net::TcpListener { /// Consumes value, returning the mio I/O object. /// - /// See [`tokio_net::PollEvented::into_inner`] for more details about + /// See [`tokio_net::util::PollEvented::into_inner`] for more details about /// resource deregistration that happens during the call. fn try_from(value: TcpListener) -> Result { value.io.into_inner() diff --git a/tokio-tcp/src/stream.rs b/tokio-tcp/src/stream.rs index c6397268e..0a527d431 100644 --- a/tokio-tcp/src/stream.rs +++ b/tokio-tcp/src/stream.rs @@ -4,7 +4,8 @@ use crate::split::{ }; use tokio_io::{AsyncRead, AsyncWrite}; -use tokio_net::{Handle, PollEvented}; +use tokio_net::driver::Handle; +use tokio_net::util::PollEvented; use bytes::{Buf, BufMut}; use futures_core::ready; @@ -125,7 +126,7 @@ impl TcpStream { /// /// ```no_run /// use tokio::net::TcpStream; - /// use tokio_net::Handle; + /// use tokio_net::driver::Handle; /// /// # fn dox() -> std::io::Result<()> { /// let std_stream = std::net::TcpStream::connect("127.0.0.1:34254")?; @@ -801,7 +802,7 @@ impl TryFrom for mio::net::TcpStream { /// Consumes value, returning the mio I/O object. /// - /// See [`tokio_net::PollEvented::into_inner`] for more details about + /// See [`tokio_net::util::PollEvented::into_inner`] for more details about /// resource deregistration that happens during the call. fn try_from(value: TcpStream) -> Result { value.io.into_inner() diff --git a/tokio-udp/src/socket.rs b/tokio-udp/src/socket.rs index 53ace17d4..52ef0645f 100644 --- a/tokio-udp/src/socket.rs +++ b/tokio-udp/src/socket.rs @@ -1,6 +1,7 @@ use super::split::{split, UdpSocketRecvHalf, UdpSocketSendHalf}; -use tokio_net::{Handle, PollEvented}; +use tokio_net::driver::Handle; +use tokio_net::util::PollEvented; use futures_core::ready; use futures_util::future::poll_fn; @@ -328,7 +329,7 @@ impl TryFrom for mio::net::UdpSocket { /// Consumes value, returning the mio I/O object. /// - /// See [`tokio_net::PollEvented::into_inner`] for more details about + /// See [`tokio_net::util::PollEvented::into_inner`] for more details about /// resource deregistration that happens during the call. fn try_from(value: UdpSocket) -> Result { value.io.into_inner() diff --git a/tokio-uds/src/datagram.rs b/tokio-uds/src/datagram.rs index c978b78ee..644ad7534 100644 --- a/tokio-uds/src/datagram.rs +++ b/tokio-uds/src/datagram.rs @@ -1,4 +1,5 @@ -use tokio_net::{Handle, PollEvented}; +use tokio_net::driver::Handle; +use tokio_net::util::PollEvented; use futures_core::ready; use futures_util::future::poll_fn; @@ -200,7 +201,7 @@ impl TryFrom for mio_uds::UnixDatagram { /// Consumes value, returning the mio I/O object. /// - /// See [`tokio_net::PollEvented::into_inner`] for more details about + /// See [`tokio_net::util::PollEvented::into_inner`] for more details about /// resource deregistration that happens during the call. fn try_from(value: UnixDatagram) -> Result { value.io.into_inner() diff --git a/tokio-uds/src/listener.rs b/tokio-uds/src/listener.rs index b3e02c508..310914359 100644 --- a/tokio-uds/src/listener.rs +++ b/tokio-uds/src/listener.rs @@ -1,6 +1,7 @@ use crate::UnixStream; -use tokio_net::{Handle, PollEvented}; +use tokio_net::driver::Handle; +use tokio_net::util::PollEvented; use futures_core::ready; use futures_util::future::poll_fn; @@ -102,7 +103,7 @@ impl TryFrom for mio_uds::UnixListener { /// Consumes value, returning the mio I/O object. /// - /// See [`tokio_net::PollEvented::into_inner`] for more details about + /// See [`tokio_net::util::PollEvented::into_inner`] for more details about /// resource deregistration that happens during the call. fn try_from(value: UnixListener) -> Result { value.io.into_inner() diff --git a/tokio-uds/src/stream.rs b/tokio-uds/src/stream.rs index 29a76ce3b..4b318a054 100644 --- a/tokio-uds/src/stream.rs +++ b/tokio-uds/src/stream.rs @@ -5,7 +5,8 @@ use crate::split::{ use crate::ucred::{self, UCred}; use tokio_io::{AsyncRead, AsyncWrite}; -use tokio_net::{Handle, PollEvented}; +use tokio_net::driver::Handle; +use tokio_net::util::PollEvented; use bytes::{Buf, BufMut}; use futures_core::ready; @@ -131,7 +132,7 @@ impl TryFrom for mio_uds::UnixStream { /// Consumes value, returning the mio I/O object. /// - /// See [`tokio_net::PollEvented::into_inner`] for more details about + /// See [`tokio_net::util::PollEvented::into_inner`] for more details about /// resource deregistration that happens during the call. fn try_from(value: UnixStream) -> Result { value.io.into_inner() diff --git a/tokio/src/lib.rs b/tokio/src/lib.rs index 81747bae8..ec90d5ec1 100644 --- a/tokio/src/lib.rs +++ b/tokio/src/lib.rs @@ -88,8 +88,6 @@ pub mod io; #[cfg(any(feature = "tcp", feature = "udp", feature = "uds"))] pub mod net; pub mod prelude; -#[cfg(feature = "tokio-net")] -pub mod reactor; pub mod stream; #[cfg(feature = "sync")] pub mod sync; diff --git a/tokio/src/runtime/current_thread/builder.rs b/tokio/src/runtime/current_thread/builder.rs index 7c9cd7979..d48136d4c 100644 --- a/tokio/src/runtime/current_thread/builder.rs +++ b/tokio/src/runtime/current_thread/builder.rs @@ -1,7 +1,7 @@ use crate::runtime::current_thread::Runtime; use tokio_executor::current_thread::CurrentThread; -use tokio_net::Reactor; +use tokio_net::driver::Reactor; use tokio_timer::clock::Clock; use tokio_timer::timer::Timer; diff --git a/tokio/src/runtime/current_thread/runtime.rs b/tokio/src/runtime/current_thread/runtime.rs index 8482ecadc..da3e6770a 100644 --- a/tokio/src/runtime/current_thread/runtime.rs +++ b/tokio/src/runtime/current_thread/runtime.rs @@ -2,7 +2,7 @@ use crate::runtime::current_thread::Builder; use tokio_executor::current_thread::Handle as ExecutorHandle; use tokio_executor::current_thread::{self, CurrentThread}; -use tokio_net::{self, Reactor}; +use tokio_net::driver::{self, Reactor}; use tokio_timer::clock::{self, Clock}; use tokio_timer::timer::{self, Timer}; @@ -19,7 +19,7 @@ use std::io; /// [mod]: index.html #[derive(Debug)] pub struct Runtime { - reactor_handle: tokio_net::Handle, + reactor_handle: driver::Handle, timer_handle: timer::Handle, clock: Clock, executor: CurrentThread, @@ -93,7 +93,7 @@ impl Runtime { } pub(super) fn new2( - reactor_handle: tokio_net::Handle, + reactor_handle: driver::Handle, timer_handle: timer::Handle, clock: Clock, executor: CurrentThread, @@ -197,7 +197,7 @@ impl Runtime { // This will set the default handle and timer to use inside the closure // and run the future. - let _reactor = tokio_net::set_default(&reactor_handle); + let _reactor = driver::set_default(&reactor_handle); clock::with_default(clock, || { let _timer = timer::set_default(&timer_handle); // The TaskExecutor is a fake executor that looks into the diff --git a/tokio/src/runtime/threadpool/background.rs b/tokio/src/runtime/threadpool/background.rs index a65e02504..3d884118d 100644 --- a/tokio/src/runtime/threadpool/background.rs +++ b/tokio/src/runtime/threadpool/background.rs @@ -2,7 +2,7 @@ //! `block_on` work. use tokio_executor::current_thread::CurrentThread; -use tokio_net::Reactor; +use tokio_net::driver::{self, Reactor}; use tokio_sync::oneshot; use tokio_timer::clock::Clock; use tokio_timer::timer::{self, Timer}; @@ -11,7 +11,7 @@ use std::{io, thread}; #[derive(Debug)] pub(crate) struct Background { - reactor_handle: tokio_net::Handle, + reactor_handle: driver::Handle, timer_handle: timer::Handle, shutdown_tx: Option>, thread: Option>, @@ -44,7 +44,7 @@ pub(crate) fn spawn(clock: &Clock) -> io::Result { } impl Background { - pub(super) fn reactor(&self) -> &tokio_net::Handle { + pub(super) fn reactor(&self) -> &driver::Handle { &self.reactor_handle } diff --git a/tokio/src/runtime/threadpool/builder.rs b/tokio/src/runtime/threadpool/builder.rs index 10beff32e..d045c398b 100644 --- a/tokio/src/runtime/threadpool/builder.rs +++ b/tokio/src/runtime/threadpool/builder.rs @@ -1,7 +1,7 @@ use super::{background, Inner, Runtime}; -use crate::reactor::Reactor; use tokio_executor::threadpool; +use tokio_net::driver::{self, Reactor}; use tokio_timer::clock::{self, Clock}; use tokio_timer::timer::{self, Timer}; @@ -343,7 +343,7 @@ impl Builder { .around_worker(move |w| { let index = w.id().to_usize(); - let _reactor = tokio_net::set_default(&reactor_handles[index]); + let _reactor = driver::set_default(&reactor_handles[index]); clock::with_default(&clock, || { let _timer = timer::set_default(&timer_handles[index]); trace::dispatcher::with_default(&dispatch, || { diff --git a/tokio/src/runtime/threadpool/mod.rs b/tokio/src/runtime/threadpool/mod.rs index 19bcdf07d..b9de101a5 100644 --- a/tokio/src/runtime/threadpool/mod.rs +++ b/tokio/src/runtime/threadpool/mod.rs @@ -10,6 +10,7 @@ use background::Background; use tokio_executor::enter; use tokio_executor::threadpool::ThreadPool; +use tokio_net::driver; use tokio_timer::timer; use tracing_core as trace; @@ -174,7 +175,7 @@ impl Runtime { let trace = &self.inner().trace; tokio_executor::with_default(&mut self.inner().pool.sender(), || { - let _reactor = tokio_net::set_default(bg.reactor()); + let _reactor = driver::set_default(bg.reactor()); let _timer = timer::set_default(bg.timer()); trace::dispatcher::with_default(trace, || { entered.block_on(future) diff --git a/tokio/tests/drop-core.rs b/tokio/tests/drop-core.rs index 32fb24447..bb41c12e7 100644 --- a/tokio/tests/drop-core.rs +++ b/tokio/tests/drop-core.rs @@ -3,7 +3,7 @@ #![cfg(feature = "default")] use tokio::net::TcpListener; -use tokio::reactor::Reactor; +use tokio_net::driver::Reactor; use tokio_test::{assert_err, assert_pending, assert_ready, task}; #[test] diff --git a/tokio/tests/reactor.rs b/tokio/tests/reactor.rs index 04107ccb6..8e93bd485 100644 --- a/tokio/tests/reactor.rs +++ b/tokio/tests/reactor.rs @@ -2,7 +2,7 @@ #![warn(rust_2018_idioms)] #![cfg(feature = "default")] -use tokio_net::Reactor; +use tokio_net::driver::Reactor; use tokio_tcp::TcpListener; use tokio_test::{assert_ok, assert_pending}; @@ -68,7 +68,7 @@ fn test_drop_on_notify() { { let handle = reactor.handle(); - let _reactor = tokio_net::set_default(&handle); + let _reactor = tokio_net::driver::set_default(&handle); let waker = waker_ref(&task); let mut cx = Context::from_waker(&waker); assert_pending!(task.future.lock().unwrap().as_mut().poll(&mut cx)); diff --git a/ui-tests/tests/features.rs b/ui-tests/tests/features.rs index 71ef87eb6..be69854c8 100644 --- a/ui-tests/tests/features.rs +++ b/ui-tests/tests/features.rs @@ -2,9 +2,6 @@ #[cfg(feature = "tokio-with-net")] #[allow(unused_imports)] fn tokio_with_net() { - // Reactor is present - use ui_tests::tokio::reactor; - // net is present use ui_tests::tokio::net; } @@ -16,7 +13,7 @@ fn compile_fail() { t.compile_fail("tests/ui/executor_without_current_thread.rs"); #[cfg(feature = "tokio-no-features")] - t.compile_fail("tests/ui/tokio_without_net_missing_reactor.rs"); + t.compile_fail("tests/ui/tokio_without_net_missing_net.rs"); drop(t); } diff --git a/ui-tests/tests/ui/tokio_without_net_missing_net.rs b/ui-tests/tests/ui/tokio_without_net_missing_net.rs new file mode 100644 index 000000000..c180712b1 --- /dev/null +++ b/ui-tests/tests/ui/tokio_without_net_missing_net.rs @@ -0,0 +1,3 @@ +use ui_tests::tokio::net; + +fn main() {} diff --git a/ui-tests/tests/ui/tokio_without_net_missing_net.stderr b/ui-tests/tests/ui/tokio_without_net_missing_net.stderr new file mode 100644 index 000000000..741b11909 --- /dev/null +++ b/ui-tests/tests/ui/tokio_without_net_missing_net.stderr @@ -0,0 +1,7 @@ +error[E0432]: unresolved import `ui_tests::tokio::net` + --> $DIR/tokio_without_net_missing_net.rs:1:5 + | +1 | use ui_tests::tokio::net; + | ^^^^^^^^^^^^^^^^^^^^ no `net` in `tokio` + +For more information about this error, try `rustc --explain E0432`. diff --git a/ui-tests/tests/ui/tokio_without_net_missing_reactor.rs b/ui-tests/tests/ui/tokio_without_net_missing_reactor.rs deleted file mode 100644 index 2ca30ed12..000000000 --- a/ui-tests/tests/ui/tokio_without_net_missing_reactor.rs +++ /dev/null @@ -1,3 +0,0 @@ -use ui_tests::tokio::reactor; - -fn main() {} diff --git a/ui-tests/tests/ui/tokio_without_net_missing_reactor.stderr b/ui-tests/tests/ui/tokio_without_net_missing_reactor.stderr deleted file mode 100644 index 3b1610b9e..000000000 --- a/ui-tests/tests/ui/tokio_without_net_missing_reactor.stderr +++ /dev/null @@ -1,7 +0,0 @@ -error[E0432]: unresolved import `ui_tests::tokio::reactor` - --> $DIR/tokio_without_net_missing_reactor.rs:1:5 - | -1 | use ui_tests::tokio::reactor; - | ^^^^^^^^^^^^^^^^^^^^^^^^ no `reactor` in `tokio` - -For more information about this error, try `rustc --explain E0432`.