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