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

tokio::mpsc::unbounded_channel消息丢失:多客户端Websocket服务器异常

问题分析与解决方案

核心问题分析

  • 阻塞式睡眠破坏Tokio调度:广播线程使用std::thread::sleep,会直接占用Tokio工作线程并阻塞调度,导致客户端消息接收、广播转发等任务无法获得执行机会,多客户端场景下消息传递异常。
  • Mutex锁持有时间过长:send_clients函数在遍历客户端、序列化数据包的全流程中一直持有客户端列表的Mutex锁,导致客户端连接、断开等操作被长时间阻塞,影响消息转发效率。
  • 代码逻辑错误:客户端消息处理中,handle_message调用时传入的&mut clients未在当前作用域定义,属于编译级错误。
  • 冗余可变状态:State中的clients已通过Arc<Mutex>实现共享可变,广播线程和send_clients无需持有可变引用,多余的mut会引发生命周期问题。

修正后的代码

use futures_util::{future, StreamExt, TryStreamExt, SinkExt};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::{
    collections::HashMap,
    net::SocketAddr,
    sync::{Arc, Mutex},
    time::Instant,
};
use tokio::{net::TcpListener, sync::mpsc::{UnboundedSender, self}, time::sleep};
use tungstenite::Message;
use uuid::Uuid;

mod packets;
use packets::*;

static REFRESH_TIME: u64 = 2000;

#[derive(Debug, Clone)]
pub struct Client {
    pub send: UnboundedSender<Message>
}

#[derive(Debug, Default, Clone)]
pub struct State {
    pub clients: Arc<Mutex<HashMap<Uuid, Client>>>
}

#[tokio::main]
async fn main() {
    let g_state = State::default();

    let state_clone = g_state.clone();
    tokio::task::spawn(async move {
        let state = state_clone;
        loop {
            send_clients(&state).await;
            sleep(std::time::Duration::from_millis(REFRESH_TIME)).await;
        }
    });

    let stream = TcpListener::bind("0.0.0.0:25342").await.unwrap();
    loop {
        let (socket, addr) = stream.accept().await.unwrap();
        let state_clone = g_state.clone();
        tokio::task::spawn(async move {
            let self_uuid = Uuid::new_v4();
            let state = state_clone;
            let (tx, mut rx) = mpsc::unbounded_channel::<Message>();

            {
                let mut clients = state.clients.lock().unwrap();
                clients.insert(
                    self_uuid,
                    Client {
                        send: tx
                    },
                );
            }

            let websocket = tokio_tungstenite::accept_async(socket).await.unwrap();
            let (mut outgoing, mut incoming) = websocket.split();

            let incoming_handle = incoming.try_for_each(|msg| {
                handle_message(msg, &state, self_uuid);
                future::ok(())
            });
            
            let forward_handle = tokio::task::spawn(async move {
                while let Some(message) = rx.recv().await {
                    if let Err(e) = outgoing.send(message).await {
                        eprintln!("Failed to send message to client {}: {}", self_uuid, e);
                        break;
                    }
                }
            });

            tokio::select! {
                res = incoming_handle => {
                    if let Err(e) = res {
                        eprintln!("Incoming message error for client {}: {}", self_uuid, e);
                    }
                }
                _ = forward_handle => {}
            }

            {
                let mut clients = state.clients.lock().unwrap();
                println!("[x] Client disconnected! @ {}", self_uuid);
                clients.remove(&self_uuid);
            }
        });
    }
}

async fn send_clients(state: &State) {
    let senders: Vec<UnboundedSender<Message>> = {
        let clients = state.clients.lock().unwrap();
        clients.values().map(|client| client.send.clone()).collect()
    };

    let clients_snapshot = {
        let clients = state.clients.lock().unwrap();
        get_clients_packet_from_clients(&clients)
    };
    let packet = serde_json::to_string(&clients_snapshot).unwrap();
    let message = Message::Text(packet);

    for sender in senders {
        if let Err(e) = sender.send(message.clone()) {
            eprintln!("Failed to send broadcast message: {}", e);
        }
    }
}

fn handle_message(msg: Message, state: &State, client_id: Uuid) {
    match msg {
        Message::Text(text) => println!("Received message from client {}: {}", client_id, text),
        Message::Binary(_) => println!("Received binary message from client {}", client_id),
        _ => {}
    }
}

关键修改说明

  1. 替换为Tokio异步睡眠:用tokio::time::sleep替代std::thread::sleep,让Tokio在睡眠期间调度其他任务,避免阻塞工作线程。
  2. 缩短锁持有时间:先复制客户端Sender列表、生成数据包快照,立即释放锁后再逐个发送消息,减少锁竞争对其他操作的影响。
  3. 修复未定义变量错误:修正handle_message的调用参数,通过state访问客户端列表。
  4. 优化任务生命周期:用tokio::select!管理消息处理和转发任务,确保客户端断开时正确清理资源;处理转发失败的情况,避免无限阻塞。
  5. 移除冗余可变引用:删除State和send_clients中不必要的mut修饰符,符合Rust所有权规则。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 07:14:56