Skip to main content

rill_ml/drift/
consensus.rs

1//! Bounded, explainable consensus over multiple drift detectors.
2//!
3//! `DriftConsensus` combines caller-supplied detector levels. It never resets
4//! a model or performs a business action; it only reports a hysteretic,
5//! auditable consensus level.
6
7use 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/// One detector's vote for a consensus window.
16#[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    /// Stable caller-defined detector identifier.
22    pub detector: String,
23    /// Detector level for this window.
24    pub level: DriftLevel,
25    /// Whether this detector had enough data to cast a meaningful vote.
26    pub available: bool,
27}
28
29impl DriftVote {
30    /// Construct an available vote.
31    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    /// Construct a temporarily unavailable detector vote.
42    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/// Configuration for [`DriftConsensus`].
68#[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    /// Warning-or-drift votes required for a warning candidate.
74    pub minimum_warning_votes: usize,
75    /// Drift votes required for a confirmed drift candidate.
76    pub minimum_drift_votes: usize,
77    /// Consecutive candidate windows required to change level.
78    pub confirmation_windows: u64,
79    /// Consecutive complete, stable windows required to clear a level.
80    pub clear_windows: u64,
81    /// Windows after an event during which another event is suppressed.
82    pub cooldown_windows: u64,
83    /// Maximum distance at which a repeated event is merged into the last one.
84    pub event_merge_horizon: u64,
85    /// Initial complete windows reported as warming.
86    pub warming_windows: u64,
87    /// Maximum votes and retained detector identifiers.
88    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    /// Validate thresholds and capacity bounds.
108    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/// Bounded summary of a consensus event.
132#[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    /// Monotonic event identifier local to this consensus instance.
137    pub event_id: u64,
138    /// First consensus window included in this event.
139    pub first_window: u64,
140    /// Most recent merged consensus window.
141    pub last_window: u64,
142    /// Highest level observed across merged occurrences.
143    pub peak_level: DriftLevel,
144    /// Sorted, deduplicated detector identifiers involved in the event.
145    pub detectors: Vec<String>,
146    /// Number of merged occurrences.
147    pub occurrences: u64,
148    /// Caller-supplied configuration/model generation context.
149    pub generation: u64,
150}
151
152/// Explainable result for one consensus window.
153#[derive(Debug, Clone, PartialEq, Eq)]
154pub struct DriftConsensusResult {
155    /// Hysteretic consensus level after this window.
156    pub level: DriftLevel,
157    /// Available warning-or-drift votes.
158    pub warning_votes: usize,
159    /// Available drift votes.
160    pub drift_votes: usize,
161    /// Number of available detector votes.
162    pub available_detectors: usize,
163    /// Detectors currently voting Warning or Drift.
164    pub triggered_detectors: Vec<String>,
165    /// Consecutive windows for the current raw candidate.
166    pub consecutive_windows: u64,
167    /// Current stable-window streak used for clearing.
168    pub clear_streak: u64,
169    /// Remaining event cooldown windows.
170    pub cooldown_remaining: u64,
171    /// Whether the consensus is still warming.
172    pub warming: bool,
173    /// Whether insufficient detector data prevented a state transition.
174    pub incomplete: bool,
175    /// Whether this call observed a generation change.
176    pub generation_changed: bool,
177    /// Whether a new (unmerged) event was created.
178    pub new_event: bool,
179    /// Whether a repeated occurrence was merged into the prior event.
180    pub merged_event: bool,
181    /// Latest event summary, if any.
182    pub event: Option<DriftEventSummary>,
183}
184
185/// Stateful drift vote consensus with hysteresis, cooldown and event merging.
186#[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    /// Create an empty consensus state.
204    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    /// Configuration used by this consensus instance.
222    pub const fn config(&self) -> &DriftConsensusConfig {
223        &self.config
224    }
225
226    /// Number of vote windows processed.
227    pub const fn windows_seen(&self) -> u64 {
228        self.windows_seen
229    }
230
231    /// Current hysteretic consensus level.
232    pub const fn level(&self) -> DriftLevel {
233        self.active_level
234    }
235
236    /// Latest bounded event summary.
237    pub const fn last_event(&self) -> Option<&DriftEventSummary> {
238        self.last_event.as_ref()
239    }
240
241    /// Reset all dynamic state while retaining configuration.
242    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    /// Incorporate one complete or partial detector-vote window.
256    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        // Cooldown suppresses the next confirmation boundary.
653        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}