如何修复Rust中「值已被移动,位于循环前一次迭代」错误?
解决WebSocket转发程序中的E0382移动值错误
问题描述
我想编写一个实现本地WebSocket与远程WebSocket消息互转的程序,但在添加while循环生成异步任务时遇到了E0382: use of moved value: write_remote错误,ws_local也出现完全相同的错误。
错误详情
error[E0382]: use of moved value: `write_remote` | 42 | let (mut write_remote, mut read_remote) = ws_remote.split(); | ---------------- move occurs because `write_remote` has type `SplitSink<WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>, Message>`, which does not implement the `Copy` trait ... 70 | let _handle_two = task::spawn(async move { | ________________________________________________^ 71 | | while let Some(msg) = read_local.next().await { 72 | | let msg = msg?; 73 | | if msg.is_text() || msg.is_binary() { 74 | | write_remote.send(msg).await; | | ------------ use occurs due to use in generator ... | 78 | | Result::(), tungstenite::Error>::Ok(()) 79 | | }); | |_______^ value moved here, in previous iteration of loop
原始代码
#![cfg_attr( all(not(debug_assertions), target_os = "windows"), windows_subsystem = "windows" )] use tokio::net::{TcpListener, TcpStream}; use futures_util::{future, SinkExt, StreamExt, TryStreamExt}; use tokio_tungstenite::{ connect_async, accept_async, tungstenite::{Result}, }; use http::Request; use tokio::sync::oneshot; use futures::{ future::FutureExt, // for `.fuse()` pin_mut, select, }; use tokio::io::AsyncWriteExt; use std::io; use std::net::SocketAddr; use std::thread; use tokio::spawn; use tokio::task; async fn client() -> Result<()> { // Client let request = Request::builder() .method("GET") .header("Host", "demo.piesocket.com") // .header("Origin", "https://example.com/") .header("Connection", "Upgrade") .header("Upgrade", "websocket") .header("Sec-WebSocket-Version", "13") .header("Sec-WebSocket-Key", tungstenite::handshake::client::generate_key()) .uri("wss://demo.piesocket.com/v3/channel_1?api_key=VCXCEuvhGcBDP7XhiJJUDvR1e1D3eiVjgZ9VRiaV¬ify_self") .body(())?; let (mut ws_remote, _) = connect_async(request).await?; let (mut write_remote, mut read_remote) = ws_remote.split(); let listener = TcpListener::bind("127.0.0.1:4444").await.expect("Can't listen"); while let Ok((stream, _)) = listener.accept().await { let mut ws_local = accept_async(stream).await.expect("Failed to accept"); let (mut write_local, mut read_local) = ws_local.split(); // read_remote.try_filter(|msg| future::ready(msg.is_text() || msg.is_binary())) // .forward(write_local) // .await // .expect("Failed to forward messages"); // read_local.try_filter(|msg| future::ready(msg.is_text() || msg.is_binary())) // .forward(write_remote) // .await // .expect("Failed to forward messages"); let _handle_one = task::spawn(async move { while let Some(msg) = read_remote.next().await { let msg = msg?; if msg.is_text() || msg.is_binary() { write_local.send(msg).await; } }; Result<(), tungstenite::Error>::Ok(()) }); let _handle_two = task::spawn(async move { while let Some(msg) = read_local.next().await { let msg = msg?; if msg.is_text() || msg.is_binary() { write_remote.send(msg).await; } }; Result<(), tungstenite::Error>::Ok(()) }); // handle_one.await.expect("The task being joined has panicked"); // handle_two.await.expect("The task being joined has panicked"); } Ok(()) } fn main() { tauri::async_runtime::spawn(client()); tauri::Builder::default() // .plugin(PluginBuilder::default().build()) .run(tauri::generate_context!()) .expect("failed to run app"); }
错误原因分析
这个错误的核心是所有权问题:
write_remote和read_remote是从远程WebSocket连接拆分出来的Sink和Stream,它们不实现Copytrait,只能拥有唯一所有权。- 在while循环的第一次迭代中,你把
read_remote移动到了_handle_one的异步任务中,把write_remote移动到了_handle_two的异步任务中。 - 当循环进入第二次迭代时,这两个变量已经被移动,无法再次使用,因此编译器抛出E0382错误。
- 另外,原始设计存在逻辑问题:一个远程WebSocket连接不能同时被多个本地连接共享,每个本地连接应该对应独立的远程连接,或者需要用共享所有权的方式处理,但前者更符合WebSocket的一对一通信模型。
解决方案
方案1:每个本地连接创建独立的远程WebSocket连接(推荐)
把远程连接的创建逻辑放到while循环内部,这样每个新的本地连接都会建立自己的远程连接,彻底避免所有权冲突:
修改后的代码:
#![cfg_attr( all(not(debug_assertions), target_os = "windows"), windows_subsystem = "windows" )] use tokio::net::{TcpListener, TcpStream}; use futures_util::{future, SinkExt, StreamExt, TryStreamExt}; use tokio_tungstenite::{ connect_async, accept_async, tungstenite::{Result, Error}, }; use http::Request; use tokio::task; async fn client() -> Result<()> { let listener = TcpListener::bind("127.0.0.1:4444").await.expect("Can't listen"); while let Ok((stream, _)) = listener.accept().await { // 每个本地连接创建独立的远程WebSocket连接 let request = Request::builder() .method("GET") .header("Host", "demo.piesocket.com") .header("Connection", "Upgrade") .header("Upgrade", "websocket") .header("Sec-WebSocket-Version", "13") .header("Sec-WebSocket-Key", tungstenite::handshake::client::generate_key()) .uri("wss://demo.piesocket.com/v3/channel_1?api_key=VCXCEuvhGcBDP7XhiJJUDvR1e1D3eiVjgZ9VRiaV¬ify_self") .body(())?; let (mut ws_remote, _) = connect_async(request).await?; let (mut write_remote, mut read_remote) = ws_remote.split(); let mut ws_local = accept_async(stream).await.expect("Failed to accept"); let (mut write_local, mut read_local) = ws_local.split(); // 转发远程消息到本地 let handle_one = task::spawn(async move { while let Some(msg) = read_remote.next().await { match msg { Ok(msg) if msg.is_text() || msg.is_binary() => { if let Err(e) = write_local.send(msg).await { eprintln!("Failed to send to local: {}", e); break; } } Err(e) => { eprintln!("Remote read error: {}", e); break; } _ => {} } } Result::<(), Error>::Ok(()) }); // 转发本地消息到远程 let handle_two = task::spawn(async move { while let Some(msg) = read_local.next().await { match msg { Ok(msg) if msg.is_text() || msg.is_binary() => { if let Err(e) = write_remote.send(msg).await { eprintln!("Failed to send to remote: {}", e); break; } } Err(e) => { eprintln!("Local read error: {}", e); break; } _ => {} } } Result::<(), Error>::Ok(()) }); // 等待两个任务完成,避免连接提前关闭 let _ = tokio::join!(handle_one, handle_two); } Ok(()) } fn main() { tauri::async_runtime::spawn(client()); tauri::Builder::default() .run(tauri::generate_context!()) .expect("failed to run app"); }
方案2:用共享所有权共享远程连接(适合单远程多本地场景)
如果确实需要一个远程连接对应多个本地连接,可以使用Arc<Mutex<...>>来包装write_remote和read_remote,让多个任务共享所有权:
// 需要添加use语句 use std::sync::Arc; use tokio::sync::Mutex; async fn client() -> Result<()> { // 建立远程连接并包装成Arc<Mutex> let request = Request::builder() .method("GET") .header("Host", "demo.piesocket.com") .header("Connection", "Upgrade") .header("Upgrade", "websocket") .header("Sec-WebSocket-Version", "13") .header("Sec-WebSocket-Key", tungstenite::handshake::client::generate_key()) .uri("wss://demo.piesocket.com/v3/channel_1?api_key=VCXCEuvhGcBDP7XhiJJUDvR1e1D3eiVjgZ9VRiaV¬ify_self") .body(())?; let (mut ws_remote, _) = connect_async(request).await?; let (write_remote, read_remote) = ws_remote.split(); // 包装成共享所有权的类型 let write_remote = Arc::new(Mutex::new(write_remote)); let read_remote = Arc::new(Mutex::new(read_remote)); let listener = TcpListener::bind("127.0.0.1:4444").await.expect("Can't listen"); while let Ok((stream, _)) = listener.accept().await { let mut ws_local = accept_async(stream).await.expect("Failed to accept"); let (mut write_local, mut read_local) = ws_local.split(); // 克隆共享引用 let write_remote_clone = Arc::clone(&write_remote); let read_remote_clone = Arc::clone(&read_remote); // 转发远程消息到本地 let handle_one = task::spawn(async move { let mut read_remote = read_remote_clone.lock().await; while let Some(msg) = read_remote.next().await { match msg { Ok(msg) if msg.is_text() || msg.is_binary() => { if let Err(e) = write_local.send(msg).await { eprintln!("Failed to send to local: {}", e); break; } } Err(e) => { eprintln!("Remote read error: {}", e); break; } _ => {} } } Result::<(), Error>::Ok(()) }); // 转发本地消息到远程 let handle_two = task::spawn(async move { while let Some(msg) = read_local.next().await { match msg { Ok(msg) if msg.is_text() || msg.is_binary() => { let mut write_remote = write_remote_clone.lock().await; if let Err(e) = write_remote.send(msg).await { eprintln!("Failed to send to remote: {}", e); break; } } Err(e) => { eprintln!("Local read error: {}", e); break; } _ => {} } } Result::<(), Error>::Ok(()) }); // 不等待任务完成,让多个本地连接同时共享远程连接 // 注意:这种方式需要处理连接关闭的情况,避免任务 panic let _ = handle_one; let _ = handle_two; } Ok(()) }
关键修改点说明
方案1:
- 将远程连接创建逻辑移入while循环,每个本地连接对应独立的远程连接,彻底避免所有权冲突。
- 添加了错误处理,避免任务因panic终止。
- 使用
tokio::join!等待转发任务完成,确保连接在消息转发结束后再关闭。
方案2:
- 用
Arc<Mutex<>>包装远程的Sink和Stream,实现多任务共享所有权。 - 每个任务克隆
Arc引用,在需要操作时通过lock().await获取互斥锁。 - 适合需要单远程连接对应多本地连接的场景,但要注意互斥锁带来的性能开销和并发问题。
- 用
内容的提问来源于stack exchange,提问作者Kenn Wong
相关产品推荐
相关产品推荐

