1use std::cmp;
28
29use std::sync::Arc;
30
31use std::collections::hash_map;
32use std::collections::HashMap;
33use std::collections::HashSet;
34
35use intrusive_collections::intrusive_adapter;
36use intrusive_collections::KeyAdapter;
37use intrusive_collections::RBTree;
38use intrusive_collections::RBTreeAtomicLink;
39
40use smallvec::SmallVec;
41
42use crate::buffers::DefaultBufFactory;
43use crate::ranges::RangeSet;
44use crate::BufFactory;
45use crate::Error;
46use crate::Result;
47
48const DEFAULT_URGENCY: u8 = 127;
49
50pub const MAX_STREAM_WINDOW: u64 = 16 * 1024 * 1024;
52
53#[derive(Default)]
58pub struct StreamIdHasher {
59 id: u64,
60}
61
62#[derive(Debug, PartialEq, Clone, Copy)]
64pub struct RecvBufResetReturn {
65 pub max_data_delta: u64,
68
69 pub consumed_flowcontrol: u64,
72}
73
74impl RecvBufResetReturn {
75 pub fn zero() -> Self {
76 Self {
77 max_data_delta: 0,
78 consumed_flowcontrol: 0,
79 }
80 }
81}
82
83pub enum RecvAction<T: bytes::BufMut> {
85 Emit { out: T },
87 Discard { len: usize },
89}
90
91impl std::hash::Hasher for StreamIdHasher {
92 #[inline]
93 fn finish(&self) -> u64 {
94 self.id
95 }
96
97 #[inline]
98 fn write_u64(&mut self, id: u64) {
99 self.id = id;
100 }
101
102 #[inline]
103 fn write(&mut self, _: &[u8]) {
104 unimplemented!()
107 }
108}
109
110type BuildStreamIdHasher = std::hash::BuildHasherDefault<StreamIdHasher>;
111
112pub type StreamIdHashMap<V> = HashMap<u64, V, BuildStreamIdHasher>;
113pub type StreamIdHashSet = HashSet<u64, BuildStreamIdHasher>;
114
115#[derive(Default)]
117struct CollectedStreams {
118 ranges: Option<Box<[RangeSet; 4]>>,
122}
123
124impl CollectedStreams {
125 fn insert(&mut self, stream_id: u64) {
126 let ranges = self.ranges.get_or_insert_with(Default::default);
129 ranges[(stream_id & 0x3) as usize].push_item(stream_id >> 2);
130 }
131
132 fn contains(&self, stream_id: u64) -> bool {
133 let Some(ranges) = &self.ranges else {
134 return false;
135 };
136
137 ranges[(stream_id & 0x3) as usize].contains(stream_id >> 2)
138 }
139}
140
141#[derive(Default)]
143pub struct StreamMap<F: BufFactory = DefaultBufFactory> {
144 streams: StreamIdHashMap<Stream<F>>,
146
147 collected: CollectedStreams,
153
154 peer_max_streams_bidi: u64,
156
157 peer_max_streams_uni: u64,
159
160 peer_opened_streams_bidi: u64,
162
163 peer_opened_streams_uni: u64,
165
166 local_max_streams_bidi: u64,
168 local_max_streams_bidi_next: u64,
169
170 initial_max_streams_bidi: u64,
172
173 local_max_streams_uni: u64,
175 local_max_streams_uni_next: u64,
176
177 initial_max_streams_uni: u64,
179
180 local_opened_streams_bidi: u64,
182
183 local_opened_streams_uni: u64,
185
186 flushable: RBTree<StreamFlushablePriorityAdapter>,
190
191 pub readable: RBTree<StreamReadablePriorityAdapter>,
195
196 pub writable: RBTree<StreamWritablePriorityAdapter>,
199
200 stopped_writable: RBTree<StreamStoppedWritablePriorityAdapter>,
202
203 almost_full: StreamIdHashSet,
208
209 blocked: StreamIdHashMap<u64>,
213
214 reset: StreamIdHashMap<(u64, u64)>,
218
219 stopped: StreamIdHashMap<u64>,
223
224 max_stream_window: u64,
226
227 tx_buffered: usize,
229}
230
231impl<F: BufFactory> StreamMap<F> {
232 pub fn new(
233 max_streams_bidi: u64, max_streams_uni: u64, max_stream_window: u64,
234 ) -> Self {
235 StreamMap {
236 local_max_streams_bidi: max_streams_bidi,
237 local_max_streams_bidi_next: max_streams_bidi,
238 initial_max_streams_bidi: max_streams_bidi,
239
240 local_max_streams_uni: max_streams_uni,
241 local_max_streams_uni_next: max_streams_uni,
242 initial_max_streams_uni: max_streams_uni,
243
244 max_stream_window,
245
246 ..StreamMap::default()
247 }
248 }
249
250 pub fn get(&self, id: u64) -> Option<&Stream<F>> {
252 self.streams.get(&id)
253 }
254
255 pub fn get_mut(&mut self, id: u64) -> Option<&mut Stream<F>> {
257 self.streams.get_mut(&id)
258 }
259
260 pub(crate) fn get_or_create(
271 &mut self, id: u64, local_params: &crate::TransportParams,
272 peer_params: &crate::TransportParams, local: bool, is_server: bool,
273 ) -> Result<&mut Stream<F>> {
274 let (stream, is_new_and_writable) = match self.streams.entry(id) {
275 hash_map::Entry::Vacant(v) => {
276 if self.collected.contains(id) {
278 return Err(Error::Done);
279 }
280
281 if local != is_local(id, is_server) {
282 return Err(Error::InvalidStreamState(id));
283 }
284
285 let (max_rx_data, max_tx_data) = match (local, is_bidi(id)) {
286 (true, true) => (
288 local_params.initial_max_stream_data_bidi_local,
289 peer_params.initial_max_stream_data_bidi_remote,
290 ),
291
292 (true, false) => (0, peer_params.initial_max_stream_data_uni),
294
295 (false, true) => (
297 local_params.initial_max_stream_data_bidi_remote,
298 peer_params.initial_max_stream_data_bidi_local,
299 ),
300
301 (false, false) =>
303 (local_params.initial_max_stream_data_uni, 0),
304 };
305
306 let stream_sequence = id >> 2;
310
311 match (is_local(id, is_server), is_bidi(id)) {
313 (true, true) => {
314 let n = cmp::max(
315 self.local_opened_streams_bidi,
316 stream_sequence + 1,
317 );
318
319 if n > self.peer_max_streams_bidi {
320 return Err(Error::StreamLimit);
321 }
322
323 self.local_opened_streams_bidi = n;
324 },
325
326 (true, false) => {
327 let n = cmp::max(
328 self.local_opened_streams_uni,
329 stream_sequence + 1,
330 );
331
332 if n > self.peer_max_streams_uni {
333 return Err(Error::StreamLimit);
334 }
335
336 self.local_opened_streams_uni = n;
337 },
338
339 (false, true) => {
340 let n = cmp::max(
341 self.peer_opened_streams_bidi,
342 stream_sequence + 1,
343 );
344
345 if n > self.local_max_streams_bidi {
346 return Err(Error::StreamLimit);
347 }
348
349 self.peer_opened_streams_bidi = n;
350 },
351
352 (false, false) => {
353 let n = cmp::max(
354 self.peer_opened_streams_uni,
355 stream_sequence + 1,
356 );
357
358 if n > self.local_max_streams_uni {
359 return Err(Error::StreamLimit);
360 }
361
362 self.peer_opened_streams_uni = n;
363 },
364 };
365
366 let initial_window = max_rx_data;
367 let s = Stream::new(
368 id,
369 max_rx_data,
370 max_tx_data,
371 local,
372 initial_window,
373 self.max_stream_window,
374 );
375
376 let is_writable = s.is_writable();
377
378 (v.insert(s), is_writable)
379 },
380
381 hash_map::Entry::Occupied(v) => (v.into_mut(), false),
382 };
383
384 if is_new_and_writable {
387 self.writable.insert(Arc::clone(&stream.priority_key));
388 }
389
390 Ok(stream)
391 }
392
393 pub fn insert_readable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
397 if !priority_key.readable.is_linked() {
398 self.readable.insert(Arc::clone(priority_key));
399 }
400 }
401
402 pub fn remove_readable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
404 if !priority_key.readable.is_linked() {
405 return;
406 }
407
408 let mut c = {
409 let ptr = Arc::as_ptr(priority_key);
410 unsafe { self.readable.cursor_mut_from_ptr(ptr) }
411 };
412
413 c.remove();
414 }
415
416 pub fn insert_writable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
423 if !priority_key.writable.is_linked() {
424 self.writable.insert(Arc::clone(priority_key));
425 }
426 }
427
428 pub fn insert_stopped_writable(
430 &mut self, priority_key: &Arc<StreamPriorityKey>,
431 ) {
432 if !priority_key.stopped_writable.is_linked() {
433 self.stopped_writable.insert(Arc::clone(priority_key));
434 }
435
436 self.insert_writable(priority_key);
437 }
438
439 pub fn remove_writable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
444 if priority_key.stopped_writable.is_linked() {
445 let ptr = Arc::as_ptr(priority_key);
446 let mut c = unsafe { self.stopped_writable.cursor_mut_from_ptr(ptr) };
447 c.remove();
448 }
449
450 if !priority_key.writable.is_linked() {
451 return;
452 }
453
454 let mut c = {
455 let ptr = Arc::as_ptr(priority_key);
456 unsafe { self.writable.cursor_mut_from_ptr(ptr) }
457 };
458
459 c.remove();
460 }
461
462 pub fn insert_flushable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
466 if !priority_key.flushable.is_linked() {
467 self.flushable.insert(Arc::clone(priority_key));
468 }
469 }
470
471 pub fn remove_flushable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
473 if !priority_key.flushable.is_linked() {
474 return;
475 }
476
477 let mut c = {
478 let ptr = Arc::as_ptr(priority_key);
479 unsafe { self.flushable.cursor_mut_from_ptr(ptr) }
480 };
481
482 c.remove();
483 }
484
485 pub fn peek_flushable(&self) -> Option<Arc<StreamPriorityKey>> {
486 self.flushable.front().clone_pointer()
487 }
488
489 pub fn update_priority(
491 &mut self, old: &Arc<StreamPriorityKey>, new: &Arc<StreamPriorityKey>,
492 ) {
493 if old.readable.is_linked() {
494 self.remove_readable(old);
495 self.readable.insert(Arc::clone(new));
496 }
497
498 if old.writable.is_linked() {
499 let stopped = old.stopped_writable.is_linked();
500 self.remove_writable(old);
501 if stopped {
502 self.insert_stopped_writable(new);
503 } else {
504 self.insert_writable(new);
505 }
506 }
507
508 if old.flushable.is_linked() {
509 self.remove_flushable(old);
510 self.flushable.insert(Arc::clone(new));
511 }
512 }
513
514 pub fn insert_almost_full(&mut self, stream_id: u64) {
518 self.almost_full.insert(stream_id);
519 }
520
521 pub fn remove_almost_full(&mut self, stream_id: u64) {
523 self.almost_full.remove(&stream_id);
524 }
525
526 pub fn insert_blocked(&mut self, stream_id: u64, off: u64) {
531 self.blocked.insert(stream_id, off);
532 }
533
534 pub fn remove_blocked(&mut self, stream_id: u64) {
536 self.blocked.remove(&stream_id);
537 }
538
539 pub fn insert_reset(
544 &mut self, stream_id: u64, error_code: u64, final_size: u64,
545 ) {
546 self.reset.insert(stream_id, (error_code, final_size));
547 }
548
549 pub fn remove_reset(&mut self, stream_id: u64) {
551 self.reset.remove(&stream_id);
552 }
553
554 pub fn insert_stopped(&mut self, stream_id: u64, error_code: u64) {
559 self.stopped.insert(stream_id, error_code);
560 }
561
562 pub fn remove_stopped(&mut self, stream_id: u64) {
564 self.stopped.remove(&stream_id);
565 }
566
567 pub fn update_peer_max_streams_bidi(&mut self, v: u64) {
569 self.peer_max_streams_bidi = cmp::max(self.peer_max_streams_bidi, v);
570 }
571
572 pub fn update_peer_max_streams_uni(&mut self, v: u64) {
574 self.peer_max_streams_uni = cmp::max(self.peer_max_streams_uni, v);
575 }
576
577 pub fn update_max_streams_bidi(&mut self) {
579 self.local_max_streams_bidi = self.local_max_streams_bidi_next;
580 }
581
582 pub fn set_max_streams_bidi(&mut self, max: u64) {
584 self.local_max_streams_bidi = max;
585 self.local_max_streams_bidi_next = max;
586 self.initial_max_streams_bidi = max;
587 }
588
589 pub fn max_streams_bidi(&self) -> u64 {
591 self.local_max_streams_bidi
592 }
593
594 pub fn max_streams_bidi_next(&mut self) -> u64 {
596 self.local_max_streams_bidi_next
597 }
598
599 pub fn update_max_streams_uni(&mut self) {
601 self.local_max_streams_uni = self.local_max_streams_uni_next;
602 }
603
604 pub fn max_streams_uni_next(&mut self) -> u64 {
606 self.local_max_streams_uni_next
607 }
608
609 pub fn peer_max_streams_bidi(&self) -> u64 {
611 self.peer_max_streams_bidi
612 }
613
614 pub fn peer_streams_left_bidi(&self) -> u64 {
617 self.peer_max_streams_bidi - self.local_opened_streams_bidi
618 }
619
620 pub fn peer_max_streams_uni(&self) -> u64 {
622 self.peer_max_streams_uni
623 }
624
625 pub fn peer_streams_left_uni(&self) -> u64 {
628 self.peer_max_streams_uni - self.local_opened_streams_uni
629 }
630
631 pub fn mark_stop_reported(&mut self, stream_id: u64) {
633 let stream = self.streams.get_mut(&stream_id).unwrap();
634 stream.send.mark_stop_reported();
635
636 if stream.is_collectable() {
637 let local = stream.local;
638 self.collect(stream_id, local);
639 } else {
640 let priority_key = Arc::clone(&stream.priority_key);
641 self.remove_writable(&priority_key);
642 }
643 }
644
645 pub fn collect(&mut self, stream_id: u64, local: bool) {
650 if !local {
651 if is_bidi(stream_id) {
654 self.local_max_streams_bidi_next =
655 self.local_max_streams_bidi_next.saturating_add(1);
656 } else {
657 self.local_max_streams_uni_next =
658 self.local_max_streams_uni_next.saturating_add(1);
659 }
660 }
661
662 let s = self.streams.remove(&stream_id).unwrap();
663
664 self.remove_readable(&s.priority_key);
665
666 self.remove_writable(&s.priority_key);
667
668 self.remove_flushable(&s.priority_key);
669
670 self.collected.insert(stream_id);
671 }
672
673 pub fn readable(&self) -> StreamIter {
675 StreamIter {
676 streams: self.readable.iter().map(|s| s.id).collect(),
677 index: 0,
678 }
679 }
680
681 pub fn writable(&self) -> StreamIter {
683 StreamIter {
684 streams: self.writable.iter().map(|s| s.id).collect(),
685 index: 0,
686 }
687 }
688
689 pub fn stopped_writable(&self) -> StreamIter {
691 if self.stopped_writable.is_empty() {
692 return StreamIter::default();
693 }
694
695 StreamIter {
696 streams: self.stopped_writable.iter().map(|key| key.id).collect(),
697 index: 0,
698 }
699 }
700
701 pub fn pop_stopped_writable(&mut self) -> Option<u64> {
703 let priority_key = self.stopped_writable.front().clone_pointer()?;
704 self.remove_writable(&priority_key);
705 Some(priority_key.id)
706 }
707
708 pub fn local_stream_opened(&self, stream_id: u64) -> bool {
713 let opened = if is_bidi(stream_id) {
714 self.local_opened_streams_bidi
715 } else {
716 self.local_opened_streams_uni
717 };
718
719 stream_id >> 2 < opened
720 }
721
722 pub fn almost_full(&self) -> StreamIter {
724 StreamIter::from(&self.almost_full)
725 }
726
727 pub fn blocked(&self) -> hash_map::Iter<'_, u64, u64> {
729 self.blocked.iter()
730 }
731
732 pub fn reset(&self) -> hash_map::Iter<'_, u64, (u64, u64)> {
734 self.reset.iter()
735 }
736
737 pub fn stopped(&self) -> hash_map::Iter<'_, u64, u64> {
739 self.stopped.iter()
740 }
741
742 pub fn is_collected(&self, stream_id: u64) -> bool {
744 self.collected.contains(stream_id)
745 }
746
747 pub fn has_flushable(&self) -> bool {
749 !self.flushable.is_empty()
750 }
751
752 pub fn has_readable(&self) -> bool {
754 !self.readable.is_empty()
755 }
756
757 pub fn has_almost_full(&self) -> bool {
760 !self.almost_full.is_empty()
761 }
762
763 pub fn has_blocked(&self) -> bool {
765 !self.blocked.is_empty()
766 }
767
768 pub fn has_reset(&self) -> bool {
770 !self.reset.is_empty()
771 }
772
773 pub fn has_stopped(&self) -> bool {
775 !self.stopped.is_empty()
776 }
777
778 pub fn should_update_max_streams_bidi(&self) -> bool {
784 let available = self
785 .local_max_streams_bidi
786 .saturating_sub(self.peer_opened_streams_bidi);
787 self.local_max_streams_bidi_next != self.local_max_streams_bidi &&
788 available <= self.initial_max_streams_bidi / 2
789 }
790
791 pub fn should_update_max_streams_uni(&self) -> bool {
797 let available = self
798 .local_max_streams_uni
799 .saturating_sub(self.peer_opened_streams_uni);
800 self.local_max_streams_uni_next != self.local_max_streams_uni &&
801 available <= self.initial_max_streams_uni / 2
802 }
803
804 #[cfg(test)]
806 pub fn len(&self) -> usize {
807 self.streams.len()
808 }
809
810 pub(crate) fn tx_buffered(&self) -> usize {
812 self.tx_buffered
813 }
814
815 fn tx_buffered_actual(&self) -> usize {
819 self.streams
820 .values()
821 .map(|s| s.send.buffered_bytes() as usize)
822 .sum()
823 }
824
825 pub(crate) fn tx_buffered_is_consistent(&self) -> bool {
828 self.tx_buffered == self.tx_buffered_actual()
829 }
830
831 pub(crate) fn add_tx_buffered(&mut self, delta: usize) {
833 self.tx_buffered += delta;
834
835 #[cfg(debug_assertions)]
836 self.debug_check_tx_buffered_consistency();
837 }
838
839 pub(crate) fn sub_tx_buffered(&mut self, delta: usize) {
841 debug_assert!(self.tx_buffered >= delta);
842 self.tx_buffered = self.tx_buffered.saturating_sub(delta);
843
844 #[cfg(debug_assertions)]
845 self.debug_check_tx_buffered_consistency();
846 }
847
848 #[cfg(debug_assertions)]
852 pub(crate) fn debug_check_tx_buffered_consistency(&self) {
853 if !self.tx_buffered_is_consistent() {
854 let buffered_per_stream = self
855 .streams
856 .iter()
857 .map(|(id, s)| (*id, s.send.buffered_bytes()))
858 .collect::<Vec<_>>();
859
860 let actual = self.tx_buffered_actual();
861 let stored = self.tx_buffered;
862 panic!(
863 "tx_buffered mismatch: stored={}, actual={}, diff={}, buffered_per_stream={:?}",
864 stored,
865 actual,
866 stored as i64 - actual as i64,
867 buffered_per_stream
868 );
869 }
870 }
871}
872
873pub struct Stream<F: BufFactory = DefaultBufFactory> {
875 pub recv: recv_buf::RecvBuf,
877
878 pub send: send_buf::SendBuf<F>,
880
881 pub send_lowat: usize,
882
883 pub bidi: bool,
885
886 pub local: bool,
888
889 pub urgency: u8,
891
892 pub incremental: bool,
894
895 pub priority_key: Arc<StreamPriorityKey>,
896}
897
898impl<F: BufFactory> Stream<F> {
899 pub fn new(
901 id: u64, max_rx_data: u64, max_tx_data: u64, local: bool,
902 initial_window: u64, max_window: u64,
903 ) -> Self {
904 let priority_key = Arc::new(StreamPriorityKey {
905 id,
906 ..Default::default()
907 });
908
909 Stream {
910 recv: recv_buf::RecvBuf::new(max_rx_data, initial_window, max_window),
911 send: send_buf::SendBuf::new(max_tx_data),
912 send_lowat: 1,
913 bidi: is_bidi(id),
914 local,
915 urgency: priority_key.urgency,
916 incremental: priority_key.incremental,
917 priority_key,
918 }
919 }
920
921 pub fn is_readable(&self) -> bool {
923 self.recv.ready()
924 }
925
926 pub fn is_writable(&self) -> bool {
929 !self.send.is_shutdown() &&
930 !self.send.is_fin() &&
931 (self.send.off_back() + self.send_lowat as u64) <
932 self.send.max_off()
933 }
934
935 pub fn is_flushable(&self) -> bool {
938 let off_front = self.send.off_front();
939
940 !self.send.is_empty() &&
941 off_front < self.send.off_back() &&
942 off_front < self.send.max_off()
943 }
944
945 pub fn is_complete(&self) -> bool {
955 match (self.bidi, self.local) {
956 (true, _) => self.recv.is_fin() && self.send.is_complete(),
959
960 (false, true) => self.send.is_complete(),
963
964 (false, false) => self.recv.is_fin(),
967 }
968 }
969
970 pub fn is_collectable(&self) -> bool {
972 self.is_complete() &&
973 !self.is_readable() &&
974 !self.send.has_unreported_stop()
975 }
976}
977
978pub fn is_local(stream_id: u64, is_server: bool) -> bool {
980 (stream_id & 0x1) == (is_server as u64)
981}
982
983pub fn is_bidi(stream_id: u64) -> bool {
985 (stream_id & 0x2) == 0
986}
987
988#[derive(Clone, Debug)]
989pub struct StreamPriorityKey {
990 pub urgency: u8,
991 pub incremental: bool,
992 pub id: u64,
993
994 pub readable: RBTreeAtomicLink,
995 pub writable: RBTreeAtomicLink,
996 pub stopped_writable: RBTreeAtomicLink,
997 pub flushable: RBTreeAtomicLink,
998}
999
1000impl Default for StreamPriorityKey {
1001 fn default() -> Self {
1002 Self {
1003 urgency: DEFAULT_URGENCY,
1004 incremental: true,
1005 id: Default::default(),
1006 readable: Default::default(),
1007 writable: Default::default(),
1008 stopped_writable: Default::default(),
1009 flushable: Default::default(),
1010 }
1011 }
1012}
1013
1014impl PartialEq for StreamPriorityKey {
1015 fn eq(&self, other: &Self) -> bool {
1016 self.id == other.id
1017 }
1018}
1019
1020impl Eq for StreamPriorityKey {}
1021
1022impl PartialOrd for StreamPriorityKey {
1023 fn partial_cmp(&self, other: &Self) -> Option<cmp::Ordering> {
1024 Some(self.cmp(other))
1025 }
1026}
1027
1028impl Ord for StreamPriorityKey {
1029 fn cmp(&self, other: &Self) -> cmp::Ordering {
1030 if self.id == other.id {
1032 return cmp::Ordering::Equal;
1033 }
1034
1035 if self.urgency != other.urgency {
1037 return self.urgency.cmp(&other.urgency);
1038 }
1039
1040 if !self.incremental && !other.incremental {
1043 return self.id.cmp(&other.id);
1044 }
1045
1046 if self.incremental && !other.incremental {
1048 return cmp::Ordering::Greater;
1049 }
1050 if !self.incremental && other.incremental {
1051 return cmp::Ordering::Less;
1052 }
1053
1054 cmp::Ordering::Greater
1058 }
1059}
1060
1061intrusive_adapter!(pub StreamWritablePriorityAdapter = Arc<StreamPriorityKey>: StreamPriorityKey { writable => RBTreeAtomicLink });
1062
1063impl KeyAdapter<'_> for StreamWritablePriorityAdapter {
1064 type Key = StreamPriorityKey;
1065
1066 fn get_key(&self, s: &StreamPriorityKey) -> Self::Key {
1067 s.clone()
1068 }
1069}
1070
1071intrusive_adapter!(pub StreamStoppedWritablePriorityAdapter = Arc<StreamPriorityKey>: StreamPriorityKey { stopped_writable => RBTreeAtomicLink });
1072
1073impl KeyAdapter<'_> for StreamStoppedWritablePriorityAdapter {
1074 type Key = StreamPriorityKey;
1075
1076 fn get_key(&self, s: &StreamPriorityKey) -> Self::Key {
1077 s.clone()
1078 }
1079}
1080
1081intrusive_adapter!(pub StreamReadablePriorityAdapter = Arc<StreamPriorityKey>: StreamPriorityKey { readable => RBTreeAtomicLink });
1082
1083impl KeyAdapter<'_> for StreamReadablePriorityAdapter {
1084 type Key = StreamPriorityKey;
1085
1086 fn get_key(&self, s: &StreamPriorityKey) -> Self::Key {
1087 s.clone()
1088 }
1089}
1090
1091intrusive_adapter!(pub StreamFlushablePriorityAdapter = Arc<StreamPriorityKey>: StreamPriorityKey { flushable => RBTreeAtomicLink });
1092
1093impl KeyAdapter<'_> for StreamFlushablePriorityAdapter {
1094 type Key = StreamPriorityKey;
1095
1096 fn get_key(&self, s: &StreamPriorityKey) -> Self::Key {
1097 s.clone()
1098 }
1099}
1100
1101#[derive(Default)]
1103pub struct StreamIter {
1104 streams: SmallVec<[u64; 8]>,
1105 index: usize,
1106}
1107
1108impl StreamIter {
1109 #[inline]
1110 fn from(streams: &StreamIdHashSet) -> Self {
1111 StreamIter {
1112 streams: streams.iter().copied().collect(),
1113 index: 0,
1114 }
1115 }
1116}
1117
1118impl Iterator for StreamIter {
1119 type Item = u64;
1120
1121 #[inline]
1122 fn next(&mut self) -> Option<Self::Item> {
1123 let v = self.streams.get(self.index)?;
1124 self.index += 1;
1125 Some(*v)
1126 }
1127}
1128
1129impl ExactSizeIterator for StreamIter {
1130 #[inline]
1131 fn len(&self) -> usize {
1132 self.streams.len() - self.index
1133 }
1134}
1135
1136#[cfg(test)]
1137mod tests {
1138 use rstest::rstest;
1139
1140 use crate::range_buf::RangeBuf;
1141 use crate::test_utils::Pipe;
1142
1143 use super::*;
1144
1145 const DEFAULT_STREAM_WINDOW: u64 = 32 * 1024;
1147
1148 #[rstest]
1149 fn collected_streams_per_type(#[values(0, 1, 2, 3)] stream_type: u64) {
1150 let mut collected = CollectedStreams::default();
1151 assert!(collected.ranges.is_none());
1152
1153 for id in 0..4 {
1154 assert!(!collected.contains(id));
1155 }
1156 assert!(collected.ranges.is_none());
1157
1158 collected.insert(stream_type);
1159 assert!(collected.ranges.is_some());
1160
1161 for id in 0..4 {
1162 assert_eq!(collected.contains(id), id == stream_type);
1163 }
1164
1165 collected.insert(8 | stream_type);
1166 assert!(!collected.contains(4 | stream_type));
1167 assert_eq!(
1168 collected.ranges.as_ref().unwrap()[stream_type as usize].len(),
1169 2
1170 );
1171
1172 collected.insert(4 | stream_type);
1173 collected.insert(4 | stream_type);
1174 assert_eq!(
1175 collected.ranges.as_ref().unwrap()[stream_type as usize],
1176 0..3
1177 );
1178
1179 for id in 0..12 {
1180 assert_eq!(collected.contains(id), id & 0x3 == stream_type);
1181 }
1182 }
1183
1184 #[test]
1185 fn collected_streams_preserve_fragmented_and_large_ids() {
1186 let mut collected = CollectedStreams::default();
1187
1188 for sequence in (0..2048).step_by(2) {
1189 for stream_type in 0..4 {
1190 collected.insert((sequence << 2) | stream_type);
1191 }
1192 }
1193
1194 for sequence in 0..2048 {
1195 for stream_type in 0..4 {
1196 assert_eq!(
1197 collected.contains((sequence << 2) | stream_type),
1198 sequence % 2 == 0
1199 );
1200 }
1201 }
1202
1203 for stream_type in 0..4 {
1204 assert_eq!(
1205 collected.ranges.as_ref().unwrap()[stream_type as usize].len(),
1206 1024
1207 );
1208
1209 let stream_id = ((1u64 << 62) - 4) | stream_type;
1210 assert!(!collected.contains(stream_id));
1211 collected.insert(stream_id);
1212 assert!(collected.contains(stream_id));
1213 assert!(!collected.contains(stream_id - 4));
1214 assert!(!collected.contains((1u64 << 40) | stream_type));
1215 }
1216 }
1217
1218 fn collect_pipe_stream(pipe: &mut Pipe, stream_id: u64) {
1221 let mut buf = [0; 1];
1222
1223 assert_eq!(pipe.client.stream_send(stream_id, b"a", true), Ok(1));
1224 assert_eq!(pipe.advance(), Ok(()));
1225 assert_eq!(pipe.server.stream_recv(stream_id, &mut buf), Ok((1, true)));
1226
1227 if is_bidi(stream_id) {
1228 assert_eq!(pipe.server.stream_send(stream_id, b"a", true), Ok(1));
1229 assert_eq!(pipe.advance(), Ok(()));
1230 assert_eq!(
1231 pipe.client.stream_recv(stream_id, &mut buf),
1232 Ok((1, true))
1233 );
1234 }
1235
1236 assert_eq!(pipe.advance(), Ok(()));
1237 }
1238
1239 #[rstest]
1240 fn collected_streams_out_of_order(
1241 #[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
1242 #[values(0, 2)] stream_type: u64,
1243 ) {
1244 let mut pipe = Pipe::new(cc_algorithm_name).unwrap();
1245 assert_eq!(pipe.handshake(), Ok(()));
1246
1247 for sequence in [2, 0, 1] {
1248 let stream_id = (sequence << 2) | stream_type;
1249 collect_pipe_stream(&mut pipe, stream_id);
1250
1251 for conn in [&mut pipe.client, &mut pipe.server] {
1252 assert!(conn.streams.get(stream_id).is_none());
1253 assert!(conn.streams.is_collected(stream_id));
1254 assert!(conn.stream_closed(stream_id));
1255 assert_eq!(
1256 conn.streams
1257 .get_or_create(
1258 stream_id,
1259 &conn.local_transport_params,
1260 &conn.peer_transport_params,
1261 !conn.is_server,
1262 conn.is_server,
1263 )
1264 .err(),
1265 Some(Error::Done)
1266 );
1267 }
1268 }
1269
1270 assert_eq!(
1271 pipe.server.streams.collected.ranges.as_ref().unwrap()
1272 [stream_type as usize],
1273 0..3
1274 );
1275
1276 let frames = [crate::frame::Frame::Stream {
1278 stream_id: 8 | stream_type,
1279 data: RangeBuf::from(b"a", 0, true),
1280 }];
1281 assert!(pipe
1282 .send_pkt_to_server(crate::Type::Short, &frames, &mut [0; 1280])
1283 .is_ok());
1284 assert_eq!(pipe.server.streams.len(), 0);
1285 }
1286
1287 #[rstest]
1288 fn collected_streams_sparse_peer_credit(
1289 #[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
1290 #[values(0, 2)] stream_type: u64, #[values(1, 8)] initial_limit: u64,
1291 ) {
1292 let mut config = Pipe::default_config(cc_algorithm_name).unwrap();
1293 config.set_initial_max_streams_bidi(initial_limit);
1294 config.set_initial_max_streams_uni(initial_limit);
1295
1296 let mut pipe = Pipe::with_config(&mut config).unwrap();
1297 assert_eq!(pipe.handshake(), Ok(()));
1298
1299 for index in 0..initial_limit {
1300 let stream_id = (index << 3) | stream_type;
1301 collect_pipe_stream(&mut pipe, stream_id);
1302
1303 assert_eq!(
1304 pipe.server.streams.collected.ranges.as_ref().unwrap()
1305 [stream_type as usize]
1306 .len(),
1307 index as usize + 1
1308 );
1309 assert!(pipe.server.stream_closed(stream_id));
1310 assert!(pipe.client.stream_closed(stream_id));
1311 assert_eq!(pipe.server.streams.len(), 0);
1312 }
1313
1314 let peer_limit = if is_bidi(stream_type) {
1315 pipe.client.streams.peer_max_streams_bidi()
1316 } else {
1317 pipe.client.streams.peer_max_streams_uni()
1318 };
1319 assert_eq!(peer_limit, 2 * initial_limit);
1320
1321 for sequence in (1..2 * initial_limit - 1).step_by(2) {
1323 let stream_id = (sequence << 2) | stream_type;
1324
1325 for conn in [&pipe.client, &pipe.server] {
1326 assert!(conn.streams.get(stream_id).is_none());
1327 assert!(!conn.streams.is_collected(stream_id));
1328 assert!(!conn.stream_closed(stream_id));
1329 }
1330 }
1331
1332 for sequence in (1..2 * initial_limit - 1).step_by(2) {
1334 collect_pipe_stream(&mut pipe, (sequence << 2) | stream_type);
1335 }
1336
1337 assert_eq!(
1338 pipe.server.streams.collected.ranges.as_ref().unwrap()
1339 [stream_type as usize],
1340 0..2 * initial_limit - 1
1341 );
1342 }
1343
1344 #[rstest]
1345 fn collected_streams_fragmented_local(
1346 #[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
1347 #[values(0, 2)] stream_type: u64,
1348 ) {
1349 let mut client_config = Pipe::default_config(cc_algorithm_name).unwrap();
1350 client_config.set_initial_max_streams_bidi(1);
1351 client_config.set_initial_max_streams_uni(1);
1352
1353 let mut server_config = Pipe::default_config(cc_algorithm_name).unwrap();
1354 server_config.set_initial_max_streams_bidi(32);
1355 server_config.set_initial_max_streams_uni(32);
1356
1357 let mut pipe = Pipe::with_client_and_server_config(
1358 &mut client_config,
1359 &mut server_config,
1360 )
1361 .unwrap();
1362 assert_eq!(pipe.handshake(), Ok(()));
1363
1364 for index in 0..32 {
1367 collect_pipe_stream(&mut pipe, (index << 3) | stream_type);
1368 assert_eq!(
1369 pipe.client.streams.collected.ranges.as_ref().unwrap()
1370 [stream_type as usize]
1371 .len(),
1372 index as usize + 1
1373 );
1374 assert_eq!(pipe.client.streams.len(), 0);
1375 }
1376 }
1377
1378 #[rstest]
1379 fn stream_limit_does_not_collect(
1380 #[values(0, 1, 2, 3)] stream_type: u64,
1381 #[values(true, false)] local: bool,
1382 ) {
1383 let params = crate::TransportParams::default();
1384 let mut streams = <StreamMap>::new(1, 1, DEFAULT_STREAM_WINDOW);
1385 streams.update_peer_max_streams_bidi(1);
1386 streams.update_peer_max_streams_uni(1);
1387
1388 let stream_id = 4 | stream_type;
1389 let is_server = (stream_type & 1 != 0) == local;
1390 assert_eq!(
1391 streams
1392 .get_or_create(stream_id, ¶ms, ¶ms, local, is_server)
1393 .err(),
1394 Some(Error::StreamLimit)
1395 );
1396 assert!(!streams.is_collected(stream_id));
1397 assert_eq!(streams.len(), 0);
1398 assert!(streams.collected.ranges.is_none());
1399
1400 if local {
1401 streams.update_peer_max_streams_bidi(2);
1402 streams.update_peer_max_streams_uni(2);
1403 } else {
1404 streams.local_max_streams_bidi_next = 2;
1405 streams.local_max_streams_uni_next = 2;
1406 streams.update_max_streams_bidi();
1407 streams.update_max_streams_uni();
1408 }
1409
1410 assert!(streams
1411 .get_or_create(stream_id, ¶ms, ¶ms, local, is_server)
1412 .is_ok());
1413 assert!(streams.collected.ranges.is_none());
1414 }
1415
1416 #[test]
1417 fn recv_flow_control() {
1418 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1419 assert!(!stream.recv.almost_full());
1420
1421 let mut buf = [0; 32];
1422
1423 let first = RangeBuf::from(b"hello", 0, false);
1424 let second = RangeBuf::from(b"world", 5, false);
1425 let third = RangeBuf::from(b"something", 10, false);
1426
1427 assert_eq!(stream.recv.write(second), Ok(()));
1428 assert_eq!(stream.recv.write(first), Ok(()));
1429 assert!(!stream.recv.almost_full());
1430
1431 assert_eq!(stream.recv.write(third), Err(Error::FlowControl));
1432
1433 let (len, fin) = stream.recv.emit(&mut buf).unwrap();
1434 assert_eq!(&buf[..len], b"helloworld");
1435 assert!(!fin);
1436
1437 assert!(stream.recv.almost_full());
1438
1439 stream.recv.update_max_data(std::time::Instant::now());
1440 assert_eq!(stream.recv.max_data_next(), 25);
1441 assert!(!stream.recv.almost_full());
1442
1443 let third = RangeBuf::from(b"something", 10, false);
1444 assert_eq!(stream.recv.write(third), Ok(()));
1445 }
1446
1447 #[test]
1448 fn recv_past_fin() {
1449 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1450 assert!(!stream.recv.almost_full());
1451
1452 let first = RangeBuf::from(b"hello", 0, true);
1453 let second = RangeBuf::from(b"world", 5, false);
1454
1455 assert_eq!(stream.recv.write(first), Ok(()));
1456 assert_eq!(stream.recv.write(second), Err(Error::FinalSize));
1457 }
1458
1459 #[test]
1460 fn recv_fin_dup() {
1461 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1462 assert!(!stream.recv.almost_full());
1463
1464 let first = RangeBuf::from(b"hello", 0, true);
1465 let second = RangeBuf::from(b"hello", 0, true);
1466
1467 assert_eq!(stream.recv.write(first), Ok(()));
1468 assert_eq!(stream.recv.write(second), Ok(()));
1469
1470 let mut buf = [0; 32];
1471
1472 let (len, fin) = stream.recv.emit(&mut buf).unwrap();
1473 assert_eq!(&buf[..len], b"hello");
1474 assert!(fin);
1475 }
1476
1477 #[test]
1478 fn recv_fin_change() {
1479 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1480 assert!(!stream.recv.almost_full());
1481
1482 let first = RangeBuf::from(b"hello", 0, true);
1483 let second = RangeBuf::from(b"world", 5, true);
1484
1485 assert_eq!(stream.recv.write(second), Ok(()));
1486 assert_eq!(stream.recv.write(first), Err(Error::FinalSize));
1487 }
1488
1489 #[test]
1490 fn recv_fin_lower_than_received() {
1491 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1492 assert!(!stream.recv.almost_full());
1493
1494 let first = RangeBuf::from(b"hello", 0, true);
1495 let second = RangeBuf::from(b"world", 5, false);
1496
1497 assert_eq!(stream.recv.write(second), Ok(()));
1498 assert_eq!(stream.recv.write(first), Err(Error::FinalSize));
1499 }
1500
1501 #[test]
1502 fn recv_fin_flow_control() {
1503 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1504 assert!(!stream.recv.almost_full());
1505
1506 let mut buf = [0; 32];
1507
1508 let first = RangeBuf::from(b"hello", 0, false);
1509 let second = RangeBuf::from(b"world", 5, true);
1510
1511 assert_eq!(stream.recv.write(first), Ok(()));
1512 assert_eq!(stream.recv.write(second), Ok(()));
1513
1514 let (len, fin) = stream.recv.emit(&mut buf).unwrap();
1515 assert_eq!(&buf[..len], b"helloworld");
1516 assert!(fin);
1517
1518 assert!(!stream.recv.almost_full());
1519 }
1520
1521 #[test]
1522 fn recv_fin_reset_mismatch() {
1523 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1524 assert!(!stream.recv.almost_full());
1525
1526 let first = RangeBuf::from(b"hello", 0, true);
1527
1528 assert_eq!(stream.recv.write(first), Ok(()));
1529 assert_eq!(stream.recv.reset(0, 10), Err(Error::FinalSize));
1530 }
1531
1532 #[test]
1533 fn recv_reset_with_gap() {
1534 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1535 assert!(!stream.recv.almost_full());
1536
1537 let first = RangeBuf::from(b"hello", 0, false);
1538
1539 assert_eq!(stream.recv.write(first), Ok(()));
1540 assert_eq!(stream.recv.emit(&mut [0; 1]), Ok((1, false)));
1542 assert_eq!(
1544 stream.recv.reset(0, 10),
1545 Ok(RecvBufResetReturn {
1546 max_data_delta: 5,
1547 consumed_flowcontrol: 9
1549 })
1550 );
1551 assert_eq!(stream.recv.reset(0, 10), Ok(RecvBufResetReturn::zero()));
1552 }
1553
1554 #[test]
1555 fn recv_reset_dup() {
1556 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1557 assert!(!stream.recv.almost_full());
1558
1559 let first = RangeBuf::from(b"hello", 0, false);
1560
1561 assert_eq!(stream.recv.write(first), Ok(()));
1562 assert_eq!(
1563 stream.recv.reset(0, 5),
1564 Ok(RecvBufResetReturn {
1565 max_data_delta: 0,
1566 consumed_flowcontrol: 5
1567 })
1568 );
1569 assert_eq!(stream.recv.reset(0, 5), Ok(RecvBufResetReturn::zero()));
1570 }
1571
1572 #[test]
1573 fn recv_reset_change() {
1574 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1575 assert!(!stream.recv.almost_full());
1576
1577 let first = RangeBuf::from(b"hello", 0, false);
1578
1579 assert_eq!(stream.recv.write(first), Ok(()));
1580 assert_eq!(
1581 stream.recv.reset(0, 5),
1582 Ok(RecvBufResetReturn {
1583 max_data_delta: 0,
1584 consumed_flowcontrol: 5
1585 })
1586 );
1587 assert_eq!(stream.recv.reset(0, 10), Err(Error::FinalSize));
1588 }
1589
1590 #[test]
1591 fn recv_reset_lower_than_received() {
1592 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1593 assert!(!stream.recv.almost_full());
1594
1595 let first = RangeBuf::from(b"hello", 0, false);
1596
1597 assert_eq!(stream.recv.write(first), Ok(()));
1598 assert_eq!(stream.recv.reset(0, 4), Err(Error::FinalSize));
1599 }
1600
1601 #[test]
1602 fn send_flow_control() {
1603 let mut buf = [0; 25];
1604
1605 let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1606
1607 let first = b"hello";
1608 let second = b"world";
1609 let third = b"something";
1610
1611 assert!(stream.send.write(first, false).is_ok());
1612 assert!(stream.send.write(second, false).is_ok());
1613 assert!(stream.send.write(third, false).is_ok());
1614
1615 assert_eq!(stream.send.off_front(), 0);
1616
1617 let (written, fin) = stream.send.emit(&mut buf[..25]).unwrap();
1618 assert_eq!(written, 15);
1619 assert!(!fin);
1620 assert_eq!(&buf[..written], b"helloworldsomet");
1621
1622 assert_eq!(stream.send.off_front(), 15);
1623
1624 let (written, fin) = stream.send.emit(&mut buf[..25]).unwrap();
1625 assert_eq!(written, 0);
1626 assert!(!fin);
1627 assert_eq!(&buf[..written], b"");
1628
1629 stream.send.retransmit(0, 15);
1630
1631 assert_eq!(stream.send.off_front(), 0);
1632
1633 let (written, fin) = stream.send.emit(&mut buf[..10]).unwrap();
1634 assert_eq!(written, 10);
1635 assert!(!fin);
1636 assert_eq!(&buf[..written], b"helloworld");
1637
1638 assert_eq!(stream.send.off_front(), 10);
1639
1640 let (written, fin) = stream.send.emit(&mut buf[..10]).unwrap();
1641 assert_eq!(written, 5);
1642 assert!(!fin);
1643 assert_eq!(&buf[..written], b"somet");
1644 }
1645
1646 #[test]
1647 fn send_past_fin() {
1648 let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1649
1650 let first = b"hello";
1651 let second = b"world";
1652 let third = b"third";
1653
1654 assert_eq!(stream.send.write(first, false), Ok(5));
1655
1656 assert_eq!(stream.send.write(second, true), Ok(5));
1657 assert!(stream.send.is_fin());
1658
1659 assert_eq!(stream.send.write(third, false), Err(Error::FinalSize));
1660 }
1661
1662 #[test]
1663 fn send_fin_dup() {
1664 let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1665
1666 assert_eq!(stream.send.write(b"hello", true), Ok(5));
1667 assert!(stream.send.is_fin());
1668
1669 assert_eq!(stream.send.write(b"", true), Ok(0));
1670 assert!(stream.send.is_fin());
1671 }
1672
1673 #[test]
1674 fn send_undo_fin() {
1675 let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1676
1677 assert_eq!(stream.send.write(b"hello", true), Ok(5));
1678 assert!(stream.send.is_fin());
1679
1680 assert_eq!(
1681 stream.send.write(b"helloworld", true),
1682 Err(Error::FinalSize)
1683 );
1684 }
1685
1686 #[test]
1687 fn send_fin_max_data_match() {
1688 let mut buf = [0; 15];
1689
1690 let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1691
1692 let slice = b"hellohellohello";
1693
1694 assert!(stream.send.write(slice, true).is_ok());
1695
1696 let (written, fin) = stream.send.emit(&mut buf[..15]).unwrap();
1697 assert_eq!(written, 15);
1698 assert!(fin);
1699 assert_eq!(&buf[..written], slice);
1700 }
1701
1702 #[test]
1703 fn send_fin_zero_length() {
1704 let mut buf = [0; 5];
1705
1706 let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1707
1708 assert_eq!(stream.send.write(b"hello", false), Ok(5));
1709 assert_eq!(stream.send.write(b"", true), Ok(0));
1710 assert!(stream.send.is_fin());
1711
1712 let (written, fin) = stream.send.emit(&mut buf[..5]).unwrap();
1713 assert_eq!(written, 5);
1714 assert!(fin);
1715 assert_eq!(&buf[..written], b"hello");
1716 }
1717
1718 #[test]
1719 fn send_ack() {
1720 let mut buf = [0; 5];
1721
1722 let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1723
1724 assert_eq!(stream.send.write(b"hello", false), Ok(5));
1725 assert_eq!(stream.send.write(b"world", false), Ok(5));
1726 assert_eq!(stream.send.write(b"", true), Ok(0));
1727 assert!(stream.send.is_fin());
1728
1729 assert_eq!(stream.send.off_front(), 0);
1730
1731 let (written, fin) = stream.send.emit(&mut buf[..5]).unwrap();
1732 assert_eq!(written, 5);
1733 assert!(!fin);
1734 assert_eq!(&buf[..written], b"hello");
1735
1736 stream.send.ack_and_drop(0, 5);
1737
1738 stream.send.retransmit(0, 5);
1739
1740 assert_eq!(stream.send.off_front(), 5);
1741
1742 let (written, fin) = stream.send.emit(&mut buf[..5]).unwrap();
1743 assert_eq!(written, 5);
1744 assert!(fin);
1745 assert_eq!(&buf[..written], b"world");
1746 }
1747
1748 #[test]
1749 fn send_ack_reordering() {
1750 let mut buf = [0; 5];
1751
1752 let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1753
1754 assert_eq!(stream.send.write(b"hello", false), Ok(5));
1755 assert_eq!(stream.send.write(b"world", false), Ok(5));
1756 assert_eq!(stream.send.write(b"", true), Ok(0));
1757 assert!(stream.send.is_fin());
1758
1759 assert_eq!(stream.send.off_front(), 0);
1760
1761 let (written, fin) = stream.send.emit(&mut buf[..5]).unwrap();
1762 assert_eq!(written, 5);
1763 assert!(!fin);
1764 assert_eq!(&buf[..written], b"hello");
1765
1766 assert_eq!(stream.send.off_front(), 5);
1767
1768 let (written, fin) = stream.send.emit(&mut buf[..1]).unwrap();
1769 assert_eq!(written, 1);
1770 assert!(!fin);
1771 assert_eq!(&buf[..written], b"w");
1772
1773 stream.send.ack_and_drop(5, 1);
1774 stream.send.ack_and_drop(0, 5);
1775
1776 stream.send.retransmit(0, 5);
1777 stream.send.retransmit(5, 1);
1778
1779 assert_eq!(stream.send.off_front(), 6);
1780
1781 let (written, fin) = stream.send.emit(&mut buf[..5]).unwrap();
1782 assert_eq!(written, 4);
1783 assert!(fin);
1784 assert_eq!(&buf[..written], b"orld");
1785 }
1786
1787 #[test]
1788 fn recv_data_below_off() {
1789 let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1790
1791 let first = RangeBuf::from(b"hello", 0, false);
1792
1793 assert_eq!(stream.recv.write(first), Ok(()));
1794
1795 let mut buf = [0; 10];
1796
1797 let (len, fin) = stream.recv.emit(&mut buf).unwrap();
1798 assert_eq!(&buf[..len], b"hello");
1799 assert!(!fin);
1800
1801 let first = RangeBuf::from(b"elloworld", 1, true);
1802 assert_eq!(stream.recv.write(first), Ok(()));
1803
1804 let (len, fin) = stream.recv.emit(&mut buf).unwrap();
1805 assert_eq!(&buf[..len], b"world");
1806 assert!(fin);
1807 }
1808
1809 #[test]
1810 fn stream_complete() {
1811 let mut stream =
1812 <Stream>::new(0, 30, 30, true, 30, DEFAULT_STREAM_WINDOW);
1813
1814 assert_eq!(stream.send.write(b"hello", false), Ok(5));
1815 assert_eq!(stream.send.write(b"world", false), Ok(5));
1816
1817 assert!(!stream.send.is_complete());
1818 assert!(!stream.send.is_fin());
1819
1820 assert_eq!(stream.send.write(b"", true), Ok(0));
1821
1822 assert!(!stream.send.is_complete());
1823 assert!(stream.send.is_fin());
1824
1825 let buf = RangeBuf::from(b"hello", 0, true);
1826 assert!(stream.recv.write(buf).is_ok());
1827 assert!(!stream.recv.is_fin());
1828
1829 stream.send.ack(6, 4);
1830 assert!(!stream.send.is_complete());
1831
1832 let mut buf = [0; 2];
1833 assert_eq!(stream.recv.emit(&mut buf), Ok((2, false)));
1834 assert!(!stream.recv.is_fin());
1835
1836 stream.send.ack(1, 5);
1837 assert!(!stream.send.is_complete());
1838
1839 stream.send.ack(0, 1);
1840 assert!(!stream.send.is_complete());
1841
1842 stream.send.ack_fin();
1843 assert!(stream.send.is_complete());
1844
1845 assert!(!stream.is_complete());
1846
1847 let mut buf = [0; 3];
1848 assert_eq!(stream.recv.emit(&mut buf), Ok((3, true)));
1849 assert!(stream.recv.is_fin());
1850
1851 assert!(stream.is_complete());
1852 }
1853
1854 #[test]
1855 fn send_fin_zero_length_output() {
1856 let mut buf = [0; 5];
1857
1858 let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1859
1860 assert_eq!(stream.send.write(b"hello", false), Ok(5));
1861 assert_eq!(stream.send.off_front(), 0);
1862 assert!(!stream.send.is_fin());
1863
1864 let (written, fin) = stream.send.emit(&mut buf).unwrap();
1865 assert_eq!(written, 5);
1866 assert!(!fin);
1867 assert_eq!(&buf[..written], b"hello");
1868
1869 assert_eq!(stream.send.write(b"", true), Ok(0));
1870 assert!(stream.send.is_fin());
1871 assert_eq!(stream.send.off_front(), 5);
1872
1873 let (written, fin) = stream.send.emit(&mut buf).unwrap();
1874 assert_eq!(written, 0);
1875 assert!(fin);
1876 assert_eq!(&buf[..written], b"");
1877 }
1878
1879 fn stream_send_ready(stream: &Stream) -> bool {
1880 !stream.send.is_empty() &&
1881 stream.send.off_front() < stream.send.off_back()
1882 }
1883
1884 #[test]
1885 fn send_emit() {
1886 let mut buf = [0; 5];
1887
1888 let mut stream = <Stream>::new(0, 0, 20, true, 0, DEFAULT_STREAM_WINDOW);
1889
1890 assert_eq!(stream.send.write(b"hello", false), Ok(5));
1891 assert_eq!(stream.send.write(b"world", false), Ok(5));
1892 assert_eq!(stream.send.write(b"olleh", false), Ok(5));
1893 assert_eq!(stream.send.write(b"dlrow", true), Ok(5));
1894 assert_eq!(stream.send.off_front(), 0);
1895 assert_eq!(stream.send.bufs_count(), 4);
1896
1897 assert!(stream.is_flushable());
1898
1899 assert!(stream_send_ready(&stream));
1900 assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
1901 assert_eq!(stream.send.off_front(), 4);
1902 assert_eq!(&buf[..4], b"hell");
1903
1904 assert!(stream_send_ready(&stream));
1905 assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
1906 assert_eq!(stream.send.off_front(), 8);
1907 assert_eq!(&buf[..4], b"owor");
1908
1909 assert!(stream_send_ready(&stream));
1910 assert_eq!(stream.send.emit(&mut buf[..2]), Ok((2, false)));
1911 assert_eq!(stream.send.off_front(), 10);
1912 assert_eq!(&buf[..2], b"ld");
1913
1914 assert!(stream_send_ready(&stream));
1915 assert_eq!(stream.send.emit(&mut buf[..1]), Ok((1, false)));
1916 assert_eq!(stream.send.off_front(), 11);
1917 assert_eq!(&buf[..1], b"o");
1918
1919 assert!(stream_send_ready(&stream));
1920 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
1921 assert_eq!(stream.send.off_front(), 16);
1922 assert_eq!(&buf[..5], b"llehd");
1923
1924 assert!(stream_send_ready(&stream));
1925 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((4, true)));
1926 assert_eq!(stream.send.off_front(), 20);
1927 assert_eq!(&buf[..4], b"lrow");
1928
1929 assert!(!stream.is_flushable());
1930
1931 assert!(!stream_send_ready(&stream));
1932 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((0, true)));
1933 assert_eq!(stream.send.off_front(), 20);
1934 }
1935
1936 #[test]
1937 fn send_emit_ack() {
1938 let mut buf = [0; 5];
1939
1940 let mut stream = <Stream>::new(0, 0, 20, true, 0, DEFAULT_STREAM_WINDOW);
1941
1942 assert_eq!(stream.send.write(b"hello", false), Ok(5));
1943 assert_eq!(stream.send.write(b"world", false), Ok(5));
1944 assert_eq!(stream.send.write(b"olleh", false), Ok(5));
1945 assert_eq!(stream.send.write(b"dlrow", true), Ok(5));
1946 assert_eq!(stream.send.off_front(), 0);
1947 assert_eq!(stream.send.bufs_count(), 4);
1948
1949 assert!(stream.is_flushable());
1950
1951 assert!(stream_send_ready(&stream));
1952 assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
1953 assert_eq!(stream.send.off_front(), 4);
1954 assert_eq!(&buf[..4], b"hell");
1955
1956 assert!(stream_send_ready(&stream));
1957 assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
1958 assert_eq!(stream.send.off_front(), 8);
1959 assert_eq!(&buf[..4], b"owor");
1960
1961 stream.send.ack_and_drop(0, 5);
1962 assert_eq!(stream.send.bufs_count(), 3);
1963
1964 assert!(stream_send_ready(&stream));
1965 assert_eq!(stream.send.emit(&mut buf[..2]), Ok((2, false)));
1966 assert_eq!(stream.send.off_front(), 10);
1967 assert_eq!(&buf[..2], b"ld");
1968
1969 stream.send.ack_and_drop(7, 5);
1970 assert_eq!(stream.send.bufs_count(), 3);
1971
1972 assert!(stream_send_ready(&stream));
1973 assert_eq!(stream.send.emit(&mut buf[..1]), Ok((1, false)));
1974 assert_eq!(stream.send.off_front(), 11);
1975 assert_eq!(&buf[..1], b"o");
1976
1977 assert!(stream_send_ready(&stream));
1978 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
1979 assert_eq!(stream.send.off_front(), 16);
1980 assert_eq!(&buf[..5], b"llehd");
1981
1982 stream.send.ack_and_drop(5, 7);
1983 assert_eq!(stream.send.bufs_count(), 2);
1984
1985 assert!(stream_send_ready(&stream));
1986 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((4, true)));
1987 assert_eq!(stream.send.off_front(), 20);
1988 assert_eq!(&buf[..4], b"lrow");
1989
1990 assert!(!stream.is_flushable());
1991
1992 assert!(!stream_send_ready(&stream));
1993 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((0, true)));
1994 assert_eq!(stream.send.off_front(), 20);
1995
1996 stream.send.ack_and_drop(22, 4);
1997 assert_eq!(stream.send.bufs_count(), 2);
1998
1999 stream.send.ack_and_drop(20, 1);
2000 assert_eq!(stream.send.bufs_count(), 2);
2001 }
2002
2003 #[test]
2004 fn send_emit_retransmit() {
2005 let mut buf = [0; 5];
2006
2007 let mut stream = <Stream>::new(
2008 0,
2009 0,
2010 20,
2011 true,
2012 DEFAULT_STREAM_WINDOW,
2013 DEFAULT_STREAM_WINDOW,
2014 );
2015
2016 assert_eq!(stream.send.write(b"hello", false), Ok(5));
2017 assert_eq!(stream.send.write(b"world", false), Ok(5));
2018 assert_eq!(stream.send.write(b"olleh", false), Ok(5));
2019 assert_eq!(stream.send.write(b"dlrow", true), Ok(5));
2020 assert_eq!(stream.send.off_front(), 0);
2021 assert_eq!(stream.send.bufs_count(), 4);
2022
2023 assert!(stream.is_flushable());
2024
2025 assert!(stream_send_ready(&stream));
2026 assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
2027 assert_eq!(stream.send.off_front(), 4);
2028 assert_eq!(&buf[..4], b"hell");
2029
2030 assert!(stream_send_ready(&stream));
2031 assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
2032 assert_eq!(stream.send.off_front(), 8);
2033 assert_eq!(&buf[..4], b"owor");
2034
2035 stream.send.retransmit(3, 3);
2036 assert_eq!(stream.send.off_front(), 3);
2037
2038 assert!(stream_send_ready(&stream));
2039 assert_eq!(stream.send.emit(&mut buf[..3]), Ok((3, false)));
2040 assert_eq!(stream.send.off_front(), 8);
2041 assert_eq!(&buf[..3], b"low");
2042
2043 assert!(stream_send_ready(&stream));
2044 assert_eq!(stream.send.emit(&mut buf[..2]), Ok((2, false)));
2045 assert_eq!(stream.send.off_front(), 10);
2046 assert_eq!(&buf[..2], b"ld");
2047
2048 stream.send.ack_and_drop(7, 2);
2049
2050 stream.send.retransmit(8, 2);
2051
2052 assert!(stream_send_ready(&stream));
2053 assert_eq!(stream.send.emit(&mut buf[..2]), Ok((2, false)));
2054 assert_eq!(stream.send.off_front(), 10);
2055 assert_eq!(&buf[..2], b"ld");
2056
2057 assert!(stream_send_ready(&stream));
2058 assert_eq!(stream.send.emit(&mut buf[..1]), Ok((1, false)));
2059 assert_eq!(stream.send.off_front(), 11);
2060 assert_eq!(&buf[..1], b"o");
2061
2062 assert!(stream_send_ready(&stream));
2063 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
2064 assert_eq!(stream.send.off_front(), 16);
2065 assert_eq!(&buf[..5], b"llehd");
2066
2067 stream.send.retransmit(12, 2);
2068
2069 assert!(stream_send_ready(&stream));
2070 assert_eq!(stream.send.emit(&mut buf[..2]), Ok((2, false)));
2071 assert_eq!(stream.send.off_front(), 16);
2072 assert_eq!(&buf[..2], b"le");
2073
2074 assert!(stream_send_ready(&stream));
2075 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((4, true)));
2076 assert_eq!(stream.send.off_front(), 20);
2077 assert_eq!(&buf[..4], b"lrow");
2078
2079 assert!(!stream.is_flushable());
2080
2081 assert!(!stream_send_ready(&stream));
2082 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((0, true)));
2083 assert_eq!(stream.send.off_front(), 20);
2084
2085 stream.send.retransmit(7, 12);
2086
2087 assert!(stream_send_ready(&stream));
2088 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
2089 assert_eq!(stream.send.off_front(), 12);
2090 assert_eq!(&buf[..5], b"rldol");
2091
2092 assert!(stream_send_ready(&stream));
2093 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
2094 assert_eq!(stream.send.off_front(), 17);
2095 assert_eq!(&buf[..5], b"lehdl");
2096
2097 assert!(stream_send_ready(&stream));
2098 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((2, false)));
2099 assert_eq!(stream.send.off_front(), 20);
2100 assert_eq!(&buf[..2], b"ro");
2101
2102 stream.send.ack_and_drop(12, 7);
2103
2104 stream.send.retransmit(7, 12);
2105
2106 assert!(stream_send_ready(&stream));
2107 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
2108 assert_eq!(stream.send.off_front(), 12);
2109 assert_eq!(&buf[..5], b"rldol");
2110
2111 assert!(stream_send_ready(&stream));
2112 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
2113 assert_eq!(stream.send.off_front(), 17);
2114 assert_eq!(&buf[..5], b"lehdl");
2115
2116 assert!(stream_send_ready(&stream));
2117 assert_eq!(stream.send.emit(&mut buf[..5]), Ok((2, false)));
2118 assert_eq!(stream.send.off_front(), 20);
2119 assert_eq!(&buf[..2], b"ro");
2120 }
2121
2122 #[test]
2123 fn rangebuf_split_off() {
2124 let mut buf = <RangeBuf>::from(b"helloworld", 5, true);
2125 assert_eq!(buf.start, 0);
2126 assert_eq!(buf.pos, 0);
2127 assert_eq!(buf.len, 10);
2128 assert_eq!(buf.off, 5);
2129 assert!(buf.fin);
2130
2131 assert_eq!(buf.len(), 10);
2132 assert_eq!(buf.off(), 5);
2133 assert!(buf.fin());
2134
2135 assert_eq!(&buf[..], b"helloworld");
2136
2137 buf.consume(5);
2139
2140 assert_eq!(buf.start, 0);
2141 assert_eq!(buf.pos, 5);
2142 assert_eq!(buf.len, 10);
2143 assert_eq!(buf.off, 5);
2144 assert!(buf.fin);
2145
2146 assert_eq!(buf.len(), 5);
2147 assert_eq!(buf.off(), 10);
2148 assert!(buf.fin());
2149
2150 assert_eq!(&buf[..], b"world");
2151
2152 let mut new_buf = buf.split_off(3);
2154
2155 assert_eq!(buf.start, 0);
2156 assert_eq!(buf.pos, 3);
2157 assert_eq!(buf.len, 3);
2158 assert_eq!(buf.off, 5);
2159 assert!(!buf.fin);
2160
2161 assert_eq!(buf.len(), 0);
2162 assert_eq!(buf.off(), 8);
2163 assert!(!buf.fin());
2164
2165 assert_eq!(&buf[..], b"");
2166
2167 assert_eq!(new_buf.start, 3);
2168 assert_eq!(new_buf.pos, 5);
2169 assert_eq!(new_buf.len, 7);
2170 assert_eq!(new_buf.off, 8);
2171 assert!(new_buf.fin);
2172
2173 assert_eq!(new_buf.len(), 5);
2174 assert_eq!(new_buf.off(), 10);
2175 assert!(new_buf.fin());
2176
2177 assert_eq!(&new_buf[..], b"world");
2178
2179 new_buf.consume(2);
2181
2182 assert_eq!(new_buf.start, 3);
2183 assert_eq!(new_buf.pos, 7);
2184 assert_eq!(new_buf.len, 7);
2185 assert_eq!(new_buf.off, 8);
2186 assert!(new_buf.fin);
2187
2188 assert_eq!(new_buf.len(), 3);
2189 assert_eq!(new_buf.off(), 12);
2190 assert!(new_buf.fin());
2191
2192 assert_eq!(&new_buf[..], b"rld");
2193
2194 let mut new_new_buf = new_buf.split_off(5);
2196
2197 assert_eq!(new_buf.start, 3);
2198 assert_eq!(new_buf.pos, 7);
2199 assert_eq!(new_buf.len, 5);
2200 assert_eq!(new_buf.off, 8);
2201 assert!(!new_buf.fin);
2202
2203 assert_eq!(new_buf.len(), 1);
2204 assert_eq!(new_buf.off(), 12);
2205 assert!(!new_buf.fin());
2206
2207 assert_eq!(&new_buf[..], b"r");
2208
2209 assert_eq!(new_new_buf.start, 8);
2210 assert_eq!(new_new_buf.pos, 8);
2211 assert_eq!(new_new_buf.len, 2);
2212 assert_eq!(new_new_buf.off, 13);
2213 assert!(new_new_buf.fin);
2214
2215 assert_eq!(new_new_buf.len(), 2);
2216 assert_eq!(new_new_buf.off(), 13);
2217 assert!(new_new_buf.fin());
2218
2219 assert_eq!(&new_new_buf[..], b"ld");
2220
2221 new_new_buf.consume(2);
2223
2224 assert_eq!(new_new_buf.start, 8);
2225 assert_eq!(new_new_buf.pos, 10);
2226 assert_eq!(new_new_buf.len, 2);
2227 assert_eq!(new_new_buf.off, 13);
2228 assert!(new_new_buf.fin);
2229
2230 assert_eq!(new_new_buf.len(), 0);
2231 assert_eq!(new_new_buf.off(), 15);
2232 assert!(new_new_buf.fin());
2233
2234 assert_eq!(&new_new_buf[..], b"");
2235 }
2236
2237 #[test]
2240 fn stream_limit_auto_open() {
2241 let local_tp = crate::TransportParams::default();
2242 let peer_tp = crate::TransportParams::default();
2243
2244 let mut streams = <StreamMap>::new(5, 5, 5);
2245
2246 let stream_id = 500;
2247 assert!(!is_local(stream_id, true), "stream id is peer initiated");
2248 assert!(is_bidi(stream_id), "stream id is bidirectional");
2249 assert_eq!(
2250 streams
2251 .get_or_create(stream_id, &local_tp, &peer_tp, false, true)
2252 .err(),
2253 Some(Error::StreamLimit),
2254 "stream limit should be exceeded"
2255 );
2256 }
2257
2258 #[test]
2261 fn stream_create_out_of_order() {
2262 let local_tp = crate::TransportParams::default();
2263 let peer_tp = crate::TransportParams::default();
2264
2265 let mut streams = <StreamMap>::new(5, 5, 5);
2266
2267 for stream_id in [8, 12, 4] {
2268 assert!(is_local(stream_id, false), "stream id is client initiated");
2269 assert!(is_bidi(stream_id), "stream id is bidirectional");
2270 assert!(streams
2271 .get_or_create(stream_id, &local_tp, &peer_tp, false, true)
2272 .is_ok());
2273 }
2274 }
2275
2276 #[test]
2278 fn stream_limit_edge() {
2279 let local_tp = crate::TransportParams::default();
2280 let peer_tp = crate::TransportParams::default();
2281
2282 let mut streams = <StreamMap>::new(3, 3, 3);
2283
2284 let stream_id = 8;
2286 assert!(streams
2287 .get_or_create(stream_id, &local_tp, &peer_tp, false, true)
2288 .is_ok());
2289
2290 let stream_id = 12;
2292 assert_eq!(
2293 streams
2294 .get_or_create(stream_id, &local_tp, &peer_tp, false, true)
2295 .err(),
2296 Some(Error::StreamLimit)
2297 );
2298 }
2299
2300 fn cycle_stream_priority(stream_id: u64, streams: &mut StreamMap) {
2301 let key = streams.get(stream_id).unwrap().priority_key.clone();
2302 streams.update_priority(&key.clone(), &key);
2303 }
2304
2305 #[test]
2306 fn writable_prioritized_default_priority() {
2307 let local_tp = crate::TransportParams::default();
2308 let peer_tp = crate::TransportParams {
2309 initial_max_stream_data_bidi_local: 100,
2310 initial_max_stream_data_uni: 100,
2311 ..Default::default()
2312 };
2313
2314 let mut streams = StreamMap::new(100, 100, 100);
2315
2316 for id in [0, 4, 8, 12] {
2317 assert!(streams
2318 .get_or_create(id, &local_tp, &peer_tp, false, true)
2319 .is_ok());
2320 }
2321
2322 let walk_1: Vec<u64> = streams.writable().collect();
2323 cycle_stream_priority(*walk_1.first().unwrap(), &mut streams);
2324 let walk_2: Vec<u64> = streams.writable().collect();
2325 cycle_stream_priority(*walk_2.first().unwrap(), &mut streams);
2326 let walk_3: Vec<u64> = streams.writable().collect();
2327 cycle_stream_priority(*walk_3.first().unwrap(), &mut streams);
2328 let walk_4: Vec<u64> = streams.writable().collect();
2329 cycle_stream_priority(*walk_4.first().unwrap(), &mut streams);
2330 let walk_5: Vec<u64> = streams.writable().collect();
2331
2332 assert_eq!(walk_1, vec![0, 4, 8, 12]);
2335 assert_eq!(walk_2, vec![4, 8, 12, 0]);
2336 assert_eq!(walk_3, vec![8, 12, 0, 4]);
2337 assert_eq!(walk_4, vec![12, 0, 4, 8,]);
2338 assert_eq!(walk_5, vec![0, 4, 8, 12]);
2339 }
2340
2341 #[test]
2342 fn writable_prioritized_insert_order() {
2343 let local_tp = crate::TransportParams::default();
2344 let peer_tp = crate::TransportParams {
2345 initial_max_stream_data_bidi_local: 100,
2346 initial_max_stream_data_uni: 100,
2347 ..Default::default()
2348 };
2349
2350 let mut streams = StreamMap::new(100, 100, 100);
2351
2352 for id in [12, 4, 8, 0] {
2355 assert!(streams
2356 .get_or_create(id, &local_tp, &peer_tp, false, true)
2357 .is_ok());
2358 }
2359
2360 let walk_1: Vec<u64> = streams.writable().collect();
2361 cycle_stream_priority(*walk_1.first().unwrap(), &mut streams);
2362 let walk_2: Vec<u64> = streams.writable().collect();
2363 cycle_stream_priority(*walk_2.first().unwrap(), &mut streams);
2364 let walk_3: Vec<u64> = streams.writable().collect();
2365 cycle_stream_priority(*walk_3.first().unwrap(), &mut streams);
2366 let walk_4: Vec<u64> = streams.writable().collect();
2367 cycle_stream_priority(*walk_4.first().unwrap(), &mut streams);
2368 let walk_5: Vec<u64> = streams.writable().collect();
2369 assert_eq!(walk_1, vec![12, 4, 8, 0]);
2370 assert_eq!(walk_2, vec![4, 8, 0, 12]);
2371 assert_eq!(walk_3, vec![8, 0, 12, 4,]);
2372 assert_eq!(walk_4, vec![0, 12, 4, 8]);
2373 assert_eq!(walk_5, vec![12, 4, 8, 0]);
2374 }
2375
2376 #[test]
2377 fn writable_prioritized_mixed_urgency() {
2378 let local_tp = crate::TransportParams::default();
2379 let peer_tp = crate::TransportParams {
2380 initial_max_stream_data_bidi_local: 100,
2381 initial_max_stream_data_uni: 100,
2382 ..Default::default()
2383 };
2384
2385 let mut streams = <StreamMap>::new(100, 100, 100);
2386
2387 let input = vec![
2390 (0, 100),
2391 (4, 90),
2392 (8, 80),
2393 (12, 70),
2394 (16, 60),
2395 (20, 50),
2396 (24, 40),
2397 (28, 30),
2398 (32, 20),
2399 (36, 10),
2400 (40, 0),
2401 ];
2402
2403 for (id, urgency) in input.clone() {
2404 let stream = streams
2407 .get_or_create(id, &local_tp, &peer_tp, false, true)
2408 .unwrap();
2409
2410 stream.urgency = urgency;
2411
2412 let new_priority_key = Arc::new(StreamPriorityKey {
2413 urgency: stream.urgency,
2414 incremental: stream.incremental,
2415 id,
2416 ..Default::default()
2417 });
2418
2419 let old_priority_key = std::mem::replace(
2420 &mut stream.priority_key,
2421 new_priority_key.clone(),
2422 );
2423
2424 streams.update_priority(&old_priority_key, &new_priority_key);
2425 }
2426
2427 let walk_1: Vec<u64> = streams.writable().collect();
2428 assert_eq!(walk_1, vec![40, 36, 32, 28, 24, 20, 16, 12, 8, 4, 0]);
2429
2430 for (id, urgency) in input {
2432 let stream = streams
2435 .get_or_create(id, &local_tp, &peer_tp, false, true)
2436 .unwrap();
2437
2438 stream.urgency = urgency;
2439
2440 let new_priority_key = Arc::new(StreamPriorityKey {
2441 urgency: stream.urgency,
2442 incremental: stream.incremental,
2443 id,
2444 ..Default::default()
2445 });
2446
2447 let old_priority_key = std::mem::replace(
2448 &mut stream.priority_key,
2449 new_priority_key.clone(),
2450 );
2451
2452 streams.update_priority(&old_priority_key, &new_priority_key);
2453 }
2454
2455 let walk_2: Vec<u64> = streams.writable().collect();
2456 assert_eq!(walk_2, vec![40, 36, 32, 28, 24, 20, 16, 12, 8, 4, 0]);
2457
2458 streams.collect(24, true);
2460
2461 let walk_3: Vec<u64> = streams.writable().collect();
2462 assert_eq!(walk_3, vec![40, 36, 32, 28, 20, 16, 12, 8, 4, 0]);
2463
2464 streams.collect(40, true);
2465 streams.collect(0, true);
2466
2467 let walk_4: Vec<u64> = streams.writable().collect();
2468 assert_eq!(walk_4, vec![36, 32, 28, 20, 16, 12, 8, 4]);
2469
2470 streams
2472 .get_or_create(44, &local_tp, &peer_tp, false, true)
2473 .unwrap();
2474
2475 let walk_5: Vec<u64> = streams.writable().collect();
2476 assert_eq!(walk_5, vec![36, 32, 28, 20, 16, 12, 8, 4, 44]);
2477 }
2478
2479 #[test]
2480 fn writable_prioritized_mixed_urgencies_incrementals() {
2481 let local_tp = crate::TransportParams::default();
2482 let peer_tp = crate::TransportParams {
2483 initial_max_stream_data_bidi_local: 100,
2484 initial_max_stream_data_uni: 100,
2485 ..Default::default()
2486 };
2487
2488 let mut streams = StreamMap::new(100, 100, 100);
2489
2490 let input = vec![
2492 (0, 100),
2493 (4, 20),
2494 (8, 100),
2495 (12, 20),
2496 (16, 90),
2497 (20, 25),
2498 (24, 90),
2499 (28, 30),
2500 (32, 80),
2501 (36, 20),
2502 (40, 0),
2503 ];
2504
2505 for (id, urgency) in input.clone() {
2506 let stream = streams
2509 .get_or_create(id, &local_tp, &peer_tp, false, true)
2510 .unwrap();
2511
2512 stream.urgency = urgency;
2513
2514 let new_priority_key = Arc::new(StreamPriorityKey {
2515 urgency: stream.urgency,
2516 incremental: stream.incremental,
2517 id,
2518 ..Default::default()
2519 });
2520
2521 let old_priority_key = std::mem::replace(
2522 &mut stream.priority_key,
2523 new_priority_key.clone(),
2524 );
2525
2526 streams.update_priority(&old_priority_key, &new_priority_key);
2527 }
2528
2529 let walk_1: Vec<u64> = streams.writable().collect();
2530 cycle_stream_priority(4, &mut streams);
2531 cycle_stream_priority(16, &mut streams);
2532 cycle_stream_priority(0, &mut streams);
2533 let walk_2: Vec<u64> = streams.writable().collect();
2534 cycle_stream_priority(12, &mut streams);
2535 cycle_stream_priority(24, &mut streams);
2536 cycle_stream_priority(8, &mut streams);
2537 let walk_3: Vec<u64> = streams.writable().collect();
2538 cycle_stream_priority(36, &mut streams);
2539 cycle_stream_priority(16, &mut streams);
2540 cycle_stream_priority(0, &mut streams);
2541 let walk_4: Vec<u64> = streams.writable().collect();
2542 cycle_stream_priority(4, &mut streams);
2543 cycle_stream_priority(24, &mut streams);
2544 cycle_stream_priority(8, &mut streams);
2545 let walk_5: Vec<u64> = streams.writable().collect();
2546 cycle_stream_priority(12, &mut streams);
2547 cycle_stream_priority(16, &mut streams);
2548 cycle_stream_priority(0, &mut streams);
2549 let walk_6: Vec<u64> = streams.writable().collect();
2550 cycle_stream_priority(36, &mut streams);
2551 cycle_stream_priority(24, &mut streams);
2552 cycle_stream_priority(8, &mut streams);
2553 let walk_7: Vec<u64> = streams.writable().collect();
2554 cycle_stream_priority(4, &mut streams);
2555 cycle_stream_priority(16, &mut streams);
2556 cycle_stream_priority(0, &mut streams);
2557 let walk_8: Vec<u64> = streams.writable().collect();
2558 cycle_stream_priority(12, &mut streams);
2559 cycle_stream_priority(24, &mut streams);
2560 cycle_stream_priority(8, &mut streams);
2561 let walk_9: Vec<u64> = streams.writable().collect();
2562 cycle_stream_priority(36, &mut streams);
2563 cycle_stream_priority(16, &mut streams);
2564 cycle_stream_priority(0, &mut streams);
2565
2566 assert_eq!(walk_1, vec![40, 4, 12, 36, 20, 28, 32, 16, 24, 0, 8]);
2567 assert_eq!(walk_2, vec![40, 12, 36, 4, 20, 28, 32, 24, 16, 8, 0]);
2568 assert_eq!(walk_3, vec![40, 36, 4, 12, 20, 28, 32, 16, 24, 0, 8]);
2569 assert_eq!(walk_4, vec![40, 4, 12, 36, 20, 28, 32, 24, 16, 8, 0]);
2570 assert_eq!(walk_5, vec![40, 12, 36, 4, 20, 28, 32, 16, 24, 0, 8]);
2571 assert_eq!(walk_6, vec![40, 36, 4, 12, 20, 28, 32, 24, 16, 8, 0]);
2572 assert_eq!(walk_7, vec![40, 4, 12, 36, 20, 28, 32, 16, 24, 0, 8]);
2573 assert_eq!(walk_8, vec![40, 12, 36, 4, 20, 28, 32, 24, 16, 8, 0]);
2574 assert_eq!(walk_9, vec![40, 36, 4, 12, 20, 28, 32, 16, 24, 0, 8]);
2575
2576 streams.collect(20, true);
2578
2579 let walk_10: Vec<u64> = streams.writable().collect();
2580 assert_eq!(walk_10, vec![40, 4, 12, 36, 28, 32, 24, 16, 8, 0]);
2581
2582 let stream = streams
2584 .get_or_create(44, &local_tp, &peer_tp, false, true)
2585 .unwrap();
2586
2587 stream.urgency = 20;
2588 stream.incremental = true;
2589
2590 let new_priority_key = Arc::new(StreamPriorityKey {
2591 urgency: stream.urgency,
2592 incremental: stream.incremental,
2593 id: 44,
2594 ..Default::default()
2595 });
2596
2597 let old_priority_key =
2598 std::mem::replace(&mut stream.priority_key, new_priority_key.clone());
2599
2600 streams.update_priority(&old_priority_key, &new_priority_key);
2601
2602 let walk_11: Vec<u64> = streams.writable().collect();
2603 assert_eq!(walk_11, vec![40, 4, 12, 36, 44, 28, 32, 24, 16, 8, 0]);
2604 }
2605
2606 #[test]
2607 fn priority_tree_dupes() {
2608 let mut prioritized_writable: RBTree<StreamWritablePriorityAdapter> =
2609 Default::default();
2610
2611 for id in [0, 4, 8, 12] {
2612 let s = Arc::new(StreamPriorityKey {
2613 urgency: 0,
2614 incremental: false,
2615 id,
2616 ..Default::default()
2617 });
2618
2619 prioritized_writable.insert(s);
2620 }
2621
2622 let walk_1: Vec<u64> =
2623 prioritized_writable.iter().map(|s| s.id).collect();
2624 assert_eq!(walk_1, vec![0, 4, 8, 12]);
2625
2626 for id in [0, 4, 8, 12] {
2629 let s = Arc::new(StreamPriorityKey {
2630 urgency: 0,
2631 incremental: false,
2632 id,
2633 ..Default::default()
2634 });
2635
2636 prioritized_writable.insert(s);
2637 }
2638
2639 let walk_2: Vec<u64> =
2640 prioritized_writable.iter().map(|s| s.id).collect();
2641 assert_eq!(walk_2, vec![0, 0, 4, 4, 8, 8, 12, 12]);
2642 }
2643
2644 #[test]
2645 fn retransmit_returns_zero_when_already_acked() {
2646 let mut stream = <Stream>::new(0, 15, 15, true, 0, 15);
2647
2648 assert_eq!(stream.send.write(b"hello", false), Ok(5));
2650 assert_eq!(stream.send.buffered_bytes(), 5);
2651
2652 let mut buf = [0; 10];
2653 let (written, _) = stream.send.emit(&mut buf).unwrap();
2654 assert_eq!(written, 5);
2655 assert_eq!(stream.send.buffered_bytes(), 0);
2656
2657 let retransmitted = stream.send.retransmit(0, 5);
2659 assert_eq!(retransmitted, 5);
2660 assert_eq!(stream.send.buffered_bytes(), 5);
2661
2662 stream.send.ack_and_drop(0, 5);
2664 assert_eq!(stream.send.buffered_bytes(), 0);
2665
2666 let retransmitted = stream.send.retransmit(0, 5);
2668 assert_eq!(retransmitted, 0);
2669 assert_eq!(stream.send.buffered_bytes(), 0);
2670 }
2671
2672 #[test]
2673 fn retransmit_returns_partial_when_some_acked() {
2674 let mut stream = <Stream>::new(0, 15, 15, true, 0, 15);
2675
2676 assert_eq!(stream.send.write(b"helloworld", false), Ok(10));
2678 assert_eq!(stream.send.buffered_bytes(), 10);
2679
2680 let mut buf = [0; 10];
2681 let (written, _) = stream.send.emit(&mut buf).unwrap();
2682 assert_eq!(written, 10);
2683 assert_eq!(stream.send.buffered_bytes(), 0);
2684
2685 let retransmitted = stream.send.retransmit(0, 10);
2687 assert_eq!(retransmitted, 10);
2688 assert_eq!(stream.send.buffered_bytes(), 10);
2689
2690 let dropped = stream.send.ack_and_drop(0, 5);
2692 assert_eq!(dropped, 5);
2693 assert_eq!(stream.send.buffered_bytes(), 5);
2694
2695 let retransmitted = stream.send.retransmit(0, 10);
2698 assert_eq!(retransmitted, 0); assert_eq!(stream.send.buffered_bytes(), 5);
2700 }
2701
2702 #[test]
2703 fn ack_and_drop_decrements_len_and_returns_dropped() {
2704 let mut stream = <Stream>::new(0, 15, 15, true, 0, 15);
2705
2706 assert_eq!(stream.send.write(b"hello", false), Ok(5));
2708 assert_eq!(stream.send.buffered_bytes(), 5);
2709
2710 let mut buf = [0; 10];
2712 let (written, _) = stream.send.emit(&mut buf).unwrap();
2713 assert_eq!(written, 5);
2714 assert_eq!(stream.send.buffered_bytes(), 0);
2715
2716 let retransmitted = stream.send.retransmit(0, 5);
2718 assert_eq!(retransmitted, 5);
2719 assert_eq!(stream.send.buffered_bytes(), 5);
2720
2721 let dropped = stream.send.ack_and_drop(0, 5);
2723 assert_eq!(dropped, 5);
2724 assert_eq!(stream.send.buffered_bytes(), 0);
2725 }
2726
2727 #[test]
2728 fn ack_and_drop_partial_buffer() {
2729 let mut stream = <Stream>::new(0, 30, 30, true, 0, 30);
2730
2731 assert_eq!(stream.send.write(b"hello", false), Ok(5));
2733 assert_eq!(stream.send.write(b"world", false), Ok(5));
2734 assert_eq!(stream.send.buffered_bytes(), 10);
2735
2736 let mut buf = [0; 10];
2737 let (written, _) = stream.send.emit(&mut buf).unwrap();
2738 assert_eq!(written, 10);
2739 assert_eq!(stream.send.buffered_bytes(), 0);
2740
2741 let retransmitted = stream.send.retransmit(0, 10);
2743 assert_eq!(retransmitted, 10);
2744 assert_eq!(stream.send.buffered_bytes(), 10);
2745
2746 let dropped = stream.send.ack_and_drop(0, 5);
2748 assert_eq!(dropped, 5);
2749 assert_eq!(stream.send.buffered_bytes(), 5);
2750
2751 let dropped = stream.send.ack_and_drop(5, 5);
2753 assert_eq!(dropped, 5);
2754 assert_eq!(stream.send.buffered_bytes(), 0);
2755 }
2756
2757 #[test]
2758 fn ack_and_drop_returns_zero_when_nothing_dropped() {
2759 let mut stream = <Stream>::new(0, 15, 15, true, 0, 15);
2760
2761 assert_eq!(stream.send.write(b"hello", false), Ok(5));
2763 let mut buf = [0; 10];
2764 let (written, _) = stream.send.emit(&mut buf).unwrap();
2765 assert_eq!(written, 5);
2766
2767 let dropped = stream.send.ack_and_drop(0, 5);
2770 assert_eq!(dropped, 0);
2771 assert_eq!(stream.send.buffered_bytes(), 0);
2772 }
2773
2774 #[test]
2775 fn cache_consistency_through_full_lifecycle() {
2776 let mut streams = <StreamMap>::new(5, 5, 15);
2780
2781 let local_params = crate::TransportParams {
2783 initial_max_data: 30,
2784 initial_max_stream_data_bidi_local: 15,
2785 initial_max_stream_data_bidi_remote: 15,
2786 initial_max_stream_data_uni: 10,
2787 initial_max_streams_bidi: 5,
2788 initial_max_streams_uni: 5,
2789 ..Default::default()
2790 };
2791 let peer_params = local_params.clone();
2792
2793 streams.update_peer_max_streams_bidi(5);
2795 streams.update_peer_max_streams_uni(5);
2796
2797 let stream_id = 0u64;
2798
2799 {
2801 let stream = streams
2802 .get_or_create(
2803 stream_id,
2804 &local_params,
2805 &peer_params,
2806 true,
2807 false,
2808 )
2809 .unwrap();
2810 assert_eq!(stream.send.write(b"hello", false), Ok(5));
2811 }
2812 streams.add_tx_buffered(5);
2813 assert_eq!(streams.get(stream_id).unwrap().send.buffered_bytes(), 5);
2814 assert_eq!(streams.tx_buffered(), 5);
2815 assert!(streams.tx_buffered_is_consistent());
2816
2817 let mut buf = [0; 10];
2819 let written = {
2820 let stream = streams.get_mut(stream_id).unwrap();
2821 let (written, _) = stream.send.emit(&mut buf).unwrap();
2822 written
2823 };
2824 assert_eq!(written, 5);
2825 streams.sub_tx_buffered(5);
2826 assert_eq!(streams.get(stream_id).unwrap().send.buffered_bytes(), 0);
2827 assert_eq!(streams.tx_buffered(), 0);
2828 assert!(streams.tx_buffered_is_consistent());
2829
2830 let retransmitted = {
2832 let stream = streams.get_mut(stream_id).unwrap();
2833 stream.send.retransmit(0, 5)
2834 };
2835 assert_eq!(retransmitted, 5);
2836 streams.add_tx_buffered(retransmitted);
2837 assert_eq!(streams.get(stream_id).unwrap().send.buffered_bytes(), 5);
2838 assert_eq!(streams.tx_buffered(), 5);
2839 assert!(streams.tx_buffered_is_consistent());
2840
2841 let dropped = {
2844 let stream = streams.get_mut(stream_id).unwrap();
2845 stream.send.ack_and_drop(0, 5)
2846 };
2847 assert_eq!(dropped, 5);
2848 streams.sub_tx_buffered(dropped);
2849 assert_eq!(streams.get(stream_id).unwrap().send.buffered_bytes(), 0);
2850 assert_eq!(streams.tx_buffered(), 0);
2851 assert!(streams.tx_buffered_is_consistent());
2852 }
2853
2854 #[test]
2855 fn send_buf_len_reflects_buffered_data() {
2856 let mut stream = <Stream>::new(0, 15, 15, true, 0, 15);
2857
2858 assert_eq!(stream.send.buffered_bytes(), 0);
2860
2861 assert_eq!(stream.send.write(b"hello", false), Ok(5));
2863 assert_eq!(stream.send.buffered_bytes(), 5);
2864
2865 let mut buf = [0; 10];
2867 let (written, _) = stream.send.emit(&mut buf).unwrap();
2868 assert_eq!(written, 5);
2869 assert_eq!(stream.send.buffered_bytes(), 0);
2870
2871 let retransmitted = stream.send.retransmit(0, 5);
2873 assert_eq!(retransmitted, 5);
2874 assert_eq!(stream.send.buffered_bytes(), 5);
2875
2876 let dropped = stream.send.ack_and_drop(0, 5);
2878 assert_eq!(dropped, 5);
2879 assert_eq!(stream.send.buffered_bytes(), 0);
2880 }
2881}
2882
2883mod recv_buf;
2884mod send_buf;