use std::task::{Context, Poll};
use flowscope::{PacketView, Timestamp};
use crate::AsyncCapture;
use crate::error::{Error, Result};
use crate::stats::CaptureStats;
#[allow(clippy::large_enum_variant)]
pub(crate) enum AnyBackend {
AfPacket(AsyncCapture<crate::Capture>),
#[cfg(feature = "af-xdp")]
Xdp(crate::AsyncXdpSocket),
#[cfg(all(feature = "af-xdp", feature = "xdp-loader"))]
XdpMq(crate::AsyncXdpCapture),
}
impl AnyBackend {
pub(crate) fn poll_read_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<()>> {
match self {
AnyBackend::AfPacket(cap) => cap.poll_read_ready_mut(cx).map(|r| match r {
Ok(_guard) => Ok(()),
Err(e) => Err(Error::Io(e)),
}),
#[cfg(feature = "af-xdp")]
AnyBackend::Xdp(xdp) => xdp.poll_read_ready_mut(cx).map(|r| match r {
Ok(_guard) => Ok(()),
Err(e) => Err(Error::Io(e)),
}),
#[cfg(all(feature = "af-xdp", feature = "xdp-loader"))]
AnyBackend::XdpMq(cap) => cap.poll_read_ready(cx),
}
}
pub(crate) async fn drain_batch(
&mut self,
mut on_packet: impl FnMut(PacketView<'_>),
) -> Result<Option<Timestamp>> {
let mut last_ts: Option<Timestamp> = None;
match self {
AnyBackend::AfPacket(cap) => {
let mut guard = cap.readable().await?;
while let Some(batch) = guard.next_batch() {
for pkt in &batch {
let ts = pkt.timestamp();
last_ts = Some(ts);
on_packet(PacketView::new(pkt.data(), ts));
}
}
}
#[cfg(feature = "af-xdp")]
AnyBackend::Xdp(xdp) => {
let mut guard = xdp.readable().await?;
while let Some(batch) = guard.next_batch() {
for pkt in &batch {
let ts = pkt.timestamp().unwrap_or_else(now_ts);
last_ts = Some(ts);
on_packet(PacketView::new(pkt.data(), ts));
}
}
}
#[cfg(all(feature = "af-xdp", feature = "xdp-loader"))]
AnyBackend::XdpMq(cap) => {
let n = cap.socket_count();
let start = cap.next_cursor();
for off in 0..n {
let i = (start + off) % n;
if !cap.socket_rx_ready(i) {
continue;
}
let mut guard = cap.socket_readable(i).await?;
while let Some(batch) = guard.next_batch() {
for pkt in &batch {
let ts = pkt.timestamp().unwrap_or_else(now_ts);
last_ts = Some(ts);
on_packet(PacketView::new(pkt.data(), ts));
}
}
}
}
}
Ok(last_ts)
}
pub(crate) fn detailed_stats(&self) -> Result<(CaptureStats, crate::stats::DropBreakdown)> {
match self {
AnyBackend::AfPacket(cap) => {
let stats = cap.cumulative_stats()?;
let detail = crate::stats::DropBreakdown::AfPacket {
freezes: stats.freeze_count as u64,
};
Ok((stats, detail))
}
#[cfg(feature = "af-xdp")]
AnyBackend::Xdp(xdp) => {
let s = xdp.statistics()?;
Ok((s.to_capture_stats(), s.into()))
}
#[cfg(all(feature = "af-xdp", feature = "xdp-loader"))]
AnyBackend::XdpMq(cap) => cap.detailed_stats(),
}
}
}
#[cfg(feature = "af-xdp")]
fn now_ts() -> Timestamp {
Timestamp::from_system_time(std::time::SystemTime::now())
}