use std::sync::Arc;
use std::time::{Duration, Instant};
use crate::error::{Error, Result};
use crate::monitor::{Monitor, MonitorBuilder};
use crate::xdp::{Queues, XdpCapture, XdpFlags};
pub struct XdpShardedRunner {
iface: String,
queues: Queues,
promiscuous: bool,
busy_poll_us: Option<u32>,
pin_cpus: bool,
numa_auto: bool,
attach_flags: XdpFlags,
build_shard: Arc<dyn Fn(usize, MonitorBuilder) -> MonitorBuilder + Send + Sync + 'static>,
}
impl XdpShardedRunner {
pub fn new<F>(iface: impl Into<String>, queues: Queues, build_shard: F) -> Self
where
F: Fn(usize, MonitorBuilder) -> MonitorBuilder + Send + Sync + 'static,
{
Self {
iface: iface.into(),
queues,
promiscuous: false,
busy_poll_us: None,
pin_cpus: false,
numa_auto: false,
attach_flags: XdpFlags::SKB_MODE,
build_shard: Arc::new(build_shard),
}
}
pub fn numa_auto(mut self, on: bool) -> Self {
self.numa_auto = on;
self
}
pub fn promiscuous(mut self, enable: bool) -> Self {
self.promiscuous = enable;
self
}
pub fn busy_poll(mut self, us: u32) -> Self {
self.busy_poll_us = Some(us);
self
}
pub fn pin_cpus(mut self, on: bool) -> Self {
self.pin_cpus = on;
self
}
pub fn attach_flags(mut self, flags: XdpFlags) -> Self {
self.attach_flags = flags;
self
}
pub fn run_until(self, deadline: Instant) -> Result<()> {
let mut b = XdpCapture::builder()
.interface(&self.iface)
.queues(self.queues.clone())
.promiscuous(self.promiscuous)
.attach_flags(self.attach_flags);
if let Some(us) = self.busy_poll_us {
b = b.busy_poll(us).prefer_busy_poll(true);
}
if self.numa_auto {
b = b.numa_auto();
}
let capture = b.build()?;
let n = capture.socket_count();
let (sockets, _guard) = capture.into_parts();
let build_shard = self.build_shard;
let pin_cpus = self.pin_cpus;
let mut handles = Vec::with_capacity(n);
for (i, socket) in sockets.into_iter().enumerate() {
let build = Arc::clone(&build_shard);
let handle = std::thread::Builder::new()
.name(format!("netring-xdp-shard-{i}"))
.spawn(move || -> Result<()> {
if pin_cpus && !crate::monitor::shard::pin_current_thread_to_core(i) {
tracing::warn!(shard = i, "could not set CPU affinity for xdp shard");
}
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(Error::Io)?;
rt.block_on(async move {
let async_sock = crate::AsyncXdpSocket::new(socket)?;
let monitor = build(i, Monitor::builder())
.inject_xdp_backend(async_sock)
.build()?;
let dur = deadline.saturating_duration_since(Instant::now());
monitor.run_for(dur).await
})
})
.map_err(Error::Io)?;
handles.push(handle);
}
let mut first_err: Option<Error> = None;
for h in handles {
match h.join() {
Ok(Ok(())) => {}
Ok(Err(e)) => {
if first_err.is_none() {
first_err = Some(e);
}
}
Err(panic) => {
if first_err.is_none() {
first_err = Some(Error::Io(std::io::Error::other(format!(
"xdp shard thread panicked: {panic:?}"
))));
}
}
}
}
drop(_guard);
match first_err {
Some(e) => Err(e),
None => Ok(()),
}
}
pub fn run_for(self, duration: Duration) -> Result<()> {
let deadline = Instant::now() + duration;
self.run_until(deadline)
}
pub fn interface(&self) -> &str {
&self.iface
}
}
impl std::fmt::Debug for XdpShardedRunner {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("XdpShardedRunner")
.field("iface", &self.iface)
.field("queues", &self.queues)
.field("promiscuous", &self.promiscuous)
.field("busy_poll_us", &self.busy_poll_us)
.field("pin_cpus", &self.pin_cpus)
.finish()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn builder_defaults_and_setters() {
let r = XdpShardedRunner::new("eth0", Queues::Auto, |_q, b| b)
.promiscuous(true)
.busy_poll(50)
.pin_cpus(true);
assert_eq!(r.interface(), "eth0");
assert!(r.promiscuous);
assert_eq!(r.busy_poll_us, Some(50));
assert!(r.pin_cpus);
assert!(matches!(r.queues, Queues::Auto));
}
}