use std::pin::Pin;
use std::task::{Context, Poll};
use std::time::{Duration, Instant};
use flowscope::L4Proto;
use flowscope::driver::Event as FsEvent;
use flowscope::extract::FiveTuple;
use futures_core::Stream;
use crate::AsyncCapture;
use crate::OwnedPacket;
use crate::anomaly::sink::AnomalySink;
use crate::async_adapters::tokio_adapter::PacketStream;
use crate::ctx::{CounterRegistry, Ctx, SourceIdx, StateMap};
use crate::error::Result;
use crate::monitor::Monitor;
use crate::monitor::dispatcher::Dispatcher;
use crate::protocol::FlowKey;
#[cfg(feature = "icmp")]
use crate::protocol::builtin::Icmp;
use crate::protocol::builtin::{Tcp, Udp};
use crate::protocol::event_typed::{
AnyFlowAnomaly, FlowEnded, FlowEstablished, FlowPacket, FlowStarted, FlowTick, ParserClosed,
TcpRst, Tick,
};
use std::time::SystemTime;
pub(crate) enum StopCondition {
Deadline(Instant),
Signal,
Idle(Duration),
}
pub(crate) async fn run_loop(monitor: Monitor, stop: StopCondition) -> Result<()> {
let Monitor {
interfaces,
mut driver,
mut dispatcher,
mut protocol_slots,
mut state_map,
mut counters,
mut sink,
mut tick_handlers,
detector_names: _,
monitor_name,
drain_timeout,
broadcast_handles: _,
#[cfg(all(feature = "pcap", feature = "tokio"))]
pcap_source_path: _,
#[cfg(all(feature = "pcap", feature = "tokio"))]
pcap_speed_factor: _,
mut flow_states,
fanout,
label_table,
mut merge_rx,
} = monitor;
let monitor_name_borrow: Option<&str> = monitor_name.as_deref();
let mut streams: Vec<PacketStream<_>> = Vec::with_capacity(interfaces.len());
for iface in &interfaces {
let cap = match fanout {
Some((mode, group_id)) => {
let rx = crate::Capture::builder()
.interface(iface)
.fanout(mode, group_id)
.build()?;
AsyncCapture::new(rx)?
}
None => AsyncCapture::open(iface)?,
};
streams.push(cap.into_stream());
}
let mut events: Vec<FsEvent<FlowKey>> = Vec::with_capacity(64);
let mut shutdown = ShutdownSignal::new(stop);
let mut rr_anchor: usize = 0;
let mut last_event_at = Instant::now();
let mut tick_intervals: Vec<tokio::time::Interval> = tick_handlers
.iter()
.map(|t| {
let mut int =
tokio::time::interval_at(tokio::time::Instant::now() + t.period, t.period);
int.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
int
})
.collect();
loop {
let next = tokio::select! {
biased;
_ = shutdown.recv(last_event_at) => break,
packet = next_packet(&mut streams, &mut rr_anchor) => packet,
tick_idx = next_tick(&mut tick_intervals), if !tick_intervals.is_empty() => {
last_event_at = Instant::now();
fire_tick(
tick_idx,
&mut tick_handlers,
&mut dispatcher,
sink.as_mut(),
&mut state_map,
&mut counters,
monitor_name_borrow,
&mut flow_states,
&label_table,
)
.await?;
continue;
}
req = recv_merge(&mut merge_rx), if merge_rx.is_some() => {
if let Some(req) = req {
let taken = state_map.take_dyn(req.type_id);
let _ = req.reply.send(taken);
}
continue;
}
};
let (source_idx, batch) = match next {
Some((i, Ok(b))) => (i, b),
Some((_, Err(e))) => return Err(e),
None => break, };
let source = SourceIdx(source_idx as u8);
last_event_at = Instant::now();
for pkt in batch {
let view = flowscope::PacketView::new(&pkt.data, pkt.timestamp);
events.clear();
driver.track_into(view, &mut events);
for evt in events.drain(..) {
dispatch_lifecycle(
&mut dispatcher,
sink.as_mut(),
&mut state_map,
&mut counters,
evt.clone(),
source,
monitor_name_borrow,
&mut flow_states,
&label_table,
)?;
dispatch_lifecycle_async(&mut dispatcher, evt).await?;
}
for slot in &mut protocol_slots {
let mut ctx = Ctx::new(
None,
pkt.timestamp,
source,
&mut state_map,
sink.as_mut(),
&mut counters,
&mut flow_states,
);
ctx.monitor_name = monitor_name_borrow;
ctx.label_table = &label_table;
ctx.tracker = Some(driver.tracker());
slot.drain_and_dispatch(&mut dispatcher, &mut ctx)?;
}
}
}
if !drain_timeout.is_zero() {
let deadline = Instant::now() + drain_timeout;
drain_phase(
&mut driver,
&mut dispatcher,
sink.as_mut(),
&mut state_map,
&mut counters,
&mut protocol_slots,
monitor_name_borrow,
deadline,
&mut flow_states,
&label_table,
)
.await?;
}
Ok(())
}
#[cfg(all(feature = "pcap", feature = "tokio"))]
pub(crate) async fn replay_loop(
monitor: Monitor,
path: std::path::PathBuf,
config: crate::pcap_source::AsyncPcapConfig,
) -> Result<()> {
use futures_core::Stream;
let Monitor {
interfaces: _,
mut driver,
mut dispatcher,
mut protocol_slots,
mut state_map,
mut counters,
mut sink,
tick_handlers: _,
detector_names: _,
monitor_name,
drain_timeout,
broadcast_handles: _,
pcap_source_path: _,
pcap_speed_factor: _,
mut flow_states,
fanout: _,
label_table,
merge_rx: _, } = monitor;
let monitor_name_borrow: Option<&str> = monitor_name.as_deref();
let mut source = crate::pcap_source::AsyncPcapSource::open_with_config(&path, config).await?;
let mut events: Vec<FsEvent<FlowKey>> = Vec::with_capacity(64);
loop {
let next = std::future::poll_fn(|cx| Pin::new(&mut source).poll_next(cx)).await;
let pkt = match next {
Some(Ok(p)) => p,
Some(Err(e)) => return Err(e),
None => break,
};
let view = flowscope::PacketView::new(&pkt.data, pkt.timestamp);
events.clear();
driver.track_into(view, &mut events);
for evt in events.drain(..) {
dispatch_lifecycle(
&mut dispatcher,
sink.as_mut(),
&mut state_map,
&mut counters,
evt.clone(),
SourceIdx(0),
monitor_name_borrow,
&mut flow_states,
&label_table,
)?;
dispatch_lifecycle_async(&mut dispatcher, evt).await?;
}
for slot in &mut protocol_slots {
let mut ctx = Ctx::new(
None,
pkt.timestamp,
SourceIdx(0),
&mut state_map,
sink.as_mut(),
&mut counters,
&mut flow_states,
);
ctx.monitor_name = monitor_name_borrow;
ctx.label_table = &label_table;
ctx.tracker = Some(driver.tracker());
slot.drain_and_dispatch(&mut dispatcher, &mut ctx)?;
}
}
if !drain_timeout.is_zero() {
let deadline = Instant::now() + drain_timeout;
drain_phase(
&mut driver,
&mut dispatcher,
sink.as_mut(),
&mut state_map,
&mut counters,
&mut protocol_slots,
monitor_name_borrow,
deadline,
&mut flow_states,
&label_table,
)
.await?;
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
async fn drain_phase(
driver: &mut flowscope::driver::Driver<FiveTuple>,
dispatcher: &mut Dispatcher,
sink: &mut dyn AnomalySink,
state_map: &mut StateMap,
counters: &mut CounterRegistry,
protocol_slots: &mut [Box<dyn crate::monitor::ProtocolSlot>],
monitor_name: Option<&str>,
deadline: Instant,
flow_states: &mut crate::ctx::FlowStateRegistry,
label_table: &flowscope::well_known::LabelTable,
) -> Result<()> {
let mut leftover: Vec<FsEvent<FlowKey>> = Vec::new();
driver.finish_into(&mut leftover);
for evt in leftover.drain(..) {
if Instant::now() >= deadline {
return Ok(());
}
dispatch_lifecycle(
dispatcher,
sink,
state_map,
counters,
evt.clone(),
SourceIdx(0),
monitor_name,
flow_states,
label_table,
)?;
dispatch_lifecycle_async(dispatcher, evt).await?;
}
if Instant::now() >= deadline {
return Ok(());
}
let ts = flowscope::Timestamp::from_system_time(SystemTime::now());
for slot in protocol_slots.iter_mut() {
if Instant::now() >= deadline {
return Ok(());
}
let mut ctx = Ctx::new(
None,
ts,
SourceIdx(0),
state_map,
sink,
counters,
flow_states,
);
ctx.monitor_name = monitor_name;
ctx.label_table = label_table;
ctx.tracker = Some(driver.tracker());
slot.drain_and_dispatch(dispatcher, &mut ctx)?;
}
if Instant::now() >= deadline {
return Ok(());
}
sink.flush().map_err(|e| {
crate::error::Error::Io(std::io::Error::new(e.kind(), format!("sink flush: {e}")))
})?;
Ok(())
}
async fn next_packet<S>(
streams: &mut [PacketStream<S>],
anchor: &mut usize,
) -> Option<(usize, Result<Vec<OwnedPacket>>)>
where
S: crate::traits::PacketSource + std::os::unix::io::AsRawFd + Unpin,
{
std::future::poll_fn(
|cx: &mut Context<'_>| -> Poll<Option<(usize, Result<Vec<OwnedPacket>>)>> {
let n = streams.len();
if n == 0 {
return Poll::Ready(None);
}
let start = *anchor % n;
let mut all_done = true;
for offset in 0..n {
let i = (start + offset) % n;
match Pin::new(&mut streams[i]).poll_next(cx) {
Poll::Ready(Some(item)) => {
*anchor = (i + 1) % n;
return Poll::Ready(Some((i, item)));
}
Poll::Ready(None) => {
}
Poll::Pending => {
all_done = false;
}
}
}
if all_done {
Poll::Ready(None)
} else {
Poll::Pending
}
},
)
.await
}
async fn next_tick(intervals: &mut [tokio::time::Interval]) -> usize {
std::future::poll_fn(|cx: &mut Context<'_>| -> Poll<usize> {
for (i, interval) in intervals.iter_mut().enumerate() {
if interval.poll_tick(cx).is_ready() {
return Poll::Ready(i);
}
}
Poll::Pending
})
.await
}
async fn recv_merge(
rx: &mut Option<tokio::sync::mpsc::UnboundedReceiver<crate::monitor::merge::MergeRequest>>,
) -> Option<crate::monitor::merge::MergeRequest> {
match rx {
Some(r) => r.recv().await,
None => std::future::pending().await,
}
}
#[allow(clippy::too_many_arguments)]
async fn fire_tick(
tick_idx: usize,
tick_handlers: &mut [crate::monitor::tick::TickRegistration],
dispatcher: &mut Dispatcher,
sink: &mut dyn AnomalySink,
state_map: &mut StateMap,
counters: &mut CounterRegistry,
monitor_name: Option<&str>,
flow_states: &mut crate::ctx::FlowStateRegistry,
label_table: &flowscope::well_known::LabelTable,
) -> Result<()> {
let reg = &mut tick_handlers[tick_idx];
let tick = Tick {
now: flowscope::Timestamp::from_system_time(SystemTime::now()),
period: reg.period,
};
{
let mut ctx = Ctx::new(
None,
tick.now,
SourceIdx(0),
state_map,
sink,
counters,
flow_states,
);
ctx.monitor_name = monitor_name;
ctx.label_table = label_table;
(reg.handler)(&tick, &mut ctx)?;
}
{
let mut ctx = Ctx::new(
None,
tick.now,
SourceIdx(0),
state_map,
sink,
counters,
flow_states,
);
ctx.monitor_name = monitor_name;
ctx.label_table = label_table;
dispatcher.dispatch::<Tick>(&tick, &mut ctx)?;
}
dispatcher.dispatch_async::<Tick>(&tick).await?;
Ok(())
}
struct ShutdownSignal {
stop: StopCondition,
sig_int: Option<tokio::signal::unix::Signal>,
sig_term: Option<tokio::signal::unix::Signal>,
}
impl ShutdownSignal {
fn new(stop: StopCondition) -> Self {
let (sig_int, sig_term) = match &stop {
StopCondition::Signal => {
let sigint =
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::interrupt()).ok();
let sigterm =
tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()).ok();
(sigint, sigterm)
}
StopCondition::Deadline(_) | StopCondition::Idle(_) => (None, None),
};
Self {
stop,
sig_int,
sig_term,
}
}
async fn recv(&mut self, last_event_at: Instant) {
match &mut self.stop {
StopCondition::Deadline(t) => {
tokio::time::sleep_until((*t).into()).await;
}
StopCondition::Idle(window) => {
tokio::time::sleep_until((last_event_at + *window).into()).await;
}
StopCondition::Signal => match (self.sig_int.as_mut(), self.sig_term.as_mut()) {
(Some(i), Some(t)) => {
tokio::select! {
_ = i.recv() => {},
_ = t.recv() => {},
}
}
(Some(i), None) => {
let _ = i.recv().await;
}
(None, Some(t)) => {
let _ = t.recv().await;
}
(None, None) => std::future::pending::<()>().await,
},
}
}
}
async fn dispatch_lifecycle_async(
dispatcher: &mut Dispatcher,
evt: FsEvent<FlowKey>,
) -> Result<()> {
match evt {
FsEvent::FlowStarted { key, ts, l4 } => match l4 {
Some(L4Proto::Tcp) => {
dispatcher
.dispatch_async(&FlowStarted::<Tcp>::new(key, l4, ts))
.await?;
}
Some(L4Proto::Udp) => {
dispatcher
.dispatch_async(&FlowStarted::<Udp>::new(key, l4, ts))
.await?;
}
#[cfg(feature = "icmp")]
Some(L4Proto::Icmp) | Some(L4Proto::IcmpV6) => {
dispatcher
.dispatch_async(&FlowStarted::<Icmp>::new(key, l4, ts))
.await?;
}
_ => {}
},
FsEvent::FlowEnded {
key,
reason,
stats,
ts,
l4,
..
} => match l4 {
Some(L4Proto::Tcp) => {
let is_rst = reason == flowscope::EndReason::Rst;
dispatcher
.dispatch_async(&FlowEnded::<Tcp>::new(key, reason, stats.clone(), l4, ts))
.await?;
if is_rst {
dispatcher
.dispatch_async(&TcpRst::new(key, stats, ts))
.await?;
}
}
Some(L4Proto::Udp) => {
dispatcher
.dispatch_async(&FlowEnded::<Udp>::new(key, reason, stats, l4, ts))
.await?;
}
#[cfg(feature = "icmp")]
Some(L4Proto::Icmp) | Some(L4Proto::IcmpV6) => {
dispatcher
.dispatch_async(&FlowEnded::<Icmp>::new(key, reason, stats, l4, ts))
.await?;
}
_ => {}
},
FsEvent::FlowEstablished { key, ts, l4 } => {
if matches!(l4, Some(L4Proto::Tcp)) {
dispatcher
.dispatch_async(&FlowEstablished::<Tcp>::new(key, ts))
.await?;
}
}
FsEvent::FlowAnomaly { key, kind, ts } => {
dispatcher
.dispatch_async(&AnyFlowAnomaly {
key: Some(key),
kind,
ts,
})
.await?;
}
FsEvent::TrackerAnomaly { kind, ts } => {
dispatcher
.dispatch_async(&AnyFlowAnomaly {
key: None,
kind,
ts,
})
.await?;
}
FsEvent::FlowPacket {
key,
side,
len,
ts,
tcp,
} => {
dispatcher
.dispatch_async(&FlowPacket::new(key.proto, key, side, len, tcp, ts))
.await?;
}
FsEvent::FlowTick { key, stats, ts } => match key.proto {
L4Proto::Tcp => {
dispatcher
.dispatch_async(&FlowTick::<Tcp>::new(key, stats, ts))
.await?;
}
L4Proto::Udp => {
dispatcher
.dispatch_async(&FlowTick::<Udp>::new(key, stats, ts))
.await?;
}
#[cfg(feature = "icmp")]
L4Proto::Icmp | L4Proto::IcmpV6 => {
dispatcher
.dispatch_async(&FlowTick::<Icmp>::new(key, stats, ts))
.await?;
}
_ => {}
},
FsEvent::ParserClosed {
key,
parser_kind,
reason,
ts,
} => match key.proto {
L4Proto::Tcp => {
dispatcher
.dispatch_async(&ParserClosed::<Tcp>::new(key, parser_kind, reason, ts))
.await?;
}
L4Proto::Udp => {
dispatcher
.dispatch_async(&ParserClosed::<Udp>::new(key, parser_kind, reason, ts))
.await?;
}
#[cfg(feature = "icmp")]
L4Proto::Icmp | L4Proto::IcmpV6 => {
dispatcher
.dispatch_async(&ParserClosed::<Icmp>::new(key, parser_kind, reason, ts))
.await?;
}
_ => {}
},
_ => {}
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
fn dispatch_lifecycle(
dispatcher: &mut Dispatcher,
sink: &mut dyn AnomalySink,
state_map: &mut StateMap,
counters: &mut CounterRegistry,
evt: FsEvent<FlowKey>,
source: SourceIdx,
monitor_name: Option<&str>,
flow_states: &mut crate::ctx::FlowStateRegistry,
label_table: &flowscope::well_known::LabelTable,
) -> Result<()> {
macro_rules! dispatch_one {
($ty:ty, $payload:expr, $flow:expr, $ts:expr) => {{
let mut ctx = Ctx {
flow: $flow,
ts: $ts,
source,
monitor_name,
state_map: &mut *state_map,
sink: &mut *sink,
counters: &mut *counters,
flow_states: &mut *flow_states,
label_table,
tracker: None,
};
dispatcher.dispatch::<$ty>(&$payload, &mut ctx)?;
}};
}
match evt {
FsEvent::FlowStarted { key, ts, l4 } => match l4 {
Some(L4Proto::Tcp) => {
dispatch_one!(
FlowStarted<Tcp>,
FlowStarted::<Tcp>::new(key, l4, ts),
Some(key),
ts
);
}
Some(L4Proto::Udp) => {
dispatch_one!(
FlowStarted<Udp>,
FlowStarted::<Udp>::new(key, l4, ts),
Some(key),
ts
);
}
#[cfg(feature = "icmp")]
Some(L4Proto::Icmp) | Some(L4Proto::IcmpV6) => {
dispatch_one!(
FlowStarted<Icmp>,
FlowStarted::<Icmp>::new(key, l4, ts),
Some(key),
ts
);
}
_ => {}
},
FsEvent::FlowEnded {
key,
reason,
stats,
ts,
l4,
..
} => match l4 {
Some(L4Proto::Tcp) => {
let is_rst = reason == flowscope::EndReason::Rst;
dispatch_one!(
FlowEnded<Tcp>,
FlowEnded::<Tcp>::new(key, reason, stats.clone(), l4, ts),
Some(key),
ts
);
if is_rst {
dispatch_one!(TcpRst, TcpRst::new(key, stats, ts), Some(key), ts);
}
}
Some(L4Proto::Udp) => {
dispatch_one!(
FlowEnded<Udp>,
FlowEnded::<Udp>::new(key, reason, stats, l4, ts),
Some(key),
ts
);
}
#[cfg(feature = "icmp")]
Some(L4Proto::Icmp) | Some(L4Proto::IcmpV6) => {
dispatch_one!(
FlowEnded<Icmp>,
FlowEnded::<Icmp>::new(key, reason, stats, l4, ts),
Some(key),
ts
);
}
_ => {}
},
FsEvent::FlowEstablished { key, ts, l4 } => {
if matches!(l4, Some(L4Proto::Tcp)) {
dispatch_one!(
FlowEstablished<Tcp>,
FlowEstablished::<Tcp>::new(key, ts),
Some(key),
ts
);
}
}
FsEvent::FlowAnomaly { key, kind, ts } => {
dispatch_one!(
AnyFlowAnomaly,
AnyFlowAnomaly {
key: Some(key),
kind,
ts,
},
Some(key),
ts
);
}
FsEvent::TrackerAnomaly { kind, ts } => {
dispatch_one!(
AnyFlowAnomaly,
AnyFlowAnomaly {
key: None,
kind,
ts,
},
None,
ts
);
}
FsEvent::FlowPacket {
key,
side,
len,
ts,
tcp,
} => {
dispatch_one!(
FlowPacket,
FlowPacket::new(key.proto, key, side, len, tcp, ts),
Some(key),
ts
);
}
FsEvent::FlowTick { key, stats, ts } => match key.proto {
L4Proto::Tcp => {
dispatch_one!(
FlowTick<Tcp>,
FlowTick::<Tcp>::new(key, stats, ts),
Some(key),
ts
);
}
L4Proto::Udp => {
dispatch_one!(
FlowTick<Udp>,
FlowTick::<Udp>::new(key, stats, ts),
Some(key),
ts
);
}
#[cfg(feature = "icmp")]
L4Proto::Icmp | L4Proto::IcmpV6 => {
dispatch_one!(
FlowTick<Icmp>,
FlowTick::<Icmp>::new(key, stats, ts),
Some(key),
ts
);
}
_ => {}
},
FsEvent::ParserClosed {
key,
parser_kind,
reason,
ts,
} => match key.proto {
L4Proto::Tcp => {
dispatch_one!(
ParserClosed<Tcp>,
ParserClosed::<Tcp>::new(key, parser_kind, reason, ts),
Some(key),
ts
);
}
L4Proto::Udp => {
dispatch_one!(
ParserClosed<Udp>,
ParserClosed::<Udp>::new(key, parser_kind, reason, ts),
Some(key),
ts
);
}
#[cfg(feature = "icmp")]
L4Proto::Icmp | L4Proto::IcmpV6 => {
dispatch_one!(
ParserClosed<Icmp>,
ParserClosed::<Icmp>::new(key, parser_kind, reason, ts),
Some(key),
ts
);
}
_ => {}
},
_ => {}
}
Ok(())
}