diff --git a/tokio-util/Cargo.toml b/tokio-util/Cargo.toml index c2b83565b..e029de2f6 100644 --- a/tokio-util/Cargo.toml +++ b/tokio-util/Cargo.toml @@ -24,8 +24,9 @@ categories = ["asynchronous"] default = [] # Shorthand for enabling everything -full = ["codec", "udp"] +full = ["codec", "udp", "compat"] +compat = ["futures-io",] codec = ["tokio/stream"] udp = ["tokio/udp"] @@ -35,6 +36,7 @@ tokio = { version = "0.2.0", path = "../tokio" } bytes = "0.5.0" futures-core = "0.3.0" futures-sink = "0.3.0" +futures-io = { version = "0.3.0", optional = true } log = "0.4" pin-project-lite = "0.1.1" diff --git a/tokio-util/src/cfg.rs b/tokio-util/src/cfg.rs index 13fabd363..27e8c66a4 100644 --- a/tokio-util/src/cfg.rs +++ b/tokio-util/src/cfg.rs @@ -8,6 +8,16 @@ macro_rules! cfg_codec { } } +macro_rules! cfg_compat { + ($($item:item)*) => { + $( + #[cfg(feature = "compat")] + #[cfg_attr(docsrs, doc(cfg(feature = "compat")))] + $item + )* + } +} + macro_rules! cfg_udp { ($($item:item)*) => { $( diff --git a/tokio-util/src/compat.rs b/tokio-util/src/compat.rs new file mode 100644 index 000000000..769e30c2b --- /dev/null +++ b/tokio-util/src/compat.rs @@ -0,0 +1,201 @@ +//! Compatibility between the `tokio::io` and `futures-io` versions of the +//! `AsyncRead` and `AsyncWrite` traits. +use pin_project_lite::pin_project; +use std::io; +use std::pin::Pin; +use std::task::{Context, Poll}; + +pin_project! { + /// A compatibility layer that allows conversion between the + /// `tokio::io` and `futures-io` `AsyncRead` and `AsyncWrite` traits. + #[derive(Copy, Clone, Debug)] + pub struct Compat { + #[pin] + inner: T, + } +} + +/// Extension trait that allows converting a type implementing +/// `futures_io::AsyncRead` to implement `tokio::io::AsyncRead`. +pub trait FuturesAsyncReadCompatExt: futures_io::AsyncRead { + /// Wraps `self` with a compatibility layer that implements + /// `tokio_io::AsyncWrite`. + fn compat(self) -> Compat + where + Self: Sized, + { + Compat::new(self) + } +} + +impl FuturesAsyncReadCompatExt for T {} + +/// Extension trait that allows converting a type implementing +/// `futures_io::AsyncWrite` to implement `tokio::io::AsyncWrite`. +pub trait FuturesAsyncWriteCompatExt: futures_io::AsyncWrite { + /// Wraps `self` with a compatibility layer that implements + /// `tokio::io::AsyncWrite`. + fn compat_write(self) -> Compat + where + Self: Sized, + { + Compat::new(self) + } +} + +impl FuturesAsyncWriteCompatExt for T {} + +/// Extension trait that allows converting a type implementing +/// `tokio::io::AsyncRead` to implement `futures_io::AsyncRead`. +pub trait Tokio02AsyncReadCompatExt: tokio::io::AsyncRead { + /// Wraps `self` with a compatibility layer that implements + /// `futures_io::AsyncRead`. + fn compat(self) -> Compat + where + Self: Sized, + { + Compat::new(self) + } +} + +impl Tokio02AsyncReadCompatExt for T {} + +/// Extension trait that allows converting a type implementing +/// `tokio::io::AsyncWrite` to implement `futures_io::AsyncWrite`. +pub trait Tokio02AsyncWriteCompatExt: tokio::io::AsyncWrite { + /// Wraps `self` with a compatibility layer that implements + /// `futures_io::AsyncWrite`. + fn compat_write(self) -> Compat + where + Self: Sized, + { + Compat::new(self) + } +} + +impl Tokio02AsyncWriteCompatExt for T {} + +// === impl Compat === + +impl Compat { + fn new(inner: T) -> Self { + Self { inner } + } + + /// Get a reference to the `Future`, `Stream`, `AsyncRead`, or `AsyncWrite` object + /// contained within. + pub fn get_ref(&self) -> &T { + &self.inner + } + + /// Get a mutable reference to the `Future`, `Stream`, `AsyncRead`, or `AsyncWrite` object + /// contained within. + pub fn get_mut(&mut self) -> &mut T { + &mut self.inner + } + + /// Returns the wrapped item. + pub fn into_inner(self) -> T { + self.inner + } +} + +impl tokio::io::AsyncRead for Compat +where + T: futures_io::AsyncRead, +{ + fn poll_read( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + buf: &mut [u8], + ) -> Poll> { + futures_io::AsyncRead::poll_read(self.project().inner, cx, buf) + } +} + +impl futures_io::AsyncRead for Compat +where + T: tokio::io::AsyncRead, +{ + fn poll_read( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + buf: &mut [u8], + ) -> Poll> { + tokio::io::AsyncRead::poll_read(self.project().inner, cx, buf) + } +} + +impl tokio::io::AsyncBufRead for Compat +where + T: futures_io::AsyncBufRead, +{ + fn poll_fill_buf<'a>( + self: Pin<&'a mut Self>, + cx: &mut Context<'_>, + ) -> Poll> { + futures_io::AsyncBufRead::poll_fill_buf(self.project().inner, cx) + } + + fn consume(self: Pin<&mut Self>, amt: usize) { + futures_io::AsyncBufRead::consume(self.project().inner, amt) + } +} + +impl futures_io::AsyncBufRead for Compat +where + T: tokio::io::AsyncBufRead, +{ + fn poll_fill_buf<'a>( + self: Pin<&'a mut Self>, + cx: &mut Context<'_>, + ) -> Poll> { + tokio::io::AsyncBufRead::poll_fill_buf(self.project().inner, cx) + } + + fn consume(self: Pin<&mut Self>, amt: usize) { + tokio::io::AsyncBufRead::consume(self.project().inner, amt) + } +} + +impl tokio::io::AsyncWrite for Compat +where + T: futures_io::AsyncWrite, +{ + fn poll_write( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + buf: &[u8], + ) -> Poll> { + futures_io::AsyncWrite::poll_write(self.project().inner, cx, buf) + } + + fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + futures_io::AsyncWrite::poll_flush(self.project().inner, cx) + } + + fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + futures_io::AsyncWrite::poll_close(self.project().inner, cx) + } +} + +impl futures_io::AsyncWrite for Compat +where + T: tokio::io::AsyncWrite, +{ + fn poll_write( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + buf: &[u8], + ) -> Poll> { + tokio::io::AsyncWrite::poll_write(self.project().inner, cx, buf) + } + + fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + tokio::io::AsyncWrite::poll_flush(self.project().inner, cx) + } + + fn poll_close(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + tokio::io::AsyncWrite::poll_shutdown(self.project().inner, cx) + } +} diff --git a/tokio-util/src/lib.rs b/tokio-util/src/lib.rs index 48c0fd166..4554516e0 100644 --- a/tokio-util/src/lib.rs +++ b/tokio-util/src/lib.rs @@ -25,3 +25,7 @@ cfg_codec! { cfg_udp! { pub mod udp; } + +cfg_compat! { + pub mod compat; +}