use std::cmp::Reverse;
use std::collections::BinaryHeap;
use arrow::record_batch::RecordBatch;
use crate::cep::pattern::CompiledPattern;
#[derive(Debug, Clone)]
pub struct PartialMatch {
pub stage_index: usize,
pub captured_events: Vec<RecordBatch>,
pub start_time_ms: i64,
pub captured_event_count: usize,
}
#[derive(Debug, Default, Clone, serde::Serialize, serde::Deserialize)]
pub struct CepKeyState {
#[serde(skip)]
pub partial: Option<PartialMatch>,
pub last_event_ms: i64,
}
#[derive(Debug, Clone)]
pub struct SequentialPatternMatcher {
pattern: CompiledPattern,
}
impl SequentialPatternMatcher {
pub fn new(pattern: CompiledPattern) -> Self {
Self { pattern }
}
pub fn process_event(
&self,
state: &mut CepKeyState,
stage_name: &str,
batch: RecordBatch,
event_time_ms: i64,
) -> Vec<Vec<RecordBatch>> {
state.last_event_ms = event_time_ms;
if let Some(ref partial) = state.partial
&& event_time_ms.saturating_sub(partial.start_time_ms)
> i64::try_from(self.pattern.window_ms).unwrap_or(i64::MAX)
{
state.partial = None;
}
let stage_idx = self
.pattern
.stages
.iter()
.position(|s| s.name == stage_name);
let Some(stage_idx) = stage_idx else {
return Vec::new();
};
if state.partial.is_none() {
if stage_idx != 0 {
return Vec::new();
}
state.partial = Some(PartialMatch {
stage_index: 0,
captured_events: vec![batch],
start_time_ms: event_time_ms,
captured_event_count: 1,
});
if self.pattern.stages.len() == 1 {
return self.take_complete(state);
}
return Vec::new();
}
if let Some(ref mut partial) = state.partial {
let expected_next = partial.stage_index + 1;
if stage_idx != expected_next {
return Vec::new();
}
partial.captured_events.push(batch);
partial.captured_event_count = partial.captured_events.len();
partial.stage_index = stage_idx;
if partial.stage_index + 1 == self.pattern.stages.len() {
return self.take_complete(state);
}
}
Vec::new()
}
fn take_complete(&self, state: &mut CepKeyState) -> Vec<Vec<RecordBatch>> {
state
.partial
.take()
.map(|p| vec![p.captured_events])
.unwrap_or_default()
}
}
#[derive(Debug, Clone)]
pub struct PartitionedCepMatcher<K>
where
K: std::hash::Hash + Eq + Clone + Ord,
{
pattern: CompiledPattern,
states: std::collections::HashMap<K, (SequentialPatternMatcher, CepKeyState)>,
max_partitions: usize,
eviction_heap: BinaryHeap<Reverse<(i64, K)>>,
}
impl<K> PartitionedCepMatcher<K>
where
K: std::hash::Hash + Eq + Clone + Ord,
{
pub fn new(pattern: CompiledPattern) -> Self {
Self {
pattern,
states: std::collections::HashMap::new(),
max_partitions: 1024,
eviction_heap: BinaryHeap::new(),
}
}
pub fn process_event(
&mut self,
key: K,
stage_name: &str,
batch: RecordBatch,
event_time_ms: i64,
) -> Vec<Vec<RecordBatch>> {
let entry = self.states.entry(key.clone()).or_insert_with(|| {
(
SequentialPatternMatcher::new(self.pattern.clone()),
CepKeyState::default(),
)
});
let result = entry
.0
.process_event(&mut entry.1, stage_name, batch, event_time_ms);
self.eviction_heap.push(Reverse((event_time_ms, key)));
if self.states.len() > self.max_partitions {
self.evict_stalest();
}
result
}
fn evict_stalest(&mut self) {
while let Some(Reverse((ts, k))) = self.eviction_heap.pop() {
if let Some((_, state)) = self.states.get(&k)
&& state.last_event_ms == ts
{
self.states.remove(&k);
return;
}
}
}
pub fn evict_keys_before(&mut self, cutoff_ms: i64) {
self.states
.retain(|_, (_, state)| state.last_event_ms >= cutoff_ms);
}
pub fn partition_count(&self) -> usize {
self.states.len()
}
pub fn partial_signature(&self, key: &K) -> Option<(usize, i64)> {
self.states
.get(key)
.and_then(|(_, state)| state.partial.as_ref())
.map(|partial| (partial.stage_index, partial.start_time_ms))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cep::pattern::Pattern;
use arrow::array::{Int32Array, RecordBatch};
use arrow::datatypes::{DataType, Field, Schema};
use std::sync::Arc;
use std::time::Duration;
fn schema() -> Arc<Schema> {
Arc::new(Schema::new(vec![
Field::new("event_type", DataType::Utf8, false),
Field::new("timestamp", DataType::Int64, false),
Field::new("value", DataType::Int32, false),
]))
}
fn batch(v: i32) -> RecordBatch {
RecordBatch::try_new(
Arc::new(Schema::new(vec![Field::new("v", DataType::Int32, false)])),
vec![Arc::new(Int32Array::from(vec![v]))],
)
.unwrap()
}
fn rich_batch(event_type: &str, timestamp: i64, value: i32) -> RecordBatch {
RecordBatch::try_new(
schema(),
vec![
Arc::new(arrow::array::StringArray::from(vec![event_type])),
Arc::new(arrow::array::Int64Array::from(vec![timestamp])),
Arc::new(Int32Array::from(vec![value])),
],
)
.unwrap()
}
#[test]
fn two_stage_pattern_matches() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
assert!(
matcher
.process_event(&mut state, "a", batch(1), 100)
.is_empty()
);
let done = matcher.process_event(&mut state, "b", batch(2), 200);
assert_eq!(done.len(), 1);
assert_eq!(done[0].len(), 2);
}
#[test]
fn expired_partial_discarded() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_millis(50))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 0);
assert!(
matcher
.process_event(&mut state, "b", batch(2), 100)
.is_empty()
);
}
#[test]
fn empty_pattern_compile_rejected() {
let result = Pattern::begin("a")
.compile()
.unwrap() ;
assert_eq!(result.stages.len(), 1);
}
#[test]
fn single_stage_match_completes_immediately() {
let pattern = Pattern::begin("only")
.within(Duration::from_secs(1))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
let done = matcher.process_event(&mut state, "only", batch(42), 100);
assert_eq!(
done.len(),
1,
"single-stage pattern must complete on first match"
);
assert_eq!(done[0].len(), 1);
assert!(
state.partial.is_none(),
"state must be cleared after completion"
);
}
#[test]
fn boundary_event_at_exact_window_limit() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_millis(100))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 0);
let done = matcher.process_event(&mut state, "b", batch(2), 100);
assert_eq!(done.len(), 1, "event at exact window boundary must match");
}
#[test]
fn boundary_event_one_ms_past_window() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_millis(100))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 0);
let done = matcher.process_event(&mut state, "b", batch(2), 101);
assert!(done.is_empty(), "event past window must be discarded");
}
#[test]
fn partitioned_matcher_independent_keys() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let mut pm = PartitionedCepMatcher::<String>::new(pattern);
assert!(pm.process_event("k1".into(), "a", batch(1), 100).is_empty());
assert!(
pm.process_event("k2".into(), "a", batch(10), 200)
.is_empty()
);
let done = pm.process_event("k1".into(), "b", batch(2), 300);
assert_eq!(done.len(), 1);
assert!(
pm.process_event("k2".into(), "a", batch(11), 400)
.is_empty()
);
}
#[test]
fn partitioned_matcher_independent_state() {
let pattern = Pattern::begin("x")
.followed_by("y")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let mut pm = PartitionedCepMatcher::<i32>::new(pattern);
pm.process_event(1, "x", batch(1), 100);
pm.process_event(2, "x", batch(2), 200);
assert!(pm.states.contains_key(&1));
assert!(pm.states.contains_key(&2));
}
#[test]
fn wrong_stage_name_ignored() {
let pattern = Pattern::begin("a").followed_by("b").compile().unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
let result = matcher.process_event(&mut state, "c", batch(1), 100);
assert!(result.is_empty());
assert!(
state.partial.is_none(),
"no partial match should be started"
);
}
#[test]
fn partial_state_persisted_after_first_event() {
let pattern = Pattern::begin("a").followed_by("b").compile().unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 100);
assert!(
state.partial.is_some(),
"partial must exist after first stage"
);
let partial = state.partial.as_ref().unwrap();
assert_eq!(partial.stage_index, 0);
assert_eq!(partial.captured_events.len(), 1);
assert_eq!(partial.start_time_ms, 100);
}
#[test]
fn stage_ordering_enforced() {
let pattern = Pattern::begin("a")
.followed_by("b")
.followed_by("c")
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 100);
let result = matcher.process_event(&mut state, "c", batch(3), 200);
assert!(result.is_empty());
assert!(
state.partial.is_some(),
"partial should still be waiting for stage b"
);
}
#[test]
fn three_stage_sequential_match() {
let pattern = Pattern::begin("a")
.followed_by("b")
.followed_by("c")
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
assert!(
matcher
.process_event(&mut state, "a", batch(1), 100)
.is_empty()
);
assert!(
matcher
.process_event(&mut state, "b", batch(2), 200)
.is_empty()
);
let done = matcher.process_event(&mut state, "c", batch(3), 300);
assert_eq!(done.len(), 1);
assert_eq!(done[0].len(), 3);
assert_eq!(
done[0][0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::Int32Array>()
.unwrap()
.value(0),
1
);
assert_eq!(
done[0][1]
.column(0)
.as_any()
.downcast_ref::<arrow::array::Int32Array>()
.unwrap()
.value(0),
2
);
assert_eq!(
done[0][2]
.column(0)
.as_any()
.downcast_ref::<arrow::array::Int32Array>()
.unwrap()
.value(0),
3
);
}
#[test]
fn multiple_matches_on_same_key() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 100);
let done1 = matcher.process_event(&mut state, "b", batch(2), 200);
assert_eq!(done1.len(), 1);
assert!(state.partial.is_none(), "state cleared after first match");
matcher.process_event(&mut state, "a", batch(10), 300);
let done2 = matcher.process_event(&mut state, "b", batch(20), 400);
assert_eq!(done2.len(), 1);
}
#[test]
fn last_event_ms_updated() {
let pattern = Pattern::begin("a").followed_by("b").compile().unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
assert_eq!(state.last_event_ms, 0);
matcher.process_event(&mut state, "a", batch(1), 500);
assert_eq!(state.last_event_ms, 500);
matcher.process_event(&mut state, "b", batch(2), 600);
assert_eq!(state.last_event_ms, 600);
}
#[test]
fn wrong_stage_between_matches_does_not_corrupt_state() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 100);
assert!(
matcher
.process_event(&mut state, "x", batch(99), 150)
.is_empty()
);
let done = matcher.process_event(&mut state, "b", batch(2), 200);
assert_eq!(done.len(), 1);
}
#[test]
fn out_of_order_stage_after_partial_resets_correctly() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 100);
assert!(
matcher
.process_event(&mut state, "a", batch(10), 150)
.is_empty()
);
let done = matcher.process_event(&mut state, "b", batch(2), 200);
assert_eq!(done.len(), 1);
}
#[test]
fn rich_batch_sequential_match() {
let pattern = Pattern::begin("login")
.followed_by("query")
.within(Duration::from_secs(10))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
let b1 = rich_batch("login", 1000, 0);
let b2 = rich_batch("query", 2000, 42);
assert!(
matcher
.process_event(&mut state, "login", b1, 1000)
.is_empty()
);
let done = matcher.process_event(&mut state, "query", b2, 2000);
assert_eq!(done.len(), 1);
assert_eq!(done[0].len(), 2);
let col = done[0][0]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(col.value(0), "login");
let col = done[0][1]
.column(0)
.as_any()
.downcast_ref::<arrow::array::StringArray>()
.unwrap();
assert_eq!(col.value(0), "query");
}
#[test]
fn default_window_is_60s() {
let pattern = Pattern::begin("a").compile().unwrap();
assert_eq!(pattern.window_ms, 60_000);
}
#[test]
fn no_partial_match_when_no_events_processed() {
let pattern = Pattern::begin("a").followed_by("b").compile().unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
let result = matcher.process_event(&mut state, "b", batch(2), 100);
assert!(result.is_empty());
assert!(state.partial.is_none());
}
#[test]
fn partitioned_wrong_stage_ignored_per_key() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let mut pm = PartitionedCepMatcher::<String>::new(pattern);
pm.process_event("k1".into(), "a", batch(1), 100);
assert!(
pm.process_event("k1".into(), "x", batch(99), 200)
.is_empty()
);
let done = pm.process_event("k1".into(), "b", batch(2), 300);
assert_eq!(done.len(), 1);
}
#[test]
fn partitioned_multiple_matches_per_key() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let mut pm = PartitionedCepMatcher::<i32>::new(pattern);
pm.process_event(1, "a", batch(1), 100);
let done1 = pm.process_event(1, "b", batch(2), 200);
assert_eq!(done1.len(), 1);
pm.process_event(1, "a", batch(10), 300);
let done2 = pm.process_event(1, "b", batch(20), 400);
assert_eq!(done2.len(), 1);
}
#[test]
fn partitioned_three_stage_match() {
let pattern = Pattern::begin("a")
.followed_by("b")
.followed_by("c")
.within(Duration::from_secs(10))
.compile()
.unwrap();
let mut pm = PartitionedCepMatcher::<String>::new(pattern);
assert!(pm.process_event("k1".into(), "a", batch(1), 100).is_empty());
assert!(pm.process_event("k1".into(), "b", batch(2), 200).is_empty());
let done = pm.process_event("k1".into(), "c", batch(3), 300);
assert_eq!(done.len(), 1);
assert_eq!(done[0].len(), 3);
}
#[test]
fn partitioned_independent_timeout_per_key() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_millis(50))
.compile()
.unwrap();
let mut pm = PartitionedCepMatcher::<String>::new(pattern);
pm.process_event("k1".into(), "a", batch(1), 0);
pm.process_event("k2".into(), "a", batch(10), 1000);
assert!(pm.process_event("k1".into(), "b", batch(2), 60).is_empty());
let done = pm.process_event("k2".into(), "b", batch(20), 1040);
assert_eq!(done.len(), 1);
}
#[test]
fn partitioned_wrong_key_stage_not_cross_contaminated() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let mut pm = PartitionedCepMatcher::<String>::new(pattern);
pm.process_event("k1".into(), "a", batch(1), 100);
pm.process_event("k2".into(), "a", batch(2), 200);
let done = pm.process_event("k1".into(), "b", batch(3), 300);
assert_eq!(done.len(), 1);
assert!(pm.states.get("k2").unwrap().1.partial.is_some());
let done2 = pm.process_event("k2".into(), "b", batch(4), 400);
assert_eq!(done2.len(), 1);
}
#[test]
fn partitioned_rich_batch_preserves_data() {
let pattern = Pattern::begin("click")
.followed_by("purchase")
.within(Duration::from_secs(30))
.compile()
.unwrap();
let mut pm = PartitionedCepMatcher::<String>::new(pattern);
let b1 = rich_batch("click", 1000, 0);
let b2 = rich_batch("purchase", 2000, 99);
assert!(
pm.process_event("user1".into(), "click", b1, 1000)
.is_empty()
);
let done = pm.process_event("user1".into(), "purchase", b2, 2000);
assert_eq!(done.len(), 1);
assert_eq!(done[0].len(), 2);
let val_col = done[0][1]
.column(2)
.as_any()
.downcast_ref::<Int32Array>()
.unwrap();
assert_eq!(val_col.value(0), 99);
}
#[test]
fn partitioned_new_key_auto_created() {
let pattern = Pattern::begin("a").followed_by("b").compile().unwrap();
let mut pm = PartitionedCepMatcher::<String>::new(pattern);
assert!(pm.states.is_empty());
pm.process_event("new_key".into(), "a", batch(1), 100);
assert!(pm.states.contains_key("new_key"));
assert_eq!(pm.states.len(), 1);
}
#[test]
fn negative_event_timestamps() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), -1000);
let done = matcher.process_event(&mut state, "b", batch(2), -500);
assert_eq!(done.len(), 1);
}
#[test]
fn negative_timestamp_window_expired() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_millis(100))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), -200);
let done = matcher.process_event(&mut state, "b", batch(2), -50);
assert!(
done.is_empty(),
"event past window with negative timestamps must be discarded"
);
}
#[test]
fn large_window_millis() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_millis(u64::MAX / 2))
.compile()
.unwrap();
assert!(pattern.window_ms > 0);
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 0);
let done = matcher.process_event(&mut state, "b", batch(2), 1_000_000);
assert_eq!(done.len(), 1);
}
#[test]
fn multi_row_batch_preserves_all_rows() {
let pattern = Pattern::begin("a").followed_by("b").compile().unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
let multi_batch = RecordBatch::try_new(
Arc::new(Schema::new(vec![Field::new("v", DataType::Int32, false)])),
vec![Arc::new(Int32Array::from(vec![10, 20, 30]))],
)
.unwrap();
matcher.process_event(&mut state, "a", multi_batch.clone(), 100);
let done = matcher.process_event(&mut state, "b", batch(2), 200);
assert_eq!(done.len(), 1);
assert_eq!(done[0].len(), 2);
let col = done[0][0]
.column(0)
.as_any()
.downcast_ref::<Int32Array>()
.unwrap();
assert_eq!(col.len(), 3);
assert_eq!(col.value(0), 10);
assert_eq!(col.value(1), 20);
assert_eq!(col.value(2), 30);
}
#[test]
fn first_event_at_zero_time() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(1))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 0);
let done = matcher.process_event(&mut state, "b", batch(2), 0);
assert_eq!(done.len(), 1);
}
#[test]
fn exact_duplicate_stage_names_reset_partial() {
let pattern = Pattern::begin("a").followed_by("a").compile().unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 100);
let result = matcher.process_event(&mut state, "a", batch(2), 200);
assert!(result.is_empty());
}
#[test]
fn five_stage_pattern() {
let pattern = Pattern::begin("s1")
.followed_by("s2")
.followed_by("s3")
.followed_by("s4")
.followed_by("s5")
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
assert!(
matcher
.process_event(&mut state, "s1", batch(1), 100)
.is_empty()
);
assert!(
matcher
.process_event(&mut state, "s2", batch(2), 200)
.is_empty()
);
assert!(
matcher
.process_event(&mut state, "s3", batch(3), 300)
.is_empty()
);
assert!(
matcher
.process_event(&mut state, "s4", batch(4), 400)
.is_empty()
);
let done = matcher.process_event(&mut state, "s5", batch(5), 500);
assert_eq!(done.len(), 1);
assert_eq!(done[0].len(), 5);
}
#[test]
fn partitioned_many_keys() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let mut pm = PartitionedCepMatcher::<i32>::new(pattern);
for k in 0..100 {
pm.process_event(k, "a", batch(k), k as i64 * 100);
}
assert_eq!(pm.states.len(), 100);
let done = pm.process_event(50, "b", batch(50), 5000);
assert_eq!(done.len(), 1);
assert!(pm.states.get(&0).unwrap().1.partial.is_some());
assert!(pm.states.get(&99).unwrap().1.partial.is_some());
}
#[test]
fn partitioned_completed_key_can_restart() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let mut pm = PartitionedCepMatcher::<String>::new(pattern);
pm.process_event("k".into(), "a", batch(1), 100);
let done1 = pm.process_event("k".into(), "b", batch(2), 200);
assert_eq!(done1.len(), 1);
assert!(pm.states.get("k").unwrap().1.partial.is_none());
pm.process_event("k".into(), "a", batch(10), 300);
let done2 = pm.process_event("k".into(), "b", batch(20), 400);
assert_eq!(done2.len(), 1);
}
#[test]
fn cep_key_state_default_values() {
let state = CepKeyState::default();
assert!(state.partial.is_none());
assert_eq!(state.last_event_ms, 0);
}
#[test]
fn partial_match_default_values() {
let pm = PartialMatch {
stage_index: 0,
captured_events: Vec::new(),
start_time_ms: 0,
captured_event_count: 0,
};
assert_eq!(pm.stage_index, 0);
assert!(pm.captured_events.is_empty());
assert_eq!(pm.start_time_ms, 0);
}
#[test]
fn compiled_pattern_clone() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_secs(5))
.compile()
.unwrap();
let cloned = pattern.clone();
assert_eq!(cloned.stages.len(), pattern.stages.len());
assert_eq!(cloned.window_ms, pattern.window_ms);
}
#[test]
fn cep_key_state_serde_skips_partial_but_preserves_metadata() {
let state = CepKeyState {
last_event_ms: 1_234_567,
..Default::default()
};
let json = serde_json::to_string(&state).unwrap();
let restored: CepKeyState = serde_json::from_str(&json).unwrap();
assert_eq!(restored.last_event_ms, 1_234_567);
assert!(restored.partial.is_none());
}
#[test]
fn sequential_matcher_clone() {
let pattern = Pattern::begin("a").compile().unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let cloned = matcher.clone();
let mut state = CepKeyState::default();
let done = cloned.process_event(&mut state, "a", batch(1), 100);
assert_eq!(done.len(), 1);
}
#[test]
fn partitioned_matcher_clone() {
let pattern = Pattern::begin("a").followed_by("b").compile().unwrap();
let mut pm = PartitionedCepMatcher::<String>::new(pattern);
pm.process_event("k1".into(), "a", batch(1), 100);
let cloned = pm.clone();
assert!(cloned.states.contains_key("k1"));
}
#[test]
fn zero_duration_window_allows_same_time_match() {
let pattern = Pattern::begin("a")
.followed_by("b")
.within(Duration::from_millis(0))
.compile()
.unwrap();
let matcher = SequentialPatternMatcher::new(pattern);
let mut state = CepKeyState::default();
matcher.process_event(&mut state, "a", batch(1), 100);
let done = matcher.process_event(&mut state, "b", batch(2), 100);
assert_eq!(done.len(), 1);
}
#[test]
fn window_ms_default_is_60000() {
let pattern = Pattern::begin("a").compile().unwrap();
assert_eq!(pattern.window_ms, 60_000);
}
#[test]
fn pattern_stage_names_and_gap() {
let pattern = Pattern::begin("start")
.followed_by("end")
.within(Duration::from_secs(10))
.compile()
.unwrap();
assert_eq!(pattern.stages[0].name, "start");
assert!(pattern.stages[0].max_gap_ms.is_none());
assert_eq!(pattern.stages[1].name, "end");
assert!(pattern.stages[1].max_gap_ms.is_none());
}
}