如何主动关闭同步Tungstenite WebSocket连接?解决线程锁阻塞问题
解决WebSocket多线程互斥锁阻塞问题
问题描述
编写WebSocket客户端时,将连接放入子线程进行消息读取,主线程尝试关闭连接时,arc_ws_stream.lock()会永久阻塞——因为子线程中的互斥锁被持续占用,无法释放。尝试过异步方案,但无法解决无效SSL证书问题。
依赖配置(Cargo.toml)
[dependencies] tungstenite = "0.21.0" url = "2.5.0" native-tls = "0.2.11" websocket = "0.27.0"
原始代码(main.rs)
#![allow(unused_imports)] use std::{net::ToSocketAddrs, error, fmt::format, net, sync::{Arc, Mutex}, thread::{self, sleep}, time::Duration}; use tungstenite::{protocol::CloseFrame, WebSocket, Message, Error}; use url::Url; use native_tls::TlsConnector; use websocket::{result::WebSocketError, ClientBuilder, OwnedMessage}; fn main() { connect_to_websocket("wss://127.0.0.1:7878/ws", String::from("login message")); } fn connect_to_websocket(uri: &str, login_msg: String){ let ws_url = Url::parse(uri).unwrap(); // Create a TLS connector that ignores invalid certificates let tls_connector = TlsConnector::builder().danger_accept_invalid_certs(true).build().unwrap(); // Establish a TCP connection, then wrap the TCP stream with TLS and connect to the server let remote_addr = format!("{}:{}", ws_url.host().unwrap(), ws_url.port().unwrap()); let tcp_stream = std::net::TcpStream::connect(remote_addr.clone()).unwrap(); let tls_stream = tls_connector.connect(remote_addr.as_str(), tcp_stream).unwrap(); let (mut raw_ws_stream, _) = tungstenite::client(uri, tls_stream).unwrap(); raw_ws_stream.send(Message::Text(login_msg)).unwrap(); let arc_ws_stream = Arc::new(Mutex::new(raw_ws_stream)); let thread_ws_stream = arc_ws_stream.clone(); let ws_recv_thread = thread::spawn(move ||{ loop{ let msg = thread_ws_stream.lock().unwrap().read(); if !process_ws_msg(msg) { break; } } }); sleep(Duration::from_millis(5_000)); let _ = arc_ws_stream.lock().unwrap().close(None); //this line will be blocked forever let _ = ws_recv_thread.join(); } fn process_ws_msg(msg: Result<Message, Error>)->bool { match msg{ Ok(msg) => { match msg { Message::Text(msg) => { println!("WS received message: {}", msg); }, Message::Close(cls_msg) => { match cls_msg{ Some(cls_frame) =>{ println!("WS session closed message, close frame = {}", cls_frame); } None =>{ println!("WS session closed message, no close frame"); } } return false; }, _ => {} } }, Err(err) =>{ println!("WS session error: {}", err); return false; } } true }
问题原因
子线程中thread_ws_stream.lock().unwrap().read()的写法会让互斥锁在整个read()阻塞期间持续被持有。由于tungstenite的read()是阻塞式调用,只要没有消息返回,互斥锁就不会释放,导致主线程的arc_ws_stream.lock()永远无法获取锁,陷入永久阻塞。
解决方案
引入原子布尔值作为终止信号,让主线程通知子线程主动退出并释放互斥锁,同时缩小互斥锁的持有范围,避免阻塞期间占用锁。
修改后的代码(main.rs)
#![allow(unused_imports)] use std::{ net, sync::{Arc, Mutex, atomic::{AtomicBool, Ordering}}, thread::{self, sleep}, time::Duration, }; use tungstenite::{protocol::CloseFrame, WebSocket, Message, Error}; use url::Url; use native_tls::TlsConnector; fn main() { connect_to_websocket("wss://127.0.0.1:7878/ws", String::from("login message")); } fn connect_to_websocket(uri: &str, login_msg: String) { let ws_url = Url::parse(uri).unwrap(); // 创建忽略无效证书的TLS连接器 let tls_connector = TlsConnector::builder() .danger_accept_invalid_certs(true) .build() .unwrap(); // 建立TCP连接并包装TLS let remote_addr = format!("{}:{}", ws_url.host().unwrap(), ws_url.port().unwrap()); let tcp_stream = std::net::TcpStream::connect(remote_addr.clone()).unwrap(); let tls_stream = tls_connector.connect(remote_addr.as_str(), tcp_stream).unwrap(); let (mut raw_ws_stream, _) = tungstenite::client(uri, tls_stream).unwrap(); raw_ws_stream.send(Message::Text(login_msg)).unwrap(); let arc_ws_stream = Arc::new(Mutex::new(raw_ws_stream)); let thread_ws_stream = arc_ws_stream.clone(); // 创建原子布尔值作为退出信号 let should_exit = Arc::new(AtomicBool::new(false)); let thread_exit = should_exit.clone(); let ws_recv_thread = thread::spawn(move || { loop { // 先检查是否需要退出,避免无意义的锁竞争 if thread_exit.load(Ordering::Relaxed) { break; } // 仅在read时持有锁,读取完成后立即释放 let msg = { let mut ws = thread_ws_stream.lock().unwrap(); ws.read() }; if !process_ws_msg(msg) { break; } } }); sleep(Duration::from_millis(5_000)); // 设置退出标志,通知子线程退出 should_exit.store(true, Ordering::Relaxed); // 等待子线程退出后再关闭连接 let _ = ws_recv_thread.join(); // 此时子线程已释放锁,可以安全获取并关闭连接 let _ = arc_ws_stream.lock().unwrap().close(None); } fn process_ws_msg(msg: Result<Message, Error>) -> bool { match msg { Ok(msg) => match msg { Message::Text(msg) => { println!("WS received message: {}", msg); } Message::Close(cls_msg) => { match cls_msg { Some(cls_frame) => { println!("WS session closed message, close frame = {}", cls_frame); } None => { println!("WS session closed message, no close frame"); } } return false; } _ => {} }, Err(err) => { println!("WS session error: {}", err); return false; } } true }
代码说明
- 原子布尔值
should_exit:通过原子操作实现线程间安全通信,主线程设置为true后,子线程会在下一次循环检查时主动退出。 - 缩小互斥锁范围:用代码块包裹
ws.read(),让锁仅在读取操作期间被持有,读取完成后立即释放,避免阻塞期间占用锁。 - 调整关闭顺序:先等待子线程退出,再获取锁关闭连接,彻底避免锁竞争。
内容的提问来源于stack exchange,提问作者roy
相关产品推荐
相关产品推荐

