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

基于Rust实现带RX/TX时间戳的异步Socket recvmsg

基于Tokio获取Socket RX/TX时间戳的解决方案

方案一:改进Tokio AsyncFd的错误队列监听

核心思路是让AsyncFd同时监听Socket的可读事件和异常事件——错误队列存在TX时间戳时会触发异常事件(POLLERR),这样就能自动触发读取逻辑,无需手动调用recvmsg(MSG_ERRQUEUE)。

  • 初始化AsyncFd时,同时注册READABLE和ERROR兴趣事件
  • 事件触发时按类型分支处理:
    • 若为ERROR事件,循环读取错误队列中的TX时间戳(直到recvmsg返回EAGAIN),处理完成后清除错误标记
    • 若为READABLE事件,读取普通接收队列的RX数据包及时间戳,处理后清除就绪标记

示例代码片段:

use tokio::io::AsyncFd;
use libc::{recvmsg, MSG_ERRQUEUE, EAGAIN};
use std::os::unix::io::AsRawFd;

// 假设已配置好SO_TIMESTAMPING的UDP/TCP Socket
let socket = ...;
let async_fd = AsyncFd::new(socket.as_raw_fd())?;

loop {
    let mut guard = async_fd.readable().await?;
    let events = guard.interest();

    if events.is_error() {
        // 循环读取错误队列,确保清空所有TX时间戳
        loop {
            let mut msg = libc::msghdr::default();
            let mut cmsg_buf = [0u8; 128];
            msg.msg_control = cmsg_buf.as_mut_ptr() as *mut libc::c_void;
            msg.msg_controllen = cmsg_buf.len();

            let ret = unsafe { recvmsg(socket.as_raw_fd(), &mut msg, MSG_ERRQUEUE) };
            if ret < 0 {
                if unsafe { *libc::__errno_location() } == EAGAIN {
                    break; // 队列已空,退出循环
                }
                continue;
            }

            // 解析控制消息中的TX时间戳
            parse_tx_timestamp(&msg);
        }
        guard.clear_error();
    } else {
        // 读取RX数据及时间戳
        let mut msg = libc::msghdr::default();
        let mut buf = [0u8; 1500];
        let mut iov = libc::iovec {
            iov_base: buf.as_mut_ptr() as *mut libc::c_void,
            iov_len: buf.len(),
        };
        msg.msg_iov = &mut iov;
        msg.msg_iovlen = 1;

        let mut cmsg_buf = [0u8; 128];
        msg.msg_control = cmsg_buf.as_mut_ptr() as *mut libc::c_void;
        msg.msg_controllen = cmsg_buf.len();

        let ret = unsafe { recvmsg(socket.as_raw_fd(), &mut msg, 0) };
        if ret >= 0 {
            let rx_ts = parse_rx_timestamp(&msg);
            process_rx_data(&buf[0..ret], rx_ts);
        }
        guard.clear_ready();
    }
}

方案二:脱离AsyncFd的原生事件循环方案

如果不想依赖Tokio的AsyncFd,可以直接用epoll监听Socket事件,结合Tokio通道将结果传递到异步任务中,完全自主控制事件触发逻辑。

  • 用epoll_create1创建epoll实例,注册EPOLLIN(普通读)和EPOLLERR(错误队列)事件
  • 在阻塞线程中调用epoll_wait等待事件,根据事件类型分别处理RX和TX时间戳
  • 通过tokio::sync::mpsc通道将解析后的时间戳/数据发送到异步任务处理

示例代码片段:

use tokio::sync::mpsc;
use libc::{
    epoll_create1, epoll_ctl, epoll_wait, EPOLLIN, EPOLLERR, EPOLL_CTL_ADD,
    recvmsg, MSG_ERRQUEUE, EAGAIN,
};
use std::os::unix::io::AsRawFd;

// 假设已配置好时间戳的Socket
let socket = ...;
let socket_fd = socket.as_raw_fd();

// 创建通道传递时间戳和数据
let (tx_tx, rx_tx) = mpsc::unbounded_channel();
let (tx_rx, rx_rx) = mpsc::unbounded_channel();

// 在阻塞线程中运行epoll事件循环
tokio::spawn_blocking(move || {
    let epoll_fd = unsafe { epoll_create1(0) };
    if epoll_fd < 0 {
        eprintln!("Failed to create epoll fd");
        return;
    }

    let mut event = libc::epoll_event {
        events: (EPOLLIN | EPOLLERR) as u32,
        u64: socket_fd as u64,
    };
    if unsafe { epoll_ctl(epoll_fd, EPOLL_CTL_ADD, socket_fd, &mut event) } < 0 {
        eprintln!("Failed to add socket to epoll");
        return;
    }

    let mut events = vec![libc::epoll_event::default(); 8];
    loop {
        let n = unsafe { epoll_wait(epoll_fd, events.as_mut_ptr(), 8, -1) };
        if n <= 0 {
            continue;
        }

        for e in &events[0..n as usize] {
            if e.events & EPOLLERR as u32 != 0 {
                // 清空错误队列的TX时间戳
                loop {
                    let mut msg = libc::msghdr::default();
                    let mut cmsg_buf = [0u8; 128];
                    msg.msg_control = cmsg_buf.as_mut_ptr() as *mut libc::c_void;
                    msg.msg_controllen = cmsg_buf.len();

                    let ret = unsafe { recvmsg(socket_fd, &mut msg, MSG_ERRQUEUE) };
                    if ret < 0 {
                        if unsafe { *libc::__errno_location() } == EAGAIN {
                            break;
                        }
                        continue;
                    }

                    if let Some(ts) = parse_tx_timestamp(&msg) {
                        let _ = tx_tx.send(ts);
                    }
                }
            }

            if e.events & EPOLLIN as u32 != 0 {
                // 读取RX数据和时间戳
                let mut msg = libc::msghdr::default();
                let mut buf = [0u8; 1500];
                let mut iov = libc::iovec {
                    iov_base: buf.as_mut_ptr() as *mut libc::c_void,
                    iov_len: buf.len(),
                };
                msg.msg_iov = &mut iov;
                msg.msg_iovlen = 1;

                let mut cmsg_buf = [0u8; 128];
                msg.msg_control = cmsg_buf.as_mut_ptr() as *mut libc::c_void;
                msg.msg_controllen = cmsg_buf.len();

                let ret = unsafe { recvmsg(socket_fd, &mut msg, 0) };
                if ret >= 0 {
                    if let Some(ts) = parse_rx_timestamp(&msg) {
                        let _ = tx_rx.send((&buf[0..ret].to_vec(), ts));
                    }
                }
            }
        }
    }
});

// 异步处理TX时间戳
tokio::spawn(async move {
    while let Some(tx_ts) = rx_tx.recv().await {
        handle_tx_timestamp(tx_ts);
    }
});

// 异步处理RX数据和时间戳
tokio::spawn(async move {
    while let Some((rx_data, rx_ts)) = rx_rx.recv().await {
        handle_rx_data(rx_data, rx_ts);
    }
});

关键注意事项

  • Socket时间戳配置:必须正确设置SO_TIMESTAMPING选项,比如启用SOF_TIMESTAMPING_TX_HARDWARE/SOF_TIMESTAMPING_TX_SOFTWARE和对应的RX标志,确保内核生成并存储时间戳到错误队列
  • 错误队列读取完整性:每次触发错误事件时,一定要循环读取直到recvmsg返回EAGAIN,否则残留的时间戳会导致后续事件无法正常触发
  • 线程安全:Socket文件描述符不能在多个线程中同时操作,epoll事件循环要单独放在阻塞线程中,通过通道传递结果到异步任务

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 23:31:04