1use std::collections::HashMap;
32use std::fs;
33use std::path::Path;
34use std::sync::Mutex;
35use std::sync::atomic::{AtomicU64, Ordering};
36use std::time::Instant;
37
38use regex::Regex;
39use rsigma_parser::LogSource;
40use serde::{Deserialize, Serialize};
41
42use crate::event::Event;
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub enum CompareOp {
47 Gt,
49 Gte,
51 Lt,
53 Lte,
55}
56
57impl CompareOp {
58 fn apply(self, lhs: f64, rhs: f64) -> bool {
59 match self {
60 CompareOp::Gt => lhs > rhs,
61 CompareOp::Gte => lhs >= rhs,
62 CompareOp::Lt => lhs < rhs,
63 CompareOp::Lte => lhs <= rhs,
64 }
65 }
66
67 fn symbol(self) -> &'static str {
68 match self {
69 CompareOp::Gt => ">",
70 CompareOp::Gte => ">=",
71 CompareOp::Lt => "<",
72 CompareOp::Lte => "<=",
73 }
74 }
75}
76
77#[derive(Debug, Clone)]
83pub enum SchemaPredicate {
84 FieldPresent(String),
86 FieldAbsent(String),
88 AnyOf(Vec<String>),
90 Equals { field: String, value: String },
93 Matches { field: String, regex: Regex },
95 Compare {
98 field: String,
99 op: CompareOp,
100 value: f64,
101 },
102 In { field: String, values: Vec<String> },
105 FieldEqualsField { left: String, right: String },
107 Not(Box<SchemaPredicate>),
109 Any(Vec<SchemaPredicate>),
111 All(Vec<SchemaPredicate>),
114 HasAnyField,
118}
119
120impl SchemaPredicate {
121 fn eval<E: Event + ?Sized>(&self, event: &E) -> bool {
122 match self {
123 SchemaPredicate::FieldPresent(f) => event.get_field(f).is_some(),
124 SchemaPredicate::FieldAbsent(f) => event.get_field(f).is_none(),
125 SchemaPredicate::AnyOf(fields) => fields.iter().any(|f| event.get_field(f).is_some()),
126 SchemaPredicate::Equals { field, value } => event
127 .get_field(field)
128 .and_then(|v| v.as_str().map(|s| s.as_ref().eq_ignore_ascii_case(value)))
129 .unwrap_or(false),
130 SchemaPredicate::Matches { field, regex } => event
131 .get_field(field)
132 .and_then(|v| v.as_str().map(|s| regex.is_match(s.as_ref())))
133 .unwrap_or(false),
134 SchemaPredicate::Compare { field, op, value } => event
135 .get_field(field)
136 .and_then(|v| v.as_f64())
137 .map(|n| op.apply(n, *value))
138 .unwrap_or(false),
139 SchemaPredicate::In { field, values } => event
140 .get_field(field)
141 .and_then(|v| {
142 v.as_str().map(|s| {
143 values
144 .iter()
145 .any(|val| s.as_ref().eq_ignore_ascii_case(val))
146 })
147 })
148 .unwrap_or(false),
149 SchemaPredicate::FieldEqualsField { left, right } => {
150 let l = event
151 .get_field(left)
152 .and_then(|v| v.as_str().map(|s| s.into_owned()));
153 let r = event
154 .get_field(right)
155 .and_then(|v| v.as_str().map(|s| s.into_owned()));
156 matches!((l, r), (Some(a), Some(b)) if a.eq_ignore_ascii_case(&b))
157 }
158 SchemaPredicate::Not(inner) => !inner.eval(event),
159 SchemaPredicate::Any(preds) => preds.iter().any(|p| p.eval(event)),
160 SchemaPredicate::All(preds) => preds.iter().all(|p| p.eval(event)),
161 SchemaPredicate::HasAnyField => !event.field_keys().is_empty(),
162 }
163 }
164
165 fn describe(&self) -> String {
167 match self {
168 SchemaPredicate::FieldPresent(f) => format!("field_present({f})"),
169 SchemaPredicate::FieldAbsent(f) => format!("field_absent({f})"),
170 SchemaPredicate::AnyOf(fs) => format!("any_of([{}])", fs.join(", ")),
171 SchemaPredicate::Equals { field, value } => format!("{field} == \"{value}\""),
172 SchemaPredicate::Matches { field, regex } => {
173 format!("{field} matches /{}/", regex.as_str())
174 }
175 SchemaPredicate::Compare { field, op, value } => {
176 format!("{field} {} {value}", op.symbol())
177 }
178 SchemaPredicate::In { field, values } => format!("{field} in [{}]", values.join(", ")),
179 SchemaPredicate::FieldEqualsField { left, right } => format!("{left} == {right}"),
180 SchemaPredicate::Not(inner) => format!("not({})", inner.describe()),
181 SchemaPredicate::Any(ps) => format!(
182 "any({})",
183 ps.iter()
184 .map(|p| p.describe())
185 .collect::<Vec<_>>()
186 .join(" | ")
187 ),
188 SchemaPredicate::All(ps) => format!(
189 "all({})",
190 ps.iter()
191 .map(|p| p.describe())
192 .collect::<Vec<_>>()
193 .join(" & ")
194 ),
195 SchemaPredicate::HasAnyField => "has_any_field".to_string(),
196 }
197 }
198}
199
200#[derive(Debug, Clone)]
205pub struct SchemaSignature {
206 pub name: String,
208 pub predicates: Vec<SchemaPredicate>,
212 pub specificity: u32,
214}
215
216impl SchemaSignature {
217 fn matches<E: Event + ?Sized>(&self, event: &E) -> bool {
218 self.predicates.iter().all(|p| p.eval(event))
219 }
220
221 fn explain<E: Event + ?Sized>(&self, event: &E) -> SignatureExplanation {
222 let predicates: Vec<PredicateOutcome> = self
223 .predicates
224 .iter()
225 .map(|p| PredicateOutcome {
226 predicate: p.describe(),
227 matched: p.eval(event),
228 })
229 .collect();
230 let predicates_matched = predicates.iter().all(|p| p.matched);
231 SignatureExplanation {
232 name: self.name.clone(),
233 specificity: self.specificity,
234 predicates_matched,
235 predicates,
236 }
237 }
238}
239
240#[derive(Debug, Clone, Serialize)]
242pub struct PredicateOutcome {
243 pub predicate: String,
245 pub matched: bool,
247}
248
249#[derive(Debug, Clone, Serialize)]
251pub struct SignatureExplanation {
252 pub name: String,
254 pub specificity: u32,
256 pub predicates_matched: bool,
258 pub predicates: Vec<PredicateOutcome>,
260}
261
262#[derive(Debug, Clone, Serialize)]
267pub struct SchemaExplanation {
268 pub matched: Option<String>,
270 pub specificity: Option<u32>,
272 pub signature: Option<SignatureExplanation>,
275}
276
277#[derive(Debug, Clone, PartialEq, Eq)]
280pub struct SchemaMatch {
281 pub name: String,
282 pub specificity: u32,
283}
284
285#[derive(Debug, Clone)]
291pub struct SchemaClassifier {
292 signatures: Vec<SchemaSignature>,
293}
294
295impl SchemaClassifier {
296 pub fn new(mut signatures: Vec<SchemaSignature>) -> Self {
298 signatures.sort_by(|a, b| {
299 b.specificity
300 .cmp(&a.specificity)
301 .then_with(|| a.name.cmp(&b.name))
302 });
303 Self { signatures }
304 }
305
306 pub fn builtin() -> Self {
308 Self::new(builtin_signatures())
309 }
310
311 pub fn with_user_signatures(user: Vec<SchemaSignature>) -> Self {
315 let mut signatures = builtin_signatures();
316 signatures.extend(user);
317 Self::new(signatures)
318 }
319
320 pub fn classify<E: Event + ?Sized>(&self, event: &E) -> Option<SchemaMatch> {
323 self.signatures
324 .iter()
325 .find(|s| s.matches(event))
326 .map(|s| SchemaMatch {
327 name: s.name.clone(),
328 specificity: s.specificity,
329 })
330 }
331
332 pub fn classify_with_ambiguity<E: Event + ?Sized>(
338 &self,
339 event: &E,
340 ) -> (Option<SchemaMatch>, bool) {
341 let mut matching = self.signatures.iter().filter(|s| s.matches(event));
345 let Some(winner) = matching.next() else {
346 return (None, false);
347 };
348 let ambiguous = matching
349 .take_while(|s| s.specificity == winner.specificity)
350 .any(|s| s.name != winner.name);
351 (
352 Some(SchemaMatch {
353 name: winner.name.clone(),
354 specificity: winner.specificity,
355 }),
356 ambiguous,
357 )
358 }
359
360 pub fn classify_all<E: Event + ?Sized>(&self, event: &E) -> Vec<String> {
364 let mut out: Vec<String> = Vec::new();
365 for sig in self.signatures.iter().filter(|s| s.matches(event)) {
366 if !out.iter().any(|n| n == &sig.name) {
367 out.push(sig.name.clone());
368 }
369 }
370 out
371 }
372
373 pub fn explain<E: Event + ?Sized>(&self, event: &E) -> SchemaExplanation {
378 let mut best_near: Option<SignatureExplanation> = None;
379 let mut best_near_passing = 0usize;
380 for sig in &self.signatures {
381 let ex = sig.explain(event);
382 if ex.predicates_matched {
383 return SchemaExplanation {
384 matched: Some(ex.name.clone()),
385 specificity: Some(ex.specificity),
386 signature: Some(ex),
387 };
388 }
389 let passing = ex.predicates.iter().filter(|p| p.matched).count();
392 if best_near.is_none() || passing > best_near_passing {
393 best_near_passing = passing;
394 best_near = Some(ex);
395 }
396 }
397 SchemaExplanation {
398 matched: None,
399 specificity: None,
400 signature: best_near,
401 }
402 }
403
404 pub fn schema_names(&self) -> Vec<&str> {
406 let mut out: Vec<&str> = Vec::new();
407 for sig in &self.signatures {
408 if !out.contains(&sig.name.as_str()) {
409 out.push(sig.name.as_str());
410 }
411 }
412 out
413 }
414}
415
416impl Default for SchemaClassifier {
417 fn default() -> Self {
418 Self::builtin()
419 }
420}
421
422fn builtin_signatures() -> Vec<SchemaSignature> {
426 vec![
427 SchemaSignature {
432 name: "ecs_windows".to_string(),
433 specificity: 105,
434 predicates: vec![
435 SchemaPredicate::FieldPresent("ecs.version".to_string()),
436 SchemaPredicate::Any(vec![
437 SchemaPredicate::FieldPresent("winlog.channel".to_string()),
438 SchemaPredicate::FieldPresent("winlog.event_id".to_string()),
439 SchemaPredicate::Equals {
440 field: "host.os.type".to_string(),
441 value: "windows".to_string(),
442 },
443 SchemaPredicate::Equals {
444 field: "os.type".to_string(),
445 value: "windows".to_string(),
446 },
447 ]),
448 ],
449 },
450 SchemaSignature {
453 name: "ecs_linux".to_string(),
454 specificity: 105,
455 predicates: vec![
456 SchemaPredicate::FieldPresent("ecs.version".to_string()),
457 SchemaPredicate::Any(vec![
458 SchemaPredicate::Equals {
459 field: "host.os.type".to_string(),
460 value: "linux".to_string(),
461 },
462 SchemaPredicate::Equals {
463 field: "os.type".to_string(),
464 value: "linux".to_string(),
465 },
466 SchemaPredicate::FieldPresent("host.os.kernel".to_string()),
467 ]),
468 ],
469 },
470 SchemaSignature {
472 name: "ecs".to_string(),
473 specificity: 100,
474 predicates: vec![SchemaPredicate::FieldPresent("ecs.version".to_string())],
475 },
476 SchemaSignature {
478 name: "ocsf".to_string(),
479 specificity: 95,
480 predicates: vec![
481 SchemaPredicate::FieldPresent("class_uid".to_string()),
482 SchemaPredicate::FieldPresent("metadata.version".to_string()),
483 ],
484 },
485 SchemaSignature {
487 name: "windows_eventlog".to_string(),
488 specificity: 90,
489 predicates: vec![SchemaPredicate::AnyOf(vec![
490 "Event.System.EventID".to_string(),
491 "Event.System.Provider".to_string(),
492 ])],
493 },
494 SchemaSignature {
496 name: "sysmon".to_string(),
497 specificity: 88,
498 predicates: vec![SchemaPredicate::Equals {
499 field: "Channel".to_string(),
500 value: "Microsoft-Windows-Sysmon/Operational".to_string(),
501 }],
502 },
503 SchemaSignature {
505 name: "sysmon".to_string(),
506 specificity: 88,
507 predicates: vec![SchemaPredicate::Equals {
508 field: "Provider_Name".to_string(),
509 value: "Microsoft-Windows-Sysmon".to_string(),
510 }],
511 },
512 SchemaSignature {
514 name: "sysmon".to_string(),
515 specificity: 80,
516 predicates: vec![
517 SchemaPredicate::FieldPresent("EventID".to_string()),
518 SchemaPredicate::FieldPresent("ProcessGuid".to_string()),
519 SchemaPredicate::AnyOf(vec!["Image".to_string(), "CommandLine".to_string()]),
520 ],
521 },
522 SchemaSignature {
525 name: "cef".to_string(),
526 specificity: 85,
527 predicates: vec![
528 SchemaPredicate::FieldPresent("deviceVendor".to_string()),
529 SchemaPredicate::FieldPresent("deviceProduct".to_string()),
530 SchemaPredicate::FieldPresent("signatureId".to_string()),
531 ],
532 },
533 SchemaSignature {
539 name: "aws_vpcflow".to_string(),
540 specificity: 80,
541 predicates: vec![
542 SchemaPredicate::FieldPresent("srcaddr".to_string()),
543 SchemaPredicate::FieldPresent("dstaddr".to_string()),
544 SchemaPredicate::In {
545 field: "action".to_string(),
546 values: vec!["ACCEPT".to_string(), "REJECT".to_string()],
547 },
548 ],
549 },
550 SchemaSignature {
554 name: "aws_cloudtrail".to_string(),
555 specificity: 85,
556 predicates: vec![
557 SchemaPredicate::FieldPresent("eventVersion".to_string()),
558 SchemaPredicate::FieldPresent("eventSource".to_string()),
559 SchemaPredicate::FieldPresent("eventID".to_string()),
560 SchemaPredicate::FieldPresent("userIdentity".to_string()),
561 ],
562 },
563 SchemaSignature {
566 name: "onelogin_events".to_string(),
567 specificity: 85,
568 predicates: vec![
569 SchemaPredicate::FieldPresent("event_type_id".to_string()),
570 SchemaPredicate::FieldPresent("account_id".to_string()),
571 SchemaPredicate::AnyOf(vec!["user_id".to_string(), "actor_user_id".to_string()]),
572 ],
573 },
574 SchemaSignature {
578 name: "k8s_audit".to_string(),
579 specificity: 92,
580 predicates: vec![
581 SchemaPredicate::Equals {
582 field: "kind".to_string(),
583 value: "Event".to_string(),
584 },
585 SchemaPredicate::Matches {
586 field: "apiVersion".to_string(),
587 regex: regex::Regex::new("^audit\\.k8s\\.io/")
588 .expect("k8s audit apiVersion regex"),
589 },
590 SchemaPredicate::FieldPresent("auditID".to_string()),
591 ],
592 },
593 SchemaSignature {
597 name: "github_audit".to_string(),
598 specificity: 92,
599 predicates: vec![
600 SchemaPredicate::FieldPresent("action".to_string()),
601 SchemaPredicate::FieldPresent("actor".to_string()),
602 SchemaPredicate::AnyOf(vec!["org".to_string(), "repo".to_string()]),
603 SchemaPredicate::AnyOf(vec!["created_at".to_string(), "_document_id".to_string()]),
604 ],
605 },
606 SchemaSignature {
610 name: "okta_system_log".to_string(),
611 specificity: 88,
612 predicates: vec![
613 SchemaPredicate::FieldPresent("eventType".to_string()),
614 SchemaPredicate::FieldPresent("actor".to_string()),
615 SchemaPredicate::FieldPresent("published".to_string()),
616 SchemaPredicate::FieldPresent("outcome".to_string()),
617 ],
618 },
619 SchemaSignature {
622 name: "docker_events".to_string(),
623 specificity: 70,
624 predicates: vec![
625 SchemaPredicate::FieldPresent("Type".to_string()),
626 SchemaPredicate::FieldPresent("Action".to_string()),
627 SchemaPredicate::FieldPresent("Actor".to_string()),
628 ],
629 },
630 SchemaSignature {
634 name: "osquery_result".to_string(),
635 specificity: 75,
636 predicates: vec![
637 SchemaPredicate::FieldPresent("name".to_string()),
638 SchemaPredicate::In {
639 field: "action".to_string(),
640 values: vec![
641 "added".to_string(),
642 "removed".to_string(),
643 "snapshot".to_string(),
644 ],
645 },
646 SchemaPredicate::AnyOf(vec!["columns".to_string(), "snapshot".to_string()]),
647 SchemaPredicate::FieldPresent("hostIdentifier".to_string()),
648 ],
649 },
650 SchemaSignature {
653 name: "gcp_audit".to_string(),
654 specificity: 95,
655 predicates: vec![SchemaPredicate::Equals {
656 field: "protoPayload.@type".to_string(),
657 value: "type.googleapis.com/google.cloud.audit.AuditLog".to_string(),
658 }],
659 },
660 SchemaSignature {
665 name: "azure_activitylogs".to_string(),
666 specificity: 90,
667 predicates: vec![
668 SchemaPredicate::In {
669 field: "category".to_string(),
670 values: vec![
671 "Administrative".to_string(),
672 "Policy".to_string(),
673 "Security".to_string(),
674 ],
675 },
676 SchemaPredicate::Matches {
677 field: "id".to_string(),
678 regex: regex::Regex::new("(?i)^/subscriptions/")
679 .expect("Azure resourceId regex"),
680 },
681 SchemaPredicate::FieldPresent("operationName".to_string()),
682 ],
683 },
684 SchemaSignature {
687 name: "azure_auditlogs".to_string(),
688 specificity: 90,
689 predicates: vec![
690 SchemaPredicate::Equals {
691 field: "category".to_string(),
692 value: "AuditLogs".to_string(),
693 },
694 SchemaPredicate::FieldPresent("properties.activityDisplayName".to_string()),
695 ],
696 },
697 SchemaSignature {
700 name: "azure_signinlogs".to_string(),
701 specificity: 90,
702 predicates: vec![
703 SchemaPredicate::Equals {
704 field: "category".to_string(),
705 value: "SignInLogs".to_string(),
706 },
707 SchemaPredicate::FieldPresent("properties.userDisplayName".to_string()),
708 ],
709 },
710 SchemaSignature {
715 name: "azure".to_string(),
716 specificity: 65,
717 predicates: vec![SchemaPredicate::Matches {
718 field: "id".to_string(),
719 regex: regex::Regex::new("(?i)^/subscriptions/")
720 .expect("Azure subscriptionId regex"),
721 }],
722 },
723 SchemaSignature {
733 name: "m365_audit".to_string(),
734 specificity: 88,
735 predicates: vec![
736 SchemaPredicate::FieldPresent("RecordType".to_string()),
737 SchemaPredicate::FieldPresent("Operation".to_string()),
738 SchemaPredicate::FieldPresent("CreationTime".to_string()),
739 SchemaPredicate::FieldPresent("Workload".to_string()),
740 ],
741 },
742 SchemaSignature {
744 name: "generic_json".to_string(),
745 specificity: 0,
746 predicates: vec![SchemaPredicate::HasAnyField],
747 },
748 ]
749}
750
751pub fn builtin_schema_names() -> Vec<&'static str> {
755 vec![
756 "ecs_linux",
758 "ecs_windows",
759 "ecs",
761 "gcp_audit",
763 "ocsf",
764 "github_audit",
766 "k8s_audit",
767 "azure_activitylogs",
769 "azure_auditlogs",
770 "azure_signinlogs",
771 "windows_eventlog",
772 "m365_audit",
774 "okta_system_log",
775 "sysmon",
776 "aws_cloudtrail",
778 "cef",
779 "onelogin_events",
780 "aws_vpcflow",
782 "osquery_result",
784 "docker_events",
786 "azure",
788 "generic_json",
790 ]
791}
792
793fn builtin_schema_aliases() -> HashMap<String, String> {
799 HashMap::from([
800 ("ecs_windows".to_string(), "ecs".to_string()),
801 ("ecs_linux".to_string(), "ecs".to_string()),
802 ])
803}
804
805#[derive(Debug, thiserror::Error)]
811pub enum SchemaError {
812 #[error("cannot read schema signatures file '{path}': {source}")]
814 Io {
815 path: String,
816 #[source]
817 source: std::io::Error,
818 },
819 #[error("schema signatures YAML parse error: {0}")]
821 Parse(String),
822 #[error("invalid regex in schema '{name}': {error}")]
824 InvalidRegex { name: String, error: String },
825}
826
827#[derive(Debug, Clone, Deserialize)]
830#[serde(deny_unknown_fields)]
831pub struct FieldValueConfig {
832 pub field: String,
833 pub value: String,
834}
835
836#[derive(Debug, Clone, Deserialize)]
839#[serde(deny_unknown_fields)]
840pub struct FieldNumberConfig {
841 pub field: String,
842 pub value: f64,
843}
844
845#[derive(Debug, Clone, Deserialize)]
847#[serde(deny_unknown_fields)]
848pub struct FieldValuesConfig {
849 pub field: String,
850 pub values: Vec<String>,
851}
852
853#[derive(Debug, Clone, Deserialize)]
855#[serde(deny_unknown_fields)]
856pub struct FieldPairConfig {
857 pub left: String,
858 pub right: String,
859}
860
861#[derive(Debug, Clone, Default, Deserialize)]
866#[serde(deny_unknown_fields)]
867pub struct SchemaPredicateConfig {
868 #[serde(default)]
870 pub field_present: Option<String>,
871 #[serde(default)]
873 pub field_absent: Option<String>,
874 #[serde(default)]
876 pub any_of: Option<Vec<String>>,
877 #[serde(default)]
879 pub equals: Option<FieldValueConfig>,
880 #[serde(default)]
882 pub matches: Option<FieldValueConfig>,
883 #[serde(default)]
885 pub gt: Option<FieldNumberConfig>,
886 #[serde(default)]
888 pub gte: Option<FieldNumberConfig>,
889 #[serde(default)]
891 pub lt: Option<FieldNumberConfig>,
892 #[serde(default)]
894 pub lte: Option<FieldNumberConfig>,
895 #[serde(default, rename = "in")]
897 pub in_set: Option<FieldValuesConfig>,
898 #[serde(default)]
900 pub field_equals_field: Option<FieldPairConfig>,
901 #[serde(default)]
903 pub not: Option<Box<SchemaPredicateConfig>>,
904 #[serde(default)]
906 pub any: Option<Vec<SchemaPredicateConfig>>,
907 #[serde(default)]
909 pub all: Option<Vec<SchemaPredicateConfig>>,
910}
911
912impl SchemaPredicateConfig {
913 fn build(self, schema_name: &str) -> Result<SchemaPredicate, SchemaError> {
914 let mut chosen: Option<SchemaPredicate> = None;
915 let mut set = 0u32;
916 if let Some(f) = self.field_present {
917 set += 1;
918 chosen = Some(SchemaPredicate::FieldPresent(f));
919 }
920 if let Some(f) = self.field_absent {
921 set += 1;
922 chosen = Some(SchemaPredicate::FieldAbsent(f));
923 }
924 if let Some(fields) = self.any_of {
925 set += 1;
926 chosen = Some(SchemaPredicate::AnyOf(fields));
927 }
928 if let Some(fv) = self.equals {
929 set += 1;
930 chosen = Some(SchemaPredicate::Equals {
931 field: fv.field,
932 value: fv.value,
933 });
934 }
935 if let Some(fv) = self.matches {
936 set += 1;
937 chosen = Some(SchemaPredicate::Matches {
938 field: fv.field,
939 regex: Regex::new(&fv.value).map_err(|e| SchemaError::InvalidRegex {
940 name: schema_name.to_string(),
941 error: e.to_string(),
942 })?,
943 });
944 }
945 for (op, cfg) in [
946 (CompareOp::Gt, self.gt),
947 (CompareOp::Gte, self.gte),
948 (CompareOp::Lt, self.lt),
949 (CompareOp::Lte, self.lte),
950 ] {
951 if let Some(fv) = cfg {
952 set += 1;
953 chosen = Some(SchemaPredicate::Compare {
954 field: fv.field,
955 op,
956 value: fv.value,
957 });
958 }
959 }
960 if let Some(fv) = self.in_set {
961 set += 1;
962 chosen = Some(SchemaPredicate::In {
963 field: fv.field,
964 values: fv.values,
965 });
966 }
967 if let Some(fp) = self.field_equals_field {
968 set += 1;
969 chosen = Some(SchemaPredicate::FieldEqualsField {
970 left: fp.left,
971 right: fp.right,
972 });
973 }
974 if let Some(inner) = self.not {
975 set += 1;
976 chosen = Some(SchemaPredicate::Not(Box::new(inner.build(schema_name)?)));
977 }
978 if let Some(list) = self.any {
979 set += 1;
980 chosen = Some(SchemaPredicate::Any(build_group(list, schema_name, "any")?));
981 }
982 if let Some(list) = self.all {
983 set += 1;
984 chosen = Some(SchemaPredicate::All(build_group(list, schema_name, "all")?));
985 }
986 match (set, chosen) {
987 (1, Some(p)) => Ok(p),
988 (0, _) => Err(SchemaError::Parse(format!(
989 "schema '{schema_name}': a predicate has no condition (expected one of \
990 field_present, field_absent, any_of, equals, matches, gt, gte, lt, lte, \
991 in, field_equals_field, not, any, all)"
992 ))),
993 _ => Err(SchemaError::Parse(format!(
994 "schema '{schema_name}': a predicate sets multiple conditions; use one per list item"
995 ))),
996 }
997 }
998}
999
1000fn build_group(
1002 list: Vec<SchemaPredicateConfig>,
1003 schema_name: &str,
1004 kind: &str,
1005) -> Result<Vec<SchemaPredicate>, SchemaError> {
1006 if list.is_empty() {
1007 return Err(SchemaError::Parse(format!(
1008 "schema '{schema_name}': '{kind}' needs at least one sub-predicate"
1009 )));
1010 }
1011 list.into_iter().map(|p| p.build(schema_name)).collect()
1012}
1013
1014#[derive(Debug, Clone, Deserialize)]
1016pub struct SchemaSignatureConfig {
1017 pub name: String,
1019 #[serde(default = "default_user_specificity")]
1022 pub specificity: u32,
1023 #[serde(default, rename = "match")]
1025 pub predicates: Vec<SchemaPredicateConfig>,
1026}
1027
1028fn default_user_specificity() -> u32 {
1029 50
1030}
1031
1032#[derive(Debug, Clone, Default, Deserialize)]
1035pub struct SchemaSignaturesFile {
1036 #[serde(default)]
1037 pub schemas: Vec<SchemaSignatureConfig>,
1038 #[serde(default)]
1039 pub routing: Option<RoutingConfig>,
1040}
1041
1042impl SchemaSignatureConfig {
1043 fn build(self) -> Result<SchemaSignature, SchemaError> {
1044 let name = self.name;
1045 let predicates = self
1046 .predicates
1047 .into_iter()
1048 .map(|p| p.build(&name))
1049 .collect::<Result<Vec<_>, _>>()?;
1050 Ok(SchemaSignature {
1051 name,
1052 predicates,
1053 specificity: self.specificity,
1054 })
1055 }
1056}
1057
1058pub fn parse_schema_signatures(yaml: &str) -> Result<Vec<SchemaSignature>, SchemaError> {
1060 let file: SchemaSignaturesFile =
1061 yaml_serde::from_str(yaml).map_err(|e| SchemaError::Parse(e.to_string()))?;
1062 file.schemas.into_iter().map(|s| s.build()).collect()
1063}
1064
1065pub fn load_schema_signatures(path: &Path) -> Result<Vec<SchemaSignature>, SchemaError> {
1067 let content = fs::read_to_string(path).map_err(|e| SchemaError::Io {
1068 path: path.display().to_string(),
1069 source: e,
1070 })?;
1071 parse_schema_signatures(&content)
1072}
1073
1074pub fn parse_schema_config(
1077 yaml: &str,
1078) -> Result<(Vec<SchemaSignature>, Option<RoutingConfig>), SchemaError> {
1079 let file: SchemaSignaturesFile =
1080 yaml_serde::from_str(yaml).map_err(|e| SchemaError::Parse(e.to_string()))?;
1081 let signatures = file
1082 .schemas
1083 .into_iter()
1084 .map(|s| s.build())
1085 .collect::<Result<Vec<_>, _>>()?;
1086 Ok((signatures, file.routing))
1087}
1088
1089pub fn load_schema_config(
1092 path: &Path,
1093) -> Result<(Vec<SchemaSignature>, Option<RoutingConfig>), SchemaError> {
1094 let content = fs::read_to_string(path).map_err(|e| SchemaError::Io {
1095 path: path.display().to_string(),
1096 source: e,
1097 })?;
1098 parse_schema_config(&content)
1099}
1100
1101pub fn validate_schema_config(
1115 user_signatures: &[SchemaSignature],
1116 routing: Option<&RoutingConfig>,
1117) -> Vec<String> {
1118 let mut findings = Vec::new();
1119
1120 let mut all = builtin_signatures();
1122 all.extend(user_signatures.iter().cloned());
1123 let preds = |s: &SchemaSignature| -> Vec<String> {
1124 s.predicates.iter().map(|p| p.describe()).collect()
1125 };
1126
1127 for i in 0..user_signatures.len() {
1129 for j in (i + 1)..user_signatures.len() {
1130 if user_signatures[i].name == user_signatures[j].name
1131 && preds(&user_signatures[i]) == preds(&user_signatures[j])
1132 {
1133 findings.push(format!(
1134 "duplicate signature '{}' with identical predicates",
1135 user_signatures[i].name
1136 ));
1137 }
1138 }
1139 }
1140
1141 for b in &all {
1143 let b_preds = preds(b);
1144 for a in &all {
1145 if a.name != b.name
1146 && a.specificity > b.specificity
1147 && !a.predicates.is_empty()
1148 && preds(a).iter().all(|p| b_preds.contains(p))
1149 {
1150 findings.push(format!(
1151 "signature '{}' (specificity {}) is unreachable: shadowed by '{}' (specificity {}) whose predicates are a subset",
1152 b.name, b.specificity, a.name, a.specificity
1153 ));
1154 break;
1155 }
1156 }
1157 }
1158
1159 if let Some(routing) = routing {
1161 let mut known: std::collections::HashSet<&str> =
1162 builtin_schema_names().into_iter().collect();
1163 for s in user_signatures {
1164 known.insert(s.name.as_str());
1165 }
1166 let mut seen: std::collections::HashSet<&str> = std::collections::HashSet::new();
1167 for binding in &routing.bindings {
1168 if !known.contains(binding.schema.as_str()) {
1169 findings.push(format!(
1170 "routing binding references unknown schema '{}' (no built-in or user signature produces it)",
1171 binding.schema
1172 ));
1173 }
1174 if !seen.insert(binding.schema.as_str()) {
1175 findings.push(format!(
1176 "duplicate routing binding for schema '{}'",
1177 binding.schema
1178 ));
1179 }
1180 }
1181 for (alias, canonical) in &routing.aliases {
1182 if !known.contains(canonical.as_str()) {
1183 findings.push(format!(
1184 "alias '{alias}' targets unknown schema '{canonical}' (no built-in or user signature produces it)"
1185 ));
1186 }
1187 }
1188 }
1189
1190 findings
1191}
1192
1193#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Deserialize)]
1199#[serde(rename_all = "snake_case")]
1200pub enum OnUnknown {
1201 #[default]
1203 Warn,
1204 Drop,
1206 Passthrough,
1208 Error,
1210}
1211
1212#[derive(Debug, Clone, Default, Deserialize)]
1216#[serde(deny_unknown_fields)]
1217pub struct SchemaLogsource {
1218 #[serde(default)]
1219 pub product: Option<String>,
1220 #[serde(default)]
1221 pub service: Option<String>,
1222 #[serde(default)]
1223 pub category: Option<String>,
1224 #[serde(default)]
1225 pub custom: HashMap<String, String>,
1226}
1227
1228impl SchemaLogsource {
1229 fn to_logsource(&self) -> LogSource {
1230 LogSource {
1231 product: self.product.clone(),
1232 service: self.service.clone(),
1233 category: self.category.clone(),
1234 custom: self.custom.clone(),
1235 ..LogSource::default()
1236 }
1237 }
1238}
1239
1240#[derive(Debug, Clone, Deserialize)]
1243pub struct SchemaBinding {
1244 pub schema: String,
1245 #[serde(default)]
1247 pub pipelines: Vec<String>,
1248 #[serde(default)]
1251 pub logsource: Option<SchemaLogsource>,
1252}
1253
1254pub fn builtin_schema_logsource() -> HashMap<String, LogSource> {
1267 fn ls(product: &str, service: Option<&str>) -> LogSource {
1268 LogSource {
1269 product: Some(product.to_string()),
1270 service: service.map(str::to_string),
1271 ..LogSource::default()
1272 }
1273 }
1274 fn ls_custom(product: Option<&str>, custom: HashMap<&str, String>) -> LogSource {
1275 LogSource {
1276 product: product.map(str::to_string),
1277 service: None,
1278 category: None,
1279 custom: custom
1280 .into_iter()
1281 .map(|(k, v)| (k.to_string(), v))
1282 .collect(),
1283 ..LogSource::default()
1284 }
1285 }
1286 let mut map = HashMap::new();
1287
1288 map.insert("sysmon".to_string(), ls("windows", Some("sysmon")));
1290 map.insert("windows_eventlog".to_string(), ls("windows", None));
1291 map.insert("ecs_windows".to_string(), ls("windows", None));
1292 map.insert("ecs_linux".to_string(), ls("linux", None));
1293
1294 map.insert("aws_cloudtrail".to_string(), ls("aws", Some("cloudtrail")));
1296 map.insert(
1298 "aws_vpcflow".to_string(),
1299 ls_custom(
1300 Some("aws"),
1301 HashMap::from([("source", "vpcflow".to_string())]),
1302 ),
1303 );
1304
1305 map.insert(
1307 "azure_activitylogs".to_string(),
1308 ls("azure", Some("activitylogs")),
1309 );
1310 map.insert(
1311 "azure_auditlogs".to_string(),
1312 ls("azure", Some("auditlogs")),
1313 );
1314 map.insert(
1315 "azure_signinlogs".to_string(),
1316 ls("azure", Some("signinlogs")),
1317 );
1318
1319 map.insert("gcp_audit".to_string(), ls("gcp", Some("gcp.audit")));
1321
1322 map.insert("m365_audit".to_string(), ls("m365", Some("audit")));
1324
1325 map.insert("github_audit".to_string(), ls("github", Some("audit")));
1327 map.insert("okta_system_log".to_string(), ls("okta", Some("okta")));
1328 map.insert(
1329 "onelogin_events".to_string(),
1330 ls("onelogin", Some("onelogin.events")),
1331 );
1332
1333 map.insert(
1335 "k8s_audit".to_string(),
1336 ls_custom(
1337 None,
1338 HashMap::from([
1339 ("platform", "kubernetes".to_string()),
1340 ("source", "k8s.audit".to_string()),
1341 ]),
1342 ),
1343 );
1344 map.insert(
1345 "docker_events".to_string(),
1346 ls_custom(
1347 None,
1348 HashMap::from([
1349 ("platform", "docker".to_string()),
1350 ("source", "docker.events".to_string()),
1351 ]),
1352 ),
1353 );
1354 map.insert(
1355 "osquery_result".to_string(),
1356 ls_custom(
1357 None,
1358 HashMap::from([
1359 ("platform", "osquery".to_string()),
1360 ("source", "osquery.result".to_string()),
1361 ]),
1362 ),
1363 );
1364
1365 map
1366}
1367
1368#[derive(Debug, Clone, Default, Deserialize)]
1370pub struct RoutingConfig {
1371 #[serde(default)]
1372 pub on_unknown: OnUnknown,
1373 #[serde(default)]
1374 pub bindings: Vec<SchemaBinding>,
1375 #[serde(default)]
1378 pub default_pipelines: Vec<String>,
1379 #[serde(default)]
1384 pub aliases: HashMap<String, String>,
1385}
1386
1387#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1389pub enum RouteDecision {
1390 Evaluate { set: usize, unknown: bool },
1393 Drop,
1395 Error,
1397}
1398
1399#[derive(Debug, Clone)]
1406pub struct RoutingPlan {
1407 pipeline_sets: Vec<Vec<String>>,
1409 schema_to_set: HashMap<String, usize>,
1411 schema_logsource: HashMap<String, LogSource>,
1415 aliases: HashMap<String, String>,
1419 on_unknown: OnUnknown,
1420}
1421
1422impl RoutingPlan {
1423 pub fn from_config(config: &RoutingConfig) -> Self {
1426 let mut pipeline_sets: Vec<Vec<String>> = vec![config.default_pipelines.clone()];
1428 let mut schema_to_set: HashMap<String, usize> = HashMap::new();
1429 let mut schema_logsource = builtin_schema_logsource();
1432 let mut aliases = builtin_schema_aliases();
1434 for (alias, canonical) in &config.aliases {
1435 aliases.insert(alias.clone(), canonical.clone());
1436 }
1437
1438 for binding in &config.bindings {
1439 let idx = pipeline_sets
1440 .iter()
1441 .position(|s| s == &binding.pipelines)
1442 .unwrap_or_else(|| {
1443 pipeline_sets.push(binding.pipelines.clone());
1444 pipeline_sets.len() - 1
1445 });
1446 schema_to_set.insert(binding.schema.clone(), idx);
1447 if let Some(ls) = &binding.logsource {
1448 schema_logsource.insert(binding.schema.clone(), ls.to_logsource());
1449 }
1450 }
1451
1452 RoutingPlan {
1453 pipeline_sets,
1454 schema_to_set,
1455 schema_logsource,
1456 aliases,
1457 on_unknown: config.on_unknown,
1458 }
1459 }
1460
1461 pub fn pipeline_sets(&self) -> &[Vec<String>] {
1464 &self.pipeline_sets
1465 }
1466
1467 pub fn on_unknown(&self) -> OnUnknown {
1469 self.on_unknown
1470 }
1471
1472 pub fn schema_logsource(&self, schema: &str) -> Option<&LogSource> {
1476 self.schema_logsource.get(schema)
1477 }
1478
1479 pub fn schemas_with_logsource(&self) -> Vec<String> {
1482 let mut names: Vec<String> = self.schema_logsource.keys().cloned().collect();
1483 names.sort();
1484 names
1485 }
1486
1487 pub fn set_product_partition(&self) -> Vec<Option<std::collections::HashSet<String>>> {
1499 use std::collections::HashSet;
1500 let n = self.pipeline_sets.len();
1501 let mut out: Vec<Option<HashSet<String>>> = (0..n).map(|_| Some(HashSet::new())).collect();
1502 if let Some(first) = out.get_mut(0) {
1503 *first = None;
1504 }
1505
1506 let mut routes: Vec<(usize, &str)> = self
1509 .schema_to_set
1510 .iter()
1511 .map(|(s, &set)| (set, s.as_str()))
1512 .collect();
1513 for (alias, canonical) in &self.aliases {
1514 if !self.schema_to_set.contains_key(alias)
1515 && let Some(&set) = self.schema_to_set.get(canonical)
1516 {
1517 routes.push((set, alias.as_str()));
1518 }
1519 }
1520
1521 for (set, schema) in routes {
1522 if set == 0 {
1523 continue;
1524 }
1525 let product = self
1526 .schema_logsource
1527 .get(schema)
1528 .and_then(|ls| ls.product.as_deref());
1529 let Some(slot) = out.get_mut(set) else {
1530 continue;
1531 };
1532 match product {
1533 Some(p) => {
1534 if let Some(products) = slot {
1535 products.insert(p.to_ascii_lowercase());
1536 }
1537 }
1538 None => *slot = None,
1539 }
1540 }
1541 out
1542 }
1543
1544 pub fn decide(&self, schema: Option<&str>) -> RouteDecision {
1547 match schema {
1548 Some(s) if self.schema_to_set.contains_key(s) => RouteDecision::Evaluate {
1550 set: self.schema_to_set[s],
1551 unknown: false,
1552 },
1553 Some(s)
1557 if self
1558 .aliases
1559 .get(s)
1560 .and_then(|canonical| self.schema_to_set.get(canonical))
1561 .is_some() =>
1562 {
1563 let canonical = &self.aliases[s];
1564 RouteDecision::Evaluate {
1565 set: self.schema_to_set[canonical],
1566 unknown: false,
1567 }
1568 }
1569 Some(_) => RouteDecision::Evaluate {
1571 set: 0,
1572 unknown: false,
1573 },
1574 None => match self.on_unknown {
1576 OnUnknown::Warn | OnUnknown::Passthrough => RouteDecision::Evaluate {
1577 set: 0,
1578 unknown: true,
1579 },
1580 OnUnknown::Drop => RouteDecision::Drop,
1581 OnUnknown::Error => RouteDecision::Error,
1582 },
1583 }
1584 }
1585}
1586
1587#[derive(Debug, Clone, PartialEq, Eq)]
1593pub struct SchemaCountEntry {
1594 pub schema: String,
1596 pub count: u64,
1598}
1599
1600#[derive(Debug, Clone, PartialEq, Eq)]
1602pub struct UnknownShapeEntry {
1603 pub keys: Vec<String>,
1606 pub count: u64,
1608}
1609
1610const UNKNOWN_SHAPE_CAP: usize = 200;
1612const UNKNOWN_SHAPE_MAX_KEYS: usize = 64;
1614
1615#[derive(Debug, Clone, Default)]
1617pub struct SchemaObservation {
1618 pub by_schema: Vec<SchemaCountEntry>,
1620 pub classified: u64,
1622 pub unknown: u64,
1624 pub ambiguous: u64,
1627 pub unknown_shapes: Vec<UnknownShapeEntry>,
1630 pub unrecognized_shapes: Vec<UnknownShapeEntry>,
1634 pub events_observed: u64,
1636 pub lifetime_classified: u64,
1639 pub lifetime_unknown: u64,
1641 pub lifetime_ambiguous: u64,
1643 pub uptime_seconds: f64,
1645}
1646
1647pub struct SchemaObserver {
1653 classifier: SchemaClassifier,
1654 counts: Mutex<HashMap<String, u64>>,
1655 unknown: AtomicU64,
1656 ambiguous: AtomicU64,
1657 unknown_shapes: Mutex<HashMap<Vec<String>, u64>>,
1660 discovery_sampling: bool,
1667 unrecognized_shapes: Mutex<HashMap<Vec<String>, u64>>,
1670 lifetime_classified: AtomicU64,
1671 lifetime_unknown: AtomicU64,
1672 lifetime_ambiguous: AtomicU64,
1673 start: Mutex<Instant>,
1674}
1675
1676impl SchemaObserver {
1677 pub fn new(classifier: SchemaClassifier) -> Self {
1680 Self::new_with_discovery(classifier, false)
1681 }
1682
1683 pub fn new_with_discovery(classifier: SchemaClassifier, discovery_sampling: bool) -> Self {
1687 Self {
1688 classifier,
1689 counts: Mutex::new(HashMap::new()),
1690 unknown: AtomicU64::new(0),
1691 ambiguous: AtomicU64::new(0),
1692 unknown_shapes: Mutex::new(HashMap::new()),
1693 discovery_sampling,
1694 unrecognized_shapes: Mutex::new(HashMap::new()),
1695 lifetime_classified: AtomicU64::new(0),
1696 lifetime_unknown: AtomicU64::new(0),
1697 lifetime_ambiguous: AtomicU64::new(0),
1698 start: Mutex::new(Instant::now()),
1699 }
1700 }
1701
1702 pub fn discovery_sampling(&self) -> bool {
1705 self.discovery_sampling
1706 }
1707
1708 pub fn builtin() -> Self {
1710 Self::new(SchemaClassifier::builtin())
1711 }
1712
1713 pub fn observe<E: Event + ?Sized>(&self, event: &E) {
1716 let (matched, ambiguous) = self.classifier.classify_with_ambiguity(event);
1717 if ambiguous {
1718 self.ambiguous.fetch_add(1, Ordering::Relaxed);
1719 self.lifetime_ambiguous.fetch_add(1, Ordering::Relaxed);
1720 }
1721 let discovery_unrecognized = match &matched {
1725 None => true,
1726 Some(m) => m.name == "generic_json",
1727 };
1728 if self.discovery_sampling && discovery_unrecognized {
1729 self.record_unrecognized_shape(event);
1730 }
1731
1732 match matched {
1733 Some(m) => {
1734 self.lifetime_classified.fetch_add(1, Ordering::Relaxed);
1735 let mut counts = self.counts.lock().expect("schema observer mutex poisoned");
1736 *counts.entry(m.name).or_insert(0) += 1;
1737 }
1738 None => {
1739 self.unknown.fetch_add(1, Ordering::Relaxed);
1740 self.lifetime_unknown.fetch_add(1, Ordering::Relaxed);
1741 self.record_unknown_shape(event);
1742 }
1743 }
1744 }
1745
1746 fn record_unknown_shape<E: Event + ?Sized>(&self, event: &E) {
1749 let mut keys: Vec<String> = event.field_keys().iter().map(|k| k.to_string()).collect();
1750 keys.sort();
1751 keys.dedup();
1752 keys.truncate(UNKNOWN_SHAPE_MAX_KEYS);
1753 let mut shapes = self
1754 .unknown_shapes
1755 .lock()
1756 .expect("schema observer shapes mutex poisoned");
1757 if shapes.contains_key(&keys) || shapes.len() < UNKNOWN_SHAPE_CAP {
1759 *shapes.entry(keys).or_insert(0) += 1;
1760 }
1761 }
1762
1763 fn record_unrecognized_shape<E: Event + ?Sized>(&self, event: &E) {
1767 let mut keys: Vec<String> = event.field_keys().iter().map(|k| k.to_string()).collect();
1768 keys.sort();
1769 keys.dedup();
1770 keys.truncate(UNKNOWN_SHAPE_MAX_KEYS);
1771 if keys.is_empty() {
1772 return;
1773 }
1774 let mut shapes = self
1775 .unrecognized_shapes
1776 .lock()
1777 .expect("schema observer shapes mutex poisoned");
1778 if shapes.contains_key(&keys) || shapes.len() < UNKNOWN_SHAPE_CAP {
1779 *shapes.entry(keys).or_insert(0) += 1;
1780 }
1781 }
1782
1783 pub fn snapshot(&self) -> SchemaObservation {
1785 let counts = self.counts.lock().expect("schema observer mutex poisoned");
1786 let mut by_schema: Vec<SchemaCountEntry> = counts
1787 .iter()
1788 .map(|(schema, count)| SchemaCountEntry {
1789 schema: schema.clone(),
1790 count: *count,
1791 })
1792 .collect();
1793 let classified: u64 = counts.values().sum();
1794 drop(counts);
1795 by_schema.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.schema.cmp(&b.schema)));
1796
1797 let shapes = self
1798 .unknown_shapes
1799 .lock()
1800 .expect("schema observer shapes mutex poisoned");
1801 let mut unknown_shapes: Vec<UnknownShapeEntry> = shapes
1802 .iter()
1803 .map(|(keys, count)| UnknownShapeEntry {
1804 keys: keys.clone(),
1805 count: *count,
1806 })
1807 .collect();
1808 drop(shapes);
1809 unknown_shapes.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.keys.cmp(&b.keys)));
1810
1811 let unrec = self
1812 .unrecognized_shapes
1813 .lock()
1814 .expect("schema observer shapes mutex poisoned");
1815 let mut unrecognized_shapes: Vec<UnknownShapeEntry> = unrec
1816 .iter()
1817 .map(|(keys, count)| UnknownShapeEntry {
1818 keys: keys.clone(),
1819 count: *count,
1820 })
1821 .collect();
1822 drop(unrec);
1823 unrecognized_shapes.sort_by(|a, b| b.count.cmp(&a.count).then_with(|| a.keys.cmp(&b.keys)));
1824
1825 let unknown = self.unknown.load(Ordering::Relaxed);
1826 SchemaObservation {
1827 by_schema,
1828 classified,
1829 unknown,
1830 ambiguous: self.ambiguous.load(Ordering::Relaxed),
1831 unknown_shapes,
1832 unrecognized_shapes,
1833 events_observed: classified + unknown,
1838 lifetime_classified: self.lifetime_classified.load(Ordering::Relaxed),
1839 lifetime_unknown: self.lifetime_unknown.load(Ordering::Relaxed),
1840 lifetime_ambiguous: self.lifetime_ambiguous.load(Ordering::Relaxed),
1841 uptime_seconds: self
1842 .start
1843 .lock()
1844 .expect("schema observer start mutex poisoned")
1845 .elapsed()
1846 .as_secs_f64(),
1847 }
1848 }
1849
1850 pub fn reset(&self) -> (u64, u64) {
1853 let mut counts = self.counts.lock().expect("schema observer mutex poisoned");
1854 let previous_classified: u64 = counts.values().sum();
1855 counts.clear();
1856 drop(counts);
1857 self.unknown_shapes
1858 .lock()
1859 .expect("schema observer shapes mutex poisoned")
1860 .clear();
1861 self.unrecognized_shapes
1862 .lock()
1863 .expect("schema observer shapes mutex poisoned")
1864 .clear();
1865 let previous_unknown = self.unknown.swap(0, Ordering::Relaxed);
1866 self.ambiguous.store(0, Ordering::Relaxed);
1867 *self
1868 .start
1869 .lock()
1870 .expect("schema observer start mutex poisoned") = Instant::now();
1871 (previous_classified, previous_unknown)
1872 }
1873
1874 pub fn lifetime_classified(&self) -> u64 {
1876 self.lifetime_classified.load(Ordering::Relaxed)
1877 }
1878
1879 pub fn lifetime_unknown(&self) -> u64 {
1881 self.lifetime_unknown.load(Ordering::Relaxed)
1882 }
1883
1884 pub fn lifetime_ambiguous(&self) -> u64 {
1886 self.lifetime_ambiguous.load(Ordering::Relaxed)
1887 }
1888}
1889
1890#[cfg(test)]
1891mod tests {
1892 use super::*;
1893 use crate::event::JsonEvent;
1894 use serde_json::json;
1895
1896 fn classify(value: &serde_json::Value) -> Option<String> {
1897 SchemaClassifier::builtin()
1898 .classify(&JsonEvent::borrow(value))
1899 .map(|m| m.name)
1900 }
1901
1902 #[test]
1903 fn recognizes_ecs_by_version_marker() {
1904 let v = json!({"ecs": {"version": "8.11.0"}, "process": {"command_line": "whoami"}});
1905 assert_eq!(classify(&v).as_deref(), Some("ecs"));
1906 }
1907
1908 #[test]
1909 fn recognizes_ecs_with_flattened_keys() {
1910 let v = json!({"ecs.version": "8.11.0", "process.command_line": "whoami"});
1911 assert_eq!(classify(&v).as_deref(), Some("ecs"));
1912 }
1913
1914 #[test]
1915 fn recognizes_ocsf_by_class_and_metadata() {
1916 let v = json!({"class_uid": 1001, "category_uid": 1, "metadata": {"version": "1.1.0"}});
1917 assert_eq!(classify(&v).as_deref(), Some("ocsf"));
1918 }
1919
1920 #[test]
1921 fn recognizes_rendered_windows_event_log() {
1922 let v = json!({"Event": {"System": {"EventID": 4688, "Provider": "Microsoft-Windows-Security-Auditing"}}});
1923 assert_eq!(classify(&v).as_deref(), Some("windows_eventlog"));
1924 }
1925
1926 #[test]
1927 fn recognizes_sysmon_by_channel() {
1928 let v = json!({"Channel": "Microsoft-Windows-Sysmon/Operational", "EventID": 1, "Image": "C:/cmd.exe"});
1929 assert_eq!(classify(&v).as_deref(), Some("sysmon"));
1930 }
1931
1932 #[test]
1933 fn recognizes_sysmon_by_provider() {
1934 let v = json!({"Provider_Name": "Microsoft-Windows-Sysmon", "EventID": 3});
1935 assert_eq!(classify(&v).as_deref(), Some("sysmon"));
1936 }
1937
1938 #[test]
1939 fn recognizes_flat_sysmon_by_field_shape() {
1940 let v = json!({"EventID": 1, "ProcessGuid": "{abc}", "CommandLine": "cmd /c whoami"});
1941 assert_eq!(classify(&v).as_deref(), Some("sysmon"));
1942 }
1943
1944 #[test]
1945 fn recognizes_cef_structured_fields() {
1946 let v = json!({"deviceVendor": "Security", "deviceProduct": "IDS", "signatureId": "100", "src": "10.0.0.1"});
1947 assert_eq!(classify(&v).as_deref(), Some("cef"));
1948 }
1949
1950 #[test]
1951 fn falls_back_to_generic_json_for_unrecognized_structured_events() {
1952 let v = json!({"some_vendor_field": "x", "another": 1});
1953 assert_eq!(classify(&v).as_deref(), Some("generic_json"));
1954 }
1955
1956 #[test]
1957 fn fieldless_events_are_unknown() {
1958 assert_eq!(classify(&json!({})), None);
1960 assert_eq!(classify(&json!("just a string")), None);
1962 }
1963
1964 #[test]
1965 fn specificity_prefers_specific_schema_over_generic() {
1966 let v = json!({"ecs.version": "8.0.0", "vendor_blob": {"x": 1}});
1968 let cls = SchemaClassifier::builtin();
1969 let m = cls.classify(&JsonEvent::borrow(&v)).unwrap();
1970 assert_eq!(m.name, "ecs");
1971 assert_eq!(m.specificity, 100);
1972 let all = cls.classify_all(&JsonEvent::borrow(&v));
1974 assert_eq!(all.first().map(String::as_str), Some("ecs"));
1975 assert!(all.iter().any(|n| n == "generic_json"));
1976 }
1977
1978 #[test]
1979 fn schema_names_lists_builtins_most_specific_first() {
1980 let classifier = SchemaClassifier::builtin();
1981 let names = classifier.schema_names();
1982 assert_eq!(names.first(), Some(&"ecs_linux"));
1985 assert!(names.contains(&"ecs_windows"));
1986 assert!(names.contains(&"ecs"));
1987 assert!(names.contains(&"generic_json"));
1988 assert_eq!(names.last(), Some(&"generic_json"));
1990 }
1991
1992 #[test]
1993 fn ecs_windows_specialization_classifies_and_aliases_to_ecs() {
1994 let v = json!({"ecs.version": "8.11.0", "winlog": {"channel": "Security"}});
1997 assert_eq!(classify(&v).as_deref(), Some("ecs_windows"));
1998 let plain = json!({"ecs.version": "8.11.0", "process": {"command_line": "whoami"}});
2000 assert_eq!(classify(&plain).as_deref(), Some("ecs"));
2001
2002 let config = RoutingConfig {
2005 on_unknown: OnUnknown::Warn,
2006 default_pipelines: vec![],
2007 aliases: HashMap::new(),
2008 bindings: vec![SchemaBinding {
2009 schema: "ecs".to_string(),
2010 pipelines: vec!["ecs_windows".to_string()],
2011 logsource: None,
2012 }],
2013 };
2014 let plan = RoutingPlan::from_config(&config);
2015 let ecs_set = match plan.decide(Some("ecs")) {
2016 RouteDecision::Evaluate { set, .. } => set,
2017 other => panic!("unexpected: {other:?}"),
2018 };
2019 assert_eq!(plan.decide(Some("ecs_windows")), plan.decide(Some("ecs")));
2021 assert_ne!(ecs_set, 0, "ecs binding is a non-default set");
2022 assert_eq!(
2023 plan.schema_logsource("ecs_windows")
2024 .and_then(|l| l.product.as_deref()),
2025 Some("windows")
2026 );
2027 }
2028
2029 #[test]
2030 fn user_alias_routes_as_canonical() {
2031 let yaml = r#"
2032schemas:
2033 - name: my_win
2034 specificity: 70
2035 match:
2036 - field_present: vendor.win_marker
2037routing:
2038 aliases:
2039 my_win: ecs
2040 bindings:
2041 - schema: ecs
2042 pipelines: [ecs_windows]
2043"#;
2044 let (_sigs, routing) = parse_schema_config(yaml).unwrap();
2045 let plan = RoutingPlan::from_config(&routing.expect("routing"));
2046 assert_eq!(plan.decide(Some("my_win")), plan.decide(Some("ecs")));
2048 assert!(matches!(
2049 plan.decide(Some("my_win")),
2050 RouteDecision::Evaluate { unknown: false, .. }
2051 ));
2052 }
2053
2054 #[test]
2055 fn set_product_partition_only_for_platform_locked_sets() {
2056 let config = RoutingConfig {
2057 on_unknown: OnUnknown::Warn,
2058 default_pipelines: vec![],
2059 aliases: HashMap::new(),
2060 bindings: vec![
2061 SchemaBinding {
2062 schema: "sysmon".to_string(),
2063 pipelines: vec!["p_sysmon".to_string()],
2064 logsource: None,
2065 },
2066 SchemaBinding {
2067 schema: "ecs".to_string(),
2068 pipelines: vec!["p_ecs".to_string()],
2069 logsource: None,
2070 },
2071 ],
2072 };
2073 let plan = RoutingPlan::from_config(&config);
2074 let part = plan.set_product_partition();
2075 assert!(part[0].is_none(), "default set is never partitioned");
2076
2077 let set_of = |schema| match plan.decide(Some(schema)) {
2078 RouteDecision::Evaluate { set, .. } => set,
2079 other => panic!("unexpected: {other:?}"),
2080 };
2081 let sysmon_set = set_of("sysmon");
2083 assert_eq!(
2084 part[sysmon_set].as_ref().map(|s| s.contains("windows")),
2085 Some(true)
2086 );
2087 assert!(part[set_of("ecs")].is_none());
2089 }
2090
2091 #[test]
2092 fn parses_user_signatures_from_yaml() {
2093 let yaml = r#"
2094schemas:
2095 - name: my_vendor
2096 specificity: 70
2097 match:
2098 - field_present: vendor.product
2099 - equals:
2100 field: event_type
2101 value: alert
2102 - any_of: [a, b]
2103"#;
2104 let sigs = parse_schema_signatures(yaml).expect("parse");
2105 assert_eq!(sigs.len(), 1);
2106 assert_eq!(sigs[0].name, "my_vendor");
2107 assert_eq!(sigs[0].specificity, 70);
2108 assert_eq!(sigs[0].predicates.len(), 3);
2109
2110 let cls = SchemaClassifier::with_user_signatures(sigs);
2111 let v = json!({"vendor": {"product": "X"}, "event_type": "ALERT", "a": 1});
2112 assert_eq!(
2113 cls.classify(&JsonEvent::borrow(&v))
2114 .map(|m| m.name)
2115 .as_deref(),
2116 Some("my_vendor")
2117 );
2118 }
2119
2120 #[test]
2121 fn user_signature_with_invalid_regex_is_rejected() {
2122 let yaml = r#"
2123schemas:
2124 - name: bad
2125 match:
2126 - matches:
2127 field: msg
2128 value: "([unclosed"
2129"#;
2130 let err = parse_schema_signatures(yaml).unwrap_err();
2131 assert!(matches!(err, SchemaError::InvalidRegex { .. }));
2132 }
2133
2134 #[test]
2135 fn user_regex_signature_matches_field_value() {
2136 let yaml = r#"
2137schemas:
2138 - name: cef_raw
2139 specificity: 60
2140 match:
2141 - matches:
2142 field: message
2143 value: "^CEF:\\d"
2144"#;
2145 let sigs = parse_schema_signatures(yaml).expect("parse");
2146 let cls = SchemaClassifier::with_user_signatures(sigs);
2147 let v = json!({"message": "CEF:0|Vendor|Product|1.0|100|Name|9|src=1.2.3.4"});
2148 assert_eq!(
2149 cls.classify(&JsonEvent::borrow(&v))
2150 .map(|m| m.name)
2151 .as_deref(),
2152 Some("cef_raw")
2153 );
2154 }
2155
2156 fn classifier_from_match(match_body: &str) -> SchemaClassifier {
2158 let yaml = format!("schemas:\n - name: t\n specificity: 70\n match:\n{match_body}");
2159 let sigs = parse_schema_signatures(&yaml).expect("parse");
2160 SchemaClassifier::new(sigs)
2161 }
2162
2163 fn matches_t(match_body: &str, event: &serde_json::Value) -> bool {
2164 classifier_from_match(match_body)
2165 .classify(&JsonEvent::borrow(event))
2166 .is_some()
2167 }
2168
2169 #[test]
2170 fn numeric_comparisons() {
2171 let body = " - gte: { field: EventID, value: 4000 }\n";
2172 assert!(matches_t(body, &json!({"EventID": 4688})));
2173 assert!(matches_t(body, &json!({"EventID": 4000})));
2174 assert!(!matches_t(body, &json!({"EventID": 1})));
2175 assert!(matches_t(body, &json!({"EventID": "4688"})));
2177 assert!(!matches_t(body, &json!({"EventID": "not-a-number"})));
2179 assert!(matches_t(
2181 " - lt: { field: score, value: 10 }\n",
2182 &json!({"score": 9.5})
2183 ));
2184 assert!(matches_t(
2185 " - gt: { field: score, value: 10 }\n",
2186 &json!({"score": 10.1})
2187 ));
2188 }
2189
2190 #[test]
2191 fn in_set_membership_is_case_insensitive() {
2192 let body = " - in: { field: event_type, values: [alert, alarm] }\n";
2193 assert!(matches_t(body, &json!({"event_type": "ALERT"})));
2194 assert!(matches_t(body, &json!({"event_type": "alarm"})));
2195 assert!(!matches_t(body, &json!({"event_type": "info"})));
2196 assert!(!matches_t(body, &json!({})));
2197 }
2198
2199 #[test]
2200 fn field_equals_field_compares_two_fields() {
2201 let body = " - field_equals_field: { left: a, right: b }\n";
2202 assert!(matches_t(body, &json!({"a": "X", "b": "x"})));
2203 assert!(!matches_t(body, &json!({"a": "X", "b": "y"})));
2204 assert!(!matches_t(body, &json!({"a": "X"})));
2206 }
2207
2208 #[test]
2209 fn recursive_not_any_all_groups() {
2210 let any_body = " - any:\n - field_present: winlog.channel\n - equals: { field: host.os.type, value: windows }\n";
2212 assert!(matches_t(
2213 any_body,
2214 &json!({"winlog": {"channel": "Security"}})
2215 ));
2216 assert!(matches_t(
2217 any_body,
2218 &json!({"host": {"os": {"type": "windows"}}})
2219 ));
2220 assert!(!matches_t(any_body, &json!({"unrelated": 1})));
2221
2222 let not_body = " - not: { field_present: ecs.version }\n";
2224 assert!(matches_t(not_body, &json!({"CommandLine": "whoami"})));
2225 assert!(!matches_t(not_body, &json!({"ecs.version": "8.0.0"})));
2226
2227 let all_body = " - all:\n - field_present: a\n - field_present: b\n";
2229 assert!(matches_t(all_body, &json!({"a": 1, "b": 2})));
2230 assert!(!matches_t(all_body, &json!({"a": 1})));
2231 }
2232
2233 #[test]
2234 fn empty_group_is_rejected() {
2235 let yaml = "schemas:\n - name: t\n match:\n - any: []\n";
2236 let err = parse_schema_signatures(yaml).unwrap_err();
2237 assert!(
2238 matches!(&err, SchemaError::Parse(m) if m.contains("'any' needs at least one")),
2239 "got: {err}"
2240 );
2241 }
2242
2243 #[test]
2244 fn predicate_with_two_conditions_is_rejected() {
2245 let yaml = "schemas:\n - name: t\n match:\n - field_present: a\n field_absent: b\n";
2246 let err = parse_schema_signatures(yaml).unwrap_err();
2247 assert!(
2248 matches!(&err, SchemaError::Parse(m) if m.contains("multiple conditions")),
2249 "got: {err}"
2250 );
2251 }
2252
2253 #[test]
2254 fn explain_reports_matched_signature() {
2255 let cls = SchemaClassifier::builtin();
2256 let v = json!({"ecs.version": "8.0.0"});
2257 let ex = cls.explain(&JsonEvent::borrow(&v));
2258 assert_eq!(ex.matched.as_deref(), Some("ecs"));
2259 let sig = ex.signature.expect("signature");
2260 assert!(sig.predicates_matched);
2261 assert!(sig.predicates.iter().all(|p| p.matched));
2262 }
2263
2264 #[test]
2265 fn explain_reports_near_miss_for_unknown() {
2266 let sigs = builtin_signatures()
2268 .into_iter()
2269 .filter(|s| s.name != "generic_json")
2270 .collect();
2271 let cls = SchemaClassifier::new(sigs);
2272 let v = json!({"EventID": 1, "Image": "C:/cmd.exe"});
2274 let ex = cls.explain(&JsonEvent::borrow(&v));
2275 assert_eq!(ex.matched, None);
2276 let sig = ex.signature.expect("near-miss");
2277 assert_eq!(sig.name, "sysmon");
2278 assert!(!sig.predicates_matched);
2279 assert!(sig.predicates.iter().any(|p| !p.matched));
2280 }
2281
2282 #[test]
2283 fn validate_flags_unknown_binding_and_shadow() {
2284 let yaml = r#"
2285schemas:
2286 - name: shadowed
2287 specificity: 40
2288 match:
2289 - field_present: ecs.version
2290 - field_present: extra.marker
2291routing:
2292 bindings:
2293 - schema: ecs
2294 pipelines: [ecs_windows]
2295 - schema: nonexistent
2296 pipelines: [x]
2297"#;
2298 let (sigs, routing) = parse_schema_config(yaml).unwrap();
2299 let findings = validate_schema_config(&sigs, routing.as_ref());
2300 assert!(
2301 findings
2302 .iter()
2303 .any(|f| f.contains("unknown schema 'nonexistent'")),
2304 "findings: {findings:?}"
2305 );
2306 assert!(
2309 findings
2310 .iter()
2311 .any(|f| f.contains("'shadowed'") && f.contains("unreachable")),
2312 "findings: {findings:?}"
2313 );
2314 }
2315
2316 #[test]
2317 fn observer_counts_per_schema_and_unknown() {
2318 let observer = SchemaObserver::builtin();
2319 observer.observe(&JsonEvent::borrow(&json!({"ecs.version": "8.0.0"})));
2320 observer.observe(&JsonEvent::borrow(&json!({"ecs.version": "8.1.0"})));
2321 observer.observe(&JsonEvent::borrow(
2322 &json!({"class_uid": 1001, "metadata": {"version": "1.1.0"}}),
2323 ));
2324 observer.observe(&JsonEvent::borrow(&json!({})));
2325
2326 let snap = observer.snapshot();
2327 assert_eq!(snap.events_observed, 4);
2328 assert_eq!(snap.classified, 3);
2329 assert_eq!(snap.unknown, 1);
2330 assert_eq!(snap.by_schema[0].schema, "ecs");
2332 assert_eq!(snap.by_schema[0].count, 2);
2333 let ocsf = snap.by_schema.iter().find(|e| e.schema == "ocsf").unwrap();
2334 assert_eq!(ocsf.count, 1);
2335 }
2336
2337 #[test]
2338 fn routing_plan_dedups_pipeline_sets() {
2339 let config = RoutingConfig {
2340 on_unknown: OnUnknown::Warn,
2341 default_pipelines: vec![],
2342 aliases: HashMap::new(),
2343 bindings: vec![
2344 SchemaBinding {
2345 schema: "ecs".to_string(),
2346 pipelines: vec!["ecs_windows".to_string()],
2347 logsource: None,
2348 },
2349 SchemaBinding {
2350 schema: "winlogbeat".to_string(),
2351 pipelines: vec!["ecs_windows".to_string()],
2352 logsource: None,
2353 },
2354 SchemaBinding {
2355 schema: "sysmon".to_string(),
2356 pipelines: vec!["sysmon".to_string()],
2357 logsource: None,
2358 },
2359 ],
2360 };
2361 let plan = RoutingPlan::from_config(&config);
2362 assert_eq!(plan.pipeline_sets().len(), 3);
2364 let ecs = plan.decide(Some("ecs"));
2366 let win = plan.decide(Some("winlogbeat"));
2367 assert_eq!(ecs, win);
2368 assert!(matches!(
2369 ecs,
2370 RouteDecision::Evaluate { unknown: false, .. }
2371 ));
2372 assert_ne!(plan.decide(Some("sysmon")), ecs);
2374 }
2375
2376 #[test]
2377 fn routing_decides_bound_unbound_and_unknown() {
2378 let config = RoutingConfig {
2379 on_unknown: OnUnknown::Warn,
2380 default_pipelines: vec![],
2381 aliases: HashMap::new(),
2382 bindings: vec![SchemaBinding {
2383 schema: "ecs".to_string(),
2384 pipelines: vec!["ecs_windows".to_string()],
2385 logsource: None,
2386 }],
2387 };
2388 let plan = RoutingPlan::from_config(&config);
2389 assert!(matches!(
2391 plan.decide(Some("ecs")),
2392 RouteDecision::Evaluate { unknown: false, .. }
2393 ));
2394 assert_eq!(
2396 plan.decide(Some("cef")),
2397 RouteDecision::Evaluate {
2398 set: 0,
2399 unknown: false
2400 }
2401 );
2402 assert_eq!(
2404 plan.decide(None),
2405 RouteDecision::Evaluate {
2406 set: 0,
2407 unknown: true
2408 }
2409 );
2410 }
2411
2412 #[test]
2413 fn routing_on_unknown_policies() {
2414 let base = |policy| RoutingConfig {
2415 on_unknown: policy,
2416 default_pipelines: vec![],
2417 aliases: HashMap::new(),
2418 bindings: vec![],
2419 };
2420 assert_eq!(
2421 RoutingPlan::from_config(&base(OnUnknown::Drop)).decide(None),
2422 RouteDecision::Drop
2423 );
2424 assert_eq!(
2425 RoutingPlan::from_config(&base(OnUnknown::Error)).decide(None),
2426 RouteDecision::Error
2427 );
2428 assert_eq!(
2429 RoutingPlan::from_config(&base(OnUnknown::Passthrough)).decide(None),
2430 RouteDecision::Evaluate {
2431 set: 0,
2432 unknown: true
2433 }
2434 );
2435 }
2436
2437 #[test]
2438 fn parses_routing_section_from_yaml() {
2439 let yaml = r#"
2440schemas:
2441 - name: my_vendor
2442 match:
2443 - field_present: vendor.id
2444routing:
2445 on_unknown: drop
2446 default_pipelines: [base]
2447 bindings:
2448 - schema: ecs
2449 pipelines: [ecs_windows]
2450 - schema: my_vendor
2451 pipelines: [vendor_map, base]
2452"#;
2453 let (sigs, routing) = parse_schema_config(yaml).expect("parse");
2454 assert_eq!(sigs.len(), 1);
2455 let routing = routing.expect("routing present");
2456 assert_eq!(routing.on_unknown, OnUnknown::Drop);
2457 assert_eq!(routing.default_pipelines, vec!["base".to_string()]);
2458 assert_eq!(routing.bindings.len(), 2);
2459 let plan = RoutingPlan::from_config(&routing);
2460 assert_eq!(plan.pipeline_sets().len(), 3);
2462 assert_eq!(plan.decide(None), RouteDecision::Drop);
2463 }
2464
2465 #[test]
2466 fn schema_logsource_builtin_defaults_and_overrides() {
2467 let plan = RoutingPlan::from_config(&RoutingConfig::default());
2469 let sysmon = plan.schema_logsource("sysmon").expect("sysmon default");
2470 assert_eq!(sysmon.product.as_deref(), Some("windows"));
2471 assert_eq!(sysmon.service.as_deref(), Some("sysmon"));
2472 assert_eq!(
2473 plan.schema_logsource("windows_eventlog")
2474 .and_then(|l| l.product.as_deref()),
2475 Some("windows")
2476 );
2477 assert!(plan.schema_logsource("ecs").is_none());
2479 assert!(plan.schema_logsource("cef").is_none());
2480
2481 let yaml = r#"
2483schemas:
2484 - name: ecs_windows
2485 match:
2486 - field_present: ecs.version
2487 - field_present: winlog.channel
2488routing:
2489 bindings:
2490 - schema: ecs_windows
2491 pipelines: [ecs_windows]
2492 logsource:
2493 product: windows
2494 - schema: sysmon
2495 pipelines: [sysmon]
2496 logsource:
2497 product: windows
2498 service: sysmon
2499 custom:
2500 tenant: acme
2501"#;
2502 let (_sigs, routing) = parse_schema_config(yaml).expect("parse");
2503 let plan = RoutingPlan::from_config(&routing.expect("routing"));
2504 assert_eq!(
2505 plan.schema_logsource("ecs_windows")
2506 .and_then(|l| l.product.as_deref()),
2507 Some("windows")
2508 );
2509 let sysmon = plan.schema_logsource("sysmon").expect("sysmon override");
2510 assert_eq!(
2511 sysmon.custom.get("tenant").map(String::as_str),
2512 Some("acme")
2513 );
2514 }
2515
2516 #[test]
2517 fn observer_reset_preserves_lifetime_counters() {
2518 let observer = SchemaObserver::builtin();
2519 observer.observe(&JsonEvent::borrow(&json!({"ecs.version": "8.0.0"})));
2520 observer.observe(&JsonEvent::borrow(&json!({})));
2521 let (classified, unknown) = observer.reset();
2522 assert_eq!(classified, 1);
2523 assert_eq!(unknown, 1);
2524
2525 let snap = observer.snapshot();
2526 assert_eq!(snap.classified, 0);
2527 assert_eq!(snap.unknown, 0);
2528 assert_eq!(snap.events_observed, 0);
2529 assert_eq!(snap.lifetime_classified, 1);
2531 assert_eq!(snap.lifetime_unknown, 1);
2532 }
2533
2534 #[test]
2535 fn classify_with_ambiguity_flags_equal_specificity_ties() {
2536 let sigs = vec![
2538 SchemaSignature {
2539 name: "alpha".to_string(),
2540 specificity: 70,
2541 predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
2542 },
2543 SchemaSignature {
2544 name: "beta".to_string(),
2545 specificity: 70,
2546 predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
2547 },
2548 ];
2549 let cls = SchemaClassifier::new(sigs);
2550 let (m, ambiguous) = cls.classify_with_ambiguity(&JsonEvent::borrow(&json!({"a": 1})));
2551 assert!(m.is_some());
2552 assert!(
2553 ambiguous,
2554 "equal-specificity different-name match is ambiguous"
2555 );
2556 let cls = SchemaClassifier::builtin();
2558 let (_, ambiguous) =
2559 cls.classify_with_ambiguity(&JsonEvent::borrow(&json!({"ecs.version": "8.0.0"})));
2560 assert!(!ambiguous);
2561 }
2562
2563 #[test]
2564 fn observer_records_ambiguity_and_unknown_shapes() {
2565 let sigs = vec![
2566 SchemaSignature {
2567 name: "alpha".to_string(),
2568 specificity: 70,
2569 predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
2570 },
2571 SchemaSignature {
2572 name: "beta".to_string(),
2573 specificity: 70,
2574 predicates: vec![SchemaPredicate::FieldPresent("a".to_string())],
2575 },
2576 ];
2577 let observer = SchemaObserver::new(SchemaClassifier::new(sigs));
2578 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();
2583 assert_eq!(snap.ambiguous, 1);
2584 assert_eq!(snap.unknown, 2);
2585 assert_eq!(snap.unknown_shapes.len(), 1);
2587 assert_eq!(snap.unknown_shapes[0].count, 2);
2588 assert_eq!(snap.unknown_shapes[0].keys, vec!["shape", "weird"]);
2589 }
2590
2591 #[test]
2596 fn cloud_signatures_have_higher_specificity_than_generic_json() {
2597 let cls = SchemaClassifier::builtin();
2599 let names = cls.schema_names();
2600 assert_eq!(names.first(), Some(&"ecs_linux"));
2601 assert_eq!(names.last(), Some(&"generic_json"));
2603 let cloud_order: Vec<&str> = vec![
2605 "gcp_audit",
2606 "aws_cloudtrail",
2607 "azure_signinlogs",
2608 "github_audit",
2609 "k8s_audit",
2610 "aws_vpcflow",
2611 "docker_events",
2612 "osquery_result",
2613 ];
2614 let cloud_indices: Vec<usize> = cloud_order
2615 .iter()
2616 .filter_map(|&n| names.iter().position(|&x| x == n))
2617 .collect();
2618 if cloud_indices.len() == cloud_order.len() {
2619 let max_cloud = cloud_indices.iter().max().copied().unwrap_or(0);
2621 let generic_idx = names
2622 .last()
2623 .copied()
2624 .map(|n| names.iter().position(|&x| x == n).unwrap_or(0))
2625 .unwrap_or(0);
2626 assert!(
2627 generic_idx > max_cloud,
2628 "generic_json ({}) should be after all cloud schemas, but cloud max = {max_cloud}",
2629 names[generic_idx]
2630 );
2631 }
2632 }
2633
2634 #[test]
2635 fn cloud_signatures_dont_shadow_each_other() {
2636 let gcp =
2642 json!({"protoPayload": {"@type": "type.googleapis.com/google.cloud.audit.AuditLog"}});
2643 assert_eq!(
2644 SchemaClassifier::builtin()
2645 .classify(&JsonEvent::borrow(&gcp))
2646 .as_ref()
2647 .map(|m| m.name.as_str()),
2648 Some("gcp_audit")
2649 );
2650
2651 let cloudtrail = json!({"eventVersion": "1.05", "eventSource": "s3.amazonaws.com", "eventID": "abc", "userIdentity": {"type": "IAMUser"}});
2653 assert_eq!(
2654 SchemaClassifier::builtin()
2655 .classify(&JsonEvent::borrow(&cloudtrail))
2656 .as_ref()
2657 .map(|m| m.name.as_str()),
2658 Some("aws_cloudtrail")
2659 );
2660
2661 let github = json!({"action": "repo.create", "actor": "admin", "org": {"id": 123}, "created_at": "2024-01-01T00:00:00Z"});
2663 assert_eq!(
2664 SchemaClassifier::builtin()
2665 .classify(&JsonEvent::borrow(&github))
2666 .as_ref()
2667 .map(|m| m.name.as_str()),
2668 Some("github_audit")
2669 );
2670
2671 let k8s = json!({"kind": "Event", "apiVersion": "audit.k8s.io/v1", "auditID": "abc"});
2673 assert_eq!(
2674 SchemaClassifier::builtin()
2675 .classify(&JsonEvent::borrow(&k8s))
2676 .as_ref()
2677 .map(|m| m.name.as_str()),
2678 Some("k8s_audit")
2679 );
2680
2681 let azure_act = json!({"category": "Administrative", "id": "/SUBSCRIPTIONS/abc", "operationName": {"value": "test"}});
2683 assert_eq!(
2684 SchemaClassifier::builtin()
2685 .classify(&JsonEvent::borrow(&azure_act))
2686 .as_ref()
2687 .map(|m| m.name.as_str()),
2688 Some("azure_activitylogs")
2689 );
2690
2691 let m365 = json!({"RecordType": 15, "Workload": "AzureActiveDirectory", "Operation": "UserLoggedIn", "CreationTime": "2024-01-01T00:00:00Z"});
2693 assert_eq!(
2694 SchemaClassifier::builtin()
2695 .classify(&JsonEvent::borrow(&m365))
2696 .as_ref()
2697 .map(|m| m.name.as_str()),
2698 Some("m365_audit")
2699 );
2700 }
2701
2702 #[test]
2703 fn off_taxonomy_signatures_use_custom_logsource() {
2704 let map = builtin_schema_logsource();
2707
2708 for schema in [
2709 "k8s_audit",
2710 "docker_events",
2711 "osquery_result",
2712 "aws_vpcflow",
2713 ] {
2714 let ls = map
2715 .get(schema)
2716 .unwrap_or_else(|| panic!("logsource mapping for {schema}"));
2717 if schema != "aws_vpcflow" {
2719 assert!(
2720 ls.product.is_none(),
2721 "{schema} must not have a product (off-taxonomy uses custom only)"
2722 );
2723 }
2724 assert!(
2726 !ls.custom.is_empty(),
2727 "{schema} must have custom dimensions for pruning"
2728 );
2729 }
2730
2731 let vpc = map.get("aws_vpcflow").expect("vpcflow mapping");
2733 assert_eq!(vpc.product.as_deref(), Some("aws"));
2734 assert_eq!(vpc.custom.get("source"), Some(&"vpcflow".to_string()));
2735 }
2736
2737 #[test]
2738 fn builtin_schema_names_match_signatures() {
2739 use std::collections::{HashMap, HashSet};
2742
2743 let mut spec_by_name: HashMap<String, u32> = HashMap::new();
2746 for sig in builtin_signatures() {
2747 let entry = spec_by_name
2748 .entry(sig.name.clone())
2749 .or_insert(sig.specificity);
2750 *entry = (*entry).max(sig.specificity);
2751 }
2752
2753 let names = builtin_schema_names();
2754
2755 let listed: HashSet<String> = names.iter().map(|s| s.to_string()).collect();
2757 let produced: HashSet<String> = spec_by_name.keys().cloned().collect();
2758 assert_eq!(
2759 listed, produced,
2760 "builtin_schema_names() is out of sync with builtin_signatures()"
2761 );
2762
2763 let mut prev = u32::MAX;
2765 for name in &names {
2766 let spec = spec_by_name[*name];
2767 assert!(
2768 spec <= prev,
2769 "builtin_schema_names() not ordered by non-increasing specificity at '{name}' ({spec} > {prev})"
2770 );
2771 prev = spec;
2772 }
2773 }
2774
2775 #[test]
2776 fn okta_and_onelogin_signatures() {
2777 let okta = json!({
2779 "eventType": "user.lifecycle.activate.pre_auth",
2780 "actor": {"id": "abc"},
2781 "published": "2024-01-01T00:00:00Z",
2782 "outcome": {"result": "SUCCESS"}
2783 });
2784 let classifier = SchemaClassifier::builtin();
2785 assert_eq!(
2786 classifier
2787 .classify(&JsonEvent::borrow(&okta))
2788 .as_ref()
2789 .map(|m| m.name.as_str()),
2790 Some("okta_system_log")
2791 );
2792
2793 let onelogin = json!({
2795 "event_type_id": 123,
2796 "account_id": 456,
2797 "user_id": 789,
2798 "created_at": "2024-01-01T00:00:00Z"
2799 });
2800 assert_eq!(
2801 classifier
2802 .classify(&JsonEvent::borrow(&onelogin))
2803 .as_ref()
2804 .map(|m| m.name.as_str()),
2805 Some("onelogin_events")
2806 );
2807 }
2808
2809 #[test]
2810 fn m365_unified_audit_log_maps_to_audit_service() {
2811 let exchange = json!({
2815 "CreationTime": "2024-01-01T00:00:00Z",
2816 "RecordType": 1,
2817 "Workload": "Exchange",
2818 "Operation": "New-RemoteDomain"
2819 });
2820 let classifier = SchemaClassifier::builtin();
2821
2822 let m = classifier
2823 .classify(&JsonEvent::borrow(&exchange))
2824 .expect("matched");
2825 assert_eq!(m.name, "m365_audit");
2826
2827 let map = builtin_schema_logsource();
2828 let ls = map.get("m365_audit").expect("m365_audit mapping");
2829 assert_eq!(ls.product.as_deref(), Some("m365"));
2830 assert_eq!(ls.service.as_deref(), Some("audit"));
2831 }
2832
2833 #[test]
2834 fn docker_and_osquery_signatures() {
2835 let docker = json!({"Type": "container", "Action": "start", "Actor": {"ID": "abc"}});
2837 let classifier = SchemaClassifier::builtin();
2838 assert_eq!(
2839 classifier
2840 .classify(&JsonEvent::borrow(&docker))
2841 .as_ref()
2842 .map(|m| m.name.as_str()),
2843 Some("docker_events")
2844 );
2845
2846 let osquery = json!({
2848 "name": "users",
2849 "action": "added",
2850 "columns": {"uid": "1000", "username": "admin"},
2851 "hostIdentifier": "workstation-01"
2852 });
2853 assert_eq!(
2854 classifier
2855 .classify(&JsonEvent::borrow(&osquery))
2856 .as_ref()
2857 .map(|m| m.name.as_str()),
2858 Some("osquery_result")
2859 );
2860 }
2861
2862 #[test]
2863 fn aws_vpcflow_classification() {
2864 let vpc = json!({
2866 "version": 2,
2867 "srcaddr": "10.0.1.100",
2868 "dstaddr": "10.0.2.50",
2869 "srcport": 45678,
2870 "dstport": 443,
2871 "protocol": 6,
2872 "action": "ACCEPT",
2873 "log_status": "OK"
2874 });
2875 let classifier = SchemaClassifier::builtin();
2876 assert_eq!(
2877 classifier
2878 .classify(&JsonEvent::borrow(&vpc))
2879 .as_ref()
2880 .map(|m| m.name.as_str()),
2881 Some("aws_vpcflow")
2882 );
2883
2884 let all = classifier.classify_all(&JsonEvent::borrow(&vpc));
2886 assert!(
2887 !all.iter().any(|s| s == "aws_cloudtrail"),
2888 "VPC flow event should not match CloudTrail"
2889 );
2890
2891 let map = builtin_schema_logsource();
2893 let vpc_ls = map.get("aws_vpcflow").expect("vpcflow mapping");
2894 assert_eq!(vpc_ls.product.as_deref(), Some("aws"));
2895 assert_eq!(vpc_ls.custom.get("source"), Some(&"vpcflow".to_string()));
2896 }
2897}