You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Tokio::select循环中服务器无法检测到传入流的问题

问题根源分析

1. 变量遮蔽导致读取缓冲区未清空

在reader.read_line的处理分支中,你重新定义了let mut line = format!(...),这会遮蔽外层循环中用来接收输入的line变量。后续调用的line.clear()仅清除了这个新变量,外层的line会持续累积之前的输入内容,导致下一次read_line时新内容被追加到旧内容后,甚至因缓冲区有内容但未触发换行而持续阻塞。

2. read_line依赖换行符触发返回

read_line会一直阻塞,直到读取到换行符(\n)或者流关闭。如果客户端发送的消息没有以换行结尾,read_line永远不会返回,导致tokio::select一直卡在这个分支,无法处理其他事件,直到连接关闭。


修复后的代码

async fn main() {
    let listener = TcpListener::bind("localhost:3001").await.unwrap();
    let (tx, _rx) = broadcast::channel(10);
    loop {
        let (mut stream, socket_addr) = listener.accept().await.unwrap();
        let tx = tx.clone();
        let mut rx = tx.subscribe();
        tokio::spawn(async move {
            let mut name_buf = [0; 16];
            // 读取名称时只取实际读取的字节,避免空字节干扰
            let n = stream.read(&mut name_buf).await.unwrap();
            let name = String::from_utf8_lossy(&name_buf[..n]).trim().to_string();
            println!("{name} Connected.");

            let (reader, mut writer) = stream.split();
            let mut reader = BufReader::new(reader);
            let mut line = String::new();

            loop {
                tokio::select! {
                    result = reader.read_line(&mut line) => {
                        match result {
                            Ok(0) => {
                                // 读取到0字节,说明流已关闭
                                println!("Connection lost with {name}");
                                break;
                            }
                            Ok(_) => {
                                // 修剪换行符,处理空消息
                                let msg = line.trim_end();
                                if msg.is_empty() {
                                    line.clear();
                                    continue;
                                }
                                let broadcast_msg = format!("{}: {}\n", name, msg);
                                println!("{}", broadcast_msg.trim_end());
                                tx.send((broadcast_msg.clone(), socket_addr)).unwrap();
                                // 清空缓冲区,准备下一次读取
                                line.clear();
                            }
                            Err(_) => {
                                println!("Connection lost with {name}");
                                break;
                            }
                        }
                    }
                    result = rx.recv() => {
                        let (msg, other_addr) = result.unwrap();
                        if other_addr != socket_addr {
                            writer.write_all(msg.as_bytes()).await.unwrap();
                        }
                    }
                }
            }
        });
    }
}

关键修复点

  • 消除变量遮蔽:复用外层的line缓冲区,处理完成后调用line.clear()清空,确保每次读取都是干净状态。
  • 优化名称读取:读取名称时获取实际字节数n,仅转换有效字节为字符串,避免数组中空字节的干扰。
  • 明确read_line返回值处理:Ok(0)直接判定为流关闭;读取到内容后先修剪换行符,再处理消息。
  • 保证消息格式统一:广播消息末尾添加\n,确保客户端使用read_line时能正确识别消息边界。

内容的提问来源于stack exchange,提问作者lucianopinochet

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 01:37:36