如何实现WebSocket连接全局化,接收多客户端消息并查看日志
问题
我搭建了一个WebSocket服务,客户端发送消息后服务器会将消息回传。使用两个浏览器测试时,两个客户端都能收到各自的消息,但服务器日志中无法看到两条消息。请问如何实现WebSocket连接全局化以查看所有客户端消息?是否需要创建另一个WebSocket来接收并转发这些消息?
前端代码
import { isPlatformBrowser } from '@angular/common'; import { Inject, Injectable } from '@angular/core'; import { PLATFORM_ID } from '@angular/core'; @Injectable({ providedIn: 'root', }) export class WebsocketService { messages: Array<any> = []; chatWebSocket: WebSocket | undefined = undefined; createSocket() { if (isPlatformBrowser(this.platformId)) { this.chatWebSocket = new WebSocket('ws://localhost:3000/ws'); } } send(msg: string) { console.log(this.chatWebSocket); (this.chatWebSocket as WebSocket).send(msg); } recv() { if (!this.chatWebSocket) return; this.chatWebSocket.addEventListener('message', (e) => { this.messages.push(e.data); }); } constructor(@Inject(PLATFORM_ID) public platformId: Object) {} }
服务端代码
pub async fn chat_ws( ws: WebSocketUpgrade, user_agent: Option<TypedHeader<UserAgent>>, ) -> impl IntoResponse { if let Some(TypedHeader(user_agent)) = user_agent { println!("`{}` connected", user_agent.as_str()); } ws.on_upgrade(handle_chat_socket) } pub async fn handle_chat_socket(mut socket: WebSocket) { loop { if let Some(msg) = socket.recv().await { if let Ok(msg) = msg { match msg { Message::Text(a) => { socket.send(Message::Text(String::from(a.clone()))).await; }, _ => println!("Other"), } } } } }
解决方案
不需要额外创建WebSocket服务,你只需要在服务端维护一个全局的客户端连接集合,就能追踪所有活跃连接并查看所有客户端消息,具体实现步骤如下:
用线程安全容器管理连接集合
在Rust中,使用Arc<Mutex<Vec<WebSocket>>>来存储所有活跃的WebSocket连接,确保多异步任务下的安全访问。修改服务端逻辑
- 新客户端连接时,将其加入全局集合;
- 收到消息时先打印日志,这样就能看到所有客户端的消息;
- 可选:将消息广播给所有客户端(如果需要多客户端通信);
- 客户端断开连接时,从集合中移除对应的连接,避免内存泄漏。
修改后的服务端示例代码
use std::sync::{Arc, Mutex}; use actix_web::{web, IntoResponse}; use actix_web_actors::ws; use actix_web::http::header::TypedHeader; use actix_web::http::header::UserAgent; // 定义全局连接存储类型 type WsConnections = Arc<Mutex<Vec<ws::WebSocket>>>; // 服务配置:初始化连接集合并注入 pub fn configure_services(cfg: &mut web::ServiceConfig) { let connections = Arc::new(Mutex::new(Vec::new())); cfg.app_data(connections.clone()) .route("/ws", web::get().to(chat_ws)); } pub async fn chat_ws( ws: ws::WebSocketUpgrade, user_agent: Option<TypedHeader<UserAgent>>, connections: web::Data<WsConnections>, ) -> impl IntoResponse { if let Some(TypedHeader(user_agent)) = user_agent { println!("`{}` connected", user_agent.as_str()); } // 将连接集合传递给处理函数 ws.on_upgrade(|socket| handle_chat_socket(socket, connections)) } pub async fn handle_chat_socket(mut socket: ws::WebSocket, connections: web::Data<WsConnections>) { // 新连接加入全局集合 { let mut conns = connections.lock().unwrap(); conns.push(socket.clone()); } loop { match socket.recv().await { Some(Ok(msg)) => { match msg { ws::Message::Text(text) => { // 打印所有客户端消息到服务器日志 println!("收到客户端消息: {}", text); // 回传给发送者(原有逻辑) if let Err(e) = socket.send(ws::Message::Text(text.clone())).await { eprintln!("回传消息失败: {}", e); } // 可选:广播消息给所有其他客户端 let mut conns = connections.lock().unwrap(); for conn in conns.iter_mut() { if *conn != socket { if let Err(e) = conn.send(ws::Message::Text(text.clone())).await { eprintln!("广播消息失败: {}", e); } } } }, ws::Message::Binary(bin) => { println!("收到二进制消息,长度: {}", bin.len()); if let Err(e) = socket.send(ws::Message::Binary(bin)).await { eprintln!("回传二进制消息失败: {}", e); } }, _ => println!("收到其他类型消息"), } }, Some(Err(e)) => { eprintln!("连接错误: {}", e); break; }, None => { println!("客户端断开连接"); break; } } } // 连接断开后从全局集合移除 let mut conns = connections.lock().unwrap(); conns.retain(|conn| *conn != socket); }
注意事项
- 必须使用线程安全的容器(
Arc<Mutex>),因为每个WebSocket连接运行在独立的异步任务中,多任务访问集合需要同步; - 处理消息发送时的错误,避免单个连接的异常影响整个服务;
- 前端代码无需修改,保持现有连接逻辑即可。
内容的提问来源于stack exchange,提问作者Haider ali
相关产品推荐
相关产品推荐

