Skip to main content

quiche/stream/
recv_buf.rs

1// Copyright (C) 2023, Cloudflare, Inc.
2// All rights reserved.
3//
4// Redistribution and use in source and binary forms, with or without
5// modification, are permitted provided that the following conditions are
6// met:
7//
8//     * Redistributions of source code must retain the above copyright notice,
9//       this list of conditions and the following disclaimer.
10//
11//     * Redistributions in binary form must reproduce the above copyright
12//       notice, this list of conditions and the following disclaimer in the
13//       documentation and/or other materials provided with the distribution.
14//
15// THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS
16// IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO,
17// THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR
18// PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR
19// CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL,
20// EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO,
21// PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR
22// PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF
23// LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING
24// NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
25// SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
26
27use 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/// Receive-side stream buffer.
45///
46/// Stream data received by the peer is buffered in a list of data chunks
47/// ordered by offset in ascending order. Contiguous data can then be read
48/// into a slice.
49#[derive(Debug, Default)]
50pub struct RecvBuf {
51    /// Chunks of data received from the peer that have not yet been read by
52    /// the application, ordered by offset.
53    data: BTreeMap<u64, RangeBuf>,
54
55    /// The lowest data offset that has yet to be read by the application.
56    off: u64,
57
58    /// The total length of data received on this stream.
59    len: u64,
60
61    /// Receiver flow controller.
62    flow_control: flowcontrol::FlowControl,
63
64    /// The final stream offset received from the peer, if any.
65    fin_off: Option<u64>,
66
67    /// The error code received via RESET_STREAM.
68    error: Option<u64>,
69
70    /// Whether incoming data is validated but not buffered.
71    drain: bool,
72}
73
74impl RecvBuf {
75    /// Creates a new receive buffer.
76    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    /// Inserts the given chunk of data in the buffer.
88    ///
89    /// This also takes care of enforcing stream flow control limits, as well
90    /// as handling incoming data that overlaps data that is already in the
91    /// buffer.
92    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            // Stream's size is known, forbid data beyond that point.
99            if buf.max_off() > fin_off {
100                return Err(Error::FinalSize);
101            }
102
103            // Stream's size is already known, forbid changing it.
104            if buf.fin() && fin_off != buf.max_off() {
105                return Err(Error::FinalSize);
106            }
107        }
108
109        // Stream's known size is lower than data already received.
110        if buf.fin() && buf.max_off() < self.len {
111            return Err(Error::FinalSize);
112        }
113
114        // We already saved the final offset, so there's nothing else we
115        // need to keep from the RangeBuf if it's empty.
116        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        // No need to store empty buffer that doesn't carry the fin flag, but
125        // its offset still advances the largest received offset used by flow
126        // control (RFC 9000 Section 19.8).
127        if !buf.fin() && buf.is_empty() {
128            self.len = cmp::max(self.len, buf.max_off());
129
130            if self.drain {
131                // we are not storing any data, off == len
132                self.off = self.len;
133            }
134
135            return Ok(());
136        }
137
138        // Check if data is fully duplicate, that is the buffer's max offset is
139        // lower or equal to the offset already stored in the recv buffer.
140        if self.off >= buf.max_off() {
141            // An exception is applied to empty range buffers, because an empty
142            // buffer's max offset matches the max offset of the recv buffer.
143            //
144            // By this point all spurious empty buffers should have already been
145            // discarded, so allowing empty buffers here should be safe.
146            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            // Discard incoming data below current stream offset. Bytes up to
156            // `self.off` have already been received so we should not buffer
157            // them again. This is also important to make sure `ready()` doesn't
158            // get stuck when a buffer with lower offset than the stream's is
159            // buffered.
160            if self.off_front() > buf.off() {
161                buf = buf.split_off((self.off_front() - buf.off()) as usize);
162            }
163
164            // Handle overlapping data. If the incoming data's starting offset
165            // is above the previous maximum received offset, there is clearly
166            // no overlap so this logic can be skipped. However do still try to
167            // merge an empty final buffer (i.e. an empty buffer with the fin
168            // flag set, which is the only kind of empty buffer that should
169            // reach this point).
170            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                    // We are past the current buffer.
175                    if b.off() > buf.max_off() {
176                        break;
177                    }
178
179                    // New buffer is fully contained in existing buffer.
180                    if off >= b.off() && buf.max_off() <= b.max_off() {
181                        continue 'tmp;
182                    }
183
184                    // New buffer's start overlaps existing buffer.
185                    if off >= b.off() && off < b.max_off() {
186                        buf = buf.split_off((b.max_off() - off) as usize);
187                    }
188
189                    // New buffer's end overlaps existing buffer.
190                    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                // we are not storing any data, off == len
203                self.off = self.len;
204            }
205        }
206
207        Ok(())
208    }
209
210    /// Reads contiguous data from the receive buffer.
211    ///
212    /// Data is written into the given `out` buffer, up to the length of `out`.
213    ///
214    /// Only contiguous data is removed, starting from offset 0. The offset is
215    /// incremented as data is taken out of the receive buffer. If there is no
216    /// data at the expected read offset, the `Done` error is returned.
217    ///
218    /// On success the amount of data read and a flag indicating
219    /// if there is no more data in the buffer, are returned as a tuple.
220    #[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    /// Reads or discards contiguous data from the receive buffer.
226    ///
227    /// Passing an `action` of `StreamRecvAction::Emit` results in data being
228    /// written into the provided buffer, up to its length.
229    ///
230    /// Passing an `action` of `StreamRecvAction::Discard` results in up to
231    /// the indicated number of bytes being discarded without copying.
232    ///
233    /// Only contiguous data is removed, starting from offset 0. The offset is
234    /// incremented as data is taken out of the receive buffer. If there is no
235    /// data at the expected read offset, the `Done` error is returned.
236    ///
237    /// On success the amount of data read or discarded, and a flag indicating
238    /// if there is no more data in the buffer, are returned as a tuple.
239    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        // The stream was reset, so clear its data and return the error code
253        // instead.
254        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            // Only copy data if we're emitting, not discarding.
274            if let RecvAction::Emit { ref mut out } = action {
275                // Note: `BufMut::remaining_mut()` cannot "shrink", but BufMut
276                // impls are allowed to grow the buffer, so we
277                // check here that we still have at least
278                // `cap` bytes, but we can't require equality
279                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                // We reached the maximum capacity, so end here.
295                break;
296            }
297
298            entry.remove();
299        }
300
301        // Update consumed bytes for flow control.
302        self.flow_control.add_consumed(len as u64);
303
304        Ok((len, self.is_fin()))
305    }
306
307    /// Resets the stream at the given offset.
308    pub fn reset(
309        &mut self, error_code: u64, final_size: u64,
310    ) -> Result<RecvBufResetReturn> {
311        // Stream's size is already known, forbid changing it.
312        if let Some(fin_off) = self.fin_off {
313            if fin_off != final_size {
314                return Err(Error::FinalSize);
315            }
316        }
317
318        // Stream's known size is lower than data already received.
319        if final_size < self.len {
320            return Err(Error::FinalSize);
321        }
322
323        if self.error.is_some() {
324            // We already verified that the final size matches
325            return Ok(RecvBufResetReturn::zero());
326        }
327
328        // Calculate how many bytes need to be removed from the connection flow
329        // control.
330        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        // Clear all data already buffered.
338        self.off = final_size;
339
340        self.data.clear();
341
342        // In order to ensure the application is notified when the stream is
343        // reset, enqueue a zero-length buffer at the final size offset.
344        let buf = RangeBuf::from(b"", final_size, true);
345        self.write(buf)?;
346
347        Ok(result)
348    }
349
350    /// Commits the new max_data limit.
351    pub fn update_max_data(&mut self, now: Instant) {
352        self.flow_control.update_max_data(now);
353    }
354
355    /// Return the new max_data limit.
356    pub fn max_data_next(&mut self) -> u64 {
357        self.flow_control.max_data_next()
358    }
359
360    /// Return the current flow control limit.
361    pub fn max_data(&self) -> u64 {
362        self.flow_control.max_data()
363    }
364
365    /// Return the current window.
366    pub fn window(&self) -> u64 {
367        self.flow_control.window()
368    }
369
370    /// Autotune the window size.
371    pub fn autotune_window(&mut self, now: Instant, rtt: Duration) {
372        self.flow_control.autotune_window(now, rtt);
373    }
374
375    /// Shuts down receiving data and returns the number of bytes
376    /// that should be returned to the connection level flow
377    /// control
378    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    /// Returns the lowest offset of data buffered.
394    pub fn off_front(&self) -> u64 {
395        self.off
396    }
397
398    /// Returns true if we need to update the local flow control limit.
399    pub fn almost_full(&self) -> bool {
400        self.fin_off.is_none() && self.flow_control.should_update_max_data()
401    }
402
403    /// Returns the largest offset ever received.
404    pub fn max_off(&self) -> u64 {
405        self.len
406    }
407
408    /// Returns true if the receive-side of the stream is complete.
409    ///
410    /// This happens when the stream's receive final size is known, and the
411    /// application has read all data from the stream.
412    pub fn is_fin(&self) -> bool {
413        if self.fin_off == Some(self.off) {
414            return true;
415        }
416
417        false
418    }
419
420    /// Returns true if the stream is not storing incoming data.
421    pub fn is_draining(&self) -> bool {
422        self.drain
423    }
424
425    /// Returns true if the stream has data to be read.
426    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    /// Returns the number of bytes that can be read contiguously from the
436    /// current read offset, up to `max_len`.
437    ///
438    /// Data buffered behind a gap (received out of order) is not counted, so
439    /// this never reports bytes that are not yet readable. The cost is
440    /// proportional to the number of contiguous buffered chunks at the front
441    /// of the buffer, up to `max_len`; no data is copied.
442    pub fn readable_len(&self, max_len: usize) -> usize {
443        let mut contiguous = 0usize;
444        let mut next_off = self.off;
445
446        // `data` is ordered by offset, so walk from the front and stop at the
447        // first gap (a chunk that does not start where the contiguous run so
448        // far leaves off).
449        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    /// The default size of the receiver stream flow control window.
476    const DEFAULT_STREAM_WINDOW: u64 = 32 * 1024;
477    use bytes::BufMut as _;
478    use rstest::rstest;
479
480    // Helper function for testing either buffer emit or discard.
481    //
482    // The `emit` parameter controls whether data is emitted or discarded from
483    // `recv`.
484    //
485    // The `target_len` parameter controls the maximum amount of bytes that
486    // could be read, up to the capacity of `recv`. The `result_len` is the
487    // actual number of bytes that were taken out of `recv`. An assert is
488    // performed on `result_len` to ensure the number of bytes read meets the
489    // caller expectations.
490    //
491    // The `is_fin` parameter relates to the buffer's finished status. An assert
492    // is performed on it to ensure the status meet the caller expectations.
493    //
494    // The `test_bytes` parameter carries an optional slice of bytes. Is set, an
495    // assert is performed against the bytes that were read out of the buffer,
496    // to ensure caller expectations are met.
497    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    // Helper function for testing buffer status for either emit or discard.
523    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        // Don't store non-fin empty buffer, but track its offset.
559        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        // Check flow control for empty buffer.
566        let buf = RangeBuf::from(b"", 16, false);
567        assert_eq!(recv.write(buf), Err(Error::FlowControl));
568
569        // A final size below the advanced largest offset is an error
570        // (RFC 9000 Section 4.5).
571        let buf = RangeBuf::from(b"", 5, true);
572        assert_eq!(recv.write(buf), Err(Error::FinalSize));
573
574        // Store fin empty buffer at the largest received offset.
575        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        // Don't store additional fin empty buffers.
582        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        // Accept another fin buffer with the same final size.
589        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        // A fin buffer whose end disagrees with the known final size errors.
596        let buf = RangeBuf::from(b"aa", 3, true);
597        assert_eq!(recv.write(buf), Err(Error::FinalSize));
598
599        // Validate final size with fin empty buffers.
600        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        // The range (5..10) was never received, so the stream cannot reach
606        // its fin and nothing further is readable.
607        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    /// `readable_len` counts only contiguous in-order data, up to its limit.
652    fn readable_len() {
653        let mut recv =
654            RecvBuf::new(u64::MAX, DEFAULT_STREAM_WINDOW, DEFAULT_STREAM_WINDOW);
655
656        // Empty buffer: nothing readable.
657        assert_eq!(recv.readable_len(64 * 1024), 0);
658
659        // Data buffered behind a gap is not readable.
660        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        // Filling the gap makes the full range readable, bounded by the limit.
665        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    /// Test shutdown behavior
671    #[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        // shutdown the buffer. Buffer is dropped.
688        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        // subsequent writes are validated but not added to the buffer
696        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        // the max offset of received data can increase and
702        // the recv.off must increase with it
703        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        // Send a reset
709        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    /// An empty non-fin buffer advances the largest received offset, which the
725    /// connection charges to flow control on arrival, so a draining stream must
726    /// consume it too and a later reset must not credit it again.
727    #[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}