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

基于Tokio的TCP服务器无法实时推送用户离开消息求助

TCP聊天服务器用户离开通知延迟问题

我开发了一款可检测用户断开连接的TCP服务器,需求是当用户离开聊天时,向所有连接的TCP客户端推送该用户离开的消息。但当前功能异常,无法实时收到用户离开的通知,只有在其他用户发送新消息后,离开提示才会显示。

客户端代码

use {
    std::{io::BufRead, thread},
    tokio::{
        io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
        net::TcpStream,
        sync::mpsc,
    },
};

#[tokio::main]
async fn main() {
    let mut stream = TcpStream::connect("localhost:8080").await.unwrap();

    let (reader, mut writer) = stream.split();

    let mut reader = BufReader::new(reader);

    let mut nickname: String = String::new();
    let mut response = " ".to_string();

    while response.trim() != "200" {
        nickname.clear();
        if response != " " {
            println!("Already exist in a chat🤨");
        }

        println!("{}", "Enter your nickname: ");
        std::io::stdin().read_line(&mut nickname).unwrap();

        let request = format!("/register){}", nickname);

        writer.write_all(request.as_bytes()).await.unwrap();
        response.clear();
        reader.read_line(&mut response).await.unwrap();
    }

    let (tx, mut rx) = mpsc::channel(10);

    thread::spawn(move || {
        let mut stdin = std::io::stdin().lock();
        loop {
            let mut line = String::new();
            stdin.read_line(&mut line).unwrap();
            tx.blocking_send(line).unwrap();
        }
    });

    loop {
        let mut line = String::new();
        tokio::select! {
            result = reader.read_line(&mut line) => {
                result.unwrap();
                println!("{:?}", line);

                let mut split2 = line.splitn(2, ':');

                let author = split2.next().unwrap();
                let msg = split2.next().unwrap();

                print!("){}:{}", author, msg);
            }

            result = rx.recv() => {
                let line_write = result.unwrap();

                if line_write.trim().len() > 0 {
                    let output: String;

                    output = format!("){}: {}", nickname.trim(), line_write);

                    writer.write_all(output.as_bytes()).await.unwrap();
                } else {
                    println!("message is too short");
                }
            }
        }
    }
}

服务器代码

use tokio::{
    io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
    net::TcpListener,
    sync::{broadcast, Mutex},
};

use std::{
    collections::{HashMap, HashSet},
    sync::Arc,
};

#[tokio::main]
async fn main() {
    let tcp_listener = TcpListener::bind("localhost:8080").await.unwrap();
    let nicknames: Arc<Mutex<HashSet<String>>> = Arc::new(Mutex::new(HashSet::new()));
    let addr_nickname: Arc<Mutex<HashMap<String, String>>> = Arc::new(Mutex::new(HashMap::new()));

    let (tx, _rx) = broadcast::channel(10);

    loop {
        let (mut socket, addr) = tcp_listener.accept().await.unwrap();

        let tx: broadcast::Sender<(String, std::net::SocketAddr)> = tx.clone();
        let mut rx: broadcast::Receiver<(String, std::net::SocketAddr)> = tx.subscribe();
        let nicknames: Arc<Mutex<HashSet<String>>> = Arc::clone(&nicknames);
        let addr_nickname: Arc<Mutex<HashMap<String, String>>> = Arc::clone(&addr_nickname);

        tokio::spawn(async move {
            let (reader, mut writer) = socket.split();

            let mut reader = BufReader::new(reader);

            let mut line = String::new();

            loop {
                tokio::select! {
                    result = reader.read_line(&mut line) => {
                        match result {
                            Ok(_) => {
                                tx.send((line.clone(), addr)).unwrap();
                                line.clear();
                            },
                            Err(e) if e.kind() == std::io::ErrorKind::ConnectionReset => {
                                println!("Connection was reset by the client: {:?}, by {:?}", e, addr);
                                let mut nicknames = nicknames.lock().await;
                                let mut addr_nickname = addr_nickname.lock().await;
                                let nickname_to_delete = addr_nickname.get(&addr.to_string()).unwrap();
                                line = format!("/left){} has left", &nickname_to_delete);
                                nicknames.remove(nickname_to_delete);
                                addr_nickname.remove(&addr.to_string());

                                println!("{:?}, {:?}", nicknames, addr_nickname);

                                tx.send((line.clone(), addr)).unwrap();

                                line.clear();

                                break;
                            },
                            Err(e) => {
                                println!("An unexpected error occurred: {:?}", e);
                            },
                        }
                    }
                    result = rx.recv() => {
                        let (msg, other_addr) = result.unwrap();

                        println!("{:?}, {:?}", msg, other_addr);

                        let mut split = msg.splitn(2, ')');

                        let command = split.next().unwrap();
                        let nickname = split.next().unwrap();

                        if command == "/register" && addr == other_addr {
                            let mut nicknames = nicknames.lock().await;
                            let mut addr_nickname = addr_nickname.lock().await;
                            println!("{:?}, {}, {:?}", nicknames, nickname, addr_nickname);
                            if nicknames.contains(nickname.trim()) {
                                writer.write_all(b"406\n").await.unwrap();
                            } else {
                                nicknames.insert(nickname.trim().to_string());
                                addr_nickname.insert(addr.to_string(), nickname.trim().to_string());
                                writer.write_all(b"200\n").await.unwrap();
                                println!("{:?}", nicknames);
                            }
                        } else if command == "/users" && addr == other_addr {
                            let addr_nickname = addr_nickname.lock().await;
                            if addr_nickname.is_empty() {
                                writer.write_all(b"No users found\n").await.unwrap();
                            } else {
                                let users: String = addr_nickname.values().cloned().collect::<Vec<_>>().join(", ");
                                match writer.write_all(users.as_bytes()).await {
                                    Ok(_) => println!("Response sent successfully"),
                                    Err(e) => println!("Failed to send response: {}", e),
                                }
                            }
                        } else if addr != other_addr && command != "/register" && command != "/users" {
                            println!("{:?}, {:?}, {}, {:?}, {:?}", command, nickname, command=="/register", addr, other_addr);
                            writer.write_all(msg.as_bytes()).await.unwrap();
                        }
                    }
                }
            }
        });
    }
}

问题原因及修复方案

核心问题

服务器发送用户离开消息时未添加换行符,而客户端使用read_line方法读取数据——该方法会阻塞直到读取到换行符才返回。因此离开消息会被客户端缓存,直到后续有带换行符的消息发送时才会被一并读取处理,导致延迟显示。

此外,客户端未处理/left类型的消息格式,直接按:分割会触发panic,需要补充格式判断逻辑。

修复步骤

  1. 修改服务器离开消息格式:
    在服务器代码中,将离开消息添加换行符:

    // 原代码
    line = format!("/left){} has left", &nickname_to_delete);
    // 修改后
    line = format!("/left){} has left\n", &nickname_to_delete);
    
  2. 客户端添加消息类型判断:
    在客户端的reader.read_line分支中,先解析消息类型,再分别处理:

    result = reader.read_line(&mut line) => {
        result.unwrap();
        let line_content = line.trim_end(); // 去除末尾换行符
        if line_content.starts_with("/left)") {
            // 处理用户离开消息
            if let Some(msg) = line_content.splitn(2, ')').nth(1) {
                println!("{}", msg);
            }
        } else {
            // 处理普通聊天消息
            let mut split2 = line_content.splitn(2, ':');
            if let (Some(author), Some(msg)) = (split2.next(), split2.next()) {
                print!("){}:{}", author, msg);
            } else {
                println!("{}", line_content);
            }
        }
        line.clear();
    }
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 22:25:53