mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-28 00:00:11 +02:00
Track futures tokio-reform branch (#88)
This patch also updates tests and examples to remove deprecated API usage.
This commit is contained in:
+2
-1
@@ -8,6 +8,7 @@ use std::thread;
|
||||
use std::io::{Read, Write, BufReader, BufWriter};
|
||||
|
||||
use futures::Future;
|
||||
use futures::future::blocking;
|
||||
use futures::stream::Stream;
|
||||
use tokio_io::io::copy;
|
||||
use tokio::net::TcpListener;
|
||||
@@ -54,7 +55,7 @@ fn echo_server() {
|
||||
copy(a, b)
|
||||
});
|
||||
|
||||
let (amt, _, _) = t!(copied.wait());
|
||||
let (amt, _, _) = t!(blocking(copied).wait());
|
||||
let (expected, t2) = t.join().unwrap();
|
||||
let actual = t2.join().unwrap();
|
||||
|
||||
|
||||
+2
-1
@@ -7,6 +7,7 @@ use std::thread;
|
||||
use std::io::{Write, Read};
|
||||
|
||||
use futures::Future;
|
||||
use futures::future::blocking;
|
||||
use futures::stream::Stream;
|
||||
use tokio_io::io::read_to_end;
|
||||
use tokio::net::TcpListener;
|
||||
@@ -42,7 +43,7 @@ fn chain_clients() {
|
||||
read_to_end(a.chain(b).chain(c), Vec::new())
|
||||
});
|
||||
|
||||
let (_, data) = t!(copied.wait());
|
||||
let (_, data) = t!(blocking(copied).wait());
|
||||
t.join().unwrap();
|
||||
|
||||
assert_eq!(data, b"foo bar baz");
|
||||
|
||||
+4
-4
@@ -4,7 +4,7 @@ extern crate futures;
|
||||
use std::thread;
|
||||
use std::net;
|
||||
|
||||
use futures::future;
|
||||
use futures::{future, stream};
|
||||
use futures::prelude::*;
|
||||
use futures::sync::oneshot;
|
||||
use tokio::net::TcpListener;
|
||||
@@ -17,7 +17,7 @@ fn tcp_doesnt_block() {
|
||||
let listener = net::TcpListener::bind("127.0.0.1:0").unwrap();
|
||||
let listener = TcpListener::from_std(listener, &handle).unwrap();
|
||||
drop(core);
|
||||
assert!(listener.incoming().wait().next().unwrap().is_err());
|
||||
assert!(stream::blocking(listener.incoming()).next().unwrap().is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -34,9 +34,9 @@ fn drop_wakes() {
|
||||
drop(tx);
|
||||
future::ok(())
|
||||
});
|
||||
assert!(new_socket.join(drop_tx).wait().is_err());
|
||||
assert!(future::blocking(new_socket.join(drop_tx)).wait().is_err());
|
||||
});
|
||||
drop(rx.wait());
|
||||
drop(future::blocking(rx).wait());
|
||||
drop(core);
|
||||
t.join().unwrap();
|
||||
}
|
||||
|
||||
+2
-1
@@ -8,6 +8,7 @@ use std::net::TcpStream;
|
||||
use std::thread;
|
||||
|
||||
use futures::Future;
|
||||
use futures::future::blocking;
|
||||
use futures::stream::Stream;
|
||||
use tokio::net::TcpListener;
|
||||
use tokio_io::AsyncRead;
|
||||
@@ -44,7 +45,7 @@ fn echo_server() {
|
||||
let halves = client.map(|s| s.split());
|
||||
let copied = halves.and_then(|(a, b)| copy(a, b));
|
||||
|
||||
let (amt, _, _) = t!(copied.wait());
|
||||
let (amt, _, _) = t!(blocking(copied).wait());
|
||||
t.join().unwrap();
|
||||
|
||||
assert_eq!(amt, msg.len() as u64 * 1024);
|
||||
|
||||
+2
-1
@@ -3,6 +3,7 @@ extern crate tokio;
|
||||
|
||||
use std::thread;
|
||||
|
||||
use futures::future::blocking;
|
||||
use futures::prelude::*;
|
||||
use tokio::net::{TcpStream, TcpListener};
|
||||
|
||||
@@ -23,7 +24,7 @@ fn hammer() {
|
||||
let theirs = srv.incoming().into_future()
|
||||
.map(|(s, _)| s.unwrap())
|
||||
.map_err(|(s, _)| s);
|
||||
let (mine, theirs) = t!(mine.join(theirs).wait());
|
||||
let (mine, theirs) = t!(blocking(mine.join(theirs)).wait());
|
||||
|
||||
assert_eq!(t!(mine.local_addr()), t!(theirs.peer_addr()));
|
||||
assert_eq!(t!(theirs.local_addr()), t!(mine.peer_addr()));
|
||||
|
||||
+2
-1
@@ -7,6 +7,7 @@ use std::thread;
|
||||
use std::io::{Write, Read};
|
||||
|
||||
use futures::Future;
|
||||
use futures::future::blocking;
|
||||
use futures::stream::Stream;
|
||||
use tokio_io::io::read_to_end;
|
||||
use tokio::net::TcpListener;
|
||||
@@ -36,7 +37,7 @@ fn limit() {
|
||||
read_to_end(a.take(4), Vec::new())
|
||||
});
|
||||
|
||||
let (_, data) = t!(copied.wait());
|
||||
let (_, data) = t!(blocking(copied).wait());
|
||||
t.join().unwrap();
|
||||
|
||||
assert_eq!(data, b"foo ");
|
||||
|
||||
@@ -10,7 +10,7 @@ use std::net::Shutdown;
|
||||
|
||||
use bytes::{BytesMut, BufMut};
|
||||
use futures::{Future, Stream, Sink};
|
||||
use futures::future::Executor;
|
||||
use futures::future::{blocking, Executor};
|
||||
use futures_cpupool::CpuPool;
|
||||
use tokio::net::{TcpListener, TcpStream};
|
||||
use tokio_io::codec::{Encoder, Decoder};
|
||||
@@ -68,20 +68,20 @@ fn echo() {
|
||||
pool.execute(srv.map_err(|e| panic!("srv error: {}", e))).unwrap();
|
||||
|
||||
let client = TcpStream::connect(&addr);
|
||||
let client = client.wait().unwrap();
|
||||
let (client, _) = write_all(client, b"a\n").wait().unwrap();
|
||||
let (client, buf, amt) = read(client, vec![0; 1024]).wait().unwrap();
|
||||
let client = blocking(client).wait().unwrap();
|
||||
let (client, _) = blocking(write_all(client, b"a\n")).wait().unwrap();
|
||||
let (client, buf, amt) = blocking(read(client, vec![0; 1024])).wait().unwrap();
|
||||
assert_eq!(amt, 2);
|
||||
assert_eq!(&buf[..2], b"a\n");
|
||||
|
||||
let (client, _) = write_all(client, b"\n").wait().unwrap();
|
||||
let (client, buf, amt) = read(client, buf).wait().unwrap();
|
||||
let (client, _) = blocking(write_all(client, b"\n")).wait().unwrap();
|
||||
let (client, buf, amt) = blocking(read(client, buf)).wait().unwrap();
|
||||
assert_eq!(amt, 1);
|
||||
assert_eq!(&buf[..1], b"\n");
|
||||
|
||||
let (client, _) = write_all(client, b"b").wait().unwrap();
|
||||
let (client, _) = blocking(write_all(client, b"b")).wait().unwrap();
|
||||
client.shutdown(Shutdown::Write).unwrap();
|
||||
let (_client, buf, amt) = read(client, buf).wait().unwrap();
|
||||
let (_client, buf, amt) = blocking(read(client, buf)).wait().unwrap();
|
||||
assert_eq!(amt, 1);
|
||||
assert_eq!(&buf[..1], b"b");
|
||||
}
|
||||
|
||||
+2
-2
@@ -13,7 +13,7 @@ use std::os::unix::io::{AsRawFd, FromRawFd};
|
||||
use std::thread;
|
||||
use std::time::Duration;
|
||||
|
||||
use futures::prelude::*;
|
||||
use futures::future::blocking;
|
||||
use mio::event::Evented;
|
||||
use mio::unix::{UnixReady, EventedFd};
|
||||
use mio::{PollOpt, Ready, Token};
|
||||
@@ -81,7 +81,7 @@ fn hup() {
|
||||
let source = PollEvented::new(MyFile::new(read), &handle).unwrap();
|
||||
|
||||
let reader = read_to_end(source, Vec::new());
|
||||
let (_, content) = t!(reader.wait());
|
||||
let (_, content) = t!(blocking(reader).wait());
|
||||
assert_eq!(&b"Hello!\nGood bye!\n"[..], &content[..]);
|
||||
t.join().unwrap();
|
||||
}
|
||||
|
||||
@@ -8,6 +8,7 @@ use std::net::TcpStream;
|
||||
use std::thread;
|
||||
|
||||
use futures::Future;
|
||||
use futures::future::blocking;
|
||||
use futures::stream::Stream;
|
||||
use tokio_io::io::copy;
|
||||
use tokio_io::AsyncRead;
|
||||
@@ -48,7 +49,7 @@ fn echo_server() {
|
||||
.take(2)
|
||||
.collect();
|
||||
|
||||
t!(future.wait());
|
||||
t!(blocking(future).wait());
|
||||
|
||||
t.join().unwrap();
|
||||
}
|
||||
|
||||
+4
-3
@@ -7,6 +7,7 @@ use std::sync::mpsc::channel;
|
||||
use std::thread;
|
||||
|
||||
use futures::Future;
|
||||
use futures::future::blocking;
|
||||
use futures::stream::Stream;
|
||||
use tokio::net::{TcpListener, TcpStream};
|
||||
|
||||
@@ -27,7 +28,7 @@ fn connect() {
|
||||
});
|
||||
|
||||
let stream = TcpStream::connect(&addr);
|
||||
let mine = t!(stream.wait());
|
||||
let mine = t!(blocking(stream).wait());
|
||||
let theirs = t.join().unwrap();
|
||||
|
||||
assert_eq!(t!(mine.local_addr()), t!(theirs.peer_addr()));
|
||||
@@ -50,7 +51,7 @@ fn accept() {
|
||||
net::TcpStream::connect(&addr).unwrap()
|
||||
});
|
||||
|
||||
let (mine, _remaining) = t!(client.wait());
|
||||
let (mine, _remaining) = t!(blocking(client).wait());
|
||||
let mine = mine.unwrap();
|
||||
let theirs = t.join().unwrap();
|
||||
|
||||
@@ -75,7 +76,7 @@ fn accept2() {
|
||||
}).into_future().map_err(|e| e.0);
|
||||
assert!(rx.try_recv().is_err());
|
||||
|
||||
let (mine, _remaining) = t!(client.wait());
|
||||
let (mine, _remaining) = t!(blocking(client).wait());
|
||||
mine.unwrap();
|
||||
t.join().unwrap();
|
||||
}
|
||||
|
||||
+7
-6
@@ -7,6 +7,7 @@ use std::io;
|
||||
use std::net::SocketAddr;
|
||||
|
||||
use futures::{Future, Poll, Stream, Sink};
|
||||
use futures::future::blocking;
|
||||
use tokio::net::{UdpSocket, UdpCodec};
|
||||
|
||||
macro_rules! t {
|
||||
@@ -25,7 +26,7 @@ fn send_messages<S: SendFn + Clone, R: RecvFn + Clone>(send: S, recv: R) {
|
||||
{
|
||||
let send = SendMessage::new(a, send.clone(), b_addr, b"1234");
|
||||
let recv = RecvMessage::new(b, recv.clone(), a_addr, b"1234");
|
||||
let (sendt, received) = t!(send.join(recv).wait());
|
||||
let (sendt, received) = t!(blocking(send.join(recv)).wait());
|
||||
a = sendt;
|
||||
b = received;
|
||||
}
|
||||
@@ -33,7 +34,7 @@ fn send_messages<S: SendFn + Clone, R: RecvFn + Clone>(send: S, recv: R) {
|
||||
{
|
||||
let send = SendMessage::new(a, send, b_addr, b"");
|
||||
let recv = RecvMessage::new(b, recv, a_addr, b"");
|
||||
t!(send.join(recv).wait());
|
||||
t!(blocking(send.join(recv)).wait());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -172,7 +173,7 @@ fn send_dgrams() {
|
||||
{
|
||||
let send = a.send_dgram(&b"4321"[..], &b_addr);
|
||||
let recv = b.recv_dgram(&mut buf[..]);
|
||||
let (sendt, received) = t!(send.join(recv).wait());
|
||||
let (sendt, received) = t!(blocking(send.join(recv)).wait());
|
||||
assert_eq!(received.2, 4);
|
||||
assert_eq!(&received.1[..4], b"4321");
|
||||
a = sendt.0;
|
||||
@@ -182,7 +183,7 @@ fn send_dgrams() {
|
||||
{
|
||||
let send = a.send_dgram(&b""[..], &b_addr);
|
||||
let recv = b.recv_dgram(&mut buf[..]);
|
||||
let received = t!(send.join(recv).wait()).1;
|
||||
let received = t!(blocking(send.join(recv)).wait()).1;
|
||||
assert_eq!(received.2, 0);
|
||||
}
|
||||
}
|
||||
@@ -225,7 +226,7 @@ fn send_framed() {
|
||||
|
||||
let send = a.send(&b"4567"[..]);
|
||||
let recv = b.into_future().map_err(|e| e.0);
|
||||
let (sendt, received) = t!(send.join(recv).wait());
|
||||
let (sendt, received) = t!(blocking(send.join(recv)).wait());
|
||||
assert_eq!(received.0, Some(()));
|
||||
|
||||
a_soc = sendt.into_inner();
|
||||
@@ -238,7 +239,7 @@ fn send_framed() {
|
||||
|
||||
let send = a.send(&b""[..]);
|
||||
let recv = b.into_future().map_err(|e| e.0);
|
||||
let received = t!(send.join(recv).wait()).1;
|
||||
let received = t!(blocking(send.join(recv)).wait()).1;
|
||||
assert_eq!(received.0, Some(()));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user