use std::borrow::Cow;
use std::io::Write;
use flowscope::Timestamp;
use flowscope::emit::{EveJsonWriter, EveOptions};
use crate::anomaly::Severity;
use crate::anomaly::key::Key;
use crate::anomaly::sink::AnomalySink;
pub struct EveSink<W: Write + Send> {
inner: EveJsonWriter<W>,
}
impl<W: Write + Send> EveSink<W> {
pub fn new(sink: W, options: EveOptions) -> Self {
Self {
inner: EveJsonWriter::with_options(sink, options),
}
}
pub fn writer(&self) -> &EveJsonWriter<W> {
&self.inner
}
pub fn finish(self) -> std::io::Result<W> {
self.inner.finish()
}
}
impl EveSink<std::io::Stdout> {
pub fn stdout(options: EveOptions) -> Self {
Self::new(std::io::stdout(), options)
}
}
impl<W: Write + Send> AnomalySink for EveSink<W> {
fn write(
&mut self,
kind: &'static str,
severity: Severity,
ts: Timestamp,
key: Option<&dyn Key>,
observations: &[(&'static str, Cow<'_, str>)],
metrics: &[(&'static str, f64)],
) {
let mut owned = flowscope::OwnedAnomaly::new(kind, severity.into(), ts);
if let Some(k) = key
&& let Some(fkey) = k
.as_any()
.downcast_ref::<flowscope::extract::FiveTupleKey>()
{
owned = owned.with_key(fkey);
}
for (label, value) in observations {
owned = owned.with_observation(label, value.to_string());
}
for (label, value) in metrics {
owned = owned.with_metric(label, *value);
}
let _ = self.inner.write_owned_anomaly(&owned);
}
fn flush(&mut self) -> Result<(), std::io::Error> {
self.inner.flush()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::anomaly::sink::AnomalySinkExt;
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
fn eve_options() -> EveOptions {
let mut o = EveOptions::default();
o.in_iface = "eth0".into();
o
}
#[test]
fn eve_sink_emits_valid_json_with_5tuple_when_key_is_five_tuple() {
let mut sink = EveSink::new(Vec::<u8>::new(), eve_options());
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)), 443),
};
sink.begin(
"PortScan",
Severity::Warning,
Timestamp::new(1_700_000_000, 0),
)
.with_key(&key)
.with("note", "rapid SYN")
.with_metric("rate_pps", 320.0)
.emit();
let bytes = sink.finish().expect("EveJsonWriter::finish recovers Vec");
let line = std::str::from_utf8(&bytes).unwrap();
let parsed: serde_json::Value =
serde_json::from_str(line.trim()).expect("EveSink emits valid JSON");
assert_eq!(parsed["event_type"], "anomaly");
assert_eq!(parsed["anomaly"]["event"], "PortScan");
assert_eq!(parsed["in_iface"], "eth0");
assert_eq!(parsed["src_ip"], "10.0.0.1");
assert_eq!(parsed["src_port"], 12345);
assert_eq!(parsed["dest_ip"], "10.0.0.2");
assert_eq!(parsed["dest_port"], 443);
assert_eq!(parsed["proto"], "TCP");
assert_eq!(parsed["anomaly"]["labels"]["note"], "rapid SYN");
assert_eq!(parsed["anomaly"]["metrics"]["rate_pps"], 320.0);
}
#[test]
fn eve_sink_emits_record_without_5tuple_when_key_is_not_five_tuple() {
let mut sink = EveSink::new(Vec::<u8>::new(), eve_options());
let key: u32 = 42;
sink.begin("DgaQuery", Severity::Info, Timestamp::new(1_700_000_000, 0))
.with_key(&key)
.with("qname", "kjasdfkasdf.example")
.emit();
let bytes = sink.finish().expect("finish ok");
let line = std::str::from_utf8(&bytes).unwrap();
let parsed: serde_json::Value = serde_json::from_str(line.trim()).unwrap();
assert_eq!(parsed["event_type"], "anomaly");
assert_eq!(parsed["anomaly"]["event"], "DgaQuery");
assert!(
parsed.get("src_ip").is_none() || parsed["src_ip"].is_null(),
"non-FiveTupleKey key leaves src_ip None"
);
assert_eq!(parsed["anomaly"]["labels"]["qname"], "kjasdfkasdf.example");
}
}