use std::task::{Context, Poll};
use flowscope::PacketView;
use crate::packet::PacketDirection;
#[cfg(all(feature = "af-xdp", feature = "xdp-loader"))]
use crate::packet::Timestamp;
pub(crate) struct SourcePacket<'a> {
pub view: PacketView<'a>,
pub data: &'a [u8],
pub direction: PacketDirection,
#[cfg_attr(not(feature = "pcap"), allow(dead_code))]
pub original_len: usize,
}
pub(crate) enum DrainOutcome {
Drained,
Idle,
}
pub(crate) trait AsyncFlowSource {
fn poll_drain(
&mut self,
cx: &mut Context<'_>,
sink: &mut dyn FnMut(SourcePacket<'_>),
) -> Poll<std::io::Result<DrainOutcome>>;
}
use std::os::unix::io::AsRawFd;
use crate::async_adapters::tokio_adapter::AsyncCapture;
use crate::traits::PacketSource;
impl<S> AsyncFlowSource for AsyncCapture<S>
where
S: PacketSource + AsRawFd,
{
fn poll_drain(
&mut self,
cx: &mut Context<'_>,
sink: &mut dyn FnMut(SourcePacket<'_>),
) -> Poll<std::io::Result<DrainOutcome>> {
let mut guard = match self.poll_read_ready_mut(cx) {
Poll::Ready(Ok(g)) => g,
Poll::Ready(Err(e)) => return Poll::Ready(Err(e)),
Poll::Pending => return Poll::Pending,
};
let inner = guard.get_inner_mut();
if let Some(batch) = inner.next_batch() {
for pkt in &batch {
sink(SourcePacket {
view: pkt.view(),
data: pkt.data(),
direction: pkt.direction(),
original_len: pkt.original_len(),
});
}
drop(batch);
Poll::Ready(Ok(DrainOutcome::Drained))
} else {
guard.clear_ready();
Poll::Ready(Ok(DrainOutcome::Idle))
}
}
}
#[cfg(all(feature = "af-xdp", feature = "xdp-loader"))]
#[inline]
pub(crate) fn view_from_parts(
data: &[u8],
ts: Option<Timestamp>,
rx_metadata: flowscope::RxMetadata,
) -> PacketView<'_> {
let ts = ts.unwrap_or_else(crate::async_adapters::flow_stream::current_timestamp);
PacketView::new(data, ts).with_rx_metadata(rx_metadata)
}
#[cfg(all(feature = "af-xdp", feature = "xdp-loader"))]
impl AsyncFlowSource for crate::AsyncXdpCapture {
fn poll_drain(
&mut self,
cx: &mut Context<'_>,
sink: &mut dyn FnMut(SourcePacket<'_>),
) -> Poll<std::io::Result<DrainOutcome>> {
match self.poll_drain_views(cx, sink) {
Poll::Ready(Ok(outcome)) => Poll::Ready(Ok(outcome)),
Poll::Ready(Err(e)) => Poll::Ready(Err(match e {
crate::error::Error::Io(io) => io,
other => std::io::Error::other(other),
})),
Poll::Pending => Poll::Pending,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[allow(dead_code)]
fn _assert_impls() {
fn is_flow_source<T: AsyncFlowSource>() {}
fn is_stream<T: futures_core::Stream>() {}
type Ext = flowscope::extract::FiveTuple;
is_flow_source::<crate::AsyncCapture<crate::Capture>>();
is_stream::<
crate::async_adapters::flow_stream::FlowStream<
crate::AsyncCapture<crate::Capture>,
Ext,
>,
>();
#[cfg(all(feature = "af-xdp", feature = "xdp-loader"))]
{
use flowscope::{DatagramParser, FlowSide, SessionParser, Timestamp};
#[derive(Default, Clone)]
struct SParser;
impl SessionParser for SParser {
type Message = ();
fn feed_initiator(&mut self, _: &[u8], _: Timestamp, _: &mut Vec<()>) {}
fn feed_responder(&mut self, _: &[u8], _: Timestamp, _: &mut Vec<()>) {}
}
#[derive(Default, Clone)]
struct DParser;
impl DatagramParser for DParser {
type Message = ();
fn parse(&mut self, _: &[u8], _: FlowSide, _: Timestamp, _: &mut Vec<()>) {}
}
is_flow_source::<crate::AsyncXdpCapture>();
is_stream::<crate::async_adapters::flow_stream::FlowStream<crate::AsyncXdpCapture, Ext>>(
);
is_stream::<crate::async_adapters::multi_streams::XdpMultiFlowStream<Ext>>();
is_stream::<
crate::async_adapters::multi_streams::MergedFlowStream<crate::AsyncXdpCapture, Ext>,
>();
is_stream::<
crate::async_adapters::session_stream::SessionStream<
crate::AsyncXdpCapture,
Ext,
SParser,
>,
>();
is_stream::<
crate::async_adapters::datagram_stream::DatagramStream<
crate::AsyncXdpCapture,
Ext,
DParser,
>,
>();
is_stream::<crate::async_adapters::multi_streams::XdpMultiSessionStream<Ext, SParser>>(
);
is_stream::<crate::async_adapters::multi_streams::XdpMultiDatagramStream<Ext, DParser>>(
);
}
}
}