use std::collections::HashMap;
use std::fs;
use std::path::Path;
use std::sync::Mutex;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Instant;
use regex::Regex;
use rsigma_parser::LogSource;
use serde::{Deserialize, Serialize};
use crate::event::Event;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CompareOp {
Gt,
Gte,
Lt,
Lte,
}
impl CompareOp {
fn apply(self, lhs: f64, rhs: f64) -> bool {
match self {
CompareOp::Gt => lhs > rhs,
CompareOp::Gte => lhs >= rhs,
CompareOp::Lt => lhs < rhs,
CompareOp::Lte => lhs <= rhs,
}
}
fn symbol(self) -> &'static str {
match self {
CompareOp::Gt => ">",
CompareOp::Gte => ">=",
CompareOp::Lt => "<",
CompareOp::Lte => "<=",
}
}
}
#[derive(Debug, Clone)]
pub enum SchemaPredicate {
FieldPresent(String),
FieldAbsent(String),
AnyOf(Vec<String>),
Equals { field: String, value: String },
Matches { field: String, regex: Regex },
Compare {
field: String,
op: CompareOp,
value: f64,
},
In { field: String, values: Vec<String> },
FieldEqualsField { left: String, right: String },
Not(Box<SchemaPredicate>),
Any(Vec<SchemaPredicate>),
All(Vec<SchemaPredicate>),
HasAnyField,
}
impl SchemaPredicate {
fn eval<E: Event + ?Sized>(&self, event: &E) -> bool {
match self {
SchemaPredicate::FieldPresent(f) => event.get_field(f).is_some(),
SchemaPredicate::FieldAbsent(f) => event.get_field(f).is_none(),
SchemaPredicate::AnyOf(fields) => fields.iter().any(|f| event.get_field(f).is_some()),
SchemaPredicate::Equals { field, value } => event
.get_field(field)
.and_then(|v| v.as_str().map(|s| s.as_ref().eq_ignore_ascii_case(value)))
.unwrap_or(false),
SchemaPredicate::Matches { field, regex } => event
.get_field(field)
.and_then(|v| v.as_str().map(|s| regex.is_match(s.as_ref())))
.unwrap_or(false),
SchemaPredicate::Compare { field, op, value } => event
.get_field(field)
.and_then(|v| v.as_f64())
.map(|n| op.apply(n, *value))
.unwrap_or(false),
SchemaPredicate::In { field, values } => event
.get_field(field)
.and_then(|v| {
v.as_str().map(|s| {
values
.iter()
.any(|val| s.as_ref().eq_ignore_ascii_case(val))
})
})
.unwrap_or(false),
SchemaPredicate::FieldEqualsField { left, right } => {
let l = event
.get_field(left)
.and_then(|v| v.as_str().map(|s| s.into_owned()));
let r = event
.get_field(right)
.and_then(|v| v.as_str().map(|s| s.into_owned()));
matches!((l, r), (Some(a), Some(b)) if a.eq_ignore_ascii_case(&b))
}
SchemaPredicate::Not(inner) => !inner.eval(event),
SchemaPredicate::Any(preds) => preds.iter().any(|p| p.eval(event)),
SchemaPredicate::All(preds) => preds.iter().all(|p| p.eval(event)),
SchemaPredicate::HasAnyField => !event.field_keys().is_empty(),
}
}
fn describe(&self) -> String {
match self {
SchemaPredicate::FieldPresent(f) => format!("field_present({f})"),
SchemaPredicate::FieldAbsent(f) => format!("field_absent({f})"),
SchemaPredicate::AnyOf(fs) => format!("any_of([{}])", fs.join(", ")),
SchemaPredicate::Equals { field, value } => format!("{field} == \"{value}\""),
SchemaPredicate::Matches { field, regex } => {
format!("{field} matches /{}/", regex.as_str())
}
SchemaPredicate::Compare { field, op, value } => {
format!("{field} {} {value}", op.symbol())
}
SchemaPredicate::In { field, values } => format!("{field} in [{}]", values.join(", ")),
SchemaPredicate::FieldEqualsField { left, right } => format!("{left} == {right}"),
SchemaPredicate::Not(inner) => format!("not({})", inner.describe()),
SchemaPredicate::Any(ps) => format!(
"any({})",
ps.iter()
.map(|p| p.describe())
.collect::<Vec<_>>()
.join(" | ")
),
SchemaPredicate::All(ps) => format!(
"all({})",
ps.iter()
.map(|p| p.describe())
.collect::<Vec<_>>()
.join(" & ")
),
SchemaPredicate::HasAnyField => "has_any_field".to_string(),
}
}
}
#[derive(Debug, Clone)]
pub struct SchemaSignature {
pub name: String,
pub predicates: Vec<SchemaPredicate>,
pub specificity: u32,
}
impl SchemaSignature {
fn matches<E: Event + ?Sized>(&self, event: &E) -> bool {
self.predicates.iter().all(|p| p.eval(event))
}
fn explain<E: Event + ?Sized>(&self, event: &E) -> SignatureExplanation {
let predicates: Vec<PredicateOutcome> = self
.predicates
.iter()
.map(|p| PredicateOutcome {
predicate: p.describe(),
matched: p.eval(event),
})
.collect();
let predicates_matched = predicates.iter().all(|p| p.matched);
SignatureExplanation {
name: self.name.clone(),
specificity: self.specificity,
predicates_matched,
predicates,
}
}
}
#[derive(Debug, Clone, Serialize)]
pub struct PredicateOutcome {
pub predicate: String,
pub matched: bool,
}
#[derive(Debug, Clone, Serialize)]
pub struct SignatureExplanation {
pub name: String,
pub specificity: u32,
pub predicates_matched: bool,
pub predicates: Vec<PredicateOutcome>,
}
#[derive(Debug, Clone, Serialize)]
pub struct SchemaExplanation {
pub matched: Option<String>,
pub specificity: Option<u32>,
pub signature: Option<SignatureExplanation>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SchemaMatch {
pub name: String,
pub specificity: u32,
}
#[derive(Debug, Clone)]
pub struct SchemaClassifier {
signatures: Vec<SchemaSignature>,
}
impl SchemaClassifier {
pub fn new(mut signatures: Vec<SchemaSignature>) -> Self {
signatures.sort_by(|a, b| {
b.specificity
.cmp(&a.specificity)
.then_with(|| a.name.cmp(&b.name))
});
Self { signatures }
}
pub fn builtin() -> Self {
Self::new(builtin_signatures())
}
pub fn with_user_signatures(user: Vec<SchemaSignature>) -> Self {
let mut signatures = builtin_signatures();
signatures.extend(user);
Self::new(signatures)
}
pub fn classify<E: Event + ?Sized>(&self, event: &E) -> Option<SchemaMatch> {
self.signatures
.iter()
.find(|s| s.matches(event))
.map(|s| SchemaMatch {
name: s.name.clone(),
specificity: s.specificity,
})
}
pub fn classify_with_ambiguity<E: Event + ?Sized>(
&self,
event: &E,
) -> (Option<SchemaMatch>, bool) {
let mut matching = self.signatures.iter().filter(|s| s.matches(event));
let Some(winner) = matching.next() else {
return (None, false);
};
let ambiguous = matching
.take_while(|s| s.specificity == winner.specificity)
.any(|s| s.name != winner.name);
(
Some(SchemaMatch {
name: winner.name.clone(),
specificity: winner.specificity,
}),
ambiguous,
)
}
pub fn classify_all<E: Event + ?Sized>(&self, event: &E) -> Vec<String> {
let mut out: Vec<String> = Vec::new();
for sig in self.signatures.iter().filter(|s| s.matches(event)) {
if !out.iter().any(|n| n == &sig.name) {
out.push(sig.name.clone());
}
}
out
}
pub fn explain<E: Event + ?Sized>(&self, event: &E) -> SchemaExplanation {
let mut best_near: Option<SignatureExplanation> = None;
let mut best_near_passing = 0usize;
for sig in &self.signatures {
let ex = sig.explain(event);
if ex.predicates_matched {
return SchemaExplanation {
matched: Some(ex.name.clone()),
specificity: Some(ex.specificity),
signature: Some(ex),
};
}
let passing = ex.predicates.iter().filter(|p| p.matched).count();
if best_near.is_none() || passing > best_near_passing {
best_near_passing = passing;
best_near = Some(ex);
}
}
SchemaExplanation {
matched: None,
specificity: None,
signature: best_near,
}
}
pub fn schema_names(&self) -> Vec<&str> {
let mut out: Vec<&str> = Vec::new();
for sig in &self.signatures {
if !out.contains(&sig.name.as_str()) {
out.push(sig.name.as_str());
}
}
out
}
}
impl Default for SchemaClassifier {
fn default() -> Self {
Self::builtin()
}
}
fn builtin_signatures() -> Vec<SchemaSignature> {
vec![
SchemaSignature {
name: "ecs_windows".to_string(),
specificity: 105,
predicates: vec![
SchemaPredicate::FieldPresent("ecs.version".to_string()),
SchemaPredicate::Any(vec![
SchemaPredicate::FieldPresent("winlog.channel".to_string()),
SchemaPredicate::FieldPresent("winlog.event_id".to_string()),
SchemaPredicate::Equals {
field: "host.os.type".to_string(),
value: "windows".to_string(),
},
SchemaPredicate::Equals {
field: "os.type".to_string(),
value: "windows".to_string(),
},
]),
],
},
SchemaSignature {
name: "ecs_linux".to_string(),
specificity: 105,
predicates: vec![
SchemaPredicate::FieldPresent("ecs.version".to_string()),
SchemaPredicate::Any(vec![
SchemaPredicate::Equals {
field: "host.os.type".to_string(),
value: "linux".to_string(),
},
SchemaPredicate::Equals {
field: "os.type".to_string(),
value: "linux".to_string(),
},
SchemaPredicate::FieldPresent("host.os.kernel".to_string()),
]),
],
},
SchemaSignature {
name: "ecs".to_string(),
specificity: 100,
predicates: vec![SchemaPredicate::FieldPresent("ecs.version".to_string())],
},
SchemaSignature {
name: "ocsf".to_string(),
specificity: 95,
predicates: vec![
SchemaPredicate::FieldPresent("class_uid".to_string()),
SchemaPredicate::FieldPresent("metadata.version".to_string()),
],
},
SchemaSignature {
name: "windows_eventlog".to_string(),
specificity: 90,
predicates: vec![SchemaPredicate::AnyOf(vec![
"Event.System.EventID".to_string(),
"Event.System.Provider".to_string(),
])],
},
SchemaSignature {
name: "sysmon".to_string(),
specificity: 88,
predicates: vec![SchemaPredicate::Equals {
field: "Channel".to_string(),
value: "Microsoft-Windows-Sysmon/Operational".to_string(),
}],
},
SchemaSignature {
name: "sysmon".to_string(),
specificity: 88,
predicates: vec![SchemaPredicate::Equals {
field: "Provider_Name".to_string(),
value: "Microsoft-Windows-Sysmon".to_string(),
}],
},
SchemaSignature {
name: "sysmon".to_string(),
specificity: 80,
predicates: vec![
SchemaPredicate::FieldPresent("EventID".to_string()),
SchemaPredicate::FieldPresent("ProcessGuid".to_string()),
SchemaPredicate::AnyOf(vec!["Image".to_string(), "CommandLine".to_string()]),
],
},
SchemaSignature {
name: "cef".to_string(),
specificity: 85,
predicates: vec![
SchemaPredicate::FieldPresent("deviceVendor".to_string()),
SchemaPredicate::FieldPresent("deviceProduct".to_string()),
SchemaPredicate::FieldPresent("signatureId".to_string()),
],
},
SchemaSignature {
name: "generic_json".to_string(),
specificity: 0,
predicates: vec![SchemaPredicate::HasAnyField],
},
]
}
pub fn builtin_schema_names() -> Vec<&'static str> {
vec![
"ecs_windows",
"ecs_linux",
"ecs",
"ocsf",
"windows_eventlog",
"sysmon",
"cef",
"generic_json",
]
}
fn builtin_schema_aliases() -> HashMap<String, String> {
HashMap::from([
("ecs_windows".to_string(), "ecs".to_string()),
("ecs_linux".to_string(), "ecs".to_string()),
])
}
#[derive(Debug, thiserror::Error)]
pub enum SchemaError {
#[error("cannot read schema signatures file '{path}': {source}")]
Io {
path: String,
#[source]
source: std::io::Error,
},
#[error("schema signatures YAML parse error: {0}")]
Parse(String),
#[error("invalid regex in schema '{name}': {error}")]
InvalidRegex { name: String, error: String },
}
#[derive(Debug, Clone, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct FieldValueConfig {
pub field: String,
pub value: String,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct FieldNumberConfig {
pub field: String,
pub value: f64,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct FieldValuesConfig {
pub field: String,
pub values: Vec<String>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct FieldPairConfig {
pub left: String,
pub right: String,
}
#[derive(Debug, Clone, Default, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct SchemaPredicateConfig {
#[serde(default)]
pub field_present: Option<String>,
#[serde(default)]
pub field_absent: Option<String>,
#[serde(default)]
pub any_of: Option<Vec<String>>,
#[serde(default)]
pub equals: Option<FieldValueConfig>,
#[serde(default)]
pub matches: Option<FieldValueConfig>,
#[serde(default)]
pub gt: Option<FieldNumberConfig>,
#[serde(default)]
pub gte: Option<FieldNumberConfig>,
#[serde(default)]
pub lt: Option<FieldNumberConfig>,
#[serde(default)]
pub lte: Option<FieldNumberConfig>,
#[serde(default, rename = "in")]
pub in_set: Option<FieldValuesConfig>,
#[serde(default)]
pub field_equals_field: Option<FieldPairConfig>,
#[serde(default)]
pub not: Option<Box<SchemaPredicateConfig>>,
#[serde(default)]
pub any: Option<Vec<SchemaPredicateConfig>>,
#[serde(default)]
pub all: Option<Vec<SchemaPredicateConfig>>,
}
impl SchemaPredicateConfig {
fn build(self, schema_name: &str) -> Result<SchemaPredicate, SchemaError> {
let mut chosen: Option<SchemaPredicate> = None;
let mut set = 0u32;
if let Some(f) = self.field_present {
set += 1;
chosen = Some(SchemaPredicate::FieldPresent(f));
}
if let Some(f) = self.field_absent {
set += 1;
chosen = Some(SchemaPredicate::FieldAbsent(f));
}
if let Some(fields) = self.any_of {
set += 1;
chosen = Some(SchemaPredicate::AnyOf(fields));
}
if let Some(fv) = self.equals {
set += 1;
chosen = Some(SchemaPredicate::Equals {
field: fv.field,
value: fv.value,
});
}
if let Some(fv) = self.matches {
set += 1;
chosen = Some(SchemaPredicate::Matches {
field: fv.field,
regex: Regex::new(&fv.value).map_err(|e| SchemaError::InvalidRegex {
name: schema_name.to_string(),
error: e.to_string(),
})?,
});
}
for (op, cfg) in [
(CompareOp::Gt, self.gt),
(CompareOp::Gte, self.gte),
(CompareOp::Lt, self.lt),
(CompareOp::Lte, self.lte),
] {
if let Some(fv) = cfg {
set += 1;
chosen = Some(SchemaPredicate::Compare {
field: fv.field,
op,
value: fv.value,
});
}
}
if let Some(fv) = self.in_set {
set += 1;
chosen = Some(SchemaPredicate::In {
field: fv.field,
values: fv.values,
});
}
if let Some(fp) = self.field_equals_field {
set += 1;
chosen = Some(SchemaPredicate::FieldEqualsField {
left: fp.left,
right: fp.right,
});
}
if let Some(inner) = self.not {
set += 1;
chosen = Some(SchemaPredicate::Not(Box::new(inner.build(schema_name)?)));
}
if let Some(list) = self.any {
set += 1;
chosen = Some(SchemaPredicate::Any(build_group(list, schema_name, "any")?));
}
if let Some(list) = self.all {
set += 1;
chosen = Some(SchemaPredicate::All(build_group(list, schema_name, "all")?));
}
match (set, chosen) {
(1, Some(p)) => Ok(p),
(0, _) => Err(SchemaError::Parse(format!(
"schema '{schema_name}': a predicate has no condition (expected one of \
field_present, field_absent, any_of, equals, matches, gt, gte, lt, lte, \
in, field_equals_field, not, any, all)"
))),
_ => Err(SchemaError::Parse(format!(
"schema '{schema_name}': a predicate sets multiple conditions; use one per list item"
))),
}
}
}
fn build_group(
list: Vec<SchemaPredicateConfig>,
schema_name: &str,
kind: &str,
) -> Result<Vec<SchemaPredicate>, SchemaError> {
if list.is_empty() {
return Err(SchemaError::Parse(format!(
"schema '{schema_name}': '{kind}' needs at least one sub-predicate"
)));
}
list.into_iter().map(|p| p.build(schema_name)).collect()
}
#[derive(Debug, Clone, Deserialize)]
pub struct SchemaSignatureConfig {
pub name: String,
#[serde(default = "default_user_specificity")]
pub specificity: u32,
#[serde(default, rename = "match")]
pub predicates: Vec<SchemaPredicateConfig>,
}
fn default_user_specificity() -> u32 {
50
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct SchemaSignaturesFile {
#[serde(default)]
pub schemas: Vec<SchemaSignatureConfig>,
#[serde(default)]
pub routing: Option<RoutingConfig>,
}
impl SchemaSignatureConfig {
fn build(self) -> Result<SchemaSignature, SchemaError> {
let name = self.name;
let predicates = self
.predicates
.into_iter()
.map(|p| p.build(&name))
.collect::<Result<Vec<_>, _>>()?;
Ok(SchemaSignature {
name,
predicates,
specificity: self.specificity,
})
}
}
pub fn parse_schema_signatures(yaml: &str) -> Result<Vec<SchemaSignature>, SchemaError> {
let file: SchemaSignaturesFile =
yaml_serde::from_str(yaml).map_err(|e| SchemaError::Parse(e.to_string()))?;
file.schemas.into_iter().map(|s| s.build()).collect()
}
pub fn load_schema_signatures(path: &Path) -> Result<Vec<SchemaSignature>, SchemaError> {
let content = fs::read_to_string(path).map_err(|e| SchemaError::Io {
path: path.display().to_string(),
source: e,
})?;
parse_schema_signatures(&content)
}
pub fn parse_schema_config(
yaml: &str,
) -> Result<(Vec<SchemaSignature>, Option<RoutingConfig>), SchemaError> {
let file: SchemaSignaturesFile =
yaml_serde::from_str(yaml).map_err(|e| SchemaError::Parse(e.to_string()))?;
let signatures = file
.schemas
.into_iter()
.map(|s| s.build())
.collect::<Result<Vec<_>, _>>()?;
Ok((signatures, file.routing))
}
pub fn load_schema_config(
path: &Path,
) -> Result<(Vec<SchemaSignature>, Option<RoutingConfig>), SchemaError> {
let content = fs::read_to_string(path).map_err(|e| SchemaError::Io {
path: path.display().to_string(),
source: e,
})?;
parse_schema_config(&content)
}
pub fn validate_schema_config(
user_signatures: &[SchemaSignature],
routing: Option<&RoutingConfig>,
) -> Vec<String> {
let mut findings = Vec::new();
let mut all = builtin_signatures();
all.extend(user_signatures.iter().cloned());
let preds = |s: &SchemaSignature| -> Vec<String> {
s.predicates.iter().map(|p| p.describe()).collect()
};
for i in 0..user_signatures.len() {
for j in (i + 1)..user_signatures.len() {
if user_signatures[i].name == user_signatures[j].name
&& preds(&user_signatures[i]) == preds(&user_signatures[j])
{
findings.push(format!(
"duplicate signature '{}' with identical predicates",
user_signatures[i].name
));
}
}
}
for b in &all {
let b_preds = preds(b);
for a in &all {
if a.name != b.name
&& a.specificity > b.specificity
&& !a.predicates.is_empty()
&& preds(a).iter().all(|p| b_preds.contains(p))
{
findings.push(format!(
"signature '{}' (specificity {}) is unreachable: shadowed by '{}' (specificity {}) whose predicates are a subset",
b.name, b.specificity, a.name, a.specificity
));
break;
}
}
}
if let Some(routing) = routing {
let mut known: std::collections::HashSet<&str> =
builtin_schema_names().into_iter().collect();
for s in user_signatures {
known.insert(s.name.as_str());
}
let mut seen: std::collections::HashSet<&str> = std::collections::HashSet::new();
for binding in &routing.bindings {
if !known.contains(binding.schema.as_str()) {
findings.push(format!(
"routing binding references unknown schema '{}' (no built-in or user signature produces it)",
binding.schema
));
}
if !seen.insert(binding.schema.as_str()) {
findings.push(format!(
"duplicate routing binding for schema '{}'",
binding.schema
));
}
}
for (alias, canonical) in &routing.aliases {
if !known.contains(canonical.as_str()) {
findings.push(format!(
"alias '{alias}' targets unknown schema '{canonical}' (no built-in or user signature produces it)"
));
}
}
}
findings
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum OnUnknown {
#[default]
Warn,
Drop,
Passthrough,
Error,
}
#[derive(Debug, Clone, Default, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct SchemaLogsource {
#[serde(default)]
pub product: Option<String>,
#[serde(default)]
pub service: Option<String>,
#[serde(default)]
pub category: Option<String>,
#[serde(default)]
pub custom: HashMap<String, String>,
}
impl SchemaLogsource {
fn to_logsource(&self) -> LogSource {
LogSource {
product: self.product.clone(),
service: self.service.clone(),
category: self.category.clone(),
custom: self.custom.clone(),
..LogSource::default()
}
}
}
#[derive(Debug, Clone, Deserialize)]
pub struct SchemaBinding {
pub schema: String,
#[serde(default)]
pub pipelines: Vec<String>,
#[serde(default)]
pub logsource: Option<SchemaLogsource>,
}
fn builtin_schema_logsource() -> HashMap<String, LogSource> {
fn ls(product: &str, service: Option<&str>) -> LogSource {
LogSource {
product: Some(product.to_string()),
service: service.map(str::to_string),
..LogSource::default()
}
}
HashMap::from([
("sysmon".to_string(), ls("windows", Some("sysmon"))),
("windows_eventlog".to_string(), ls("windows", None)),
("ecs_windows".to_string(), ls("windows", None)),
("ecs_linux".to_string(), ls("linux", None)),
])
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct RoutingConfig {
#[serde(default)]
pub on_unknown: OnUnknown,
#[serde(default)]
pub bindings: Vec<SchemaBinding>,
#[serde(default)]
pub default_pipelines: Vec<String>,
#[serde(default)]
pub aliases: HashMap<String, String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RouteDecision {
Evaluate { set: usize, unknown: bool },
Drop,
Error,
}
#[derive(Debug, Clone)]
pub struct RoutingPlan {
pipeline_sets: Vec<Vec<String>>,
schema_to_set: HashMap<String, usize>,
schema_logsource: HashMap<String, LogSource>,
aliases: HashMap<String, String>,
on_unknown: OnUnknown,
}
impl RoutingPlan {
pub fn from_config(config: &RoutingConfig) -> Self {
let mut pipeline_sets: Vec<Vec<String>> = vec![config.default_pipelines.clone()];
let mut schema_to_set: HashMap<String, usize> = HashMap::new();
let mut schema_logsource = builtin_schema_logsource();
let mut aliases = builtin_schema_aliases();
for (alias, canonical) in &config.aliases {
aliases.insert(alias.clone(), canonical.clone());
}
for binding in &config.bindings {
let idx = pipeline_sets
.iter()
.position(|s| s == &binding.pipelines)
.unwrap_or_else(|| {
pipeline_sets.push(binding.pipelines.clone());
pipeline_sets.len() - 1
});
schema_to_set.insert(binding.schema.clone(), idx);
if let Some(ls) = &binding.logsource {
schema_logsource.insert(binding.schema.clone(), ls.to_logsource());
}
}
RoutingPlan {
pipeline_sets,
schema_to_set,
schema_logsource,
aliases,
on_unknown: config.on_unknown,
}
}
pub fn pipeline_sets(&self) -> &[Vec<String>] {
&self.pipeline_sets
}
pub fn on_unknown(&self) -> OnUnknown {
self.on_unknown
}
pub fn schema_logsource(&self, schema: &str) -> Option<&LogSource> {
self.schema_logsource.get(schema)
}
pub fn schemas_with_logsource(&self) -> Vec<String> {
let mut names: Vec<String> = self.schema_logsource.keys().cloned().collect();
names.sort();
names
}
pub fn set_product_partition(&self) -> Vec<Option<std::collections::HashSet<String>>> {
use std::collections::HashSet;
let n = self.pipeline_sets.len();
let mut out: Vec<Option<HashSet<String>>> = (0..n).map(|_| Some(HashSet::new())).collect();
if let Some(first) = out.get_mut(0) {
*first = None;
}
let mut routes: Vec<(usize, &str)> = self
.schema_to_set
.iter()
.map(|(s, &set)| (set, s.as_str()))
.collect();
for (alias, canonical) in &self.aliases {
if !self.schema_to_set.contains_key(alias)
&& let Some(&set) = self.schema_to_set.get(canonical)
{
routes.push((set, alias.as_str()));
}
}
for (set, schema) in routes {
if set == 0 {
continue;
}
let product = self
.schema_logsource
.get(schema)
.and_then(|ls| ls.product.as_deref());
let Some(slot) = out.get_mut(set) else {
continue;
};
match product {
Some(p) => {
if let Some(products) = slot {
products.insert(p.to_ascii_lowercase());
}
}
None => *slot = None,
}
}
out
}
pub fn decide(&self, schema: Option<&str>) -> RouteDecision {
match schema {
Some(s) if self.schema_to_set.contains_key(s) => RouteDecision::Evaluate {
set: self.schema_to_set[s],
unknown: false,
},
Some(s)
if self
.aliases
.get(s)
.and_then(|canonical| self.schema_to_set.get(canonical))
.is_some() =>
{
let canonical = &self.aliases[s];
RouteDecision::Evaluate {
set: self.schema_to_set[canonical],
unknown: false,
}
}
Some(_) => RouteDecision::Evaluate {
set: 0,
unknown: false,
},
None => match self.on_unknown {
OnUnknown::Warn | OnUnknown::Passthrough => RouteDecision::Evaluate {
set: 0,
unknown: true,
},
OnUnknown::Drop => RouteDecision::Drop,
OnUnknown::Error => RouteDecision::Error,
},
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SchemaCountEntry {
pub schema: String,
pub count: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UnknownShapeEntry {
pub keys: Vec<String>,
pub count: u64,
}
const UNKNOWN_SHAPE_CAP: usize = 200;
const UNKNOWN_SHAPE_MAX_KEYS: usize = 64;
#[derive(Debug, Clone, Default)]
pub struct SchemaObservation {
pub by_schema: Vec<SchemaCountEntry>,
pub classified: u64,
pub unknown: u64,
pub ambiguous: u64,
pub unknown_shapes: Vec<UnknownShapeEntry>,
pub unrecognized_shapes: Vec<UnknownShapeEntry>,
pub events_observed: u64,
pub lifetime_classified: u64,
pub lifetime_unknown: u64,
pub lifetime_ambiguous: u64,
pub uptime_seconds: f64,
}
pub struct SchemaObserver {
classifier: SchemaClassifier,
counts: Mutex<HashMap<String, u64>>,
unknown: AtomicU64,
ambiguous: AtomicU64,
unknown_shapes: Mutex<HashMap<Vec<String>, u64>>,
discovery_sampling: bool,
unrecognized_shapes: Mutex<HashMap<Vec<String>, u64>>,
lifetime_classified: AtomicU64,
lifetime_unknown: AtomicU64,
lifetime_ambiguous: AtomicU64,
start: Mutex<Instant>,
}
impl SchemaObserver {
pub fn new(classifier: SchemaClassifier) -> Self {
Self::new_with_discovery(classifier, false)
}
pub fn new_with_discovery(classifier: SchemaClassifier, discovery_sampling: bool) -> Self {
Self {
classifier,
counts: Mutex::new(HashMap::new()),
unknown: AtomicU64::new(0),
ambiguous: AtomicU64::new(0),
unknown_shapes: Mutex::new(HashMap::new()),
discovery_sampling,
unrecognized_shapes: Mutex::new(HashMap::new()),
lifetime_classified: AtomicU64::new(0),
lifetime_unknown: AtomicU64::new(0),
lifetime_ambiguous: AtomicU64::new(0),
start: Mutex::new(Instant::now()),
}
}
pub fn discovery_sampling(&self) -> bool {
self.discovery_sampling
}
pub fn builtin() -> Self {
Self::new(SchemaClassifier::builtin())
}
pub fn observe<E: Event + ?Sized>(&self, event: &E) {
let (matched, ambiguous) = self.classifier.classify_with_ambiguity(event);
if ambiguous {
self.ambiguous.fetch_add(1, Ordering::Relaxed);
self.lifetime_ambiguous.fetch_add(1, Ordering::Relaxed);
}
let discovery_unrecognized = match &matched {
None => true,
Some(m) => m.name == "generic_json",
};
if self.discovery_sampling && discovery_unrecognized {
self.record_unrecognized_shape(event);
}
match matched {
Some(m) => {
self.lifetime_classified.fetch_add(1, Ordering::Relaxed);
let mut counts = self.counts.lock().expect("schema observer mutex poisoned");
*counts.entry(m.name).or_insert(0) += 1;
}
None => {
self.unknown.fetch_add(1, Ordering::Relaxed);
self.lifetime_unknown.fetch_add(1, Ordering::Relaxed);
self.record_unknown_shape(event);
}
}
}
fn record_unknown_shape<E: Event + ?Sized>(&self, event: &E) {
let mut keys: Vec<String> = event.field_keys().iter().map(|k| k.to_string()).collect();
keys.sort();
keys.dedup();
keys.truncate(UNKNOWN_SHAPE_MAX_KEYS);
let mut shapes = self
.unknown_shapes
.lock()
.expect("schema observer shapes mutex poisoned");
if shapes.contains_key(&keys) || shapes.len() < UNKNOWN_SHAPE_CAP {
*shapes.entry(keys).or_insert(0) += 1;
}
}
fn record_unrecognized_shape<E: Event + ?Sized>(&self, event: &E) {
let mut keys: Vec<String> = event.field_keys().iter().map(|k| k.to_string()).collect();
keys.sort();
keys.dedup();
keys.truncate(UNKNOWN_SHAPE_MAX_KEYS);
if keys.is_empty() {
return;
}
let mut shapes = self
.unrecognized_shapes
.lock()
.expect("schema observer shapes mutex poisoned");
if shapes.contains_key(&keys) || shapes.len() < UNKNOWN_SHAPE_CAP {
*shapes.entry(keys).or_insert(0) += 1;
}
}
pub fn snapshot(&self) -> SchemaObservation {
let counts = self.counts.lock().expect("schema observer mutex poisoned");
let mut by_schema: Vec<SchemaCountEntry> = counts
.iter()
.map(|(schema, count)| SchemaCountEntry {
schema: schema.clone(),
count: *count,
})
.collect();
let classified: u64 = counts.values().sum();
drop(counts);
by_schema.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.schema.cmp(&b.schema)));
let shapes = self
.unknown_shapes
.lock()
.expect("schema observer shapes mutex poisoned");
let mut unknown_shapes: Vec<UnknownShapeEntry> = shapes
.iter()
.map(|(keys, count)| UnknownShapeEntry {
keys: keys.clone(),
count: *count,
})
.collect();
drop(shapes);
unknown_shapes.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.keys.cmp(&b.keys)));
let unrec = self
.unrecognized_shapes
.lock()
.expect("schema observer shapes mutex poisoned");
let mut unrecognized_shapes: Vec<UnknownShapeEntry> = unrec
.iter()
.map(|(keys, count)| UnknownShapeEntry {
keys: keys.clone(),
count: *count,
})
.collect();
drop(unrec);
unrecognized_shapes.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.keys.cmp(&b.keys)));
let unknown = self.unknown.load(Ordering::Relaxed);
SchemaObservation {
by_schema,
classified,
unknown,
ambiguous: self.ambiguous.load(Ordering::Relaxed),
unknown_shapes,
unrecognized_shapes,
events_observed: classified + unknown,
lifetime_classified: self.lifetime_classified.load(Ordering::Relaxed),
lifetime_unknown: self.lifetime_unknown.load(Ordering::Relaxed),
lifetime_ambiguous: self.lifetime_ambiguous.load(Ordering::Relaxed),
uptime_seconds: self
.start
.lock()
.expect("schema observer start mutex poisoned")
.elapsed()
.as_secs_f64(),
}
}
pub fn reset(&self) -> (u64, u64) {
let mut counts = self.counts.lock().expect("schema observer mutex poisoned");
let previous_classified: u64 = counts.values().sum();
counts.clear();
drop(counts);
self.unknown_shapes
.lock()
.expect("schema observer shapes mutex poisoned")
.clear();
self.unrecognized_shapes
.lock()
.expect("schema observer shapes mutex poisoned")
.clear();
let previous_unknown = self.unknown.swap(0, Ordering::Relaxed);
self.ambiguous.store(0, Ordering::Relaxed);
*self
.start
.lock()
.expect("schema observer start mutex poisoned") = Instant::now();
(previous_classified, previous_unknown)
}
pub fn lifetime_classified(&self) -> u64 {
self.lifetime_classified.load(Ordering::Relaxed)
}
pub fn lifetime_unknown(&self) -> u64 {
self.lifetime_unknown.load(Ordering::Relaxed)
}
pub fn lifetime_ambiguous(&self) -> u64 {
self.lifetime_ambiguous.load(Ordering::Relaxed)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::event::JsonEvent;
use serde_json::json;
fn classify(value: &serde_json::Value) -> Option<String> {
SchemaClassifier::builtin()
.classify(&JsonEvent::borrow(value))
.map(|m| m.name)
}
#[test]
fn recognizes_ecs_by_version_marker() {
let v = json!({"ecs": {"version": "8.11.0"}, "process": {"command_line": "whoami"}});
assert_eq!(classify(&v).as_deref(), Some("ecs"));
}
#[test]
fn recognizes_ecs_with_flattened_keys() {
let v = json!({"ecs.version": "8.11.0", "process.command_line": "whoami"});
assert_eq!(classify(&v).as_deref(), Some("ecs"));
}
#[test]
fn recognizes_ocsf_by_class_and_metadata() {
let v = json!({"class_uid": 1001, "category_uid": 1, "metadata": {"version": "1.1.0"}});
assert_eq!(classify(&v).as_deref(), Some("ocsf"));
}
#[test]
fn recognizes_rendered_windows_event_log() {
let v = json!({"Event": {"System": {"EventID": 4688, "Provider": "Microsoft-Windows-Security-Auditing"}}});
assert_eq!(classify(&v).as_deref(), Some("windows_eventlog"));
}
#[test]
fn recognizes_sysmon_by_channel() {
let v = json!({"Channel": "Microsoft-Windows-Sysmon/Operational", "EventID": 1, "Image": "C:/cmd.exe"});
assert_eq!(classify(&v).as_deref(), Some("sysmon"));
}
#[test]
fn recognizes_sysmon_by_provider() {
let v = json!({"Provider_Name": "Microsoft-Windows-Sysmon", "EventID": 3});
assert_eq!(classify(&v).as_deref(), Some("sysmon"));
}
#[test]
fn recognizes_flat_sysmon_by_field_shape() {
let v = json!({"EventID": 1, "ProcessGuid": "{abc}", "CommandLine": "cmd /c whoami"});
assert_eq!(classify(&v).as_deref(), Some("sysmon"));
}
#[test]
fn recognizes_cef_structured_fields() {
let v = json!({"deviceVendor": "Security", "deviceProduct": "IDS", "signatureId": "100", "src": "10.0.0.1"});
assert_eq!(classify(&v).as_deref(), Some("cef"));
}
#[test]
fn falls_back_to_generic_json_for_unrecognized_structured_events() {
let v = json!({"some_vendor_field": "x", "another": 1});
assert_eq!(classify(&v).as_deref(), Some("generic_json"));
}
#[test]
fn fieldless_events_are_unknown() {
assert_eq!(classify(&json!({})), None);
assert_eq!(classify(&json!("just a string")), None);
}
#[test]
fn specificity_prefers_specific_schema_over_generic() {
let v = json!({"ecs.version": "8.0.0", "vendor_blob": {"x": 1}});
let cls = SchemaClassifier::builtin();
let m = cls.classify(&JsonEvent::borrow(&v)).unwrap();
assert_eq!(m.name, "ecs");
assert_eq!(m.specificity, 100);
let all = cls.classify_all(&JsonEvent::borrow(&v));
assert_eq!(all.first().map(String::as_str), Some("ecs"));
assert!(all.iter().any(|n| n == "generic_json"));
}
#[test]
fn schema_names_lists_builtins_most_specific_first() {
let classifier = SchemaClassifier::builtin();
let names = classifier.schema_names();
assert_eq!(names.first(), Some(&"ecs_linux"));
assert!(names.contains(&"ecs_windows"));
assert!(names.contains(&"ecs"));
assert!(names.contains(&"generic_json"));
assert_eq!(names.last(), Some(&"generic_json"));
}
#[test]
fn ecs_windows_specialization_classifies_and_aliases_to_ecs() {
let v = json!({"ecs.version": "8.11.0", "winlog": {"channel": "Security"}});
assert_eq!(classify(&v).as_deref(), Some("ecs_windows"));
let plain = json!({"ecs.version": "8.11.0", "process": {"command_line": "whoami"}});
assert_eq!(classify(&plain).as_deref(), Some("ecs"));
let config = RoutingConfig {
on_unknown: OnUnknown::Warn,
default_pipelines: vec![],
aliases: HashMap::new(),
bindings: vec![SchemaBinding {
schema: "ecs".to_string(),
pipelines: vec!["ecs_windows".to_string()],
logsource: None,
}],
};
let plan = RoutingPlan::from_config(&config);
let ecs_set = match plan.decide(Some("ecs")) {
RouteDecision::Evaluate { set, .. } => set,
other => panic!("unexpected: {other:?}"),
};
assert_eq!(plan.decide(Some("ecs_windows")), plan.decide(Some("ecs")));
assert_ne!(ecs_set, 0, "ecs binding is a non-default set");
assert_eq!(
plan.schema_logsource("ecs_windows")
.and_then(|l| l.product.as_deref()),
Some("windows")
);
}
#[test]
fn user_alias_routes_as_canonical() {
let yaml = r#"
schemas:
- name: my_win
specificity: 70
match:
- field_present: vendor.win_marker
routing:
aliases:
my_win: ecs
bindings:
- schema: ecs
pipelines: [ecs_windows]
"#;
let (_sigs, routing) = parse_schema_config(yaml).unwrap();
let plan = RoutingPlan::from_config(&routing.expect("routing"));
assert_eq!(plan.decide(Some("my_win")), plan.decide(Some("ecs")));
assert!(matches!(
plan.decide(Some("my_win")),
RouteDecision::Evaluate { unknown: false, .. }
));
}
#[test]
fn set_product_partition_only_for_platform_locked_sets() {
let config = RoutingConfig {
on_unknown: OnUnknown::Warn,
default_pipelines: vec![],
aliases: HashMap::new(),
bindings: vec![
SchemaBinding {
schema: "sysmon".to_string(),
pipelines: vec!["p_sysmon".to_string()],
logsource: None,
},
SchemaBinding {
schema: "ecs".to_string(),
pipelines: vec!["p_ecs".to_string()],
logsource: None,
},
],
};
let plan = RoutingPlan::from_config(&config);
let part = plan.set_product_partition();
assert!(part[0].is_none(), "default set is never partitioned");
let set_of = |schema| match plan.decide(Some(schema)) {
RouteDecision::Evaluate { set, .. } => set,
other => panic!("unexpected: {other:?}"),
};
let sysmon_set = set_of("sysmon");
assert_eq!(
part[sysmon_set].as_ref().map(|s| s.contains("windows")),
Some(true)
);
assert!(part[set_of("ecs")].is_none());
}
#[test]
fn parses_user_signatures_from_yaml() {
let yaml = r#"
schemas:
- name: my_vendor
specificity: 70
match:
- field_present: vendor.product
- equals:
field: event_type
value: alert
- any_of: [a, b]
"#;
let sigs = parse_schema_signatures(yaml).expect("parse");
assert_eq!(sigs.len(), 1);
assert_eq!(sigs[0].name, "my_vendor");
assert_eq!(sigs[0].specificity, 70);
assert_eq!(sigs[0].predicates.len(), 3);
let cls = SchemaClassifier::with_user_signatures(sigs);
let v = json!({"vendor": {"product": "X"}, "event_type": "ALERT", "a": 1});
assert_eq!(
cls.classify(&JsonEvent::borrow(&v))
.map(|m| m.name)
.as_deref(),
Some("my_vendor")
);
}
#[test]
fn user_signature_with_invalid_regex_is_rejected() {
let yaml = r#"
schemas:
- name: bad
match:
- matches:
field: msg
value: "([unclosed"
"#;
let err = parse_schema_signatures(yaml).unwrap_err();
assert!(matches!(err, SchemaError::InvalidRegex { .. }));
}
#[test]
fn user_regex_signature_matches_field_value() {
let yaml = r#"
schemas:
- name: cef_raw
specificity: 60
match:
- matches:
field: message
value: "^CEF:\\d"
"#;
let sigs = parse_schema_signatures(yaml).expect("parse");
let cls = SchemaClassifier::with_user_signatures(sigs);
let v = json!({"message": "CEF:0|Vendor|Product|1.0|100|Name|9|src=1.2.3.4"});
assert_eq!(
cls.classify(&JsonEvent::borrow(&v))
.map(|m| m.name)
.as_deref(),
Some("cef_raw")
);
}
fn classifier_from_match(match_body: &str) -> SchemaClassifier {
let yaml = format!("schemas:\n - name: t\n specificity: 70\n match:\n{match_body}");
let sigs = parse_schema_signatures(&yaml).expect("parse");
SchemaClassifier::new(sigs)
}
fn matches_t(match_body: &str, event: &serde_json::Value) -> bool {
classifier_from_match(match_body)
.classify(&JsonEvent::borrow(event))
.is_some()
}
#[test]
fn numeric_comparisons() {
let body = " - gte: { field: EventID, value: 4000 }\n";
assert!(matches_t(body, &json!({"EventID": 4688})));
assert!(matches_t(body, &json!({"EventID": 4000})));
assert!(!matches_t(body, &json!({"EventID": 1})));
assert!(matches_t(body, &json!({"EventID": "4688"})));
assert!(!matches_t(body, &json!({"EventID": "not-a-number"})));
assert!(matches_t(
" - lt: { field: score, value: 10 }\n",
&json!({"score": 9.5})
));
assert!(matches_t(
" - gt: { field: score, value: 10 }\n",
&json!({"score": 10.1})
));
}
#[test]
fn in_set_membership_is_case_insensitive() {
let body = " - in: { field: event_type, values: [alert, alarm] }\n";
assert!(matches_t(body, &json!({"event_type": "ALERT"})));
assert!(matches_t(body, &json!({"event_type": "alarm"})));
assert!(!matches_t(body, &json!({"event_type": "info"})));
assert!(!matches_t(body, &json!({})));
}
#[test]
fn field_equals_field_compares_two_fields() {
let body = " - field_equals_field: { left: a, right: b }\n";
assert!(matches_t(body, &json!({"a": "X", "b": "x"})));
assert!(!matches_t(body, &json!({"a": "X", "b": "y"})));
assert!(!matches_t(body, &json!({"a": "X"})));
}
#[test]
fn recursive_not_any_all_groups() {
let any_body = " - any:\n - field_present: winlog.channel\n - equals: { field: host.os.type, value: windows }\n";
assert!(matches_t(
any_body,
&json!({"winlog": {"channel": "Security"}})
));
assert!(matches_t(
any_body,
&json!({"host": {"os": {"type": "windows"}}})
));
assert!(!matches_t(any_body, &json!({"unrelated": 1})));
let not_body = " - not: { field_present: ecs.version }\n";
assert!(matches_t(not_body, &json!({"CommandLine": "whoami"})));
assert!(!matches_t(not_body, &json!({"ecs.version": "8.0.0"})));
let all_body = " - all:\n - field_present: a\n - field_present: b\n";
assert!(matches_t(all_body, &json!({"a": 1, "b": 2})));
assert!(!matches_t(all_body, &json!({"a": 1})));
}
#[test]
fn empty_group_is_rejected() {
let yaml = "schemas:\n - name: t\n match:\n - any: []\n";
let err = parse_schema_signatures(yaml).unwrap_err();
assert!(
matches!(&err, SchemaError::Parse(m) if m.contains("'any' needs at least one")),
"got: {err}"
);
}
#[test]
fn predicate_with_two_conditions_is_rejected() {
let yaml = "schemas:\n - name: t\n match:\n - field_present: a\n field_absent: b\n";
let err = parse_schema_signatures(yaml).unwrap_err();
assert!(
matches!(&err, SchemaError::Parse(m) if m.contains("multiple conditions")),
"got: {err}"
);
}
#[test]
fn explain_reports_matched_signature() {
let cls = SchemaClassifier::builtin();
let v = json!({"ecs.version": "8.0.0"});
let ex = cls.explain(&JsonEvent::borrow(&v));
assert_eq!(ex.matched.as_deref(), Some("ecs"));
let sig = ex.signature.expect("signature");
assert!(sig.predicates_matched);
assert!(sig.predicates.iter().all(|p| p.matched));
}
#[test]
fn explain_reports_near_miss_for_unknown() {
let sigs = builtin_signatures()
.into_iter()
.filter(|s| s.name != "generic_json")
.collect();
let cls = SchemaClassifier::new(sigs);
let v = json!({"EventID": 1, "Image": "C:/cmd.exe"});
let ex = cls.explain(&JsonEvent::borrow(&v));
assert_eq!(ex.matched, None);
let sig = ex.signature.expect("near-miss");
assert_eq!(sig.name, "sysmon");
assert!(!sig.predicates_matched);
assert!(sig.predicates.iter().any(|p| !p.matched));
}
#[test]
fn validate_flags_unknown_binding_and_shadow() {
let yaml = r#"
schemas:
- name: shadowed
specificity: 40
match:
- field_present: ecs.version
- field_present: extra.marker
routing:
bindings:
- schema: ecs
pipelines: [ecs_windows]
- schema: nonexistent
pipelines: [x]
"#;
let (sigs, routing) = parse_schema_config(yaml).unwrap();
let findings = validate_schema_config(&sigs, routing.as_ref());
assert!(
findings
.iter()
.any(|f| f.contains("unknown schema 'nonexistent'")),
"findings: {findings:?}"
);
assert!(
findings
.iter()
.any(|f| f.contains("'shadowed'") && f.contains("unreachable")),
"findings: {findings:?}"
);
}
#[test]
fn observer_counts_per_schema_and_unknown() {
let observer = SchemaObserver::builtin();
observer.observe(&JsonEvent::borrow(&json!({"ecs.version": "8.0.0"})));
observer.observe(&JsonEvent::borrow(&json!({"ecs.version": "8.1.0"})));
observer.observe(&JsonEvent::borrow(
&json!({"class_uid": 1001, "metadata": {"version": "1.1.0"}}),
));
observer.observe(&JsonEvent::borrow(&json!({})));
let snap = observer.snapshot();
assert_eq!(snap.events_observed, 4);
assert_eq!(snap.classified, 3);
assert_eq!(snap.unknown, 1);
assert_eq!(snap.by_schema[0].schema, "ecs");
assert_eq!(snap.by_schema[0].count, 2);
let ocsf = snap.by_schema.iter().find(|e| e.schema == "ocsf").unwrap();
assert_eq!(ocsf.count, 1);
}
#[test]
fn routing_plan_dedups_pipeline_sets() {
let config = RoutingConfig {
on_unknown: OnUnknown::Warn,
default_pipelines: vec![],
aliases: HashMap::new(),
bindings: vec![
SchemaBinding {
schema: "ecs".to_string(),
pipelines: vec!["ecs_windows".to_string()],
logsource: None,
},
SchemaBinding {
schema: "winlogbeat".to_string(),
pipelines: vec!["ecs_windows".to_string()],
logsource: None,
},
SchemaBinding {
schema: "sysmon".to_string(),
pipelines: vec!["sysmon".to_string()],
logsource: None,
},
],
};
let plan = RoutingPlan::from_config(&config);
assert_eq!(plan.pipeline_sets().len(), 3);
let ecs = plan.decide(Some("ecs"));
let win = plan.decide(Some("winlogbeat"));
assert_eq!(ecs, win);
assert!(matches!(
ecs,
RouteDecision::Evaluate { unknown: false, .. }
));
assert_ne!(plan.decide(Some("sysmon")), ecs);
}
#[test]
fn routing_decides_bound_unbound_and_unknown() {
let config = RoutingConfig {
on_unknown: OnUnknown::Warn,
default_pipelines: vec![],
aliases: HashMap::new(),
bindings: vec![SchemaBinding {
schema: "ecs".to_string(),
pipelines: vec!["ecs_windows".to_string()],
logsource: None,
}],
};
let plan = RoutingPlan::from_config(&config);
assert!(matches!(
plan.decide(Some("ecs")),
RouteDecision::Evaluate { unknown: false, .. }
));
assert_eq!(
plan.decide(Some("cef")),
RouteDecision::Evaluate {
set: 0,
unknown: false
}
);
assert_eq!(
plan.decide(None),
RouteDecision::Evaluate {
set: 0,
unknown: true
}
);
}
#[test]
fn routing_on_unknown_policies() {
let base = |policy| RoutingConfig {
on_unknown: policy,
default_pipelines: vec![],
aliases: HashMap::new(),
bindings: vec![],
};
assert_eq!(
RoutingPlan::from_config(&base(OnUnknown::Drop)).decide(None),
RouteDecision::Drop
);
assert_eq!(
RoutingPlan::from_config(&base(OnUnknown::Error)).decide(None),
RouteDecision::Error
);
assert_eq!(
RoutingPlan::from_config(&base(OnUnknown::Passthrough)).decide(None),
RouteDecision::Evaluate {
set: 0,
unknown: true
}
);
}
#[test]
fn parses_routing_section_from_yaml() {
let yaml = r#"
schemas:
- name: my_vendor
match:
- field_present: vendor.id
routing:
on_unknown: drop
default_pipelines: [base]
bindings:
- schema: ecs
pipelines: [ecs_windows]
- schema: my_vendor
pipelines: [vendor_map, base]
"#;
let (sigs, routing) = parse_schema_config(yaml).expect("parse");
assert_eq!(sigs.len(), 1);
let routing = routing.expect("routing present");
assert_eq!(routing.on_unknown, OnUnknown::Drop);
assert_eq!(routing.default_pipelines, vec!["base".to_string()]);
assert_eq!(routing.bindings.len(), 2);
let plan = RoutingPlan::from_config(&routing);
assert_eq!(plan.pipeline_sets().len(), 3);
assert_eq!(plan.decide(None), RouteDecision::Drop);
}
#[test]
fn schema_logsource_builtin_defaults_and_overrides() {
let plan = RoutingPlan::from_config(&RoutingConfig::default());
let sysmon = plan.schema_logsource("sysmon").expect("sysmon default");
assert_eq!(sysmon.product.as_deref(), Some("windows"));
assert_eq!(sysmon.service.as_deref(), Some("sysmon"));
assert_eq!(
plan.schema_logsource("windows_eventlog")
.and_then(|l| l.product.as_deref()),
Some("windows")
);
assert!(plan.schema_logsource("ecs").is_none());
assert!(plan.schema_logsource("cef").is_none());
let yaml = r#"
schemas:
- name: ecs_windows
match:
- field_present: ecs.version
- field_present: winlog.channel
routing:
bindings:
- schema: ecs_windows
pipelines: [ecs_windows]
logsource:
product: windows
- schema: sysmon
pipelines: [sysmon]
logsource:
product: windows
service: sysmon
custom:
tenant: acme
"#;
let (_sigs, routing) = parse_schema_config(yaml).expect("parse");
let plan = RoutingPlan::from_config(&routing.expect("routing"));
assert_eq!(
plan.schema_logsource("ecs_windows")
.and_then(|l| l.product.as_deref()),
Some("windows")
);
let sysmon = plan.schema_logsource("sysmon").expect("sysmon override");
assert_eq!(
sysmon.custom.get("tenant").map(String::as_str),
Some("acme")
);
}
#[test]
fn observer_reset_preserves_lifetime_counters() {
let observer = SchemaObserver::builtin();
observer.observe(&JsonEvent::borrow(&json!({"ecs.version": "8.0.0"})));
observer.observe(&JsonEvent::borrow(&json!({})));
let (classified, unknown) = observer.reset();
assert_eq!(classified, 1);
assert_eq!(unknown, 1);
let snap = observer.snapshot();
assert_eq!(snap.classified, 0);
assert_eq!(snap.unknown, 0);
assert_eq!(snap.events_observed, 0);
assert_eq!(snap.lifetime_classified, 1);
assert_eq!(snap.lifetime_unknown, 1);
}
#[test]
fn classify_with_ambiguity_flags_equal_specificity_ties() {
let sigs = vec![
SchemaSignature {
name: "alpha".to_string(),
specificity: 70,
predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
},
SchemaSignature {
name: "beta".to_string(),
specificity: 70,
predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
},
];
let cls = SchemaClassifier::new(sigs);
let (m, ambiguous) = cls.classify_with_ambiguity(&JsonEvent::borrow(&json!({"a": 1})));
assert!(m.is_some());
assert!(
ambiguous,
"equal-specificity different-name match is ambiguous"
);
let cls = SchemaClassifier::builtin();
let (_, ambiguous) =
cls.classify_with_ambiguity(&JsonEvent::borrow(&json!({"ecs.version": "8.0.0"})));
assert!(!ambiguous);
}
#[test]
fn observer_records_ambiguity_and_unknown_shapes() {
let sigs = vec![
SchemaSignature {
name: "alpha".to_string(),
specificity: 70,
predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
},
SchemaSignature {
name: "beta".to_string(),
specificity: 70,
predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
},
];
let observer = SchemaObserver::new(SchemaClassifier::new(sigs));
observer.observe(&JsonEvent::borrow(&json!({"a": 1}))); observer.observe(&JsonEvent::borrow(&json!({"weird": 1, "shape": 2}))); observer.observe(&JsonEvent::borrow(&json!({"shape": 3, "weird": 4})));
let snap = observer.snapshot();
assert_eq!(snap.ambiguous, 1);
assert_eq!(snap.unknown, 2);
assert_eq!(snap.unknown_shapes.len(), 1);
assert_eq!(snap.unknown_shapes[0].count, 2);
assert_eq!(snap.unknown_shapes[0].keys, vec!["shape", "weird"]);
}
}