Skip to main content

ironfix_session/
sequence.rs

1/******************************************************************************
2   Author: Joaquín Béjar García
3   Email: jb@taunais.com
4   Date: 27/1/26
5******************************************************************************/
6
7//! Sequence number management.
8//!
9//! This module provides atomic sequence number management for FIX sessions.
10
11use ironfix_core::types::SeqNum;
12use std::sync::atomic::{AtomicU64, Ordering};
13
14/// Error returned when a sequence counter has reached its maximum value.
15///
16/// FIX sequence numbers are unbounded in the specification, but this
17/// implementation stores them as `u64`. Once a counter reaches `u64::MAX`
18/// no further numbers can be allocated: the session must perform a
19/// sequence reset (Logon with `ResetSeqNumFlag(141)=Y`, or an out-of-band
20/// reset agreed with the counterparty) and then call
21/// [`SequenceManager::reset`] before continuing.
22#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
23#[error(
24    "sequence counter exhausted: {counter} reached u64::MAX, session requires a sequence reset"
25)]
26pub struct SequenceExhausted {
27    /// Which counter was exhausted.
28    pub counter: SequenceCounter,
29}
30
31/// Identifies one of the two sequence counters of a session.
32#[derive(Debug, Clone, Copy, PartialEq, Eq)]
33pub enum SequenceCounter {
34    /// Outgoing (sender) sequence counter.
35    Sender,
36    /// Incoming (target) sequence counter.
37    Target,
38}
39
40impl std::fmt::Display for SequenceCounter {
41    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
42        match self {
43            Self::Sender => write!(f, "sender"),
44            Self::Target => write!(f, "target"),
45        }
46    }
47}
48
49/// Manages sequence numbers for a FIX session.
50///
51/// Uses atomic operations for thread-safe access without locks.
52#[derive(Debug)]
53pub struct SequenceManager {
54    /// Next outgoing sequence number.
55    next_sender_seq: AtomicU64,
56    /// Next expected incoming sequence number.
57    next_target_seq: AtomicU64,
58}
59
60impl SequenceManager {
61    /// Creates a new sequence manager with sequence numbers starting at 1.
62    #[must_use]
63    pub fn new() -> Self {
64        Self {
65            next_sender_seq: AtomicU64::new(1),
66            next_target_seq: AtomicU64::new(1),
67        }
68    }
69
70    /// Creates a new sequence manager with specified starting values.
71    ///
72    /// # Arguments
73    /// * `sender_seq` - Initial sender sequence number
74    /// * `target_seq` - Initial target sequence number
75    #[must_use]
76    pub fn with_initial(sender_seq: u64, target_seq: u64) -> Self {
77        Self {
78            next_sender_seq: AtomicU64::new(sender_seq),
79            next_target_seq: AtomicU64::new(target_seq),
80        }
81    }
82
83    /// Returns the next sender sequence number without incrementing.
84    #[inline]
85    #[must_use]
86    pub fn next_sender_seq(&self) -> SeqNum {
87        SeqNum::new(self.next_sender_seq.load(Ordering::SeqCst))
88    }
89
90    /// Returns the next target sequence number without incrementing.
91    #[inline]
92    #[must_use]
93    pub fn next_target_seq(&self) -> SeqNum {
94        SeqNum::new(self.next_target_seq.load(Ordering::SeqCst))
95    }
96
97    /// Allocates and returns the next sender sequence number.
98    ///
99    /// This atomically increments the sequence number and returns the
100    /// value before the increment.
101    ///
102    /// Note: wraps silently on `u64` overflow. Prefer
103    /// [`try_allocate_sender_seq`](Self::try_allocate_sender_seq) for
104    /// venue-grade sessions where exhaustion must be an explicit error.
105    #[inline]
106    pub fn allocate_sender_seq(&self) -> SeqNum {
107        SeqNum::new(self.next_sender_seq.fetch_add(1, Ordering::SeqCst))
108    }
109
110    /// Allocates and returns the next sender sequence number, failing
111    /// instead of wrapping when the counter is exhausted.
112    ///
113    /// On success this atomically increments the counter and returns the
114    /// value before the increment. On exhaustion the counter is left
115    /// untouched; the session must perform a sequence reset (see
116    /// [`SequenceExhausted`]) before more numbers can be allocated.
117    ///
118    /// # Errors
119    /// Returns [`SequenceExhausted`] if the counter has reached `u64::MAX`.
120    #[inline]
121    pub fn try_allocate_sender_seq(&self) -> Result<SeqNum, SequenceExhausted> {
122        self.next_sender_seq
123            .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |current| {
124                current.checked_add(1)
125            })
126            .map(SeqNum::new)
127            .map_err(|_| SequenceExhausted {
128                counter: SequenceCounter::Sender,
129            })
130    }
131
132    /// Increments the target sequence number.
133    ///
134    /// Call this after successfully processing an incoming message.
135    ///
136    /// Note: wraps silently on `u64` overflow. Prefer
137    /// [`try_increment_target_seq`](Self::try_increment_target_seq) for
138    /// venue-grade sessions where exhaustion must be an explicit error.
139    #[inline]
140    pub fn increment_target_seq(&self) {
141        self.next_target_seq.fetch_add(1, Ordering::SeqCst);
142    }
143
144    /// Increments the target sequence number, failing instead of wrapping
145    /// when the counter is exhausted.
146    ///
147    /// On success returns the new next expected target sequence number.
148    /// On exhaustion the counter is left untouched; the session must
149    /// perform a sequence reset (see [`SequenceExhausted`]) before more
150    /// messages can be accepted.
151    ///
152    /// # Errors
153    /// Returns [`SequenceExhausted`] if the counter has reached `u64::MAX`.
154    #[inline]
155    pub fn try_increment_target_seq(&self) -> Result<SeqNum, SequenceExhausted> {
156        self.next_target_seq
157            .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |current| {
158                current.checked_add(1)
159            })
160            .map(|previous| SeqNum::new(previous + 1))
161            .map_err(|_| SequenceExhausted {
162                counter: SequenceCounter::Target,
163            })
164    }
165
166    /// Sets the next sender sequence number.
167    ///
168    /// # Arguments
169    /// * `seq` - The new sequence number
170    #[inline]
171    pub fn set_sender_seq(&self, seq: u64) {
172        self.next_sender_seq.store(seq, Ordering::SeqCst);
173    }
174
175    /// Sets the next target sequence number.
176    ///
177    /// # Arguments
178    /// * `seq` - The new sequence number
179    #[inline]
180    pub fn set_target_seq(&self, seq: u64) {
181        self.next_target_seq.store(seq, Ordering::SeqCst);
182    }
183
184    /// Resets both sequence numbers to 1.
185    #[inline]
186    pub fn reset(&self) {
187        self.next_sender_seq.store(1, Ordering::SeqCst);
188        self.next_target_seq.store(1, Ordering::SeqCst);
189    }
190
191    /// Validates an incoming sequence number.
192    ///
193    /// # Arguments
194    /// * `received` - The received sequence number
195    ///
196    /// # Returns
197    /// - `Ok(())` if the sequence number matches expected
198    /// - `Err(SequenceResult::TooLow)` if it's a possible duplicate
199    /// - `Err(SequenceResult::Gap)` if there's a gap
200    #[must_use]
201    pub fn validate_incoming(&self, received: u64) -> SequenceResult {
202        let expected = self.next_target_seq.load(Ordering::SeqCst);
203
204        if received == expected {
205            SequenceResult::Ok
206        } else if received < expected {
207            SequenceResult::TooLow { expected, received }
208        } else {
209            SequenceResult::Gap { expected, received }
210        }
211    }
212}
213
214impl Default for SequenceManager {
215    fn default() -> Self {
216        Self::new()
217    }
218}
219
220/// Result of sequence number validation.
221#[derive(Debug, Clone, Copy, PartialEq, Eq)]
222pub enum SequenceResult {
223    /// Sequence number is as expected.
224    Ok,
225    /// Sequence number is lower than expected (possible duplicate).
226    TooLow {
227        /// Expected sequence number.
228        expected: u64,
229        /// Received sequence number.
230        received: u64,
231    },
232    /// Sequence number is higher than expected (gap detected).
233    Gap {
234        /// Expected sequence number.
235        expected: u64,
236        /// Received sequence number.
237        received: u64,
238    },
239}
240
241impl SequenceResult {
242    /// Returns true if the sequence is valid.
243    #[must_use]
244    pub const fn is_ok(&self) -> bool {
245        matches!(self, Self::Ok)
246    }
247
248    /// Returns true if there's a gap.
249    #[must_use]
250    pub const fn is_gap(&self) -> bool {
251        matches!(self, Self::Gap { .. })
252    }
253
254    /// Returns true if the sequence is too low.
255    #[must_use]
256    pub const fn is_too_low(&self) -> bool {
257        matches!(self, Self::TooLow { .. })
258    }
259}
260
261#[cfg(test)]
262mod tests {
263    use super::*;
264
265    #[test]
266    fn test_sequence_manager_new() {
267        let mgr = SequenceManager::new();
268        assert_eq!(mgr.next_sender_seq().value(), 1);
269        assert_eq!(mgr.next_target_seq().value(), 1);
270    }
271
272    #[test]
273    fn test_allocate_sender_seq() {
274        let mgr = SequenceManager::new();
275
276        let seq1 = mgr.allocate_sender_seq();
277        assert_eq!(seq1.value(), 1);
278        assert_eq!(mgr.next_sender_seq().value(), 2);
279
280        let seq2 = mgr.allocate_sender_seq();
281        assert_eq!(seq2.value(), 2);
282        assert_eq!(mgr.next_sender_seq().value(), 3);
283    }
284
285    #[test]
286    fn test_increment_target_seq() {
287        let mgr = SequenceManager::new();
288
289        mgr.increment_target_seq();
290        assert_eq!(mgr.next_target_seq().value(), 2);
291
292        mgr.increment_target_seq();
293        assert_eq!(mgr.next_target_seq().value(), 3);
294    }
295
296    #[test]
297    fn test_validate_incoming() {
298        let mgr = SequenceManager::new();
299
300        assert!(mgr.validate_incoming(1).is_ok());
301
302        mgr.set_target_seq(5);
303        assert!(mgr.validate_incoming(4).is_too_low());
304        assert!(mgr.validate_incoming(5).is_ok());
305        assert!(mgr.validate_incoming(10).is_gap());
306    }
307
308    #[test]
309    fn test_try_allocate_sender_seq() {
310        let mgr = SequenceManager::new();
311
312        assert_eq!(mgr.try_allocate_sender_seq().unwrap().value(), 1);
313        assert_eq!(mgr.try_allocate_sender_seq().unwrap().value(), 2);
314        assert_eq!(mgr.next_sender_seq().value(), 3);
315    }
316
317    #[test]
318    fn test_try_allocate_sender_seq_exhausted() {
319        let mgr = SequenceManager::with_initial(u64::MAX, 1);
320
321        let err = mgr.try_allocate_sender_seq().unwrap_err();
322        assert_eq!(err.counter, SequenceCounter::Sender);
323        // Counter untouched: still exhausted, no wraparound.
324        assert_eq!(mgr.next_sender_seq().value(), u64::MAX);
325        assert!(mgr.try_allocate_sender_seq().is_err());
326
327        // Reset restores a usable session.
328        mgr.reset();
329        assert_eq!(mgr.try_allocate_sender_seq().unwrap().value(), 1);
330    }
331
332    #[test]
333    fn test_try_increment_target_seq() {
334        let mgr = SequenceManager::new();
335
336        assert_eq!(mgr.try_increment_target_seq().unwrap().value(), 2);
337        assert_eq!(mgr.try_increment_target_seq().unwrap().value(), 3);
338        assert_eq!(mgr.next_target_seq().value(), 3);
339    }
340
341    #[test]
342    fn test_try_increment_target_seq_exhausted() {
343        let mgr = SequenceManager::with_initial(1, u64::MAX);
344
345        let err = mgr.try_increment_target_seq().unwrap_err();
346        assert_eq!(err.counter, SequenceCounter::Target);
347        assert_eq!(mgr.next_target_seq().value(), u64::MAX);
348        assert!(mgr.try_increment_target_seq().is_err());
349    }
350
351    #[test]
352    fn test_reset() {
353        let mgr = SequenceManager::with_initial(100, 200);
354        assert_eq!(mgr.next_sender_seq().value(), 100);
355        assert_eq!(mgr.next_target_seq().value(), 200);
356
357        mgr.reset();
358        assert_eq!(mgr.next_sender_seq().value(), 1);
359        assert_eq!(mgr.next_target_seq().value(), 1);
360    }
361}