基于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数据包及时间戳,处理后清除就绪标记
- 若为ERROR事件,循环读取错误队列中的TX时间戳(直到
示例代码片段:
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
相关产品推荐
相关产品推荐

