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

Rust中WebSocket连接异常中断:无法进入local.spawn_local代码块求助

Rust + Tokio + Hyper WebSocket连接中断问题解决

问题现象

使用Rust结合Tokio和Hyper实现WebSocket服务时,连接建立1-2毫秒后突然中断。终端日志显示代码未进入local.spawn_local代码块,WebSocket连接无法完成。尝试在返回响应前添加local.await后,出现std::marker::Send trait未实现的编译错误,同时main函数中server.await处提示trait绑定不满足。

相关代码

#[tokio::main]
async fn main() {
    println!("Starting Tokio runtime...");
    if tokio::runtime::Handle::try_current().is_ok() {
        println!("Tokio runtime is active.");
    } else {
        println!("No active Tokio runtime.");
    }
    let local = tokio::task::LocalSet::new(); 
    println!("LocalSet created.");
    dotenv().ok();
    let services = Arc::new(services::new_services());
    let connected_clients: ConnectedClients = Arc::new(Mutex::new(HashMap::new()));

    let make_svc = make_service_fn(|_conn| {
        let services = Arc::clone(&services);
        let clients = Arc::clone(&connected_clients);
        async move {
            Ok::<_, hyper::Error>(service_fn(move |req: Request<Body>| {
                handle_request(req, services.clone(), clients.clone())
            }))
        }
    });

    let port: u16 = env::var("PORT")
        .unwrap_or_else(|_| "8080".to_string())
        .parse()
        .expect("PORT must be a valid number");

    let addr = format!("127.0.0.1:{}", port).parse().expect("Invalid address");

    local.run_until(async {
        let server = Server::bind(&addr).serve(make_svc);

        println!("Server listening on {}", addr);
        if let Err(e) = server.await {
            eprintln!("Server error: {}", e);
        }
    }).await;
}


async fn handle_request(
    mut req: Request<Body>, 
    _services: Arc<Services>, 
    clients: ConnectedClients
) -> Result<Response<Body>, hyper::Error> {
    if req.uri().path() == "/ws" {
        println!("/ws path entered");

        // Check if it's a WebSocket request
        if !is_websocket_request(&req) {
            return Ok(Response::builder()
                .status(StatusCode::BAD_REQUEST)
                .body(Body::from("Not a WebSocket request"))
                .unwrap());
        }
        println!("ITS A WS CONNECTION!");
        // Extract token from query parameters
        let token = extract_token_from_query(req.uri());
        let (mut user_id, 
            mut org_id, 
            mut stores, 
            mut first_name, 
            mut last_name, 
            mut is_admin, 
            mut route,
            mut device_id
        ) = (
            String::new(),
            String::new(),
            Vec::new(),
            String::new(),
            String::new(),
            false,
            String::new(),
            None,
        );

        match token {
            Some(ref t) => {
                // Token validation logic
                println!("Entered into match token");
                let secret_key = std::env::var("JWT_SECRET").unwrap_or_else(|_| "default_secret".to_string());
                let decoding_key = DecodingKey::from_secret(secret_key.as_ref());
                let validation = Validation::default();

                // TOKEN VALIDATION AND PARSING
            }
            _ => {}
        }
                
        // Get the WebSocket key from headers
        let ws_key = req.headers().get(SEC_WEBSOCKET_KEY).and_then(|v| v.to_str().ok()).unwrap();

        // Generate accept key
        let accept_key = derive_accept_key(ws_key.as_bytes());

        println!("WebSocket Key: {:?}", ws_key);
        println!("Accept Key: {}", accept_key);

        let mut response = Response::builder()
            .status(StatusCode::SWITCHING_PROTOCOLS)
            .header(CONNECTION, "Upgrade")
            .header(UPGRADE, "websocket")
            .header(SEC_WEBSOCKET_ACCEPT, accept_key);

        // Add protocol if present
        if let Some(protocol) = req.headers().get(SEC_WEBSOCKET_PROTOCOL) {
            response = response.header(SEC_WEBSOCKET_PROTOCOL, protocol);
        }

        let response = response.body(Body::empty()).unwrap();
        println!("response : {:#?}", response);

        // Perform WebSocket upgrade and wrapping outside `tokio::spawn`
        let device_id_clone = device_id.clone(); 
        
        let local = LocalSet::new();    

        println!("local {:#?}", local);
        if tokio::runtime::Handle::try_current().is_err() {
            println!("No active Tokio runtime");
        }
        println!("Incoming request: {:?}", req);
        
        local.spawn_local(async move{
            println!("Entered into local spawn");
            match on(req).await {
                Ok(upgraded) => {
                    println!("Connection upgraded");
            
                    let ws_stream: WebSocketStream<Upgraded> = WebSocketStream::from_raw_socket(
                        upgraded,
                        Role::Server,
                        Some(WebSocketConfig::default())
                    ).await;

                    let shared_ws_stream = Arc::new(RwLock::new(ws_stream));
                    let shared_ws_stream_clone = Arc::clone(&shared_ws_stream);
                    
                    let user_connection: Connection = Connection::new(
                        None, 
                        None, 
                        shared_ws_stream,
                        device_id.clone(), 
                        None, 
                        None,
                    );
            
                    println!("WebSocket stream created");
                
                    if let Err(e) = handle_websocket_stream(shared_ws_stream_clone, clients).await {
                        eprintln!("WebSocket handler error: {}", e);
                    }
                }
                Err(e) => eprintln!("Failed to upgrade connection: {}", e),
            }
        });
       
        Ok(response)
    } else {
        Ok(Response::new(Body::from("404 not found!")))
    }
}

终端日志

Starting Tokio runtime...
Tokio runtime is active.
LocalSet created.
Server listening on 127.0.0.1:8080
/ws path entered
WebSocket request validation:
  Upgrade header present: true
  Connection upgrade: true
  Has WebSocket key: true
  Has correct version: true
ITS A WS CONNECTION!
Entered into match token
Token validated successfully
CALLED
WebSocket Key: "rTNBeHrwfawHZ9uhlw3M1Q=="
Accept Key: pJhyS7Aq49YTSEt5997ylXf8Z7g=
response : Response {
    status: 101,
    version: HTTP/1.1,
    headers: {
        "connection": "Upgrade",
        "upgrade": "websocket",
        "sec-websocket-accept": "pJhyS7Aq49YTSEt5997ylXf8Z7g=",
    },
    body: Body(
        Empty,
    ),
}
local LocalSet
Incoming request: Request { method: GET, uri: /ws?token=JWT_TOKEN, version: HTTP/1.1, headers: {"sec-websocket-version": "13", "sec-websocket-key": "rTNBeHrwfawHZ9uhlw3M1Q==", "connection": "Upgrade", "upgrade": "websocket", "sec-websocket-extensions": "permessage-deflate; client_max_window_bits", "host": "localhost:8080"}, body: Body(Empty) }

添加local.await后的编译报错

future cannot be sent between threads safely
within `hyper::proto::h2::server::H2Stream<impl futures::Future<Output = Result<hyper::Response<Body>, hyper::Error>>, Body>`, the trait `std::marker::Send` is not implemented for `Rc<tokio::task::local::Context>`, which is required by `hyper::common::exec::Exec: hyper::common::exec::ConnStreamExec<_, _>`
the trait `hyper::common::exec::ConnStreamExec<F, B>` is implemented for `hyper::common::exec::Exec`

main函数中server.await处的报错

the trait bound `hyper::common::exec::Exec: hyper::common::exec::ConnStreamExec<impl futures::Future<Output = Result<hyper::Response<Body>, hyper::Error>>, Body>` is not satisfied
the trait `hyper::common::exec::ConnStreamExec<F, B>` is implemented for `hyper::common::exec::Exec`
required for `NewSvcTask<AddrStream, {async block@src/main.rs:114:9: 118:10}, ServiceFn<..., ...>, ..., ...>` to implement `futures::Future`
required for `hyper::common::exec::Exec` to implement `hyper::common::exec::NewSvcExec<AddrStream, {async block@src/main.rs:114:9: 118:10}, hyper::service::util::ServiceFn<{closure@src/main.rs:115:46: 115:71}, Body>, hyper::common::exec::Exec, hyper::server::server::NoopWatcher>`
1 redundant requirement hidden
required for `hyper::Server<AddrIncoming, hyper::service::make::MakeServiceFn<{closure@src/main.rs:111:36: 111:43}>>` to implement `futures::Future`
required for `hyper::Server<AddrIncoming, hyper::service::make::MakeServiceFn<{closure@src/main.rs:111:36: 111:43}>>` to implement `std::future::IntoFuture`
consider using `--verbose` to print the full type name to the console

问题原因与解决方案

原因分析

  1. local.spawn_local任务未执行:在handle_request中新建的LocalSet没有被调用run_until或run方法启动,导致内部spawn的任务永远不会被Tokio调度执行,WebSocket升级逻辑无法运行,连接自然中断。
  2. 添加local.await后的编译错误:LocalSet内部依赖非Send的Rc类型,而Hyper的请求处理流程要求返回的Future必须实现Send trait,因此出现编译失败。

修复步骤

  1. 移除不必要的LocalSet:代码中使用的WebSocketStream<Upgraded>及相关共享状态(Arc、RwLock)均实现了Send trait,完全可以使用Tokio的多线程任务调度,无需LocalSet。
  2. 替换local.spawn_local为tokio::spawn:直接使用Tokio的全局任务调度器来启动WebSocket处理任务。

修改后的handle_request中相关代码片段:

// 移除所有LocalSet相关代码
// let local = LocalSet::new();    
// println!("local {:#?}", local);

// 替换为tokio::spawn
tokio::spawn(async move{
    println!("Entered into spawn");
    match on(req).await {
        Ok(upgraded) => {
            println!("Connection upgraded");
    
            let ws_stream: WebSocketStream<Upgraded> = WebSocketStream::from_raw_socket(
                upgraded,
                Role::Server,
                Some(WebSocketConfig::default())
            ).await;

            let shared_ws_stream = Arc::new(RwLock::new(ws_stream));
            let shared_ws_stream_clone = Arc::clone(&shared_ws_stream);
            
            let user_connection: Connection = Connection::new(
                None, 
                None, 
                shared_ws_stream,
                device_id.clone(), 
                None, 
                None,
            );
    
            println!("WebSocket stream created");
        
            if let Err(e) = handle_websocket_stream(shared_ws_stream_clone, clients).await {
                eprintln!("WebSocket handler error: {}", e);
            }
        }
        Err(e) => eprintln!("Failed to upgrade connection: {}", e),
    }
});
  1. 验证Connection结构体的Send实现:确保Connection结构体中的所有字段都实现了Send trait,若使用了非Send类型(如Rc),需替换为Arc等线程安全的替代方案。

完成上述修改后,WebSocket任务会被正常调度执行,连接不会再中断,同时编译错误也会消失。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 05:09:51