use libc::{EAGAIN, EBUSY, ENETDOWN, ENOBUFS, MSG_DONTWAIT};
use std::{io, os::unix::prelude::AsRawFd, ptr};
use crate::{ring::XskRingProd, umem::frame::FrameDesc, util};
use super::{Socket, fd::Fd};
#[derive(Debug)]
pub struct TxQueue {
ring: XskRingProd,
socket: Socket,
}
impl TxQueue {
pub(super) fn new(ring: XskRingProd, socket: Socket) -> Self {
Self { ring, socket }
}
#[inline]
pub unsafe fn produce(&mut self, descs: &[FrameDesc]) -> usize {
let nb = descs.len() as u32;
if nb == 0 {
return 0;
}
let mut idx = 0;
let cnt = unsafe { libxdp_sys::xsk_ring_prod__reserve(self.ring.as_ptr(), nb, &mut idx) };
if cnt > 0 {
for desc in descs.iter().take(cnt as usize) {
let send_pkt_desc =
unsafe { libxdp_sys::xsk_ring_prod__tx_desc(self.ring.as_ptr(), idx) };
unsafe { desc.write_xdp_desc(&mut *send_pkt_desc) };
idx = idx.wrapping_add(1);
}
unsafe { libxdp_sys::xsk_ring_prod__submit(self.ring.as_ptr(), cnt) };
}
cnt as usize
}
#[inline]
pub unsafe fn produce_one(&mut self, desc: &FrameDesc) -> usize {
let mut idx = 0;
let cnt = unsafe { libxdp_sys::xsk_ring_prod__reserve(self.ring.as_ptr(), 1, &mut idx) };
if cnt > 0 {
let send_pkt_desc =
unsafe { libxdp_sys::xsk_ring_prod__tx_desc(self.ring.as_ptr(), idx) };
unsafe { desc.write_xdp_desc(&mut *send_pkt_desc) };
unsafe { libxdp_sys::xsk_ring_prod__submit(self.ring.as_ptr(), cnt) };
}
cnt as usize
}
#[inline]
pub unsafe fn produce_and_wakeup(&mut self, descs: &[FrameDesc]) -> io::Result<usize> {
let cnt = unsafe { self.produce(descs) };
if self.needs_wakeup() {
self.wakeup()?;
}
Ok(cnt)
}
#[inline]
pub unsafe fn produce_one_and_wakeup(&mut self, desc: &FrameDesc) -> io::Result<usize> {
let cnt = unsafe { self.produce_one(desc) };
if self.needs_wakeup() {
self.wakeup()?;
}
Ok(cnt)
}
#[inline]
pub fn wakeup(&self) -> io::Result<()> {
let ret = unsafe {
libc::sendto(
self.socket.fd.as_raw_fd(),
ptr::null(),
0,
MSG_DONTWAIT,
ptr::null(),
0,
)
};
if ret < 0 {
match util::get_errno() {
ENOBUFS | EAGAIN | EBUSY | ENETDOWN => (),
_ => return Err(io::Error::last_os_error()),
}
}
Ok(())
}
#[inline]
pub fn needs_wakeup(&self) -> bool {
unsafe { libxdp_sys::xsk_ring_prod__needs_wakeup(self.ring.as_ptr()) != 0 }
}
#[inline]
pub fn poll(&mut self, poll_timeout: i32) -> io::Result<bool> {
self.socket.fd.poll_write(poll_timeout)
}
#[inline]
pub fn fd(&self) -> &Fd {
&self.socket.fd
}
#[inline]
pub fn fd_mut(&mut self) -> &mut Fd {
&mut self.socket.fd
}
}