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

如何在Rust中为自定义AsyncUtpStream实现AsyncRead与AsyncWrite

基于tokio的类uTP协议AsyncUtpStream实现AsyncRead/AsyncWrite的问题

现状说明

目前基于tokio::net::UdpSocket实现的类uTP协议底层UDP通信正常,AsyncUtpSocket的异步send_to和recv_from方法可正常工作。服务器与客户端交互示例输出如下:

服务器输出

UDP server listening on 127.0.0.1:8080...
Received 20 bytes from 127.0.0.1:37776. Presumably a SYN packet.
---> Received SYN packet
Sent 20 bytes back to client.
Accepted connection from 127.0.0.1:37776
Received 5 bytes from 127.0.0.1:37776: "hello"
Received 18 bytes from 127.0.0.1:37776: "my name is shivank"

客户端输出

UDP client sending packets...
Sent SYN packet to 127.0.0.1:8080
Received SYN-ACK packet
Connection established with the server
Sent 5 bytes to server: "hello"
Sent 18 bytes to server: "my name is shivank"

相关代码

1. AsyncUtpSocket(底层UDP通信封装)

#[derive(Debug, Clone)]
pub struct AsyncUtpSocket {
    pub(crate) socket: Arc<Mutex<UdpSocket>>,
    pub(crate) connected_to: Option<SocketAddr>,
    pub(crate) congestion_timeout: u64,
    pub(crate) state: SocketState,
}

impl AsyncUtpSocket {
    pub async fn send_to(&self, buf: &[u8]) -> io::Result<usize> { /* 已实现且正常工作 */ }
    pub async fn recv_from(&self, buf: &mut [u8]) -> io::Result<(usize, SocketAddr)> { /* 已实现且正常工作 */ }
}

2. AsyncUtpStream(待实现AsyncRead/AsyncWrite)

pub struct AsyncUtpStream {
    pub socket: AsyncUtpSocket,
    // 需添加内部缓冲区存储未读取的uTP数据包内容
    pub(crate) read_buf: Vec<u8>,
}

// 省略bind和connect方法

impl AsyncRead for AsyncUtpStream {
    fn poll_read(
        mut self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        buf: &mut ReadBuf<'_>
    ) -> Poll<Result<()>> {
        // 待实现:处理异步recv_from并填充ReadBuf
    }
}

impl AsyncWrite for AsyncUtpStream {
    fn poll_write(
        mut self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        data: &[u8]
    ) -> Poll<Result<usize>> {
        // Todo!()
    }

    fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
        // uTP是可靠协议,需等待所有已发送数据包被确认后再返回Ready
        Poll::Ready(Ok(()))
    }

    fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
        // 发送uTP的FIN包,等待对方确认后关闭连接
        Poll::Ready(Ok(()))
    }
}

遇到的问题

  • 无法正确实现poll_read,处理异步数据包并填充ReadBuf
  • 不清楚如何实现poll_write以异步发送UDP数据
  • 想了解基于UDP的协议实现AsyncRead/AsyncWrite的模式与最佳实践

解决方案指导

1. poll_read实现要点

由于uTP是基于数据包的可靠协议,AsyncUtpStream需要维护内部读缓冲区暂存已接收但未被上层读取的数据包内容。实现步骤:

  1. 优先处理内部缓冲区的剩余数据,直接填充到ReadBuf并更新缓冲区
  2. 缓冲区为空时,异步接收新数据包,提取payload后填充ReadBuf,剩余数据存入内部缓冲区

示例实现:

impl AsyncRead for AsyncUtpStream {
    fn poll_read(
        mut self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        buf: &mut ReadBuf<'_>
    ) -> Poll<Result<()>> {
        // 先处理内部缓冲区剩余数据
        if !self.read_buf.is_empty() {
            let copy_len = buf.remaining().min(self.read_buf.len());
            buf.put_slice(&self.read_buf[..copy_len]);
            self.read_buf.drain(..copy_len);
            return Poll::Ready(Ok(()));
        }

        // 缓冲区为空,尝试接收新数据包
        let mut recv_buf = vec![0; 1500]; // UDP MTU大小
        match self.socket.recv_from(&mut recv_buf).poll_unpin(cx) {
            Poll::Ready(Ok((len, _addr))) => {
                // 过滤uTP协议头,提取实际payload(需根据你的协议实现调整)
                let payload = &recv_buf[..len];
                let copy_len = buf.remaining().min(payload.len());
                buf.put_slice(&payload[..copy_len]);
                // 剩余payload存入内部缓冲区
                if copy_len < payload.len() {
                    self.read_buf.extend_from_slice(&payload[copy_len..]);
                }
                Poll::Ready(Ok(()))
            }
            Poll::Ready(Err(e)) => Poll::Ready(Err(e)),
            Poll::Pending => Poll::Pending,
        }
    }
}

2. poll_write实现要点

poll_write需将上层流式数据按uTP协议的数据包大小拆分,封装协议头后异步发送。利用poll_unpin适配已实现的异步send_to方法:

示例实现:

impl AsyncWrite for AsyncUtpStream {
    fn poll_write(
        mut self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        data: &[u8]
    ) -> Poll<Result<usize>> {
        // uTP数据包最大payload大小(根据协议定义调整,如1400字节)
        const MAX_PAYLOAD_SIZE: usize = 1400;
        let send_len = data.len().min(MAX_PAYLOAD_SIZE);
        
        // 封装uTP协议头(添加序列号、标志位等字段,需根据你的实现调整)
        let mut utp_packet = Vec::with_capacity(20 + send_len); // 假设协议头20字节
        // utp_packet.extend_from_slice(&protocol_header);
        utp_packet.extend_from_slice(&data[..send_len]);

        match self.socket.send_to(&utp_packet).poll_unpin(cx) {
            Poll::Ready(Ok(_)) => Poll::Ready(Ok(send_len)),
            Poll::Ready(Err(e)) => Poll::Ready(Err(e)),
            Poll::Pending => Poll::Pending,
        }
    }

    fn poll_flush(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
        // 需结合uTP的确认机制,等待所有未确认的数据包被对方接收
        // 可维护未确认数据包队列,队列空则返回Ready,否则注册waker等待
        Poll::Ready(Ok(()))
    }

    fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<()>> {
        // 发送uTP FIN包,等待对方返回FIN-ACK后完成关闭
        // 需实现FIN包发送与确认逻辑
        Poll::Ready(Ok(()))
    }
}

3. 基于UDP的可靠协议实现AsyncRead/AsyncWrite的最佳实践

  • 缓冲区管理:必须维护内部读/写缓冲区,实现流式接口与数据包协议的适配——将数据包拼接为流,或将流拆分为数据包
  • Waker注册:当poll操作返回Pending时,确保底层异步操作(如recv_from)注册了当前waker,保证操作完成时能唤醒任务
  • 状态同步:在AsyncUtpStream中维护连接状态(已建立、关闭中),确保poll方法仅在合法状态下执行
  • 错误映射:将UDP层面错误(如连接重置)与uTP协议层面错误(如超时、丢包)转换为标准io::Error返回
  • 协议兼容:严格遵循uTP协议的可靠传输机制(重传、流量控制、确认),与Async trait的poll模型结合,避免数据丢失或乱序

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 08:04:51