mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-09-08 00:00:13 +02:00
Touch up examples to ensure consistency
This commit is contained in:
+13
-1
@@ -3,6 +3,19 @@
|
|||||||
//! This is a simple line-based server which accepts connections, reads lines
|
//! This is a simple line-based server which accepts connections, reads lines
|
||||||
//! from those connections, and broadcasts the lines to all other connected
|
//! from those connections, and broadcasts the lines to all other connected
|
||||||
//! clients. In a sense this is a bit of a "poor man's chat server".
|
//! clients. In a sense this is a bit of a "poor man's chat server".
|
||||||
|
//!
|
||||||
|
//! You can test this out by running:
|
||||||
|
//!
|
||||||
|
//! cargo run --example chat
|
||||||
|
//!
|
||||||
|
//! And then in another window run:
|
||||||
|
//!
|
||||||
|
//! nc -4 localhost 8080
|
||||||
|
//!
|
||||||
|
//! You can run the second command in multiple windows and then chat between the
|
||||||
|
//! two, seeing the messages from the other client as they're received. For all
|
||||||
|
//! connected clients they'll all join the same room and see everyone else's
|
||||||
|
//! messages.
|
||||||
|
|
||||||
extern crate tokio_core;
|
extern crate tokio_core;
|
||||||
extern crate futures;
|
extern crate futures;
|
||||||
@@ -118,4 +131,3 @@ fn main() {
|
|||||||
// execute server
|
// execute server
|
||||||
core.run(srv).unwrap();
|
core.run(srv).unwrap();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -2,15 +2,11 @@
|
|||||||
//!
|
//!
|
||||||
//! If you're on unix you can test this out by in one terminal executing:
|
//! If you're on unix you can test this out by in one terminal executing:
|
||||||
//!
|
//!
|
||||||
//! ```sh
|
//! cargo run --example echo-udp
|
||||||
//! $ cargo run --example echo-udp
|
|
||||||
//! ```
|
|
||||||
//!
|
//!
|
||||||
//! and in another terminal you can run:
|
//! and in another terminal you can run:
|
||||||
//!
|
//!
|
||||||
//! ```sh
|
//! nc -4u localhost 8080
|
||||||
//! $ nc -4u localhost 8080
|
|
||||||
//! ```
|
|
||||||
//!
|
//!
|
||||||
//! Each line you type in to the `nc` terminal should be echo'd back to you!
|
//! Each line you type in to the `nc` terminal should be echo'd back to you!
|
||||||
|
|
||||||
|
|||||||
+2
-6
@@ -2,15 +2,11 @@
|
|||||||
//!
|
//!
|
||||||
//! If you're on unix you can test this out by in one terminal executing:
|
//! If you're on unix you can test this out by in one terminal executing:
|
||||||
//!
|
//!
|
||||||
//! ```sh
|
//! cargo run --example echo
|
||||||
//! $ cargo run --example echo
|
|
||||||
//! ```
|
|
||||||
//!
|
//!
|
||||||
//! and in another terminal you can run:
|
//! and in another terminal you can run:
|
||||||
//!
|
//!
|
||||||
//! ```sh
|
//! nc -4 localhost 8080
|
||||||
//! $ nc localhost 8080
|
|
||||||
//! ```
|
|
||||||
//!
|
//!
|
||||||
//! Each line you type in to the `nc` terminal should be echo'd back to you!
|
//! Each line you type in to the `nc` terminal should be echo'd back to you!
|
||||||
|
|
||||||
|
|||||||
+22
-2
@@ -1,14 +1,34 @@
|
|||||||
|
//! A small example of a server that accepts TCP connections and writes out
|
||||||
|
//! `Hello!` to them, afterwards closing the connection.
|
||||||
|
//!
|
||||||
|
//! You can test this out by running:
|
||||||
|
//!
|
||||||
|
//! cargo run --example hello
|
||||||
|
//!
|
||||||
|
//! and then in another terminal executing
|
||||||
|
//!
|
||||||
|
//! nc -4 localhost 8080
|
||||||
|
//!
|
||||||
|
//! You should see `Hello!` printed out and then the `nc` program will exit.
|
||||||
|
|
||||||
extern crate futures;
|
extern crate futures;
|
||||||
extern crate tokio_core;
|
extern crate tokio_core;
|
||||||
|
extern crate env_logger;
|
||||||
|
|
||||||
|
use std::env;
|
||||||
|
use std::net::SocketAddr;
|
||||||
|
|
||||||
use futures::stream::Stream;
|
use futures::stream::Stream;
|
||||||
use tokio_core::reactor::Core;
|
use tokio_core::reactor::Core;
|
||||||
use tokio_core::net::TcpListener;
|
use tokio_core::net::TcpListener;
|
||||||
|
|
||||||
fn main() {
|
fn main() {
|
||||||
|
env_logger::init().unwrap();
|
||||||
|
let addr = env::args().nth(1).unwrap_or("127.0.0.1:8080".to_string());
|
||||||
|
let addr = addr.parse::<SocketAddr>().unwrap();
|
||||||
|
|
||||||
let mut core = Core::new().unwrap();
|
let mut core = Core::new().unwrap();
|
||||||
let address = "127.0.0.1:8080".parse().unwrap();
|
let listener = TcpListener::bind(&addr, &core.handle()).unwrap();
|
||||||
let listener = TcpListener::bind(&address, &core.handle()).unwrap();
|
|
||||||
|
|
||||||
let addr = listener.local_addr().unwrap();
|
let addr = listener.local_addr().unwrap();
|
||||||
println!("Listening for connections on {}", addr);
|
println!("Listening for connections on {}", addr);
|
||||||
|
|||||||
+13
-1
@@ -1,7 +1,19 @@
|
|||||||
//! A small server that writes as many nul bytes on all connections it receives.
|
//! A small server that writes as many nul bytes on all connections it receives.
|
||||||
//!
|
//!
|
||||||
//! There is no concurrency in this server, only one connection is written to at
|
//! There is no concurrency in this server, only one connection is written to at
|
||||||
//! a time.
|
//! a time. You can use this as a benchmark for the raw performance of writing
|
||||||
|
//! data to a socket by measuring how much data is being written on each
|
||||||
|
//! connection.
|
||||||
|
//!
|
||||||
|
//! Typically you'll want to run this example with:
|
||||||
|
//!
|
||||||
|
//! cargo run --example sink --release
|
||||||
|
//!
|
||||||
|
//! And then you can connect to it via:
|
||||||
|
//!
|
||||||
|
//! nc -4 localhost 8080 > /dev/null
|
||||||
|
//!
|
||||||
|
//! You should see your CPUs light up as data's being shove into the ether.
|
||||||
|
|
||||||
extern crate env_logger;
|
extern crate env_logger;
|
||||||
extern crate futures;
|
extern crate futures;
|
||||||
|
|||||||
+53
-90
@@ -1,59 +1,38 @@
|
|||||||
|
//! This is a basic example of leveraging `UdpCodec` to create a simple UDP
|
||||||
|
//! client and server which speak a custom protocol.
|
||||||
|
//!
|
||||||
|
//! Here we're using the a custom codec to convert a UDP socket to a stream of
|
||||||
|
//! client messages. These messages are then processed and returned back as a
|
||||||
|
//! new message with a new destination. Overall, we then use this to construct a
|
||||||
|
//! "ping pong" pair where two sockets are sending messages back and forth.
|
||||||
|
|
||||||
extern crate tokio_core;
|
extern crate tokio_core;
|
||||||
extern crate env_logger;
|
extern crate env_logger;
|
||||||
extern crate futures;
|
extern crate futures;
|
||||||
|
|
||||||
#[macro_use]
|
|
||||||
extern crate log;
|
|
||||||
|
|
||||||
use std::io;
|
use std::io;
|
||||||
use std::net::{SocketAddr};
|
use std::net::SocketAddr;
|
||||||
use futures::{future, Future, Stream, Sink};
|
|
||||||
use tokio_core::net::{UdpSocket, UdpCodec};
|
|
||||||
use tokio_core::reactor::{Core, Timeout};
|
|
||||||
use std::time::Duration;
|
|
||||||
use std::str;
|
use std::str;
|
||||||
|
|
||||||
/// This is a basic example of leveraging `FramedUdp` to create
|
use futures::{Future, Stream, Sink};
|
||||||
/// a simple UDP client and server which speak a custom Protocol.
|
use tokio_core::net::{UdpSocket, UdpCodec};
|
||||||
/// `FramedUdp` applies a `Codec` to the input and output of an
|
use tokio_core::reactor::Core;
|
||||||
/// `Evented`
|
|
||||||
|
|
||||||
/// Simple Newline based parser,
|
pub struct LineCodec;
|
||||||
/// This is for a connectionless server, it must keep track
|
|
||||||
/// of the Socket address of the last peer to contact it
|
|
||||||
/// so that it can respond back.
|
|
||||||
/// In the real world, one would probably
|
|
||||||
/// want an associative of remote peers to their state
|
|
||||||
///
|
|
||||||
/// Note that this takes a pretty draconian stance by returning
|
|
||||||
/// an error if it can't find a newline in the datagram it received
|
|
||||||
pub struct LineCodec {
|
|
||||||
addr : Option<SocketAddr>
|
|
||||||
}
|
|
||||||
|
|
||||||
impl UdpCodec for LineCodec {
|
impl UdpCodec for LineCodec {
|
||||||
type In = Vec<Vec<u8>>;
|
type In = (SocketAddr, Vec<u8>);
|
||||||
type Out = Vec<u8>;
|
type Out = (SocketAddr, Vec<u8>);
|
||||||
|
|
||||||
fn decode(&mut self, addr : &SocketAddr, buf: &[u8]) -> Result<Self::In, io::Error> {
|
fn decode(&mut self, addr: &SocketAddr, buf: &[u8]) -> io::Result<Self::In> {
|
||||||
trace!("decoding {} - {}", str::from_utf8(buf).unwrap(), addr);
|
Ok((*addr, buf.to_vec()))
|
||||||
self.addr = Some(*addr);
|
|
||||||
let res : Vec<Vec<u8>> = buf.split(|c| *c == b'\n').map(|s| s.into()).collect();
|
|
||||||
if res.len() > 0 {
|
|
||||||
Ok(res)
|
|
||||||
}
|
|
||||||
else {
|
|
||||||
Err(io::Error::new(io::ErrorKind::Other,
|
|
||||||
"failed to find newline in datagram"))
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn encode(&mut self, item: Vec<u8>, into: &mut Vec<u8>) -> SocketAddr {
|
fn encode(&mut self,
|
||||||
trace!("encoding {}", str::from_utf8(item.as_slice()).unwrap());
|
(addr, buf): (SocketAddr, Vec<u8>),
|
||||||
into.extend_from_slice(item.as_slice());
|
into: &mut Vec<u8>) -> SocketAddr {
|
||||||
into.push('\n' as u8);
|
into.extend(buf);
|
||||||
|
return addr
|
||||||
self.addr.unwrap()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -63,56 +42,40 @@ fn main() {
|
|||||||
let mut core = Core::new().unwrap();
|
let mut core = Core::new().unwrap();
|
||||||
let handle = core.handle();
|
let handle = core.handle();
|
||||||
|
|
||||||
//create the line codec parser for each
|
let addr: SocketAddr = "127.0.0.1:0".parse().unwrap();
|
||||||
let srvcodec = LineCodec { addr : None };
|
|
||||||
let clicodec = LineCodec { addr : None };
|
|
||||||
|
|
||||||
let srvaddr : SocketAddr = "127.0.0.1:31999".parse().unwrap();
|
// Bind both our sockets and then figure out what ports we got.
|
||||||
let clientaddr : SocketAddr = "127.0.0.1:32000".parse().unwrap();
|
let a = UdpSocket::bind(&addr, &handle).unwrap();
|
||||||
|
let b = UdpSocket::bind(&addr, &handle).unwrap();
|
||||||
|
let b_addr = b.local_addr().unwrap();
|
||||||
|
|
||||||
//We bind each socket to a specific port
|
// We're parsing each socket with the `LineCodec` defined above, and then we
|
||||||
let server = UdpSocket::bind(&srvaddr, &handle).unwrap();
|
// `split` each codec into the sink/stream halves.
|
||||||
let client = UdpSocket::bind(&clientaddr, &handle).unwrap();
|
let (a_sink, a_stream) = a.framed(LineCodec).split();
|
||||||
|
let (b_sink, b_stream) = b.framed(LineCodec).split();
|
||||||
|
|
||||||
//start things off by sending a ping from the client to the server
|
// Start off by sending a ping from a to b, afterwards we just print out
|
||||||
//This doesn't utilize the codec to encode the message, but rather
|
// what they send us and continually send pings
|
||||||
//it sends raw data directly to the remote peer with the send_dgram future
|
// let pings = stream::iter((0..5).map(Ok));
|
||||||
let job = client.send_dgram(b"PING\n", srvaddr);
|
let a = a_sink.send((b_addr, b"PING".to_vec())).and_then(|a_sink| {
|
||||||
let (client, _buf) = core.run(job).unwrap();
|
let mut i = 0;
|
||||||
|
let a_stream = a_stream.take(4).map(move |(addr, msg)| {
|
||||||
|
i += 1;
|
||||||
|
println!("[a] recv: {}", String::from_utf8_lossy(&msg));
|
||||||
|
(addr, format!("PING {}", i).into_bytes())
|
||||||
|
});
|
||||||
|
a_sink.send_all(a_stream)
|
||||||
|
});
|
||||||
|
|
||||||
//We create a FramedUdp instance, which associates a socket
|
// The second client we have will receive the pings from `a` and then send
|
||||||
//with a codec. We then immediate split that into the
|
// back pongs.
|
||||||
//receiving side `Stream` and the writing side `Sink`
|
let b_stream = b_stream.map(|(addr, msg)| {
|
||||||
let (srvsink, srvstream) = server.framed(srvcodec).split();
|
println!("[b] recv: {}", String::from_utf8_lossy(&msg));
|
||||||
|
(addr, b"PONG".to_vec())
|
||||||
|
});
|
||||||
|
let b = b_sink.send_all(b_stream);
|
||||||
|
|
||||||
//`Stream::fold` runs once per every received datagram.
|
// Spawn the sender of pongs and then wait for our pinger to finish.
|
||||||
//Note that we pass srvsink into fold, so that it can be
|
handle.spawn(b.then(|_| Ok(())));
|
||||||
//supplied to every iteration. The reason for this is
|
drop(core.run(a));
|
||||||
//sink.send moves itself into `send` and then returns itself
|
|
||||||
let srvloop = srvstream.fold(srvsink, move |sink, lines| {
|
|
||||||
println!("{}", str::from_utf8(lines[0].as_slice()).unwrap());
|
|
||||||
sink.send(b"PONG".to_vec())
|
|
||||||
}).map(|_| ());
|
|
||||||
|
|
||||||
//We create another FramedUdp instance, this time for the client socket
|
|
||||||
let (clisink, clistream) = client.framed(clicodec).split();
|
|
||||||
|
|
||||||
//And another infinite iteration
|
|
||||||
let cliloop = clistream.fold(clisink, move |sink, lines| {
|
|
||||||
println!("{}", str::from_utf8(lines[0].as_slice()).unwrap());
|
|
||||||
sink.send(b"PING".to_vec())
|
|
||||||
}).map(|_| ());
|
|
||||||
|
|
||||||
let timeout = Timeout::new(Duration::from_millis(500), &handle).unwrap();
|
|
||||||
|
|
||||||
//`select_all` takes an `Iterable` of `Future` and returns a future itself
|
|
||||||
//This future waits until the first `Future` completes, it then returns
|
|
||||||
//that result.
|
|
||||||
let wait = future::select_all(vec![timeout.boxed(), srvloop.boxed(), cliloop.boxed()]);
|
|
||||||
|
|
||||||
//Now we instruct `reactor::Core` to iterate, processing events until its future, `SelectAll`
|
|
||||||
//has completed
|
|
||||||
if let Err(e) = core.run(wait) {
|
|
||||||
error!("{}", e.0);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user