mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-28 00:00:11 +02:00
Space is made to add `tcp`, `udp`, `uds`, ... modules.
This commit is contained in:
@@ -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<dyn std::error::Error>> {
|
||||
|
||||
@@ -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;
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
@@ -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<Inner>,
|
||||
|
||||
_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<HandlePriv>,
|
||||
}
|
||||
|
||||
/// Like `Handle`, but never `None`.
|
||||
#[derive(Clone)]
|
||||
pub(crate) struct HandlePriv {
|
||||
inner: Weak<Inner>,
|
||||
}
|
||||
|
||||
/// 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::<Handle>(), mem::size_of::<HandlePriv>());
|
||||
}
|
||||
|
||||
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<Slab<ScheduledIo>>,
|
||||
|
||||
/// 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<Option<HandlePriv>> = 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<T: Send + Sync>() {}
|
||||
|
||||
_assert::<Handle>();
|
||||
}
|
||||
|
||||
// ===== 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<Reactor> {
|
||||
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<Duration>) -> io::Result<Turn> {
|
||||
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<Duration>) -> 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<HandlePriv> {
|
||||
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<Arc<Inner>> {
|
||||
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<usize> {
|
||||
// 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(),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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];
|
||||
+2
-567
@@ -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<Inner>,
|
||||
|
||||
_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<HandlePriv>,
|
||||
}
|
||||
|
||||
/// Like `Handle`, but never `None`.
|
||||
#[derive(Clone)]
|
||||
struct HandlePriv {
|
||||
inner: Weak<Inner>,
|
||||
}
|
||||
|
||||
/// 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::<Handle>(), mem::size_of::<HandlePriv>());
|
||||
}
|
||||
|
||||
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<Slab<ScheduledIo>>,
|
||||
|
||||
/// 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<Option<HandlePriv>> = 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<T: Send + Sync>() {}
|
||||
|
||||
_assert::<Handle>();
|
||||
}
|
||||
|
||||
// ===== 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<Reactor> {
|
||||
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<Duration>) -> io::Result<Turn> {
|
||||
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<Duration>) -> 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<HandlePriv> {
|
||||
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<Arc<Inner>> {
|
||||
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<usize> {
|
||||
// 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;
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
//! Utilities for implementing networking types.
|
||||
|
||||
mod poll_evented;
|
||||
|
||||
pub use self::poll_evented::PollEvented;
|
||||
@@ -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
|
||||
@@ -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;
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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::*;
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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};
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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<TcpListener> 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<Self, Self::Error> {
|
||||
value.io.into_inner()
|
||||
|
||||
@@ -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<TcpStream> 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<Self, Self::Error> {
|
||||
value.io.into_inner()
|
||||
|
||||
@@ -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<UdpSocket> 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<Self, Self::Error> {
|
||||
value.io.into_inner()
|
||||
|
||||
@@ -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<UnixDatagram> 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<Self, Self::Error> {
|
||||
value.io.into_inner()
|
||||
|
||||
@@ -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<UnixListener> 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<Self, Self::Error> {
|
||||
value.io.into_inner()
|
||||
|
||||
@@ -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<UnixStream> 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<Self, Self::Error> {
|
||||
value.io.into_inner()
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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;
|
||||
|
||||
|
||||
@@ -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<Parker>,
|
||||
@@ -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<Parker>,
|
||||
@@ -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
|
||||
|
||||
@@ -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<oneshot::Sender<()>>,
|
||||
thread: Option<thread::JoinHandle<()>>,
|
||||
@@ -44,7 +44,7 @@ pub(crate) fn spawn(clock: &Clock) -> io::Result<Background> {
|
||||
}
|
||||
|
||||
impl Background {
|
||||
pub(super) fn reactor(&self) -> &tokio_net::Handle {
|
||||
pub(super) fn reactor(&self) -> &driver::Handle {
|
||||
&self.reactor_handle
|
||||
}
|
||||
|
||||
|
||||
@@ -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, || {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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));
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,3 @@
|
||||
use ui_tests::tokio::net;
|
||||
|
||||
fn main() {}
|
||||
@@ -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`.
|
||||
@@ -1,3 +0,0 @@
|
||||
use ui_tests::tokio::reactor;
|
||||
|
||||
fn main() {}
|
||||
@@ -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`.
|
||||
Reference in New Issue
Block a user