1mod introspect;
16#[cfg(test)]
17mod tests;
18mod types;
19
20pub use introspect::{CorrelationInfo, CorrelationStateSnapshot, GroupKeyPart, GroupStateInfo};
21pub use types::*;
22
23use std::collections::HashMap;
24
25use chrono::{DateTime, TimeZone, Utc};
26
27use rsigma_parser::{CorrelationRule, CorrelationType, SigmaCollection, SigmaRule, WindowMode};
28
29use crate::correlation::{
30 CompiledCorrelation, EventBuffer, EventRefBuffer, GroupKey, WindowDecision, WindowState,
31 apply_window_open, compile_correlation,
32};
33use crate::engine::Engine;
34use crate::error::{EvalError, Result};
35use crate::event::{Event, EventValue};
36use crate::pipeline::{Pipeline, apply_pipelines, apply_pipelines_to_correlation};
37use crate::result::{CorrelationBody, EvaluationResult, ResultBody, RuleHeader};
38use crate::rule_metadata::{RuleBundleMetadata, RuleMetadataLookup};
39
40const SNAPSHOT_VERSION: u32 = 1;
46
47const MAX_CHAIN_DEPTH: usize = 10;
53
54pub struct CorrelationEngine {
59 engine: Engine,
61 correlations: Vec<CompiledCorrelation>,
63 rule_index: HashMap<String, Vec<usize>>,
66 rule_ids: Vec<(Option<String>, Option<String>)>,
69 state: HashMap<(usize, GroupKey), WindowState>,
71 last_alert: HashMap<(usize, GroupKey), i64>,
73 event_buffers: HashMap<(usize, GroupKey), EventBuffer>,
75 event_ref_buffers: HashMap<(usize, GroupKey), EventRefBuffer>,
77 correlation_only_rules: std::collections::HashSet<String>,
81 config: CorrelationConfig,
83 pipelines: Vec<Pipeline>,
85}
86
87impl CorrelationEngine {
88 pub fn new(config: CorrelationConfig) -> Self {
90 CorrelationEngine {
91 engine: Engine::new(),
92 correlations: Vec::new(),
93 rule_index: HashMap::new(),
94 rule_ids: Vec::new(),
95 state: HashMap::new(),
96 last_alert: HashMap::new(),
97 event_buffers: HashMap::new(),
98 event_ref_buffers: HashMap::new(),
99 correlation_only_rules: std::collections::HashSet::new(),
100 config,
101 pipelines: Vec::new(),
102 }
103 }
104
105 pub fn add_pipeline(&mut self, pipeline: Pipeline) {
109 self.pipelines.push(pipeline);
110 self.pipelines.sort_by_key(|p| p.priority);
111 }
112
113 pub fn set_include_event(&mut self, include: bool) {
115 self.engine.set_include_event(include);
116 }
117
118 pub fn set_match_detail(&mut self, level: crate::result::MatchDetailLevel) {
121 self.engine.set_match_detail(level);
122 }
123
124 pub fn set_bloom_prefilter(&mut self, enabled: bool) {
128 self.engine.set_bloom_prefilter(enabled);
129 }
130
131 pub fn set_logsource_extractor(
135 &mut self,
136 extractor: Option<crate::logsource::LogSourceExtractor>,
137 ) {
138 self.engine.set_logsource_extractor(extractor);
139 }
140
141 pub fn logsource_pruned_total(&self) -> u64 {
143 self.engine.logsource_pruned_total()
144 }
145
146 pub fn logsource_absent_total(&self) -> u64 {
148 self.engine.logsource_absent_total()
149 }
150
151 pub fn set_bloom_max_bytes(&mut self, max_bytes: usize) {
154 self.engine.set_bloom_max_bytes(max_bytes);
155 }
156
157 #[cfg(feature = "daachorse-index")]
161 pub fn set_cross_rule_ac(&mut self, enabled: bool) {
162 self.engine.set_cross_rule_ac(enabled);
163 }
164
165 pub fn set_correlation_event_mode(&mut self, mode: CorrelationEventMode) {
171 self.config.correlation_event_mode = mode;
172 }
173
174 pub fn set_max_correlation_events(&mut self, max: usize) {
177 self.config.max_correlation_events = max;
178 }
179
180 pub fn add_rule(&mut self, rule: &SigmaRule) -> Result<()> {
186 if self.pipelines.is_empty() {
187 self.apply_custom_attributes(&rule.custom_attributes);
188 self.rule_ids.push((rule.id.clone(), rule.name.clone()));
189 self.engine.add_rule(rule)?;
190 } else {
191 let mut transformed = rule.clone();
192 apply_pipelines(&self.pipelines, &mut transformed)?;
193 self.apply_custom_attributes(&transformed.custom_attributes);
194 self.rule_ids
195 .push((transformed.id.clone(), transformed.name.clone()));
196 let compiled = crate::compiler::compile_rule(&transformed)?;
198 self.engine.add_compiled_rule(compiled);
199 }
200 Ok(())
201 }
202
203 fn apply_custom_attributes(
217 &mut self,
218 attrs: &std::collections::HashMap<String, yaml_serde::Value>,
219 ) {
220 if let Some(field) = attrs.get("rsigma.timestamp_field").and_then(|v| v.as_str())
222 && !self.config.timestamp_fields.iter().any(|f| f == field)
223 {
224 self.config.timestamp_fields.insert(0, field.to_string());
225 }
226
227 if let Some(val) = attrs.get("rsigma.suppress").and_then(|v| v.as_str())
229 && self.config.suppress.is_none()
230 && let Ok(ts) = rsigma_parser::Timespan::parse(val)
231 {
232 self.config.suppress = Some(ts.seconds);
233 }
234
235 if let Some(val) = attrs.get("rsigma.action").and_then(|v| v.as_str())
237 && self.config.action_on_match == CorrelationAction::Alert
238 && let Ok(a) = val.parse::<CorrelationAction>()
239 {
240 self.config.action_on_match = a;
241 }
242 }
243
244 pub fn add_correlation(&mut self, corr: &CorrelationRule) -> Result<()> {
246 let owned;
247 let effective = if self.pipelines.is_empty() {
248 corr
249 } else {
250 owned = {
251 let mut c = corr.clone();
252 apply_pipelines_to_correlation(&self.pipelines, &mut c)?;
253 c
254 };
255 &owned
256 };
257
258 self.apply_custom_attributes(&effective.custom_attributes);
261
262 let compiled = compile_correlation(effective)?;
263 let idx = self.correlations.len();
264
265 for rule_ref in &compiled.rule_refs {
267 self.rule_index
268 .entry(rule_ref.clone())
269 .or_default()
270 .push(idx);
271 }
272
273 if !compiled.generate {
275 for rule_ref in &compiled.rule_refs {
276 self.correlation_only_rules.insert(rule_ref.clone());
277 }
278 }
279
280 self.correlations.push(compiled);
281 Ok(())
282 }
283
284 pub fn add_collection(&mut self, collection: &SigmaCollection) -> Result<()> {
293 let mut compiled_batch = Vec::with_capacity(collection.rules.len());
294 if self.pipelines.is_empty() {
295 for rule in &collection.rules {
296 self.apply_custom_attributes(&rule.custom_attributes);
297 self.rule_ids.push((rule.id.clone(), rule.name.clone()));
298 compiled_batch.push(crate::compiler::compile_rule(rule)?);
299 }
300 } else {
301 for rule in &collection.rules {
302 let mut transformed = rule.clone();
303 apply_pipelines(&self.pipelines, &mut transformed)?;
304 self.apply_custom_attributes(&transformed.custom_attributes);
305 self.rule_ids
306 .push((transformed.id.clone(), transformed.name.clone()));
307 compiled_batch.push(crate::compiler::compile_rule(&transformed)?);
309 }
310 }
311 self.engine.extend_compiled_rules(compiled_batch);
312 for filter in &collection.filters {
314 self.engine.apply_filter(filter)?;
315 }
316 for corr in &collection.correlations {
317 self.add_correlation(corr)?;
318 }
319 self.validate_rule_refs()?;
320 self.detect_correlation_cycles()?;
321 Ok(())
322 }
323
324 fn validate_rule_refs(&self) -> Result<()> {
327 let mut known: std::collections::HashSet<&str> = std::collections::HashSet::new();
328
329 for (id, name) in &self.rule_ids {
330 if let Some(id) = id {
331 known.insert(id.as_str());
332 }
333 if let Some(name) = name {
334 known.insert(name.as_str());
335 }
336 }
337 for corr in &self.correlations {
338 if let Some(ref id) = corr.id {
339 known.insert(id.as_str());
340 }
341 if let Some(ref name) = corr.name {
342 known.insert(name.as_str());
343 }
344 }
345
346 for corr in &self.correlations {
347 for rule_ref in &corr.rule_refs {
348 if !known.contains(rule_ref.as_str()) {
349 return Err(EvalError::UnknownRuleRef(rule_ref.clone()));
350 }
351 }
352 }
353 Ok(())
354 }
355
356 fn detect_correlation_cycles(&self) -> Result<()> {
364 let mut corr_identifiers: HashMap<&str, usize> = HashMap::new();
366 for (idx, corr) in self.correlations.iter().enumerate() {
367 if let Some(ref id) = corr.id {
368 corr_identifiers.insert(id.as_str(), idx);
369 }
370 if let Some(ref name) = corr.name {
371 corr_identifiers.insert(name.as_str(), idx);
372 }
373 }
374
375 let mut adj: Vec<Vec<usize>> = vec![Vec::new(); self.correlations.len()];
377 for (idx, corr) in self.correlations.iter().enumerate() {
378 for rule_ref in &corr.rule_refs {
379 if let Some(&target_idx) = corr_identifiers.get(rule_ref.as_str()) {
380 adj[idx].push(target_idx);
381 }
382 }
383 }
384
385 let mut state = vec![0u8; self.correlations.len()]; let mut path: Vec<usize> = Vec::new();
388
389 for start in 0..self.correlations.len() {
390 if state[start] == 0
391 && let Some(cycle) = Self::dfs_find_cycle(start, &adj, &mut state, &mut path)
392 {
393 let names: Vec<String> = cycle
394 .iter()
395 .map(|&i| {
396 self.correlations[i]
397 .id
398 .as_deref()
399 .or(self.correlations[i].name.as_deref())
400 .unwrap_or(&self.correlations[i].title)
401 .to_string()
402 })
403 .collect();
404 return Err(crate::error::EvalError::CorrelationCycle(
405 names.join(" -> "),
406 ));
407 }
408 }
409 Ok(())
410 }
411
412 fn dfs_find_cycle(
414 node: usize,
415 adj: &[Vec<usize>],
416 state: &mut [u8],
417 path: &mut Vec<usize>,
418 ) -> Option<Vec<usize>> {
419 state[node] = 1; path.push(node);
421
422 for &next in &adj[node] {
423 if state[next] == 1 {
424 if let Some(pos) = path.iter().position(|&n| n == next) {
426 let mut cycle = path[pos..].to_vec();
427 cycle.push(next); return Some(cycle);
429 }
430 }
431 if state[next] == 0
432 && let Some(cycle) = Self::dfs_find_cycle(next, adj, state, path)
433 {
434 return Some(cycle);
435 }
436 }
437
438 path.pop();
439 state[node] = 2; None
441 }
442
443 pub fn process_event(&mut self, event: &impl Event) -> ProcessResult {
449 let all_detections = self.engine.evaluate(event);
450 self.correlate_detections(event, all_detections)
451 }
452
453 pub fn correlate_detections(
462 &mut self,
463 event: &impl Event,
464 all_detections: Vec<EvaluationResult>,
465 ) -> ProcessResult {
466 let ts = match self.extract_event_timestamp(event) {
467 Some(ts) => ts,
468 None => match self.config.timestamp_fallback {
469 TimestampFallback::WallClock => Utc::now().timestamp(),
470 TimestampFallback::Skip => {
471 return self.filter_detections(all_detections);
473 }
474 },
475 };
476 self.process_with_detections(event, all_detections, ts)
477 }
478
479 pub fn process_event_at(&mut self, event: &impl Event, timestamp_secs: i64) -> ProcessResult {
484 let all_detections = self.engine.evaluate(event);
485 self.process_with_detections(event, all_detections, timestamp_secs)
486 }
487
488 pub fn process_with_detections(
494 &mut self,
495 event: &impl Event,
496 all_detections: Vec<EvaluationResult>,
497 timestamp_secs: i64,
498 ) -> ProcessResult {
499 let timestamp_secs = timestamp_secs.clamp(0, i64::MAX / 2);
500
501 if self.state.len() >= self.config.max_state_entries {
503 self.evict_all(timestamp_secs);
504 }
505
506 let mut correlations: Vec<EvaluationResult> = Vec::new();
508 self.feed_detections(event, &all_detections, timestamp_secs, &mut correlations);
509
510 let mut chained = Vec::new();
513 self.chain_correlations(&correlations, timestamp_secs, &mut chained);
514 correlations.extend(chained);
515
516 let mut out = self.filter_detections(all_detections);
518 out.extend(correlations);
519 out
520 }
521
522 pub fn evaluate(&self, event: &impl Event) -> Vec<EvaluationResult> {
529 self.engine.evaluate(event)
530 }
531
532 pub fn process_batch<E: Event + Sync>(&mut self, events: &[&E]) -> Vec<ProcessResult> {
540 let engine = &self.engine;
543 let ts_fields = &self.config.timestamp_fields;
544
545 let batch_results: Vec<(Vec<EvaluationResult>, Option<i64>)> = {
546 #[cfg(feature = "parallel")]
547 {
548 use rayon::prelude::*;
549 events
550 .par_iter()
551 .map(|e| {
552 let detections = engine.evaluate(e);
553 let ts = extract_event_ts(e, ts_fields);
554 (detections, ts)
555 })
556 .collect()
557 }
558 #[cfg(not(feature = "parallel"))]
559 {
560 events
561 .iter()
562 .map(|e| {
563 let detections = engine.evaluate(e);
564 let ts = extract_event_ts(e, ts_fields);
565 (detections, ts)
566 })
567 .collect()
568 }
569 };
570
571 let mut results = Vec::with_capacity(events.len());
573 for ((detections, ts_opt), event) in batch_results.into_iter().zip(events) {
574 match ts_opt {
575 Some(ts) => {
576 results.push(self.process_with_detections(event, detections, ts));
577 }
578 None => match self.config.timestamp_fallback {
579 TimestampFallback::WallClock => {
580 let ts = Utc::now().timestamp();
581 results.push(self.process_with_detections(event, detections, ts));
582 }
583 TimestampFallback::Skip => {
584 results.push(self.filter_detections(detections));
586 }
587 },
588 }
589 }
590 results
591 }
592
593 fn filter_detections(&self, all_detections: Vec<EvaluationResult>) -> Vec<EvaluationResult> {
598 if !self.config.emit_detections && !self.correlation_only_rules.is_empty() {
599 all_detections
600 .into_iter()
601 .filter(|m| {
602 let id_match = m
603 .header
604 .rule_id
605 .as_ref()
606 .is_some_and(|id| self.correlation_only_rules.contains(id));
607 !id_match
608 })
609 .collect()
610 } else {
611 all_detections
612 }
613 }
614
615 fn feed_detections(
617 &mut self,
618 event: &impl Event,
619 detections: &[EvaluationResult],
620 ts: i64,
621 out: &mut Vec<EvaluationResult>,
622 ) {
623 let mut work: Vec<(usize, Option<String>, Option<String>)> = Vec::new();
626
627 for det in detections {
628 let (rule_id, rule_name) = self.find_rule_identity(det);
631
632 let mut corr_indices = Vec::new();
634 if let Some(ref id) = rule_id
635 && let Some(indices) = self.rule_index.get(id)
636 {
637 corr_indices.extend(indices);
638 }
639 if let Some(ref name) = rule_name
640 && let Some(indices) = self.rule_index.get(name)
641 {
642 corr_indices.extend(indices);
643 }
644
645 corr_indices.sort_unstable();
646 corr_indices.dedup();
647
648 for &corr_idx in &corr_indices {
649 work.push((corr_idx, rule_id.clone(), rule_name.clone()));
650 }
651 }
652
653 for (corr_idx, rule_id, rule_name) in work {
654 self.update_correlation(corr_idx, event, ts, &rule_id, &rule_name, out);
655 }
656 }
657
658 fn find_rule_identity(&self, det: &EvaluationResult) -> (Option<String>, Option<String>) {
660 if let Some(ref match_id) = det.header.rule_id {
662 for (id, name) in &self.rule_ids {
663 if id.as_deref() == Some(match_id.as_str()) {
664 return (id.clone(), name.clone());
665 }
666 }
667 }
668 (det.header.rule_id.clone(), None)
670 }
671
672 fn resolve_event_mode(&self, corr_idx: usize) -> CorrelationEventMode {
674 let corr = &self.correlations[corr_idx];
675 corr.event_mode
676 .unwrap_or(self.config.correlation_event_mode)
677 }
678
679 fn resolve_max_events(&self, corr_idx: usize) -> usize {
681 let corr = &self.correlations[corr_idx];
682 corr.max_events
683 .unwrap_or(self.config.max_correlation_events)
684 }
685
686 fn resolve_max_group_entries(&self, corr_idx: usize) -> Option<usize> {
689 let corr = &self.correlations[corr_idx];
690 corr.max_group_entries.or(self.config.max_group_entries)
691 }
692
693 fn update_correlation(
695 &mut self,
696 corr_idx: usize,
697 event: &impl Event,
698 ts: i64,
699 rule_id: &Option<String>,
700 rule_name: &Option<String>,
701 out: &mut Vec<EvaluationResult>,
702 ) {
703 let corr = &self.correlations[corr_idx];
707 let corr_type = corr.correlation_type;
708 let timespan = corr.timespan_secs;
709 let window_mode = corr.window_mode;
710 let gap_secs = corr.gap_secs;
711 let level = corr.level;
712 let suppress_secs = corr.suppress_secs.or(self.config.suppress);
713 let action = corr.action.unwrap_or(self.config.action_on_match);
714 let event_mode = self.resolve_event_mode(corr_idx);
715 let max_events = self.resolve_max_events(corr_idx);
716 let max_group_entries = self.resolve_max_group_entries(corr_idx);
717
718 let mut ref_strs: Vec<&str> = Vec::new();
720 if let Some(id) = rule_id.as_deref() {
721 ref_strs.push(id);
722 }
723 if let Some(name) = rule_name.as_deref() {
724 ref_strs.push(name);
725 }
726 let rule_ref = ref_strs
727 .iter()
728 .copied()
729 .find(|identity| corr.rule_refs.iter().any(|rule_ref| rule_ref == identity))
730 .unwrap_or("");
731
732 let group_key = GroupKey::extract(event, &corr.group_by, &ref_strs);
734
735 let state_key = (corr_idx, group_key.clone());
737 let state = self
738 .state
739 .entry(state_key.clone())
740 .or_insert_with(|| WindowState::new_for(corr_type));
741
742 let cutoff = ts - timespan as i64;
748 let decision = apply_window_open(state, ts, timespan, window_mode, gap_secs);
749 if decision == WindowDecision::Discard {
750 return;
751 }
752 let reset = decision == WindowDecision::Reset;
753
754 match corr_type {
756 CorrelationType::EventCount => {
757 state.push_event_count(ts);
758 }
759 CorrelationType::ValueCount => {
760 if let Some(ref fields) = corr.condition.field
761 && let Some(key) = composite_value_count_key(event, fields)
762 {
763 state.push_value_count(ts, key);
764 }
765 }
766 CorrelationType::Temporal | CorrelationType::TemporalOrdered => {
767 state.push_temporal(ts, rule_ref);
768 }
769 CorrelationType::ValueSum
770 | CorrelationType::ValueAvg
771 | CorrelationType::ValuePercentile
772 | CorrelationType::ValueMedian => {
773 if let Some(ref fields) = corr.condition.field
774 && let Some(field_name) = fields.first()
775 && let Some(val) = event.get_field(field_name)
776 && let Some(n) = value_to_f64_ev(&val)
777 {
778 state.push_numeric(ts, n);
779 }
780 }
781 }
782
783 if let Some(cap) = max_group_entries {
787 state.truncate_oldest(cap, window_mode == WindowMode::Session);
788 }
789
790 match event_mode {
794 CorrelationEventMode::Full => {
795 let buf = self
796 .event_buffers
797 .entry(state_key.clone())
798 .or_insert_with(|| EventBuffer::new(max_events));
799 if window_mode == rsigma_parser::WindowMode::Sliding {
800 buf.evict(cutoff);
801 } else if reset {
802 buf.clear();
803 }
804 let json = event.to_json();
805 buf.push(ts, &json);
806 }
807 CorrelationEventMode::Refs => {
808 let buf = self
809 .event_ref_buffers
810 .entry(state_key.clone())
811 .or_insert_with(|| EventRefBuffer::new(max_events));
812 if window_mode == rsigma_parser::WindowMode::Sliding {
813 buf.evict(cutoff);
814 } else if reset {
815 buf.clear();
816 }
817 let json = event.to_json();
818 buf.push(ts, &json);
819 }
820 CorrelationEventMode::None => {}
821 }
822
823 let fired = state.check_condition(
825 &corr.condition,
826 corr_type,
827 &corr.rule_refs,
828 corr.extended_expr.as_ref(),
829 );
830
831 if let Some(agg_value) = fired {
832 let alert_key = (corr_idx, group_key.clone());
833
834 let suppressed = if let Some(suppress) = suppress_secs {
836 if let Some(&last_ts) = self.last_alert.get(&alert_key) {
837 (ts - last_ts) < suppress as i64
838 } else {
839 false
840 }
841 } else {
842 false
843 };
844
845 if !suppressed {
846 let (events, event_refs) = match event_mode {
848 CorrelationEventMode::Full => {
849 let stored = self
850 .event_buffers
851 .get(&alert_key)
852 .map(|buf| buf.decompress_all())
853 .unwrap_or_default();
854 (Some(stored), None)
855 }
856 CorrelationEventMode::Refs => {
857 let stored = self
858 .event_ref_buffers
859 .get(&alert_key)
860 .map(|buf| buf.refs())
861 .unwrap_or_default();
862 (None, Some(stored))
863 }
864 CorrelationEventMode::None => (None, None),
865 };
866
867 let corr = &self.correlations[corr_idx];
869 let result = EvaluationResult {
870 header: RuleHeader {
871 rule_title: corr.title.clone(),
872 rule_id: corr.id.clone(),
873 level,
874 tags: corr.tags.clone(),
875 custom_attributes: corr.custom_attributes.clone(),
876 enrichments: None,
877 },
878 body: ResultBody::Correlation(CorrelationBody {
879 correlation_type: corr_type,
880 group_key: group_key.to_pairs(&corr.group_by),
881 aggregated_value: agg_value,
882 timespan_secs: timespan,
883 events,
884 event_refs,
885 }),
886 };
887 out.push(result);
888
889 self.last_alert.insert(alert_key.clone(), ts);
891
892 if action == CorrelationAction::Reset {
894 if let Some(state) = self.state.get_mut(&alert_key) {
895 state.clear();
896 }
897 if let Some(buf) = self.event_buffers.get_mut(&alert_key) {
898 buf.clear();
899 }
900 if let Some(buf) = self.event_ref_buffers.get_mut(&alert_key) {
901 buf.clear();
902 }
903 }
904 }
905 }
906 }
907
908 fn chain_lookup_keys(&self, result: &EvaluationResult) -> Vec<String> {
914 let Some(id) = result.header.rule_id.as_deref() else {
915 return Vec::new();
916 };
917 let mut keys = vec![id.to_string()];
918 if let Some(name) = self
919 .correlations
920 .iter()
921 .find(|c| c.id.as_deref() == Some(id))
922 .and_then(|c| c.name.as_deref())
923 && name != id
924 {
925 keys.push(name.to_string());
926 }
927 keys
928 }
929
930 fn chain_correlations(
936 &mut self,
937 fired: &[EvaluationResult],
938 ts: i64,
939 out: &mut Vec<EvaluationResult>,
940 ) {
941 let mut pending: Vec<EvaluationResult> = fired.to_vec();
942 let mut depth = 0;
943
944 while !pending.is_empty() && depth < MAX_CHAIN_DEPTH {
945 depth += 1;
946
947 #[allow(clippy::type_complexity)]
949 let mut work: Vec<(usize, Vec<(String, String)>, String)> = Vec::new();
950 let mut seen = std::collections::HashSet::<(usize, String)>::new();
951 for result in &pending {
952 let Some(body) = result.as_correlation() else {
954 continue;
955 };
956 for key in self.chain_lookup_keys(result) {
957 if let Some(indices) = self.rule_index.get(&key) {
958 for &corr_idx in indices {
959 if seen.insert((corr_idx, key.clone())) {
960 work.push((corr_idx, body.group_key.clone(), key.clone()));
961 }
962 }
963 }
964 }
965 }
966
967 let mut next_pending = Vec::new();
968 for (corr_idx, group_key_pairs, fired_ref) in work {
969 let corr = &self.correlations[corr_idx];
970 let corr_type = corr.correlation_type;
971 let timespan = corr.timespan_secs;
972 let window_mode = corr.window_mode;
973 let gap_secs = corr.gap_secs;
974 let level = corr.level;
975 let suppress_secs = corr.suppress_secs.or(self.config.suppress);
976 let action = corr.action.unwrap_or(self.config.action_on_match);
977
978 let group_key = GroupKey::from_pairs(&group_key_pairs, &corr.group_by);
979 let state_key = (corr_idx, group_key.clone());
980 let state = self
981 .state
982 .entry(state_key.clone())
983 .or_insert_with(|| WindowState::new_for(corr_type));
984
985 if apply_window_open(state, ts, timespan, window_mode, gap_secs)
989 == WindowDecision::Discard
990 {
991 continue;
992 }
993
994 match corr_type {
995 CorrelationType::EventCount => {
996 state.push_event_count(ts);
997 }
998 CorrelationType::Temporal | CorrelationType::TemporalOrdered => {
999 state.push_temporal(ts, &fired_ref);
1000 }
1001 _ => {
1002 state.push_event_count(ts);
1003 }
1004 }
1005
1006 if let Some(cap) = corr.max_group_entries.or(self.config.max_group_entries) {
1009 state.truncate_oldest(cap, window_mode == WindowMode::Session);
1010 }
1011
1012 let fired = state.check_condition(
1013 &corr.condition,
1014 corr_type,
1015 &corr.rule_refs,
1016 corr.extended_expr.as_ref(),
1017 );
1018
1019 if let Some(agg_value) = fired {
1020 let alert_key = state_key;
1021 let suppressed = if let Some(suppress) = suppress_secs {
1022 self.last_alert
1023 .get(&alert_key)
1024 .is_some_and(|&last_ts| (ts - last_ts) < suppress as i64)
1025 } else {
1026 false
1027 };
1028 if suppressed {
1029 continue;
1030 }
1031
1032 let corr = &self.correlations[corr_idx];
1033 let result = EvaluationResult {
1034 header: RuleHeader {
1035 rule_title: corr.title.clone(),
1036 rule_id: corr.id.clone(),
1037 level,
1038 tags: corr.tags.clone(),
1039 custom_attributes: corr.custom_attributes.clone(),
1040 enrichments: None,
1041 },
1042 body: ResultBody::Correlation(CorrelationBody {
1043 correlation_type: corr_type,
1044 group_key: group_key.to_pairs(&corr.group_by),
1045 aggregated_value: agg_value,
1046 timespan_secs: timespan,
1047 events: None,
1051 event_refs: None,
1052 }),
1053 };
1054 next_pending.push(result.clone());
1055 out.push(result);
1056 self.last_alert.insert(alert_key.clone(), ts);
1057
1058 if action == CorrelationAction::Reset
1059 && let Some(state) = self.state.get_mut(&alert_key)
1060 {
1061 state.clear();
1062 }
1063 }
1064 }
1065
1066 pending = next_pending;
1067 }
1068
1069 if !pending.is_empty() {
1070 log::warn!(
1071 "Correlation chain depth limit reached ({MAX_CHAIN_DEPTH}); \
1072 {} pending result(s) were not propagated further. \
1073 This may indicate a cycle in correlation references.",
1074 pending.len()
1075 );
1076 }
1077 }
1078
1079 fn extract_event_timestamp(&self, event: &impl Event) -> Option<i64> {
1091 for field_name in &self.config.timestamp_fields {
1092 if let Some(val) = event.get_field(field_name)
1093 && let Some(ts) = parse_timestamp_value(&val)
1094 {
1095 return Some(ts);
1096 }
1097 }
1098 None
1099 }
1100
1101 pub fn evict_expired(&mut self, now_secs: i64) {
1107 self.evict_all(now_secs);
1108 }
1109
1110 fn evict_all(&mut self, now_secs: i64) {
1112 let specs: Vec<(u64, WindowMode, Option<u64>)> = self
1123 .correlations
1124 .iter()
1125 .map(|c| (c.timespan_secs, c.window_mode, c.gap_secs))
1126 .collect();
1127
1128 self.state.retain(|&(corr_idx, _), state| {
1129 if let Some(&(timespan, mode, gap)) = specs.get(corr_idx) {
1130 match mode {
1131 WindowMode::Sliding => {
1132 state.evict(now_secs - timespan as i64);
1133 }
1134 WindowMode::Tumbling | WindowMode::Session => {
1135 let staleness = if mode == WindowMode::Session {
1136 gap.unwrap_or(timespan)
1137 } else {
1138 timespan
1139 } as i64;
1140 if state
1141 .latest_timestamp()
1142 .is_some_and(|last| now_secs - last > staleness)
1143 {
1144 state.clear();
1145 }
1146 }
1147 }
1148 }
1149 !state.is_empty()
1150 });
1151
1152 let state = &self.state;
1156 self.event_buffers.retain(|key, buf| {
1157 if let Some(&(timespan, mode, _)) = specs.get(key.0) {
1158 match mode {
1159 WindowMode::Sliding => buf.evict(now_secs - timespan as i64),
1160 WindowMode::Tumbling | WindowMode::Session => {
1161 if !state.contains_key(key) {
1162 return false;
1163 }
1164 }
1165 }
1166 }
1167 !buf.is_empty()
1168 });
1169 self.event_ref_buffers.retain(|key, buf| {
1170 if let Some(&(timespan, mode, _)) = specs.get(key.0) {
1171 match mode {
1172 WindowMode::Sliding => buf.evict(now_secs - timespan as i64),
1173 WindowMode::Tumbling | WindowMode::Session => {
1174 if !state.contains_key(key) {
1175 return false;
1176 }
1177 }
1178 }
1179 }
1180 !buf.is_empty()
1181 });
1182
1183 if self.state.len() >= self.config.max_state_entries {
1187 let target = self.config.max_state_entries * 9 / 10;
1188 let excess = self.state.len() - target;
1189
1190 log::warn!(
1191 "Correlation state hard cap reached ({} entries, max {}); \
1192 evicting {} stalest entries to {} (90% capacity). \
1193 This indicates high-cardinality traffic; consider raising \
1194 max_state_entries or shortening correlation windows.",
1195 self.state.len(),
1196 self.config.max_state_entries,
1197 excess,
1198 target,
1199 );
1200
1201 let mut by_staleness: Vec<_> = self
1203 .state
1204 .iter()
1205 .map(|(k, v)| (k.clone(), v.latest_timestamp().unwrap_or(i64::MIN)))
1206 .collect();
1207 by_staleness.sort_unstable_by_key(|&(_, ts)| ts);
1208
1209 for (key, _) in by_staleness.into_iter().take(excess) {
1211 self.state.remove(&key);
1212 self.last_alert.remove(&key);
1213 self.event_buffers.remove(&key);
1214 self.event_ref_buffers.remove(&key);
1215 }
1216 }
1217
1218 self.last_alert.retain(|key, &mut alert_ts| {
1221 let suppress = if key.0 < self.correlations.len() {
1222 self.correlations[key.0]
1223 .suppress_secs
1224 .or(self.config.suppress)
1225 .unwrap_or(0)
1226 } else {
1227 0
1228 };
1229 (now_secs - alert_ts) < suppress as i64
1230 });
1231 }
1232
1233 pub fn state_count(&self) -> usize {
1235 self.state.len()
1236 }
1237
1238 pub fn detection_rule_count(&self) -> usize {
1240 self.engine.rule_count()
1241 }
1242
1243 pub fn correlation_rule_count(&self) -> usize {
1245 self.correlations.len()
1246 }
1247
1248 pub fn event_buffer_count(&self) -> usize {
1250 self.event_buffers.len()
1251 }
1252
1253 pub fn event_buffer_bytes(&self) -> usize {
1255 self.event_buffers
1256 .values()
1257 .map(|b| b.compressed_bytes())
1258 .sum()
1259 }
1260
1261 pub fn event_ref_buffer_count(&self) -> usize {
1263 self.event_ref_buffers.len()
1264 }
1265
1266 pub fn engine(&self) -> &Engine {
1268 &self.engine
1269 }
1270
1271 pub fn rule_metadata(&self, key: &str) -> RuleMetadataLookup {
1275 let mut variants = Vec::new();
1276 self.collect_rule_metadata(key, &mut variants);
1277 RuleMetadataLookup::from_variants(variants)
1278 }
1279
1280 pub(crate) fn collect_rule_metadata(&self, key: &str, out: &mut Vec<RuleBundleMetadata>) {
1281 self.engine.collect_rule_metadata(key, out);
1282 crate::rule_metadata::matching_correlations(&self.correlations, key, out);
1283 }
1284
1285 pub fn export_state(&self) -> CorrelationSnapshot {
1291 let mut windows: HashMap<String, Vec<(GroupKey, WindowState)>> = HashMap::new();
1292 for ((idx, gk), ws) in &self.state {
1293 let corr_id = self.correlation_stable_id(*idx);
1294 windows
1295 .entry(corr_id)
1296 .or_default()
1297 .push((gk.clone(), ws.clone()));
1298 }
1299
1300 let mut last_alert: HashMap<String, Vec<(GroupKey, i64)>> = HashMap::new();
1301 for ((idx, gk), ts) in &self.last_alert {
1302 let corr_id = self.correlation_stable_id(*idx);
1303 last_alert
1304 .entry(corr_id)
1305 .or_default()
1306 .push((gk.clone(), *ts));
1307 }
1308
1309 let mut event_buffers: HashMap<String, Vec<(GroupKey, EventBuffer)>> = HashMap::new();
1310 for ((idx, gk), buf) in &self.event_buffers {
1311 let corr_id = self.correlation_stable_id(*idx);
1312 event_buffers
1313 .entry(corr_id)
1314 .or_default()
1315 .push((gk.clone(), buf.clone()));
1316 }
1317
1318 let mut event_ref_buffers: HashMap<String, Vec<(GroupKey, EventRefBuffer)>> =
1319 HashMap::new();
1320 for ((idx, gk), buf) in &self.event_ref_buffers {
1321 let corr_id = self.correlation_stable_id(*idx);
1322 event_ref_buffers
1323 .entry(corr_id)
1324 .or_default()
1325 .push((gk.clone(), buf.clone()));
1326 }
1327
1328 CorrelationSnapshot {
1329 version: SNAPSHOT_VERSION,
1330 windows,
1331 last_alert,
1332 event_buffers,
1333 event_ref_buffers,
1334 }
1335 }
1336
1337 pub fn import_state(&mut self, snapshot: CorrelationSnapshot) -> bool {
1344 if snapshot.version != SNAPSHOT_VERSION {
1345 return false;
1346 }
1347 let id_to_idx = self.build_id_to_index_map();
1348
1349 for (corr_id, groups) in snapshot.windows {
1350 if let Some(&idx) = id_to_idx.get(&corr_id) {
1351 for (gk, ws) in groups {
1352 self.state.insert((idx, gk), ws);
1353 }
1354 }
1355 }
1356
1357 for (corr_id, groups) in snapshot.last_alert {
1358 if let Some(&idx) = id_to_idx.get(&corr_id) {
1359 for (gk, ts) in groups {
1360 self.last_alert.insert((idx, gk), ts);
1361 }
1362 }
1363 }
1364
1365 for (corr_id, groups) in snapshot.event_buffers {
1366 if let Some(&idx) = id_to_idx.get(&corr_id) {
1367 for (gk, buf) in groups {
1368 self.event_buffers.insert((idx, gk), buf);
1369 }
1370 }
1371 }
1372
1373 for (corr_id, groups) in snapshot.event_ref_buffers {
1374 if let Some(&idx) = id_to_idx.get(&corr_id) {
1375 for (gk, buf) in groups {
1376 self.event_ref_buffers.insert((idx, gk), buf);
1377 }
1378 }
1379 }
1380
1381 true
1382 }
1383
1384 fn correlation_stable_id(&self, idx: usize) -> String {
1386 let corr = &self.correlations[idx];
1387 corr.id
1388 .clone()
1389 .or_else(|| corr.name.clone())
1390 .unwrap_or_else(|| corr.title.clone())
1391 }
1392
1393 fn build_id_to_index_map(&self) -> HashMap<String, usize> {
1395 self.correlations
1396 .iter()
1397 .enumerate()
1398 .map(|(idx, _)| (self.correlation_stable_id(idx), idx))
1399 .collect()
1400 }
1401}
1402
1403impl Default for CorrelationEngine {
1404 fn default() -> Self {
1405 Self::new(CorrelationConfig::default())
1406 }
1407}
1408
1409fn extract_event_ts(event: &impl Event, timestamp_fields: &[String]) -> Option<i64> {
1418 for field_name in timestamp_fields {
1419 if let Some(val) = event.get_field(field_name)
1420 && let Some(ts) = parse_timestamp_value(&val)
1421 {
1422 return Some(ts);
1423 }
1424 }
1425 None
1426}
1427
1428fn parse_timestamp_value(val: &EventValue) -> Option<i64> {
1430 match val {
1431 EventValue::Int(i) => Some(normalize_epoch(*i)),
1432 EventValue::Float(f) => Some(normalize_epoch(*f as i64)),
1433 EventValue::Str(s) => parse_timestamp_string(s),
1434 _ => None,
1435 }
1436}
1437
1438fn normalize_epoch(v: i64) -> i64 {
1441 if v > 1_000_000_000_000 { v / 1000 } else { v }
1442}
1443
1444fn parse_timestamp_string(s: &str) -> Option<i64> {
1446 if let Ok(dt) = DateTime::parse_from_rfc3339(s) {
1448 return Some(dt.timestamp());
1449 }
1450
1451 if let Ok(naive) = chrono::NaiveDateTime::parse_from_str(s, "%Y-%m-%dT%H:%M:%S") {
1454 return Some(Utc.from_utc_datetime(&naive).timestamp());
1455 }
1456 if let Ok(naive) = chrono::NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S") {
1457 return Some(Utc.from_utc_datetime(&naive).timestamp());
1458 }
1459
1460 if let Ok(naive) = chrono::NaiveDateTime::parse_from_str(s, "%Y-%m-%dT%H:%M:%S%.f") {
1462 return Some(Utc.from_utc_datetime(&naive).timestamp());
1463 }
1464 if let Ok(naive) = chrono::NaiveDateTime::parse_from_str(s, "%Y-%m-%d %H:%M:%S%.f") {
1465 return Some(Utc.from_utc_datetime(&naive).timestamp());
1466 }
1467
1468 None
1469}
1470
1471fn value_to_string_for_count(v: &EventValue) -> Option<String> {
1473 match v {
1474 EventValue::Str(s) => Some(s.to_string()),
1475 EventValue::Int(n) => Some(n.to_string()),
1476 EventValue::Float(f) => Some(f.to_string()),
1477 EventValue::Bool(b) => Some(b.to_string()),
1478 EventValue::Null => Some("null".to_string()),
1479 _ => None,
1480 }
1481}
1482
1483fn composite_value_count_key(event: &impl Event, fields: &[String]) -> Option<String> {
1492 if let [field_name] = fields {
1494 let val = event.get_field(field_name)?;
1495 return value_to_string_for_count(&val);
1496 }
1497
1498 let mut parts = Vec::with_capacity(fields.len());
1499 for field_name in fields {
1500 let val = event.get_field(field_name)?;
1501 let rendered = value_to_string_for_count(&val)?;
1502 parts.push(rendered);
1503 }
1504 Some(parts.join("\u{1f}"))
1505}
1506
1507fn value_to_f64_ev(v: &EventValue) -> Option<f64> {
1509 v.as_f64()
1510}