use crate::afxdp::ffi;
use crate::error::Error;
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum Queues {
Single(u32),
Range(std::ops::Range<u32>),
Auto,
}
impl Default for Queues {
fn default() -> Self {
Queues::Single(0)
}
}
impl Queues {
pub fn single(q: u32) -> Self {
Queues::Single(q)
}
pub fn range(r: std::ops::Range<u32>) -> Self {
Queues::Range(r)
}
pub fn resolve(&self, iface: &str) -> Result<Vec<u32>, Error> {
match self {
Queues::Single(q) => Ok(vec![*q]),
Queues::Range(r) => {
if r.start >= r.end {
return Err(Error::Config(format!(
"empty queue range {}..{}",
r.start, r.end
)));
}
Ok(r.clone().collect())
}
Queues::Auto => match queue_count(iface) {
Ok(n) => Ok((0..n).collect()),
Err(e) => {
tracing::warn!(
"Queues::Auto: queue_count({iface}) failed ({e}); \
falling back to queue 0 only"
);
Ok(vec![0])
}
},
}
}
}
pub fn queue_count(iface: &str) -> Result<u32, Error> {
use std::os::fd::{AsRawFd, FromRawFd, OwnedFd};
if iface.len() >= libc::IFNAMSIZ {
return Err(Error::Config(format!("interface name too long: {iface}")));
}
let raw = unsafe { libc::socket(libc::AF_INET, libc::SOCK_DGRAM, 0) };
if raw < 0 {
return Err(Error::Socket(std::io::Error::last_os_error()));
}
let fd = unsafe { OwnedFd::from_raw_fd(raw) };
let mut channels = ffi::ethtool_channels {
cmd: ffi::ETHTOOL_GCHANNELS,
..Default::default()
};
let mut ifr = ffi::ethtool_ifreq {
ifr_name: [0; libc::IFNAMSIZ],
ifr_data: (&mut channels as *mut ffi::ethtool_channels).cast(),
};
for (dst, &b) in ifr.ifr_name.iter_mut().zip(iface.as_bytes()) {
*dst = b as libc::c_char;
}
let rc = unsafe { libc::ioctl(fd.as_raw_fd(), ffi::SIOCETHTOOL, &mut ifr as *mut _) };
if rc != 0 {
return Err(Error::Io(std::io::Error::last_os_error()));
}
let n = if channels.combined_count > 0 {
channels.combined_count
} else if channels.rx_count > 0 {
channels.rx_count
} else {
1
};
Ok(n)
}
pub fn interface_numa_node(iface: &str) -> Option<u32> {
let path = format!("/sys/class/net/{iface}/device/numa_node");
let raw = std::fs::read_to_string(path).ok()?;
let node = raw.trim().parse::<i32>().ok()?;
(node >= 0).then_some(node as u32)
}
#[cfg(feature = "xdp-loader")]
pub use mq::{XdpCapture, XdpCaptureBuilder, XdpCaptureGuard};
#[cfg(feature = "xdp-loader")]
mod mq {
use std::time::Duration;
use super::{Error, Queues};
use crate::afxdp::loader::{XdpAttachment, XdpFlags, XdpProgram, default_program};
use crate::afxdp::{XdpBatch, XdpMode, XdpSocket, XdpSocketBuilder};
#[must_use]
pub struct XdpCaptureBuilder {
interface: Option<String>,
queues: Queues,
frame_size: usize,
frame_count: usize,
mode: XdpMode,
promiscuous: bool,
hugepages: bool,
numa_node: Option<u32>,
numa_auto: bool,
busy_poll_us: Option<u32>,
prefer_busy_poll: Option<bool>,
busy_poll_budget: Option<u16>,
attach_flags: XdpFlags,
program: Option<XdpProgram>,
}
impl Default for XdpCaptureBuilder {
fn default() -> Self {
Self {
interface: None,
queues: Queues::default(),
frame_size: 4096,
frame_count: 4096,
mode: XdpMode::Rx, promiscuous: false,
hugepages: false,
numa_node: None,
numa_auto: false,
busy_poll_us: None,
prefer_busy_poll: None,
busy_poll_budget: None,
attach_flags: XdpFlags::SKB_MODE, program: None,
}
}
}
impl XdpCaptureBuilder {
pub fn interface(mut self, name: &str) -> Self {
self.interface = Some(name.to_string());
self
}
pub fn queues(mut self, queues: Queues) -> Self {
self.queues = queues;
self
}
pub fn frame_size(mut self, size: usize) -> Self {
self.frame_size = size;
self
}
pub fn frame_count(mut self, count: usize) -> Self {
self.frame_count = count;
self
}
pub fn mode(mut self, mode: XdpMode) -> Self {
self.mode = mode;
self
}
pub fn promiscuous(mut self, enable: bool) -> Self {
self.promiscuous = enable;
self
}
pub fn hugepages(mut self, enable: bool) -> Self {
self.hugepages = enable;
self
}
pub fn numa_node(mut self, node: u32) -> Self {
self.numa_node = Some(node);
self
}
pub fn numa_auto(mut self) -> Self {
self.numa_auto = true;
self
}
pub fn busy_poll(mut self, us: u32) -> Self {
self.busy_poll_us = Some(us);
self
}
pub fn prefer_busy_poll(mut self, enable: bool) -> Self {
self.prefer_busy_poll = Some(enable);
self
}
pub fn busy_poll_budget(mut self, budget: u16) -> Self {
self.busy_poll_budget = Some(budget);
self
}
pub fn attach_flags(mut self, flags: XdpFlags) -> Self {
self.attach_flags = flags;
self
}
pub fn with_program(mut self, prog: XdpProgram) -> Self {
self.program = Some(prog);
self
}
pub fn build(mut self) -> Result<XdpCapture, Error> {
let iface = self
.interface
.clone()
.ok_or_else(|| Error::Config("interface is required".into()))?;
let ifindex = crate::afpacket::socket::resolve_interface(&iface)? as u32;
let queue_ids = self.queues.resolve(&iface)?;
let numa_node = self.numa_node.or_else(|| {
if self.numa_auto {
super::interface_numa_node(&iface)
} else {
None
}
});
let map_size = queue_ids.iter().max().map(|&m| m + 1).unwrap_or(1);
let mut prog = match self.program.take() {
Some(p) => p,
None => default_program(map_size)?,
};
let mut sockets = Vec::with_capacity(queue_ids.len());
for &q in &queue_ids {
let mut b = XdpSocketBuilder::default()
.interface(&iface)
.queue_id(q)
.mode(self.mode)
.frame_size(self.frame_size)
.frame_count(self.frame_count)
.hugepages(self.hugepages);
if let Some(node) = numa_node {
b = b.numa_node(node);
}
if let Some(us) = self.busy_poll_us {
b = b.busy_poll_us(us);
}
if let Some(prefer) = self.prefer_busy_poll {
b = b.prefer_busy_poll(prefer);
}
if let Some(budget) = self.busy_poll_budget {
b = b.busy_poll_budget(budget);
}
let sock = b.build()?;
prog.register(q, &sock)?;
sockets.push(sock);
}
let attachment = prog.attach(&iface, self.attach_flags)?;
let promisc = if self.promiscuous {
Some(crate::afpacket::socket::PromiscGuard::enable(
ifindex as i32,
)?)
} else {
None
};
Ok(XdpCapture {
sockets,
queue_ids,
cursor: 0,
_attachment: attachment,
_promisc: promisc,
})
}
}
pub struct XdpCapture {
sockets: Vec<XdpSocket>,
queue_ids: Vec<u32>,
cursor: usize,
_attachment: XdpAttachment,
_promisc: Option<crate::afpacket::socket::PromiscGuard>,
}
impl XdpCapture {
pub fn builder() -> XdpCaptureBuilder {
XdpCaptureBuilder::default()
}
pub fn open(iface: &str) -> Result<Self, Error> {
Self::builder()
.interface(iface)
.queues(Queues::Auto)
.promiscuous(true)
.build()
}
pub fn queue_ids(&self) -> &[u32] {
&self.queue_ids
}
pub fn socket_count(&self) -> usize {
self.sockets.len()
}
pub fn is_zerocopy(&self) -> bool {
self.sockets.iter().all(|s| s.is_zerocopy())
}
pub fn sockets_mut(&mut self) -> &mut [XdpSocket] {
&mut self.sockets
}
pub fn next_batch(&mut self) -> Option<(u32, XdpBatch<'_>)> {
let n = self.sockets.len();
if n == 0 {
return None;
}
let mut target = None;
for off in 0..n {
let i = (self.cursor + off) % n;
if self.sockets[i].rx_poll_ready() {
target = Some(i);
break;
}
}
let i = target?;
self.cursor = (i + 1) % n;
let qid = self.queue_ids[i];
self.sockets[i].next_batch().map(|b| (qid, b))
}
pub fn next_batch_blocking(
&mut self,
timeout: Duration,
) -> Result<Option<(u32, XdpBatch<'_>)>, Error> {
if self.sockets.iter_mut().any(|s| s.rx_poll_ready()) {
return Ok(self.next_batch());
}
let mut pfds: Vec<nix::poll::PollFd> = self
.sockets
.iter()
.map(|s| nix::poll::PollFd::new(s.poll_fd(), nix::poll::PollFlags::POLLIN))
.collect();
crate::syscall::poll_eintr_safe(&mut pfds, timeout).map_err(Error::Io)?;
Ok(self.next_batch())
}
pub fn into_parts(self) -> (Vec<XdpSocket>, XdpCaptureGuard) {
(
self.sockets,
XdpCaptureGuard {
_attachment: self._attachment,
_promisc: self._promisc,
},
)
}
}
pub struct XdpCaptureGuard {
_attachment: XdpAttachment,
_promisc: Option<crate::afpacket::socket::PromiscGuard>,
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn queues_default_is_single_zero() {
assert!(matches!(Queues::default(), Queues::Single(0)));
}
#[test]
fn queues_resolve_single_and_range() {
assert_eq!(Queues::Single(2).resolve("lo").unwrap(), vec![2]);
assert_eq!(Queues::range(0..4).resolve("lo").unwrap(), vec![0, 1, 2, 3]);
}
#[test]
fn queues_resolve_empty_range_errors() {
assert!(Queues::range(3..3).resolve("lo").is_err());
}
#[test]
fn queues_auto_never_errors_falls_back_to_zero() {
let resolved = Queues::Auto.resolve("lo").unwrap();
assert!(resolved.contains(&0));
}
#[test]
fn interface_numa_node_is_none_for_lo() {
assert_eq!(interface_numa_node("lo"), None);
assert_eq!(interface_numa_node("netring-no-such-if0"), None);
}
#[test]
fn queue_count_does_not_panic() {
if let Ok(n) = queue_count("lo") {
assert!(n >= 1);
}
}
}