use std::io;
use crate::{ring::XskRingCons, umem::frame::FrameDesc};
use super::{Socket, fd::Fd};
#[derive(Debug)]
pub struct RxQueue {
ring: XskRingCons,
socket: Socket,
}
impl RxQueue {
pub(super) fn new(ring: XskRingCons, socket: Socket) -> Self {
Self { ring, socket }
}
#[inline]
pub unsafe fn consume(&mut self, descs: &mut [FrameDesc]) -> usize {
let nb = descs.len() as u32;
if nb == 0 {
return 0;
}
let mut idx = 0;
let cnt = unsafe { libxdp_sys::xsk_ring_cons__peek(self.ring.as_ptr(), nb, &mut idx) };
if cnt > 0 {
for desc in descs.iter_mut().take(cnt as usize) {
let recv_pkt_desc =
unsafe { libxdp_sys::xsk_ring_cons__rx_desc(self.ring.as_ptr(), idx) };
unsafe {
desc.addr = (*recv_pkt_desc).addr as usize;
desc.lengths.data = (*recv_pkt_desc).len as usize;
desc.lengths.headroom = 0;
desc.options = (*recv_pkt_desc).options;
}
idx = idx.wrapping_add(1);
}
unsafe { libxdp_sys::xsk_ring_cons__release(self.ring.as_ptr(), cnt) };
}
cnt as usize
}
#[inline]
pub unsafe fn consume_one(&mut self, desc: &mut FrameDesc) -> usize {
let mut idx = 0;
let cnt = unsafe { libxdp_sys::xsk_ring_cons__peek(self.ring.as_ptr(), 1, &mut idx) };
if cnt > 0 {
let recv_pkt_desc =
unsafe { libxdp_sys::xsk_ring_cons__rx_desc(self.ring.as_ptr(), idx) };
unsafe {
desc.addr = (*recv_pkt_desc).addr as usize;
desc.lengths.data = (*recv_pkt_desc).len as usize;
desc.lengths.headroom = 0;
desc.options = (*recv_pkt_desc).options;
}
unsafe { libxdp_sys::xsk_ring_cons__release(self.ring.as_ptr(), cnt) };
}
cnt as usize
}
#[inline]
pub unsafe fn poll_and_consume(
&mut self,
descs: &mut [FrameDesc],
poll_timeout: i32,
) -> io::Result<usize> {
match self.poll(poll_timeout)? {
true => Ok(unsafe { self.consume(descs) }),
false => Ok(0),
}
}
#[inline]
pub unsafe fn poll_and_consume_one(
&mut self,
desc: &mut FrameDesc,
poll_timeout: i32,
) -> io::Result<usize> {
match self.poll(poll_timeout)? {
true => Ok(unsafe { self.consume_one(desc) }),
false => Ok(0),
}
}
#[inline]
pub fn poll(&mut self, poll_timeout: i32) -> io::Result<bool> {
self.socket.fd.poll_read(poll_timeout)
}
#[inline]
pub fn fd(&self) -> &Fd {
&self.socket.fd
}
#[inline]
pub fn fd_mut(&mut self) -> &mut Fd {
&mut self.socket.fd
}
}