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

如何主动关闭同步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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 22:46:11