1use std::cmp;
28
29use std::collections::BTreeMap;
30use std::collections::VecDeque;
31
32use std::time::Duration;
33use std::time::Instant;
34
35use crate::stream::RecvAction;
36use crate::stream::RecvBufResetReturn;
37use crate::Error;
38use crate::Result;
39
40use crate::flowcontrol;
41
42use crate::range_buf::RangeBuf;
43
44#[derive(Debug, Default)]
50pub struct RecvBuf {
51 data: BTreeMap<u64, RangeBuf>,
54
55 off: u64,
57
58 len: u64,
60
61 flow_control: flowcontrol::FlowControl,
63
64 fin_off: Option<u64>,
66
67 error: Option<u64>,
69
70 drain: bool,
72}
73
74impl RecvBuf {
75 pub fn new(max_data: u64, initial_window: u64, max_window: u64) -> RecvBuf {
77 RecvBuf {
78 flow_control: flowcontrol::FlowControl::new(
79 max_data,
80 initial_window,
81 max_window,
82 ),
83 ..RecvBuf::default()
84 }
85 }
86
87 pub fn write(&mut self, buf: RangeBuf) -> Result<()> {
93 if buf.max_off() > self.max_data() {
94 return Err(Error::FlowControl);
95 }
96
97 if let Some(fin_off) = self.fin_off {
98 if buf.max_off() > fin_off {
100 return Err(Error::FinalSize);
101 }
102
103 if buf.fin() && fin_off != buf.max_off() {
105 return Err(Error::FinalSize);
106 }
107 }
108
109 if buf.fin() && buf.max_off() < self.len {
111 return Err(Error::FinalSize);
112 }
113
114 if self.fin_off.is_some() && buf.is_empty() {
117 return Ok(());
118 }
119
120 if buf.fin() {
121 self.fin_off = Some(buf.max_off());
122 }
123
124 if !buf.fin() && buf.is_empty() {
128 self.len = cmp::max(self.len, buf.max_off());
129
130 if self.drain {
131 self.off = self.len;
133 }
134
135 return Ok(());
136 }
137
138 if self.off >= buf.max_off() {
141 if !buf.is_empty() {
147 return Ok(());
148 }
149 }
150
151 let mut tmp_bufs = VecDeque::with_capacity(2);
152 tmp_bufs.push_back(buf);
153
154 'tmp: while let Some(mut buf) = tmp_bufs.pop_front() {
155 if self.off_front() > buf.off() {
161 buf = buf.split_off((self.off_front() - buf.off()) as usize);
162 }
163
164 if buf.off() < self.max_off() || buf.is_empty() {
171 for (_, b) in self.data.range(buf.off()..) {
172 let off = buf.off();
173
174 if b.off() > buf.max_off() {
176 break;
177 }
178
179 if off >= b.off() && buf.max_off() <= b.max_off() {
181 continue 'tmp;
182 }
183
184 if off >= b.off() && off < b.max_off() {
186 buf = buf.split_off((b.max_off() - off) as usize);
187 }
188
189 if off < b.off() && buf.max_off() > b.off() {
191 tmp_bufs
192 .push_back(buf.split_off((b.off() - off) as usize));
193 }
194 }
195 }
196
197 self.len = cmp::max(self.len, buf.max_off());
198
199 if !self.drain {
200 self.data.insert(buf.max_off(), buf);
201 } else {
202 self.off = self.len;
204 }
205 }
206
207 Ok(())
208 }
209
210 #[inline]
221 pub fn emit(&mut self, mut out: &mut [u8]) -> Result<(usize, bool)> {
222 self.emit_or_discard(RecvAction::Emit { out: &mut out })
223 }
224
225 pub fn emit_or_discard<B: bytes::BufMut>(
240 &mut self, mut action: RecvAction<B>,
241 ) -> Result<(usize, bool)> {
242 let mut len = 0;
243 let mut cap = match &action {
244 RecvAction::Emit { out } => out.remaining_mut(),
245 RecvAction::Discard { len } => *len,
246 };
247
248 if !self.ready() {
249 return Err(Error::Done);
250 }
251
252 if let Some(e) = self.error {
255 self.data.clear();
256 return Err(Error::StreamReset(e));
257 }
258
259 while self.ready() {
260 let mut entry = match self.data.first_entry() {
261 Some(entry) => entry,
262 None => break,
263 };
264
265 let buf = entry.get_mut();
266
267 if cap == 0 && !buf.is_empty() {
268 break;
269 }
270
271 let buf_len = cmp::min(buf.len(), cap);
272
273 if let RecvAction::Emit { ref mut out } = action {
275 debug_assert!(
280 cap <= out.remaining_mut(),
281 "We updated `cap` incorrectly"
282 );
283 out.put_slice(&buf[..buf_len])
284 }
285
286 self.off += buf_len as u64;
287
288 len += buf_len;
289 cap -= buf_len;
290
291 if buf_len < buf.len() {
292 buf.consume(buf_len);
293
294 break;
296 }
297
298 entry.remove();
299 }
300
301 self.flow_control.add_consumed(len as u64);
303
304 Ok((len, self.is_fin()))
305 }
306
307 pub fn reset(
309 &mut self, error_code: u64, final_size: u64,
310 ) -> Result<RecvBufResetReturn> {
311 if let Some(fin_off) = self.fin_off {
313 if fin_off != final_size {
314 return Err(Error::FinalSize);
315 }
316 }
317
318 if final_size < self.len {
320 return Err(Error::FinalSize);
321 }
322
323 if self.error.is_some() {
324 return Ok(RecvBufResetReturn::zero());
326 }
327
328 let result = RecvBufResetReturn {
331 max_data_delta: final_size - self.len,
332 consumed_flowcontrol: final_size - self.off,
333 };
334
335 self.error = Some(error_code);
336
337 self.off = final_size;
339
340 self.data.clear();
341
342 let buf = RangeBuf::from(b"", final_size, true);
345 self.write(buf)?;
346
347 Ok(result)
348 }
349
350 pub fn update_max_data(&mut self, now: Instant) {
352 self.flow_control.update_max_data(now);
353 }
354
355 pub fn max_data_next(&mut self) -> u64 {
357 self.flow_control.max_data_next()
358 }
359
360 pub fn max_data(&self) -> u64 {
362 self.flow_control.max_data()
363 }
364
365 pub fn window(&self) -> u64 {
367 self.flow_control.window()
368 }
369
370 pub fn autotune_window(&mut self, now: Instant, rtt: Duration) {
372 self.flow_control.autotune_window(now, rtt);
373 }
374
375 pub fn shutdown(&mut self) -> Result<u64> {
379 if self.drain {
380 return Err(Error::Done);
381 }
382
383 self.drain = true;
384
385 self.data.clear();
386
387 let consumed = self.max_off() - self.off;
388 self.off = self.max_off();
389
390 Ok(consumed)
391 }
392
393 pub fn off_front(&self) -> u64 {
395 self.off
396 }
397
398 pub fn almost_full(&self) -> bool {
400 self.fin_off.is_none() && self.flow_control.should_update_max_data()
401 }
402
403 pub fn max_off(&self) -> u64 {
405 self.len
406 }
407
408 pub fn is_fin(&self) -> bool {
413 if self.fin_off == Some(self.off) {
414 return true;
415 }
416
417 false
418 }
419
420 pub fn is_draining(&self) -> bool {
422 self.drain
423 }
424
425 pub fn ready(&self) -> bool {
427 let (_, buf) = match self.data.first_key_value() {
428 Some(v) => v,
429 None => return false,
430 };
431
432 buf.off() == self.off
433 }
434
435 pub fn readable_len(&self, max_len: usize) -> usize {
443 let mut contiguous = 0usize;
444 let mut next_off = self.off;
445
446 for buf in self.data.values() {
450 if buf.off() != next_off {
451 break;
452 }
453
454 contiguous = contiguous.saturating_add(buf.len()).min(max_len);
455 next_off = buf.max_off();
456
457 if contiguous == max_len {
458 break;
459 }
460 }
461
462 contiguous
463 }
464
465 #[cfg(test)]
466 pub(crate) fn flow_control_for_tests(&self) -> &flowcontrol::FlowControl {
467 &self.flow_control
468 }
469}
470
471#[cfg(test)]
472mod tests {
473 use super::*;
474
475 const DEFAULT_STREAM_WINDOW: u64 = 32 * 1024;
477 use bytes::BufMut as _;
478 use rstest::rstest;
479
480 fn assert_emit_discard(
498 recv: &mut RecvBuf, emit: bool, target_len: usize, result_len: usize,
499 is_fin: bool, test_bytes: Option<&[u8]>,
500 ) {
501 let mut buf = Vec::<u8>::with_capacity(512).limit(target_len);
502 let action = if emit {
503 RecvAction::Emit { out: &mut buf }
504 } else {
505 RecvAction::Discard { len: target_len }
506 };
507
508 let (read, fin) = recv.emit_or_discard(action).unwrap();
509
510 let buf = buf.into_inner();
511 if emit {
512 assert_eq!(buf.len(), read);
513 if let Some(v) = test_bytes {
514 assert_eq!(&buf, v);
515 }
516 }
517
518 assert_eq!(read, result_len);
519 assert_eq!(is_fin, fin);
520 }
521
522 fn assert_emit_discard_done(recv: &mut RecvBuf, emit: bool) {
524 let mut buf = [0u8; 32];
525 let action = if emit {
526 RecvAction::Emit {
527 out: &mut buf.as_mut_slice(),
528 }
529 } else {
530 RecvAction::Discard { len: 32 }
531 };
532 assert_eq!(recv.emit_or_discard(action), Err(Error::Done));
533 }
534
535 #[rstest]
536 fn empty_read(#[values(true, false)] emit: bool) {
537 let mut recv =
538 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
539 assert_eq!(recv.len, 0);
540
541 assert_emit_discard_done(&mut recv, emit);
542 }
543
544 #[rstest]
545 fn empty_stream_frame(#[values(true, false)] emit: bool) {
546 let mut recv =
547 RecvBuf::new(15, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
548 assert_eq!(recv.len, 0);
549
550 let buf = RangeBuf::from(b"hello", 0, false);
551 assert!(recv.write(buf).is_ok());
552 assert_eq!(recv.len, 5);
553 assert_eq!(recv.off, 0);
554 assert_eq!(recv.data.len(), 1);
555
556 assert_emit_discard(&mut recv, emit, 32, 5, false, None);
557
558 let buf = RangeBuf::from(b"", 10, false);
560 assert!(recv.write(buf).is_ok());
561 assert_eq!(recv.len, 10);
562 assert_eq!(recv.off, 5);
563 assert_eq!(recv.data.len(), 0);
564
565 let buf = RangeBuf::from(b"", 16, false);
567 assert_eq!(recv.write(buf), Err(Error::FlowControl));
568
569 let buf = RangeBuf::from(b"", 5, true);
572 assert_eq!(recv.write(buf), Err(Error::FinalSize));
573
574 let buf = RangeBuf::from(b"", 10, true);
576 assert!(recv.write(buf).is_ok());
577 assert_eq!(recv.len, 10);
578 assert_eq!(recv.off, 5);
579 assert_eq!(recv.data.len(), 1);
580
581 let buf = RangeBuf::from(b"", 10, true);
583 assert!(recv.write(buf).is_ok());
584 assert_eq!(recv.len, 10);
585 assert_eq!(recv.off, 5);
586 assert_eq!(recv.data.len(), 1);
587
588 let buf = RangeBuf::from(b"aa", 8, true);
590 assert!(recv.write(buf).is_ok());
591 assert_eq!(recv.len, 10);
592 assert_eq!(recv.off, 5);
593 assert_eq!(recv.data.len(), 1);
594
595 let buf = RangeBuf::from(b"aa", 3, true);
597 assert_eq!(recv.write(buf), Err(Error::FinalSize));
598
599 let buf = RangeBuf::from(b"", 11, true);
601 assert_eq!(recv.write(buf), Err(Error::FinalSize));
602 let buf = RangeBuf::from(b"", 9, true);
603 assert_eq!(recv.write(buf), Err(Error::FinalSize));
604
605 assert_emit_discard_done(&mut recv, emit);
608 }
609
610 #[rstest]
611 fn ordered_read(#[values(true, false)] emit: bool) {
612 let mut recv =
613 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
614 assert_eq!(recv.len, 0);
615
616 let first = RangeBuf::from(b"hello", 0, false);
617 let second = RangeBuf::from(b"world", 5, false);
618 let third = RangeBuf::from(b"something", 10, true);
619
620 assert!(recv.write(second).is_ok());
621 assert_eq!(recv.len, 10);
622 assert_eq!(recv.off, 0);
623
624 assert_emit_discard_done(&mut recv, emit);
625
626 assert!(recv.write(third).is_ok());
627 assert_eq!(recv.len, 19);
628 assert_eq!(recv.off, 0);
629
630 assert_emit_discard_done(&mut recv, emit);
631
632 assert!(recv.write(first).is_ok());
633 assert_eq!(recv.len, 19);
634 assert_eq!(recv.off, 0);
635
636 assert_emit_discard(
637 &mut recv,
638 emit,
639 32,
640 19,
641 true,
642 Some(b"helloworldsomething"),
643 );
644 assert_eq!(recv.len, 19);
645 assert_eq!(recv.off, 19);
646
647 assert_emit_discard_done(&mut recv, emit);
648 }
649
650 #[test]
651 fn readable_len() {
653 let mut recv =
654 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
655
656 assert_eq!(recv.readable_len(64 * 1024), 0);
658
659 assert!(recv.write(RangeBuf::from(b"hello", 0, false)).is_ok());
661 assert!(recv.write(RangeBuf::from(b"something", 10, false)).is_ok());
662 assert_eq!(recv.readable_len(64 * 1024), 5);
663
664 assert!(recv.write(RangeBuf::from(b"world", 5, false)).is_ok());
666 assert_eq!(recv.readable_len(64 * 1024), 19);
667 assert_eq!(recv.readable_len(10), 10);
668 }
669
670 #[rstest]
672 fn shutdown(#[values(true, false)] emit: bool) {
673 let mut recv =
674 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
675 assert_eq!(recv.len, 0);
676
677 let first = RangeBuf::from(b"hello", 0, false);
678 let second = RangeBuf::from(b"world", 5, false);
679 let third = RangeBuf::from(b"something", 10, false);
680
681 assert!(recv.write(second).is_ok());
682 assert_eq!(recv.len, 10);
683 assert_eq!(recv.off, 0);
684
685 assert_emit_discard_done(&mut recv, emit);
686
687 assert_eq!(recv.shutdown(), Ok(10));
689 assert_eq!(recv.len, 10);
690 assert_eq!(recv.off, 10);
691 assert_eq!(recv.data.len(), 0);
692
693 assert_emit_discard_done(&mut recv, emit);
694
695 assert!(recv.write(first).is_ok());
697 assert_eq!(recv.len, 10);
698 assert_eq!(recv.off, 10);
699 assert_eq!(recv.data.len(), 0);
700
701 assert!(recv.write(third).is_ok());
704 assert_eq!(recv.len, 19);
705 assert_eq!(recv.off, 19);
706 assert_eq!(recv.data.len(), 0);
707
708 assert_emit_discard_done(&mut recv, emit);
710 assert_eq!(
711 recv.reset(42, 123),
712 Ok(RecvBufResetReturn {
713 max_data_delta: 104,
714 consumed_flowcontrol: 104,
715 })
716 );
717 assert_eq!(recv.len, 123);
718 assert_eq!(recv.off, 123);
719 assert_eq!(recv.data.len(), 0);
720
721 assert_emit_discard_done(&mut recv, emit);
722 }
723
724 #[test]
728 fn shutdown_empty_stream_frame() {
729 let mut recv =
730 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
731
732 assert!(recv.write(RangeBuf::from(b"hello", 0, false)).is_ok());
733 assert_eq!(recv.shutdown(), Ok(5));
734
735 assert!(recv.write(RangeBuf::from(b"", 10, false)).is_ok());
736 assert_eq!(recv.len, 10);
737 assert_eq!(recv.off, 10);
738 assert_eq!(recv.data.len(), 0);
739
740 assert_eq!(
741 recv.reset(42, 10),
742 Ok(RecvBufResetReturn {
743 max_data_delta: 0,
744 consumed_flowcontrol: 0,
745 })
746 );
747 }
748
749 #[rstest]
750 fn split_read(#[values(true, false)] emit: bool) {
751 let mut recv =
752 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
753 assert_eq!(recv.len, 0);
754
755 let first = RangeBuf::from(b"something", 0, false);
756 let second = RangeBuf::from(b"helloworld", 9, true);
757
758 assert!(recv.write(first).is_ok());
759 assert_eq!(recv.len, 9);
760 assert_eq!(recv.off, 0);
761
762 assert!(recv.write(second).is_ok());
763 assert_eq!(recv.len, 19);
764 assert_eq!(recv.off, 0);
765
766 assert_emit_discard(&mut recv, emit, 10, 10, false, Some(b"somethingh"));
767 assert_eq!(recv.len, 19);
768 assert_eq!(recv.off, 10);
769
770 assert_emit_discard(&mut recv, emit, 5, 5, false, Some(b"ellow"));
771 assert_eq!(recv.len, 19);
772 assert_eq!(recv.off, 15);
773
774 assert_emit_discard(&mut recv, emit, 5, 4, true, Some(b"orld"));
775 assert_eq!(recv.len, 19);
776 assert_eq!(recv.off, 19);
777 }
778
779 #[test]
780 fn split_read_incremental_buf() {
781 let mut recv =
782 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
783 assert_eq!(recv.len, 0);
784
785 let first = RangeBuf::from(b"something", 0, false);
786 let second = RangeBuf::from(b"helloworld", 9, true);
787
788 assert!(recv.write(first).is_ok());
789 assert_eq!(recv.len, 9);
790 assert_eq!(recv.off, 0);
791
792 assert!(recv.write(second).is_ok());
793 assert_eq!(recv.len, 19);
794 assert_eq!(recv.off, 0);
795
796 let mut buf = Vec::new().limit(10);
797 assert_eq!(
798 recv.emit_or_discard(RecvAction::Emit { out: &mut buf }),
799 Ok((10, false))
800 );
801 assert_eq!(recv.len, 19);
802 assert_eq!(recv.off, 10);
803 assert_eq!(buf.get_ref().len(), 10);
804 assert_eq!(buf.get_ref().as_slice(), b"somethingh");
805
806 buf.set_limit(5);
807 assert_eq!(
808 recv.emit_or_discard(RecvAction::Emit { out: &mut buf }),
809 Ok((5, false))
810 );
811 assert_eq!(recv.len, 19);
812 assert_eq!(recv.off, 15);
813 assert_eq!(buf.get_ref().len(), 15);
814 assert_eq!(buf.get_ref().as_slice(), b"somethinghellow");
815
816 buf.set_limit(42);
817 assert_eq!(
818 recv.emit_or_discard(RecvAction::Emit { out: &mut buf }),
819 Ok((4, true))
820 );
821 assert_eq!(recv.len, 19);
822 assert_eq!(recv.off, 19);
823 assert_eq!(buf.get_ref().len(), 19);
824 assert_eq!(buf.get_ref().as_slice(), b"somethinghelloworld");
825 }
826
827 #[rstest]
828 fn incomplete_read(#[values(true, false)] emit: bool) {
829 let mut recv =
830 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
831 assert_eq!(recv.len, 0);
832
833 let mut buf = [0u8; 32];
834
835 let first = RangeBuf::from(b"something", 0, false);
836 let second = RangeBuf::from(b"helloworld", 9, true);
837
838 assert!(recv.write(second).is_ok());
839 assert_eq!(recv.len, 19);
840 assert_eq!(recv.off, 0);
841
842 let action = if emit {
843 RecvAction::Emit {
844 out: &mut buf.as_mut_slice(),
845 }
846 } else {
847 RecvAction::Discard { len: 32 }
848 };
849 assert_eq!(recv.emit_or_discard(action), Err(Error::Done));
850
851 assert!(recv.write(first).is_ok());
852 assert_eq!(recv.len, 19);
853 assert_eq!(recv.off, 0);
854
855 assert_emit_discard(
856 &mut recv,
857 emit,
858 32,
859 19,
860 true,
861 Some(b"somethinghelloworld"),
862 );
863 assert_eq!(recv.len, 19);
864 assert_eq!(recv.off, 19);
865 }
866
867 #[rstest]
868 fn zero_len_read(#[values(true, false)] emit: bool) {
869 let mut recv =
870 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
871 assert_eq!(recv.len, 0);
872
873 let first = RangeBuf::from(b"something", 0, false);
874 let second = RangeBuf::from(b"", 9, true);
875
876 assert!(recv.write(first).is_ok());
877 assert_eq!(recv.len, 9);
878 assert_eq!(recv.off, 0);
879 assert_eq!(recv.data.len(), 1);
880
881 assert!(recv.write(second).is_ok());
882 assert_eq!(recv.len, 9);
883 assert_eq!(recv.off, 0);
884 assert_eq!(recv.data.len(), 1);
885
886 assert_emit_discard(&mut recv, emit, 32, 9, true, Some(b"something"));
887 assert_eq!(recv.len, 9);
888 assert_eq!(recv.off, 9);
889 }
890
891 #[rstest]
892 fn past_read(#[values(true, false)] emit: bool) {
893 let mut recv =
894 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
895 assert_eq!(recv.len, 0);
896
897 let first = RangeBuf::from(b"something", 0, false);
898 let second = RangeBuf::from(b"hello", 3, false);
899 let third = RangeBuf::from(b"ello", 4, true);
900 let fourth = RangeBuf::from(b"ello", 5, true);
901
902 assert!(recv.write(first).is_ok());
903 assert_eq!(recv.len, 9);
904 assert_eq!(recv.off, 0);
905 assert_eq!(recv.data.len(), 1);
906
907 assert_emit_discard(&mut recv, emit, 32, 9, false, Some(b"something"));
908 assert_eq!(recv.len, 9);
909 assert_eq!(recv.off, 9);
910
911 assert!(recv.write(second).is_ok());
912 assert_eq!(recv.len, 9);
913 assert_eq!(recv.off, 9);
914 assert_eq!(recv.data.len(), 0);
915
916 assert_eq!(recv.write(third), Err(Error::FinalSize));
917
918 assert!(recv.write(fourth).is_ok());
919 assert_eq!(recv.len, 9);
920 assert_eq!(recv.off, 9);
921 assert_eq!(recv.data.len(), 0);
922
923 assert_emit_discard_done(&mut recv, emit);
924 }
925
926 #[rstest]
927 fn fully_overlapping_read(#[values(true, false)] emit: bool) {
928 let mut recv =
929 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
930 assert_eq!(recv.len, 0);
931
932 let first = RangeBuf::from(b"something", 0, false);
933 let second = RangeBuf::from(b"hello", 4, false);
934
935 assert!(recv.write(first).is_ok());
936 assert_eq!(recv.len, 9);
937 assert_eq!(recv.off, 0);
938 assert_eq!(recv.data.len(), 1);
939
940 assert!(recv.write(second).is_ok());
941 assert_eq!(recv.len, 9);
942 assert_eq!(recv.off, 0);
943 assert_eq!(recv.data.len(), 1);
944
945 assert_emit_discard(&mut recv, emit, 32, 9, false, Some(b"something"));
946 assert_eq!(recv.len, 9);
947 assert_eq!(recv.off, 9);
948 assert_eq!(recv.data.len(), 0);
949
950 assert_emit_discard_done(&mut recv, emit);
951 }
952
953 #[rstest]
954 fn fully_overlapping_read2(#[values(true, false)] emit: bool) {
955 let mut recv =
956 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
957 assert_eq!(recv.len, 0);
958
959 let first = RangeBuf::from(b"something", 0, false);
960 let second = RangeBuf::from(b"hello", 4, false);
961
962 assert!(recv.write(second).is_ok());
963 assert_eq!(recv.len, 9);
964 assert_eq!(recv.off, 0);
965 assert_eq!(recv.data.len(), 1);
966
967 assert!(recv.write(first).is_ok());
968 assert_eq!(recv.len, 9);
969 assert_eq!(recv.off, 0);
970 assert_eq!(recv.data.len(), 2);
971
972 assert_emit_discard(&mut recv, emit, 32, 9, false, Some(b"somehello"));
973 assert_eq!(recv.len, 9);
974 assert_eq!(recv.off, 9);
975 assert_eq!(recv.data.len(), 0);
976
977 assert_emit_discard_done(&mut recv, emit);
978 }
979
980 #[rstest]
981 fn fully_overlapping_read3(#[values(true, false)] emit: bool) {
982 let mut recv =
983 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
984 assert_eq!(recv.len, 0);
985
986 let first = RangeBuf::from(b"something", 0, false);
987 let second = RangeBuf::from(b"hello", 3, false);
988
989 assert!(recv.write(second).is_ok());
990 assert_eq!(recv.len, 8);
991 assert_eq!(recv.off, 0);
992 assert_eq!(recv.data.len(), 1);
993
994 assert!(recv.write(first).is_ok());
995 assert_eq!(recv.len, 9);
996 assert_eq!(recv.off, 0);
997 assert_eq!(recv.data.len(), 3);
998
999 assert_emit_discard(&mut recv, emit, 32, 9, false, Some(b"somhellog"));
1000 assert_eq!(recv.len, 9);
1001 assert_eq!(recv.off, 9);
1002 assert_eq!(recv.data.len(), 0);
1003
1004 assert_emit_discard_done(&mut recv, emit);
1005 }
1006
1007 #[rstest]
1008 fn fully_overlapping_read_multi(#[values(true, false)] emit: bool) {
1009 let mut recv =
1010 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
1011 assert_eq!(recv.len, 0);
1012
1013 let first = RangeBuf::from(b"somethingsomething", 0, false);
1014 let second = RangeBuf::from(b"hello", 3, false);
1015 let third = RangeBuf::from(b"hello", 12, false);
1016
1017 assert!(recv.write(second).is_ok());
1018 assert_eq!(recv.len, 8);
1019 assert_eq!(recv.off, 0);
1020 assert_eq!(recv.data.len(), 1);
1021
1022 assert!(recv.write(third).is_ok());
1023 assert_eq!(recv.len, 17);
1024 assert_eq!(recv.off, 0);
1025 assert_eq!(recv.data.len(), 2);
1026
1027 assert!(recv.write(first).is_ok());
1028 assert_eq!(recv.len, 18);
1029 assert_eq!(recv.off, 0);
1030 assert_eq!(recv.data.len(), 5);
1031
1032 assert_emit_discard(
1033 &mut recv,
1034 emit,
1035 32,
1036 18,
1037 false,
1038 Some(b"somhellogsomhellog"),
1039 );
1040 assert_eq!(recv.len, 18);
1041 assert_eq!(recv.off, 18);
1042 assert_eq!(recv.data.len(), 0);
1043
1044 assert_emit_discard_done(&mut recv, emit);
1045 }
1046
1047 #[rstest]
1048 fn overlapping_start_read(#[values(true, false)] emit: bool) {
1049 let mut recv =
1050 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
1051 assert_eq!(recv.len, 0);
1052
1053 let first = RangeBuf::from(b"something", 0, false);
1054 let second = RangeBuf::from(b"hello", 8, true);
1055
1056 assert!(recv.write(first).is_ok());
1057 assert_eq!(recv.len, 9);
1058 assert_eq!(recv.off, 0);
1059 assert_eq!(recv.data.len(), 1);
1060
1061 assert!(recv.write(second).is_ok());
1062 assert_eq!(recv.len, 13);
1063 assert_eq!(recv.off, 0);
1064 assert_eq!(recv.data.len(), 2);
1065
1066 assert_emit_discard(
1067 &mut recv,
1068 emit,
1069 32,
1070 13,
1071 true,
1072 Some(b"somethingello"),
1073 );
1074
1075 assert_eq!(recv.len, 13);
1076 assert_eq!(recv.off, 13);
1077
1078 assert_emit_discard_done(&mut recv, emit);
1079 }
1080
1081 #[rstest]
1082 fn overlapping_end_read(#[values(true, false)] emit: bool) {
1083 let mut recv =
1084 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
1085 assert_eq!(recv.len, 0);
1086
1087 let first = RangeBuf::from(b"hello", 0, false);
1088 let second = RangeBuf::from(b"something", 3, true);
1089
1090 assert!(recv.write(second).is_ok());
1091 assert_eq!(recv.len, 12);
1092 assert_eq!(recv.off, 0);
1093 assert_eq!(recv.data.len(), 1);
1094
1095 assert!(recv.write(first).is_ok());
1096 assert_eq!(recv.len, 12);
1097 assert_eq!(recv.off, 0);
1098 assert_eq!(recv.data.len(), 2);
1099
1100 assert_emit_discard(&mut recv, emit, 32, 12, true, Some(b"helsomething"));
1101 assert_eq!(recv.len, 12);
1102 assert_eq!(recv.off, 12);
1103
1104 assert_emit_discard_done(&mut recv, emit);
1105 }
1106
1107 #[rstest]
1108 fn overlapping_end_twice_read(#[values(true, false)] emit: bool) {
1109 let mut recv =
1110 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
1111 assert_eq!(recv.len, 0);
1112
1113 let first = RangeBuf::from(b"he", 0, false);
1114 let second = RangeBuf::from(b"ow", 4, false);
1115 let third = RangeBuf::from(b"rl", 7, false);
1116 let fourth = RangeBuf::from(b"helloworld", 0, true);
1117
1118 assert!(recv.write(third).is_ok());
1119 assert_eq!(recv.len, 9);
1120 assert_eq!(recv.off, 0);
1121 assert_eq!(recv.data.len(), 1);
1122
1123 assert!(recv.write(second).is_ok());
1124 assert_eq!(recv.len, 9);
1125 assert_eq!(recv.off, 0);
1126 assert_eq!(recv.data.len(), 2);
1127
1128 assert!(recv.write(first).is_ok());
1129 assert_eq!(recv.len, 9);
1130 assert_eq!(recv.off, 0);
1131 assert_eq!(recv.data.len(), 3);
1132
1133 assert!(recv.write(fourth).is_ok());
1134 assert_eq!(recv.len, 10);
1135 assert_eq!(recv.off, 0);
1136 assert_eq!(recv.data.len(), 6);
1137
1138 assert_emit_discard(&mut recv, emit, 32, 10, true, Some(b"helloworld"));
1139 assert_eq!(recv.len, 10);
1140 assert_eq!(recv.off, 10);
1141
1142 assert_emit_discard_done(&mut recv, emit);
1143 }
1144
1145 #[rstest]
1146 fn overlapping_end_twice_and_contained_read(
1147 #[values(true, false)] emit: bool,
1148 ) {
1149 let mut recv =
1150 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
1151 assert_eq!(recv.len, 0);
1152
1153 let first = RangeBuf::from(b"hellow", 0, false);
1154 let second = RangeBuf::from(b"barfoo", 10, true);
1155 let third = RangeBuf::from(b"rl", 7, false);
1156 let fourth = RangeBuf::from(b"elloworldbarfoo", 1, true);
1157
1158 assert!(recv.write(third).is_ok());
1159 assert_eq!(recv.len, 9);
1160 assert_eq!(recv.off, 0);
1161 assert_eq!(recv.data.len(), 1);
1162
1163 assert!(recv.write(second).is_ok());
1164 assert_eq!(recv.len, 16);
1165 assert_eq!(recv.off, 0);
1166 assert_eq!(recv.data.len(), 2);
1167
1168 assert!(recv.write(first).is_ok());
1169 assert_eq!(recv.len, 16);
1170 assert_eq!(recv.off, 0);
1171 assert_eq!(recv.data.len(), 3);
1172
1173 assert!(recv.write(fourth).is_ok());
1174 assert_eq!(recv.len, 16);
1175 assert_eq!(recv.off, 0);
1176 assert_eq!(recv.data.len(), 5);
1177
1178 assert_emit_discard(
1179 &mut recv,
1180 emit,
1181 32,
1182 16,
1183 true,
1184 Some(b"helloworldbarfoo"),
1185 );
1186 assert_eq!(recv.len, 16);
1187 assert_eq!(recv.off, 16);
1188
1189 assert_emit_discard_done(&mut recv, emit);
1190 }
1191
1192 #[rstest]
1193 fn partially_multi_overlapping_reordered_read(
1194 #[values(true, false)] emit: bool,
1195 ) {
1196 let mut recv =
1197 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
1198 assert_eq!(recv.len, 0);
1199
1200 let first = RangeBuf::from(b"hello", 8, false);
1201 let second = RangeBuf::from(b"something", 0, false);
1202 let third = RangeBuf::from(b"moar", 11, true);
1203
1204 assert!(recv.write(first).is_ok());
1205 assert_eq!(recv.len, 13);
1206 assert_eq!(recv.off, 0);
1207 assert_eq!(recv.data.len(), 1);
1208
1209 assert!(recv.write(second).is_ok());
1210 assert_eq!(recv.len, 13);
1211 assert_eq!(recv.off, 0);
1212 assert_eq!(recv.data.len(), 2);
1213
1214 assert!(recv.write(third).is_ok());
1215 assert_eq!(recv.len, 15);
1216 assert_eq!(recv.off, 0);
1217 assert_eq!(recv.data.len(), 3);
1218
1219 assert_emit_discard(
1220 &mut recv,
1221 emit,
1222 32,
1223 15,
1224 true,
1225 Some(b"somethinhelloar"),
1226 );
1227 assert_eq!(recv.len, 15);
1228 assert_eq!(recv.off, 15);
1229 assert_eq!(recv.data.len(), 0);
1230
1231 assert_emit_discard_done(&mut recv, emit);
1232 }
1233
1234 #[rstest]
1235 fn partially_multi_overlapping_reordered_read2(
1236 #[values(true, false)] emit: bool,
1237 ) {
1238 let mut recv =
1239 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
1240 assert_eq!(recv.len, 0);
1241
1242 let first = RangeBuf::from(b"aaa", 0, false);
1243 let second = RangeBuf::from(b"bbb", 2, false);
1244 let third = RangeBuf::from(b"ccc", 4, false);
1245 let fourth = RangeBuf::from(b"ddd", 6, false);
1246 let fifth = RangeBuf::from(b"eee", 9, false);
1247 let sixth = RangeBuf::from(b"fff", 11, false);
1248
1249 assert!(recv.write(second).is_ok());
1250 assert_eq!(recv.len, 5);
1251 assert_eq!(recv.off, 0);
1252 assert_eq!(recv.data.len(), 1);
1253
1254 assert!(recv.write(fourth).is_ok());
1255 assert_eq!(recv.len, 9);
1256 assert_eq!(recv.off, 0);
1257 assert_eq!(recv.data.len(), 2);
1258
1259 assert!(recv.write(third).is_ok());
1260 assert_eq!(recv.len, 9);
1261 assert_eq!(recv.off, 0);
1262 assert_eq!(recv.data.len(), 3);
1263
1264 assert!(recv.write(first).is_ok());
1265 assert_eq!(recv.len, 9);
1266 assert_eq!(recv.off, 0);
1267 assert_eq!(recv.data.len(), 4);
1268
1269 assert!(recv.write(sixth).is_ok());
1270 assert_eq!(recv.len, 14);
1271 assert_eq!(recv.off, 0);
1272 assert_eq!(recv.data.len(), 5);
1273
1274 assert!(recv.write(fifth).is_ok());
1275 assert_eq!(recv.len, 14);
1276 assert_eq!(recv.off, 0);
1277 assert_eq!(recv.data.len(), 6);
1278
1279 assert_emit_discard(
1280 &mut recv,
1281 emit,
1282 32,
1283 14,
1284 false,
1285 Some(b"aabbbcdddeefff"),
1286 );
1287 assert_eq!(recv.len, 14);
1288 assert_eq!(recv.off, 14);
1289 assert_eq!(recv.data.len(), 0);
1290
1291 assert_emit_discard_done(&mut recv, emit);
1292 }
1293
1294 #[test]
1295 fn mixed_read_actions() {
1296 let mut recv =
1297 RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
1298 assert_eq!(recv.len, 0);
1299
1300 let first = RangeBuf::from(b"hello", 0, false);
1301 let second = RangeBuf::from(b"world", 5, false);
1302 let third = RangeBuf::from(b"something", 10, true);
1303
1304 assert!(recv.write(second).is_ok());
1305 assert_eq!(recv.len, 10);
1306 assert_eq!(recv.off, 0);
1307
1308 assert_emit_discard_done(&mut recv, true);
1309 assert_emit_discard_done(&mut recv, false);
1310
1311 assert!(recv.write(third).is_ok());
1312 assert_eq!(recv.len, 19);
1313 assert_eq!(recv.off, 0);
1314
1315 assert_emit_discard_done(&mut recv, true);
1316 assert_emit_discard_done(&mut recv, false);
1317
1318 assert!(recv.write(first).is_ok());
1319 assert_eq!(recv.len, 19);
1320 assert_eq!(recv.off, 0);
1321
1322 assert_emit_discard(&mut recv, true, 5, 5, false, Some(b"hello"));
1323 assert_eq!(recv.len, 19);
1324 assert_eq!(recv.off, 5);
1325
1326 assert_emit_discard(&mut recv, false, 5, 5, false, None);
1327 assert_eq!(recv.len, 19);
1328 assert_eq!(recv.off, 10);
1329
1330 assert_emit_discard(&mut recv, true, 9, 9, true, Some(b"something"));
1331 assert_eq!(recv.len, 19);
1332 assert_eq!(recv.off, 19);
1333
1334 assert_emit_discard_done(&mut recv, true);
1335 assert_emit_discard_done(&mut recv, false);
1336 }
1337}