1use std::fmt;
2use std::sync::Arc;
3use std::sync::atomic::{AtomicU16, AtomicU64, Ordering};
4
5pub trait Sequencer: Send + Sync + fmt::Debug {
7 fn next_sequence_number(&self) -> u16;
9 fn roll_over_count(&self) -> u64;
13 fn clone_to(&self) -> Box<dyn Sequencer>;
15}
16
17impl Clone for Box<dyn Sequencer> {
18 fn clone(&self) -> Box<dyn Sequencer> {
19 self.clone_to()
20 }
21}
22
23pub fn new_random_sequencer() -> impl Sequencer {
26 let c = Counters {
27 sequence_number: Arc::new(AtomicU16::new(rand::random::<u16>())),
28 roll_over_count: Arc::new(AtomicU64::new(0)),
29 };
30 SequencerImpl(c)
31}
32
33pub fn new_fixed_sequencer(s: u16) -> impl Sequencer {
36 let sequence_number = if s == 0 { u16::MAX } else { s - 1 };
37
38 let c = Counters {
39 sequence_number: Arc::new(AtomicU16::new(sequence_number)),
40 roll_over_count: Arc::new(AtomicU64::new(0)),
41 };
42
43 SequencerImpl(c)
44}
45
46#[derive(Debug, Clone)]
47struct SequencerImpl(Counters);
48
49#[derive(Debug, Clone)]
50struct Counters {
51 sequence_number: Arc<AtomicU16>,
52 roll_over_count: Arc<AtomicU64>,
53}
54
55impl Sequencer for SequencerImpl {
56 fn next_sequence_number(&self) -> u16 {
59 if self.0.sequence_number.load(Ordering::SeqCst) == u16::MAX {
60 self.0.roll_over_count.fetch_add(1, Ordering::SeqCst);
61 self.0.sequence_number.store(0, Ordering::SeqCst);
62 0
63 } else {
64 self.0.sequence_number.fetch_add(1, Ordering::SeqCst) + 1
65 }
66 }
67
68 fn roll_over_count(&self) -> u64 {
71 self.0.roll_over_count.load(Ordering::SeqCst)
72 }
73
74 fn clone_to(&self) -> Box<dyn Sequencer> {
75 Box::new(self.clone())
76 }
77}