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(
crate::anomaly::sink::detector_kind_for(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(feature = "tls")]
pub fn eve_tls_record(
fp: &crate::monitor::TlsFingerprint,
ts: Timestamp,
in_iface: &str,
) -> serde_json::Value {
use serde_json::json;
let mut tls = serde_json::Map::new();
if let Some(sni) = &fp.sni {
tls.insert("sni".into(), json!(sni));
}
if let Some(alpn) = &fp.alpn {
tls.insert("alpn".into(), json!(alpn));
}
if let Some(ja3) = &fp.ja3 {
tls.insert("ja3".into(), json!({ "hash": ja3 }));
}
if let Some(ja4) = &fp.ja4 {
tls.insert("ja4".into(), json!(ja4));
}
#[cfg(feature = "ja4plus")]
if let Some(ja4s) = &fp.ja4s {
tls.insert("ja4s".into(), json!(ja4s));
}
if fp.app_protocol != flowscope::app_proto::AppProtocol::Unknown {
tls.insert("app_proto".into(), json!(fp.app_protocol.as_str()));
}
let mut obj = serde_json::Map::new();
obj.insert("timestamp".into(), json!(ts.to_iso8601()));
obj.insert("event_type".into(), json!("tls"));
if !in_iface.is_empty() {
obj.insert("in_iface".into(), json!(in_iface));
}
if let Some(key) = &fp.key {
obj.insert("src_ip".into(), json!(key.a.ip().to_string()));
obj.insert("src_port".into(), json!(key.a.port()));
obj.insert("dest_ip".into(), json!(key.b.ip().to_string()));
obj.insert("dest_port".into(), json!(key.b.port()));
obj.insert("proto".into(), json!(l4_proto_str(key.proto)));
if let Some(cid) = flowscope::KeyFields::community_id(key) {
obj.insert("community_id".into(), json!(cid));
}
}
obj.insert("tls".into(), serde_json::Value::Object(tls));
serde_json::Value::Object(obj)
}
#[cfg(feature = "tls")]
fn l4_proto_str(proto: flowscope::L4Proto) -> &'static str {
use flowscope::L4Proto::*;
match proto {
Tcp => "TCP",
Udp => "UDP",
Icmp => "ICMP",
IcmpV6 => "IPv6-ICMP",
_ => "TCP",
}
}
#[cfg(feature = "tls")]
pub struct EveTlsSink<W: Write + Send> {
writer: W,
in_iface: String,
}
#[cfg(feature = "tls")]
impl<W: Write + Send> EveTlsSink<W> {
pub fn new(writer: W, in_iface: impl Into<String>) -> Self {
Self {
writer,
in_iface: in_iface.into(),
}
}
pub fn write_tls(
&mut self,
fp: &crate::monitor::TlsFingerprint,
ts: Timestamp,
) -> std::io::Result<()> {
let rec = eve_tls_record(fp, ts, &self.in_iface);
writeln!(self.writer, "{rec}")
}
pub fn flush(&mut self) -> std::io::Result<()> {
self.writer.flush()
}
pub fn into_inner(self) -> W {
self.writer
}
}
#[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
}
#[cfg(feature = "tls")]
#[test]
fn eve_tls_record_carries_event_type_fivetuple_and_fingerprint() {
use crate::monitor::TlsFingerprint;
let mut hs = flowscope::tls::TlsHandshake::default();
hs.sni = Some("example.com".into());
hs.server_alpn = Some("h2".into());
hs.ja3 = Some("deadbeef".into());
hs.ja4 = Some("t13d1516h2_test".into());
let key = flowscope::extract::FiveTupleKey::new(
flowscope::L4Proto::Tcp,
SocketAddr::new(IpAddr::V4(Ipv4Addr::new(10, 0, 0, 1)), 1234),
SocketAddr::new(IpAddr::V4(Ipv4Addr::new(10, 0, 0, 2)), 443),
);
let fp = TlsFingerprint::from_handshake(&hs, Some(key));
let v = eve_tls_record(&fp, Timestamp::from_unix_f64(1000.0), "eth0");
assert_eq!(v["event_type"], "tls");
assert_eq!(v["in_iface"], "eth0");
assert_eq!(v["src_ip"], "10.0.0.1");
assert_eq!(v["dest_port"], 443);
assert_eq!(v["proto"], "TCP");
assert_eq!(v["tls"]["sni"], "example.com");
assert_eq!(v["tls"]["ja3"]["hash"], "deadbeef");
assert_eq!(v["tls"]["ja4"], "t13d1516h2_test");
assert_eq!(v["tls"]["alpn"], "h2");
assert!(v["timestamp"].is_string());
let mut sink = EveTlsSink::new(Vec::<u8>::new(), "eth0");
sink.write_tls(&fp, Timestamp::from_unix_f64(1000.0))
.unwrap();
let out = String::from_utf8(sink.into_inner()).unwrap();
assert_eq!(out.lines().count(), 1);
assert!(out.ends_with('\n'));
assert!(out.contains("\"event_type\":\"tls\""), "out = {out}");
}
#[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::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)), 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");
}
}