1use std::collections::HashMap;
27use std::fs;
28use std::path::Path;
29use std::sync::Mutex;
30use std::sync::atomic::{AtomicU64, Ordering};
31use std::time::Instant;
32
33use regex::Regex;
34use rsigma_parser::LogSource;
35use serde::{Deserialize, Serialize};
36
37use crate::event::Event;
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq)]
41pub enum CompareOp {
42 Gt,
44 Gte,
46 Lt,
48 Lte,
50}
51
52impl CompareOp {
53 fn apply(self, lhs: f64, rhs: f64) -> bool {
54 match self {
55 CompareOp::Gt => lhs > rhs,
56 CompareOp::Gte => lhs >= rhs,
57 CompareOp::Lt => lhs < rhs,
58 CompareOp::Lte => lhs <= rhs,
59 }
60 }
61
62 fn symbol(self) -> &'static str {
63 match self {
64 CompareOp::Gt => ">",
65 CompareOp::Gte => ">=",
66 CompareOp::Lt => "<",
67 CompareOp::Lte => "<=",
68 }
69 }
70}
71
72#[derive(Debug, Clone)]
78pub enum SchemaPredicate {
79 FieldPresent(String),
81 FieldAbsent(String),
83 AnyOf(Vec<String>),
85 Equals { field: String, value: String },
88 Matches { field: String, regex: Regex },
90 Compare {
93 field: String,
94 op: CompareOp,
95 value: f64,
96 },
97 In { field: String, values: Vec<String> },
100 FieldEqualsField { left: String, right: String },
102 Not(Box<SchemaPredicate>),
104 Any(Vec<SchemaPredicate>),
106 All(Vec<SchemaPredicate>),
109 HasAnyField,
113}
114
115impl SchemaPredicate {
116 fn eval<E: Event + ?Sized>(&self, event: &E) -> bool {
117 match self {
118 SchemaPredicate::FieldPresent(f) => event.get_field(f).is_some(),
119 SchemaPredicate::FieldAbsent(f) => event.get_field(f).is_none(),
120 SchemaPredicate::AnyOf(fields) => fields.iter().any(|f| event.get_field(f).is_some()),
121 SchemaPredicate::Equals { field, value } => event
122 .get_field(field)
123 .and_then(|v| v.as_str().map(|s| s.as_ref().eq_ignore_ascii_case(value)))
124 .unwrap_or(false),
125 SchemaPredicate::Matches { field, regex } => event
126 .get_field(field)
127 .and_then(|v| v.as_str().map(|s| regex.is_match(s.as_ref())))
128 .unwrap_or(false),
129 SchemaPredicate::Compare { field, op, value } => event
130 .get_field(field)
131 .and_then(|v| v.as_f64())
132 .map(|n| op.apply(n, *value))
133 .unwrap_or(false),
134 SchemaPredicate::In { field, values } => event
135 .get_field(field)
136 .and_then(|v| {
137 v.as_str().map(|s| {
138 values
139 .iter()
140 .any(|val| s.as_ref().eq_ignore_ascii_case(val))
141 })
142 })
143 .unwrap_or(false),
144 SchemaPredicate::FieldEqualsField { left, right } => {
145 let l = event
146 .get_field(left)
147 .and_then(|v| v.as_str().map(|s| s.into_owned()));
148 let r = event
149 .get_field(right)
150 .and_then(|v| v.as_str().map(|s| s.into_owned()));
151 matches!((l, r), (Some(a), Some(b)) if a.eq_ignore_ascii_case(&b))
152 }
153 SchemaPredicate::Not(inner) => !inner.eval(event),
154 SchemaPredicate::Any(preds) => preds.iter().any(|p| p.eval(event)),
155 SchemaPredicate::All(preds) => preds.iter().all(|p| p.eval(event)),
156 SchemaPredicate::HasAnyField => !event.field_keys().is_empty(),
157 }
158 }
159
160 fn describe(&self) -> String {
162 match self {
163 SchemaPredicate::FieldPresent(f) => format!("field_present({f})"),
164 SchemaPredicate::FieldAbsent(f) => format!("field_absent({f})"),
165 SchemaPredicate::AnyOf(fs) => format!("any_of([{}])", fs.join(", ")),
166 SchemaPredicate::Equals { field, value } => format!("{field} == \"{value}\""),
167 SchemaPredicate::Matches { field, regex } => {
168 format!("{field} matches /{}/", regex.as_str())
169 }
170 SchemaPredicate::Compare { field, op, value } => {
171 format!("{field} {} {value}", op.symbol())
172 }
173 SchemaPredicate::In { field, values } => format!("{field} in [{}]", values.join(", ")),
174 SchemaPredicate::FieldEqualsField { left, right } => format!("{left} == {right}"),
175 SchemaPredicate::Not(inner) => format!("not({})", inner.describe()),
176 SchemaPredicate::Any(ps) => format!(
177 "any({})",
178 ps.iter()
179 .map(|p| p.describe())
180 .collect::<Vec<_>>()
181 .join(" | ")
182 ),
183 SchemaPredicate::All(ps) => format!(
184 "all({})",
185 ps.iter()
186 .map(|p| p.describe())
187 .collect::<Vec<_>>()
188 .join(" & ")
189 ),
190 SchemaPredicate::HasAnyField => "has_any_field".to_string(),
191 }
192 }
193}
194
195#[derive(Debug, Clone)]
200pub struct SchemaSignature {
201 pub name: String,
203 pub predicates: Vec<SchemaPredicate>,
207 pub specificity: u32,
209}
210
211impl SchemaSignature {
212 fn matches<E: Event + ?Sized>(&self, event: &E) -> bool {
213 self.predicates.iter().all(|p| p.eval(event))
214 }
215
216 fn explain<E: Event + ?Sized>(&self, event: &E) -> SignatureExplanation {
217 let predicates: Vec<PredicateOutcome> = self
218 .predicates
219 .iter()
220 .map(|p| PredicateOutcome {
221 predicate: p.describe(),
222 matched: p.eval(event),
223 })
224 .collect();
225 let predicates_matched = predicates.iter().all(|p| p.matched);
226 SignatureExplanation {
227 name: self.name.clone(),
228 specificity: self.specificity,
229 predicates_matched,
230 predicates,
231 }
232 }
233}
234
235#[derive(Debug, Clone, Serialize)]
237pub struct PredicateOutcome {
238 pub predicate: String,
240 pub matched: bool,
242}
243
244#[derive(Debug, Clone, Serialize)]
246pub struct SignatureExplanation {
247 pub name: String,
249 pub specificity: u32,
251 pub predicates_matched: bool,
253 pub predicates: Vec<PredicateOutcome>,
255}
256
257#[derive(Debug, Clone, Serialize)]
262pub struct SchemaExplanation {
263 pub matched: Option<String>,
265 pub specificity: Option<u32>,
267 pub signature: Option<SignatureExplanation>,
270}
271
272#[derive(Debug, Clone, PartialEq, Eq)]
275pub struct SchemaMatch {
276 pub name: String,
277 pub specificity: u32,
278}
279
280#[derive(Debug, Clone)]
286pub struct SchemaClassifier {
287 signatures: Vec<SchemaSignature>,
288}
289
290impl SchemaClassifier {
291 pub fn new(mut signatures: Vec<SchemaSignature>) -> Self {
293 signatures.sort_by(|a, b| {
294 b.specificity
295 .cmp(&a.specificity)
296 .then_with(|| a.name.cmp(&b.name))
297 });
298 Self { signatures }
299 }
300
301 pub fn builtin() -> Self {
303 Self::new(builtin_signatures())
304 }
305
306 pub fn with_user_signatures(user: Vec<SchemaSignature>) -> Self {
310 let mut signatures = builtin_signatures();
311 signatures.extend(user);
312 Self::new(signatures)
313 }
314
315 pub fn classify<E: Event + ?Sized>(&self, event: &E) -> Option<SchemaMatch> {
318 self.signatures
319 .iter()
320 .find(|s| s.matches(event))
321 .map(|s| SchemaMatch {
322 name: s.name.clone(),
323 specificity: s.specificity,
324 })
325 }
326
327 pub fn classify_with_ambiguity<E: Event + ?Sized>(
333 &self,
334 event: &E,
335 ) -> (Option<SchemaMatch>, bool) {
336 let mut matching = self.signatures.iter().filter(|s| s.matches(event));
340 let Some(winner) = matching.next() else {
341 return (None, false);
342 };
343 let ambiguous = matching
344 .take_while(|s| s.specificity == winner.specificity)
345 .any(|s| s.name != winner.name);
346 (
347 Some(SchemaMatch {
348 name: winner.name.clone(),
349 specificity: winner.specificity,
350 }),
351 ambiguous,
352 )
353 }
354
355 pub fn classify_all<E: Event + ?Sized>(&self, event: &E) -> Vec<String> {
359 let mut out: Vec<String> = Vec::new();
360 for sig in self.signatures.iter().filter(|s| s.matches(event)) {
361 if !out.iter().any(|n| n == &sig.name) {
362 out.push(sig.name.clone());
363 }
364 }
365 out
366 }
367
368 pub fn explain<E: Event + ?Sized>(&self, event: &E) -> SchemaExplanation {
373 let mut best_near: Option<SignatureExplanation> = None;
374 let mut best_near_passing = 0usize;
375 for sig in &self.signatures {
376 let ex = sig.explain(event);
377 if ex.predicates_matched {
378 return SchemaExplanation {
379 matched: Some(ex.name.clone()),
380 specificity: Some(ex.specificity),
381 signature: Some(ex),
382 };
383 }
384 let passing = ex.predicates.iter().filter(|p| p.matched).count();
387 if best_near.is_none() || passing > best_near_passing {
388 best_near_passing = passing;
389 best_near = Some(ex);
390 }
391 }
392 SchemaExplanation {
393 matched: None,
394 specificity: None,
395 signature: best_near,
396 }
397 }
398
399 pub fn schema_names(&self) -> Vec<&str> {
401 let mut out: Vec<&str> = Vec::new();
402 for sig in &self.signatures {
403 if !out.contains(&sig.name.as_str()) {
404 out.push(sig.name.as_str());
405 }
406 }
407 out
408 }
409}
410
411impl Default for SchemaClassifier {
412 fn default() -> Self {
413 Self::builtin()
414 }
415}
416
417fn builtin_signatures() -> Vec<SchemaSignature> {
421 vec![
422 SchemaSignature {
427 name: "ecs_windows".to_string(),
428 specificity: 105,
429 predicates: vec![
430 SchemaPredicate::FieldPresent("ecs.version".to_string()),
431 SchemaPredicate::Any(vec![
432 SchemaPredicate::FieldPresent("winlog.channel".to_string()),
433 SchemaPredicate::FieldPresent("winlog.event_id".to_string()),
434 SchemaPredicate::Equals {
435 field: "host.os.type".to_string(),
436 value: "windows".to_string(),
437 },
438 SchemaPredicate::Equals {
439 field: "os.type".to_string(),
440 value: "windows".to_string(),
441 },
442 ]),
443 ],
444 },
445 SchemaSignature {
448 name: "ecs_linux".to_string(),
449 specificity: 105,
450 predicates: vec![
451 SchemaPredicate::FieldPresent("ecs.version".to_string()),
452 SchemaPredicate::Any(vec![
453 SchemaPredicate::Equals {
454 field: "host.os.type".to_string(),
455 value: "linux".to_string(),
456 },
457 SchemaPredicate::Equals {
458 field: "os.type".to_string(),
459 value: "linux".to_string(),
460 },
461 SchemaPredicate::FieldPresent("host.os.kernel".to_string()),
462 ]),
463 ],
464 },
465 SchemaSignature {
467 name: "ecs".to_string(),
468 specificity: 100,
469 predicates: vec![SchemaPredicate::FieldPresent("ecs.version".to_string())],
470 },
471 SchemaSignature {
473 name: "ocsf".to_string(),
474 specificity: 95,
475 predicates: vec![
476 SchemaPredicate::FieldPresent("class_uid".to_string()),
477 SchemaPredicate::FieldPresent("metadata.version".to_string()),
478 ],
479 },
480 SchemaSignature {
482 name: "windows_eventlog".to_string(),
483 specificity: 90,
484 predicates: vec![SchemaPredicate::AnyOf(vec![
485 "Event.System.EventID".to_string(),
486 "Event.System.Provider".to_string(),
487 ])],
488 },
489 SchemaSignature {
491 name: "sysmon".to_string(),
492 specificity: 88,
493 predicates: vec![SchemaPredicate::Equals {
494 field: "Channel".to_string(),
495 value: "Microsoft-Windows-Sysmon/Operational".to_string(),
496 }],
497 },
498 SchemaSignature {
500 name: "sysmon".to_string(),
501 specificity: 88,
502 predicates: vec![SchemaPredicate::Equals {
503 field: "Provider_Name".to_string(),
504 value: "Microsoft-Windows-Sysmon".to_string(),
505 }],
506 },
507 SchemaSignature {
509 name: "sysmon".to_string(),
510 specificity: 80,
511 predicates: vec![
512 SchemaPredicate::FieldPresent("EventID".to_string()),
513 SchemaPredicate::FieldPresent("ProcessGuid".to_string()),
514 SchemaPredicate::AnyOf(vec!["Image".to_string(), "CommandLine".to_string()]),
515 ],
516 },
517 SchemaSignature {
520 name: "cef".to_string(),
521 specificity: 85,
522 predicates: vec![
523 SchemaPredicate::FieldPresent("deviceVendor".to_string()),
524 SchemaPredicate::FieldPresent("deviceProduct".to_string()),
525 SchemaPredicate::FieldPresent("signatureId".to_string()),
526 ],
527 },
528 SchemaSignature {
530 name: "generic_json".to_string(),
531 specificity: 0,
532 predicates: vec![SchemaPredicate::HasAnyField],
533 },
534 ]
535}
536
537pub fn builtin_schema_names() -> Vec<&'static str> {
539 vec![
540 "ecs_windows",
541 "ecs_linux",
542 "ecs",
543 "ocsf",
544 "windows_eventlog",
545 "sysmon",
546 "cef",
547 "generic_json",
548 ]
549}
550
551fn builtin_schema_aliases() -> HashMap<String, String> {
557 HashMap::from([
558 ("ecs_windows".to_string(), "ecs".to_string()),
559 ("ecs_linux".to_string(), "ecs".to_string()),
560 ])
561}
562
563#[derive(Debug, thiserror::Error)]
569pub enum SchemaError {
570 #[error("cannot read schema signatures file '{path}': {source}")]
572 Io {
573 path: String,
574 #[source]
575 source: std::io::Error,
576 },
577 #[error("schema signatures YAML parse error: {0}")]
579 Parse(String),
580 #[error("invalid regex in schema '{name}': {error}")]
582 InvalidRegex { name: String, error: String },
583}
584
585#[derive(Debug, Clone, Deserialize)]
588#[serde(deny_unknown_fields)]
589pub struct FieldValueConfig {
590 pub field: String,
591 pub value: String,
592}
593
594#[derive(Debug, Clone, Deserialize)]
597#[serde(deny_unknown_fields)]
598pub struct FieldNumberConfig {
599 pub field: String,
600 pub value: f64,
601}
602
603#[derive(Debug, Clone, Deserialize)]
605#[serde(deny_unknown_fields)]
606pub struct FieldValuesConfig {
607 pub field: String,
608 pub values: Vec<String>,
609}
610
611#[derive(Debug, Clone, Deserialize)]
613#[serde(deny_unknown_fields)]
614pub struct FieldPairConfig {
615 pub left: String,
616 pub right: String,
617}
618
619#[derive(Debug, Clone, Default, Deserialize)]
624#[serde(deny_unknown_fields)]
625pub struct SchemaPredicateConfig {
626 #[serde(default)]
628 pub field_present: Option<String>,
629 #[serde(default)]
631 pub field_absent: Option<String>,
632 #[serde(default)]
634 pub any_of: Option<Vec<String>>,
635 #[serde(default)]
637 pub equals: Option<FieldValueConfig>,
638 #[serde(default)]
640 pub matches: Option<FieldValueConfig>,
641 #[serde(default)]
643 pub gt: Option<FieldNumberConfig>,
644 #[serde(default)]
646 pub gte: Option<FieldNumberConfig>,
647 #[serde(default)]
649 pub lt: Option<FieldNumberConfig>,
650 #[serde(default)]
652 pub lte: Option<FieldNumberConfig>,
653 #[serde(default, rename = "in")]
655 pub in_set: Option<FieldValuesConfig>,
656 #[serde(default)]
658 pub field_equals_field: Option<FieldPairConfig>,
659 #[serde(default)]
661 pub not: Option<Box<SchemaPredicateConfig>>,
662 #[serde(default)]
664 pub any: Option<Vec<SchemaPredicateConfig>>,
665 #[serde(default)]
667 pub all: Option<Vec<SchemaPredicateConfig>>,
668}
669
670impl SchemaPredicateConfig {
671 fn build(self, schema_name: &str) -> Result<SchemaPredicate, SchemaError> {
672 let mut chosen: Option<SchemaPredicate> = None;
673 let mut set = 0u32;
674 if let Some(f) = self.field_present {
675 set += 1;
676 chosen = Some(SchemaPredicate::FieldPresent(f));
677 }
678 if let Some(f) = self.field_absent {
679 set += 1;
680 chosen = Some(SchemaPredicate::FieldAbsent(f));
681 }
682 if let Some(fields) = self.any_of {
683 set += 1;
684 chosen = Some(SchemaPredicate::AnyOf(fields));
685 }
686 if let Some(fv) = self.equals {
687 set += 1;
688 chosen = Some(SchemaPredicate::Equals {
689 field: fv.field,
690 value: fv.value,
691 });
692 }
693 if let Some(fv) = self.matches {
694 set += 1;
695 chosen = Some(SchemaPredicate::Matches {
696 field: fv.field,
697 regex: Regex::new(&fv.value).map_err(|e| SchemaError::InvalidRegex {
698 name: schema_name.to_string(),
699 error: e.to_string(),
700 })?,
701 });
702 }
703 for (op, cfg) in [
704 (CompareOp::Gt, self.gt),
705 (CompareOp::Gte, self.gte),
706 (CompareOp::Lt, self.lt),
707 (CompareOp::Lte, self.lte),
708 ] {
709 if let Some(fv) = cfg {
710 set += 1;
711 chosen = Some(SchemaPredicate::Compare {
712 field: fv.field,
713 op,
714 value: fv.value,
715 });
716 }
717 }
718 if let Some(fv) = self.in_set {
719 set += 1;
720 chosen = Some(SchemaPredicate::In {
721 field: fv.field,
722 values: fv.values,
723 });
724 }
725 if let Some(fp) = self.field_equals_field {
726 set += 1;
727 chosen = Some(SchemaPredicate::FieldEqualsField {
728 left: fp.left,
729 right: fp.right,
730 });
731 }
732 if let Some(inner) = self.not {
733 set += 1;
734 chosen = Some(SchemaPredicate::Not(Box::new(inner.build(schema_name)?)));
735 }
736 if let Some(list) = self.any {
737 set += 1;
738 chosen = Some(SchemaPredicate::Any(build_group(list, schema_name, "any")?));
739 }
740 if let Some(list) = self.all {
741 set += 1;
742 chosen = Some(SchemaPredicate::All(build_group(list, schema_name, "all")?));
743 }
744 match (set, chosen) {
745 (1, Some(p)) => Ok(p),
746 (0, _) => Err(SchemaError::Parse(format!(
747 "schema '{schema_name}': a predicate has no condition (expected one of \
748 field_present, field_absent, any_of, equals, matches, gt, gte, lt, lte, \
749 in, field_equals_field, not, any, all)"
750 ))),
751 _ => Err(SchemaError::Parse(format!(
752 "schema '{schema_name}': a predicate sets multiple conditions; use one per list item"
753 ))),
754 }
755 }
756}
757
758fn build_group(
760 list: Vec<SchemaPredicateConfig>,
761 schema_name: &str,
762 kind: &str,
763) -> Result<Vec<SchemaPredicate>, SchemaError> {
764 if list.is_empty() {
765 return Err(SchemaError::Parse(format!(
766 "schema '{schema_name}': '{kind}' needs at least one sub-predicate"
767 )));
768 }
769 list.into_iter().map(|p| p.build(schema_name)).collect()
770}
771
772#[derive(Debug, Clone, Deserialize)]
774pub struct SchemaSignatureConfig {
775 pub name: String,
777 #[serde(default = "default_user_specificity")]
780 pub specificity: u32,
781 #[serde(default, rename = "match")]
783 pub predicates: Vec<SchemaPredicateConfig>,
784}
785
786fn default_user_specificity() -> u32 {
787 50
788}
789
790#[derive(Debug, Clone, Default, Deserialize)]
793pub struct SchemaSignaturesFile {
794 #[serde(default)]
795 pub schemas: Vec<SchemaSignatureConfig>,
796 #[serde(default)]
797 pub routing: Option<RoutingConfig>,
798}
799
800impl SchemaSignatureConfig {
801 fn build(self) -> Result<SchemaSignature, SchemaError> {
802 let name = self.name;
803 let predicates = self
804 .predicates
805 .into_iter()
806 .map(|p| p.build(&name))
807 .collect::<Result<Vec<_>, _>>()?;
808 Ok(SchemaSignature {
809 name,
810 predicates,
811 specificity: self.specificity,
812 })
813 }
814}
815
816pub fn parse_schema_signatures(yaml: &str) -> Result<Vec<SchemaSignature>, SchemaError> {
818 let file: SchemaSignaturesFile =
819 yaml_serde::from_str(yaml).map_err(|e| SchemaError::Parse(e.to_string()))?;
820 file.schemas.into_iter().map(|s| s.build()).collect()
821}
822
823pub fn load_schema_signatures(path: &Path) -> Result<Vec<SchemaSignature>, SchemaError> {
825 let content = fs::read_to_string(path).map_err(|e| SchemaError::Io {
826 path: path.display().to_string(),
827 source: e,
828 })?;
829 parse_schema_signatures(&content)
830}
831
832pub fn parse_schema_config(
835 yaml: &str,
836) -> Result<(Vec<SchemaSignature>, Option<RoutingConfig>), SchemaError> {
837 let file: SchemaSignaturesFile =
838 yaml_serde::from_str(yaml).map_err(|e| SchemaError::Parse(e.to_string()))?;
839 let signatures = file
840 .schemas
841 .into_iter()
842 .map(|s| s.build())
843 .collect::<Result<Vec<_>, _>>()?;
844 Ok((signatures, file.routing))
845}
846
847pub fn load_schema_config(
850 path: &Path,
851) -> Result<(Vec<SchemaSignature>, Option<RoutingConfig>), SchemaError> {
852 let content = fs::read_to_string(path).map_err(|e| SchemaError::Io {
853 path: path.display().to_string(),
854 source: e,
855 })?;
856 parse_schema_config(&content)
857}
858
859pub fn validate_schema_config(
873 user_signatures: &[SchemaSignature],
874 routing: Option<&RoutingConfig>,
875) -> Vec<String> {
876 let mut findings = Vec::new();
877
878 let mut all = builtin_signatures();
880 all.extend(user_signatures.iter().cloned());
881 let preds = |s: &SchemaSignature| -> Vec<String> {
882 s.predicates.iter().map(|p| p.describe()).collect()
883 };
884
885 for i in 0..user_signatures.len() {
887 for j in (i + 1)..user_signatures.len() {
888 if user_signatures[i].name == user_signatures[j].name
889 && preds(&user_signatures[i]) == preds(&user_signatures[j])
890 {
891 findings.push(format!(
892 "duplicate signature '{}' with identical predicates",
893 user_signatures[i].name
894 ));
895 }
896 }
897 }
898
899 for b in &all {
901 let b_preds = preds(b);
902 for a in &all {
903 if a.name != b.name
904 && a.specificity > b.specificity
905 && !a.predicates.is_empty()
906 && preds(a).iter().all(|p| b_preds.contains(p))
907 {
908 findings.push(format!(
909 "signature '{}' (specificity {}) is unreachable: shadowed by '{}' (specificity {}) whose predicates are a subset",
910 b.name, b.specificity, a.name, a.specificity
911 ));
912 break;
913 }
914 }
915 }
916
917 if let Some(routing) = routing {
919 let mut known: std::collections::HashSet<&str> =
920 builtin_schema_names().into_iter().collect();
921 for s in user_signatures {
922 known.insert(s.name.as_str());
923 }
924 let mut seen: std::collections::HashSet<&str> = std::collections::HashSet::new();
925 for binding in &routing.bindings {
926 if !known.contains(binding.schema.as_str()) {
927 findings.push(format!(
928 "routing binding references unknown schema '{}' (no built-in or user signature produces it)",
929 binding.schema
930 ));
931 }
932 if !seen.insert(binding.schema.as_str()) {
933 findings.push(format!(
934 "duplicate routing binding for schema '{}'",
935 binding.schema
936 ));
937 }
938 }
939 for (alias, canonical) in &routing.aliases {
940 if !known.contains(canonical.as_str()) {
941 findings.push(format!(
942 "alias '{alias}' targets unknown schema '{canonical}' (no built-in or user signature produces it)"
943 ));
944 }
945 }
946 }
947
948 findings
949}
950
951#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Deserialize)]
957#[serde(rename_all = "snake_case")]
958pub enum OnUnknown {
959 #[default]
961 Warn,
962 Drop,
964 Passthrough,
966 Error,
968}
969
970#[derive(Debug, Clone, Default, Deserialize)]
974#[serde(deny_unknown_fields)]
975pub struct SchemaLogsource {
976 #[serde(default)]
977 pub product: Option<String>,
978 #[serde(default)]
979 pub service: Option<String>,
980 #[serde(default)]
981 pub category: Option<String>,
982 #[serde(default)]
983 pub custom: HashMap<String, String>,
984}
985
986impl SchemaLogsource {
987 fn to_logsource(&self) -> LogSource {
988 LogSource {
989 product: self.product.clone(),
990 service: self.service.clone(),
991 category: self.category.clone(),
992 custom: self.custom.clone(),
993 ..LogSource::default()
994 }
995 }
996}
997
998#[derive(Debug, Clone, Deserialize)]
1001pub struct SchemaBinding {
1002 pub schema: String,
1003 #[serde(default)]
1005 pub pipelines: Vec<String>,
1006 #[serde(default)]
1009 pub logsource: Option<SchemaLogsource>,
1010}
1011
1012fn builtin_schema_logsource() -> HashMap<String, LogSource> {
1021 fn ls(product: &str, service: Option<&str>) -> LogSource {
1022 LogSource {
1023 product: Some(product.to_string()),
1024 service: service.map(str::to_string),
1025 ..LogSource::default()
1026 }
1027 }
1028 HashMap::from([
1029 ("sysmon".to_string(), ls("windows", Some("sysmon"))),
1030 ("windows_eventlog".to_string(), ls("windows", None)),
1031 ("ecs_windows".to_string(), ls("windows", None)),
1032 ("ecs_linux".to_string(), ls("linux", None)),
1033 ])
1034}
1035
1036#[derive(Debug, Clone, Default, Deserialize)]
1038pub struct RoutingConfig {
1039 #[serde(default)]
1040 pub on_unknown: OnUnknown,
1041 #[serde(default)]
1042 pub bindings: Vec<SchemaBinding>,
1043 #[serde(default)]
1046 pub default_pipelines: Vec<String>,
1047 #[serde(default)]
1052 pub aliases: HashMap<String, String>,
1053}
1054
1055#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1057pub enum RouteDecision {
1058 Evaluate { set: usize, unknown: bool },
1061 Drop,
1063 Error,
1065}
1066
1067#[derive(Debug, Clone)]
1074pub struct RoutingPlan {
1075 pipeline_sets: Vec<Vec<String>>,
1077 schema_to_set: HashMap<String, usize>,
1079 schema_logsource: HashMap<String, LogSource>,
1083 aliases: HashMap<String, String>,
1087 on_unknown: OnUnknown,
1088}
1089
1090impl RoutingPlan {
1091 pub fn from_config(config: &RoutingConfig) -> Self {
1094 let mut pipeline_sets: Vec<Vec<String>> = vec![config.default_pipelines.clone()];
1096 let mut schema_to_set: HashMap<String, usize> = HashMap::new();
1097 let mut schema_logsource = builtin_schema_logsource();
1100 let mut aliases = builtin_schema_aliases();
1102 for (alias, canonical) in &config.aliases {
1103 aliases.insert(alias.clone(), canonical.clone());
1104 }
1105
1106 for binding in &config.bindings {
1107 let idx = pipeline_sets
1108 .iter()
1109 .position(|s| s == &binding.pipelines)
1110 .unwrap_or_else(|| {
1111 pipeline_sets.push(binding.pipelines.clone());
1112 pipeline_sets.len() - 1
1113 });
1114 schema_to_set.insert(binding.schema.clone(), idx);
1115 if let Some(ls) = &binding.logsource {
1116 schema_logsource.insert(binding.schema.clone(), ls.to_logsource());
1117 }
1118 }
1119
1120 RoutingPlan {
1121 pipeline_sets,
1122 schema_to_set,
1123 schema_logsource,
1124 aliases,
1125 on_unknown: config.on_unknown,
1126 }
1127 }
1128
1129 pub fn pipeline_sets(&self) -> &[Vec<String>] {
1132 &self.pipeline_sets
1133 }
1134
1135 pub fn on_unknown(&self) -> OnUnknown {
1137 self.on_unknown
1138 }
1139
1140 pub fn schema_logsource(&self, schema: &str) -> Option<&LogSource> {
1144 self.schema_logsource.get(schema)
1145 }
1146
1147 pub fn schemas_with_logsource(&self) -> Vec<String> {
1150 let mut names: Vec<String> = self.schema_logsource.keys().cloned().collect();
1151 names.sort();
1152 names
1153 }
1154
1155 pub fn set_product_partition(&self) -> Vec<Option<std::collections::HashSet<String>>> {
1167 use std::collections::HashSet;
1168 let n = self.pipeline_sets.len();
1169 let mut out: Vec<Option<HashSet<String>>> = (0..n).map(|_| Some(HashSet::new())).collect();
1170 if let Some(first) = out.get_mut(0) {
1171 *first = None;
1172 }
1173
1174 let mut routes: Vec<(usize, &str)> = self
1177 .schema_to_set
1178 .iter()
1179 .map(|(s, &set)| (set, s.as_str()))
1180 .collect();
1181 for (alias, canonical) in &self.aliases {
1182 if !self.schema_to_set.contains_key(alias)
1183 && let Some(&set) = self.schema_to_set.get(canonical)
1184 {
1185 routes.push((set, alias.as_str()));
1186 }
1187 }
1188
1189 for (set, schema) in routes {
1190 if set == 0 {
1191 continue;
1192 }
1193 let product = self
1194 .schema_logsource
1195 .get(schema)
1196 .and_then(|ls| ls.product.as_deref());
1197 let Some(slot) = out.get_mut(set) else {
1198 continue;
1199 };
1200 match product {
1201 Some(p) => {
1202 if let Some(products) = slot {
1203 products.insert(p.to_ascii_lowercase());
1204 }
1205 }
1206 None => *slot = None,
1207 }
1208 }
1209 out
1210 }
1211
1212 pub fn decide(&self, schema: Option<&str>) -> RouteDecision {
1215 match schema {
1216 Some(s) if self.schema_to_set.contains_key(s) => RouteDecision::Evaluate {
1218 set: self.schema_to_set[s],
1219 unknown: false,
1220 },
1221 Some(s)
1225 if self
1226 .aliases
1227 .get(s)
1228 .and_then(|canonical| self.schema_to_set.get(canonical))
1229 .is_some() =>
1230 {
1231 let canonical = &self.aliases[s];
1232 RouteDecision::Evaluate {
1233 set: self.schema_to_set[canonical],
1234 unknown: false,
1235 }
1236 }
1237 Some(_) => RouteDecision::Evaluate {
1239 set: 0,
1240 unknown: false,
1241 },
1242 None => match self.on_unknown {
1244 OnUnknown::Warn | OnUnknown::Passthrough => RouteDecision::Evaluate {
1245 set: 0,
1246 unknown: true,
1247 },
1248 OnUnknown::Drop => RouteDecision::Drop,
1249 OnUnknown::Error => RouteDecision::Error,
1250 },
1251 }
1252 }
1253}
1254
1255#[derive(Debug, Clone, PartialEq, Eq)]
1261pub struct SchemaCountEntry {
1262 pub schema: String,
1264 pub count: u64,
1266}
1267
1268#[derive(Debug, Clone, PartialEq, Eq)]
1270pub struct UnknownShapeEntry {
1271 pub keys: Vec<String>,
1274 pub count: u64,
1276}
1277
1278const UNKNOWN_SHAPE_CAP: usize = 200;
1280const UNKNOWN_SHAPE_MAX_KEYS: usize = 64;
1282
1283#[derive(Debug, Clone, Default)]
1285pub struct SchemaObservation {
1286 pub by_schema: Vec<SchemaCountEntry>,
1288 pub classified: u64,
1290 pub unknown: u64,
1292 pub ambiguous: u64,
1295 pub unknown_shapes: Vec<UnknownShapeEntry>,
1298 pub unrecognized_shapes: Vec<UnknownShapeEntry>,
1302 pub events_observed: u64,
1304 pub lifetime_classified: u64,
1307 pub lifetime_unknown: u64,
1309 pub lifetime_ambiguous: u64,
1311 pub uptime_seconds: f64,
1313}
1314
1315pub struct SchemaObserver {
1321 classifier: SchemaClassifier,
1322 counts: Mutex<HashMap<String, u64>>,
1323 unknown: AtomicU64,
1324 ambiguous: AtomicU64,
1325 unknown_shapes: Mutex<HashMap<Vec<String>, u64>>,
1328 discovery_sampling: bool,
1335 unrecognized_shapes: Mutex<HashMap<Vec<String>, u64>>,
1338 lifetime_classified: AtomicU64,
1339 lifetime_unknown: AtomicU64,
1340 lifetime_ambiguous: AtomicU64,
1341 start: Mutex<Instant>,
1342}
1343
1344impl SchemaObserver {
1345 pub fn new(classifier: SchemaClassifier) -> Self {
1348 Self::new_with_discovery(classifier, false)
1349 }
1350
1351 pub fn new_with_discovery(classifier: SchemaClassifier, discovery_sampling: bool) -> Self {
1355 Self {
1356 classifier,
1357 counts: Mutex::new(HashMap::new()),
1358 unknown: AtomicU64::new(0),
1359 ambiguous: AtomicU64::new(0),
1360 unknown_shapes: Mutex::new(HashMap::new()),
1361 discovery_sampling,
1362 unrecognized_shapes: Mutex::new(HashMap::new()),
1363 lifetime_classified: AtomicU64::new(0),
1364 lifetime_unknown: AtomicU64::new(0),
1365 lifetime_ambiguous: AtomicU64::new(0),
1366 start: Mutex::new(Instant::now()),
1367 }
1368 }
1369
1370 pub fn discovery_sampling(&self) -> bool {
1373 self.discovery_sampling
1374 }
1375
1376 pub fn builtin() -> Self {
1378 Self::new(SchemaClassifier::builtin())
1379 }
1380
1381 pub fn observe<E: Event + ?Sized>(&self, event: &E) {
1384 let (matched, ambiguous) = self.classifier.classify_with_ambiguity(event);
1385 if ambiguous {
1386 self.ambiguous.fetch_add(1, Ordering::Relaxed);
1387 self.lifetime_ambiguous.fetch_add(1, Ordering::Relaxed);
1388 }
1389 let discovery_unrecognized = match &matched {
1393 None => true,
1394 Some(m) => m.name == "generic_json",
1395 };
1396 if self.discovery_sampling && discovery_unrecognized {
1397 self.record_unrecognized_shape(event);
1398 }
1399
1400 match matched {
1401 Some(m) => {
1402 self.lifetime_classified.fetch_add(1, Ordering::Relaxed);
1403 let mut counts = self.counts.lock().expect("schema observer mutex poisoned");
1404 *counts.entry(m.name).or_insert(0) += 1;
1405 }
1406 None => {
1407 self.unknown.fetch_add(1, Ordering::Relaxed);
1408 self.lifetime_unknown.fetch_add(1, Ordering::Relaxed);
1409 self.record_unknown_shape(event);
1410 }
1411 }
1412 }
1413
1414 fn record_unknown_shape<E: Event + ?Sized>(&self, event: &E) {
1417 let mut keys: Vec<String> = event.field_keys().iter().map(|k| k.to_string()).collect();
1418 keys.sort();
1419 keys.dedup();
1420 keys.truncate(UNKNOWN_SHAPE_MAX_KEYS);
1421 let mut shapes = self
1422 .unknown_shapes
1423 .lock()
1424 .expect("schema observer shapes mutex poisoned");
1425 if shapes.contains_key(&keys) || shapes.len() < UNKNOWN_SHAPE_CAP {
1427 *shapes.entry(keys).or_insert(0) += 1;
1428 }
1429 }
1430
1431 fn record_unrecognized_shape<E: Event + ?Sized>(&self, event: &E) {
1435 let mut keys: Vec<String> = event.field_keys().iter().map(|k| k.to_string()).collect();
1436 keys.sort();
1437 keys.dedup();
1438 keys.truncate(UNKNOWN_SHAPE_MAX_KEYS);
1439 if keys.is_empty() {
1440 return;
1441 }
1442 let mut shapes = self
1443 .unrecognized_shapes
1444 .lock()
1445 .expect("schema observer shapes mutex poisoned");
1446 if shapes.contains_key(&keys) || shapes.len() < UNKNOWN_SHAPE_CAP {
1447 *shapes.entry(keys).or_insert(0) += 1;
1448 }
1449 }
1450
1451 pub fn snapshot(&self) -> SchemaObservation {
1453 let counts = self.counts.lock().expect("schema observer mutex poisoned");
1454 let mut by_schema: Vec<SchemaCountEntry> = counts
1455 .iter()
1456 .map(|(schema, count)| SchemaCountEntry {
1457 schema: schema.clone(),
1458 count: *count,
1459 })
1460 .collect();
1461 let classified: u64 = counts.values().sum();
1462 drop(counts);
1463 by_schema.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.schema.cmp(&b.schema)));
1464
1465 let shapes = self
1466 .unknown_shapes
1467 .lock()
1468 .expect("schema observer shapes mutex poisoned");
1469 let mut unknown_shapes: Vec<UnknownShapeEntry> = shapes
1470 .iter()
1471 .map(|(keys, count)| UnknownShapeEntry {
1472 keys: keys.clone(),
1473 count: *count,
1474 })
1475 .collect();
1476 drop(shapes);
1477 unknown_shapes.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.keys.cmp(&b.keys)));
1478
1479 let unrec = self
1480 .unrecognized_shapes
1481 .lock()
1482 .expect("schema observer shapes mutex poisoned");
1483 let mut unrecognized_shapes: Vec<UnknownShapeEntry> = unrec
1484 .iter()
1485 .map(|(keys, count)| UnknownShapeEntry {
1486 keys: keys.clone(),
1487 count: *count,
1488 })
1489 .collect();
1490 drop(unrec);
1491 unrecognized_shapes.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.keys.cmp(&b.keys)));
1492
1493 let unknown = self.unknown.load(Ordering::Relaxed);
1494 SchemaObservation {
1495 by_schema,
1496 classified,
1497 unknown,
1498 ambiguous: self.ambiguous.load(Ordering::Relaxed),
1499 unknown_shapes,
1500 unrecognized_shapes,
1501 events_observed: classified + unknown,
1506 lifetime_classified: self.lifetime_classified.load(Ordering::Relaxed),
1507 lifetime_unknown: self.lifetime_unknown.load(Ordering::Relaxed),
1508 lifetime_ambiguous: self.lifetime_ambiguous.load(Ordering::Relaxed),
1509 uptime_seconds: self
1510 .start
1511 .lock()
1512 .expect("schema observer start mutex poisoned")
1513 .elapsed()
1514 .as_secs_f64(),
1515 }
1516 }
1517
1518 pub fn reset(&self) -> (u64, u64) {
1521 let mut counts = self.counts.lock().expect("schema observer mutex poisoned");
1522 let previous_classified: u64 = counts.values().sum();
1523 counts.clear();
1524 drop(counts);
1525 self.unknown_shapes
1526 .lock()
1527 .expect("schema observer shapes mutex poisoned")
1528 .clear();
1529 self.unrecognized_shapes
1530 .lock()
1531 .expect("schema observer shapes mutex poisoned")
1532 .clear();
1533 let previous_unknown = self.unknown.swap(0, Ordering::Relaxed);
1534 self.ambiguous.store(0, Ordering::Relaxed);
1535 *self
1536 .start
1537 .lock()
1538 .expect("schema observer start mutex poisoned") = Instant::now();
1539 (previous_classified, previous_unknown)
1540 }
1541
1542 pub fn lifetime_classified(&self) -> u64 {
1544 self.lifetime_classified.load(Ordering::Relaxed)
1545 }
1546
1547 pub fn lifetime_unknown(&self) -> u64 {
1549 self.lifetime_unknown.load(Ordering::Relaxed)
1550 }
1551
1552 pub fn lifetime_ambiguous(&self) -> u64 {
1554 self.lifetime_ambiguous.load(Ordering::Relaxed)
1555 }
1556}
1557
1558#[cfg(test)]
1559mod tests {
1560 use super::*;
1561 use crate::event::JsonEvent;
1562 use serde_json::json;
1563
1564 fn classify(value: &serde_json::Value) -> Option<String> {
1565 SchemaClassifier::builtin()
1566 .classify(&JsonEvent::borrow(value))
1567 .map(|m| m.name)
1568 }
1569
1570 #[test]
1571 fn recognizes_ecs_by_version_marker() {
1572 let v = json!({"ecs": {"version": "8.11.0"}, "process": {"command_line": "whoami"}});
1573 assert_eq!(classify(&v).as_deref(), Some("ecs"));
1574 }
1575
1576 #[test]
1577 fn recognizes_ecs_with_flattened_keys() {
1578 let v = json!({"ecs.version": "8.11.0", "process.command_line": "whoami"});
1579 assert_eq!(classify(&v).as_deref(), Some("ecs"));
1580 }
1581
1582 #[test]
1583 fn recognizes_ocsf_by_class_and_metadata() {
1584 let v = json!({"class_uid": 1001, "category_uid": 1, "metadata": {"version": "1.1.0"}});
1585 assert_eq!(classify(&v).as_deref(), Some("ocsf"));
1586 }
1587
1588 #[test]
1589 fn recognizes_rendered_windows_event_log() {
1590 let v = json!({"Event": {"System": {"EventID": 4688, "Provider": "Microsoft-Windows-Security-Auditing"}}});
1591 assert_eq!(classify(&v).as_deref(), Some("windows_eventlog"));
1592 }
1593
1594 #[test]
1595 fn recognizes_sysmon_by_channel() {
1596 let v = json!({"Channel": "Microsoft-Windows-Sysmon/Operational", "EventID": 1, "Image": "C:/cmd.exe"});
1597 assert_eq!(classify(&v).as_deref(), Some("sysmon"));
1598 }
1599
1600 #[test]
1601 fn recognizes_sysmon_by_provider() {
1602 let v = json!({"Provider_Name": "Microsoft-Windows-Sysmon", "EventID": 3});
1603 assert_eq!(classify(&v).as_deref(), Some("sysmon"));
1604 }
1605
1606 #[test]
1607 fn recognizes_flat_sysmon_by_field_shape() {
1608 let v = json!({"EventID": 1, "ProcessGuid": "{abc}", "CommandLine": "cmd /c whoami"});
1609 assert_eq!(classify(&v).as_deref(), Some("sysmon"));
1610 }
1611
1612 #[test]
1613 fn recognizes_cef_structured_fields() {
1614 let v = json!({"deviceVendor": "Security", "deviceProduct": "IDS", "signatureId": "100", "src": "10.0.0.1"});
1615 assert_eq!(classify(&v).as_deref(), Some("cef"));
1616 }
1617
1618 #[test]
1619 fn falls_back_to_generic_json_for_unrecognized_structured_events() {
1620 let v = json!({"some_vendor_field": "x", "another": 1});
1621 assert_eq!(classify(&v).as_deref(), Some("generic_json"));
1622 }
1623
1624 #[test]
1625 fn fieldless_events_are_unknown() {
1626 assert_eq!(classify(&json!({})), None);
1628 assert_eq!(classify(&json!("just a string")), None);
1630 }
1631
1632 #[test]
1633 fn specificity_prefers_specific_schema_over_generic() {
1634 let v = json!({"ecs.version": "8.0.0", "vendor_blob": {"x": 1}});
1636 let cls = SchemaClassifier::builtin();
1637 let m = cls.classify(&JsonEvent::borrow(&v)).unwrap();
1638 assert_eq!(m.name, "ecs");
1639 assert_eq!(m.specificity, 100);
1640 let all = cls.classify_all(&JsonEvent::borrow(&v));
1642 assert_eq!(all.first().map(String::as_str), Some("ecs"));
1643 assert!(all.iter().any(|n| n == "generic_json"));
1644 }
1645
1646 #[test]
1647 fn schema_names_lists_builtins_most_specific_first() {
1648 let classifier = SchemaClassifier::builtin();
1649 let names = classifier.schema_names();
1650 assert_eq!(names.first(), Some(&"ecs_linux"));
1653 assert!(names.contains(&"ecs_windows"));
1654 assert!(names.contains(&"ecs"));
1655 assert!(names.contains(&"generic_json"));
1656 assert_eq!(names.last(), Some(&"generic_json"));
1658 }
1659
1660 #[test]
1661 fn ecs_windows_specialization_classifies_and_aliases_to_ecs() {
1662 let v = json!({"ecs.version": "8.11.0", "winlog": {"channel": "Security"}});
1665 assert_eq!(classify(&v).as_deref(), Some("ecs_windows"));
1666 let plain = json!({"ecs.version": "8.11.0", "process": {"command_line": "whoami"}});
1668 assert_eq!(classify(&plain).as_deref(), Some("ecs"));
1669
1670 let config = RoutingConfig {
1673 on_unknown: OnUnknown::Warn,
1674 default_pipelines: vec![],
1675 aliases: HashMap::new(),
1676 bindings: vec![SchemaBinding {
1677 schema: "ecs".to_string(),
1678 pipelines: vec!["ecs_windows".to_string()],
1679 logsource: None,
1680 }],
1681 };
1682 let plan = RoutingPlan::from_config(&config);
1683 let ecs_set = match plan.decide(Some("ecs")) {
1684 RouteDecision::Evaluate { set, .. } => set,
1685 other => panic!("unexpected: {other:?}"),
1686 };
1687 assert_eq!(plan.decide(Some("ecs_windows")), plan.decide(Some("ecs")));
1689 assert_ne!(ecs_set, 0, "ecs binding is a non-default set");
1690 assert_eq!(
1691 plan.schema_logsource("ecs_windows")
1692 .and_then(|l| l.product.as_deref()),
1693 Some("windows")
1694 );
1695 }
1696
1697 #[test]
1698 fn user_alias_routes_as_canonical() {
1699 let yaml = r#"
1700schemas:
1701 - name: my_win
1702 specificity: 70
1703 match:
1704 - field_present: vendor.win_marker
1705routing:
1706 aliases:
1707 my_win: ecs
1708 bindings:
1709 - schema: ecs
1710 pipelines: [ecs_windows]
1711"#;
1712 let (_sigs, routing) = parse_schema_config(yaml).unwrap();
1713 let plan = RoutingPlan::from_config(&routing.expect("routing"));
1714 assert_eq!(plan.decide(Some("my_win")), plan.decide(Some("ecs")));
1716 assert!(matches!(
1717 plan.decide(Some("my_win")),
1718 RouteDecision::Evaluate { unknown: false, .. }
1719 ));
1720 }
1721
1722 #[test]
1723 fn set_product_partition_only_for_platform_locked_sets() {
1724 let config = RoutingConfig {
1725 on_unknown: OnUnknown::Warn,
1726 default_pipelines: vec![],
1727 aliases: HashMap::new(),
1728 bindings: vec![
1729 SchemaBinding {
1730 schema: "sysmon".to_string(),
1731 pipelines: vec!["p_sysmon".to_string()],
1732 logsource: None,
1733 },
1734 SchemaBinding {
1735 schema: "ecs".to_string(),
1736 pipelines: vec!["p_ecs".to_string()],
1737 logsource: None,
1738 },
1739 ],
1740 };
1741 let plan = RoutingPlan::from_config(&config);
1742 let part = plan.set_product_partition();
1743 assert!(part[0].is_none(), "default set is never partitioned");
1744
1745 let set_of = |schema| match plan.decide(Some(schema)) {
1746 RouteDecision::Evaluate { set, .. } => set,
1747 other => panic!("unexpected: {other:?}"),
1748 };
1749 let sysmon_set = set_of("sysmon");
1751 assert_eq!(
1752 part[sysmon_set].as_ref().map(|s| s.contains("windows")),
1753 Some(true)
1754 );
1755 assert!(part[set_of("ecs")].is_none());
1757 }
1758
1759 #[test]
1760 fn parses_user_signatures_from_yaml() {
1761 let yaml = r#"
1762schemas:
1763 - name: my_vendor
1764 specificity: 70
1765 match:
1766 - field_present: vendor.product
1767 - equals:
1768 field: event_type
1769 value: alert
1770 - any_of: [a, b]
1771"#;
1772 let sigs = parse_schema_signatures(yaml).expect("parse");
1773 assert_eq!(sigs.len(), 1);
1774 assert_eq!(sigs[0].name, "my_vendor");
1775 assert_eq!(sigs[0].specificity, 70);
1776 assert_eq!(sigs[0].predicates.len(), 3);
1777
1778 let cls = SchemaClassifier::with_user_signatures(sigs);
1779 let v = json!({"vendor": {"product": "X"}, "event_type": "ALERT", "a": 1});
1780 assert_eq!(
1781 cls.classify(&JsonEvent::borrow(&v))
1782 .map(|m| m.name)
1783 .as_deref(),
1784 Some("my_vendor")
1785 );
1786 }
1787
1788 #[test]
1789 fn user_signature_with_invalid_regex_is_rejected() {
1790 let yaml = r#"
1791schemas:
1792 - name: bad
1793 match:
1794 - matches:
1795 field: msg
1796 value: "([unclosed"
1797"#;
1798 let err = parse_schema_signatures(yaml).unwrap_err();
1799 assert!(matches!(err, SchemaError::InvalidRegex { .. }));
1800 }
1801
1802 #[test]
1803 fn user_regex_signature_matches_field_value() {
1804 let yaml = r#"
1805schemas:
1806 - name: cef_raw
1807 specificity: 60
1808 match:
1809 - matches:
1810 field: message
1811 value: "^CEF:\\d"
1812"#;
1813 let sigs = parse_schema_signatures(yaml).expect("parse");
1814 let cls = SchemaClassifier::with_user_signatures(sigs);
1815 let v = json!({"message": "CEF:0|Vendor|Product|1.0|100|Name|9|src=1.2.3.4"});
1816 assert_eq!(
1817 cls.classify(&JsonEvent::borrow(&v))
1818 .map(|m| m.name)
1819 .as_deref(),
1820 Some("cef_raw")
1821 );
1822 }
1823
1824 fn classifier_from_match(match_body: &str) -> SchemaClassifier {
1826 let yaml = format!("schemas:\n - name: t\n specificity: 70\n match:\n{match_body}");
1827 let sigs = parse_schema_signatures(&yaml).expect("parse");
1828 SchemaClassifier::new(sigs)
1829 }
1830
1831 fn matches_t(match_body: &str, event: &serde_json::Value) -> bool {
1832 classifier_from_match(match_body)
1833 .classify(&JsonEvent::borrow(event))
1834 .is_some()
1835 }
1836
1837 #[test]
1838 fn numeric_comparisons() {
1839 let body = " - gte: { field: EventID, value: 4000 }\n";
1840 assert!(matches_t(body, &json!({"EventID": 4688})));
1841 assert!(matches_t(body, &json!({"EventID": 4000})));
1842 assert!(!matches_t(body, &json!({"EventID": 1})));
1843 assert!(matches_t(body, &json!({"EventID": "4688"})));
1845 assert!(!matches_t(body, &json!({"EventID": "not-a-number"})));
1847 assert!(matches_t(
1849 " - lt: { field: score, value: 10 }\n",
1850 &json!({"score": 9.5})
1851 ));
1852 assert!(matches_t(
1853 " - gt: { field: score, value: 10 }\n",
1854 &json!({"score": 10.1})
1855 ));
1856 }
1857
1858 #[test]
1859 fn in_set_membership_is_case_insensitive() {
1860 let body = " - in: { field: event_type, values: [alert, alarm] }\n";
1861 assert!(matches_t(body, &json!({"event_type": "ALERT"})));
1862 assert!(matches_t(body, &json!({"event_type": "alarm"})));
1863 assert!(!matches_t(body, &json!({"event_type": "info"})));
1864 assert!(!matches_t(body, &json!({})));
1865 }
1866
1867 #[test]
1868 fn field_equals_field_compares_two_fields() {
1869 let body = " - field_equals_field: { left: a, right: b }\n";
1870 assert!(matches_t(body, &json!({"a": "X", "b": "x"})));
1871 assert!(!matches_t(body, &json!({"a": "X", "b": "y"})));
1872 assert!(!matches_t(body, &json!({"a": "X"})));
1874 }
1875
1876 #[test]
1877 fn recursive_not_any_all_groups() {
1878 let any_body = " - any:\n - field_present: winlog.channel\n - equals: { field: host.os.type, value: windows }\n";
1880 assert!(matches_t(
1881 any_body,
1882 &json!({"winlog": {"channel": "Security"}})
1883 ));
1884 assert!(matches_t(
1885 any_body,
1886 &json!({"host": {"os": {"type": "windows"}}})
1887 ));
1888 assert!(!matches_t(any_body, &json!({"unrelated": 1})));
1889
1890 let not_body = " - not: { field_present: ecs.version }\n";
1892 assert!(matches_t(not_body, &json!({"CommandLine": "whoami"})));
1893 assert!(!matches_t(not_body, &json!({"ecs.version": "8.0.0"})));
1894
1895 let all_body = " - all:\n - field_present: a\n - field_present: b\n";
1897 assert!(matches_t(all_body, &json!({"a": 1, "b": 2})));
1898 assert!(!matches_t(all_body, &json!({"a": 1})));
1899 }
1900
1901 #[test]
1902 fn empty_group_is_rejected() {
1903 let yaml = "schemas:\n - name: t\n match:\n - any: []\n";
1904 let err = parse_schema_signatures(yaml).unwrap_err();
1905 assert!(
1906 matches!(&err, SchemaError::Parse(m) if m.contains("'any' needs at least one")),
1907 "got: {err}"
1908 );
1909 }
1910
1911 #[test]
1912 fn predicate_with_two_conditions_is_rejected() {
1913 let yaml = "schemas:\n - name: t\n match:\n - field_present: a\n field_absent: b\n";
1914 let err = parse_schema_signatures(yaml).unwrap_err();
1915 assert!(
1916 matches!(&err, SchemaError::Parse(m) if m.contains("multiple conditions")),
1917 "got: {err}"
1918 );
1919 }
1920
1921 #[test]
1922 fn explain_reports_matched_signature() {
1923 let cls = SchemaClassifier::builtin();
1924 let v = json!({"ecs.version": "8.0.0"});
1925 let ex = cls.explain(&JsonEvent::borrow(&v));
1926 assert_eq!(ex.matched.as_deref(), Some("ecs"));
1927 let sig = ex.signature.expect("signature");
1928 assert!(sig.predicates_matched);
1929 assert!(sig.predicates.iter().all(|p| p.matched));
1930 }
1931
1932 #[test]
1933 fn explain_reports_near_miss_for_unknown() {
1934 let sigs = builtin_signatures()
1936 .into_iter()
1937 .filter(|s| s.name != "generic_json")
1938 .collect();
1939 let cls = SchemaClassifier::new(sigs);
1940 let v = json!({"EventID": 1, "Image": "C:/cmd.exe"});
1942 let ex = cls.explain(&JsonEvent::borrow(&v));
1943 assert_eq!(ex.matched, None);
1944 let sig = ex.signature.expect("near-miss");
1945 assert_eq!(sig.name, "sysmon");
1946 assert!(!sig.predicates_matched);
1947 assert!(sig.predicates.iter().any(|p| !p.matched));
1948 }
1949
1950 #[test]
1951 fn validate_flags_unknown_binding_and_shadow() {
1952 let yaml = r#"
1953schemas:
1954 - name: shadowed
1955 specificity: 40
1956 match:
1957 - field_present: ecs.version
1958 - field_present: extra.marker
1959routing:
1960 bindings:
1961 - schema: ecs
1962 pipelines: [ecs_windows]
1963 - schema: nonexistent
1964 pipelines: [x]
1965"#;
1966 let (sigs, routing) = parse_schema_config(yaml).unwrap();
1967 let findings = validate_schema_config(&sigs, routing.as_ref());
1968 assert!(
1969 findings
1970 .iter()
1971 .any(|f| f.contains("unknown schema 'nonexistent'")),
1972 "findings: {findings:?}"
1973 );
1974 assert!(
1977 findings
1978 .iter()
1979 .any(|f| f.contains("'shadowed'") && f.contains("unreachable")),
1980 "findings: {findings:?}"
1981 );
1982 }
1983
1984 #[test]
1985 fn observer_counts_per_schema_and_unknown() {
1986 let observer = SchemaObserver::builtin();
1987 observer.observe(&JsonEvent::borrow(&json!({"ecs.version": "8.0.0"})));
1988 observer.observe(&JsonEvent::borrow(&json!({"ecs.version": "8.1.0"})));
1989 observer.observe(&JsonEvent::borrow(
1990 &json!({"class_uid": 1001, "metadata": {"version": "1.1.0"}}),
1991 ));
1992 observer.observe(&JsonEvent::borrow(&json!({})));
1993
1994 let snap = observer.snapshot();
1995 assert_eq!(snap.events_observed, 4);
1996 assert_eq!(snap.classified, 3);
1997 assert_eq!(snap.unknown, 1);
1998 assert_eq!(snap.by_schema[0].schema, "ecs");
2000 assert_eq!(snap.by_schema[0].count, 2);
2001 let ocsf = snap.by_schema.iter().find(|e| e.schema == "ocsf").unwrap();
2002 assert_eq!(ocsf.count, 1);
2003 }
2004
2005 #[test]
2006 fn routing_plan_dedups_pipeline_sets() {
2007 let config = RoutingConfig {
2008 on_unknown: OnUnknown::Warn,
2009 default_pipelines: vec![],
2010 aliases: HashMap::new(),
2011 bindings: vec![
2012 SchemaBinding {
2013 schema: "ecs".to_string(),
2014 pipelines: vec!["ecs_windows".to_string()],
2015 logsource: None,
2016 },
2017 SchemaBinding {
2018 schema: "winlogbeat".to_string(),
2019 pipelines: vec!["ecs_windows".to_string()],
2020 logsource: None,
2021 },
2022 SchemaBinding {
2023 schema: "sysmon".to_string(),
2024 pipelines: vec!["sysmon".to_string()],
2025 logsource: None,
2026 },
2027 ],
2028 };
2029 let plan = RoutingPlan::from_config(&config);
2030 assert_eq!(plan.pipeline_sets().len(), 3);
2032 let ecs = plan.decide(Some("ecs"));
2034 let win = plan.decide(Some("winlogbeat"));
2035 assert_eq!(ecs, win);
2036 assert!(matches!(
2037 ecs,
2038 RouteDecision::Evaluate { unknown: false, .. }
2039 ));
2040 assert_ne!(plan.decide(Some("sysmon")), ecs);
2042 }
2043
2044 #[test]
2045 fn routing_decides_bound_unbound_and_unknown() {
2046 let config = RoutingConfig {
2047 on_unknown: OnUnknown::Warn,
2048 default_pipelines: vec![],
2049 aliases: HashMap::new(),
2050 bindings: vec![SchemaBinding {
2051 schema: "ecs".to_string(),
2052 pipelines: vec!["ecs_windows".to_string()],
2053 logsource: None,
2054 }],
2055 };
2056 let plan = RoutingPlan::from_config(&config);
2057 assert!(matches!(
2059 plan.decide(Some("ecs")),
2060 RouteDecision::Evaluate { unknown: false, .. }
2061 ));
2062 assert_eq!(
2064 plan.decide(Some("cef")),
2065 RouteDecision::Evaluate {
2066 set: 0,
2067 unknown: false
2068 }
2069 );
2070 assert_eq!(
2072 plan.decide(None),
2073 RouteDecision::Evaluate {
2074 set: 0,
2075 unknown: true
2076 }
2077 );
2078 }
2079
2080 #[test]
2081 fn routing_on_unknown_policies() {
2082 let base = |policy| RoutingConfig {
2083 on_unknown: policy,
2084 default_pipelines: vec![],
2085 aliases: HashMap::new(),
2086 bindings: vec![],
2087 };
2088 assert_eq!(
2089 RoutingPlan::from_config(&base(OnUnknown::Drop)).decide(None),
2090 RouteDecision::Drop
2091 );
2092 assert_eq!(
2093 RoutingPlan::from_config(&base(OnUnknown::Error)).decide(None),
2094 RouteDecision::Error
2095 );
2096 assert_eq!(
2097 RoutingPlan::from_config(&base(OnUnknown::Passthrough)).decide(None),
2098 RouteDecision::Evaluate {
2099 set: 0,
2100 unknown: true
2101 }
2102 );
2103 }
2104
2105 #[test]
2106 fn parses_routing_section_from_yaml() {
2107 let yaml = r#"
2108schemas:
2109 - name: my_vendor
2110 match:
2111 - field_present: vendor.id
2112routing:
2113 on_unknown: drop
2114 default_pipelines: [base]
2115 bindings:
2116 - schema: ecs
2117 pipelines: [ecs_windows]
2118 - schema: my_vendor
2119 pipelines: [vendor_map, base]
2120"#;
2121 let (sigs, routing) = parse_schema_config(yaml).expect("parse");
2122 assert_eq!(sigs.len(), 1);
2123 let routing = routing.expect("routing present");
2124 assert_eq!(routing.on_unknown, OnUnknown::Drop);
2125 assert_eq!(routing.default_pipelines, vec!["base".to_string()]);
2126 assert_eq!(routing.bindings.len(), 2);
2127 let plan = RoutingPlan::from_config(&routing);
2128 assert_eq!(plan.pipeline_sets().len(), 3);
2130 assert_eq!(plan.decide(None), RouteDecision::Drop);
2131 }
2132
2133 #[test]
2134 fn schema_logsource_builtin_defaults_and_overrides() {
2135 let plan = RoutingPlan::from_config(&RoutingConfig::default());
2137 let sysmon = plan.schema_logsource("sysmon").expect("sysmon default");
2138 assert_eq!(sysmon.product.as_deref(), Some("windows"));
2139 assert_eq!(sysmon.service.as_deref(), Some("sysmon"));
2140 assert_eq!(
2141 plan.schema_logsource("windows_eventlog")
2142 .and_then(|l| l.product.as_deref()),
2143 Some("windows")
2144 );
2145 assert!(plan.schema_logsource("ecs").is_none());
2147 assert!(plan.schema_logsource("cef").is_none());
2148
2149 let yaml = r#"
2151schemas:
2152 - name: ecs_windows
2153 match:
2154 - field_present: ecs.version
2155 - field_present: winlog.channel
2156routing:
2157 bindings:
2158 - schema: ecs_windows
2159 pipelines: [ecs_windows]
2160 logsource:
2161 product: windows
2162 - schema: sysmon
2163 pipelines: [sysmon]
2164 logsource:
2165 product: windows
2166 service: sysmon
2167 custom:
2168 tenant: acme
2169"#;
2170 let (_sigs, routing) = parse_schema_config(yaml).expect("parse");
2171 let plan = RoutingPlan::from_config(&routing.expect("routing"));
2172 assert_eq!(
2173 plan.schema_logsource("ecs_windows")
2174 .and_then(|l| l.product.as_deref()),
2175 Some("windows")
2176 );
2177 let sysmon = plan.schema_logsource("sysmon").expect("sysmon override");
2178 assert_eq!(
2179 sysmon.custom.get("tenant").map(String::as_str),
2180 Some("acme")
2181 );
2182 }
2183
2184 #[test]
2185 fn observer_reset_preserves_lifetime_counters() {
2186 let observer = SchemaObserver::builtin();
2187 observer.observe(&JsonEvent::borrow(&json!({"ecs.version": "8.0.0"})));
2188 observer.observe(&JsonEvent::borrow(&json!({})));
2189 let (classified, unknown) = observer.reset();
2190 assert_eq!(classified, 1);
2191 assert_eq!(unknown, 1);
2192
2193 let snap = observer.snapshot();
2194 assert_eq!(snap.classified, 0);
2195 assert_eq!(snap.unknown, 0);
2196 assert_eq!(snap.events_observed, 0);
2197 assert_eq!(snap.lifetime_classified, 1);
2199 assert_eq!(snap.lifetime_unknown, 1);
2200 }
2201
2202 #[test]
2203 fn classify_with_ambiguity_flags_equal_specificity_ties() {
2204 let sigs = vec![
2206 SchemaSignature {
2207 name: "alpha".to_string(),
2208 specificity: 70,
2209 predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
2210 },
2211 SchemaSignature {
2212 name: "beta".to_string(),
2213 specificity: 70,
2214 predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
2215 },
2216 ];
2217 let cls = SchemaClassifier::new(sigs);
2218 let (m, ambiguous) = cls.classify_with_ambiguity(&JsonEvent::borrow(&json!({"a": 1})));
2219 assert!(m.is_some());
2220 assert!(
2221 ambiguous,
2222 "equal-specificity different-name match is ambiguous"
2223 );
2224 let cls = SchemaClassifier::builtin();
2226 let (_, ambiguous) =
2227 cls.classify_with_ambiguity(&JsonEvent::borrow(&json!({"ecs.version": "8.0.0"})));
2228 assert!(!ambiguous);
2229 }
2230
2231 #[test]
2232 fn observer_records_ambiguity_and_unknown_shapes() {
2233 let sigs = vec![
2234 SchemaSignature {
2235 name: "alpha".to_string(),
2236 specificity: 70,
2237 predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
2238 },
2239 SchemaSignature {
2240 name: "beta".to_string(),
2241 specificity: 70,
2242 predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
2243 },
2244 ];
2245 let observer = SchemaObserver::new(SchemaClassifier::new(sigs));
2246 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();
2251 assert_eq!(snap.ambiguous, 1);
2252 assert_eq!(snap.unknown, 2);
2253 assert_eq!(snap.unknown_shapes.len(), 1);
2255 assert_eq!(snap.unknown_shapes[0].count, 2);
2256 assert_eq!(snap.unknown_shapes[0].keys, vec!["shape", "weird"]);
2257 }
2258}