3.3 接收消息

原文链接: https://book.async.rs/tutorial/receiving_messages.html

接收消息

让我们实现协议的接收部分。 我们需要:

  1. 将传入的 TcpStream 按 \n 分割,并将字节解码为 utf-8
  2. 将第一行解释为登录名
  3. 将其余行解析为 login: message
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
# use async_std::{
#     net::{TcpListener, ToSocketAddrs},
#     prelude::*,
#     task,
# };
#
# type Result<T> = std::result::Result<T, Box<dyn std::error::Error + Send + Sync>>;
#
use async_std::{
    io::BufReader,
    net::TcpStream,
};

async fn accept_loop(addr: impl ToSocketAddrs) -> Result<()> {
    let listener = TcpListener::bind(addr).await?;
    let mut incoming = listener.incoming();
    while let Some(stream) = incoming.next().await {
        let stream = stream?;
        println!("Accepting from: {}", stream.peer_addr()?);
        let _handle = task::spawn(connection_loop(stream)); // 1
    }
    Ok(())
}

async fn connection_loop(stream: TcpStream) -> Result<()> {
    let reader = BufReader::new(&stream); // 2
    let mut lines = reader.lines();

    let name = match lines.next().await { // 3
        None => Err("peer disconnected immediately")?,
        Some(line) => line?,
    };
    println!("name = {}", name);

    while let Some(line) = lines.next().await { // 4
        let line = line?;
        let (dest, msg) = match line.find(':') { // 5
            None => continue,
            Some(idx) => (&line[..idx], line[idx + 1 ..].trim()),
        };
        let dest: Vec<String> = dest.split(',').map(|name| name.trim().to_string()).collect();
        let msg: String = msg.to_string();
    }
    Ok(())
}
  1. 我们使用 task::spawn 函数为每个客户端的工作生成独立任务。 也就是说,在接受客户端后,accept_loop 立即开始等待下一个。 这是事件驱动架构的核心优势:我们并发服务许多客户端,而无需消耗大量硬件线程。

  2. 幸运的是,「将字节流分割成行」的功能已经实现。 .lines() 调用返回 String 的流。

  3. 我们获取第一行——登录名

  4. 再一次,我们实现手动的 async for 循环。

  5. 最后,我们将每行解析为目标登录名列表和消息本身。

管理错误

上述解决方案中的一个严重问题是,虽然我们在 connection_loop 中正确地传播了错误,但之后我们只是把错误丢在一边! 也就是说,task::spawn 不会立即返回错误(它不能,需要先运行 future 直到完成),只有在 join 之后才会。 我们可以像这样等待任务 join 来「修复」它:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
# #![feature(async_closure)]
# extern crate async_std;
# use async_std::{
#     io::BufReader,
#     net::{TcpListener, TcpStream, ToSocketAddrs},
#     prelude::*,
#     task,
# };
#
# type Result<T> = std::result::Result<T, Box<dyn std::error::Error + Send + Sync>>;
#
# async fn connection_loop(stream: TcpStream) -> Result<()> {
#     let reader = BufReader::new(&stream); // 2
#     let mut lines = reader.lines();
#
#     let name = match lines.next().await { // 3
#         None => Err("peer disconnected immediately")?,
#         Some(line) => line?,
#     };
#     println!("name = {}", name);
#
#     while let Some(line) = lines.next().await { // 4
#         let line = line?;
#         let (dest, msg) = match line.find(':') { // 5
#             None => continue,
#             Some(idx) => (&line[..idx], line[idx + 1 ..].trim()),
#         };
#         let dest: Vec<String> = dest.split(',').map(|name| name.trim().to_string()).collect();
#         let msg: String = msg.trim().to_string();
#     }
#     Ok(())
# }
#
# async move |stream| {
let handle = task::spawn(connection_loop(stream));
handle.await?
# };

.await 等待客户端完成,? 传播结果。

然而,这个解决方案有两个问题! 首先,因为我们立即 await 客户端,我们一次只能处理一个客户端,这完全违背了 async 的目的! 其次,如果客户端遇到 IO 错误,整个服务器会立即退出。 也就是说,一个对等方不稳定的网络连接会导致整个聊天室崩溃!

在这种情况下,正确处理客户端错误的方法是记录它们,并继续服务其他客户端。 所以让我们为此使用一个辅助函数:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
# extern crate async_std;
# use async_std::{
#     io,
#     prelude::*,
#     task,
# };
fn spawn_and_log_error<F>(fut: F) -> task::JoinHandle<()>
where
    F: Future<Output = Result<()>> + Send + 'static,
{
    task::spawn(async move {
        if let Err(e) = fut.await {
            eprintln!("{}", e)
        }
    })
}
最后修改 August 23, 2026: 更新 (499855b16)