mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-17 00:00:11 +02:00
We know that the future will never persist beyond this stack frame, so we can just leave its ownership on the stack frame itself and receive notifications off the event loop that we need to poll it.
86 lines
2.2 KiB
Rust
86 lines
2.2 KiB
Rust
extern crate env_logger;
|
|
extern crate futures;
|
|
extern crate futures_mio;
|
|
|
|
use std::net::{TcpListener, TcpStream};
|
|
use std::sync::mpsc::channel;
|
|
use std::thread;
|
|
|
|
use futures::Future;
|
|
use futures::stream::Stream;
|
|
|
|
macro_rules! t {
|
|
($e:expr) => (match $e {
|
|
Ok(e) => e,
|
|
Err(e) => panic!("{} failed with {:?}", stringify!($e), e),
|
|
})
|
|
}
|
|
|
|
#[test]
|
|
fn connect() {
|
|
drop(env_logger::init());
|
|
let mut l = t!(futures_mio::Loop::new());
|
|
let srv = t!(TcpListener::bind("127.0.0.1:0"));
|
|
let addr = t!(srv.local_addr());
|
|
let t = thread::spawn(move || {
|
|
t!(srv.accept()).0
|
|
});
|
|
|
|
let stream = l.handle().tcp_connect(&addr);
|
|
let mine = t!(l.run(stream));
|
|
let theirs = t.join().unwrap();
|
|
|
|
assert_eq!(t!(mine.local_addr()), t!(theirs.peer_addr()));
|
|
assert_eq!(t!(theirs.local_addr()), t!(mine.peer_addr()));
|
|
}
|
|
|
|
#[test]
|
|
fn accept() {
|
|
drop(env_logger::init());
|
|
let mut l = t!(futures_mio::Loop::new());
|
|
let srv = l.handle().tcp_listen(&"127.0.0.1:0".parse().unwrap());
|
|
let srv = t!(l.run(srv));
|
|
let addr = t!(srv.local_addr());
|
|
|
|
let (tx, rx) = channel();
|
|
let client = srv.incoming().map(move |t| {
|
|
tx.send(()).unwrap();
|
|
t.0
|
|
}).into_future().map_err(|e| e.0);
|
|
assert!(rx.try_recv().is_err());
|
|
let t = thread::spawn(move || {
|
|
TcpStream::connect(&addr).unwrap()
|
|
});
|
|
|
|
let (mine, _remaining) = t!(l.run(client));
|
|
let mine = mine.unwrap();
|
|
let theirs = t.join().unwrap();
|
|
|
|
assert_eq!(t!(mine.local_addr()), t!(theirs.peer_addr()));
|
|
assert_eq!(t!(theirs.local_addr()), t!(mine.peer_addr()));
|
|
}
|
|
|
|
#[test]
|
|
fn accept2() {
|
|
drop(env_logger::init());
|
|
let mut l = t!(futures_mio::Loop::new());
|
|
let srv = l.handle().tcp_listen(&"127.0.0.1:0".parse().unwrap());
|
|
let srv = t!(l.run(srv));
|
|
let addr = t!(srv.local_addr());
|
|
|
|
let t = thread::spawn(move || {
|
|
TcpStream::connect(&addr).unwrap()
|
|
});
|
|
|
|
let (tx, rx) = channel();
|
|
let client = srv.incoming().map(move |t| {
|
|
tx.send(()).unwrap();
|
|
t.0
|
|
}).into_future().map_err(|e| e.0);
|
|
assert!(rx.try_recv().is_err());
|
|
|
|
let (mine, _remaining) = t!(l.run(client));
|
|
mine.unwrap();
|
|
t.join().unwrap();
|
|
}
|