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
问题原因与解决方案
原因分析
local.spawn_local任务未执行:在handle_request中新建的LocalSet没有被调用run_until或run方法启动,导致内部spawn的任务永远不会被Tokio调度执行,WebSocket升级逻辑无法运行,连接自然中断。- 添加
local.await后的编译错误:LocalSet内部依赖非Send的Rc类型,而Hyper的请求处理流程要求返回的Future必须实现Sendtrait,因此出现编译失败。
修复步骤
- 移除不必要的
LocalSet:代码中使用的WebSocketStream<Upgraded>及相关共享状态(Arc、RwLock)均实现了Sendtrait,完全可以使用Tokio的多线程任务调度,无需LocalSet。 - 替换
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), } });
- 验证
Connection结构体的Send实现:确保Connection结构体中的所有字段都实现了Sendtrait,若使用了非Send类型(如Rc),需替换为Arc等线程安全的替代方案。
完成上述修改后,WebSocket任务会被正常调度执行,连接不会再中断,同时编译错误也会消失。
内容的提问来源于stack exchange,提问作者Sagar Maheshwari
相关产品推荐
相关产品推荐

