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

如何解决TCP连接聊天中的消息重叠问题

TCP聊天程序输入与接收消息重叠问题

我开发了一个TCP客户端和TCP服务器,遇到以下问题:当用户正在输入消息时,若其他用户发送消息,两条消息会发生重叠。我正在寻找解决该问题的方法,可借助外部库实现。

重叠示例

user1 正在输入消息:
user1 writing message
user2 发送消息:
user2 sends message
user1 接收消息:
user1 receiving message
user2 接收user1发送的消息:
user2 receives message that user1 has sent

客户端代码(client.rs)

use std::io::{stdout, Write};

mod message;
mod user;

use {
    crate::{message::Message, user::User},
    chrono::Local,
    colored::Colorize,
    crossterm::{cursor::MoveToPreviousLine, execute, QueueableCommand},
    rusty_audio::Audio,
    std::{io::BufRead, thread},
    tokio::{
        io::{AsyncReadExt, AsyncWriteExt, BufReader},
        net::TcpStream,
        sync::mpsc,
    },
};

#[tokio::main]
async fn main() {
    // let mut adr_url = String::new();
    // println!("{}", "Enter chat address:".green());
    // std::io::stdin().read_line(&mut adr_url).unwrap();
    // let mut stream = TcpStream::connect(adr_url.trim()).await.unwrap();
    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 received_message: Message;

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

        let user = User::new(nickname.trim(), (0, 0, 0));
        let message = Message::new("/register", user, 1, "");
        let output = bincode::serialize(&message).unwrap();

        writer.write_all(&output).await.unwrap();

        let mut buffer = vec![0; 1024];

        reader.read(&mut buffer).await.unwrap();

        received_message = bincode::deserialize(&buffer).unwrap();

        if received_message.status == 1 {
            break;
        }
        nickname.clear();
    }

    let me = User::new(&received_message.user.nickname, received_message.user.color);

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

    thread::spawn(move || {
        let mut stdin = std::io::stdin().lock();
        loop {
            let mut stdout = stdout();
            let mut line = String::new();
            stdin.read_line(&mut line).unwrap();
            execute!(stdout, MoveToPreviousLine(1)).unwrap();
            if line.trim().len() < 206 && line.trim().len() > 0 {
                print!(
                    "{}: {} - {} ✓\n",
                    "You".red(),
                    line.trim(),
                    Local::now().format("%H:%M").to_string().blue()
                );
                std::io::stdout().flush().unwrap();
                tx.blocking_send(line).unwrap();
            } else {
                stdout.queue(crossterm::style::Print("\x1b[2K\r")).unwrap();
                stdout.flush().unwrap();
                println!("too long or too short");
            }
        }
    });

    let mut audio = Audio::new();
    audio.add("message", "message.wav");
    audio.add("leave", "leave.wav");

    loop {
        let mut buffer = vec![0; 1024];
        tokio::select! {
            result = reader.read(&mut buffer) => {
                result.unwrap();
                let message: Message = bincode::deserialize(&buffer).unwrap();
                let (command, author, color, msg) = (message.command, message.user.nickname, message.user.color, message.message);


                print!("{}: {} - {}\n", author.truecolor(color.0, color.1, color.2), msg, Local::now().format("%H:%M").to_string().blue());

                if command == "/left" {
                    audio.play("leave");
                } else {
                    audio.play("message");
                }
            }

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

                let message: Message;
                if line_write.chars().nth(0).unwrap() == '/' {
                    message = Message::new(line_write.trim(), me.clone(), 1, "");
                } else {
                    message = Message::new("", me.clone(), 1, line_write.trim());
                }
                let output = bincode::serialize(&message).unwrap();
                writer.write_all(&output).await.unwrap();
            }
        }
    }
}

服务端代码(server.rs)

mod message;
mod user;

use {
    crate::message::Message,
    rand::prelude::*,
    std::{collections::HashMap, sync::Arc},
    tokio::{
        io::{AsyncReadExt, AsyncWriteExt, BufReader},
        net::TcpListener,
        sync::{broadcast, Mutex},
    },
    user::User,
};

#[tokio::main]
async fn main() {
    let tcp_listener = TcpListener::bind("localhost:8080").await.unwrap();
    let clients: Arc<Mutex<HashMap<String, User>>> = 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 clients: Arc<Mutex<HashMap<String, User>>> = Arc::clone(&clients);

        tokio::spawn(async move {
            let (reader, mut writer) = socket.split();
            let mut reader = BufReader::new(reader);
            let mut buffer = vec![0; 1024];

            loop {
                tokio::select! {
                    result = reader.read(&mut buffer) => {
                        match result {
                            Ok(_) => {
                                let received_message: Message = bincode::deserialize(&buffer).unwrap();
                                let json_string = serde_json::to_string(&received_message).unwrap();
                                tx.send((json_string, addr)).unwrap();
                            },
                            Err(e) if e.kind() == std::io::ErrorKind::ConnectionReset => {
                                println!("Connection was reset by {:?}", addr);
                                let mut clients = clients.lock().await;
                                println!("{:?}", clients);
                                match clients.get(&addr.to_string()) {
                                    Some(user_left) => {
                                        println!("{:?}", clients);
                                        let server = User::new("server", (100, 100, 100));
                                        let message = Message::new("/left", server, 1, &format!("{} has left", user_left.nickname));
                                        clients.remove(&addr.to_string());

                                        let json_string = serde_json::to_string(&message).unwrap();

                                        tx.send((json_string, addr)).unwrap();
                                    },
                                    None => {
                                        break;
                                    }
                                }
                            },
                            Err(e) => {
                                println!("An unexpected error occurred: {:?}", e);
                            },
                        }
                    }

                    result = rx.recv() => {
                        let (msg, other_addr) = result.unwrap();
                        let message: Message = serde_json::from_str(&msg).unwrap();
                        let vec = bincode::serialize(&message).unwrap();
                        let (command, user) = (message.command, message.user);

                        println!("{:?}, {:?}", command, user);

                        if command.len() > 0 && addr == other_addr{
                            let mut clients = clients.lock().await;
                            let message: Message;
                            let mut output: Vec<u8> = Vec::new();
                            if command == "/register" {
                                if clients.values().any(|val| val.nickname == user.nickname) || user.nickname == "server" || user.nickname == "You" || user.nickname == "you" {
                                    let user = User::new("", (0, 0 ,0));
                                    message = Message::new("", user, 0, "");
                                    output = bincode::serialize(&message).unwrap();

                                } else {
                                    let mut rng = rand::thread_rng();
                                    let color = (
                                        rng.gen_range(0..=255),
                                        rng.gen_range(0..=255),
                                        rng.gen_range(0..=255),
                                    );

                                    let user = User::new(&user.nickname, color);

                                    clients.insert(addr.to_string(), user.clone());

                                    message = Message::new("", user, 1, "");
                                    output = bincode::serialize(&message).unwrap();

                                }
                            } else if command == "/users" {
                                let server = User::new("server", (100, 100, 100));
                                let nicknames: String = clients.iter()
                                    .filter_map(|(_, user_iter)| if user_iter.nickname != user.nickname { Some(user_iter.nickname.clone()) } else { None })
                                    .collect::<Vec<String>>()
                                    .join(", ");
                                println!("{:?}", nicknames);
                                if nicknames.len() == 0 {
                                    message = Message::new("", server.clone(), 0, "server is empty");
                                    output = bincode::serialize(&message).unwrap();
                                } else {
                                    message = Message::new("", server.clone(), 0, &nicknames);
                                    output = bincode::serialize(&message).unwrap();
                                }
                            }
                            writer.write_all(&output).await.unwrap();
                        } else if addr != other_addr && (command.len() == 0 || command=="/left") {
                            writer.write_all(&vec).await.unwrap();
                        }
                    }
                }
            }
        });
    }
}

user.rs代码

use serde::{Deserialize, Serialize};

#[derive(Serialize, Deserialize, Debug, Clone)]
pub struct User {
    pub nickname: String,
    pub color: (u8, u8, u8),
}

impl User {
    pub fn new(nickname: &str, color: (u8, u8, u8)) -> Self {
        Self {
            nickname: nickname.to_string(),
            color,
        }
    }
}

message.rs代码

use serde::{Deserialize, Serialize};

use crate::user::User;

#[derive(Serialize, Deserialize, Debug)]
pub struct Message {
    pub command: String,
    pub user: User,
    pub status: u8,
    pub message: String,
}

impl Message {
    pub fn new(command: &str, user: User, status: u8, message: &str) -> Self {
        Self {
            command: command.to_string(),
            user,
            status,
            message: message.to_string(),
        }
    }
}

我曾尝试用crossterm保存光标位置后删除消息,但未成功,可能是使用方式有误。我的期望是当其他用户发送消息时,不会与当前用户正在输入的消息重叠,可通过清除用户输入、打印收到的消息后恢复输入的方式实现。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 16:18:13