mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-29 00:00:11 +02:00
tls: update to std-future (#1224)
This commit is contained in:
committed by
Carl Lerche
parent
0de3a69eb4
commit
448d9d2eab
+1
-1
@@ -17,7 +17,7 @@ members = [
|
|||||||
"tokio-threadpool",
|
"tokio-threadpool",
|
||||||
"tokio-timer",
|
"tokio-timer",
|
||||||
"tokio-tcp",
|
"tokio-tcp",
|
||||||
# "tokio-tls",
|
"tokio-tls",
|
||||||
"tokio-udp",
|
"tokio-udp",
|
||||||
"tokio-uds",
|
"tokio-uds",
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -26,14 +26,15 @@ publish = false
|
|||||||
travis-ci = { repository = "tokio-rs/tokio-tls" }
|
travis-ci = { repository = "tokio-rs/tokio-tls" }
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
futures = "0.1.23"
|
|
||||||
native-tls = "0.2"
|
native-tls = "0.2"
|
||||||
tokio-io = { version = "0.2.0", path = "../tokio-io" }
|
tokio-io = { version = "0.2.0", path = "../tokio-io" }
|
||||||
|
|
||||||
[dev-dependencies]
|
[dev-dependencies]
|
||||||
tokio = { version = "0.2.0", path = "../tokio" }
|
tokio = { version = "0.2.0", path = "../tokio" }
|
||||||
|
tokio-tcp = { version = "0.2.0", path = "../tokio-tcp", features = ["async-traits"] }
|
||||||
cfg-if = "0.1"
|
cfg-if = "0.1"
|
||||||
env_logger = { version = "0.5", default-features = false }
|
env_logger = { version = "0.5", default-features = false }
|
||||||
|
futures-preview = { version = "0.3.0-alpha.17", features = ["async-await", "nightly"] }
|
||||||
|
|
||||||
[target.'cfg(all(not(target_os = "macos"), not(windows), not(target_os = "ios")))'.dev-dependencies]
|
[target.'cfg(all(not(target_os = "macos"), not(windows), not(target_os = "ios")))'.dev-dependencies]
|
||||||
openssl = "0.10"
|
openssl = "0.10"
|
||||||
|
|||||||
@@ -1,32 +1,28 @@
|
|||||||
#![deny(warnings, rust_2018_idioms)]
|
// #![deny(warnings, rust_2018_idioms)]
|
||||||
|
#![feature(async_await)]
|
||||||
|
|
||||||
use futures::Future;
|
|
||||||
use native_tls::TlsConnector;
|
use native_tls::TlsConnector;
|
||||||
use std::io;
|
use std::error::Error;
|
||||||
use std::net::ToSocketAddrs;
|
use std::net::ToSocketAddrs;
|
||||||
|
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||||
use tokio::net::TcpStream;
|
use tokio::net::TcpStream;
|
||||||
use tokio::runtime::Runtime;
|
|
||||||
use tokio_io;
|
|
||||||
use tokio_tls;
|
use tokio_tls;
|
||||||
|
|
||||||
fn main() -> Result<(), Box<dyn std::error::Error>> {
|
#[tokio::main]
|
||||||
let runtime = Runtime::new()?;
|
async fn main() -> Result<(), Box<dyn Error + Send + Sync>> {
|
||||||
let addr = "www.rust-lang.org:443"
|
let addr = "www.rust-lang.org:443"
|
||||||
.to_socket_addrs()?
|
.to_socket_addrs()?
|
||||||
.next()
|
.next()
|
||||||
.ok_or("failed to resolve www.rust-lang.org")?;
|
.ok_or("failed to resolve www.rust-lang.org")?;
|
||||||
|
|
||||||
let socket = TcpStream::connect(&addr);
|
let socket = TcpStream::connect(&addr).await?;
|
||||||
let cx = TlsConnector::builder().build()?;
|
let cx = TlsConnector::builder().build()?;
|
||||||
let cx = tokio_tls::TlsConnector::from(cx);
|
let cx = tokio_tls::TlsConnector::from(cx);
|
||||||
|
|
||||||
let tls_handshake = socket.and_then(move |socket| {
|
let mut socket = cx.connect("www.rust-lang.org", socket).await?;
|
||||||
cx.connect("www.rust-lang.org", socket)
|
|
||||||
.map_err(|e| io::Error::new(io::ErrorKind::Other, e))
|
socket
|
||||||
});
|
.write_all(
|
||||||
let request = tls_handshake.and_then(|socket| {
|
|
||||||
tokio_io::io::write_all(
|
|
||||||
socket,
|
|
||||||
"\
|
"\
|
||||||
GET / HTTP/1.0\r\n\
|
GET / HTTP/1.0\r\n\
|
||||||
Host: www.rust-lang.org\r\n\
|
Host: www.rust-lang.org\r\n\
|
||||||
@@ -34,10 +30,12 @@ fn main() -> Result<(), Box<dyn std::error::Error>> {
|
|||||||
"
|
"
|
||||||
.as_bytes(),
|
.as_bytes(),
|
||||||
)
|
)
|
||||||
});
|
.await?;
|
||||||
let response = request.and_then(|(socket, _)| tokio_io::io::read_to_end(socket, Vec::new()));
|
|
||||||
|
|
||||||
let (_, data) = runtime.block_on(response)?;
|
let mut data = Vec::new();
|
||||||
println!("{}", String::from_utf8_lossy(&data));
|
socket.read_to_end(&mut data).await?;
|
||||||
|
|
||||||
|
// println!("data: {:?}", &data);
|
||||||
|
println!("{}", String::from_utf8_lossy(&data[..]));
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
+207
-93
@@ -2,6 +2,7 @@
|
|||||||
#![deny(rust_2018_idioms)]
|
#![deny(rust_2018_idioms)]
|
||||||
#![cfg_attr(test, deny(warnings))]
|
#![cfg_attr(test, deny(warnings))]
|
||||||
#![doc(test(no_crate_inject, attr(deny(rust_2018_idioms))))]
|
#![doc(test(no_crate_inject, attr(deny(rust_2018_idioms))))]
|
||||||
|
#![feature(async_await)]
|
||||||
|
|
||||||
//! Async TLS streams
|
//! Async TLS streams
|
||||||
//!
|
//!
|
||||||
@@ -20,10 +21,19 @@
|
|||||||
//! built. Configuration of TLS parameters is still primarily done through the
|
//! built. Configuration of TLS parameters is still primarily done through the
|
||||||
//! `native-tls` crate.
|
//! `native-tls` crate.
|
||||||
|
|
||||||
use futures::{Async, Future, Poll};
|
use native_tls::{Error, HandshakeError, MidHandshakeTlsStream};
|
||||||
use native_tls::{Error, HandshakeError};
|
use std::future::Future;
|
||||||
use std::io::{self, Read, Write};
|
use std::io::{self, Read, Write};
|
||||||
use tokio_io::{try_nb, AsyncRead, AsyncWrite};
|
use std::marker::Unpin;
|
||||||
|
use std::pin::Pin;
|
||||||
|
use std::task::{Context, Poll};
|
||||||
|
use tokio_io::{AsyncRead, AsyncWrite};
|
||||||
|
|
||||||
|
#[derive(Debug)]
|
||||||
|
struct AllowStd<S> {
|
||||||
|
inner: S,
|
||||||
|
context: *mut (),
|
||||||
|
}
|
||||||
|
|
||||||
/// A wrapper around an underlying raw stream which implements the TLS or SSL
|
/// A wrapper around an underlying raw stream which implements the TLS or SSL
|
||||||
/// protocol.
|
/// protocol.
|
||||||
@@ -33,76 +43,206 @@ use tokio_io::{try_nb, AsyncRead, AsyncWrite};
|
|||||||
/// data. Bytes read from a `TlsStream` are decrypted from `S` and bytes written
|
/// data. Bytes read from a `TlsStream` are decrypted from `S` and bytes written
|
||||||
/// to a `TlsStream` are encrypted when passing through to `S`.
|
/// to a `TlsStream` are encrypted when passing through to `S`.
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub struct TlsStream<S> {
|
pub struct TlsStream<S>(native_tls::TlsStream<AllowStd<S>>);
|
||||||
inner: native_tls::TlsStream<S>,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A wrapper around a `native_tls::TlsConnector`, providing an async `connect`
|
/// A wrapper around a `native_tls::TlsConnector`, providing an async `connect`
|
||||||
/// method.
|
/// method.
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
pub struct TlsConnector {
|
pub struct TlsConnector(native_tls::TlsConnector);
|
||||||
inner: native_tls::TlsConnector,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A wrapper around a `native_tls::TlsAcceptor`, providing an async `accept`
|
/// A wrapper around a `native_tls::TlsAcceptor`, providing an async `accept`
|
||||||
/// method.
|
/// method.
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
pub struct TlsAcceptor {
|
pub struct TlsAcceptor(native_tls::TlsAcceptor);
|
||||||
inner: native_tls::TlsAcceptor,
|
|
||||||
|
struct MidHandshake<S>(Option<MidHandshakeTlsStream<AllowStd<S>>>);
|
||||||
|
|
||||||
|
enum StartedHandshake<S> {
|
||||||
|
Done(TlsStream<S>),
|
||||||
|
Mid(MidHandshakeTlsStream<AllowStd<S>>),
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Future returned from `TlsConnector::connect` which will resolve
|
struct StartedHandshakeFuture<F, S>(Option<StartedHandshakeFutureInner<F, S>>);
|
||||||
/// once the connection handshake has finished.
|
struct StartedHandshakeFutureInner<F, S> {
|
||||||
pub struct Connect<S> {
|
f: F,
|
||||||
inner: MidHandshake<S>,
|
stream: S,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Future returned from `TlsAcceptor::accept` which will resolve
|
struct Guard<'a, S>(&'a mut TlsStream<S>)
|
||||||
/// once the accept handshake has finished.
|
where
|
||||||
pub struct Accept<S> {
|
AllowStd<S>: Read + Write;
|
||||||
inner: MidHandshake<S>,
|
|
||||||
}
|
|
||||||
|
|
||||||
struct MidHandshake<S> {
|
impl<'a, S> Drop for Guard<'a, S>
|
||||||
inner: Option<Result<native_tls::TlsStream<S>, HandshakeError<S>>>,
|
where
|
||||||
}
|
AllowStd<S>: Read + Write,
|
||||||
|
{
|
||||||
impl<S> TlsStream<S> {
|
fn drop(&mut self) {
|
||||||
/// Get access to the internal `native_tls::TlsStream` stream which also
|
(self.0).0.get_mut().context = 0 as *mut ();
|
||||||
/// transitively allows access to `S`.
|
|
||||||
pub fn get_ref(&self) -> &native_tls::TlsStream<S> {
|
|
||||||
&self.inner
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Get mutable access to the internal `native_tls::TlsStream` stream which
|
|
||||||
/// also transitively allows mutable access to `S`.
|
|
||||||
pub fn get_mut(&mut self) -> &mut native_tls::TlsStream<S> {
|
|
||||||
&mut self.inner
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<S: Read + Write> Read for TlsStream<S> {
|
impl<S> AllowStd<S>
|
||||||
|
where
|
||||||
|
S: Unpin,
|
||||||
|
{
|
||||||
|
fn with_context<F, R>(&mut self, f: F) -> R
|
||||||
|
where
|
||||||
|
F: FnOnce(&mut Context<'_>, Pin<&mut S>) -> R,
|
||||||
|
{
|
||||||
|
unsafe {
|
||||||
|
assert!(!self.context.is_null());
|
||||||
|
let waker = &mut *(self.context as *mut _);
|
||||||
|
f(waker, Pin::new(&mut self.inner))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<S> Read for AllowStd<S>
|
||||||
|
where
|
||||||
|
S: AsyncRead + Unpin,
|
||||||
|
{
|
||||||
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
|
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
|
||||||
self.inner.read(buf)
|
match self.with_context(|ctx, stream| stream.poll_read(ctx, buf)) {
|
||||||
|
Poll::Ready(r) => r,
|
||||||
|
Poll::Pending => Err(io::Error::from(io::ErrorKind::WouldBlock)),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<S: Read + Write> Write for TlsStream<S> {
|
impl<S> Write for AllowStd<S>
|
||||||
|
where
|
||||||
|
S: AsyncWrite + Unpin,
|
||||||
|
{
|
||||||
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
|
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
|
||||||
self.inner.write(buf)
|
match self.with_context(|ctx, stream| stream.poll_write(ctx, buf)) {
|
||||||
|
Poll::Ready(r) => r,
|
||||||
|
Poll::Pending => Err(io::Error::from(io::ErrorKind::WouldBlock)),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn flush(&mut self) -> io::Result<()> {
|
fn flush(&mut self) -> io::Result<()> {
|
||||||
self.inner.flush()
|
match self.with_context(|ctx, stream| stream.poll_flush(ctx)) {
|
||||||
|
Poll::Ready(r) => r,
|
||||||
|
Poll::Pending => Err(io::Error::from(io::ErrorKind::WouldBlock)),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<S: AsyncRead + AsyncWrite> AsyncRead for TlsStream<S> {}
|
fn cvt<T>(r: io::Result<T>) -> Poll<io::Result<T>> {
|
||||||
|
match r {
|
||||||
|
Ok(v) => Poll::Ready(Ok(v)),
|
||||||
|
Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => Poll::Pending,
|
||||||
|
Err(e) => Poll::Ready(Err(e)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl<S: AsyncRead + AsyncWrite> AsyncWrite for TlsStream<S> {
|
impl<S> TlsStream<S> {
|
||||||
fn shutdown(&mut self) -> Poll<(), io::Error> {
|
fn with_context<F, R>(&mut self, ctx: &mut Context<'_>, f: F) -> R
|
||||||
try_nb!(self.inner.shutdown());
|
where
|
||||||
self.inner.get_mut().shutdown()
|
F: FnOnce(&mut native_tls::TlsStream<AllowStd<S>>) -> R,
|
||||||
|
AllowStd<S>: Read + Write,
|
||||||
|
{
|
||||||
|
self.0.get_mut().context = ctx as *mut _ as *mut ();
|
||||||
|
let g = Guard(self);
|
||||||
|
let r = f(&mut (g.0).0);
|
||||||
|
r
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<S> AsyncRead for TlsStream<S>
|
||||||
|
where
|
||||||
|
S: AsyncRead + AsyncWrite + Unpin,
|
||||||
|
{
|
||||||
|
unsafe fn prepare_uninitialized_buffer(&self, _: &mut [u8]) -> bool {
|
||||||
|
// Note that this does not forward to `S` because the buffer is
|
||||||
|
// unconditionally filled in by OpenSSL, not the actual object `S`.
|
||||||
|
// We're decrypting bytes from `S` into the buffer above!
|
||||||
|
false
|
||||||
|
}
|
||||||
|
|
||||||
|
fn poll_read(
|
||||||
|
mut self: Pin<&mut Self>,
|
||||||
|
ctx: &mut Context<'_>,
|
||||||
|
buf: &mut [u8],
|
||||||
|
) -> Poll<io::Result<usize>> {
|
||||||
|
self.with_context(ctx, |s| cvt(s.read(buf)))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<S> AsyncWrite for TlsStream<S>
|
||||||
|
where
|
||||||
|
S: AsyncRead + AsyncWrite + Unpin,
|
||||||
|
{
|
||||||
|
fn poll_write(
|
||||||
|
mut self: Pin<&mut Self>,
|
||||||
|
ctx: &mut Context<'_>,
|
||||||
|
buf: &[u8],
|
||||||
|
) -> Poll<io::Result<usize>> {
|
||||||
|
self.with_context(ctx, |s| cvt(s.write(buf)))
|
||||||
|
}
|
||||||
|
|
||||||
|
fn poll_flush(mut self: Pin<&mut Self>, ctx: &mut Context<'_>) -> Poll<io::Result<()>> {
|
||||||
|
self.with_context(ctx, |s| cvt(s.flush()))
|
||||||
|
}
|
||||||
|
|
||||||
|
fn poll_shutdown(mut self: Pin<&mut Self>, ctx: &mut Context<'_>) -> Poll<io::Result<()>> {
|
||||||
|
match self.with_context(ctx, |s| s.shutdown()) {
|
||||||
|
Ok(()) => Poll::Ready(Ok(())),
|
||||||
|
Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => return Poll::Pending,
|
||||||
|
Err(e) => return Poll::Ready(Err(e.into())),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn handshake<F, S>(f: F, stream: S) -> Result<TlsStream<S>, Error>
|
||||||
|
where
|
||||||
|
F: FnOnce(
|
||||||
|
AllowStd<S>,
|
||||||
|
) -> Result<native_tls::TlsStream<AllowStd<S>>, HandshakeError<AllowStd<S>>>
|
||||||
|
+ Unpin,
|
||||||
|
S: AsyncRead + AsyncWrite + Unpin,
|
||||||
|
{
|
||||||
|
let start = StartedHandshakeFuture(Some(StartedHandshakeFutureInner { f, stream }));
|
||||||
|
|
||||||
|
match start.await {
|
||||||
|
Err(e) => Err(e),
|
||||||
|
Ok(StartedHandshake::Done(s)) => Ok(s),
|
||||||
|
Ok(StartedHandshake::Mid(s)) => MidHandshake(Some(s)).await,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<F, S> Future for StartedHandshakeFuture<F, S>
|
||||||
|
where
|
||||||
|
F: FnOnce(
|
||||||
|
AllowStd<S>,
|
||||||
|
) -> Result<native_tls::TlsStream<AllowStd<S>>, HandshakeError<AllowStd<S>>>
|
||||||
|
+ Unpin,
|
||||||
|
S: Unpin,
|
||||||
|
AllowStd<S>: Read + Write,
|
||||||
|
{
|
||||||
|
type Output = Result<StartedHandshake<S>, Error>;
|
||||||
|
|
||||||
|
fn poll(
|
||||||
|
mut self: Pin<&mut Self>,
|
||||||
|
ctx: &mut Context<'_>,
|
||||||
|
) -> Poll<Result<StartedHandshake<S>, Error>> {
|
||||||
|
let inner = self.0.take().expect("future polled after completion");
|
||||||
|
let stream = AllowStd {
|
||||||
|
inner: inner.stream,
|
||||||
|
context: ctx as *mut _ as *mut (),
|
||||||
|
};
|
||||||
|
|
||||||
|
match (inner.f)(stream) {
|
||||||
|
Ok(mut s) => {
|
||||||
|
s.get_mut().context = 0 as *mut ();
|
||||||
|
Poll::Ready(Ok(StartedHandshake::Done(TlsStream(s))))
|
||||||
|
}
|
||||||
|
Err(HandshakeError::WouldBlock(mut s)) => {
|
||||||
|
s.get_mut().context = 0 as *mut ();
|
||||||
|
Poll::Ready(Ok(StartedHandshake::Mid(s)))
|
||||||
|
}
|
||||||
|
Err(HandshakeError::Failure(e)) => Poll::Ready(Err(e)),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -119,21 +259,17 @@ impl TlsConnector {
|
|||||||
/// example, a TCP connection to a remote server. That stream is then
|
/// example, a TCP connection to a remote server. That stream is then
|
||||||
/// provided here to perform the client half of a connection to a
|
/// provided here to perform the client half of a connection to a
|
||||||
/// TLS-powered server.
|
/// TLS-powered server.
|
||||||
pub fn connect<S>(&self, domain: &str, stream: S) -> Connect<S>
|
pub async fn connect<S>(&self, domain: &str, stream: S) -> Result<TlsStream<S>, Error>
|
||||||
where
|
where
|
||||||
S: AsyncRead + AsyncWrite,
|
S: AsyncRead + AsyncWrite + Unpin,
|
||||||
{
|
{
|
||||||
Connect {
|
handshake(|s| self.0.connect(domain, s), stream).await
|
||||||
inner: MidHandshake {
|
|
||||||
inner: Some(self.inner.connect(domain, stream)),
|
|
||||||
},
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<native_tls::TlsConnector> for TlsConnector {
|
impl From<native_tls::TlsConnector> for TlsConnector {
|
||||||
fn from(inner: native_tls::TlsConnector) -> TlsConnector {
|
fn from(inner: native_tls::TlsConnector) -> TlsConnector {
|
||||||
TlsConnector { inner }
|
TlsConnector(inner)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -148,58 +284,36 @@ impl TlsAcceptor {
|
|||||||
/// This is typically used after a new socket has been accepted from a
|
/// This is typically used after a new socket has been accepted from a
|
||||||
/// `TcpListener`. That socket is then passed to this function to perform
|
/// `TcpListener`. That socket is then passed to this function to perform
|
||||||
/// the server half of accepting a client connection.
|
/// the server half of accepting a client connection.
|
||||||
pub fn accept<S>(&self, stream: S) -> Accept<S>
|
pub async fn accept<S>(&self, stream: S) -> Result<TlsStream<S>, Error>
|
||||||
where
|
where
|
||||||
S: AsyncRead + AsyncWrite,
|
S: AsyncRead + AsyncWrite + Unpin,
|
||||||
{
|
{
|
||||||
Accept {
|
handshake(|s| self.0.accept(s), stream).await
|
||||||
inner: MidHandshake {
|
|
||||||
inner: Some(self.inner.accept(stream)),
|
|
||||||
},
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl From<native_tls::TlsAcceptor> for TlsAcceptor {
|
impl From<native_tls::TlsAcceptor> for TlsAcceptor {
|
||||||
fn from(inner: native_tls::TlsAcceptor) -> TlsAcceptor {
|
fn from(inner: native_tls::TlsAcceptor) -> TlsAcceptor {
|
||||||
TlsAcceptor { inner }
|
TlsAcceptor(inner)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<S: AsyncRead + AsyncWrite> Future for Connect<S> {
|
impl<S: AsyncRead + AsyncWrite + Unpin> Future for MidHandshake<S> {
|
||||||
type Item = TlsStream<S>;
|
type Output = Result<TlsStream<S>, Error>;
|
||||||
type Error = Error;
|
|
||||||
|
|
||||||
fn poll(&mut self) -> Poll<TlsStream<S>, Error> {
|
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
|
||||||
self.inner.poll()
|
let mut_self = self.get_mut();
|
||||||
}
|
let mut s = mut_self.0.take().expect("future polled after completion");
|
||||||
}
|
|
||||||
|
|
||||||
impl<S: AsyncRead + AsyncWrite> Future for Accept<S> {
|
s.get_mut().context = cx as *mut _ as *mut ();
|
||||||
type Item = TlsStream<S>;
|
match s.handshake() {
|
||||||
type Error = Error;
|
Ok(stream) => Poll::Ready(Ok(TlsStream(stream))),
|
||||||
|
Err(HandshakeError::Failure(e)) => Poll::Ready(Err(e)),
|
||||||
fn poll(&mut self) -> Poll<TlsStream<S>, Error> {
|
Err(HandshakeError::WouldBlock(mut s)) => {
|
||||||
self.inner.poll()
|
s.get_mut().context = 0 as *mut ();
|
||||||
}
|
mut_self.0 = Some(s);
|
||||||
}
|
Poll::Pending
|
||||||
|
}
|
||||||
impl<S: AsyncRead + AsyncWrite> Future for MidHandshake<S> {
|
|
||||||
type Item = TlsStream<S>;
|
|
||||||
type Error = Error;
|
|
||||||
|
|
||||||
fn poll(&mut self) -> Poll<TlsStream<S>, Error> {
|
|
||||||
match self.inner.take().expect("cannot poll MidHandshake twice") {
|
|
||||||
Ok(stream) => Ok(TlsStream { inner: stream }.into()),
|
|
||||||
Err(HandshakeError::Failure(e)) => Err(e),
|
|
||||||
Err(HandshakeError::WouldBlock(s)) => match s.handshake() {
|
|
||||||
Ok(stream) => Ok(TlsStream { inner: stream }.into()),
|
|
||||||
Err(HandshakeError::Failure(e)) => Err(e),
|
|
||||||
Err(HandshakeError::WouldBlock(s)) => {
|
|
||||||
self.inner = Some(Err(HandshakeError::WouldBlock(s)));
|
|
||||||
Ok(Async::NotReady)
|
|
||||||
}
|
|
||||||
},
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+22
-25
@@ -1,13 +1,12 @@
|
|||||||
#![deny(warnings, rust_2018_idioms)]
|
#![deny(warnings, rust_2018_idioms)]
|
||||||
|
#![feature(async_await)]
|
||||||
|
|
||||||
use cfg_if::cfg_if;
|
use cfg_if::cfg_if;
|
||||||
use env_logger;
|
use env_logger;
|
||||||
use futures::Future;
|
|
||||||
use native_tls::TlsConnector;
|
use native_tls::TlsConnector;
|
||||||
use std::io::{self, Error};
|
use std::io::{self, Error};
|
||||||
use std::net::ToSocketAddrs;
|
use std::net::ToSocketAddrs;
|
||||||
use tokio::net::TcpStream;
|
use tokio::net::TcpStream;
|
||||||
use tokio::runtime::Runtime;
|
|
||||||
use tokio_tls;
|
use tokio_tls;
|
||||||
|
|
||||||
macro_rules! t {
|
macro_rules! t {
|
||||||
@@ -83,46 +82,44 @@ cfg_if! {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn get_host(host: &'static str) -> Error {
|
async fn get_host(host: &'static str) -> Error {
|
||||||
drop(env_logger::try_init());
|
drop(env_logger::try_init());
|
||||||
|
|
||||||
let addr = format!("{}:443", host);
|
let addr = format!("{}:443", host);
|
||||||
let addr = t!(addr.to_socket_addrs()).next().unwrap();
|
let addr = t!(addr.to_socket_addrs()).next().unwrap();
|
||||||
|
|
||||||
let l = t!(Runtime::new());
|
let socket = t!(TcpStream::connect(&addr).await);
|
||||||
let client = TcpStream::connect(&addr);
|
let builder = TlsConnector::builder();
|
||||||
let data = client.and_then(move |socket| {
|
let cx = t!(builder.build());
|
||||||
let builder = TlsConnector::builder();
|
let cx = tokio_tls::TlsConnector::from(cx);
|
||||||
let cx = builder.build().unwrap();
|
let res = cx
|
||||||
let cx = tokio_tls::TlsConnector::from(cx);
|
.connect(host, socket)
|
||||||
cx.connect(host, socket)
|
.await
|
||||||
.map_err(|e| Error::new(io::ErrorKind::Other, e))
|
.map_err(|e| Error::new(io::ErrorKind::Other, e));
|
||||||
});
|
|
||||||
|
|
||||||
let res = l.block_on(data);
|
|
||||||
assert!(res.is_err());
|
assert!(res.is_err());
|
||||||
res.err().unwrap()
|
res.err().unwrap()
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[tokio::test]
|
||||||
fn expired() {
|
async fn expired() {
|
||||||
assert_expired_error(&get_host("expired.badssl.com"))
|
assert_expired_error(&get_host("expired.badssl.com").await)
|
||||||
}
|
}
|
||||||
|
|
||||||
// TODO: the OSX builders on Travis apparently fail this tests spuriously?
|
// TODO: the OSX builders on Travis apparently fail this tests spuriously?
|
||||||
// passes locally though? Seems... bad!
|
// passes locally though? Seems... bad!
|
||||||
#[test]
|
#[tokio::test]
|
||||||
#[cfg_attr(all(target_os = "macos", feature = "force-openssl"), ignore)]
|
#[cfg_attr(all(target_os = "macos", feature = "force-openssl"), ignore)]
|
||||||
fn wrong_host() {
|
async fn wrong_host() {
|
||||||
assert_wrong_host(&get_host("wrong.host.badssl.com"))
|
assert_wrong_host(&get_host("wrong.host.badssl.com").await)
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[tokio::test]
|
||||||
fn self_signed() {
|
async fn self_signed() {
|
||||||
assert_self_signed(&get_host("self-signed.badssl.com"))
|
assert_self_signed(&get_host("self-signed.badssl.com").await)
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[tokio::test]
|
||||||
fn untrusted_root() {
|
async fn untrusted_root() {
|
||||||
assert_untrusted_root(&get_host("untrusted-root.badssl.com"))
|
assert_untrusted_root(&get_host("untrusted-root.badssl.com").await)
|
||||||
}
|
}
|
||||||
|
|||||||
+26
-37
@@ -1,15 +1,14 @@
|
|||||||
#![deny(warnings, rust_2018_idioms)]
|
#![deny(warnings, rust_2018_idioms)]
|
||||||
|
#![feature(async_await)]
|
||||||
|
|
||||||
use cfg_if::cfg_if;
|
use cfg_if::cfg_if;
|
||||||
use env_logger;
|
use env_logger;
|
||||||
use futures::Future;
|
|
||||||
use native_tls;
|
use native_tls;
|
||||||
use native_tls::TlsConnector;
|
use native_tls::TlsConnector;
|
||||||
use std::io;
|
use std::io;
|
||||||
use std::net::ToSocketAddrs;
|
use std::net::ToSocketAddrs;
|
||||||
|
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||||
use tokio::net::TcpStream;
|
use tokio::net::TcpStream;
|
||||||
use tokio::runtime::Runtime;
|
|
||||||
use tokio_io::io::{flush, read_to_end, write_all};
|
|
||||||
use tokio_tls;
|
use tokio_tls;
|
||||||
|
|
||||||
macro_rules! t {
|
macro_rules! t {
|
||||||
@@ -51,35 +50,24 @@ cfg_if! {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn native2io(e: native_tls::Error) -> io::Error {
|
#[tokio::test]
|
||||||
io::Error::new(io::ErrorKind::Other, e)
|
async fn fetch_google() {
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn fetch_google() {
|
|
||||||
drop(env_logger::try_init());
|
drop(env_logger::try_init());
|
||||||
|
|
||||||
// First up, resolve google.com
|
// First up, resolve google.com
|
||||||
let addr = t!("google.com:443".to_socket_addrs()).next().unwrap();
|
let addr = t!("google.com:443".to_socket_addrs()).next().unwrap();
|
||||||
|
|
||||||
// Create an event loop and connect a socket to our resolved address.c
|
let socket = TcpStream::connect(&addr).await.unwrap();
|
||||||
let l = t!(Runtime::new());
|
|
||||||
let client = TcpStream::connect(&addr);
|
|
||||||
|
|
||||||
// Send off the request by first negotiating an SSL handshake, then writing
|
// Send off the request by first negotiating an SSL handshake, then writing
|
||||||
// of our request, then flushing, then finally read off the response.
|
// of our request, then flushing, then finally read off the response.
|
||||||
let data = client
|
let builder = TlsConnector::builder();
|
||||||
.and_then(move |socket| {
|
let connector = t!(builder.build());
|
||||||
let builder = TlsConnector::builder();
|
let connector = tokio_tls::TlsConnector::from(connector);
|
||||||
let connector = t!(builder.build());
|
let mut socket = t!(connector.connect("google.com", socket).await);
|
||||||
let connector = tokio_tls::TlsConnector::from(connector);
|
t!(socket.write_all(b"GET / HTTP/1.0\r\n\r\n").await);
|
||||||
connector.connect("google.com", socket).map_err(native2io)
|
let mut data = Vec::new();
|
||||||
})
|
t!(socket.read_to_end(&mut data).await);
|
||||||
.and_then(|socket| write_all(socket, b"GET / HTTP/1.0\r\n\r\n"))
|
|
||||||
.and_then(|(socket, _)| flush(socket))
|
|
||||||
.and_then(|socket| read_to_end(socket, Vec::new()));
|
|
||||||
|
|
||||||
let (_, data) = t!(l.block_on(data));
|
|
||||||
|
|
||||||
// any response code is fine
|
// any response code is fine
|
||||||
assert!(data.starts_with(b"HTTP/1.0 "));
|
assert!(data.starts_with(b"HTTP/1.0 "));
|
||||||
@@ -89,26 +77,27 @@ fn fetch_google() {
|
|||||||
assert!(data.ends_with("</html>") || data.ends_with("</HTML>"));
|
assert!(data.ends_with("</html>") || data.ends_with("</HTML>"));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn native2io(e: native_tls::Error) -> io::Error {
|
||||||
|
io::Error::new(io::ErrorKind::Other, e)
|
||||||
|
}
|
||||||
|
|
||||||
// see comment in bad.rs for ignore reason
|
// see comment in bad.rs for ignore reason
|
||||||
#[cfg_attr(all(target_os = "macos", feature = "force-openssl"), ignore)]
|
#[cfg_attr(all(target_os = "macos", feature = "force-openssl"), ignore)]
|
||||||
#[test]
|
#[tokio::test]
|
||||||
fn wrong_hostname_error() {
|
async fn wrong_hostname_error() {
|
||||||
drop(env_logger::try_init());
|
drop(env_logger::try_init());
|
||||||
|
|
||||||
let addr = t!("google.com:443".to_socket_addrs()).next().unwrap();
|
let addr = t!("google.com:443".to_socket_addrs()).next().unwrap();
|
||||||
|
|
||||||
let l = t!(Runtime::new());
|
let socket = t!(TcpStream::connect(&addr).await);
|
||||||
let client = TcpStream::connect(&addr);
|
let builder = TlsConnector::builder();
|
||||||
let data = client.and_then(move |socket| {
|
let connector = t!(builder.build());
|
||||||
let builder = TlsConnector::builder();
|
let connector = tokio_tls::TlsConnector::from(connector);
|
||||||
let connector = t!(builder.build());
|
let res = connector
|
||||||
let connector = tokio_tls::TlsConnector::from(connector);
|
.connect("rust-lang.org", socket)
|
||||||
connector
|
.await
|
||||||
.connect("rust-lang.org", socket)
|
.map_err(native2io);
|
||||||
.map_err(native2io)
|
|
||||||
});
|
|
||||||
|
|
||||||
let res = l.block_on(data);
|
|
||||||
assert!(res.is_err());
|
assert!(res.is_err());
|
||||||
assert_bad_hostname_error(&res.err().unwrap());
|
assert_bad_hostname_error(&res.err().unwrap());
|
||||||
}
|
}
|
||||||
|
|||||||
+92
-92
@@ -1,17 +1,17 @@
|
|||||||
#![deny(warnings, rust_2018_idioms)]
|
#![deny(warnings, rust_2018_idioms)]
|
||||||
|
#![feature(async_await)]
|
||||||
|
|
||||||
use cfg_if::cfg_if;
|
use cfg_if::cfg_if;
|
||||||
use env_logger;
|
use env_logger;
|
||||||
use futures::stream::Stream;
|
use futures::join;
|
||||||
use futures::{Future, Poll};
|
use futures::stream::StreamExt;
|
||||||
use native_tls;
|
use native_tls;
|
||||||
use native_tls::{Identity, TlsAcceptor, TlsConnector};
|
use native_tls::{Identity, TlsAcceptor, TlsConnector};
|
||||||
use std::io::{self, Read, Write};
|
use std::io::Write;
|
||||||
|
use std::marker::Unpin;
|
||||||
use std::process::Command;
|
use std::process::Command;
|
||||||
|
use tokio::io::{AsyncReadExt, AsyncWrite, AsyncWriteExt, Error, ErrorKind};
|
||||||
use tokio::net::{TcpListener, TcpStream};
|
use tokio::net::{TcpListener, TcpStream};
|
||||||
use tokio::runtime::Runtime;
|
|
||||||
use tokio_io::io::{copy, read_to_end, shutdown};
|
|
||||||
use tokio_io::{AsyncRead, AsyncWrite};
|
|
||||||
use tokio_tls;
|
use tokio_tls;
|
||||||
|
|
||||||
macro_rules! t {
|
macro_rules! t {
|
||||||
@@ -498,16 +498,30 @@ test suite later.
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn native2io(e: native_tls::Error) -> io::Error {
|
const AMT: usize = 128 * 1024;
|
||||||
io::Error::new(io::ErrorKind::Other, e)
|
|
||||||
|
async fn copy_data<W: AsyncWrite + Unpin>(mut w: W) -> Result<usize, Error> {
|
||||||
|
let mut data = vec![9; AMT as usize];
|
||||||
|
let mut amt = 0;
|
||||||
|
while !data.is_empty() {
|
||||||
|
let written = w.write(&data).await?;
|
||||||
|
if written <= data.len() {
|
||||||
|
amt += written;
|
||||||
|
data.resize(data.len() - written, 0);
|
||||||
|
} else {
|
||||||
|
w.write_all(&data).await?;
|
||||||
|
amt += data.len();
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
println!("remaining: {}", data.len());
|
||||||
|
}
|
||||||
|
Ok(amt)
|
||||||
}
|
}
|
||||||
|
|
||||||
const AMT: u64 = 128 * 1024;
|
#[tokio::test]
|
||||||
|
async fn client_to_server() {
|
||||||
#[test]
|
|
||||||
fn client_to_server() {
|
|
||||||
drop(env_logger::try_init());
|
drop(env_logger::try_init());
|
||||||
let l = t!(Runtime::new());
|
|
||||||
|
|
||||||
// Create a server listening on a port, then figure out what that port is
|
// Create a server listening on a port, then figure out what that port is
|
||||||
let srv = t!(TcpListener::bind(&t!("127.0.0.1:0".parse())));
|
let srv = t!(TcpListener::bind(&t!("127.0.0.1:0".parse())));
|
||||||
@@ -517,30 +531,32 @@ fn client_to_server() {
|
|||||||
|
|
||||||
// Create a future to accept one socket, connect the ssl stream, and then
|
// Create a future to accept one socket, connect the ssl stream, and then
|
||||||
// read all the data from it.
|
// read all the data from it.
|
||||||
let socket = srv.incoming().take(1).collect();
|
let server = async move {
|
||||||
let received = socket
|
let mut incoming = srv.incoming();
|
||||||
.map(|mut socket| socket.remove(0))
|
let socket = t!(incoming.next().await.unwrap());
|
||||||
.and_then(move |socket| server_cx.accept(socket).map_err(native2io))
|
let mut socket = t!(server_cx.accept(socket).await);
|
||||||
.and_then(|socket| read_to_end(socket, Vec::new()));
|
let mut data = Vec::new();
|
||||||
|
t!(socket.read_to_end(&mut data).await);
|
||||||
|
data
|
||||||
|
};
|
||||||
|
|
||||||
// Create a future to connect to our server, connect the ssl stream, and
|
// Create a future to connect to our server, connect the ssl stream, and
|
||||||
// then write a bunch of data to it.
|
// then write a bunch of data to it.
|
||||||
let client = TcpStream::connect(&addr);
|
let client = async move {
|
||||||
let sent = client
|
let socket = t!(TcpStream::connect(&addr).await);
|
||||||
.and_then(move |socket| client_cx.connect("localhost", socket).map_err(native2io))
|
let socket = t!(client_cx.connect("localhost", socket).await);
|
||||||
.and_then(|socket| copy(io::repeat(9).take(AMT), socket))
|
copy_data(socket).await
|
||||||
.and_then(|(amt, _repeat, socket)| shutdown(socket).map(move |_| amt));
|
};
|
||||||
|
|
||||||
// Finally, run everything!
|
// Finally, run everything!
|
||||||
let (amt, (_, data)) = t!(l.block_on(sent.join(received)));
|
let (data, _) = join!(server, client);
|
||||||
assert_eq!(amt, AMT);
|
// assert_eq!(amt, AMT);
|
||||||
assert!(data == vec![9; amt as usize]);
|
assert!(data == vec![9; AMT]);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[tokio::test]
|
||||||
fn server_to_client() {
|
async fn server_to_client() {
|
||||||
drop(env_logger::try_init());
|
drop(env_logger::try_init());
|
||||||
let l = t!(Runtime::new());
|
|
||||||
|
|
||||||
// Create a server listening on a port, then figure out what that port is
|
// Create a server listening on a port, then figure out what that port is
|
||||||
let srv = t!(TcpListener::bind(&t!("127.0.0.1:0".parse())));
|
let srv = t!(TcpListener::bind(&t!("127.0.0.1:0".parse())));
|
||||||
@@ -548,82 +564,66 @@ fn server_to_client() {
|
|||||||
|
|
||||||
let (server_cx, client_cx) = contexts();
|
let (server_cx, client_cx) = contexts();
|
||||||
|
|
||||||
let socket = srv.incoming().take(1).collect();
|
let server = async move {
|
||||||
let sent = socket
|
let mut incoming = srv.incoming();
|
||||||
.map(|mut socket| socket.remove(0))
|
let socket = t!(incoming.next().await.unwrap());
|
||||||
.and_then(move |socket| server_cx.accept(socket).map_err(native2io))
|
let socket = t!(server_cx.accept(socket).await);
|
||||||
.and_then(|socket| copy(io::repeat(9).take(AMT), socket))
|
copy_data(socket).await
|
||||||
.and_then(|(amt, _repeat, socket)| shutdown(socket).map(move |_| amt));
|
};
|
||||||
|
|
||||||
let client = TcpStream::connect(&addr);
|
let client = async move {
|
||||||
let received = client
|
let socket = t!(TcpStream::connect(&addr).await);
|
||||||
.and_then(move |socket| client_cx.connect("localhost", socket).map_err(native2io))
|
let mut socket = t!(client_cx.connect("localhost", socket).await);
|
||||||
.and_then(|socket| read_to_end(socket, Vec::new()));
|
let mut data = Vec::new();
|
||||||
|
t!(socket.read_to_end(&mut data).await);
|
||||||
|
data
|
||||||
|
};
|
||||||
|
|
||||||
// Finally, run everything!
|
// Finally, run everything!
|
||||||
let (amt, (_, data)) = t!(l.block_on(sent.join(received)));
|
let (_, data) = join!(server, client);
|
||||||
assert_eq!(amt, AMT);
|
// assert_eq!(amt, AMT);
|
||||||
assert!(data == vec![9; amt as usize]);
|
assert!(data == vec![9; AMT]);
|
||||||
}
|
}
|
||||||
|
|
||||||
struct OneByte<S> {
|
#[tokio::test]
|
||||||
inner: S,
|
async fn one_byte_at_a_time() {
|
||||||
}
|
const AMT: usize = 1024;
|
||||||
|
|
||||||
impl<S: Read> Read for OneByte<S> {
|
|
||||||
fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
|
|
||||||
self.inner.read(&mut buf[..1])
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl<S: Write> Write for OneByte<S> {
|
|
||||||
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
|
|
||||||
self.inner.write(&buf[..1])
|
|
||||||
}
|
|
||||||
|
|
||||||
fn flush(&mut self) -> io::Result<()> {
|
|
||||||
self.inner.flush()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl<S: AsyncRead> AsyncRead for OneByte<S> {}
|
|
||||||
impl<S: AsyncWrite> AsyncWrite for OneByte<S> {
|
|
||||||
fn shutdown(&mut self) -> Poll<(), io::Error> {
|
|
||||||
self.inner.shutdown()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn one_byte_at_a_time() {
|
|
||||||
const AMT: u64 = 1024;
|
|
||||||
drop(env_logger::try_init());
|
drop(env_logger::try_init());
|
||||||
let l = t!(Runtime::new());
|
|
||||||
|
|
||||||
let srv = t!(TcpListener::bind(&t!("127.0.0.1:0".parse())));
|
let srv = t!(TcpListener::bind(&t!("127.0.0.1:0".parse())));
|
||||||
let addr = t!(srv.local_addr());
|
let addr = t!(srv.local_addr());
|
||||||
|
|
||||||
let (server_cx, client_cx) = contexts();
|
let (server_cx, client_cx) = contexts();
|
||||||
|
|
||||||
let socket = srv.incoming().take(1).collect();
|
let server = async move {
|
||||||
let sent = socket
|
let mut incoming = srv.incoming();
|
||||||
.map(|mut socket| socket.remove(0))
|
let socket = t!(incoming.next().await.unwrap());
|
||||||
.and_then(move |socket| {
|
let mut socket = t!(server_cx.accept(socket).await);
|
||||||
server_cx
|
let mut amt = 0;
|
||||||
.accept(OneByte { inner: socket })
|
for b in std::iter::repeat(9).take(AMT) {
|
||||||
.map_err(native2io)
|
let data = [b as u8];
|
||||||
})
|
t!(socket.write_all(&data).await);
|
||||||
.and_then(|socket| copy(io::repeat(9).take(AMT), socket))
|
amt += 1;
|
||||||
.and_then(|(amt, _repeat, socket)| shutdown(socket).map(move |_| amt));
|
}
|
||||||
|
amt
|
||||||
|
};
|
||||||
|
|
||||||
let client = TcpStream::connect(&addr);
|
let client = async move {
|
||||||
let received = client
|
let socket = t!(TcpStream::connect(&addr).await);
|
||||||
.and_then(move |socket| {
|
let mut socket = t!(client_cx.connect("localhost", socket).await);
|
||||||
let socket = OneByte { inner: socket };
|
let mut data = Vec::new();
|
||||||
client_cx.connect("localhost", socket).map_err(native2io)
|
loop {
|
||||||
})
|
let mut buf = [0; 1];
|
||||||
.and_then(|socket| read_to_end(socket, Vec::new()));
|
match socket.read_exact(&mut buf).await {
|
||||||
|
Ok(_) => data.extend_from_slice(&buf),
|
||||||
|
Err(ref err) if err.kind() == ErrorKind::UnexpectedEof => break,
|
||||||
|
Err(err) => panic!(err),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
data
|
||||||
|
};
|
||||||
|
|
||||||
let (amt, (_, data)) = t!(l.block_on(sent.join(received)));
|
let (amt, data) = join!(server, client);
|
||||||
assert_eq!(amt, AMT);
|
assert_eq!(amt, AMT);
|
||||||
assert!(data == vec![9; amt as usize]);
|
assert!(data == vec![9; AMT as usize]);
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user