基于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,需要补充格式判断逻辑。
修复步骤
修改服务器离开消息格式:
在服务器代码中,将离开消息添加换行符:// 原代码 line = format!("/left){} has left", &nickname_to_delete); // 修改后 line = format!("/left){} has left\n", &nickname_to_delete);客户端添加消息类型判断:
在客户端的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
相关产品推荐
相关产品推荐

