1use std::{fmt::Debug, ops::Range};
2
3use crate::{ConnectionId, ResetToken, frame::NewConnectionId};
4
5#[derive(Debug, Clone, Copy)]
7struct CidData(ConnectionId, Option<ResetToken>);
8
9#[derive(Debug)]
24pub(crate) struct CidQueue {
25 buffer: [Option<CidData>; Self::LEN],
27 cursor: usize,
29 offset: u64,
33 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 pub(crate) fn insert(
61 &mut self,
62 cid: NewConnectionId,
63 ) -> Result<Option<(Range<u64>, ResetToken)>, InsertError> {
64 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 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 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 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 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 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 #[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 #[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 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 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 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 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 pub(crate) fn active(&self) -> ConnectionId {
187 self.buffer[self.cursor].unwrap().0
188 }
189
190 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 Retired,
202 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 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}