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

如何修复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&notify_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,它们不实现Copy trait,只能拥有唯一所有权。
  • 在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&notify_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&notify_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. 方案1:

    • 将远程连接创建逻辑移入while循环,每个本地连接对应独立的远程连接,彻底避免所有权冲突。
    • 添加了错误处理,避免任务因panic终止。
    • 使用tokio::join!等待转发任务完成,确保连接在消息转发结束后再关闭。
  2. 方案2:

    • 用Arc<Mutex<>>包装远程的Sink和Stream,实现多任务共享所有权。
    • 每个任务克隆Arc引用,在需要操作时通过lock().await获取互斥锁。
    • 适合需要单远程连接对应多本地连接的场景,但要注意互斥锁带来的性能开销和并发问题。

内容的提问来源于stack exchange,提问作者Kenn Wong

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 07:06:21