Skip to main content

quiche/stream/
mod.rs

1// Copyright (C) 2018-2019, 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::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
50/// The maximum size of the receiver stream flow control window.
51pub const MAX_STREAM_WINDOW: u64 = 16 * 1024 * 1024;
52
53/// A simple no-op hasher for Stream IDs.
54///
55/// The QUIC protocol and quiche library guarantees stream ID uniqueness, so
56/// we can save effort by avoiding using a more complicated algorithm.
57#[derive(Default)]
58pub struct StreamIdHasher {
59    id: u64,
60}
61
62/// Return value type of `RecvBuf::reset()`
63#[derive(Debug, PartialEq, Clone, Copy)]
64pub struct RecvBufResetReturn {
65    /// Returns the difference between the previous max_data offset
66    /// received and the final size reported by the reset
67    pub max_data_delta: u64,
68
69    /// The amount of flow control credit that should be returned to the
70    /// connection level flow control.
71    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
83/// Action to perform when reading from a stream's receive buffer.
84pub enum RecvAction<T: bytes::BufMut> {
85    /// Emit data by copying it into the provided buffer.
86    Emit { out: T },
87    /// Discard up to the specified number of bytes without copying.
88    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        // We need a default write() for the trait but stream IDs will always
105        // be a u64 so we just delegate to write_u64.
106        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/// Tracks collected stream sequences separately for each stream type.
116#[derive(Default)]
117struct CollectedStreams {
118    // Defer allocation until the first stream is collected. The range capacity
119    // is unlimited because evicting a tombstone would allow a collected stream
120    // to be recreated.
121    ranges: Option<Box<[RangeSet; 4]>>,
122}
123
124impl CollectedStreams {
125    fn insert(&mut self, stream_id: u64) {
126        // Same-type stream IDs advance by four, so store their sequences to
127        // allow adjacent collected streams to merge into a single range.
128        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/// Keeps track of QUIC streams and enforces stream limits.
142#[derive(Default)]
143pub struct StreamMap<F: BufFactory = DefaultBufFactory> {
144    /// Map of streams indexed by stream ID.
145    streams: StreamIdHashMap<Stream<F>>,
146
147    /// Set of streams that were completed and garbage collected.
148    ///
149    /// Instead of keeping the full stream state forever, we collect completed
150    /// streams to save memory, but we still need to keep track of previously
151    /// created streams, to prevent peers from re-creating them.
152    collected: CollectedStreams,
153
154    /// Peer's maximum bidirectional stream count limit.
155    peer_max_streams_bidi: u64,
156
157    /// Peer's maximum unidirectional stream count limit.
158    peer_max_streams_uni: u64,
159
160    /// The total number of bidirectional streams opened by the peer.
161    peer_opened_streams_bidi: u64,
162
163    /// The total number of unidirectional streams opened by the peer.
164    peer_opened_streams_uni: u64,
165
166    /// Local maximum bidirectional stream count limit.
167    local_max_streams_bidi: u64,
168    local_max_streams_bidi_next: u64,
169
170    /// Initial maximum bidirectional stream count.
171    initial_max_streams_bidi: u64,
172
173    /// Local maximum unidirectional stream count limit.
174    local_max_streams_uni: u64,
175    local_max_streams_uni_next: u64,
176
177    /// Initial maximum unidirectional stream count.
178    initial_max_streams_uni: u64,
179
180    /// The total number of bidirectional streams opened by the local endpoint.
181    local_opened_streams_bidi: u64,
182
183    /// The total number of unidirectional streams opened by the local endpoint.
184    local_opened_streams_uni: u64,
185
186    /// Queue of stream IDs corresponding to streams that have buffered data
187    /// ready to be sent to the peer. This also implies that the stream has
188    /// enough flow control credits to send at least some of that data.
189    flushable: RBTree<StreamFlushablePriorityAdapter>,
190
191    /// Set of stream IDs corresponding to streams that have outstanding data
192    /// to read. This is used to generate a `StreamIter` of streams without
193    /// having to iterate over the full list of streams.
194    pub readable: RBTree<StreamReadablePriorityAdapter>,
195
196    /// Set of stream IDs corresponding to streams that have enough flow control
197    /// capacity to be written to, and is not finished. This is used to generate
198    /// a `StreamIter` of streams without having to iterate over the full list
199    /// of streams.
200    pub writable: RBTree<StreamWritablePriorityAdapter>,
201
202    /// Set of stream IDs corresponding to streams that are almost out of flow
203    /// control credit and need to send MAX_STREAM_DATA. This is used to
204    /// generate a `StreamIter` of streams without having to iterate over the
205    /// full list of streams.
206    almost_full: StreamIdHashSet,
207
208    /// Set of stream IDs corresponding to streams that are blocked. The value
209    /// of the map elements represents the offset of the stream at which the
210    /// blocking occurred.
211    blocked: StreamIdHashMap<u64>,
212
213    /// Set of stream IDs corresponding to streams that are reset. The value
214    /// of the map elements is a tuple of the error code and final size values
215    /// to include in the RESET_STREAM frame.
216    reset: StreamIdHashMap<(u64, u64)>,
217
218    /// Set of stream IDs corresponding to streams that are shutdown on the
219    /// receive side, and need to send a STOP_SENDING frame. The value of the
220    /// map elements is the error code to include in the STOP_SENDING frame.
221    stopped: StreamIdHashMap<u64>,
222
223    /// The maximum size of a stream window.
224    max_stream_window: u64,
225
226    /// Total number of bytes in send buffers across all streams.
227    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    /// Returns the stream with the given ID if it exists.
250    pub fn get(&self, id: u64) -> Option<&Stream<F>> {
251        self.streams.get(&id)
252    }
253
254    /// Returns the mutable stream with the given ID if it exists.
255    pub fn get_mut(&mut self, id: u64) -> Option<&mut Stream<F>> {
256        self.streams.get_mut(&id)
257    }
258
259    /// Returns the mutable stream with the given ID if it exists, or creates
260    /// a new one otherwise.
261    ///
262    /// The `local` parameter indicates whether the stream's creation was
263    /// requested by the local application rather than the peer, and is
264    /// used to validate the requested stream ID, and to select the initial
265    /// flow control values from the local and remote transport parameters
266    /// (also passed as arguments).
267    ///
268    /// This also takes care of enforcing both local and the peer's stream
269    /// count limits. If one of these limits is violated, the `StreamLimit`
270    /// error is returned.
271    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                // Stream has already been closed and garbage collected.
278                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                    // Locally-initiated bidirectional stream.
288                    (true, true) => (
289                        local_params.initial_max_stream_data_bidi_local,
290                        peer_params.initial_max_stream_data_bidi_remote,
291                    ),
292
293                    // Locally-initiated unidirectional stream.
294                    (true, false) => (0, peer_params.initial_max_stream_data_uni),
295
296                    // Remotely-initiated bidirectional stream.
297                    (false, true) => (
298                        local_params.initial_max_stream_data_bidi_remote,
299                        peer_params.initial_max_stream_data_bidi_local,
300                    ),
301
302                    // Remotely-initiated unidirectional stream.
303                    (false, false) =>
304                        (local_params.initial_max_stream_data_uni, 0),
305                };
306
307                // The two least significant bits from a stream id identify the
308                // type of stream. Truncate those bits to get the sequence for
309                // that stream type.
310                let stream_sequence = id >> 2;
311
312                // Enforce stream count limits.
313                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        // Newly created stream might already be writable due to initial flow
386        // control limits.
387        if is_new_and_writable {
388            self.writable.insert(Arc::clone(&stream.priority_key));
389        }
390
391        Ok(stream)
392    }
393
394    /// Adds the stream ID to the readable streams set.
395    ///
396    /// If the stream was already in the list, this does nothing.
397    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    /// Removes the stream ID from the readable streams set.
404    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    /// Adds the stream ID to the writable streams set.
418    ///
419    /// This should also be called anytime a new stream is created, in addition
420    /// to when an existing stream becomes writable.
421    ///
422    /// If the stream was already in the list, this does nothing.
423    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    /// Removes the stream ID from the writable streams set.
430    ///
431    /// This should also be called anytime an existing stream stops being
432    /// writable.
433    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    /// Adds the stream ID to the flushable streams set.
447    ///
448    /// If the stream was already in the list, this does nothing.
449    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    /// Removes the stream ID from the flushable streams set.
456    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    /// Updates the priorities of a stream.
474    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    /// Adds the stream ID to the almost full streams set.
494    ///
495    /// If the stream was already in the list, this does nothing.
496    pub fn insert_almost_full(&mut self, stream_id: u64) {
497        self.almost_full.insert(stream_id);
498    }
499
500    /// Removes the stream ID from the almost full streams set.
501    pub fn remove_almost_full(&mut self, stream_id: u64) {
502        self.almost_full.remove(&stream_id);
503    }
504
505    /// Adds the stream ID to the blocked streams set with the
506    /// given offset value.
507    ///
508    /// If the stream was already in the list, this does nothing.
509    pub fn insert_blocked(&mut self, stream_id: u64, off: u64) {
510        self.blocked.insert(stream_id, off);
511    }
512
513    /// Removes the stream ID from the blocked streams set.
514    pub fn remove_blocked(&mut self, stream_id: u64) {
515        self.blocked.remove(&stream_id);
516    }
517
518    /// Adds the stream ID to the reset streams set with the
519    /// given error code and final size values.
520    ///
521    /// If the stream was already in the list, this does nothing.
522    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    /// Removes the stream ID from the reset streams set.
529    pub fn remove_reset(&mut self, stream_id: u64) {
530        self.reset.remove(&stream_id);
531    }
532
533    /// Adds the stream ID to the stopped streams set with the
534    /// given error code.
535    ///
536    /// If the stream was already in the list, this does nothing.
537    pub fn insert_stopped(&mut self, stream_id: u64, error_code: u64) {
538        self.stopped.insert(stream_id, error_code);
539    }
540
541    /// Removes the stream ID from the stopped streams set.
542    pub fn remove_stopped(&mut self, stream_id: u64) {
543        self.stopped.remove(&stream_id);
544    }
545
546    /// Updates the peer's maximum bidirectional stream count limit.
547    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    /// Updates the peer's maximum unidirectional stream count limit.
552    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    /// Commits the new max_streams_bidi limit.
557    pub fn update_max_streams_bidi(&mut self) {
558        self.local_max_streams_bidi = self.local_max_streams_bidi_next;
559    }
560
561    /// Sets the max_streams_bidi limit to the given value.
562    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    /// Returns the current max_streams_bidi limit.
569    pub fn max_streams_bidi(&self) -> u64 {
570        self.local_max_streams_bidi
571    }
572
573    /// Returns the new max_streams_bidi limit.
574    pub fn max_streams_bidi_next(&mut self) -> u64 {
575        self.local_max_streams_bidi_next
576    }
577
578    /// Commits the new max_streams_uni limit.
579    pub fn update_max_streams_uni(&mut self) {
580        self.local_max_streams_uni = self.local_max_streams_uni_next;
581    }
582
583    /// Returns the new max_streams_uni limit.
584    pub fn max_streams_uni_next(&mut self) -> u64 {
585        self.local_max_streams_uni_next
586    }
587
588    /// Returns the peer's current maximum bidirectional stream count limit.
589    pub fn peer_max_streams_bidi(&self) -> u64 {
590        self.peer_max_streams_bidi
591    }
592
593    /// Returns the number of bidirectional streams that can be created
594    /// before the peer's stream count limit is reached.
595    pub fn peer_streams_left_bidi(&self) -> u64 {
596        self.peer_max_streams_bidi - self.local_opened_streams_bidi
597    }
598
599    /// Returns the peer's current maximum unidirectional stream count limit.
600    pub fn peer_max_streams_uni(&self) -> u64 {
601        self.peer_max_streams_uni
602    }
603
604    /// Returns the number of unidirectional streams that can be created
605    /// before the peer's stream count limit is reached.
606    pub fn peer_streams_left_uni(&self) -> u64 {
607        self.peer_max_streams_uni - self.local_opened_streams_uni
608    }
609
610    /// Drops completed stream.
611    ///
612    /// This should only be called when Stream::is_complete() returns true for
613    /// the given stream.
614    pub fn collect(&mut self, stream_id: u64, local: bool) {
615        if !local {
616            // If the stream was created by the peer, give back a max streams
617            // credit.
618            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    /// Creates an iterator over streams that have outstanding data to read.
639    pub fn readable(&self) -> StreamIter {
640        StreamIter {
641            streams: self.readable.iter().map(|s| s.id).collect(),
642            index: 0,
643        }
644    }
645
646    /// Creates an iterator over streams that can be written to.
647    pub fn writable(&self) -> StreamIter {
648        StreamIter {
649            streams: self.writable.iter().map(|s| s.id).collect(),
650            index: 0,
651        }
652    }
653
654    /// Creates an iterator over streams that need to send MAX_STREAM_DATA.
655    pub fn almost_full(&self) -> StreamIter {
656        StreamIter::from(&self.almost_full)
657    }
658
659    /// Creates an iterator over streams that need to send STREAM_DATA_BLOCKED.
660    pub fn blocked(&self) -> hash_map::Iter<'_, u64, u64> {
661        self.blocked.iter()
662    }
663
664    /// Creates an iterator over streams that need to send RESET_STREAM.
665    pub fn reset(&self) -> hash_map::Iter<'_, u64, (u64, u64)> {
666        self.reset.iter()
667    }
668
669    /// Creates an iterator over streams that need to send STOP_SENDING.
670    pub fn stopped(&self) -> hash_map::Iter<'_, u64, u64> {
671        self.stopped.iter()
672    }
673
674    /// Returns true if the stream has been collected.
675    pub fn is_collected(&self, stream_id: u64) -> bool {
676        self.collected.contains(stream_id)
677    }
678
679    /// Returns true if there are any streams that have data to write.
680    pub fn has_flushable(&self) -> bool {
681        !self.flushable.is_empty()
682    }
683
684    /// Returns true if there are any streams that have data to read.
685    pub fn has_readable(&self) -> bool {
686        !self.readable.is_empty()
687    }
688
689    /// Returns true if there are any streams that need to update the local
690    /// flow control limit.
691    pub fn has_almost_full(&self) -> bool {
692        !self.almost_full.is_empty()
693    }
694
695    /// Returns true if there are any streams that are blocked.
696    pub fn has_blocked(&self) -> bool {
697        !self.blocked.is_empty()
698    }
699
700    /// Returns true if there are any streams that are reset.
701    pub fn has_reset(&self) -> bool {
702        !self.reset.is_empty()
703    }
704
705    /// Returns true if there are any streams that need to send STOP_SENDING.
706    pub fn has_stopped(&self) -> bool {
707        !self.stopped.is_empty()
708    }
709
710    /// Returns true if the max bidirectional streams count needs to be updated
711    /// by sending a MAX_STREAMS frame to the peer.
712    ///
713    /// This only sends MAX_STREAMS when available capacity is at or below 50%
714    /// of the initial maximum streams target.
715    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    /// Returns true if the max unidirectional streams count needs to be updated
724    /// by sending a MAX_STREAMS frame to the peer.
725    ///
726    /// This only send MAX_STREAMS when available capacity is at or below 50% of
727    /// the initial maximum streams target.
728    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    /// Returns the number of active streams in the map.
737    #[cfg(test)]
738    pub fn len(&self) -> usize {
739        self.streams.len()
740    }
741
742    /// Returns the total number of bytes buffered across all streams.
743    pub(crate) fn tx_buffered(&self) -> usize {
744        self.tx_buffered
745    }
746
747    /// Computes the actual number of bytes in send buffers by summing across
748    /// all streams. This is used for debugging to verify that tx_buffered
749    /// is accurate.
750    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    /// Checks if the stored tx_buffered matches the actual value.
758    /// Returns true if they match, false otherwise.
759    pub(crate) fn tx_buffered_is_consistent(&self) -> bool {
760        self.tx_buffered == self.tx_buffered_actual()
761    }
762
763    /// Updates the tx_buffered value by adding the delta.
764    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    /// Updates the tx_buffered value by subtracting the delta.
772    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    /// Verifies that the stored tx_buffered value matches the actual bytes in
781    /// send buffers across all streams. Enabled in debug builds to catch
782    /// inconsistencies early.
783    #[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
805/// A QUIC stream.
806pub struct Stream<F: BufFactory = DefaultBufFactory> {
807    /// Receive-side stream buffer.
808    pub recv: recv_buf::RecvBuf,
809
810    /// Send-side stream buffer.
811    pub send: send_buf::SendBuf<F>,
812
813    pub send_lowat: usize,
814
815    /// Whether the stream is bidirectional.
816    pub bidi: bool,
817
818    /// Whether the stream was created by the local endpoint.
819    pub local: bool,
820
821    /// The stream's urgency (lower is better). Default is `DEFAULT_URGENCY`.
822    pub urgency: u8,
823
824    /// Whether the stream can be flushed incrementally. Default is `true`.
825    pub incremental: bool,
826
827    pub priority_key: Arc<StreamPriorityKey>,
828}
829
830impl<F: BufFactory> Stream<F> {
831    /// Creates a new stream with the given flow control limits.
832    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    /// Returns true if the stream has data to read.
854    pub fn is_readable(&self) -> bool {
855        self.recv.ready()
856    }
857
858    /// Returns true if the stream has enough flow control capacity to be
859    /// written to, and is not finished.
860    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    /// Returns true if the stream has data to send and is allowed to send at
868    /// least some of it.
869    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    /// Returns true if the stream is complete.
878    ///
879    /// For bidirectional streams this happens when both the receive and send
880    /// sides are complete. That is when all incoming data has been read by the
881    /// application, and when all outgoing data has been acked by the peer.
882    ///
883    /// For unidirectional streams this happens when either the receive or send
884    /// side is complete, depending on whether the stream was created locally
885    /// or not.
886    pub fn is_complete(&self) -> bool {
887        match (self.bidi, self.local) {
888            // For bidirectional streams we need to check both receive and send
889            // sides for completion.
890            (true, _) => self.recv.is_fin() && self.send.is_complete(),
891
892            // For unidirectional streams generated locally, we only need to
893            // check the send side for completion.
894            (false, true) => self.send.is_complete(),
895
896            // For unidirectional streams generated by the peer, we only need
897            // to check the receive side for completion.
898            (false, false) => self.recv.is_fin(),
899        }
900    }
901}
902
903/// Returns true if the stream was created locally.
904pub fn is_local(stream_id: u64, is_server: bool) -> bool {
905    (stream_id & 0x1) == (is_server as u64)
906}
907
908/// Returns true if the stream is bidirectional.
909pub 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        // Ignore priority if ID matches.
954        if self.id == other.id {
955            return cmp::Ordering::Equal;
956        }
957
958        // First, order by urgency...
959        if self.urgency != other.urgency {
960            return self.urgency.cmp(&other.urgency);
961        }
962
963        // ...when the urgency is the same, and both are not incremental, order
964        // by stream ID...
965        if !self.incremental && !other.incremental {
966            return self.id.cmp(&other.id);
967        }
968
969        // ...non-incremental takes priority over incremental...
970        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        // ...finally, when both are incremental, `other` takes precedence (so
978        // `self` is always sorted after other same-urgency incremental
979        // entries).
980        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/// An iterator over QUIC streams.
1015#[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    /// The default size of the receiver stream flow control window.
1059    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    /// Completes a client-initiated stream and processes returned stream
1132    /// credit.
1133    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        // Late STREAM frames must not materialize a collected stream again.
1190        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        // The odd sequences consume credit but have never had Stream objects.
1235        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        // Filling the implicit gaps is still permitted and merges all ranges.
1246        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        // Local fragmentation follows the peer's allowance and application
1278        // behavior, not the client's initial incoming stream limit of one.
1279        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, &params, &params, 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, &params, &params, 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        // Read one byte.
1454        assert_eq!(stream.recv.emit(&mut [0; 1]), Ok((1, false)));
1455        // Reset with a final size > than max previously received
1456        assert_eq!(
1457            stream.recv.reset(0, 10),
1458            Ok(RecvBufResetReturn {
1459                max_data_delta: 5,
1460                // consumed_flowcontrol is 9, since we already read 1 byte
1461                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        // Advance buffer.
2048        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        // Split buffer before position.
2063        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        // Advance buffer.
2090        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        // Split buffer after position.
2105        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        // Advance buffer.
2132        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    /// RFC9000 2.1: A stream ID that is used out of order results in all
2148    /// streams of that type with lower-numbered stream IDs also being opened.
2149    #[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    /// Stream limit should be satisfied regardless of what order we open
2169    /// streams
2170    #[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    /// Check stream limit boundary cases
2187    #[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        // Highest permitted
2195        let stream_id = 8;
2196        assert!(streams
2197            .get_or_create(stream_id, &local_tp, &peer_tp, false, true)
2198            .is_ok());
2199
2200        // One more than highest permitted
2201        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        // All streams are non-incremental and same urgency by default. Multiple
2243        // visits shuffle their order.
2244        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        // Inserting same-urgency incremental streams in a "random" order yields
2263        // same order to start with.
2264        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        // Streams where the urgency descends (becomes more important). No
2298        // stream shares an urgency.
2299        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            // this duplicates some code from stream_priority in order to access
2315            // streams and the collection they're in
2316            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        // Re-applying priority to a stream does not cause duplication.
2341        for (id, urgency) in input {
2342            // this duplicates some code from stream_priority in order to access
2343            // streams and the collection they're in
2344            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        // Removing streams doesn't break expected ordering.
2369        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        // Adding streams doesn't break expected ordering.
2381        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        // Streams that share some urgency level
2401        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            // this duplicates some code from stream_priority in order to access
2417            // streams and the collection they're in
2418            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        // Removing streams doesn't break expected ordering.
2487        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        // Adding streams doesn't break expected ordering.
2493        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        // Default keys could cause duplicate entries. `StreamMap` normally
2537        // prevents this.
2538        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        // Write and emit some data.
2559        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        // Mark data for retransmission.
2568        let retransmitted = stream.send.retransmit(0, 5);
2569        assert_eq!(retransmitted, 5);
2570        assert_eq!(stream.send.buffered_bytes(), 5);
2571
2572        // Ack the data.
2573        stream.send.ack_and_drop(0, 5);
2574        assert_eq!(stream.send.buffered_bytes(), 0);
2575
2576        // Try to retransmit again - should return 0 since data is acked.
2577        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        // Write and emit 10 bytes.
2587        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        // Mark all data for retransmission.
2596        let retransmitted = stream.send.retransmit(0, 10);
2597        assert_eq!(retransmitted, 10);
2598        assert_eq!(stream.send.buffered_bytes(), 10);
2599
2600        // Ack first 5 bytes and drop them.
2601        let dropped = stream.send.ack_and_drop(0, 5);
2602        assert_eq!(dropped, 5);
2603        assert_eq!(stream.send.buffered_bytes(), 5);
2604
2605        // Try to retransmit all 10 bytes - should return 5 since first 5 are
2606        // acked.
2607        let retransmitted = stream.send.retransmit(0, 10);
2608        assert_eq!(retransmitted, 0); // Already marked, so no change
2609        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        // Write some data.
2617        assert_eq!(stream.send.write(b"hello", false), Ok(5));
2618        assert_eq!(stream.send.buffered_bytes(), 5);
2619
2620        // Emit it.
2621        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        // Mark for retransmission.
2627        let retransmitted = stream.send.retransmit(0, 5);
2628        assert_eq!(retransmitted, 5);
2629        assert_eq!(stream.send.buffered_bytes(), 5);
2630
2631        // Ack and drop - should decrement len and return dropped amount.
2632        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        // Write and emit two chunks.
2642        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        // Mark both chunks for retransmission.
2652        let retransmitted = stream.send.retransmit(0, 10);
2653        assert_eq!(retransmitted, 10);
2654        assert_eq!(stream.send.buffered_bytes(), 10);
2655
2656        // Ack and drop only first chunk.
2657        let dropped = stream.send.ack_and_drop(0, 5);
2658        assert_eq!(dropped, 5);
2659        assert_eq!(stream.send.buffered_bytes(), 5);
2660
2661        // Ack and drop second chunk.
2662        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        // Write and emit data.
2672        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        // Ack data that's already been fully emitted and not retransmitted.
2678        // Nothing should be dropped since there's no buffered data.
2679        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        // This test verifies that StreamMap.tx_buffered stays in sync with
2687        // actual buffered data through a full lifecycle: write → emit →
2688        // retransmit → ack.
2689        let mut streams = <StreamMap>::new(5, 5, 15);
2690
2691        // Create a stream using low-level StreamMap interface.
2692        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        // Update peer stream limits to allow locally-initiated streams.
2704        streams.update_peer_max_streams_bidi(5);
2705        streams.update_peer_max_streams_uni(5);
2706
2707        let stream_id = 0u64;
2708
2709        // Writing raises `stream.send.buffered_bytes()` and `tx_buffered`.
2710        {
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        // Emitting lowers `stream.send.buffered_bytes()` and `tx_buffered`.
2728        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        // Retransmitting raises both values by the amount retransmitted.
2741        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        // Ack and drop: both stream.send.buffered_bytes() and tx_buffered
2752        // decrease by actual amount dropped.
2753        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        // Initially empty.
2769        assert_eq!(stream.send.buffered_bytes(), 0);
2770
2771        // After write.
2772        assert_eq!(stream.send.write(b"hello", false), Ok(5));
2773        assert_eq!(stream.send.buffered_bytes(), 5);
2774
2775        // After emit.
2776        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        // After retransmit.
2782        let retransmitted = stream.send.retransmit(0, 5);
2783        assert_eq!(retransmitted, 5);
2784        assert_eq!(stream.send.buffered_bytes(), 5);
2785
2786        // After ack_and_drop.
2787        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;