1use std::collections::BTreeSet;
8
9use crate::drift::DriftLevel;
10use crate::error::{RillError, checked_increment};
11use crate::persistence::ValidateState;
12
13const MAX_DETECTOR_NAME_BYTES: usize = 128;
14
15#[derive(Debug, Clone, PartialEq, Eq)]
17#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
18#[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
19#[non_exhaustive]
20pub struct DriftVote {
21 pub detector: String,
23 pub level: DriftLevel,
25 pub available: bool,
27}
28
29impl DriftVote {
30 pub fn new(detector: impl Into<String>, level: DriftLevel) -> Result<Self, RillError> {
32 let vote = Self {
33 detector: detector.into(),
34 level,
35 available: true,
36 };
37 vote.validate()?;
38 Ok(vote)
39 }
40
41 pub fn unavailable(detector: impl Into<String>) -> Result<Self, RillError> {
43 let vote = Self {
44 detector: detector.into(),
45 level: DriftLevel::None,
46 available: false,
47 };
48 vote.validate()?;
49 Ok(vote)
50 }
51
52 fn validate(&self) -> Result<(), RillError> {
53 if self.detector.is_empty() || self.detector.len() > MAX_DETECTOR_NAME_BYTES {
54 return Err(RillError::InvalidState(format!(
55 "drift detector name length must be in 1..={MAX_DETECTOR_NAME_BYTES} bytes"
56 )));
57 }
58 if !self.available && self.level != DriftLevel::None {
59 return Err(RillError::InvalidState(
60 "an unavailable detector cannot cast a warning or drift vote".to_owned(),
61 ));
62 }
63 Ok(())
64 }
65}
66
67#[derive(Debug, Clone, PartialEq, Eq)]
69#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
70#[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
71#[non_exhaustive]
72pub struct DriftConsensusConfig {
73 pub minimum_warning_votes: usize,
75 pub minimum_drift_votes: usize,
77 pub confirmation_windows: u64,
79 pub clear_windows: u64,
81 pub cooldown_windows: u64,
83 pub event_merge_horizon: u64,
85 pub warming_windows: u64,
87 pub max_detectors: usize,
89}
90
91impl Default for DriftConsensusConfig {
92 fn default() -> Self {
93 Self {
94 minimum_warning_votes: 1,
95 minimum_drift_votes: 2,
96 confirmation_windows: 2,
97 clear_windows: 3,
98 cooldown_windows: 5,
99 event_merge_horizon: 20,
100 warming_windows: 1,
101 max_detectors: 16,
102 }
103 }
104}
105
106impl DriftConsensusConfig {
107 pub fn validate(&self) -> Result<(), RillError> {
109 if self.max_detectors == 0 || self.max_detectors > 128 {
110 return Err(RillError::InvalidCapacity(self.max_detectors));
111 }
112 if self.minimum_warning_votes == 0
113 || self.minimum_warning_votes > self.max_detectors
114 || self.minimum_drift_votes == 0
115 || self.minimum_drift_votes > self.max_detectors
116 || self.minimum_warning_votes > self.minimum_drift_votes
117 {
118 return Err(RillError::InvalidState(
119 "drift consensus vote thresholds are inconsistent".to_owned(),
120 ));
121 }
122 if self.confirmation_windows == 0 || self.clear_windows == 0 {
123 return Err(RillError::InvalidState(
124 "confirmation_windows and clear_windows must be positive".to_owned(),
125 ));
126 }
127 Ok(())
128 }
129}
130
131#[derive(Debug, Clone, PartialEq, Eq)]
133#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
134#[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
135pub struct DriftEventSummary {
136 pub event_id: u64,
138 pub first_window: u64,
140 pub last_window: u64,
142 pub peak_level: DriftLevel,
144 pub detectors: Vec<String>,
146 pub occurrences: u64,
148 pub generation: u64,
150}
151
152#[derive(Debug, Clone, PartialEq, Eq)]
154pub struct DriftConsensusResult {
155 pub level: DriftLevel,
157 pub warning_votes: usize,
159 pub drift_votes: usize,
161 pub available_detectors: usize,
163 pub triggered_detectors: Vec<String>,
165 pub consecutive_windows: u64,
167 pub clear_streak: u64,
169 pub cooldown_remaining: u64,
171 pub warming: bool,
173 pub incomplete: bool,
175 pub generation_changed: bool,
177 pub new_event: bool,
179 pub merged_event: bool,
181 pub event: Option<DriftEventSummary>,
183}
184
185#[derive(Debug, Clone, PartialEq, Eq)]
187#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
188pub struct DriftConsensus {
189 config: DriftConsensusConfig,
190 windows_seen: u64,
191 generation_windows: u64,
192 active_generation: Option<u64>,
193 active_level: DriftLevel,
194 warning_streak: u64,
195 drift_streak: u64,
196 clear_streak: u64,
197 cooldown_remaining: u64,
198 next_event_id: u64,
199 last_event: Option<DriftEventSummary>,
200}
201
202impl DriftConsensus {
203 pub fn new(config: DriftConsensusConfig) -> Result<Self, RillError> {
205 config.validate()?;
206 Ok(Self {
207 config,
208 windows_seen: 0,
209 generation_windows: 0,
210 active_generation: None,
211 active_level: DriftLevel::None,
212 warning_streak: 0,
213 drift_streak: 0,
214 clear_streak: 0,
215 cooldown_remaining: 0,
216 next_event_id: 1,
217 last_event: None,
218 })
219 }
220
221 pub const fn config(&self) -> &DriftConsensusConfig {
223 &self.config
224 }
225
226 pub const fn windows_seen(&self) -> u64 {
228 self.windows_seen
229 }
230
231 pub const fn level(&self) -> DriftLevel {
233 self.active_level
234 }
235
236 pub const fn last_event(&self) -> Option<&DriftEventSummary> {
238 self.last_event.as_ref()
239 }
240
241 pub fn reset(&mut self) {
243 self.windows_seen = 0;
244 self.generation_windows = 0;
245 self.active_generation = None;
246 self.active_level = DriftLevel::None;
247 self.warning_streak = 0;
248 self.drift_streak = 0;
249 self.clear_streak = 0;
250 self.cooldown_remaining = 0;
251 self.next_event_id = 1;
252 self.last_event = None;
253 }
254
255 pub fn update(
257 &mut self,
258 votes: &[DriftVote],
259 generation: u64,
260 ) -> Result<DriftConsensusResult, RillError> {
261 self.validate_votes(votes)?;
262 let mut next = self.clone();
263 let result = next.update_inner(votes, generation)?;
264 *self = next;
265 Ok(result)
266 }
267
268 fn update_inner(
269 &mut self,
270 votes: &[DriftVote],
271 generation: u64,
272 ) -> Result<DriftConsensusResult, RillError> {
273 let next_windows = checked_increment(self.windows_seen, "consensus windows_seen")?;
274 let generation_changed = self.active_generation.is_some_and(|old| old != generation);
275 if generation_changed {
276 self.generation_windows = 0;
277 self.active_level = DriftLevel::None;
278 self.warning_streak = 0;
279 self.drift_streak = 0;
280 self.clear_streak = 0;
281 self.cooldown_remaining = 0;
282 }
283 self.active_generation = Some(generation);
284 self.windows_seen = next_windows;
285 self.generation_windows =
286 checked_increment(self.generation_windows, "consensus generation_windows")?;
287 let event_allowed = self.cooldown_remaining == 0;
288 if self.cooldown_remaining > 0 {
289 self.cooldown_remaining -= 1;
290 }
291
292 let available_detectors = votes.iter().filter(|vote| vote.available).count();
293 let drift_votes = votes
294 .iter()
295 .filter(|vote| vote.available && vote.level == DriftLevel::Drift)
296 .count();
297 let warning_votes = votes
298 .iter()
299 .filter(|vote| vote.available && vote.level.is_change())
300 .count();
301 let triggered_detectors: Vec<String> = votes
302 .iter()
303 .filter(|vote| vote.available && vote.level.is_change())
304 .map(|vote| vote.detector.clone())
305 .collect();
306 let warming = self.generation_windows <= self.config.warming_windows;
307 let incomplete = available_detectors < self.config.minimum_warning_votes;
308 if warming || incomplete {
309 self.warning_streak = 0;
310 self.drift_streak = 0;
311 return Ok(self.result(
312 warning_votes,
313 drift_votes,
314 available_detectors,
315 triggered_detectors,
316 0,
317 warming,
318 incomplete,
319 generation_changed,
320 false,
321 false,
322 ));
323 }
324
325 let raw_level = if drift_votes >= self.config.minimum_drift_votes {
326 DriftLevel::Drift
327 } else if warning_votes >= self.config.minimum_warning_votes {
328 DriftLevel::Warning
329 } else {
330 DriftLevel::None
331 };
332
333 let mut new_event = false;
334 let mut merged_event = false;
335 let consecutive_windows;
336 match raw_level {
337 DriftLevel::Drift => {
338 self.clear_streak = 0;
339 self.warning_streak = 0;
340 self.drift_streak = checked_increment(self.drift_streak, "drift streak")?;
341 consecutive_windows = self.drift_streak;
342 if self.drift_streak >= self.config.confirmation_windows
343 && (self.active_level != DriftLevel::Drift
344 || self
345 .drift_streak
346 .is_multiple_of(self.config.confirmation_windows))
347 {
348 self.active_level = DriftLevel::Drift;
349 if event_allowed {
350 (new_event, merged_event) =
351 self.record_event(DriftLevel::Drift, &triggered_detectors, generation)?;
352 }
353 }
354 }
355 DriftLevel::Warning => {
356 self.clear_streak = 0;
357 self.drift_streak = 0;
358 self.warning_streak = checked_increment(self.warning_streak, "warning streak")?;
359 consecutive_windows = self.warning_streak;
360 if self.active_level == DriftLevel::None
361 && self.warning_streak >= self.config.confirmation_windows
362 {
363 self.active_level = DriftLevel::Warning;
364 if event_allowed {
365 (new_event, merged_event) = self.record_event(
366 DriftLevel::Warning,
367 &triggered_detectors,
368 generation,
369 )?;
370 }
371 }
372 }
373 DriftLevel::None => {
374 self.warning_streak = 0;
375 self.drift_streak = 0;
376 self.clear_streak = checked_increment(self.clear_streak, "clear streak")?;
377 consecutive_windows = self.clear_streak;
378 if self.active_level != DriftLevel::None
379 && self.clear_streak >= self.config.clear_windows
380 {
381 self.active_level = DriftLevel::None;
382 self.clear_streak = 0;
383 }
384 }
385 }
386
387 Ok(self.result(
388 warning_votes,
389 drift_votes,
390 available_detectors,
391 triggered_detectors,
392 consecutive_windows,
393 false,
394 false,
395 generation_changed,
396 new_event,
397 merged_event,
398 ))
399 }
400
401 #[allow(clippy::too_many_arguments)]
402 fn result(
403 &self,
404 warning_votes: usize,
405 drift_votes: usize,
406 available_detectors: usize,
407 triggered_detectors: Vec<String>,
408 consecutive_windows: u64,
409 warming: bool,
410 incomplete: bool,
411 generation_changed: bool,
412 new_event: bool,
413 merged_event: bool,
414 ) -> DriftConsensusResult {
415 DriftConsensusResult {
416 level: self.active_level,
417 warning_votes,
418 drift_votes,
419 available_detectors,
420 triggered_detectors,
421 consecutive_windows,
422 clear_streak: self.clear_streak,
423 cooldown_remaining: self.cooldown_remaining,
424 warming,
425 incomplete,
426 generation_changed,
427 new_event,
428 merged_event,
429 event: self.last_event.clone(),
430 }
431 }
432
433 fn validate_votes(&self, votes: &[DriftVote]) -> Result<(), RillError> {
434 if votes.len() > self.config.max_detectors {
435 return Err(RillError::InvalidCapacity(votes.len()));
436 }
437 let mut names = BTreeSet::new();
438 for vote in votes {
439 vote.validate()?;
440 if !names.insert(vote.detector.as_str()) {
441 return Err(RillError::InvalidState(format!(
442 "duplicate drift detector vote: {}",
443 vote.detector
444 )));
445 }
446 }
447 Ok(())
448 }
449
450 fn record_event(
451 &mut self,
452 level: DriftLevel,
453 detectors: &[String],
454 generation: u64,
455 ) -> Result<(bool, bool), RillError> {
456 let can_merge = self.last_event.as_ref().is_some_and(|event| {
457 event.generation == generation
458 && self.windows_seen.saturating_sub(event.last_window)
459 <= self.config.event_merge_horizon
460 });
461 if can_merge && let Some(event) = self.last_event.as_mut() {
462 event.last_window = self.windows_seen;
463 event.occurrences = checked_increment(event.occurrences, "event occurrences")?;
464 if level_rank(level) > level_rank(event.peak_level) {
465 event.peak_level = level;
466 }
467 let mut merged: BTreeSet<String> = event.detectors.iter().cloned().collect();
468 merged.extend(detectors.iter().cloned());
469 event.detectors = merged.into_iter().take(self.config.max_detectors).collect();
470 self.cooldown_remaining = self.config.cooldown_windows;
471 return Ok((false, true));
472 }
473
474 let event_id = self.next_event_id;
475 self.next_event_id = checked_increment(self.next_event_id, "next_event_id")?;
476 let unique: BTreeSet<String> = detectors.iter().cloned().collect();
477 self.last_event = Some(DriftEventSummary {
478 event_id,
479 first_window: self.windows_seen,
480 last_window: self.windows_seen,
481 peak_level: level,
482 detectors: unique.into_iter().take(self.config.max_detectors).collect(),
483 occurrences: 1,
484 generation,
485 });
486 self.cooldown_remaining = self.config.cooldown_windows;
487 Ok((true, false))
488 }
489}
490
491impl ValidateState for DriftConsensus {
492 fn validate_state(&self) -> Result<(), RillError> {
493 self.config.validate()?;
494 if self.generation_windows > self.windows_seen
495 || self.warning_streak > self.generation_windows
496 || self.drift_streak > self.generation_windows
497 || self.clear_streak > self.generation_windows
498 || self.cooldown_remaining > self.config.cooldown_windows
499 || self.next_event_id == 0
500 {
501 return Err(RillError::InvalidState(
502 "drift consensus counters are inconsistent".to_owned(),
503 ));
504 }
505 if self.active_generation.is_none() && self.generation_windows != 0 {
506 return Err(RillError::InvalidState(
507 "drift consensus generation counter has no generation".to_owned(),
508 ));
509 }
510 if let Some(event) = &self.last_event {
511 if event.event_id >= self.next_event_id
512 || event.first_window == 0
513 || event.first_window > event.last_window
514 || event.last_window > self.windows_seen
515 || event.occurrences == 0
516 || event.peak_level == DriftLevel::None
517 || event.detectors.len() > self.config.max_detectors
518 {
519 return Err(RillError::InvalidState(
520 "drift consensus event summary is inconsistent".to_owned(),
521 ));
522 }
523 let mut unique = BTreeSet::new();
524 for detector in &event.detectors {
525 DriftVote::new(detector.clone(), DriftLevel::None)?;
526 if !unique.insert(detector) {
527 return Err(RillError::InvalidState(
528 "drift consensus event has duplicate detector names".to_owned(),
529 ));
530 }
531 }
532 }
533 Ok(())
534 }
535}
536
537fn level_rank(level: DriftLevel) -> u8 {
538 match level {
539 DriftLevel::None => 0,
540 DriftLevel::Warning => 1,
541 DriftLevel::Drift => 2,
542 }
543}
544
545#[cfg(test)]
546mod tests {
547 use super::*;
548
549 fn config() -> DriftConsensusConfig {
550 DriftConsensusConfig {
551 minimum_warning_votes: 1,
552 minimum_drift_votes: 2,
553 confirmation_windows: 2,
554 clear_windows: 2,
555 cooldown_windows: 2,
556 event_merge_horizon: 10,
557 warming_windows: 0,
558 max_detectors: 3,
559 }
560 }
561
562 fn votes(levels: [DriftLevel; 3]) -> Vec<DriftVote> {
563 levels
564 .into_iter()
565 .enumerate()
566 .map(|(index, level)| DriftVote::new(format!("d{index}"), level).unwrap())
567 .collect()
568 }
569
570 #[test]
571 fn one_two_and_three_of_three_votes_are_explainable() {
572 let mut consensus = DriftConsensus::new(config()).unwrap();
573 let one = consensus
574 .update(
575 &votes([DriftLevel::Drift, DriftLevel::None, DriftLevel::None]),
576 1,
577 )
578 .unwrap();
579 assert_eq!(one.warning_votes, 1);
580 assert_eq!(one.drift_votes, 1);
581 assert_eq!(one.level, DriftLevel::None);
582
583 let two_votes = votes([DriftLevel::Drift, DriftLevel::Drift, DriftLevel::None]);
584 let first = consensus.update(&two_votes, 1).unwrap();
585 assert_eq!(first.level, DriftLevel::None);
586 let second = consensus.update(&two_votes, 1).unwrap();
587 assert_eq!(second.level, DriftLevel::Drift);
588 assert!(second.new_event);
589
590 let three = consensus.update(&votes([DriftLevel::Drift; 3]), 1).unwrap();
591 assert_eq!(three.drift_votes, 3);
592 assert_eq!(three.level, DriftLevel::Drift);
593 }
594
595 #[test]
596 fn warning_escalates_to_drift_and_clear_uses_hysteresis() {
597 let mut consensus = DriftConsensus::new(config()).unwrap();
598 let warning_votes = votes([DriftLevel::Warning, DriftLevel::None, DriftLevel::None]);
599 consensus.update(&warning_votes, 1).unwrap();
600 assert_eq!(
601 consensus.update(&warning_votes, 1).unwrap().level,
602 DriftLevel::Warning
603 );
604
605 let drift_votes = votes([DriftLevel::Drift, DriftLevel::Drift, DriftLevel::None]);
606 consensus.update(&drift_votes, 1).unwrap();
607 assert_eq!(
608 consensus.update(&drift_votes, 1).unwrap().level,
609 DriftLevel::Drift
610 );
611
612 let stable = votes([DriftLevel::None; 3]);
613 assert_eq!(
614 consensus.update(&stable, 1).unwrap().level,
615 DriftLevel::Drift
616 );
617 assert_eq!(
618 consensus.update(&stable, 1).unwrap().level,
619 DriftLevel::None
620 );
621 }
622
623 #[test]
624 fn interrupted_confirmation_and_incomplete_data_do_not_trigger_or_clear() {
625 let mut consensus = DriftConsensus::new(config()).unwrap();
626 let drift_votes = votes([DriftLevel::Drift, DriftLevel::Drift, DriftLevel::None]);
627 consensus.update(&drift_votes, 1).unwrap();
628 consensus.update(&votes([DriftLevel::None; 3]), 1).unwrap();
629 assert_eq!(
630 consensus.update(&drift_votes, 1).unwrap().level,
631 DriftLevel::None
632 );
633
634 let unavailable = vec![
635 DriftVote::unavailable("d0").unwrap(),
636 DriftVote::unavailable("d1").unwrap(),
637 ];
638 let result = consensus.update(&unavailable, 1).unwrap();
639 assert!(result.incomplete);
640 assert_eq!(result.level, DriftLevel::None);
641 }
642
643 #[test]
644 fn cooldown_and_merge_horizon_merge_repeated_events() {
645 let mut consensus = DriftConsensus::new(config()).unwrap();
646 let drift_votes = votes([DriftLevel::Drift, DriftLevel::Drift, DriftLevel::None]);
647 consensus.update(&drift_votes, 1).unwrap();
648 let first = consensus.update(&drift_votes, 1).unwrap();
649 assert!(first.new_event);
650 assert_eq!(first.event.as_ref().unwrap().occurrences, 1);
651
652 consensus.update(&drift_votes, 1).unwrap();
654 let suppressed = consensus.update(&drift_votes, 1).unwrap();
655 assert!(!suppressed.new_event && !suppressed.merged_event);
656
657 consensus.update(&drift_votes, 1).unwrap();
658 let merged = consensus.update(&drift_votes, 1).unwrap();
659 assert!(merged.merged_event);
660 assert_eq!(merged.event.unwrap().occurrences, 2);
661 }
662
663 #[test]
664 fn generation_change_resets_confirmation_context() {
665 let mut consensus = DriftConsensus::new(config()).unwrap();
666 let drift_votes = votes([DriftLevel::Drift, DriftLevel::Drift, DriftLevel::None]);
667 consensus.update(&drift_votes, 1).unwrap();
668 let changed = consensus.update(&drift_votes, 2).unwrap();
669 assert!(changed.generation_changed);
670 assert_eq!(changed.level, DriftLevel::None);
671 assert_eq!(
672 consensus.update(&drift_votes, 2).unwrap().level,
673 DriftLevel::Drift
674 );
675 }
676
677 #[test]
678 fn serde_restore_preserves_continuity_and_validation_rejects_corruption() {
679 let mut consensus = DriftConsensus::new(config()).unwrap();
680 consensus
681 .update(
682 &votes([DriftLevel::Warning, DriftLevel::None, DriftLevel::None]),
683 7,
684 )
685 .unwrap();
686 #[cfg(feature = "serde")]
687 {
688 let json = serde_json::to_string(&consensus).unwrap();
689 let mut restored: DriftConsensus = serde_json::from_str(&json).unwrap();
690 restored.validate_state().unwrap();
691 let next = votes([DriftLevel::Warning, DriftLevel::None, DriftLevel::None]);
692 assert_eq!(
693 consensus.update(&next, 7).unwrap(),
694 restored.update(&next, 7).unwrap()
695 );
696 }
697
698 let mut corrupt = consensus.clone();
699 corrupt.cooldown_remaining = corrupt.config.cooldown_windows + 1;
700 assert!(corrupt.validate_state().is_err());
701 }
702
703 #[test]
704 fn duplicate_votes_and_capacity_are_rejected_atomically() {
705 let mut consensus = DriftConsensus::new(config()).unwrap();
706 let duplicate = vec![
707 DriftVote::new("same", DriftLevel::None).unwrap(),
708 DriftVote::new("same", DriftLevel::Drift).unwrap(),
709 ];
710 assert!(consensus.update(&duplicate, 1).is_err());
711 assert_eq!(consensus.windows_seen(), 0);
712
713 let too_many = (0..4)
714 .map(|i| DriftVote::new(format!("d{i}"), DriftLevel::None).unwrap())
715 .collect::<Vec<_>>();
716 assert!(consensus.update(&too_many, 1).is_err());
717 assert_eq!(consensus.windows_seen(), 0);
718 }
719
720 #[test]
721 fn counter_overflow_and_corrupt_event_state_are_failure_atomic() {
722 let warning_votes = votes([DriftLevel::Warning, DriftLevel::None, DriftLevel::None]);
723 let mut consensus = DriftConsensus::new(config()).unwrap();
724 consensus.active_generation = Some(1);
725 consensus.windows_seen = 10;
726 consensus.generation_windows = 10;
727 consensus.warning_streak = u64::MAX;
728 let before = consensus.clone();
729 assert!(consensus.update(&warning_votes, 1).is_err());
730 assert_eq!(consensus, before);
731
732 let drift_votes = votes([DriftLevel::Drift, DriftLevel::Drift, DriftLevel::None]);
733 let mut corrupt = DriftConsensus::new(config()).unwrap();
734 corrupt.active_generation = Some(1);
735 corrupt.windows_seen = 1;
736 corrupt.generation_windows = 1;
737 corrupt.drift_streak = 1;
738 corrupt.next_event_id = u64::MAX;
739 let before = corrupt.clone();
740 assert!(corrupt.update(&drift_votes, 1).is_err());
741 assert_eq!(corrupt, before);
742 }
743}