3.5 连接读写端

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

连接读取端与写入端

那么,我们如何确保在 connection_loop 中读取的消息流入相应的 connection_writer_loop? 我们应该以某种方式维护一个 peers: HashMap<String, Sender<String>> 映射,使客户端能够找到目标 channel。 然而,这个映射将是一些共享可变状态,因此我们必须在其上包装 RwLock,并回答棘手的问题:如果客户端在接收消息的同时加入,应该发生什么。

简化状态推理的一个技巧来自 actor 模型。 我们可以创建一个专用的 broker 任务,它拥有 peers 映射,并使用 channel 与其他任务通信。 通过将 peers 隐藏在这样的「actor」任务内部,我们消除了对互斥锁的需求,也使序列化点变得明确。 「Bob 向 Alice 发送消息」和「Alice 加入」这两个事件的顺序,由 broker 事件队列中相应事件的顺序决定。

 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
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
# extern crate futures;
# use async_std::{
#     net::TcpStream,
#     prelude::*,
#     task,
# };
# use futures::channel::mpsc;
# 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>;
# type Receiver<T> = mpsc::UnboundedReceiver<T>;
#
# async fn connection_writer_loop(
#     mut messages: Receiver<String>,
#     stream: Arc<TcpStream>,
# ) -> Result<()> {
#     let mut stream = &*stream;
#     while let Some(msg) = messages.next().await {
#         stream.write_all(msg.as_bytes()).await?;
#     }
#     Ok(())
# }
#
# 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)
#         }
#     })
# }
#
use std::collections::hash_map::{Entry, HashMap};

#[derive(Debug)]
enum Event { // 1
    NewPeer {
        name: String,
        stream: Arc<TcpStream>,
    },
    Message {
        from: String,
        to: Vec<String>,
        msg: String,
    },
}

async fn broker_loop(mut events: Receiver<Event>) -> Result<()> {
    let mut peers: HashMap<String, Sender<String>> = HashMap::new(); // 2

    while let Some(event) = events.next().await {
        match event {
            Event::Message { from, to, msg } => {  // 3
                for addr in to {
                    if let Some(peer) = peers.get_mut(&addr) {
                        let msg = format!("from {}: {}\n", from, msg);
                        peer.send(msg).await?
                    }
                }
            }
            Event::NewPeer { name, stream } => {
                match peers.entry(name) {
                    Entry::Occupied(..) => (),
                    Entry::Vacant(entry) => {
                        let (client_sender, client_receiver) = mpsc::unbounded();
                        entry.insert(client_sender); // 4
                        spawn_and_log_error(connection_writer_loop(client_receiver, stream)); // 5
                    }
                }
            }
        }
    }
    Ok(())
}
  1. broker 任务应处理两种类型的事件:消息或新对等方到达。
  2. broker 的内部状态是一个 HashMap。 注意,我们这里不需要 Mutex,并且可以在 broker 循环的每次迭代中自信地说出当前的对等方集合
  3. 要处理消息,我们通过 channel 将其发送给每个目标
  4. 要处理新对等方,我们首先在 peer 映射中注册它……
  5. ……然后生成专用任务,将消息实际写入套接字。
最后修改 August 23, 2026: 更新 (499855b16)