mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-09-08 00:00:13 +02:00
No more need for lazy in chat example
This commit is contained in:
+3
-19
@@ -37,6 +37,7 @@ fn main() {
|
|||||||
|
|
||||||
let srv = socket.incoming().for_each(move |(stream, addr)| {
|
let srv = socket.incoming().for_each(move |(stream, addr)| {
|
||||||
println!("New Connection: {}", addr);
|
println!("New Connection: {}", addr);
|
||||||
|
let (reader, writer) = stream.split();
|
||||||
|
|
||||||
// Create a channel for our stream, which other sockets will use to
|
// Create a channel for our stream, which other sockets will use to
|
||||||
// send us messages. Then register our address with the stream to send
|
// send us messages. Then register our address with the stream to send
|
||||||
@@ -44,19 +45,10 @@ fn main() {
|
|||||||
let (tx, rx) = tokio_core::channel::channel(&handle).unwrap();
|
let (tx, rx) = tokio_core::channel::channel(&handle).unwrap();
|
||||||
connections.borrow_mut().insert(addr, tx);
|
connections.borrow_mut().insert(addr, tx);
|
||||||
|
|
||||||
// Note that below we're calling `spawn` to spawn a new future for this
|
|
||||||
// connection. As a result we use `futures::lazy` here to ensure that
|
|
||||||
// the call to `.split()` happens on the right task.
|
|
||||||
//
|
|
||||||
// This `split` will give us a read/write half to work with each portion
|
|
||||||
// of the socket separately.
|
|
||||||
let pair = futures::lazy(|| Ok(stream.split()));
|
|
||||||
|
|
||||||
// Define here what we do for the actual I/O. That is, read a bunch of
|
// Define here what we do for the actual I/O. That is, read a bunch of
|
||||||
// lines from the socket and dispatch them while we also write any lines
|
// lines from the socket and dispatch them while we also write any lines
|
||||||
// from other sockets.
|
// from other sockets.
|
||||||
let connections_inner = connections.clone();
|
let connections_inner = connections.clone();
|
||||||
let pair = pair.map(move |(reader, writer)| {
|
|
||||||
let reader = BufReader::new(reader);
|
let reader = BufReader::new(reader);
|
||||||
|
|
||||||
// Model the read portion of this socket by mapping an infinite
|
// Model the read portion of this socket by mapping an infinite
|
||||||
@@ -108,20 +100,12 @@ fn main() {
|
|||||||
amt
|
amt
|
||||||
});
|
});
|
||||||
|
|
||||||
(socket_reader, socket_writer)
|
|
||||||
});
|
|
||||||
|
|
||||||
// Now that we've got futures representing each half of the socket, we
|
// Now that we've got futures representing each half of the socket, we
|
||||||
// use the `select` combinator to wait for either half to be done to
|
// use the `select` combinator to wait for either half to be done to
|
||||||
// tear down the other. Then we spawn off the result.
|
// tear down the other. Then we spawn off the result.
|
||||||
let connections = connections.clone();
|
let connections = connections.clone();
|
||||||
let addr = addr;
|
let connection = socket_reader.map(|_| ()).select(socket_writer.map(|_| ()));
|
||||||
handle.spawn(pair.and_then(|(reader, writer)| {
|
handle.spawn(connection.then(move |_| {
|
||||||
let reader = reader.map(|_| ());
|
|
||||||
let writer = writer.map(|_| ());
|
|
||||||
|
|
||||||
reader.select(writer)
|
|
||||||
}).then(move |_| {
|
|
||||||
connections.borrow_mut().remove(&addr);
|
connections.borrow_mut().remove(&addr);
|
||||||
println!("Connection {} closed.", addr);
|
println!("Connection {} closed.", addr);
|
||||||
Ok(())
|
Ok(())
|
||||||
|
|||||||
Reference in New Issue
Block a user