Skip to main content

fxcp_core/consistency/
sequencer.rs

1// SPDX-License-Identifier: GPL-2.0-or-later
2// Copyright (C) 2025 Joel Wirāmu Pauling <aenertia@aenertia.net>
3//
4//! Monotonic sequence generator for event ordering.
5
6use std::sync::atomic::{AtomicU64, Ordering};
7use tokio::sync::watch;
8use std::collections::BTreeSet;
9use std::sync::Mutex;
10use tracing::warn;
11
12/// A simple monotonic counter for issuing tickets.
13#[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    /// Increments and returns the *next* sequence number (1-based if started at 0).
24    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/// A barrier that waits for a specific sequence number to be reached.
34/// It handles out-of-order completions by buffering them and advancing
35/// the publicly visible "finished" sequence monotonically.
36#[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    /// Wait until the barrier reaches at least `seq`.
54    pub async fn wait_for(&self, seq: u64) {
55        let mut rx = self.receiver.clone();
56        loop {
57            // Check current value first
58            let current = *rx.borrow_and_update();
59            if current >= seq {
60                return;
61            }
62            // Wait for change
63            if rx.changed().await.is_err() {
64                // Sender dropped, we can't wait anymore. 
65                // In a production system this might mean shutdown.
66                return; 
67            }
68        }
69    }
70
71    /// Mark a specific sequence number as complete.
72    /// If this completes the next expected sequence, the barrier advances.
73    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; // Already processed
82        }
83
84        pending.insert(seq);
85        
86        let mut next_expected = current + 1;
87        let mut updated = false;
88        
89        // Advance monotonically as far as possible
90        while pending.remove(&next_expected) {
91            next_expected += 1;
92            updated = true;
93        }
94
95        if updated {
96            // next_expected is now 1 greater than the last processed item
97            let _ = self.finished.send(next_expected - 1);
98        }
99    }
100    
101    pub fn current(&self) -> u64 {
102        *self.finished.borrow()
103    }
104}