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

使用Tokio copy_bidirectional实现代理时数据接收异常及流读取问题

问题:Tokio代理中server2无法接收客户端消息及代理读取流消息的疑问

我想用Rust的Tokio库中的tokio::io::copy_bidirectional实现一个简单代理,流程是client1连接代理server1,server1再连接server2,但现在遇到server2始终无法接收客户端消息的问题。另外,我还想了解如何让代理读取流中的消息。以下是我的代码和终端输出,希望找出问题原因:

use std::net::{SocketAddr, ToSocketAddrs};
use tokio::net::{TcpListener, TcpStream};
use tokio::io::{AsyncReadExt, AsyncWriteExt};

#[tokio::main]
async fn main() {
    // Initialize the server address
    let server1_addr = "127.0.0.1:8000".to_socket_addrs().unwrap().next().unwrap();

    let server2_addr = "127.0.0.1:9000".to_socket_addrs().unwrap().next().unwrap();

    //Init the server and client
    let server1_handle = tokio::spawn(server1(server1_addr, server2_addr));

    let server2_handle = tokio::spawn(server2(server2_addr));

    let client_handle = tokio::spawn(client(server1_addr));

    // Wait for the server to complete its work
    server1_handle.await.unwrap();
    server2_handle.await.unwrap();
    client_handle.await.unwrap();
}

async fn server1(server1_addr: SocketAddr, server2_addr: SocketAddr) {
    let listener = TcpListener::bind(&server1_addr).await.unwrap();
    while let Ok((mut stream1, _)) = listener.accept().await {
        let mut stream2 = TcpStream::connect(&server2_addr).await.unwrap();
        stream2.write_all(b"hello\n").await.unwrap();
        tokio::spawn(async move {
            match tokio::io::copy_bidirectional(&mut stream1, &mut stream2).await {
                Ok((n1, n2)) => {
                    println!("Server sent {} bytes and received {} bytes", n1, n2);
                }
                Err(e) => {
                    println!("Server error: {}", e);
                }
            }
            
        });
        
       
    }
    
}

async fn server2(server2_addr: SocketAddr) {
    let listener = TcpListener::bind(&server2_addr).await.unwrap();
    while let Ok((mut socket, _)) = listener.accept().await {
        println!("Server 2 accepted connection");
        let mut buf = Vec::new();
        let n = socket.read_to_end(&mut buf).await.unwrap();
        println!("Server 2 received {} bytes", n);
        let message = "world\n";
        socket.write_all(message.as_bytes()).await.unwrap();
        println!(" Server 2 sent message: {:?}", message);
    }
}

async fn client(server1_addr: SocketAddr) {
    let mut socket = TcpStream::connect(&server1_addr).await.unwrap();
    let message = "hello\n";
    socket.write_all(message.as_bytes()).await.unwrap();
    println!("Client sent message: {:?}", message);
    let mut buf = Vec::new();
    socket.read_to_end(&mut buf).await.unwrap();
    println!("Client received message: {:?}", buf);
}

终端输出:

Client sent message: "hello\n"
Server 2 accepted connection

问题分析与解决

1. server2无法接收客户端消息的原因

核心问题出在read_to_end的使用逻辑上:

  • read_to_end会持续读取数据,直到连接被主动关闭才返回结果。但当前场景中,client、server1、server2之间的连接始终保持打开状态,导致server2卡在read_to_end调用上,无法输出已接收的消息。
  • 实际上客户端发送的消息已经通过copy_bidirectional转发到了server2,只是server2没有机会处理并输出。

2. 修复后的代码

修改server2:改用循环读取固定缓冲区

将read_to_end替换为循环读取逻辑,实时处理收到的消息,无需等待连接关闭:

async fn server2(server2_addr: SocketAddr) {
    let listener = TcpListener::bind(&server2_addr).await.unwrap();
    while let Ok((mut socket, _)) = listener.accept().await {
        println!("Server 2 accepted connection");
        // 为每个连接启动独立任务处理
        tokio::spawn(async move {
            let mut buf = [0; 1024];
            loop {
                match socket.read(&mut buf).await {
                    Ok(0) => {
                        println!("Server 2: connection closed");
                        break;
                    }
                    Ok(n) => {
                        println!("Server 2 received {} bytes: {:?}", n, String::from_utf8_lossy(&buf[..n]));
                        // 回复客户端
                        let message = "world\n";
                        if let Err(e) = socket.write_all(message.as_bytes()).await {
                            println!("Server 2 write error: {}", e);
                            break;
                        }
                        println!("Server 2 sent message: {:?}", message);
                    }
                    Err(e) => {
                        println!("Server 2 read error: {}", e);
                        break;
                    }
                }
            }
        });
    }
}

修改client:避免read_to_end阻塞

将read_to_end改为单次读取,获取server2的回复后主动关闭连接,触发代理的copy_bidirectional任务结束:

async fn client(server1_addr: SocketAddr) {
    let mut socket = TcpStream::connect(&server1_addr).await.unwrap();
    let message = "hello\n";
    socket.write_all(message.as_bytes()).await.unwrap();
    println!("Client sent message: {:?}", message);
    
    let mut buf = [0; 1024];
    match socket.read(&mut buf).await {
        Ok(n) => {
            println!("Client received message: {:?}", String::from_utf8_lossy(&buf[..n]));
        }
        Err(e) => {
            println!("Client read error: {}", e);
        }
    }
    
    // 主动关闭连接,触发代理的copy_bidirectional结束
    socket.shutdown().await.unwrap();
}

修复后运行,终端会输出完整的交互流程,包括server2接收客户端消息、回复,以及client收到回复的内容。


如何让代理读取流中的消息

tokio::io::copy_bidirectional是黑盒式的双向数据复制,无法直接读取流中的内容。要实现代理读取消息,需要手动处理两个方向的数据流,在转发的同时读取内容:

修改server1的连接处理逻辑,替换copy_bidirectional为两个独立的异步任务,分别处理client -> server2和server2 -> client的数据流:

async fn server1(server1_addr: SocketAddr, server2_addr: SocketAddr) {
    let listener = TcpListener::bind(&server1_addr).await.unwrap();
    while let Ok((mut stream1, _)) = listener.accept().await {
        let mut stream2 = TcpStream::connect(&server2_addr).await.unwrap();
        stream2.write_all(b"hello\n").await.unwrap();
        
        tokio::spawn(async move {
            // 处理 client 到 server2 的数据流
            let client_to_server = async {
                let mut buf = [0; 1024];
                loop {
                    match stream1.read(&mut buf).await {
                        Ok(0) => break,
                        Ok(n) => {
                            let msg = String::from_utf8_lossy(&buf[..n]);
                            println!("Proxy received from client: {:?}", msg);
                            // 转发到server2
                            if let Err(e) = stream2.write_all(&buf[..n]).await {
                                println!("Proxy write to server2 error: {}", e);
                                break;
                            }
                        }
                        Err(e) => {
                            println!("Proxy read from client error: {}", e);
                            break;
                        }
                    }
                }
                stream1.shutdown().await.ok();
            };
            
            // 处理 server2 到 client 的数据流
            let server_to_client = async {
                let mut buf = [0; 1024];
                loop {
                    match stream2.read(&mut buf).await {
                        Ok(0) => break,
                        Ok(n) => {
                            let msg = String::from_utf8_lossy(&buf[..n]);
                            println!("Proxy received from server2: {:?}", msg);
                            // 转发到client
                            if let Err(e) = stream1.write_all(&buf[..n]).await {
                                println!("Proxy write to client error: {}", e);
                                break;
                            }
                        }
                        Err(e) => {
                            println!("Proxy read from server2 error: {}", e);
                            break;
                        }
                    }
                }
                stream2.shutdown().await.ok();
            };
            
            // 同时等待两个方向的任务完成
            tokio::join!(client_to_server, server_to_client);
            println!("Proxy connection closed");
        });
    }
}

这样代理就能在转发数据的同时,读取并打印两边的消息内容,实现对数据流的监控和自定义处理。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 03:54:54