#![cfg(feature = "io_uring")]
#![cfg(target_os = "linux")]
#![allow(unsafe_code)]
use crate::error::{LinuxError, Result};
use crate::syscalls::pipe2;
use io_uring::{opcode, squeue::Flags as SqeFlags, types, IoUring};
use std::os::unix::io::RawFd;
#[derive(Debug, Clone, Copy)]
pub struct Completion {
pub result: i32,
pub user_data: u64,
}
pub struct IoUringBatcher {
ring: IoUring,
}
impl std::fmt::Debug for IoUringBatcher {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("IoUringBatcher").finish_non_exhaustive()
}
}
impl IoUringBatcher {
pub fn new(entries: u32) -> Result<Self> {
let ring = IoUring::new(entries).map_err(|e| LinuxError::Syscall {
syscall: "io_uring_setup",
errno: e.raw_os_error().unwrap_or(0),
})?;
Ok(Self { ring })
}
#[allow(clippy::too_many_arguments)]
pub fn push_splice(
&mut self,
fd_in: RawFd,
off_in: i64,
fd_out: RawFd,
off_out: i64,
len: u32,
splice_flags: u32,
user_data: u64,
) -> Result<()> {
let entry = opcode::Splice::new(
types::Fd(fd_in),
off_in,
types::Fd(fd_out),
off_out,
len,
)
.flags(splice_flags)
.build()
.user_data(user_data);
let mut sq = self.ring.submission();
unsafe {
sq.push(&entry).map_err(|_| {
LinuxError::InsufficientResources("io_uring SQ full".to_string())
})?;
}
Ok(())
}
pub fn push_sendfile(
&mut self,
file_fd: RawFd,
pipe_fd: RawFd,
socket_fd: RawFd,
len: u32,
splice_flags: u32,
user_data: u64,
) -> Result<()> {
if self.sq_space_left() < 2 {
return Err(LinuxError::InsufficientResources(
"io_uring SQ 空间不足 2(sendfile 需要两个链式 SQE)".to_string(),
));
}
let entry1 = opcode::Splice::new(
types::Fd(file_fd),
-1,
types::Fd(pipe_fd),
-1,
len,
)
.flags(splice_flags)
.build()
.user_data(user_data)
.flags(SqeFlags::IO_LINK);
let entry2 = opcode::Splice::new(
types::Fd(pipe_fd),
-1,
types::Fd(socket_fd),
-1,
len,
)
.flags(splice_flags)
.build()
.user_data(user_data.wrapping_add(1));
let mut sq = self.ring.submission();
unsafe {
sq.push(&entry1).map_err(|_| {
LinuxError::InsufficientResources("io_uring SQ full".to_string())
})?;
sq.push(&entry2).map_err(|_| {
LinuxError::InsufficientResources("io_uring SQ full".to_string())
})?;
}
Ok(())
}
pub fn push_nop(&mut self, user_data: u64) -> Result<()> {
let entry = opcode::Nop::new().build().user_data(user_data);
let mut sq = self.ring.submission();
unsafe {
sq.push(&entry).map_err(|_| {
LinuxError::InsufficientResources("io_uring SQ full".to_string())
})?;
}
Ok(())
}
pub fn submit_and_wait(&mut self, min_complete: u32) -> Result<usize> {
self.ring
.submit_and_wait(min_complete as usize)
.map_err(|e| LinuxError::Syscall {
syscall: "io_uring_enter",
errno: e.raw_os_error().unwrap_or(0),
})
}
pub fn submit(&mut self) -> Result<usize> {
self.ring
.submit()
.map_err(|e| LinuxError::Syscall {
syscall: "io_uring_enter",
errno: e.raw_os_error().unwrap_or(0),
})
}
pub fn collect_completions(&mut self) -> impl Iterator<Item = Completion> + '_ {
self.ring.completion().map(|cqe| Completion {
result: cqe.result(),
user_data: cqe.user_data(),
})
}
pub fn sq_space_left(&mut self) -> usize {
let cap = self.ring.submission().capacity();
let len = self.ring.submission().len();
cap - len
}
pub fn cq_ready(&mut self) -> usize {
self.ring.completion().len()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[repr(u8)]
enum Dir {
C2UIn = 0, C2UOut = 1, U2CIn = 2, U2COut = 3, }
fn on_out_completion(pending: u32, drained: u32) -> (u32, bool) {
let new_pending = pending.saturating_sub(drained);
(new_pending, new_pending > 0)
}
struct Session {
client_fd: i32,
upstream_fd: i32,
c2u_read: i32,
c2u_write: i32,
u2c_read: i32,
u2c_write: i32,
bytes_c2u: usize,
bytes_u2c: usize,
c2u_eof: bool,
u2c_eof: bool,
c2u_pipe_pending: u32,
u2c_pipe_pending: u32,
}
impl Drop for Session {
fn drop(&mut self) {
unsafe {
libc::close(self.c2u_read);
libc::close(self.c2u_write);
libc::close(self.u2c_read);
libc::close(self.u2c_write);
}
}
}
pub struct SpliceBatcher {
ring: IoUringBatcher,
pipe_buf_size: usize,
}
impl SpliceBatcher {
pub fn new(sq_entries: u32, pipe_buf_size: usize) -> Result<Self> {
Ok(Self {
ring: IoUringBatcher::new(sq_entries)?,
pipe_buf_size,
})
}
#[inline]
fn enc_user_data(session_idx: u64, dir: Dir) -> u64 {
(session_idx << 4) | (dir as u64)
}
#[inline]
fn dec_user_data(user_data: u64) -> Result<(usize, Dir)> {
let session_idx = (user_data >> 4) as usize;
let dir_bits = (user_data & 0xF) as u8;
let dir = match dir_bits {
0 => Dir::C2UIn,
1 => Dir::C2UOut,
2 => Dir::U2CIn,
3 => Dir::U2COut,
other => {
return Err(LinuxError::Unsupported(format!(
"io_uring CQE user_data 含保留方向位 {other}(编码损坏)"
)))
}
};
Ok((session_idx, dir))
}
pub fn drive_all(mut self, pairs: &[(i32, i32)]) -> Result<Vec<(usize, usize)>> {
use std::os::unix::io::RawFd;
let n_sessions = pairs.len();
if n_sessions == 0 {
return Ok(Vec::new());
}
let mut sessions: Vec<Session> = Vec::with_capacity(n_sessions);
for (client_fd, upstream_fd) in pairs.iter().copied() {
let (c2u_read, c2u_write) = pipe2(0)?;
let (u2c_read, u2c_write) = pipe2(0)?;
sessions.push(Session {
client_fd,
upstream_fd,
c2u_read,
c2u_write,
u2c_read,
u2c_write,
bytes_c2u: 0,
bytes_u2c: 0,
c2u_eof: false,
u2c_eof: false,
c2u_pipe_pending: 0,
u2c_pipe_pending: 0,
});
}
let mut in_flight: u32 = 0;
let flags_splice: u32 = libc::SPLICE_F_MOVE | libc::SPLICE_F_NONBLOCK;
let len = self.pipe_buf_size as u32;
for (i, s) in sessions.iter().enumerate() {
let remaining = self.ring.sq_space_left();
if remaining < 2 {
self.ring.submit_and_wait(in_flight.min(1))?;
}
self.ring.push_splice(
s.client_fd as RawFd, -1,
s.c2u_write as RawFd, -1,
len, flags_splice,
Self::enc_user_data(i as u64, Dir::C2UIn),
)?;
self.ring.push_splice(
s.upstream_fd as RawFd, -1,
s.u2c_write as RawFd, -1,
len, flags_splice,
Self::enc_user_data(i as u64, Dir::U2CIn),
)?;
in_flight = in_flight.saturating_add(2);
}
let mut active_sessions = n_sessions;
while active_sessions > 0 {
let waited = self.ring.submit_and_wait(1)?;
let _ = waited;
let completions: Vec<Completion> = self.ring.collect_completions().collect();
for comp in completions {
in_flight = in_flight.saturating_sub(1);
let (s_idx, dir) = Self::dec_user_data(comp.user_data)?;
let s = sessions.get_mut(s_idx).ok_or_else(|| {
LinuxError::Unsupported(format!(
"io_uring CQE session 索引 {s_idx} 越界(共 {n_sessions} 个会话)"
))
})?;
let (transferred, eagain): (u32, bool) = if comp.result < 0 {
let err = -comp.result;
match err {
libc::EPIPE | libc::ECONNRESET | libc::ESHUTDOWN => {
match dir {
Dir::C2UIn | Dir::C2UOut => {
if !s.c2u_eof { s.c2u_eof = true; }
if matches!(dir, Dir::C2UOut) {
s.c2u_pipe_pending = 0;
}
}
Dir::U2CIn | Dir::U2COut => {
if !s.u2c_eof { s.u2c_eof = true; }
if matches!(dir, Dir::U2COut) {
s.u2c_pipe_pending = 0;
}
}
}
(0, false)
}
libc::EAGAIN => {
(0, true)
}
_ => {
return Err(LinuxError::Syscall {
syscall: "io_uring SPLICE",
errno: err,
});
}
}
} else {
(comp.result as u32, false)
};
match dir {
Dir::C2UIn => {
if transferred > 0 {
s.bytes_c2u = s.bytes_c2u.checked_add(transferred as usize)
.ok_or_else(|| LinuxError::InsufficientResources(
"c2u splice byte count overflow".to_string()
))?;
s.c2u_pipe_pending = s.c2u_pipe_pending.saturating_add(transferred);
} else if !eagain {
s.c2u_eof = true;
}
if !s.c2u_eof && self.ring.sq_space_left() > 0 {
let _ = self.ring.push_splice(
s.client_fd as RawFd, -1,
s.c2u_write as RawFd, -1,
len, flags_splice,
Self::enc_user_data(s_idx as u64, Dir::C2UIn),
);
in_flight = in_flight.saturating_add(1);
}
if s.c2u_pipe_pending > 0 && self.ring.sq_space_left() > 0 {
let to_drain = s.c2u_pipe_pending.min(len);
let _ = self.ring.push_splice(
s.c2u_read as RawFd, -1,
s.upstream_fd as RawFd, -1,
to_drain, flags_splice,
Self::enc_user_data(s_idx as u64, Dir::C2UOut),
);
in_flight = in_flight.saturating_add(1);
}
}
Dir::C2UOut => {
let (new_pending, requeue) =
on_out_completion(s.c2u_pipe_pending, transferred);
s.c2u_pipe_pending = new_pending;
if requeue && self.ring.sq_space_left() > 0 {
let to_drain = new_pending.min(len);
let _ = self.ring.push_splice(
s.c2u_read as RawFd, -1,
s.upstream_fd as RawFd, -1,
to_drain, flags_splice,
Self::enc_user_data(s_idx as u64, Dir::C2UOut),
);
in_flight = in_flight.saturating_add(1);
}
if s.c2u_eof && s.c2u_pipe_pending == 0 {
unsafe { let _ = libc::shutdown(s.upstream_fd, libc::SHUT_WR); }
}
}
Dir::U2CIn => {
if transferred > 0 {
s.bytes_u2c = s.bytes_u2c.checked_add(transferred as usize)
.ok_or_else(|| LinuxError::InsufficientResources(
"u2c splice byte count overflow".to_string()
))?;
s.u2c_pipe_pending = s.u2c_pipe_pending.saturating_add(transferred);
} else if !eagain {
s.u2c_eof = true;
}
if !s.u2c_eof && self.ring.sq_space_left() > 0 {
let _ = self.ring.push_splice(
s.upstream_fd as RawFd, -1,
s.u2c_write as RawFd, -1,
len, flags_splice,
Self::enc_user_data(s_idx as u64, Dir::U2CIn),
);
in_flight = in_flight.saturating_add(1);
}
if s.u2c_pipe_pending > 0 && self.ring.sq_space_left() > 0 {
let to_drain = s.u2c_pipe_pending.min(len);
let _ = self.ring.push_splice(
s.u2c_read as RawFd, -1,
s.client_fd as RawFd, -1,
to_drain, flags_splice,
Self::enc_user_data(s_idx as u64, Dir::U2COut),
);
in_flight = in_flight.saturating_add(1);
}
}
Dir::U2COut => {
let (new_pending, requeue) =
on_out_completion(s.u2c_pipe_pending, transferred);
s.u2c_pipe_pending = new_pending;
if requeue && self.ring.sq_space_left() > 0 {
let to_drain = new_pending.min(len);
let _ = self.ring.push_splice(
s.u2c_read as RawFd, -1,
s.client_fd as RawFd, -1,
to_drain, flags_splice,
Self::enc_user_data(s_idx as u64, Dir::U2COut),
);
in_flight = in_flight.saturating_add(1);
}
if s.u2c_eof && s.u2c_pipe_pending == 0 {
unsafe { let _ = libc::shutdown(s.client_fd, libc::SHUT_WR); }
}
}
}
let s_done = s.c2u_eof
&& s.u2c_eof
&& s.c2u_pipe_pending == 0
&& s.u2c_pipe_pending == 0;
if s_done {
active_sessions = active_sessions.saturating_sub(1);
}
}
}
let mut results = Vec::with_capacity(sessions.len());
for s in sessions.drain(..) {
results.push((s.bytes_c2u, s.bytes_u2c));
}
Ok(results)
}
}
impl std::fmt::Debug for SpliceBatcher {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SpliceBatcher")
.field("pipe_buf_size", &self.pipe_buf_size)
.finish_non_exhaustive()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_io_uring_create() {
let result = IoUringBatcher::new(8);
if let Ok(mut batcher) = result {
assert!(batcher.sq_space_left() > 0);
}
}
#[test]
fn test_io_uring_nop() {
let mut batcher = match IoUringBatcher::new(8) {
Ok(b) => b,
Err(_) => return, };
batcher.push_nop(100).unwrap();
batcher.push_nop(200).unwrap();
batcher.submit_and_wait(2).unwrap();
let completions: Vec<_> = batcher.collect_completions().collect();
assert_eq!(completions.len(), 2);
for c in &completions {
assert_eq!(c.result, 0);
}
let mut user_data_set: Vec<u64> = completions.iter().map(|c| c.user_data).collect();
user_data_set.sort_unstable();
assert_eq!(user_data_set, vec![100, 200]);
}
#[test]
fn test_dec_user_data_roundtrip_and_invalid() {
let (idx, dir) = SpliceBatcher::dec_user_data(SpliceBatcher::enc_user_data(7, Dir::C2UOut))
.unwrap();
assert_eq!(idx, 7);
assert!(matches!(dir, Dir::C2UOut));
let (idx, dir) = SpliceBatcher::dec_user_data(SpliceBatcher::enc_user_data(42, Dir::U2CIn))
.unwrap();
assert_eq!(idx, 42);
assert!(matches!(dir, Dir::U2CIn));
for bad in 4u64..=15 {
assert!(
SpliceBatcher::dec_user_data((1 << 4) | bad).is_err(),
"方向位 {bad} 必须 Fail-Closed"
);
}
}
#[test]
fn test_io_uring_sq_space() {
let mut batcher = match IoUringBatcher::new(4) {
Ok(b) => b,
Err(_) => return,
};
let initial_space = batcher.sq_space_left();
assert!(initial_space >= 4);
batcher.push_nop(1).unwrap();
assert_eq!(batcher.sq_space_left(), initial_space - 1);
batcher.push_nop(2).unwrap();
assert_eq!(batcher.sq_space_left(), initial_space - 2);
batcher.submit().unwrap();
}
}
use std::net::{Ipv4Addr, Ipv6Addr, SocketAddr};
const UD_SEND_FLAG: u64 = 1u64 << 63;
#[derive(Debug)]
struct UdpSlot {
msg: libc::msghdr,
iov: libc::iovec,
data: Vec<u8>,
addr: libc::sockaddr_storage,
in_use: bool,
}
impl UdpSlot {
fn new(packet_size: usize) -> Box<Self> {
Box::new(Self {
msg: unsafe { std::mem::zeroed() },
iov: libc::iovec {
iov_base: std::ptr::null_mut(),
iov_len: 0,
},
data: vec![0u8; packet_size],
addr: unsafe { std::mem::zeroed() },
in_use: false,
})
}
}
fn addr_to_storage(addr: &SocketAddr) -> (libc::sockaddr_storage, u32) {
let mut storage: libc::sockaddr_storage = unsafe { std::mem::zeroed() };
match addr {
SocketAddr::V4(v4) => {
let sa = libc::sockaddr_in {
sin_family: libc::AF_INET as libc::sa_family_t,
sin_port: v4.port().to_be(),
sin_addr: libc::in_addr {
s_addr: u32::from_ne_bytes(v4.ip().octets()),
},
sin_zero: [0; 8],
};
let len = std::mem::size_of::<libc::sockaddr_in>() as u32;
unsafe {
std::ptr::copy_nonoverlapping(
&sa as *const libc::sockaddr_in as *const u8,
&mut storage as *mut libc::sockaddr_storage as *mut u8,
len as usize,
);
}
(storage, len)
}
SocketAddr::V6(v6) => {
let sa = libc::sockaddr_in6 {
sin6_family: libc::AF_INET6 as libc::sa_family_t,
sin6_port: v6.port().to_be(),
sin6_flowinfo: v6.flowinfo(),
sin6_addr: libc::in6_addr {
s6_addr: v6.ip().octets(),
},
sin6_scope_id: v6.scope_id(),
};
let len = std::mem::size_of::<libc::sockaddr_in6>() as u32;
unsafe {
std::ptr::copy_nonoverlapping(
&sa as *const libc::sockaddr_in6 as *const u8,
&mut storage as *mut libc::sockaddr_storage as *mut u8,
len as usize,
);
}
(storage, len)
}
}
}
fn storage_to_addr(storage: &libc::sockaddr_storage) -> Option<SocketAddr> {
match storage.ss_family as i32 {
libc::AF_INET => {
let sa = unsafe { &*(storage as *const _ as *const libc::sockaddr_in) };
Some(SocketAddr::new(
std::net::IpAddr::V4(Ipv4Addr::from(u32::from_be(sa.sin_addr.s_addr))),
u16::from_be(sa.sin_port),
))
}
libc::AF_INET6 => {
let sa = unsafe { &*(storage as *const _ as *const libc::sockaddr_in6) };
Some(SocketAddr::new(
std::net::IpAddr::V6(Ipv6Addr::from(sa.sin6_addr.s6_addr)),
u16::from_be(sa.sin6_port),
))
}
_ => None,
}
}
#[allow(clippy::vec_box)]
pub struct UdpBatchIo {
ring: IoUring,
fd: RawFd,
recv_slots: Vec<Box<UdpSlot>>,
send_slots: Vec<Box<UdpSlot>>,
pending_send: u32,
early_recv: Vec<(Vec<u8>, SocketAddr)>,
}
impl std::fmt::Debug for UdpBatchIo {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("UdpBatchIo")
.field("fd", &self.fd)
.field("recv_slots", &self.recv_slots.len())
.field("send_slots", &self.send_slots.len())
.field("pending_send", &self.pending_send)
.finish()
}
}
impl UdpBatchIo {
pub fn new(fd: RawFd, recv_count: usize, send_count: usize, packet_size: usize) -> Result<Self> {
let entries = ((recv_count + send_count) as u32)
.checked_next_power_of_two()
.ok_or_else(|| LinuxError::InsufficientResources("SQ depth overflow".into()))?;
let ring = IoUring::new(entries).map_err(|e| LinuxError::Syscall {
syscall: "io_uring_setup",
errno: e.raw_os_error().unwrap_or(0),
})?;
let mut recv_slots = Vec::with_capacity(recv_count);
for _ in 0..recv_count {
recv_slots.push(UdpSlot::new(packet_size));
}
let mut send_slots = Vec::with_capacity(send_count);
for _ in 0..send_count {
send_slots.push(UdpSlot::new(packet_size));
}
Ok(Self {
ring,
fd,
recv_slots,
send_slots,
pending_send: 0,
early_recv: Vec::new(),
})
}
pub fn arm_recv(&mut self) -> Result<usize> {
let mut armed = 0usize;
for (idx, slot) in self.recv_slots.iter_mut().enumerate() {
if slot.in_use {
continue;
}
slot.iov = libc::iovec {
iov_base: slot.data.as_mut_ptr() as *mut libc::c_void,
iov_len: slot.data.len(),
};
slot.msg = libc::msghdr {
msg_name: &mut slot.addr as *mut libc::sockaddr_storage as *mut libc::c_void,
msg_namelen: std::mem::size_of::<libc::sockaddr_storage>() as u32,
msg_iov: &mut slot.iov as *mut libc::iovec,
msg_iovlen: 1,
msg_control: std::ptr::null_mut(),
msg_controllen: 0,
msg_flags: 0,
};
let entry = opcode::RecvMsg::new(types::Fd(self.fd), &mut slot.msg as *mut libc::msghdr)
.build()
.user_data(idx as u64); let mut sq = self.ring.submission();
unsafe {
sq.push(&entry).map_err(|_| {
LinuxError::InsufficientResources("io_uring SQ full (recv)".into())
})?;
}
slot.in_use = true;
armed += 1;
}
Ok(armed)
}
pub fn submit_and_wait(&mut self, min_complete: u32) -> Result<usize> {
self.ring
.submit_and_wait(min_complete as usize)
.map_err(|e| LinuxError::Syscall {
syscall: "io_uring_enter",
errno: e.raw_os_error().unwrap_or(0),
})
}
pub fn collect_recv(&mut self) -> Vec<(Vec<u8>, SocketAddr)> {
let mut out = std::mem::take(&mut self.early_recv);
let cq: Vec<(u64, i32)> = self
.ring
.completion()
.map(|cqe| (cqe.user_data(), cqe.result()))
.collect();
for (user_data, result) in cq {
if user_data & UD_SEND_FLAG != 0 {
let idx = (user_data & !UD_SEND_FLAG) as usize;
if let Some(slot) = self.send_slots.get_mut(idx) {
slot.in_use = false;
}
self.pending_send = self.pending_send.saturating_sub(1);
if result < 0 {
tracing::warn!(errno = -result, "io_uring sendmsg failed");
}
continue;
}
let idx = user_data as usize;
let Some(slot) = self.recv_slots.get_mut(idx) else {
continue;
};
slot.in_use = false;
if result <= 0 {
continue;
}
let n = (result as usize).min(slot.data.len());
let Some(from) = storage_to_addr(&slot.addr) else {
continue; };
out.push((slot.data[..n].to_vec(), from));
}
out
}
pub fn push_send(&mut self, data: &[u8], addr: &SocketAddr) -> Result<bool> {
let Some((idx, slot)) = self
.send_slots
.iter_mut()
.enumerate()
.find(|(_, s)| !s.in_use)
else {
return Ok(false);
};
if data.len() > slot.data.len() {
return Err(LinuxError::InsufficientResources(format!(
"send data {}B > slot buffer {}B",
data.len(),
slot.data.len()
)));
}
slot.data[..data.len()].copy_from_slice(data);
let (storage, addr_len) = addr_to_storage(addr);
slot.addr = storage;
slot.iov = libc::iovec {
iov_base: slot.data.as_mut_ptr() as *mut libc::c_void,
iov_len: data.len(),
};
slot.msg = libc::msghdr {
msg_name: &mut slot.addr as *mut libc::sockaddr_storage as *mut libc::c_void,
msg_namelen: addr_len,
msg_iov: &mut slot.iov as *mut libc::iovec,
msg_iovlen: 1,
msg_control: std::ptr::null_mut(),
msg_controllen: 0,
msg_flags: 0,
};
let entry = opcode::SendMsg::new(types::Fd(self.fd), &slot.msg as *const libc::msghdr)
.build()
.user_data(UD_SEND_FLAG | idx as u64);
let mut sq = self.ring.submission();
unsafe {
sq.push(&entry).map_err(|_| {
LinuxError::InsufficientResources("io_uring SQ full (send)".into())
})?;
}
slot.in_use = true;
self.pending_send = self.pending_send.saturating_add(1);
Ok(true)
}
pub fn flush_send(&mut self) -> Result<usize> {
let mut sent = 0usize;
while self.pending_send > 0 {
self.ring
.submit_and_wait(1)
.map_err(|e| LinuxError::Syscall {
syscall: "io_uring_enter",
errno: e.raw_os_error().unwrap_or(0),
})?;
let cq: Vec<(u64, i32)> = self
.ring
.completion()
.map(|cqe| (cqe.user_data(), cqe.result()))
.collect();
for (user_data, result) in cq {
if user_data & UD_SEND_FLAG != 0 {
let idx = (user_data & !UD_SEND_FLAG) as usize;
if let Some(slot) = self.send_slots.get_mut(idx) {
slot.in_use = false;
}
self.pending_send = self.pending_send.saturating_sub(1);
if result >= 0 {
sent += 1;
} else {
tracing::warn!(errno = -result, "io_uring sendmsg failed");
}
} else {
let idx = user_data as usize;
if let Some(slot) = self.recv_slots.get_mut(idx) {
slot.in_use = false;
if result > 0 {
let n = (result as usize).min(slot.data.len());
if let Some(from) = storage_to_addr(&slot.addr) {
self.early_recv.push((slot.data[..n].to_vec(), from));
}
}
}
}
}
}
Ok(sent)
}
pub fn send_slot_free(&self) -> usize {
self.send_slots.iter().filter(|s| !s.in_use).count()
}
pub fn early_recv_empty(&self) -> bool {
self.early_recv.is_empty()
}
}