fxcp_core/consistency/
sequencer.rs1use std::sync::atomic::{AtomicU64, Ordering};
7use tokio::sync::watch;
8use std::collections::BTreeSet;
9use std::sync::Mutex;
10use tracing::warn;
11
12#[derive(Debug)]
14pub struct GlobalSequencer {
15 counter: AtomicU64,
16}
17
18impl GlobalSequencer {
19 pub fn new(start: u64) -> Self {
20 Self { counter: AtomicU64::new(start) }
21 }
22
23 pub fn next(&self) -> u64 {
25 self.counter.fetch_add(1, Ordering::SeqCst) + 1
26 }
27
28 pub fn current(&self) -> u64 {
29 self.counter.load(Ordering::SeqCst)
30 }
31}
32
33#[derive(Debug)]
37pub struct SequenceBarrier {
38 finished: watch::Sender<u64>,
39 receiver: watch::Receiver<u64>,
40 pending_completions: Mutex<BTreeSet<u64>>,
41}
42
43impl SequenceBarrier {
44 pub fn new(start: u64) -> Self {
45 let (tx, rx) = watch::channel(start);
46 Self {
47 finished: tx,
48 receiver: rx,
49 pending_completions: Mutex::new(BTreeSet::new()),
50 }
51 }
52
53 pub async fn wait_for(&self, seq: u64) {
55 let mut rx = self.receiver.clone();
56 loop {
57 let current = *rx.borrow_and_update();
59 if current >= seq {
60 return;
61 }
62 if rx.changed().await.is_err() {
64 return;
67 }
68 }
69 }
70
71 pub fn complete(&self, seq: u64) {
74 let mut pending = self.pending_completions.lock().unwrap_or_else(|poisoned| {
75 warn!("Mutex poisoned in SequenceBarrier::complete, recovering: {}", poisoned);
76 poisoned.into_inner()
77 });
78
79 let current = *self.finished.borrow();
80 if seq <= current {
81 return; }
83
84 pending.insert(seq);
85
86 let mut next_expected = current + 1;
87 let mut updated = false;
88
89 while pending.remove(&next_expected) {
91 next_expected += 1;
92 updated = true;
93 }
94
95 if updated {
96 let _ = self.finished.send(next_expected - 1);
98 }
99 }
100
101 pub fn current(&self) -> u64 {
102 *self.finished.borrow()
103 }
104}