From a79483750ff550d95b95cb19fb211637f67c1745 Mon Sep 17 00:00:00 2001 From: Carl Lerche Date: Wed, 10 Jul 2019 14:21:20 -0700 Subject: [PATCH] tokio: update echo example (#1283) --- tokio/examples/echo.rs | 73 +++++++++++++++++++++++++++++++++++ tokio/examples/hello_world.rs | 1 + 2 files changed, 74 insertions(+) create mode 100644 tokio/examples/echo.rs diff --git a/tokio/examples/echo.rs b/tokio/examples/echo.rs new file mode 100644 index 000000000..54ae870b5 --- /dev/null +++ b/tokio/examples/echo.rs @@ -0,0 +1,73 @@ +//! A "hello world" echo server with Tokio +//! +//! This server will create a TCP listener, accept connections in a loop, and +//! write back everything that's read off of each TCP connection. +//! +//! Because the Tokio runtime uses a thread pool, each TCP connection is +//! processed concurrently with all other TCP connections across multiple +//! threads. +//! +//! To see this server in action, you can run this in one terminal: +//! +//! cargo run --example echo +//! +//! and in another terminal you can run: +//! +//! cargo run --example connect 127.0.0.1:8080 +//! +//! Each line you type in to the `connect` terminal should be echo'd back to +//! you! If you open up multiple terminals running the `connect` example you +//! should be able to see them all make progress simultaneously. + +#![feature(async_await)] +#![deny(warnings, rust_2018_idioms)] + +use tokio; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::TcpListener; + +use std::env; +use std::net::SocketAddr; + +#[tokio::main] +async fn main() { + // Allow passing an address to listen on as the first argument of this + // program, but otherwise we'll just set up our TCP listener on + // 127.0.0.1:8080 for connections. + let addr = env::args().nth(1).unwrap_or("127.0.0.1:8080".to_string()); + let addr = addr.parse::().unwrap(); + + // Next up we create a TCP listener which will listen for incoming + // connections. This TCP listener is bound to the address we determined + // above and must be associated with an event loop. + let mut listener = TcpListener::bind(&addr).unwrap(); + println!("Listening on: {}", addr); + + loop { + // Asynchronously wait for an inbound socket. + let (mut socket, _) = listener.accept().await.unwrap(); + + // And this is where much of the magic of this server happens. We + // crucially want all clients to make progress concurrently, rather than + // blocking one on completion of another. To achieve this we use the + // `tokio::spawn` function to execute the work in the background. + // + // Essentially here we're executing a new task to run concurrently, + // which will allow all of our clients to be processed concurrently. + + tokio::spawn(async move { + let mut buf = [0; 1024]; + + // In a loop, read data from the socket and write the data back. + loop { + let n = socket.read(&mut buf).await.unwrap(); + + if n == 0 { + return; + } + + socket.write_all(&buf[0..n]).await.unwrap(); + } + }); + } +} diff --git a/tokio/examples/hello_world.rs b/tokio/examples/hello_world.rs index 67afa9d1e..8e77c93ee 100644 --- a/tokio/examples/hello_world.rs +++ b/tokio/examples/hello_world.rs @@ -27,6 +27,7 @@ pub async fn main() { // Note that this is the Tokio TcpStream, which is fully async. let mut stream = TcpStream::connect(&addr).await.unwrap(); println!("created stream"); + let result = stream.write(b"hello world\n").await; println!("wrote to stream; success={:?}", result.is_ok()); }