pub struct Detector<E, F> {
pub handler: F,
pub name: &'static str,
pub declared_counters: Vec<&'static str>,
pub _marker: ::std::marker::PhantomData<fn() -> E>,
}
impl<E, F> Detector<E, F> {
pub fn new(handler: F) -> Self {
Self {
handler,
name: "unnamed",
declared_counters: Vec::new(),
_marker: ::std::marker::PhantomData,
}
}
pub fn with_name(mut self, name: &'static str) -> Self {
self.name = name;
self
}
pub fn with_declared_counters(mut self, slugs: Vec<&'static str>) -> Self {
self.declared_counters = slugs;
self
}
}
impl<E, F> crate::monitor::handler::Handler<E, crate::monitor::handler::PayloadCtx>
for Detector<E, F>
where
E: crate::protocol::event_typed::Event,
F: crate::monitor::handler::Handler<E, crate::monitor::handler::PayloadCtx>,
{
#[inline]
fn call(
&self,
payload: &E::Payload,
ctx: &mut crate::ctx::Ctx<'_>,
) -> crate::error::Result<()> {
self.handler.call(payload, ctx)
}
}
#[macro_export]
macro_rules! pattern_detector {
(
name: $name:literal,
event: $ev:ty,
detector: $detector_expr:expr,
feed: |$evt_pat:pat_param, $det_pat:pat_param| $feed_body:expr,
verdict: |$evt_pat2:pat_param, $det_ref_pat:pat_param| $verdict_body:expr $(,)?
) => {{
let detector = ::std::sync::Arc::new(::std::sync::Mutex::new($detector_expr));
let __handler =
move |__payload: &<$ev as $crate::protocol::event_typed::Event>::Payload,
__ctx: &mut $crate::ctx::Ctx<'_>|
-> $crate::error::Result<()> {
let mut guard = detector
.lock()
.expect("pattern_detector! mutex poisoned by a panicking detector body");
{
let $evt_pat = __payload;
let $det_pat = &mut *guard;
$feed_body;
}
let score_opt = {
let $evt_pat2 = __payload;
let $det_ref_pat = &*guard;
$verdict_body
};
drop(guard);
if let ::std::option::Option::Some(score) = score_opt {
let owned =
<_ as $crate::anomaly::DetectorScore>::into_anomaly(score, __ctx.ts);
$crate::anomaly::sink::publish_owned(__ctx.sink_mut(), &owned);
}
::std::result::Result::Ok(())
};
$crate::detector_macro::Detector::<$ev, _>::new(__handler).with_name($name)
}};
}
#[macro_export]
macro_rules! detector {
(
name: $name:literal,
$( counters: [ $( $counter:ty ),+ $(,)? ], )?
severity: $sev:ident,
event: $ev:ty,
$( matches: |$guard_pat:pat_param| $guard_expr:expr, )?
emit: |$payload:pat_param, $ctx:pat_param| $emit_body:expr $(,)?
) => {{
let _: $crate::anomaly::Severity = $crate::anomaly::Severity::$sev;
let __handler = move |__payload: &<$ev as $crate::protocol::event_typed::Event>::Payload,
__ctx: &mut $crate::ctx::Ctx<'_>|
-> $crate::error::Result<()> {
$(
{
let $guard_pat = __payload;
if !($guard_expr) {
return ::std::result::Result::Ok(());
}
}
)?
let $payload = __payload;
let $ctx = __ctx;
#[allow(clippy::redundant_closure_call)]
(|| -> () { $emit_body })();
::std::result::Result::Ok(())
};
let __det = $crate::detector_macro::Detector::<$ev, _>::new(__handler).with_name($name);
$(
let __declared_counters: ::std::vec::Vec<&'static str> = ::std::vec![
$( ::std::any::type_name::<$counter>() ),+
];
let __det = __det.with_declared_counters(__declared_counters);
)?
__det
}};
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use flowscope::Timestamp;
use crate::anomaly::Severity;
use crate::anomaly::sink::NoopSink;
use crate::ctx::{CounterRegistry, Ctx, SourceIdx, StateMap};
use crate::error::Result;
use crate::monitor::Handler;
use crate::protocol::builtin::Tcp;
use crate::protocol::event_typed::FlowStarted;
fn fresh_ctx<'a>(
state: &'a mut StateMap,
sink: &'a mut NoopSink,
counters: &'a mut CounterRegistry,
flow_states: &'a mut crate::ctx::FlowStateRegistry,
) -> Ctx<'a> {
Ctx {
flow: None,
ts: Timestamp::new(0, 0),
source: SourceIdx(0),
monitor_name: None,
state_map: state,
sink,
counters,
flow_states,
label_table: crate::ctx::default_label_table(),
tracker: None,
arp_table: None,
}
}
fn dummy_flow_started() -> FlowStarted<Tcp> {
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
let key = flowscope::extract::FiveTupleKey::new(
flowscope::L4Proto::Tcp,
SocketAddr::new(IpAddr::V4(Ipv4Addr::new(10, 0, 0, 1)), 12345),
SocketAddr::new(IpAddr::V4(Ipv4Addr::new(10, 0, 0, 2)), 80),
);
FlowStarted::<Tcp>::new(key, Some(flowscope::L4Proto::Tcp), Timestamp::new(0, 0))
}
#[test]
fn macro_without_guard_produces_a_handler() {
let counter = Arc::new(AtomicU32::new(0));
let c = Arc::clone(&counter);
let det = crate::detector! {
name: "NoGuard",
severity: Info,
event: FlowStarted<Tcp>,
emit: |_evt, _ctx| {
c.fetch_add(1, Ordering::Relaxed);
},
};
let mut s = StateMap::default();
let mut k = NoopSink;
let mut cr = CounterRegistry::default();
let mut fs = crate::ctx::FlowStateRegistry::default();
let mut ctx = fresh_ctx(&mut s, &mut k, &mut cr, &mut fs);
let evt = dummy_flow_started();
Handler::<FlowStarted<Tcp>, crate::monitor::PayloadCtx>::call(&det, &evt, &mut ctx)
.unwrap();
assert_eq!(counter.load(Ordering::Relaxed), 1);
}
#[test]
fn macro_guard_filters_non_matching_events() {
let counter = Arc::new(AtomicU32::new(0));
let c = Arc::clone(&counter);
let det = crate::detector! {
name: "EvenL4Only",
severity: Warning,
event: FlowStarted<Tcp>,
matches: |evt| matches!(evt.l4, Some(flowscope::L4Proto::Udp)),
emit: |_evt, _ctx| {
c.fetch_add(1, Ordering::Relaxed);
},
};
let mut s = StateMap::default();
let mut k = NoopSink;
let mut cr = CounterRegistry::default();
let mut fs = crate::ctx::FlowStateRegistry::default();
let mut ctx = fresh_ctx(&mut s, &mut k, &mut cr, &mut fs);
let evt = dummy_flow_started();
Handler::<FlowStarted<Tcp>, crate::monitor::PayloadCtx>::call(&det, &evt, &mut ctx)
.unwrap();
assert_eq!(counter.load(Ordering::Relaxed), 0);
}
#[test]
fn macro_emit_body_can_call_ctx_sink_mut() {
fn _det_compiles<F>(_: F)
where
F: Handler<FlowStarted<Tcp>, crate::monitor::PayloadCtx>,
{
}
let det = crate::detector! {
name: "Emits",
severity: Critical,
event: FlowStarted<Tcp>,
emit: |_evt, ctx| {
let now = ctx.ts;
ctx.sink_mut()
.begin("Emits", Severity::Critical, now)
.with("note", "hi")
.emit();
},
};
_det_compiles(det);
}
#[test]
fn macro_used_with_real_dispatch_path() -> Result<()> {
let counter = Arc::new(AtomicU32::new(0));
let c = Arc::clone(&counter);
let det = crate::detector! {
name: "DispatchPath",
severity: Info,
event: FlowStarted<Tcp>,
emit: |_evt, _ctx| {
c.fetch_add(1, Ordering::Relaxed);
},
};
let mut reg = crate::monitor::HandlerRegistry::default();
reg.register::<FlowStarted<Tcp>, _, _>(det);
let mut disp = reg.into_dispatcher().expect("dispatcher");
let mut s = StateMap::default();
let mut k = NoopSink;
let mut cr = CounterRegistry::default();
let mut fs = crate::ctx::FlowStateRegistry::default();
let mut ctx = fresh_ctx(&mut s, &mut k, &mut cr, &mut fs);
disp.dispatch::<FlowStarted<Tcp>>(&dummy_flow_started(), &mut ctx)?;
assert_eq!(counter.load(Ordering::Relaxed), 1);
Ok(())
}
}