use std::marker::PhantomData;
use flowscope::{AnomalyKind, EndReason, FlowSide, FlowStats, L4Proto, TcpInfo, Timestamp};
use crate::protocol::{FlowKey, Protocol};
pub trait Event: Send + Sync + 'static {
type Payload: Send + Sync + 'static;
fn protocol_marker() -> Option<std::any::TypeId> {
None
}
fn protocol_name() -> &'static str {
"unknown"
}
}
impl<P: Protocol> Event for P {
type Payload = P::Message;
fn protocol_marker() -> Option<std::any::TypeId> {
Some(std::any::TypeId::of::<P>())
}
fn protocol_name() -> &'static str {
P::NAME
}
}
#[non_exhaustive]
pub struct FlowStarted<P: Protocol> {
pub key: FlowKey,
pub l4: Option<L4Proto>,
pub ts: Timestamp,
_marker: PhantomData<fn() -> P>,
}
impl<P: Protocol> FlowStarted<P> {
#[cfg(feature = "bench-zero-alloc")]
pub fn new_for_bench(key: FlowKey, l4: Option<L4Proto>, ts: Timestamp) -> Self {
Self::new(key, l4, ts)
}
#[doc(hidden)]
pub fn new(key: FlowKey, l4: Option<L4Proto>, ts: Timestamp) -> Self {
Self {
key,
l4,
ts,
_marker: PhantomData,
}
}
}
impl<P: Protocol> Event for FlowStarted<P> {
type Payload = FlowStarted<P>;
}
impl<P: Protocol> std::fmt::Debug for FlowStarted<P> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("FlowStarted")
.field("protocol", &P::NAME)
.field("key", &self.key)
.field("l4", &self.l4)
.field("ts", &self.ts)
.finish()
}
}
#[non_exhaustive]
pub struct FlowEnded<P: Protocol> {
pub key: FlowKey,
pub reason: EndReason,
pub stats: FlowStats,
pub l4: Option<L4Proto>,
pub ts: Timestamp,
_marker: PhantomData<fn() -> P>,
}
impl<P: Protocol> FlowEnded<P> {
#[doc(hidden)]
pub fn new(
key: FlowKey,
reason: EndReason,
stats: FlowStats,
l4: Option<L4Proto>,
ts: Timestamp,
) -> Self {
Self {
key,
reason,
stats,
l4,
ts,
_marker: PhantomData,
}
}
}
impl<P: Protocol> Event for FlowEnded<P> {
type Payload = FlowEnded<P>;
}
impl<P: Protocol> std::fmt::Debug for FlowEnded<P> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("FlowEnded")
.field("protocol", &P::NAME)
.field("key", &self.key)
.field("reason", &self.reason)
.field("l4", &self.l4)
.field("ts", &self.ts)
.finish()
}
}
#[non_exhaustive]
pub struct FlowEstablished<P: Protocol> {
pub key: FlowKey,
pub ts: Timestamp,
_marker: PhantomData<fn() -> P>,
}
impl<P: Protocol> FlowEstablished<P> {
#[doc(hidden)]
pub fn new(key: FlowKey, ts: Timestamp) -> Self {
Self {
key,
ts,
_marker: PhantomData,
}
}
}
impl<P: Protocol> Event for FlowEstablished<P> {
type Payload = FlowEstablished<P>;
}
impl<P: Protocol> std::fmt::Debug for FlowEstablished<P> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("FlowEstablished")
.field("protocol", &P::NAME)
.field("key", &self.key)
.field("ts", &self.ts)
.finish()
}
}
#[non_exhaustive]
pub struct FlowPacket<P: Protocol> {
pub key: FlowKey,
pub side: FlowSide,
pub len: usize,
pub tcp: Option<TcpInfo>,
pub ts: Timestamp,
_marker: PhantomData<fn() -> P>,
}
impl<P: Protocol> FlowPacket<P> {
#[doc(hidden)]
pub fn new(
key: FlowKey,
side: FlowSide,
len: usize,
tcp: Option<TcpInfo>,
ts: Timestamp,
) -> Self {
Self {
key,
side,
len,
tcp,
ts,
_marker: PhantomData,
}
}
}
impl<P: Protocol> Event for FlowPacket<P> {
type Payload = FlowPacket<P>;
}
impl<P: Protocol> std::fmt::Debug for FlowPacket<P> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("FlowPacket")
.field("protocol", &P::NAME)
.field("key", &self.key)
.field("side", &self.side)
.field("len", &self.len)
.field("tcp", &self.tcp)
.field("ts", &self.ts)
.finish()
}
}
#[non_exhaustive]
pub struct FlowTick<P: Protocol> {
pub key: FlowKey,
pub stats: FlowStats,
pub ts: Timestamp,
_marker: PhantomData<fn() -> P>,
}
impl<P: Protocol> FlowTick<P> {
#[doc(hidden)]
pub fn new(key: FlowKey, stats: FlowStats, ts: Timestamp) -> Self {
Self {
key,
stats,
ts,
_marker: PhantomData,
}
}
}
impl<P: Protocol> Event for FlowTick<P> {
type Payload = FlowTick<P>;
}
impl<P: Protocol> std::fmt::Debug for FlowTick<P> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("FlowTick")
.field("protocol", &P::NAME)
.field("key", &self.key)
.field("stats", &self.stats)
.field("ts", &self.ts)
.finish()
}
}
#[non_exhaustive]
pub struct ParserClosed<P: Protocol> {
pub key: FlowKey,
pub parser_kind: &'static str,
pub reason: EndReason,
pub ts: Timestamp,
_marker: PhantomData<fn() -> P>,
}
impl<P: Protocol> ParserClosed<P> {
#[doc(hidden)]
pub fn new(key: FlowKey, parser_kind: &'static str, reason: EndReason, ts: Timestamp) -> Self {
Self {
key,
parser_kind,
reason,
ts,
_marker: PhantomData,
}
}
}
impl<P: Protocol> Event for ParserClosed<P> {
type Payload = ParserClosed<P>;
}
impl<P: Protocol> std::fmt::Debug for ParserClosed<P> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ParserClosed")
.field("protocol", &P::NAME)
.field("key", &self.key)
.field("parser_kind", &self.parser_kind)
.field("reason", &self.reason)
.field("ts", &self.ts)
.finish()
}
}
#[derive(Debug)]
#[non_exhaustive]
pub struct AnyFlowAnomaly {
pub key: Option<FlowKey>,
pub kind: AnomalyKind,
pub ts: Timestamp,
}
impl Event for AnyFlowAnomaly {
type Payload = AnyFlowAnomaly;
}
#[derive(Debug, Clone, Copy)]
#[non_exhaustive]
pub struct Tick {
pub now: Timestamp,
pub period: std::time::Duration,
}
impl Tick {
#[doc(hidden)]
pub fn new(now: Timestamp, period: std::time::Duration) -> Self {
Self { now, period }
}
}
impl Event for Tick {
type Payload = Tick;
}
pub use flowscope::FlowSide as Side;
#[cfg(test)]
mod tests {
use super::*;
use crate::protocol::builtin::{Tcp, Udp};
#[test]
fn flow_started_typed_by_protocol() {
fn _accept_tcp(_: &FlowStarted<Tcp>) {}
fn _accept_udp(_: &FlowStarted<Udp>) {}
}
#[test]
fn typed_events_are_send_sync() {
fn assert_send<T: Send>() {}
fn assert_sync<T: Sync>() {}
assert_send::<FlowStarted<Tcp>>();
assert_sync::<FlowStarted<Tcp>>();
assert_send::<FlowEnded<Tcp>>();
assert_sync::<FlowEnded<Tcp>>();
assert_send::<FlowEstablished<Tcp>>();
assert_send::<AnyFlowAnomaly>();
assert_send::<Tick>();
}
#[test]
fn event_trait_blanket_impl_for_protocol_markers() {
#[cfg(feature = "http")]
{
use crate::protocol::builtin::Http;
fn _accept_event<E: Event>() {}
_accept_event::<Http>();
}
}
#[test]
fn flow_started_debug_includes_protocol_name() {
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
let key = flowscope::extract::FiveTupleKey {
proto: flowscope::L4Proto::Tcp,
a: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(10, 0, 0, 1)), 12345),
b: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(10, 0, 0, 2)), 80),
};
let evt = FlowStarted::<Tcp>::new(key, Some(flowscope::L4Proto::Tcp), Timestamp::new(0, 0));
let s = format!("{evt:?}");
assert!(s.contains("tcp"));
}
}