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),
}
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)),
}),
}
}
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));
}
}
}
}
Ok(last_ts)
}
pub(crate) fn cumulative_stats(&self) -> Result<CaptureStats> {
match self {
AnyBackend::AfPacket(cap) => cap.cumulative_stats(),
#[cfg(feature = "af-xdp")]
AnyBackend::Xdp(xdp) => {
let s = xdp.statistics()?;
Ok(CaptureStats {
packets: 0,
drops: s
.rx_dropped
.saturating_add(s.rx_ring_full)
.saturating_add(s.rx_fill_ring_empty_descs)
.min(u32::MAX as u64) as u32,
freeze_count: 0,
})
}
}
}
}
#[cfg(feature = "af-xdp")]
fn now_ts() -> Timestamp {
Timestamp::from_system_time(std::time::SystemTime::now())
}