noq_proto/
cid_queue.rs

1use std::{fmt::Debug, ops::Range};
2
3use crate::{ConnectionId, ResetToken, frame::NewConnectionId};
4
5/// DataType stored in CidQueue buffer
6#[derive(Debug, Clone, Copy)]
7struct CidData(ConnectionId, Option<ResetToken>);
8
9/// Sliding window of active Connection IDs.
10///
11/// This represents a circular buffer that can contain gaps due to packet loss or reordering.
12/// The buffer has three regions:
13/// - Exactly one active CID at `self.buffer[self.cursor]`.
14/// - Zero to `Self::LEN - 1` reserved CIDs from `self.cursor` up to `self.cursor_reserved`.
15/// - More "available"/"ready" CIDs after `self.cursor_reserved`.
16///
17/// The range of reserved CIDs is grown by calling `CidQueue::next_reserved`, which takes one of
18/// the available ones and returns the CID that was reserved.
19///
20/// New available/ready CIDs are added by calling [`CidQueue::insert`].
21///
22/// May contain gaps due to packet loss or reordering.
23#[derive(Debug)]
24pub(crate) struct CidQueue {
25    /// Ring buffer indexed by `self.cursor`
26    buffer: [Option<CidData>; Self::LEN],
27    /// Index at which circular buffer addressing is based
28    cursor: usize,
29    /// Sequence number of `self.buffer[cursor]`
30    ///
31    /// The sequence number of the active CID; must be the smallest among CIDs in `buffer`.
32    offset: u64,
33    /// Circular index for the last reserved CID, i.e. a CID that is not the active CID, but was
34    /// used for probing packets on a different remote address.
35    ///
36    /// When [`Self::cursor_reserved`] and [`Self::cursor`] are equal, no CID is considered
37    /// reserved.
38    ///
39    /// The reserved CIDs section of the buffer, if non-empty will always be ahead of the active
40    /// CID.
41    cursor_reserved: usize,
42}
43
44impl CidQueue {
45    pub(crate) fn new(cid: ConnectionId) -> Self {
46        let mut buffer = [None; Self::LEN];
47        buffer[0] = Some(CidData(cid, None));
48        Self {
49            buffer,
50            cursor: 0,
51            offset: 0,
52            cursor_reserved: 0,
53        }
54    }
55
56    /// Handle a `NEW_CONNECTION_ID` frame
57    ///
58    /// Returns a non-empty range of retired sequence numbers and the reset token of the new active
59    /// CID iff any CIDs were retired.
60    pub(crate) fn insert(
61        &mut self,
62        cid: NewConnectionId,
63    ) -> Result<Option<(Range<u64>, ResetToken)>, InsertError> {
64        // Position of new CID wrt. the current active CID
65        let Some(index) = cid.sequence.checked_sub(self.offset) else {
66            return Err(InsertError::Retired);
67        };
68
69        let retired_count = cid.retire_prior_to.saturating_sub(self.offset);
70        if index >= Self::LEN as u64 + retired_count {
71            return Err(InsertError::ExceedsLimit);
72        }
73
74        // Discard retired CIDs, if any
75        for i in 0..(retired_count.min(Self::LEN as u64) as usize) {
76            self.buffer[(self.cursor + i) % Self::LEN] = None;
77        }
78
79        // Record the new CID
80        let index = ((self.cursor as u64 + index) % Self::LEN as u64) as usize;
81        self.buffer[index] = Some(CidData(cid.id, Some(cid.reset_token)));
82
83        if retired_count == 0 {
84            return Ok(None);
85        }
86
87        // The active CID was retired. Find the first known CID with sequence number of at least
88        // retire_prior_to, and inform the caller that all prior CIDs have been retired, and of
89        // the new CID's reset token.
90        self.cursor = ((self.cursor as u64 + retired_count) % Self::LEN as u64) as usize;
91        let (i, CidData(_, token)) = self
92            .iter_from_active()
93            .next()
94            .expect("it is impossible to retire a CID without supplying a new one");
95        self.cursor = (self.cursor + i) % Self::LEN;
96        self.cursor_reserved = self.cursor;
97        let orig_offset = self.offset;
98        self.offset = cid.retire_prior_to + i as u64;
99        // We don't immediately retire CIDs in the range (orig_offset +
100        // Self::LEN)..self.offset. These are CIDs that we haven't yet received from a
101        // NEW_CONNECTION_ID frame, since having previously received them would violate the
102        // connection ID limit we specified based on Self::LEN. If we do receive a such a frame
103        // in the future, e.g. due to reordering, we'll retire it then. This ensures we can't be
104        // made to buffer an arbitrarily large number of RETIRE_CONNECTION_ID frames.
105        Ok(Some((
106            orig_offset..self.offset.min(orig_offset + Self::LEN as u64),
107            token.expect("non-initial CID missing reset token"),
108        )))
109    }
110
111    /// Switch to next active CID if possible, return
112    /// 1) the corresponding ResetToken and 2) a non-empty range preceding it to retire
113    pub(crate) fn next(&mut self) -> Option<(ResetToken, Range<u64>)> {
114        let (i, cid_data) = self.iter_from_reserved().nth(1)?;
115        let reserved = self.reserved_len();
116        for j in 0..=reserved {
117            self.buffer[(self.cursor + j) % Self::LEN] = None;
118        }
119        let orig_offset = self.offset;
120        self.offset += (i + reserved) as u64;
121        self.cursor = (self.cursor_reserved + i) % Self::LEN;
122        self.cursor_reserved = self.cursor;
123
124        Some((cid_data.1.unwrap(), orig_offset..self.offset))
125    }
126
127    /// Returns a CID from the available ones and marks it as reserved.
128    ///
129    /// If there's no more CIDs in the ready set, this will return None.
130    /// CIDs marked as reserved will be skipped when the active one advances.
131    #[allow(dead_code)]
132    pub(crate) fn next_reserved(&mut self) -> Option<ConnectionId> {
133        let (i, cid_data) = self.iter_from_reserved().nth(1)?;
134
135        self.cursor_reserved = (self.cursor_reserved + i) % Self::LEN;
136        Some(cid_data.0)
137    }
138
139    /// Returns the number of unused CIDs (neither active nor reserved).
140    #[allow(unused)]
141    pub(crate) fn remaining(&self) -> usize {
142        self.iter_from_reserved()
143            .count()
144            .checked_sub(1)
145            .expect("iterator is non empty")
146    }
147
148    /// Iterate CIDs in CidQueue that are not `None`, including the active CID
149    fn iter_from_active(&self) -> impl Iterator<Item = (usize, CidData)> + '_ {
150        (0..Self::LEN).filter_map(move |step| {
151            let index = (self.cursor + step) % Self::LEN;
152            self.buffer[index].map(|cid_data| (step, cid_data))
153        })
154    }
155
156    /// Iterate CIDs in CidQueue that are not `None`, from [`Self::cursor_reserved`].
157    ///
158    /// The iterator will always have at least one item, as it will include the active CID when no
159    /// CID has been reserved, or the last reserved CID otherwise.
160    ///
161    /// Along with the CID, it returns the offset counted from [`Self::cursor_reserved`] where the
162    /// CID is stored.
163    fn iter_from_reserved(&self) -> impl Iterator<Item = (usize, CidData)> + '_ {
164        (0..(Self::LEN - self.reserved_len())).filter_map(move |step| {
165            let index = (self.cursor_reserved + step) % Self::LEN;
166            self.buffer[index].map(|cid_data| (step, cid_data))
167        })
168    }
169
170    /// The length of the internal buffer's section of CIDs that are marked as reserved.
171    fn reserved_len(&self) -> usize {
172        if self.cursor_reserved >= self.cursor {
173            self.cursor_reserved - self.cursor
174        } else {
175            self.cursor_reserved + Self::LEN - self.cursor
176        }
177    }
178
179    /// Replace the initial CID
180    pub(crate) fn update_initial_cid(&mut self, cid: ConnectionId) {
181        debug_assert_eq!(self.offset, 0);
182        self.buffer[self.cursor] = Some(CidData(cid, None));
183    }
184
185    /// Return active remote CID itself
186    pub(crate) fn active(&self) -> ConnectionId {
187        self.buffer[self.cursor].unwrap().0
188    }
189
190    /// Return the sequence number of active remote CID
191    pub(crate) fn active_seq(&self) -> u64 {
192        self.offset
193    }
194
195    pub(crate) const LEN: usize = 5;
196}
197
198#[derive(Debug, Copy, Clone, Eq, PartialEq)]
199pub(crate) enum InsertError {
200    /// CID was already retired
201    Retired,
202    /// Sequence number violates the leading edge of the window
203    ExceedsLimit,
204}
205
206#[cfg(test)]
207mod tests {
208    use super::*;
209
210    fn cid(sequence: u64, retire_prior_to: u64) -> NewConnectionId {
211        NewConnectionId {
212            path_id: None,
213            sequence,
214            id: ConnectionId::new(&sequence.to_be_bytes()),
215            reset_token: ResetToken::from([0xCD; crate::RESET_TOKEN_SIZE]),
216            retire_prior_to,
217        }
218    }
219
220    fn initial_cid() -> ConnectionId {
221        ConnectionId::new(&[0xFF; 8])
222    }
223
224    #[test]
225    fn next_dense() {
226        let mut q = CidQueue::new(initial_cid());
227        assert!(q.next().is_none());
228        assert!(q.next().is_none());
229
230        for i in 1..CidQueue::LEN as u64 {
231            q.insert(cid(i, 0)).unwrap();
232        }
233        for i in 1..CidQueue::LEN as u64 {
234            let (_, retire) = q.next().unwrap();
235            assert_eq!(q.active_seq(), i);
236            assert_eq!(retire.end - retire.start, 1);
237        }
238        assert!(q.next().is_none());
239    }
240
241    #[test]
242    fn next_sparse() {
243        let mut q = CidQueue::new(initial_cid());
244        let seqs = (1..CidQueue::LEN as u64).filter(|x| x % 2 == 0);
245        for i in seqs.clone() {
246            q.insert(cid(i, 0)).unwrap();
247        }
248        for i in seqs {
249            let (_, retire) = q.next().unwrap();
250            dbg!(&retire);
251            assert_eq!(q.active_seq(), i);
252            assert_eq!(retire, (q.active_seq().saturating_sub(2))..q.active_seq());
253        }
254        assert!(q.next().is_none());
255    }
256
257    #[test]
258    fn wrap() {
259        let mut q = CidQueue::new(initial_cid());
260
261        for i in 1..CidQueue::LEN as u64 {
262            q.insert(cid(i, 0)).unwrap();
263        }
264        for _ in 1..(CidQueue::LEN as u64 - 1) {
265            q.next().unwrap();
266        }
267        for i in CidQueue::LEN as u64..(CidQueue::LEN as u64 + 3) {
268            q.insert(cid(i, 0)).unwrap();
269        }
270        for i in (CidQueue::LEN as u64 - 1)..(CidQueue::LEN as u64 + 3) {
271            q.next().unwrap();
272            assert_eq!(q.active_seq(), i);
273        }
274        assert!(q.next().is_none());
275    }
276
277    #[test]
278    fn retire_dense() {
279        let mut q = CidQueue::new(initial_cid());
280
281        for i in 1..CidQueue::LEN as u64 {
282            q.insert(cid(i, 0)).unwrap();
283        }
284        assert_eq!(q.active_seq(), 0);
285
286        assert_eq!(q.insert(cid(4, 2)).unwrap().unwrap().0, 0..2);
287        assert_eq!(q.active_seq(), 2);
288        assert_eq!(q.insert(cid(4, 2)), Ok(None));
289
290        for i in 2..(CidQueue::LEN as u64 - 1) {
291            let _ = q.next().unwrap();
292            assert_eq!(q.active_seq(), i + 1);
293            assert_eq!(q.insert(cid(i + 1, i + 1)), Ok(None));
294        }
295
296        assert!(q.next().is_none());
297    }
298
299    #[test]
300    fn retire_sparse() {
301        // Retiring CID 0 when CID 1 is not known should retire CID 1 as we move to CID 2
302        let mut q = CidQueue::new(initial_cid());
303        q.insert(cid(2, 0)).unwrap();
304        assert_eq!(q.insert(cid(3, 1)).unwrap().unwrap().0, 0..2,);
305        assert_eq!(q.active_seq(), 2);
306    }
307
308    #[test]
309    fn retire_many() {
310        let mut q = CidQueue::new(initial_cid());
311        q.insert(cid(2, 0)).unwrap();
312        assert_eq!(
313            q.insert(cid(1_000_000, 1_000_000)).unwrap().unwrap().0,
314            0..CidQueue::LEN as u64,
315        );
316        assert_eq!(q.active_seq(), 1_000_000);
317    }
318
319    #[test]
320    fn insert_limit() {
321        let mut q = CidQueue::new(initial_cid());
322        assert_eq!(q.insert(cid(CidQueue::LEN as u64 - 1, 0)), Ok(None));
323        assert_eq!(
324            q.insert(cid(CidQueue::LEN as u64, 0)),
325            Err(InsertError::ExceedsLimit)
326        );
327    }
328
329    #[test]
330    fn insert_duplicate() {
331        let mut q = CidQueue::new(initial_cid());
332        q.insert(cid(0, 0)).unwrap();
333        q.insert(cid(0, 0)).unwrap();
334    }
335
336    #[test]
337    fn insert_retired() {
338        let mut q = CidQueue::new(initial_cid());
339        assert_eq!(
340            q.insert(cid(0, 0)),
341            Ok(None),
342            "reinserting active CID succeeds"
343        );
344        assert!(q.next().is_none(), "active CID isn't requeued");
345        q.insert(cid(1, 0)).unwrap();
346        q.next().unwrap();
347        assert_eq!(
348            q.insert(cid(0, 0)),
349            Err(InsertError::Retired),
350            "previous active CID is already retired"
351        );
352    }
353
354    #[test]
355    fn retire_then_insert_next() {
356        let mut q = CidQueue::new(initial_cid());
357        for i in 1..CidQueue::LEN as u64 {
358            q.insert(cid(i, 0)).unwrap();
359        }
360        q.next().unwrap();
361        q.insert(cid(CidQueue::LEN as u64, 0)).unwrap();
362        assert_eq!(
363            q.insert(cid(CidQueue::LEN as u64 + 1, 0)),
364            Err(InsertError::ExceedsLimit)
365        );
366    }
367
368    #[test]
369    fn always_valid() {
370        let mut q = CidQueue::new(initial_cid());
371        assert!(q.next().is_none());
372        assert_eq!(q.active(), initial_cid());
373        assert_eq!(q.active_seq(), 0);
374    }
375
376    #[test]
377    fn reserved_smoke() {
378        let mut q = CidQueue::new(initial_cid());
379        assert_eq!(q.next_reserved(), None);
380
381        let one = cid(1, 0);
382        q.insert(one).unwrap();
383        assert_eq!(q.next_reserved(), Some(one.id));
384
385        let two = cid(2, 2);
386        let (retired_range, reset_token) = q.insert(two).unwrap().unwrap();
387        assert_eq!(reset_token, two.reset_token);
388        assert_eq!(retired_range, 0..2);
389
390        assert_eq!(q.next_reserved(), None);
391
392        let four = cid(4, 2);
393        q.insert(four).unwrap();
394        println!("{q:?}");
395        assert_eq!(q.next_reserved(), Some(four.id));
396        assert_eq!(q.active(), two.id);
397
398        assert_eq!(q.next(), None);
399    }
400
401    #[test]
402    fn reserve_multiple() {
403        let mut q = CidQueue::new(initial_cid());
404        let one = cid(1, 0);
405        let two = cid(2, 0);
406        q.insert(one).unwrap();
407        q.insert(two).unwrap();
408        assert_eq!(q.next_reserved(), Some(one.id));
409        assert_eq!(q.next_reserved(), Some(two.id));
410        assert_eq!(q.next_reserved(), None);
411    }
412
413    #[test]
414    fn reserve_multiple_sparse() {
415        let mut q = CidQueue::new(initial_cid());
416        let two = cid(2, 0);
417        let four = cid(4, 0);
418        q.insert(two).unwrap();
419        q.insert(four).unwrap();
420        assert_eq!(q.next_reserved(), Some(two.id));
421        assert_eq!(q.next_reserved(), Some(four.id));
422        assert_eq!(q.next_reserved(), None);
423    }
424
425    #[test]
426    fn reserve_many_next_clears() {
427        let mut q = CidQueue::new(initial_cid());
428        for i in 1..CidQueue::LEN {
429            q.insert(cid(i as u64, 0)).unwrap();
430        }
431
432        for _ in 0..CidQueue::LEN - 2 {
433            assert!(q.next_reserved().is_some());
434        }
435
436        assert!(q.next().is_some());
437        assert_eq!(q.next(), None);
438    }
439
440    #[test]
441    fn reserve_many_next_reserved_none() {
442        let mut q = CidQueue::new(initial_cid());
443        for i in 1..CidQueue::LEN {
444            q.insert(cid(i as u64, 0)).unwrap();
445        }
446
447        for _ in 0..CidQueue::LEN - 1 {
448            assert!(q.next_reserved().is_some());
449        }
450
451        assert_eq!(q.next_reserved(), None);
452    }
453
454    #[test]
455    fn one_active_all_else_reserved_next_none() {
456        let mut q = CidQueue::new(initial_cid());
457        for i in 1..CidQueue::LEN {
458            q.insert(cid(i as u64, 0)).unwrap();
459        }
460
461        for _ in 0..CidQueue::LEN - 1 {
462            assert!(q.next_reserved().is_some());
463        }
464
465        assert_eq!(q.next(), None);
466    }
467
468    #[test]
469    fn insert_reserve_advance() {
470        let mut q = CidQueue::new(initial_cid());
471
472        let first = cid(1, 0);
473        let second = cid(2, 0);
474        let third = cid(3, 0);
475
476        q.insert(first).unwrap();
477        q.insert(second).unwrap();
478
479        assert_eq!(q.next_reserved(), Some(first.id));
480        q.insert(third).unwrap();
481        q.next();
482        assert_eq!(q.active(), second.id);
483    }
484
485    #[test]
486    fn sparse_insert_reserve_insert_advance() {
487        let mut q = CidQueue::new(initial_cid());
488
489        let one = cid(1, 0);
490        let two = cid(2, 0);
491        let three = cid(3, 0);
492
493        q.insert(two).unwrap();
494        q.insert(three).unwrap();
495        assert_eq!(q.next_reserved(), Some(two.id));
496        q.insert(one).unwrap();
497        q.next();
498        assert_eq!(q.active(), three.id);
499        assert_eq!(q.next_reserved(), None);
500    }
501
502    #[test]
503    fn reserve_many_next_clears_across_wraparound() {
504        let mut q = CidQueue::new(initial_cid());
505        for i in 1..CidQueue::LEN {
506            q.insert(cid(i as u64, 0)).unwrap();
507        }
508
509        for _ in 0..CidQueue::LEN - 2 {
510            assert!(q.next_reserved().is_some());
511        }
512
513        assert!(q.next().is_some());
514        q.insert(cid(CidQueue::LEN as u64, 0)).unwrap();
515        assert!(q.next_reserved().is_some());
516        q.insert(cid(CidQueue::LEN as u64 + 1, 0)).unwrap();
517        assert!(q.next().is_some());
518    }
519}