3.4 发送消息
原文链接: https://book.async.rs/tutorial/sending_messages.html
发送消息
现在是时候实现另一半——发送消息了。
实现发送最明显的方法是让每个 connection_loop 访问其他每个客户端 TcpStream 的写半部分。
这样,客户端可以直接向接收者 .write_all 消息。
然而,这是错误的:如果 Alice 发送 bob: foo,而 Charley 发送 bob: bar,Bob 实际上可能收到 fobaor。
通过套接字发送消息可能需要多次系统调用,因此两个并发的 .write_all 可能会相互干扰!
经验法则是,每个 TcpStream 只应由单个任务写入。
所以让我们创建一个 connection_writer_loop 任务,它通过 channel 接收消息并将它们写入套接字。
该任务将是消息序列化的节点。
如果 Alice 和 Charley 同时向 Bob 发送两条消息,Bob 将看到与它们到达 channel 时相同顺序的消息。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
| # extern crate futures;
# use async_std::{
# net::TcpStream,
# prelude::*,
# };
use futures::channel::mpsc; // 1
use futures::sink::SinkExt;
use std::sync::Arc;
# type Result<T> = std::result::Result<T, Box<dyn std::error::Error + Send + Sync>>;
type Sender<T> = mpsc::UnboundedSender<T>; // 2
type Receiver<T> = mpsc::UnboundedReceiver<T>;
async fn connection_writer_loop(
mut messages: Receiver<String>,
stream: Arc<TcpStream>, // 3
) -> Result<()> {
let mut stream = &*stream;
while let Some(msg) = messages.next().await {
stream.write_all(msg.as_bytes()).await?;
}
Ok(())
}
|
- 我们将使用
futures crate 中的 channel。 - 为简单起见,我们将使用无界 channel,本教程不会讨论背压。
- 由于
connection_loop 和 connection_writer_loop 共享同一个 TcpStream,我们需要将其放入 Arc。
注意,因为 client 只从流中读取,而 connection_writer_loop 只向流中写入,我们这里不会产生竞争。