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 streams queued for writable notification. A stopped stream can
197    /// be removed before its error is reported to the application.
198    pub writable: RBTree<StreamWritablePriorityAdapter>,
199
200    /// Stopped streams still queued for writable notification.
201    stopped_writable: RBTree<StreamStoppedWritablePriorityAdapter>,
202
203    /// Set of stream IDs corresponding to streams that are almost out of flow
204    /// control credit and need to send MAX_STREAM_DATA. This is used to
205    /// generate a `StreamIter` of streams without having to iterate over the
206    /// full list of streams.
207    almost_full: StreamIdHashSet,
208
209    /// Set of stream IDs corresponding to streams that are blocked. The value
210    /// of the map elements represents the offset of the stream at which the
211    /// blocking occurred.
212    blocked: StreamIdHashMap<u64>,
213
214    /// Set of stream IDs corresponding to streams that are reset. The value
215    /// of the map elements is a tuple of the error code and final size values
216    /// to include in the RESET_STREAM frame.
217    reset: StreamIdHashMap<(u64, u64)>,
218
219    /// Set of stream IDs corresponding to streams that are shutdown on the
220    /// receive side, and need to send a STOP_SENDING frame. The value of the
221    /// map elements is the error code to include in the STOP_SENDING frame.
222    stopped: StreamIdHashMap<u64>,
223
224    /// The maximum size of a stream window.
225    max_stream_window: u64,
226
227    /// Total number of bytes in send buffers across all streams.
228    tx_buffered: usize,
229}
230
231impl<F: BufFactory> StreamMap<F> {
232    pub fn new(
233        max_streams_bidi: u64, max_streams_uni: u64, max_stream_window: u64,
234    ) -> Self {
235        StreamMap {
236            local_max_streams_bidi: max_streams_bidi,
237            local_max_streams_bidi_next: max_streams_bidi,
238            initial_max_streams_bidi: max_streams_bidi,
239
240            local_max_streams_uni: max_streams_uni,
241            local_max_streams_uni_next: max_streams_uni,
242            initial_max_streams_uni: max_streams_uni,
243
244            max_stream_window,
245
246            ..StreamMap::default()
247        }
248    }
249
250    /// Returns the stream with the given ID if it exists.
251    pub fn get(&self, id: u64) -> Option<&Stream<F>> {
252        self.streams.get(&id)
253    }
254
255    /// Returns the mutable stream with the given ID if it exists.
256    pub fn get_mut(&mut self, id: u64) -> Option<&mut Stream<F>> {
257        self.streams.get_mut(&id)
258    }
259
260    /// Returns the mutable stream with the given ID if it exists, or creates
261    /// a new one otherwise.
262    ///
263    /// The `local` parameter indicates whether the stream is locally
264    /// initiated. It validates the stream ID and selects the initial flow
265    /// control values from the local and remote transport parameters.
266    ///
267    /// This also takes care of enforcing both local and the peer's stream
268    /// count limits. If one of these limits is violated, the `StreamLimit`
269    /// error is returned.
270    pub(crate) fn get_or_create(
271        &mut self, id: u64, local_params: &crate::TransportParams,
272        peer_params: &crate::TransportParams, local: bool, is_server: bool,
273    ) -> Result<&mut Stream<F>> {
274        let (stream, is_new_and_writable) = match self.streams.entry(id) {
275            hash_map::Entry::Vacant(v) => {
276                // Stream has already been closed and garbage collected.
277                if self.collected.contains(id) {
278                    return Err(Error::Done);
279                }
280
281                if local != is_local(id, is_server) {
282                    return Err(Error::InvalidStreamState(id));
283                }
284
285                let (max_rx_data, max_tx_data) = match (local, is_bidi(id)) {
286                    // Locally-initiated bidirectional stream.
287                    (true, true) => (
288                        local_params.initial_max_stream_data_bidi_local,
289                        peer_params.initial_max_stream_data_bidi_remote,
290                    ),
291
292                    // Locally-initiated unidirectional stream.
293                    (true, false) => (0, peer_params.initial_max_stream_data_uni),
294
295                    // Remotely-initiated bidirectional stream.
296                    (false, true) => (
297                        local_params.initial_max_stream_data_bidi_remote,
298                        peer_params.initial_max_stream_data_bidi_local,
299                    ),
300
301                    // Remotely-initiated unidirectional stream.
302                    (false, false) =>
303                        (local_params.initial_max_stream_data_uni, 0),
304                };
305
306                // The two least significant bits from a stream id identify the
307                // type of stream. Truncate those bits to get the sequence for
308                // that stream type.
309                let stream_sequence = id >> 2;
310
311                // Enforce stream count limits.
312                match (is_local(id, is_server), is_bidi(id)) {
313                    (true, true) => {
314                        let n = cmp::max(
315                            self.local_opened_streams_bidi,
316                            stream_sequence + 1,
317                        );
318
319                        if n > self.peer_max_streams_bidi {
320                            return Err(Error::StreamLimit);
321                        }
322
323                        self.local_opened_streams_bidi = n;
324                    },
325
326                    (true, false) => {
327                        let n = cmp::max(
328                            self.local_opened_streams_uni,
329                            stream_sequence + 1,
330                        );
331
332                        if n > self.peer_max_streams_uni {
333                            return Err(Error::StreamLimit);
334                        }
335
336                        self.local_opened_streams_uni = n;
337                    },
338
339                    (false, true) => {
340                        let n = cmp::max(
341                            self.peer_opened_streams_bidi,
342                            stream_sequence + 1,
343                        );
344
345                        if n > self.local_max_streams_bidi {
346                            return Err(Error::StreamLimit);
347                        }
348
349                        self.peer_opened_streams_bidi = n;
350                    },
351
352                    (false, false) => {
353                        let n = cmp::max(
354                            self.peer_opened_streams_uni,
355                            stream_sequence + 1,
356                        );
357
358                        if n > self.local_max_streams_uni {
359                            return Err(Error::StreamLimit);
360                        }
361
362                        self.peer_opened_streams_uni = n;
363                    },
364                };
365
366                let initial_window = max_rx_data;
367                let s = Stream::new(
368                    id,
369                    max_rx_data,
370                    max_tx_data,
371                    local,
372                    initial_window,
373                    self.max_stream_window,
374                );
375
376                let is_writable = s.is_writable();
377
378                (v.insert(s), is_writable)
379            },
380
381            hash_map::Entry::Occupied(v) => (v.into_mut(), false),
382        };
383
384        // Newly created stream might already be writable due to initial flow
385        // control limits.
386        if is_new_and_writable {
387            self.writable.insert(Arc::clone(&stream.priority_key));
388        }
389
390        Ok(stream)
391    }
392
393    /// Adds the stream ID to the readable streams set.
394    ///
395    /// If the stream was already in the list, this does nothing.
396    pub fn insert_readable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
397        if !priority_key.readable.is_linked() {
398            self.readable.insert(Arc::clone(priority_key));
399        }
400    }
401
402    /// Removes the stream ID from the readable streams set.
403    pub fn remove_readable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
404        if !priority_key.readable.is_linked() {
405            return;
406        }
407
408        let mut c = {
409            let ptr = Arc::as_ptr(priority_key);
410            unsafe { self.readable.cursor_mut_from_ptr(ptr) }
411        };
412
413        c.remove();
414    }
415
416    /// Adds the stream ID to the writable streams set.
417    ///
418    /// This should also be called anytime a new stream is created, in addition
419    /// to when an existing stream becomes writable.
420    ///
421    /// If the stream was already in the list, this does nothing.
422    pub fn insert_writable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
423        if !priority_key.writable.is_linked() {
424            self.writable.insert(Arc::clone(priority_key));
425        }
426    }
427
428    /// Marks a newly stopped stream as queued for writable notification.
429    pub fn insert_stopped_writable(
430        &mut self, priority_key: &Arc<StreamPriorityKey>,
431    ) {
432        if !priority_key.stopped_writable.is_linked() {
433            self.stopped_writable.insert(Arc::clone(priority_key));
434        }
435
436        self.insert_writable(priority_key);
437    }
438
439    /// Removes the stream ID from the writable streams set.
440    ///
441    /// This should also be called anytime an existing stream stops being
442    /// writable.
443    pub fn remove_writable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
444        if priority_key.stopped_writable.is_linked() {
445            let ptr = Arc::as_ptr(priority_key);
446            let mut c = unsafe { self.stopped_writable.cursor_mut_from_ptr(ptr) };
447            c.remove();
448        }
449
450        if !priority_key.writable.is_linked() {
451            return;
452        }
453
454        let mut c = {
455            let ptr = Arc::as_ptr(priority_key);
456            unsafe { self.writable.cursor_mut_from_ptr(ptr) }
457        };
458
459        c.remove();
460    }
461
462    /// Adds the stream ID to the flushable streams set.
463    ///
464    /// If the stream was already in the list, this does nothing.
465    pub fn insert_flushable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
466        if !priority_key.flushable.is_linked() {
467            self.flushable.insert(Arc::clone(priority_key));
468        }
469    }
470
471    /// Removes the stream ID from the flushable streams set.
472    pub fn remove_flushable(&mut self, priority_key: &Arc<StreamPriorityKey>) {
473        if !priority_key.flushable.is_linked() {
474            return;
475        }
476
477        let mut c = {
478            let ptr = Arc::as_ptr(priority_key);
479            unsafe { self.flushable.cursor_mut_from_ptr(ptr) }
480        };
481
482        c.remove();
483    }
484
485    pub fn peek_flushable(&self) -> Option<Arc<StreamPriorityKey>> {
486        self.flushable.front().clone_pointer()
487    }
488
489    /// Updates the priorities of a stream.
490    pub fn update_priority(
491        &mut self, old: &Arc<StreamPriorityKey>, new: &Arc<StreamPriorityKey>,
492    ) {
493        if old.readable.is_linked() {
494            self.remove_readable(old);
495            self.readable.insert(Arc::clone(new));
496        }
497
498        if old.writable.is_linked() {
499            let stopped = old.stopped_writable.is_linked();
500            self.remove_writable(old);
501            if stopped {
502                self.insert_stopped_writable(new);
503            } else {
504                self.insert_writable(new);
505            }
506        }
507
508        if old.flushable.is_linked() {
509            self.remove_flushable(old);
510            self.flushable.insert(Arc::clone(new));
511        }
512    }
513
514    /// Adds the stream ID to the almost full streams set.
515    ///
516    /// If the stream was already in the list, this does nothing.
517    pub fn insert_almost_full(&mut self, stream_id: u64) {
518        self.almost_full.insert(stream_id);
519    }
520
521    /// Removes the stream ID from the almost full streams set.
522    pub fn remove_almost_full(&mut self, stream_id: u64) {
523        self.almost_full.remove(&stream_id);
524    }
525
526    /// Adds the stream ID to the blocked streams set with the
527    /// given offset value.
528    ///
529    /// If the stream was already in the list, this does nothing.
530    pub fn insert_blocked(&mut self, stream_id: u64, off: u64) {
531        self.blocked.insert(stream_id, off);
532    }
533
534    /// Removes the stream ID from the blocked streams set.
535    pub fn remove_blocked(&mut self, stream_id: u64) {
536        self.blocked.remove(&stream_id);
537    }
538
539    /// Adds the stream ID to the reset streams set with the
540    /// given error code and final size values.
541    ///
542    /// If the stream was already in the list, this does nothing.
543    pub fn insert_reset(
544        &mut self, stream_id: u64, error_code: u64, final_size: u64,
545    ) {
546        self.reset.insert(stream_id, (error_code, final_size));
547    }
548
549    /// Removes the stream ID from the reset streams set.
550    pub fn remove_reset(&mut self, stream_id: u64) {
551        self.reset.remove(&stream_id);
552    }
553
554    /// Adds the stream ID to the stopped streams set with the
555    /// given error code.
556    ///
557    /// If the stream was already in the list, this does nothing.
558    pub fn insert_stopped(&mut self, stream_id: u64, error_code: u64) {
559        self.stopped.insert(stream_id, error_code);
560    }
561
562    /// Removes the stream ID from the stopped streams set.
563    pub fn remove_stopped(&mut self, stream_id: u64) {
564        self.stopped.remove(&stream_id);
565    }
566
567    /// Updates the peer's maximum bidirectional stream count limit.
568    pub fn update_peer_max_streams_bidi(&mut self, v: u64) {
569        self.peer_max_streams_bidi = cmp::max(self.peer_max_streams_bidi, v);
570    }
571
572    /// Updates the peer's maximum unidirectional stream count limit.
573    pub fn update_peer_max_streams_uni(&mut self, v: u64) {
574        self.peer_max_streams_uni = cmp::max(self.peer_max_streams_uni, v);
575    }
576
577    /// Commits the new max_streams_bidi limit.
578    pub fn update_max_streams_bidi(&mut self) {
579        self.local_max_streams_bidi = self.local_max_streams_bidi_next;
580    }
581
582    /// Sets the max_streams_bidi limit to the given value.
583    pub fn set_max_streams_bidi(&mut self, max: u64) {
584        self.local_max_streams_bidi = max;
585        self.local_max_streams_bidi_next = max;
586        self.initial_max_streams_bidi = max;
587    }
588
589    /// Returns the current max_streams_bidi limit.
590    pub fn max_streams_bidi(&self) -> u64 {
591        self.local_max_streams_bidi
592    }
593
594    /// Returns the new max_streams_bidi limit.
595    pub fn max_streams_bidi_next(&mut self) -> u64 {
596        self.local_max_streams_bidi_next
597    }
598
599    /// Commits the new max_streams_uni limit.
600    pub fn update_max_streams_uni(&mut self) {
601        self.local_max_streams_uni = self.local_max_streams_uni_next;
602    }
603
604    /// Returns the new max_streams_uni limit.
605    pub fn max_streams_uni_next(&mut self) -> u64 {
606        self.local_max_streams_uni_next
607    }
608
609    /// Returns the peer's current maximum bidirectional stream count limit.
610    pub fn peer_max_streams_bidi(&self) -> u64 {
611        self.peer_max_streams_bidi
612    }
613
614    /// Returns the number of bidirectional streams that can be created
615    /// before the peer's stream count limit is reached.
616    pub fn peer_streams_left_bidi(&self) -> u64 {
617        self.peer_max_streams_bidi - self.local_opened_streams_bidi
618    }
619
620    /// Returns the peer's current maximum unidirectional stream count limit.
621    pub fn peer_max_streams_uni(&self) -> u64 {
622        self.peer_max_streams_uni
623    }
624
625    /// Returns the number of unidirectional streams that can be created
626    /// before the peer's stream count limit is reached.
627    pub fn peer_streams_left_uni(&self) -> u64 {
628        self.peer_max_streams_uni - self.local_opened_streams_uni
629    }
630
631    /// Updates stream state before its STOP error is returned to the caller.
632    pub fn mark_stop_reported(&mut self, stream_id: u64) {
633        let stream = self.streams.get_mut(&stream_id).unwrap();
634        stream.send.mark_stop_reported();
635
636        if stream.is_collectable() {
637            let local = stream.local;
638            self.collect(stream_id, local);
639        } else {
640            let priority_key = Arc::clone(&stream.priority_key);
641            self.remove_writable(&priority_key);
642        }
643    }
644
645    /// Drops completed stream.
646    ///
647    /// This should only be called when Stream::is_collectable() returns true
648    /// for the given stream.
649    pub fn collect(&mut self, stream_id: u64, local: bool) {
650        if !local {
651            // If the stream was created by the peer, give back a max streams
652            // credit.
653            if is_bidi(stream_id) {
654                self.local_max_streams_bidi_next =
655                    self.local_max_streams_bidi_next.saturating_add(1);
656            } else {
657                self.local_max_streams_uni_next =
658                    self.local_max_streams_uni_next.saturating_add(1);
659            }
660        }
661
662        let s = self.streams.remove(&stream_id).unwrap();
663
664        self.remove_readable(&s.priority_key);
665
666        self.remove_writable(&s.priority_key);
667
668        self.remove_flushable(&s.priority_key);
669
670        self.collected.insert(stream_id);
671    }
672
673    /// Creates an iterator over streams that have outstanding data to read.
674    pub fn readable(&self) -> StreamIter {
675        StreamIter {
676            streams: self.readable.iter().map(|s| s.id).collect(),
677            index: 0,
678        }
679    }
680
681    /// Creates an iterator over streams that can be written to.
682    pub fn writable(&self) -> StreamIter {
683        StreamIter {
684            streams: self.writable.iter().map(|s| s.id).collect(),
685            index: 0,
686        }
687    }
688
689    /// Creates an iterator over stopped streams in the writable queue.
690    pub fn stopped_writable(&self) -> StreamIter {
691        if self.stopped_writable.is_empty() {
692            return StreamIter::default();
693        }
694
695        StreamIter {
696            streams: self.stopped_writable.iter().map(|key| key.id).collect(),
697            index: 0,
698        }
699    }
700
701    /// Returns and unlinks the next stopped stream awaiting notification.
702    pub fn pop_stopped_writable(&mut self) -> Option<u64> {
703        let priority_key = self.stopped_writable.front().clone_pointer()?;
704        self.remove_writable(&priority_key);
705        Some(priority_key.id)
706    }
707
708    /// Returns true if the given local stream has been opened.
709    ///
710    /// Opening a local stream implicitly opens all lower-numbered local streams
711    /// of the same type, even if they do not have `Stream` objects yet.
712    pub fn local_stream_opened(&self, stream_id: u64) -> bool {
713        let opened = if is_bidi(stream_id) {
714            self.local_opened_streams_bidi
715        } else {
716            self.local_opened_streams_uni
717        };
718
719        stream_id >> 2 < opened
720    }
721
722    /// Creates an iterator over streams that need to send MAX_STREAM_DATA.
723    pub fn almost_full(&self) -> StreamIter {
724        StreamIter::from(&self.almost_full)
725    }
726
727    /// Creates an iterator over streams that need to send STREAM_DATA_BLOCKED.
728    pub fn blocked(&self) -> hash_map::Iter<'_, u64, u64> {
729        self.blocked.iter()
730    }
731
732    /// Creates an iterator over streams that need to send RESET_STREAM.
733    pub fn reset(&self) -> hash_map::Iter<'_, u64, (u64, u64)> {
734        self.reset.iter()
735    }
736
737    /// Creates an iterator over streams that need to send STOP_SENDING.
738    pub fn stopped(&self) -> hash_map::Iter<'_, u64, u64> {
739        self.stopped.iter()
740    }
741
742    /// Returns true if the stream has been collected.
743    pub fn is_collected(&self, stream_id: u64) -> bool {
744        self.collected.contains(stream_id)
745    }
746
747    /// Returns true if there are any streams that have data to write.
748    pub fn has_flushable(&self) -> bool {
749        !self.flushable.is_empty()
750    }
751
752    /// Returns true if there are any streams that have data to read.
753    pub fn has_readable(&self) -> bool {
754        !self.readable.is_empty()
755    }
756
757    /// Returns true if there are any streams that need to update the local
758    /// flow control limit.
759    pub fn has_almost_full(&self) -> bool {
760        !self.almost_full.is_empty()
761    }
762
763    /// Returns true if there are any streams that are blocked.
764    pub fn has_blocked(&self) -> bool {
765        !self.blocked.is_empty()
766    }
767
768    /// Returns true if there are any streams that are reset.
769    pub fn has_reset(&self) -> bool {
770        !self.reset.is_empty()
771    }
772
773    /// Returns true if there are any streams that need to send STOP_SENDING.
774    pub fn has_stopped(&self) -> bool {
775        !self.stopped.is_empty()
776    }
777
778    /// Returns true if the max bidirectional streams count needs to be updated
779    /// by sending a MAX_STREAMS frame to the peer.
780    ///
781    /// This only sends MAX_STREAMS when available capacity is at or below 50%
782    /// of the initial maximum streams target.
783    pub fn should_update_max_streams_bidi(&self) -> bool {
784        let available = self
785            .local_max_streams_bidi
786            .saturating_sub(self.peer_opened_streams_bidi);
787        self.local_max_streams_bidi_next != self.local_max_streams_bidi &&
788            available <= self.initial_max_streams_bidi / 2
789    }
790
791    /// Returns true if the max unidirectional streams count needs to be updated
792    /// by sending a MAX_STREAMS frame to the peer.
793    ///
794    /// This only send MAX_STREAMS when available capacity is at or below 50% of
795    /// the initial maximum streams target.
796    pub fn should_update_max_streams_uni(&self) -> bool {
797        let available = self
798            .local_max_streams_uni
799            .saturating_sub(self.peer_opened_streams_uni);
800        self.local_max_streams_uni_next != self.local_max_streams_uni &&
801            available <= self.initial_max_streams_uni / 2
802    }
803
804    /// Returns the number of active streams in the map.
805    #[cfg(test)]
806    pub fn len(&self) -> usize {
807        self.streams.len()
808    }
809
810    /// Returns the total number of bytes buffered across all streams.
811    pub(crate) fn tx_buffered(&self) -> usize {
812        self.tx_buffered
813    }
814
815    /// Computes the actual number of bytes in send buffers by summing across
816    /// all streams. This is used for debugging to verify that tx_buffered
817    /// is accurate.
818    fn tx_buffered_actual(&self) -> usize {
819        self.streams
820            .values()
821            .map(|s| s.send.buffered_bytes() as usize)
822            .sum()
823    }
824
825    /// Checks if the stored tx_buffered matches the actual value.
826    /// Returns true if they match, false otherwise.
827    pub(crate) fn tx_buffered_is_consistent(&self) -> bool {
828        self.tx_buffered == self.tx_buffered_actual()
829    }
830
831    /// Updates the tx_buffered value by adding the delta.
832    pub(crate) fn add_tx_buffered(&mut self, delta: usize) {
833        self.tx_buffered += delta;
834
835        #[cfg(debug_assertions)]
836        self.debug_check_tx_buffered_consistency();
837    }
838
839    /// Updates the tx_buffered value by subtracting the delta.
840    pub(crate) fn sub_tx_buffered(&mut self, delta: usize) {
841        debug_assert!(self.tx_buffered >= delta);
842        self.tx_buffered = self.tx_buffered.saturating_sub(delta);
843
844        #[cfg(debug_assertions)]
845        self.debug_check_tx_buffered_consistency();
846    }
847
848    /// Verifies that the stored tx_buffered value matches the actual bytes in
849    /// send buffers across all streams. Enabled in debug builds to catch
850    /// inconsistencies early.
851    #[cfg(debug_assertions)]
852    pub(crate) fn debug_check_tx_buffered_consistency(&self) {
853        if !self.tx_buffered_is_consistent() {
854            let buffered_per_stream = self
855                .streams
856                .iter()
857                .map(|(id, s)| (*id, s.send.buffered_bytes()))
858                .collect::<Vec<_>>();
859
860            let actual = self.tx_buffered_actual();
861            let stored = self.tx_buffered;
862            panic!(
863                "tx_buffered mismatch: stored={}, actual={}, diff={}, buffered_per_stream={:?}",
864                stored,
865                actual,
866                stored as i64 - actual as i64,
867                buffered_per_stream
868            );
869        }
870    }
871}
872
873/// A QUIC stream.
874pub struct Stream<F: BufFactory = DefaultBufFactory> {
875    /// Receive-side stream buffer.
876    pub recv: recv_buf::RecvBuf,
877
878    /// Send-side stream buffer.
879    pub send: send_buf::SendBuf<F>,
880
881    pub send_lowat: usize,
882
883    /// Whether the stream is bidirectional.
884    pub bidi: bool,
885
886    /// Whether the stream was created by the local endpoint.
887    pub local: bool,
888
889    /// The stream's urgency (lower is better). Default is `DEFAULT_URGENCY`.
890    pub urgency: u8,
891
892    /// Whether the stream can be flushed incrementally. Default is `true`.
893    pub incremental: bool,
894
895    pub priority_key: Arc<StreamPriorityKey>,
896}
897
898impl<F: BufFactory> Stream<F> {
899    /// Creates a new stream with the given flow control limits.
900    pub fn new(
901        id: u64, max_rx_data: u64, max_tx_data: u64, local: bool,
902        initial_window: u64, max_window: u64,
903    ) -> Self {
904        let priority_key = Arc::new(StreamPriorityKey {
905            id,
906            ..Default::default()
907        });
908
909        Stream {
910            recv: recv_buf::RecvBuf::new(max_rx_data, initial_window, max_window),
911            send: send_buf::SendBuf::new(max_tx_data),
912            send_lowat: 1,
913            bidi: is_bidi(id),
914            local,
915            urgency: priority_key.urgency,
916            incremental: priority_key.incremental,
917            priority_key,
918        }
919    }
920
921    /// Returns true if the stream has data to read.
922    pub fn is_readable(&self) -> bool {
923        self.recv.ready()
924    }
925
926    /// Returns true if the stream has enough flow control capacity to be
927    /// written to, and is not finished.
928    pub fn is_writable(&self) -> bool {
929        !self.send.is_shutdown() &&
930            !self.send.is_fin() &&
931            (self.send.off_back() + self.send_lowat as u64) <
932                self.send.max_off()
933    }
934
935    /// Returns true if the stream has data to send and is allowed to send at
936    /// least some of it.
937    pub fn is_flushable(&self) -> bool {
938        let off_front = self.send.off_front();
939
940        !self.send.is_empty() &&
941            off_front < self.send.off_back() &&
942            off_front < self.send.max_off()
943    }
944
945    /// Returns true if the stream is complete.
946    ///
947    /// For bidirectional streams this happens when both the receive and send
948    /// sides are complete. That is when all incoming data has been read by the
949    /// application, and when all outgoing data has been acked by the peer.
950    ///
951    /// For unidirectional streams this happens when either the receive or send
952    /// side is complete, depending on whether the stream was created locally
953    /// or not.
954    pub fn is_complete(&self) -> bool {
955        match (self.bidi, self.local) {
956            // For bidirectional streams we need to check both receive and send
957            // sides for completion.
958            (true, _) => self.recv.is_fin() && self.send.is_complete(),
959
960            // For unidirectional streams generated locally, we only need to
961            // check the send side for completion.
962            (false, true) => self.send.is_complete(),
963
964            // For unidirectional streams generated by the peer, we only need
965            // to check the receive side for completion.
966            (false, false) => self.recv.is_fin(),
967        }
968    }
969
970    /// Returns true when no receive data or STOP error remains to deliver.
971    pub fn is_collectable(&self) -> bool {
972        self.is_complete() &&
973            !self.is_readable() &&
974            !self.send.has_unreported_stop()
975    }
976}
977
978/// Returns true if the stream was created locally.
979pub fn is_local(stream_id: u64, is_server: bool) -> bool {
980    (stream_id & 0x1) == (is_server as u64)
981}
982
983/// Returns true if the stream is bidirectional.
984pub fn is_bidi(stream_id: u64) -> bool {
985    (stream_id & 0x2) == 0
986}
987
988#[derive(Clone, Debug)]
989pub struct StreamPriorityKey {
990    pub urgency: u8,
991    pub incremental: bool,
992    pub id: u64,
993
994    pub readable: RBTreeAtomicLink,
995    pub writable: RBTreeAtomicLink,
996    pub stopped_writable: RBTreeAtomicLink,
997    pub flushable: RBTreeAtomicLink,
998}
999
1000impl Default for StreamPriorityKey {
1001    fn default() -> Self {
1002        Self {
1003            urgency: DEFAULT_URGENCY,
1004            incremental: true,
1005            id: Default::default(),
1006            readable: Default::default(),
1007            writable: Default::default(),
1008            stopped_writable: Default::default(),
1009            flushable: Default::default(),
1010        }
1011    }
1012}
1013
1014impl PartialEq for StreamPriorityKey {
1015    fn eq(&self, other: &Self) -> bool {
1016        self.id == other.id
1017    }
1018}
1019
1020impl Eq for StreamPriorityKey {}
1021
1022impl PartialOrd for StreamPriorityKey {
1023    fn partial_cmp(&self, other: &Self) -> Option<cmp::Ordering> {
1024        Some(self.cmp(other))
1025    }
1026}
1027
1028impl Ord for StreamPriorityKey {
1029    fn cmp(&self, other: &Self) -> cmp::Ordering {
1030        // Ignore priority if ID matches.
1031        if self.id == other.id {
1032            return cmp::Ordering::Equal;
1033        }
1034
1035        // First, order by urgency...
1036        if self.urgency != other.urgency {
1037            return self.urgency.cmp(&other.urgency);
1038        }
1039
1040        // ...when the urgency is the same, and both are not incremental, order
1041        // by stream ID...
1042        if !self.incremental && !other.incremental {
1043            return self.id.cmp(&other.id);
1044        }
1045
1046        // ...non-incremental takes priority over incremental...
1047        if self.incremental && !other.incremental {
1048            return cmp::Ordering::Greater;
1049        }
1050        if !self.incremental && other.incremental {
1051            return cmp::Ordering::Less;
1052        }
1053
1054        // ...finally, when both are incremental, `other` takes precedence (so
1055        // `self` is always sorted after other same-urgency incremental
1056        // entries).
1057        cmp::Ordering::Greater
1058    }
1059}
1060
1061intrusive_adapter!(pub StreamWritablePriorityAdapter = Arc<StreamPriorityKey>: StreamPriorityKey { writable => RBTreeAtomicLink });
1062
1063impl KeyAdapter<'_> for StreamWritablePriorityAdapter {
1064    type Key = StreamPriorityKey;
1065
1066    fn get_key(&self, s: &StreamPriorityKey) -> Self::Key {
1067        s.clone()
1068    }
1069}
1070
1071intrusive_adapter!(pub StreamStoppedWritablePriorityAdapter = Arc<StreamPriorityKey>: StreamPriorityKey { stopped_writable => RBTreeAtomicLink });
1072
1073impl KeyAdapter<'_> for StreamStoppedWritablePriorityAdapter {
1074    type Key = StreamPriorityKey;
1075
1076    fn get_key(&self, s: &StreamPriorityKey) -> Self::Key {
1077        s.clone()
1078    }
1079}
1080
1081intrusive_adapter!(pub StreamReadablePriorityAdapter = Arc<StreamPriorityKey>: StreamPriorityKey { readable => RBTreeAtomicLink });
1082
1083impl KeyAdapter<'_> for StreamReadablePriorityAdapter {
1084    type Key = StreamPriorityKey;
1085
1086    fn get_key(&self, s: &StreamPriorityKey) -> Self::Key {
1087        s.clone()
1088    }
1089}
1090
1091intrusive_adapter!(pub StreamFlushablePriorityAdapter = Arc<StreamPriorityKey>: StreamPriorityKey { flushable => RBTreeAtomicLink });
1092
1093impl KeyAdapter<'_> for StreamFlushablePriorityAdapter {
1094    type Key = StreamPriorityKey;
1095
1096    fn get_key(&self, s: &StreamPriorityKey) -> Self::Key {
1097        s.clone()
1098    }
1099}
1100
1101/// An iterator over QUIC streams.
1102#[derive(Default)]
1103pub struct StreamIter {
1104    streams: SmallVec<[u64; 8]>,
1105    index: usize,
1106}
1107
1108impl StreamIter {
1109    #[inline]
1110    fn from(streams: &StreamIdHashSet) -> Self {
1111        StreamIter {
1112            streams: streams.iter().copied().collect(),
1113            index: 0,
1114        }
1115    }
1116}
1117
1118impl Iterator for StreamIter {
1119    type Item = u64;
1120
1121    #[inline]
1122    fn next(&mut self) -> Option<Self::Item> {
1123        let v = self.streams.get(self.index)?;
1124        self.index += 1;
1125        Some(*v)
1126    }
1127}
1128
1129impl ExactSizeIterator for StreamIter {
1130    #[inline]
1131    fn len(&self) -> usize {
1132        self.streams.len() - self.index
1133    }
1134}
1135
1136#[cfg(test)]
1137mod tests {
1138    use rstest::rstest;
1139
1140    use crate::range_buf::RangeBuf;
1141    use crate::test_utils::Pipe;
1142
1143    use super::*;
1144
1145    /// The default size of the receiver stream flow control window.
1146    const DEFAULT_STREAM_WINDOW: u64 = 32 * 1024;
1147
1148    #[rstest]
1149    fn collected_streams_per_type(#[values(0, 1, 2, 3)] stream_type: u64) {
1150        let mut collected = CollectedStreams::default();
1151        assert!(collected.ranges.is_none());
1152
1153        for id in 0..4 {
1154            assert!(!collected.contains(id));
1155        }
1156        assert!(collected.ranges.is_none());
1157
1158        collected.insert(stream_type);
1159        assert!(collected.ranges.is_some());
1160
1161        for id in 0..4 {
1162            assert_eq!(collected.contains(id), id == stream_type);
1163        }
1164
1165        collected.insert(8 | stream_type);
1166        assert!(!collected.contains(4 | stream_type));
1167        assert_eq!(
1168            collected.ranges.as_ref().unwrap()[stream_type as usize].len(),
1169            2
1170        );
1171
1172        collected.insert(4 | stream_type);
1173        collected.insert(4 | stream_type);
1174        assert_eq!(
1175            collected.ranges.as_ref().unwrap()[stream_type as usize],
1176            0..3
1177        );
1178
1179        for id in 0..12 {
1180            assert_eq!(collected.contains(id), id & 0x3 == stream_type);
1181        }
1182    }
1183
1184    #[test]
1185    fn collected_streams_preserve_fragmented_and_large_ids() {
1186        let mut collected = CollectedStreams::default();
1187
1188        for sequence in (0..2048).step_by(2) {
1189            for stream_type in 0..4 {
1190                collected.insert((sequence << 2) | stream_type);
1191            }
1192        }
1193
1194        for sequence in 0..2048 {
1195            for stream_type in 0..4 {
1196                assert_eq!(
1197                    collected.contains((sequence << 2) | stream_type),
1198                    sequence % 2 == 0
1199                );
1200            }
1201        }
1202
1203        for stream_type in 0..4 {
1204            assert_eq!(
1205                collected.ranges.as_ref().unwrap()[stream_type as usize].len(),
1206                1024
1207            );
1208
1209            let stream_id = ((1u64 << 62) - 4) | stream_type;
1210            assert!(!collected.contains(stream_id));
1211            collected.insert(stream_id);
1212            assert!(collected.contains(stream_id));
1213            assert!(!collected.contains(stream_id - 4));
1214            assert!(!collected.contains((1u64 << 40) | stream_type));
1215        }
1216    }
1217
1218    /// Completes a client-initiated stream and processes returned stream
1219    /// credit.
1220    fn collect_pipe_stream(pipe: &mut Pipe, stream_id: u64) {
1221        let mut buf = [0; 1];
1222
1223        assert_eq!(pipe.client.stream_send(stream_id, b"a", true), Ok(1));
1224        assert_eq!(pipe.advance(), Ok(()));
1225        assert_eq!(pipe.server.stream_recv(stream_id, &mut buf), Ok((1, true)));
1226
1227        if is_bidi(stream_id) {
1228            assert_eq!(pipe.server.stream_send(stream_id, b"a", true), Ok(1));
1229            assert_eq!(pipe.advance(), Ok(()));
1230            assert_eq!(
1231                pipe.client.stream_recv(stream_id, &mut buf),
1232                Ok((1, true))
1233            );
1234        }
1235
1236        assert_eq!(pipe.advance(), Ok(()));
1237    }
1238
1239    #[rstest]
1240    fn collected_streams_out_of_order(
1241        #[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
1242        #[values(0, 2)] stream_type: u64,
1243    ) {
1244        let mut pipe = Pipe::new(cc_algorithm_name).unwrap();
1245        assert_eq!(pipe.handshake(), Ok(()));
1246
1247        for sequence in [2, 0, 1] {
1248            let stream_id = (sequence << 2) | stream_type;
1249            collect_pipe_stream(&mut pipe, stream_id);
1250
1251            for conn in [&mut pipe.client, &mut pipe.server] {
1252                assert!(conn.streams.get(stream_id).is_none());
1253                assert!(conn.streams.is_collected(stream_id));
1254                assert!(conn.stream_closed(stream_id));
1255                assert_eq!(
1256                    conn.streams
1257                        .get_or_create(
1258                            stream_id,
1259                            &conn.local_transport_params,
1260                            &conn.peer_transport_params,
1261                            !conn.is_server,
1262                            conn.is_server,
1263                        )
1264                        .err(),
1265                    Some(Error::Done)
1266                );
1267            }
1268        }
1269
1270        assert_eq!(
1271            pipe.server.streams.collected.ranges.as_ref().unwrap()
1272                [stream_type as usize],
1273            0..3
1274        );
1275
1276        // Late STREAM frames must not materialize a collected stream again.
1277        let frames = [crate::frame::Frame::Stream {
1278            stream_id: 8 | stream_type,
1279            data: RangeBuf::from(b"a", 0, true),
1280        }];
1281        assert!(pipe
1282            .send_pkt_to_server(crate::Type::Short, &frames, &mut [0; 1280])
1283            .is_ok());
1284        assert_eq!(pipe.server.streams.len(), 0);
1285    }
1286
1287    #[rstest]
1288    fn collected_streams_sparse_peer_credit(
1289        #[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
1290        #[values(0, 2)] stream_type: u64, #[values(1, 8)] initial_limit: u64,
1291    ) {
1292        let mut config = Pipe::default_config(cc_algorithm_name).unwrap();
1293        config.set_initial_max_streams_bidi(initial_limit);
1294        config.set_initial_max_streams_uni(initial_limit);
1295
1296        let mut pipe = Pipe::with_config(&mut config).unwrap();
1297        assert_eq!(pipe.handshake(), Ok(()));
1298
1299        for index in 0..initial_limit {
1300            let stream_id = (index << 3) | stream_type;
1301            collect_pipe_stream(&mut pipe, stream_id);
1302
1303            assert_eq!(
1304                pipe.server.streams.collected.ranges.as_ref().unwrap()
1305                    [stream_type as usize]
1306                    .len(),
1307                index as usize + 1
1308            );
1309            assert!(pipe.server.stream_closed(stream_id));
1310            assert!(pipe.client.stream_closed(stream_id));
1311            assert_eq!(pipe.server.streams.len(), 0);
1312        }
1313
1314        let peer_limit = if is_bidi(stream_type) {
1315            pipe.client.streams.peer_max_streams_bidi()
1316        } else {
1317            pipe.client.streams.peer_max_streams_uni()
1318        };
1319        assert_eq!(peer_limit, 2 * initial_limit);
1320
1321        // The odd sequences consume credit but have never had Stream objects.
1322        for sequence in (1..2 * initial_limit - 1).step_by(2) {
1323            let stream_id = (sequence << 2) | stream_type;
1324
1325            for conn in [&pipe.client, &pipe.server] {
1326                assert!(conn.streams.get(stream_id).is_none());
1327                assert!(!conn.streams.is_collected(stream_id));
1328                assert!(!conn.stream_closed(stream_id));
1329            }
1330        }
1331
1332        // Filling the implicit gaps is still permitted and merges all ranges.
1333        for sequence in (1..2 * initial_limit - 1).step_by(2) {
1334            collect_pipe_stream(&mut pipe, (sequence << 2) | stream_type);
1335        }
1336
1337        assert_eq!(
1338            pipe.server.streams.collected.ranges.as_ref().unwrap()
1339                [stream_type as usize],
1340            0..2 * initial_limit - 1
1341        );
1342    }
1343
1344    #[rstest]
1345    fn collected_streams_fragmented_local(
1346        #[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
1347        #[values(0, 2)] stream_type: u64,
1348    ) {
1349        let mut client_config = Pipe::default_config(cc_algorithm_name).unwrap();
1350        client_config.set_initial_max_streams_bidi(1);
1351        client_config.set_initial_max_streams_uni(1);
1352
1353        let mut server_config = Pipe::default_config(cc_algorithm_name).unwrap();
1354        server_config.set_initial_max_streams_bidi(32);
1355        server_config.set_initial_max_streams_uni(32);
1356
1357        let mut pipe = Pipe::with_client_and_server_config(
1358            &mut client_config,
1359            &mut server_config,
1360        )
1361        .unwrap();
1362        assert_eq!(pipe.handshake(), Ok(()));
1363
1364        // Local fragmentation follows the peer's allowance and application
1365        // behavior, not the client's initial incoming stream limit of one.
1366        for index in 0..32 {
1367            collect_pipe_stream(&mut pipe, (index << 3) | stream_type);
1368            assert_eq!(
1369                pipe.client.streams.collected.ranges.as_ref().unwrap()
1370                    [stream_type as usize]
1371                    .len(),
1372                index as usize + 1
1373            );
1374            assert_eq!(pipe.client.streams.len(), 0);
1375        }
1376    }
1377
1378    #[rstest]
1379    fn stream_limit_does_not_collect(
1380        #[values(0, 1, 2, 3)] stream_type: u64,
1381        #[values(true, false)] local: bool,
1382    ) {
1383        let params = crate::TransportParams::default();
1384        let mut streams = <StreamMap>::new(1, 1, DEFAULT_STREAM_WINDOW);
1385        streams.update_peer_max_streams_bidi(1);
1386        streams.update_peer_max_streams_uni(1);
1387
1388        let stream_id = 4 | stream_type;
1389        let is_server = (stream_type & 1 != 0) == local;
1390        assert_eq!(
1391            streams
1392                .get_or_create(stream_id, &params, &params, local, is_server)
1393                .err(),
1394            Some(Error::StreamLimit)
1395        );
1396        assert!(!streams.is_collected(stream_id));
1397        assert_eq!(streams.len(), 0);
1398        assert!(streams.collected.ranges.is_none());
1399
1400        if local {
1401            streams.update_peer_max_streams_bidi(2);
1402            streams.update_peer_max_streams_uni(2);
1403        } else {
1404            streams.local_max_streams_bidi_next = 2;
1405            streams.local_max_streams_uni_next = 2;
1406            streams.update_max_streams_bidi();
1407            streams.update_max_streams_uni();
1408        }
1409
1410        assert!(streams
1411            .get_or_create(stream_id, &params, &params, local, is_server)
1412            .is_ok());
1413        assert!(streams.collected.ranges.is_none());
1414    }
1415
1416    #[test]
1417    fn recv_flow_control() {
1418        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1419        assert!(!stream.recv.almost_full());
1420
1421        let mut buf = [0; 32];
1422
1423        let first = RangeBuf::from(b"hello", 0, false);
1424        let second = RangeBuf::from(b"world", 5, false);
1425        let third = RangeBuf::from(b"something", 10, false);
1426
1427        assert_eq!(stream.recv.write(second), Ok(()));
1428        assert_eq!(stream.recv.write(first), Ok(()));
1429        assert!(!stream.recv.almost_full());
1430
1431        assert_eq!(stream.recv.write(third), Err(Error::FlowControl));
1432
1433        let (len, fin) = stream.recv.emit(&mut buf).unwrap();
1434        assert_eq!(&buf[..len], b"helloworld");
1435        assert!(!fin);
1436
1437        assert!(stream.recv.almost_full());
1438
1439        stream.recv.update_max_data(std::time::Instant::now());
1440        assert_eq!(stream.recv.max_data_next(), 25);
1441        assert!(!stream.recv.almost_full());
1442
1443        let third = RangeBuf::from(b"something", 10, false);
1444        assert_eq!(stream.recv.write(third), Ok(()));
1445    }
1446
1447    #[test]
1448    fn recv_past_fin() {
1449        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1450        assert!(!stream.recv.almost_full());
1451
1452        let first = RangeBuf::from(b"hello", 0, true);
1453        let second = RangeBuf::from(b"world", 5, false);
1454
1455        assert_eq!(stream.recv.write(first), Ok(()));
1456        assert_eq!(stream.recv.write(second), Err(Error::FinalSize));
1457    }
1458
1459    #[test]
1460    fn recv_fin_dup() {
1461        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1462        assert!(!stream.recv.almost_full());
1463
1464        let first = RangeBuf::from(b"hello", 0, true);
1465        let second = RangeBuf::from(b"hello", 0, true);
1466
1467        assert_eq!(stream.recv.write(first), Ok(()));
1468        assert_eq!(stream.recv.write(second), Ok(()));
1469
1470        let mut buf = [0; 32];
1471
1472        let (len, fin) = stream.recv.emit(&mut buf).unwrap();
1473        assert_eq!(&buf[..len], b"hello");
1474        assert!(fin);
1475    }
1476
1477    #[test]
1478    fn recv_fin_change() {
1479        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1480        assert!(!stream.recv.almost_full());
1481
1482        let first = RangeBuf::from(b"hello", 0, true);
1483        let second = RangeBuf::from(b"world", 5, true);
1484
1485        assert_eq!(stream.recv.write(second), Ok(()));
1486        assert_eq!(stream.recv.write(first), Err(Error::FinalSize));
1487    }
1488
1489    #[test]
1490    fn recv_fin_lower_than_received() {
1491        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1492        assert!(!stream.recv.almost_full());
1493
1494        let first = RangeBuf::from(b"hello", 0, true);
1495        let second = RangeBuf::from(b"world", 5, false);
1496
1497        assert_eq!(stream.recv.write(second), Ok(()));
1498        assert_eq!(stream.recv.write(first), Err(Error::FinalSize));
1499    }
1500
1501    #[test]
1502    fn recv_fin_flow_control() {
1503        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1504        assert!(!stream.recv.almost_full());
1505
1506        let mut buf = [0; 32];
1507
1508        let first = RangeBuf::from(b"hello", 0, false);
1509        let second = RangeBuf::from(b"world", 5, true);
1510
1511        assert_eq!(stream.recv.write(first), Ok(()));
1512        assert_eq!(stream.recv.write(second), Ok(()));
1513
1514        let (len, fin) = stream.recv.emit(&mut buf).unwrap();
1515        assert_eq!(&buf[..len], b"helloworld");
1516        assert!(fin);
1517
1518        assert!(!stream.recv.almost_full());
1519    }
1520
1521    #[test]
1522    fn recv_fin_reset_mismatch() {
1523        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1524        assert!(!stream.recv.almost_full());
1525
1526        let first = RangeBuf::from(b"hello", 0, true);
1527
1528        assert_eq!(stream.recv.write(first), Ok(()));
1529        assert_eq!(stream.recv.reset(0, 10), Err(Error::FinalSize));
1530    }
1531
1532    #[test]
1533    fn recv_reset_with_gap() {
1534        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1535        assert!(!stream.recv.almost_full());
1536
1537        let first = RangeBuf::from(b"hello", 0, false);
1538
1539        assert_eq!(stream.recv.write(first), Ok(()));
1540        // Read one byte.
1541        assert_eq!(stream.recv.emit(&mut [0; 1]), Ok((1, false)));
1542        // Reset with a final size > than max previously received
1543        assert_eq!(
1544            stream.recv.reset(0, 10),
1545            Ok(RecvBufResetReturn {
1546                max_data_delta: 5,
1547                // consumed_flowcontrol is 9, since we already read 1 byte
1548                consumed_flowcontrol: 9
1549            })
1550        );
1551        assert_eq!(stream.recv.reset(0, 10), Ok(RecvBufResetReturn::zero()));
1552    }
1553
1554    #[test]
1555    fn recv_reset_dup() {
1556        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1557        assert!(!stream.recv.almost_full());
1558
1559        let first = RangeBuf::from(b"hello", 0, false);
1560
1561        assert_eq!(stream.recv.write(first), Ok(()));
1562        assert_eq!(
1563            stream.recv.reset(0, 5),
1564            Ok(RecvBufResetReturn {
1565                max_data_delta: 0,
1566                consumed_flowcontrol: 5
1567            })
1568        );
1569        assert_eq!(stream.recv.reset(0, 5), Ok(RecvBufResetReturn::zero()));
1570    }
1571
1572    #[test]
1573    fn recv_reset_change() {
1574        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1575        assert!(!stream.recv.almost_full());
1576
1577        let first = RangeBuf::from(b"hello", 0, false);
1578
1579        assert_eq!(stream.recv.write(first), Ok(()));
1580        assert_eq!(
1581            stream.recv.reset(0, 5),
1582            Ok(RecvBufResetReturn {
1583                max_data_delta: 0,
1584                consumed_flowcontrol: 5
1585            })
1586        );
1587        assert_eq!(stream.recv.reset(0, 10), Err(Error::FinalSize));
1588    }
1589
1590    #[test]
1591    fn recv_reset_lower_than_received() {
1592        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1593        assert!(!stream.recv.almost_full());
1594
1595        let first = RangeBuf::from(b"hello", 0, false);
1596
1597        assert_eq!(stream.recv.write(first), Ok(()));
1598        assert_eq!(stream.recv.reset(0, 4), Err(Error::FinalSize));
1599    }
1600
1601    #[test]
1602    fn send_flow_control() {
1603        let mut buf = [0; 25];
1604
1605        let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1606
1607        let first = b"hello";
1608        let second = b"world";
1609        let third = b"something";
1610
1611        assert!(stream.send.write(first, false).is_ok());
1612        assert!(stream.send.write(second, false).is_ok());
1613        assert!(stream.send.write(third, false).is_ok());
1614
1615        assert_eq!(stream.send.off_front(), 0);
1616
1617        let (written, fin) = stream.send.emit(&mut buf[..25]).unwrap();
1618        assert_eq!(written, 15);
1619        assert!(!fin);
1620        assert_eq!(&buf[..written], b"helloworldsomet");
1621
1622        assert_eq!(stream.send.off_front(), 15);
1623
1624        let (written, fin) = stream.send.emit(&mut buf[..25]).unwrap();
1625        assert_eq!(written, 0);
1626        assert!(!fin);
1627        assert_eq!(&buf[..written], b"");
1628
1629        stream.send.retransmit(0, 15);
1630
1631        assert_eq!(stream.send.off_front(), 0);
1632
1633        let (written, fin) = stream.send.emit(&mut buf[..10]).unwrap();
1634        assert_eq!(written, 10);
1635        assert!(!fin);
1636        assert_eq!(&buf[..written], b"helloworld");
1637
1638        assert_eq!(stream.send.off_front(), 10);
1639
1640        let (written, fin) = stream.send.emit(&mut buf[..10]).unwrap();
1641        assert_eq!(written, 5);
1642        assert!(!fin);
1643        assert_eq!(&buf[..written], b"somet");
1644    }
1645
1646    #[test]
1647    fn send_past_fin() {
1648        let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1649
1650        let first = b"hello";
1651        let second = b"world";
1652        let third = b"third";
1653
1654        assert_eq!(stream.send.write(first, false), Ok(5));
1655
1656        assert_eq!(stream.send.write(second, true), Ok(5));
1657        assert!(stream.send.is_fin());
1658
1659        assert_eq!(stream.send.write(third, false), Err(Error::FinalSize));
1660    }
1661
1662    #[test]
1663    fn send_fin_dup() {
1664        let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1665
1666        assert_eq!(stream.send.write(b"hello", true), Ok(5));
1667        assert!(stream.send.is_fin());
1668
1669        assert_eq!(stream.send.write(b"", true), Ok(0));
1670        assert!(stream.send.is_fin());
1671    }
1672
1673    #[test]
1674    fn send_undo_fin() {
1675        let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1676
1677        assert_eq!(stream.send.write(b"hello", true), Ok(5));
1678        assert!(stream.send.is_fin());
1679
1680        assert_eq!(
1681            stream.send.write(b"helloworld", true),
1682            Err(Error::FinalSize)
1683        );
1684    }
1685
1686    #[test]
1687    fn send_fin_max_data_match() {
1688        let mut buf = [0; 15];
1689
1690        let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1691
1692        let slice = b"hellohellohello";
1693
1694        assert!(stream.send.write(slice, true).is_ok());
1695
1696        let (written, fin) = stream.send.emit(&mut buf[..15]).unwrap();
1697        assert_eq!(written, 15);
1698        assert!(fin);
1699        assert_eq!(&buf[..written], slice);
1700    }
1701
1702    #[test]
1703    fn send_fin_zero_length() {
1704        let mut buf = [0; 5];
1705
1706        let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1707
1708        assert_eq!(stream.send.write(b"hello", false), Ok(5));
1709        assert_eq!(stream.send.write(b"", true), Ok(0));
1710        assert!(stream.send.is_fin());
1711
1712        let (written, fin) = stream.send.emit(&mut buf[..5]).unwrap();
1713        assert_eq!(written, 5);
1714        assert!(fin);
1715        assert_eq!(&buf[..written], b"hello");
1716    }
1717
1718    #[test]
1719    fn send_ack() {
1720        let mut buf = [0; 5];
1721
1722        let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1723
1724        assert_eq!(stream.send.write(b"hello", false), Ok(5));
1725        assert_eq!(stream.send.write(b"world", false), Ok(5));
1726        assert_eq!(stream.send.write(b"", true), Ok(0));
1727        assert!(stream.send.is_fin());
1728
1729        assert_eq!(stream.send.off_front(), 0);
1730
1731        let (written, fin) = stream.send.emit(&mut buf[..5]).unwrap();
1732        assert_eq!(written, 5);
1733        assert!(!fin);
1734        assert_eq!(&buf[..written], b"hello");
1735
1736        stream.send.ack_and_drop(0, 5);
1737
1738        stream.send.retransmit(0, 5);
1739
1740        assert_eq!(stream.send.off_front(), 5);
1741
1742        let (written, fin) = stream.send.emit(&mut buf[..5]).unwrap();
1743        assert_eq!(written, 5);
1744        assert!(fin);
1745        assert_eq!(&buf[..written], b"world");
1746    }
1747
1748    #[test]
1749    fn send_ack_reordering() {
1750        let mut buf = [0; 5];
1751
1752        let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1753
1754        assert_eq!(stream.send.write(b"hello", false), Ok(5));
1755        assert_eq!(stream.send.write(b"world", false), Ok(5));
1756        assert_eq!(stream.send.write(b"", true), Ok(0));
1757        assert!(stream.send.is_fin());
1758
1759        assert_eq!(stream.send.off_front(), 0);
1760
1761        let (written, fin) = stream.send.emit(&mut buf[..5]).unwrap();
1762        assert_eq!(written, 5);
1763        assert!(!fin);
1764        assert_eq!(&buf[..written], b"hello");
1765
1766        assert_eq!(stream.send.off_front(), 5);
1767
1768        let (written, fin) = stream.send.emit(&mut buf[..1]).unwrap();
1769        assert_eq!(written, 1);
1770        assert!(!fin);
1771        assert_eq!(&buf[..written], b"w");
1772
1773        stream.send.ack_and_drop(5, 1);
1774        stream.send.ack_and_drop(0, 5);
1775
1776        stream.send.retransmit(0, 5);
1777        stream.send.retransmit(5, 1);
1778
1779        assert_eq!(stream.send.off_front(), 6);
1780
1781        let (written, fin) = stream.send.emit(&mut buf[..5]).unwrap();
1782        assert_eq!(written, 4);
1783        assert!(fin);
1784        assert_eq!(&buf[..written], b"orld");
1785    }
1786
1787    #[test]
1788    fn recv_data_below_off() {
1789        let mut stream = <Stream>::new(0, 15, 0, true, 15, DEFAULT_STREAM_WINDOW);
1790
1791        let first = RangeBuf::from(b"hello", 0, false);
1792
1793        assert_eq!(stream.recv.write(first), Ok(()));
1794
1795        let mut buf = [0; 10];
1796
1797        let (len, fin) = stream.recv.emit(&mut buf).unwrap();
1798        assert_eq!(&buf[..len], b"hello");
1799        assert!(!fin);
1800
1801        let first = RangeBuf::from(b"elloworld", 1, true);
1802        assert_eq!(stream.recv.write(first), Ok(()));
1803
1804        let (len, fin) = stream.recv.emit(&mut buf).unwrap();
1805        assert_eq!(&buf[..len], b"world");
1806        assert!(fin);
1807    }
1808
1809    #[test]
1810    fn stream_complete() {
1811        let mut stream =
1812            <Stream>::new(0, 30, 30, true, 30, DEFAULT_STREAM_WINDOW);
1813
1814        assert_eq!(stream.send.write(b"hello", false), Ok(5));
1815        assert_eq!(stream.send.write(b"world", false), Ok(5));
1816
1817        assert!(!stream.send.is_complete());
1818        assert!(!stream.send.is_fin());
1819
1820        assert_eq!(stream.send.write(b"", true), Ok(0));
1821
1822        assert!(!stream.send.is_complete());
1823        assert!(stream.send.is_fin());
1824
1825        let buf = RangeBuf::from(b"hello", 0, true);
1826        assert!(stream.recv.write(buf).is_ok());
1827        assert!(!stream.recv.is_fin());
1828
1829        stream.send.ack(6, 4);
1830        assert!(!stream.send.is_complete());
1831
1832        let mut buf = [0; 2];
1833        assert_eq!(stream.recv.emit(&mut buf), Ok((2, false)));
1834        assert!(!stream.recv.is_fin());
1835
1836        stream.send.ack(1, 5);
1837        assert!(!stream.send.is_complete());
1838
1839        stream.send.ack(0, 1);
1840        assert!(!stream.send.is_complete());
1841
1842        stream.send.ack_fin();
1843        assert!(stream.send.is_complete());
1844
1845        assert!(!stream.is_complete());
1846
1847        let mut buf = [0; 3];
1848        assert_eq!(stream.recv.emit(&mut buf), Ok((3, true)));
1849        assert!(stream.recv.is_fin());
1850
1851        assert!(stream.is_complete());
1852    }
1853
1854    #[test]
1855    fn send_fin_zero_length_output() {
1856        let mut buf = [0; 5];
1857
1858        let mut stream = <Stream>::new(0, 0, 15, true, 0, DEFAULT_STREAM_WINDOW);
1859
1860        assert_eq!(stream.send.write(b"hello", false), Ok(5));
1861        assert_eq!(stream.send.off_front(), 0);
1862        assert!(!stream.send.is_fin());
1863
1864        let (written, fin) = stream.send.emit(&mut buf).unwrap();
1865        assert_eq!(written, 5);
1866        assert!(!fin);
1867        assert_eq!(&buf[..written], b"hello");
1868
1869        assert_eq!(stream.send.write(b"", true), Ok(0));
1870        assert!(stream.send.is_fin());
1871        assert_eq!(stream.send.off_front(), 5);
1872
1873        let (written, fin) = stream.send.emit(&mut buf).unwrap();
1874        assert_eq!(written, 0);
1875        assert!(fin);
1876        assert_eq!(&buf[..written], b"");
1877    }
1878
1879    fn stream_send_ready(stream: &Stream) -> bool {
1880        !stream.send.is_empty() &&
1881            stream.send.off_front() < stream.send.off_back()
1882    }
1883
1884    #[test]
1885    fn send_emit() {
1886        let mut buf = [0; 5];
1887
1888        let mut stream = <Stream>::new(0, 0, 20, true, 0, DEFAULT_STREAM_WINDOW);
1889
1890        assert_eq!(stream.send.write(b"hello", false), Ok(5));
1891        assert_eq!(stream.send.write(b"world", false), Ok(5));
1892        assert_eq!(stream.send.write(b"olleh", false), Ok(5));
1893        assert_eq!(stream.send.write(b"dlrow", true), Ok(5));
1894        assert_eq!(stream.send.off_front(), 0);
1895        assert_eq!(stream.send.bufs_count(), 4);
1896
1897        assert!(stream.is_flushable());
1898
1899        assert!(stream_send_ready(&stream));
1900        assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
1901        assert_eq!(stream.send.off_front(), 4);
1902        assert_eq!(&buf[..4], b"hell");
1903
1904        assert!(stream_send_ready(&stream));
1905        assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
1906        assert_eq!(stream.send.off_front(), 8);
1907        assert_eq!(&buf[..4], b"owor");
1908
1909        assert!(stream_send_ready(&stream));
1910        assert_eq!(stream.send.emit(&mut buf[..2]), Ok((2, false)));
1911        assert_eq!(stream.send.off_front(), 10);
1912        assert_eq!(&buf[..2], b"ld");
1913
1914        assert!(stream_send_ready(&stream));
1915        assert_eq!(stream.send.emit(&mut buf[..1]), Ok((1, false)));
1916        assert_eq!(stream.send.off_front(), 11);
1917        assert_eq!(&buf[..1], b"o");
1918
1919        assert!(stream_send_ready(&stream));
1920        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
1921        assert_eq!(stream.send.off_front(), 16);
1922        assert_eq!(&buf[..5], b"llehd");
1923
1924        assert!(stream_send_ready(&stream));
1925        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((4, true)));
1926        assert_eq!(stream.send.off_front(), 20);
1927        assert_eq!(&buf[..4], b"lrow");
1928
1929        assert!(!stream.is_flushable());
1930
1931        assert!(!stream_send_ready(&stream));
1932        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((0, true)));
1933        assert_eq!(stream.send.off_front(), 20);
1934    }
1935
1936    #[test]
1937    fn send_emit_ack() {
1938        let mut buf = [0; 5];
1939
1940        let mut stream = <Stream>::new(0, 0, 20, true, 0, DEFAULT_STREAM_WINDOW);
1941
1942        assert_eq!(stream.send.write(b"hello", false), Ok(5));
1943        assert_eq!(stream.send.write(b"world", false), Ok(5));
1944        assert_eq!(stream.send.write(b"olleh", false), Ok(5));
1945        assert_eq!(stream.send.write(b"dlrow", true), Ok(5));
1946        assert_eq!(stream.send.off_front(), 0);
1947        assert_eq!(stream.send.bufs_count(), 4);
1948
1949        assert!(stream.is_flushable());
1950
1951        assert!(stream_send_ready(&stream));
1952        assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
1953        assert_eq!(stream.send.off_front(), 4);
1954        assert_eq!(&buf[..4], b"hell");
1955
1956        assert!(stream_send_ready(&stream));
1957        assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
1958        assert_eq!(stream.send.off_front(), 8);
1959        assert_eq!(&buf[..4], b"owor");
1960
1961        stream.send.ack_and_drop(0, 5);
1962        assert_eq!(stream.send.bufs_count(), 3);
1963
1964        assert!(stream_send_ready(&stream));
1965        assert_eq!(stream.send.emit(&mut buf[..2]), Ok((2, false)));
1966        assert_eq!(stream.send.off_front(), 10);
1967        assert_eq!(&buf[..2], b"ld");
1968
1969        stream.send.ack_and_drop(7, 5);
1970        assert_eq!(stream.send.bufs_count(), 3);
1971
1972        assert!(stream_send_ready(&stream));
1973        assert_eq!(stream.send.emit(&mut buf[..1]), Ok((1, false)));
1974        assert_eq!(stream.send.off_front(), 11);
1975        assert_eq!(&buf[..1], b"o");
1976
1977        assert!(stream_send_ready(&stream));
1978        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
1979        assert_eq!(stream.send.off_front(), 16);
1980        assert_eq!(&buf[..5], b"llehd");
1981
1982        stream.send.ack_and_drop(5, 7);
1983        assert_eq!(stream.send.bufs_count(), 2);
1984
1985        assert!(stream_send_ready(&stream));
1986        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((4, true)));
1987        assert_eq!(stream.send.off_front(), 20);
1988        assert_eq!(&buf[..4], b"lrow");
1989
1990        assert!(!stream.is_flushable());
1991
1992        assert!(!stream_send_ready(&stream));
1993        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((0, true)));
1994        assert_eq!(stream.send.off_front(), 20);
1995
1996        stream.send.ack_and_drop(22, 4);
1997        assert_eq!(stream.send.bufs_count(), 2);
1998
1999        stream.send.ack_and_drop(20, 1);
2000        assert_eq!(stream.send.bufs_count(), 2);
2001    }
2002
2003    #[test]
2004    fn send_emit_retransmit() {
2005        let mut buf = [0; 5];
2006
2007        let mut stream = <Stream>::new(
2008            0,
2009            0,
2010            20,
2011            true,
2012            DEFAULT_STREAM_WINDOW,
2013            DEFAULT_STREAM_WINDOW,
2014        );
2015
2016        assert_eq!(stream.send.write(b"hello", false), Ok(5));
2017        assert_eq!(stream.send.write(b"world", false), Ok(5));
2018        assert_eq!(stream.send.write(b"olleh", false), Ok(5));
2019        assert_eq!(stream.send.write(b"dlrow", true), Ok(5));
2020        assert_eq!(stream.send.off_front(), 0);
2021        assert_eq!(stream.send.bufs_count(), 4);
2022
2023        assert!(stream.is_flushable());
2024
2025        assert!(stream_send_ready(&stream));
2026        assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
2027        assert_eq!(stream.send.off_front(), 4);
2028        assert_eq!(&buf[..4], b"hell");
2029
2030        assert!(stream_send_ready(&stream));
2031        assert_eq!(stream.send.emit(&mut buf[..4]), Ok((4, false)));
2032        assert_eq!(stream.send.off_front(), 8);
2033        assert_eq!(&buf[..4], b"owor");
2034
2035        stream.send.retransmit(3, 3);
2036        assert_eq!(stream.send.off_front(), 3);
2037
2038        assert!(stream_send_ready(&stream));
2039        assert_eq!(stream.send.emit(&mut buf[..3]), Ok((3, false)));
2040        assert_eq!(stream.send.off_front(), 8);
2041        assert_eq!(&buf[..3], b"low");
2042
2043        assert!(stream_send_ready(&stream));
2044        assert_eq!(stream.send.emit(&mut buf[..2]), Ok((2, false)));
2045        assert_eq!(stream.send.off_front(), 10);
2046        assert_eq!(&buf[..2], b"ld");
2047
2048        stream.send.ack_and_drop(7, 2);
2049
2050        stream.send.retransmit(8, 2);
2051
2052        assert!(stream_send_ready(&stream));
2053        assert_eq!(stream.send.emit(&mut buf[..2]), Ok((2, false)));
2054        assert_eq!(stream.send.off_front(), 10);
2055        assert_eq!(&buf[..2], b"ld");
2056
2057        assert!(stream_send_ready(&stream));
2058        assert_eq!(stream.send.emit(&mut buf[..1]), Ok((1, false)));
2059        assert_eq!(stream.send.off_front(), 11);
2060        assert_eq!(&buf[..1], b"o");
2061
2062        assert!(stream_send_ready(&stream));
2063        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
2064        assert_eq!(stream.send.off_front(), 16);
2065        assert_eq!(&buf[..5], b"llehd");
2066
2067        stream.send.retransmit(12, 2);
2068
2069        assert!(stream_send_ready(&stream));
2070        assert_eq!(stream.send.emit(&mut buf[..2]), Ok((2, false)));
2071        assert_eq!(stream.send.off_front(), 16);
2072        assert_eq!(&buf[..2], b"le");
2073
2074        assert!(stream_send_ready(&stream));
2075        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((4, true)));
2076        assert_eq!(stream.send.off_front(), 20);
2077        assert_eq!(&buf[..4], b"lrow");
2078
2079        assert!(!stream.is_flushable());
2080
2081        assert!(!stream_send_ready(&stream));
2082        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((0, true)));
2083        assert_eq!(stream.send.off_front(), 20);
2084
2085        stream.send.retransmit(7, 12);
2086
2087        assert!(stream_send_ready(&stream));
2088        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
2089        assert_eq!(stream.send.off_front(), 12);
2090        assert_eq!(&buf[..5], b"rldol");
2091
2092        assert!(stream_send_ready(&stream));
2093        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
2094        assert_eq!(stream.send.off_front(), 17);
2095        assert_eq!(&buf[..5], b"lehdl");
2096
2097        assert!(stream_send_ready(&stream));
2098        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((2, false)));
2099        assert_eq!(stream.send.off_front(), 20);
2100        assert_eq!(&buf[..2], b"ro");
2101
2102        stream.send.ack_and_drop(12, 7);
2103
2104        stream.send.retransmit(7, 12);
2105
2106        assert!(stream_send_ready(&stream));
2107        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
2108        assert_eq!(stream.send.off_front(), 12);
2109        assert_eq!(&buf[..5], b"rldol");
2110
2111        assert!(stream_send_ready(&stream));
2112        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((5, false)));
2113        assert_eq!(stream.send.off_front(), 17);
2114        assert_eq!(&buf[..5], b"lehdl");
2115
2116        assert!(stream_send_ready(&stream));
2117        assert_eq!(stream.send.emit(&mut buf[..5]), Ok((2, false)));
2118        assert_eq!(stream.send.off_front(), 20);
2119        assert_eq!(&buf[..2], b"ro");
2120    }
2121
2122    #[test]
2123    fn rangebuf_split_off() {
2124        let mut buf = <RangeBuf>::from(b"helloworld", 5, true);
2125        assert_eq!(buf.start, 0);
2126        assert_eq!(buf.pos, 0);
2127        assert_eq!(buf.len, 10);
2128        assert_eq!(buf.off, 5);
2129        assert!(buf.fin);
2130
2131        assert_eq!(buf.len(), 10);
2132        assert_eq!(buf.off(), 5);
2133        assert!(buf.fin());
2134
2135        assert_eq!(&buf[..], b"helloworld");
2136
2137        // Advance buffer.
2138        buf.consume(5);
2139
2140        assert_eq!(buf.start, 0);
2141        assert_eq!(buf.pos, 5);
2142        assert_eq!(buf.len, 10);
2143        assert_eq!(buf.off, 5);
2144        assert!(buf.fin);
2145
2146        assert_eq!(buf.len(), 5);
2147        assert_eq!(buf.off(), 10);
2148        assert!(buf.fin());
2149
2150        assert_eq!(&buf[..], b"world");
2151
2152        // Split buffer before position.
2153        let mut new_buf = buf.split_off(3);
2154
2155        assert_eq!(buf.start, 0);
2156        assert_eq!(buf.pos, 3);
2157        assert_eq!(buf.len, 3);
2158        assert_eq!(buf.off, 5);
2159        assert!(!buf.fin);
2160
2161        assert_eq!(buf.len(), 0);
2162        assert_eq!(buf.off(), 8);
2163        assert!(!buf.fin());
2164
2165        assert_eq!(&buf[..], b"");
2166
2167        assert_eq!(new_buf.start, 3);
2168        assert_eq!(new_buf.pos, 5);
2169        assert_eq!(new_buf.len, 7);
2170        assert_eq!(new_buf.off, 8);
2171        assert!(new_buf.fin);
2172
2173        assert_eq!(new_buf.len(), 5);
2174        assert_eq!(new_buf.off(), 10);
2175        assert!(new_buf.fin());
2176
2177        assert_eq!(&new_buf[..], b"world");
2178
2179        // Advance buffer.
2180        new_buf.consume(2);
2181
2182        assert_eq!(new_buf.start, 3);
2183        assert_eq!(new_buf.pos, 7);
2184        assert_eq!(new_buf.len, 7);
2185        assert_eq!(new_buf.off, 8);
2186        assert!(new_buf.fin);
2187
2188        assert_eq!(new_buf.len(), 3);
2189        assert_eq!(new_buf.off(), 12);
2190        assert!(new_buf.fin());
2191
2192        assert_eq!(&new_buf[..], b"rld");
2193
2194        // Split buffer after position.
2195        let mut new_new_buf = new_buf.split_off(5);
2196
2197        assert_eq!(new_buf.start, 3);
2198        assert_eq!(new_buf.pos, 7);
2199        assert_eq!(new_buf.len, 5);
2200        assert_eq!(new_buf.off, 8);
2201        assert!(!new_buf.fin);
2202
2203        assert_eq!(new_buf.len(), 1);
2204        assert_eq!(new_buf.off(), 12);
2205        assert!(!new_buf.fin());
2206
2207        assert_eq!(&new_buf[..], b"r");
2208
2209        assert_eq!(new_new_buf.start, 8);
2210        assert_eq!(new_new_buf.pos, 8);
2211        assert_eq!(new_new_buf.len, 2);
2212        assert_eq!(new_new_buf.off, 13);
2213        assert!(new_new_buf.fin);
2214
2215        assert_eq!(new_new_buf.len(), 2);
2216        assert_eq!(new_new_buf.off(), 13);
2217        assert!(new_new_buf.fin());
2218
2219        assert_eq!(&new_new_buf[..], b"ld");
2220
2221        // Advance buffer.
2222        new_new_buf.consume(2);
2223
2224        assert_eq!(new_new_buf.start, 8);
2225        assert_eq!(new_new_buf.pos, 10);
2226        assert_eq!(new_new_buf.len, 2);
2227        assert_eq!(new_new_buf.off, 13);
2228        assert!(new_new_buf.fin);
2229
2230        assert_eq!(new_new_buf.len(), 0);
2231        assert_eq!(new_new_buf.off(), 15);
2232        assert!(new_new_buf.fin());
2233
2234        assert_eq!(&new_new_buf[..], b"");
2235    }
2236
2237    /// RFC9000 2.1: A stream ID that is used out of order results in all
2238    /// streams of that type with lower-numbered stream IDs also being opened.
2239    #[test]
2240    fn stream_limit_auto_open() {
2241        let local_tp = crate::TransportParams::default();
2242        let peer_tp = crate::TransportParams::default();
2243
2244        let mut streams = <StreamMap>::new(5, 5, 5);
2245
2246        let stream_id = 500;
2247        assert!(!is_local(stream_id, true), "stream id is peer initiated");
2248        assert!(is_bidi(stream_id), "stream id is bidirectional");
2249        assert_eq!(
2250            streams
2251                .get_or_create(stream_id, &local_tp, &peer_tp, false, true)
2252                .err(),
2253            Some(Error::StreamLimit),
2254            "stream limit should be exceeded"
2255        );
2256    }
2257
2258    /// Stream limit should be satisfied regardless of what order we open
2259    /// streams
2260    #[test]
2261    fn stream_create_out_of_order() {
2262        let local_tp = crate::TransportParams::default();
2263        let peer_tp = crate::TransportParams::default();
2264
2265        let mut streams = <StreamMap>::new(5, 5, 5);
2266
2267        for stream_id in [8, 12, 4] {
2268            assert!(is_local(stream_id, false), "stream id is client initiated");
2269            assert!(is_bidi(stream_id), "stream id is bidirectional");
2270            assert!(streams
2271                .get_or_create(stream_id, &local_tp, &peer_tp, false, true)
2272                .is_ok());
2273        }
2274    }
2275
2276    /// Check stream limit boundary cases
2277    #[test]
2278    fn stream_limit_edge() {
2279        let local_tp = crate::TransportParams::default();
2280        let peer_tp = crate::TransportParams::default();
2281
2282        let mut streams = <StreamMap>::new(3, 3, 3);
2283
2284        // Highest permitted
2285        let stream_id = 8;
2286        assert!(streams
2287            .get_or_create(stream_id, &local_tp, &peer_tp, false, true)
2288            .is_ok());
2289
2290        // One more than highest permitted
2291        let stream_id = 12;
2292        assert_eq!(
2293            streams
2294                .get_or_create(stream_id, &local_tp, &peer_tp, false, true)
2295                .err(),
2296            Some(Error::StreamLimit)
2297        );
2298    }
2299
2300    fn cycle_stream_priority(stream_id: u64, streams: &mut StreamMap) {
2301        let key = streams.get(stream_id).unwrap().priority_key.clone();
2302        streams.update_priority(&key.clone(), &key);
2303    }
2304
2305    #[test]
2306    fn writable_prioritized_default_priority() {
2307        let local_tp = crate::TransportParams::default();
2308        let peer_tp = crate::TransportParams {
2309            initial_max_stream_data_bidi_local: 100,
2310            initial_max_stream_data_uni: 100,
2311            ..Default::default()
2312        };
2313
2314        let mut streams = StreamMap::new(100, 100, 100);
2315
2316        for id in [0, 4, 8, 12] {
2317            assert!(streams
2318                .get_or_create(id, &local_tp, &peer_tp, false, true)
2319                .is_ok());
2320        }
2321
2322        let walk_1: Vec<u64> = streams.writable().collect();
2323        cycle_stream_priority(*walk_1.first().unwrap(), &mut streams);
2324        let walk_2: Vec<u64> = streams.writable().collect();
2325        cycle_stream_priority(*walk_2.first().unwrap(), &mut streams);
2326        let walk_3: Vec<u64> = streams.writable().collect();
2327        cycle_stream_priority(*walk_3.first().unwrap(), &mut streams);
2328        let walk_4: Vec<u64> = streams.writable().collect();
2329        cycle_stream_priority(*walk_4.first().unwrap(), &mut streams);
2330        let walk_5: Vec<u64> = streams.writable().collect();
2331
2332        // All streams are non-incremental and same urgency by default. Multiple
2333        // visits shuffle their order.
2334        assert_eq!(walk_1, vec![0, 4, 8, 12]);
2335        assert_eq!(walk_2, vec![4, 8, 12, 0]);
2336        assert_eq!(walk_3, vec![8, 12, 0, 4]);
2337        assert_eq!(walk_4, vec![12, 0, 4, 8,]);
2338        assert_eq!(walk_5, vec![0, 4, 8, 12]);
2339    }
2340
2341    #[test]
2342    fn writable_prioritized_insert_order() {
2343        let local_tp = crate::TransportParams::default();
2344        let peer_tp = crate::TransportParams {
2345            initial_max_stream_data_bidi_local: 100,
2346            initial_max_stream_data_uni: 100,
2347            ..Default::default()
2348        };
2349
2350        let mut streams = StreamMap::new(100, 100, 100);
2351
2352        // Inserting same-urgency incremental streams in a "random" order yields
2353        // same order to start with.
2354        for id in [12, 4, 8, 0] {
2355            assert!(streams
2356                .get_or_create(id, &local_tp, &peer_tp, false, true)
2357                .is_ok());
2358        }
2359
2360        let walk_1: Vec<u64> = streams.writable().collect();
2361        cycle_stream_priority(*walk_1.first().unwrap(), &mut streams);
2362        let walk_2: Vec<u64> = streams.writable().collect();
2363        cycle_stream_priority(*walk_2.first().unwrap(), &mut streams);
2364        let walk_3: Vec<u64> = streams.writable().collect();
2365        cycle_stream_priority(*walk_3.first().unwrap(), &mut streams);
2366        let walk_4: Vec<u64> = streams.writable().collect();
2367        cycle_stream_priority(*walk_4.first().unwrap(), &mut streams);
2368        let walk_5: Vec<u64> = streams.writable().collect();
2369        assert_eq!(walk_1, vec![12, 4, 8, 0]);
2370        assert_eq!(walk_2, vec![4, 8, 0, 12]);
2371        assert_eq!(walk_3, vec![8, 0, 12, 4,]);
2372        assert_eq!(walk_4, vec![0, 12, 4, 8]);
2373        assert_eq!(walk_5, vec![12, 4, 8, 0]);
2374    }
2375
2376    #[test]
2377    fn writable_prioritized_mixed_urgency() {
2378        let local_tp = crate::TransportParams::default();
2379        let peer_tp = crate::TransportParams {
2380            initial_max_stream_data_bidi_local: 100,
2381            initial_max_stream_data_uni: 100,
2382            ..Default::default()
2383        };
2384
2385        let mut streams = <StreamMap>::new(100, 100, 100);
2386
2387        // Streams where the urgency descends (becomes more important). No
2388        // stream shares an urgency.
2389        let input = vec![
2390            (0, 100),
2391            (4, 90),
2392            (8, 80),
2393            (12, 70),
2394            (16, 60),
2395            (20, 50),
2396            (24, 40),
2397            (28, 30),
2398            (32, 20),
2399            (36, 10),
2400            (40, 0),
2401        ];
2402
2403        for (id, urgency) in input.clone() {
2404            // this duplicates some code from stream_priority in order to access
2405            // streams and the collection they're in
2406            let stream = streams
2407                .get_or_create(id, &local_tp, &peer_tp, false, true)
2408                .unwrap();
2409
2410            stream.urgency = urgency;
2411
2412            let new_priority_key = Arc::new(StreamPriorityKey {
2413                urgency: stream.urgency,
2414                incremental: stream.incremental,
2415                id,
2416                ..Default::default()
2417            });
2418
2419            let old_priority_key = std::mem::replace(
2420                &mut stream.priority_key,
2421                new_priority_key.clone(),
2422            );
2423
2424            streams.update_priority(&old_priority_key, &new_priority_key);
2425        }
2426
2427        let walk_1: Vec<u64> = streams.writable().collect();
2428        assert_eq!(walk_1, vec![40, 36, 32, 28, 24, 20, 16, 12, 8, 4, 0]);
2429
2430        // Re-applying priority to a stream does not cause duplication.
2431        for (id, urgency) in input {
2432            // this duplicates some code from stream_priority in order to access
2433            // streams and the collection they're in
2434            let stream = streams
2435                .get_or_create(id, &local_tp, &peer_tp, false, true)
2436                .unwrap();
2437
2438            stream.urgency = urgency;
2439
2440            let new_priority_key = Arc::new(StreamPriorityKey {
2441                urgency: stream.urgency,
2442                incremental: stream.incremental,
2443                id,
2444                ..Default::default()
2445            });
2446
2447            let old_priority_key = std::mem::replace(
2448                &mut stream.priority_key,
2449                new_priority_key.clone(),
2450            );
2451
2452            streams.update_priority(&old_priority_key, &new_priority_key);
2453        }
2454
2455        let walk_2: Vec<u64> = streams.writable().collect();
2456        assert_eq!(walk_2, vec![40, 36, 32, 28, 24, 20, 16, 12, 8, 4, 0]);
2457
2458        // Removing streams doesn't break expected ordering.
2459        streams.collect(24, true);
2460
2461        let walk_3: Vec<u64> = streams.writable().collect();
2462        assert_eq!(walk_3, vec![40, 36, 32, 28, 20, 16, 12, 8, 4, 0]);
2463
2464        streams.collect(40, true);
2465        streams.collect(0, true);
2466
2467        let walk_4: Vec<u64> = streams.writable().collect();
2468        assert_eq!(walk_4, vec![36, 32, 28, 20, 16, 12, 8, 4]);
2469
2470        // Adding streams doesn't break expected ordering.
2471        streams
2472            .get_or_create(44, &local_tp, &peer_tp, false, true)
2473            .unwrap();
2474
2475        let walk_5: Vec<u64> = streams.writable().collect();
2476        assert_eq!(walk_5, vec![36, 32, 28, 20, 16, 12, 8, 4, 44]);
2477    }
2478
2479    #[test]
2480    fn writable_prioritized_mixed_urgencies_incrementals() {
2481        let local_tp = crate::TransportParams::default();
2482        let peer_tp = crate::TransportParams {
2483            initial_max_stream_data_bidi_local: 100,
2484            initial_max_stream_data_uni: 100,
2485            ..Default::default()
2486        };
2487
2488        let mut streams = StreamMap::new(100, 100, 100);
2489
2490        // Streams that share some urgency level
2491        let input = vec![
2492            (0, 100),
2493            (4, 20),
2494            (8, 100),
2495            (12, 20),
2496            (16, 90),
2497            (20, 25),
2498            (24, 90),
2499            (28, 30),
2500            (32, 80),
2501            (36, 20),
2502            (40, 0),
2503        ];
2504
2505        for (id, urgency) in input.clone() {
2506            // this duplicates some code from stream_priority in order to access
2507            // streams and the collection they're in
2508            let stream = streams
2509                .get_or_create(id, &local_tp, &peer_tp, false, true)
2510                .unwrap();
2511
2512            stream.urgency = urgency;
2513
2514            let new_priority_key = Arc::new(StreamPriorityKey {
2515                urgency: stream.urgency,
2516                incremental: stream.incremental,
2517                id,
2518                ..Default::default()
2519            });
2520
2521            let old_priority_key = std::mem::replace(
2522                &mut stream.priority_key,
2523                new_priority_key.clone(),
2524            );
2525
2526            streams.update_priority(&old_priority_key, &new_priority_key);
2527        }
2528
2529        let walk_1: Vec<u64> = streams.writable().collect();
2530        cycle_stream_priority(4, &mut streams);
2531        cycle_stream_priority(16, &mut streams);
2532        cycle_stream_priority(0, &mut streams);
2533        let walk_2: Vec<u64> = streams.writable().collect();
2534        cycle_stream_priority(12, &mut streams);
2535        cycle_stream_priority(24, &mut streams);
2536        cycle_stream_priority(8, &mut streams);
2537        let walk_3: Vec<u64> = streams.writable().collect();
2538        cycle_stream_priority(36, &mut streams);
2539        cycle_stream_priority(16, &mut streams);
2540        cycle_stream_priority(0, &mut streams);
2541        let walk_4: Vec<u64> = streams.writable().collect();
2542        cycle_stream_priority(4, &mut streams);
2543        cycle_stream_priority(24, &mut streams);
2544        cycle_stream_priority(8, &mut streams);
2545        let walk_5: Vec<u64> = streams.writable().collect();
2546        cycle_stream_priority(12, &mut streams);
2547        cycle_stream_priority(16, &mut streams);
2548        cycle_stream_priority(0, &mut streams);
2549        let walk_6: Vec<u64> = streams.writable().collect();
2550        cycle_stream_priority(36, &mut streams);
2551        cycle_stream_priority(24, &mut streams);
2552        cycle_stream_priority(8, &mut streams);
2553        let walk_7: Vec<u64> = streams.writable().collect();
2554        cycle_stream_priority(4, &mut streams);
2555        cycle_stream_priority(16, &mut streams);
2556        cycle_stream_priority(0, &mut streams);
2557        let walk_8: Vec<u64> = streams.writable().collect();
2558        cycle_stream_priority(12, &mut streams);
2559        cycle_stream_priority(24, &mut streams);
2560        cycle_stream_priority(8, &mut streams);
2561        let walk_9: Vec<u64> = streams.writable().collect();
2562        cycle_stream_priority(36, &mut streams);
2563        cycle_stream_priority(16, &mut streams);
2564        cycle_stream_priority(0, &mut streams);
2565
2566        assert_eq!(walk_1, vec![40, 4, 12, 36, 20, 28, 32, 16, 24, 0, 8]);
2567        assert_eq!(walk_2, vec![40, 12, 36, 4, 20, 28, 32, 24, 16, 8, 0]);
2568        assert_eq!(walk_3, vec![40, 36, 4, 12, 20, 28, 32, 16, 24, 0, 8]);
2569        assert_eq!(walk_4, vec![40, 4, 12, 36, 20, 28, 32, 24, 16, 8, 0]);
2570        assert_eq!(walk_5, vec![40, 12, 36, 4, 20, 28, 32, 16, 24, 0, 8]);
2571        assert_eq!(walk_6, vec![40, 36, 4, 12, 20, 28, 32, 24, 16, 8, 0]);
2572        assert_eq!(walk_7, vec![40, 4, 12, 36, 20, 28, 32, 16, 24, 0, 8]);
2573        assert_eq!(walk_8, vec![40, 12, 36, 4, 20, 28, 32, 24, 16, 8, 0]);
2574        assert_eq!(walk_9, vec![40, 36, 4, 12, 20, 28, 32, 16, 24, 0, 8]);
2575
2576        // Removing streams doesn't break expected ordering.
2577        streams.collect(20, true);
2578
2579        let walk_10: Vec<u64> = streams.writable().collect();
2580        assert_eq!(walk_10, vec![40, 4, 12, 36, 28, 32, 24, 16, 8, 0]);
2581
2582        // Adding streams doesn't break expected ordering.
2583        let stream = streams
2584            .get_or_create(44, &local_tp, &peer_tp, false, true)
2585            .unwrap();
2586
2587        stream.urgency = 20;
2588        stream.incremental = true;
2589
2590        let new_priority_key = Arc::new(StreamPriorityKey {
2591            urgency: stream.urgency,
2592            incremental: stream.incremental,
2593            id: 44,
2594            ..Default::default()
2595        });
2596
2597        let old_priority_key =
2598            std::mem::replace(&mut stream.priority_key, new_priority_key.clone());
2599
2600        streams.update_priority(&old_priority_key, &new_priority_key);
2601
2602        let walk_11: Vec<u64> = streams.writable().collect();
2603        assert_eq!(walk_11, vec![40, 4, 12, 36, 44, 28, 32, 24, 16, 8, 0]);
2604    }
2605
2606    #[test]
2607    fn priority_tree_dupes() {
2608        let mut prioritized_writable: RBTree<StreamWritablePriorityAdapter> =
2609            Default::default();
2610
2611        for id in [0, 4, 8, 12] {
2612            let s = Arc::new(StreamPriorityKey {
2613                urgency: 0,
2614                incremental: false,
2615                id,
2616                ..Default::default()
2617            });
2618
2619            prioritized_writable.insert(s);
2620        }
2621
2622        let walk_1: Vec<u64> =
2623            prioritized_writable.iter().map(|s| s.id).collect();
2624        assert_eq!(walk_1, vec![0, 4, 8, 12]);
2625
2626        // Default keys could cause duplicate entries. `StreamMap` normally
2627        // prevents this.
2628        for id in [0, 4, 8, 12] {
2629            let s = Arc::new(StreamPriorityKey {
2630                urgency: 0,
2631                incremental: false,
2632                id,
2633                ..Default::default()
2634            });
2635
2636            prioritized_writable.insert(s);
2637        }
2638
2639        let walk_2: Vec<u64> =
2640            prioritized_writable.iter().map(|s| s.id).collect();
2641        assert_eq!(walk_2, vec![0, 0, 4, 4, 8, 8, 12, 12]);
2642    }
2643
2644    #[test]
2645    fn retransmit_returns_zero_when_already_acked() {
2646        let mut stream = <Stream>::new(0, 15, 15, true, 0, 15);
2647
2648        // Write and emit some data.
2649        assert_eq!(stream.send.write(b"hello", false), Ok(5));
2650        assert_eq!(stream.send.buffered_bytes(), 5);
2651
2652        let mut buf = [0; 10];
2653        let (written, _) = stream.send.emit(&mut buf).unwrap();
2654        assert_eq!(written, 5);
2655        assert_eq!(stream.send.buffered_bytes(), 0);
2656
2657        // Mark data for retransmission.
2658        let retransmitted = stream.send.retransmit(0, 5);
2659        assert_eq!(retransmitted, 5);
2660        assert_eq!(stream.send.buffered_bytes(), 5);
2661
2662        // Ack the data.
2663        stream.send.ack_and_drop(0, 5);
2664        assert_eq!(stream.send.buffered_bytes(), 0);
2665
2666        // Try to retransmit again - should return 0 since data is acked.
2667        let retransmitted = stream.send.retransmit(0, 5);
2668        assert_eq!(retransmitted, 0);
2669        assert_eq!(stream.send.buffered_bytes(), 0);
2670    }
2671
2672    #[test]
2673    fn retransmit_returns_partial_when_some_acked() {
2674        let mut stream = <Stream>::new(0, 15, 15, true, 0, 15);
2675
2676        // Write and emit 10 bytes.
2677        assert_eq!(stream.send.write(b"helloworld", false), Ok(10));
2678        assert_eq!(stream.send.buffered_bytes(), 10);
2679
2680        let mut buf = [0; 10];
2681        let (written, _) = stream.send.emit(&mut buf).unwrap();
2682        assert_eq!(written, 10);
2683        assert_eq!(stream.send.buffered_bytes(), 0);
2684
2685        // Mark all data for retransmission.
2686        let retransmitted = stream.send.retransmit(0, 10);
2687        assert_eq!(retransmitted, 10);
2688        assert_eq!(stream.send.buffered_bytes(), 10);
2689
2690        // Ack first 5 bytes and drop them.
2691        let dropped = stream.send.ack_and_drop(0, 5);
2692        assert_eq!(dropped, 5);
2693        assert_eq!(stream.send.buffered_bytes(), 5);
2694
2695        // Try to retransmit all 10 bytes - should return 5 since first 5 are
2696        // acked.
2697        let retransmitted = stream.send.retransmit(0, 10);
2698        assert_eq!(retransmitted, 0); // Already marked, so no change
2699        assert_eq!(stream.send.buffered_bytes(), 5);
2700    }
2701
2702    #[test]
2703    fn ack_and_drop_decrements_len_and_returns_dropped() {
2704        let mut stream = <Stream>::new(0, 15, 15, true, 0, 15);
2705
2706        // Write some data.
2707        assert_eq!(stream.send.write(b"hello", false), Ok(5));
2708        assert_eq!(stream.send.buffered_bytes(), 5);
2709
2710        // Emit it.
2711        let mut buf = [0; 10];
2712        let (written, _) = stream.send.emit(&mut buf).unwrap();
2713        assert_eq!(written, 5);
2714        assert_eq!(stream.send.buffered_bytes(), 0);
2715
2716        // Mark for retransmission.
2717        let retransmitted = stream.send.retransmit(0, 5);
2718        assert_eq!(retransmitted, 5);
2719        assert_eq!(stream.send.buffered_bytes(), 5);
2720
2721        // Ack and drop - should decrement len and return dropped amount.
2722        let dropped = stream.send.ack_and_drop(0, 5);
2723        assert_eq!(dropped, 5);
2724        assert_eq!(stream.send.buffered_bytes(), 0);
2725    }
2726
2727    #[test]
2728    fn ack_and_drop_partial_buffer() {
2729        let mut stream = <Stream>::new(0, 30, 30, true, 0, 30);
2730
2731        // Write and emit two chunks.
2732        assert_eq!(stream.send.write(b"hello", false), Ok(5));
2733        assert_eq!(stream.send.write(b"world", false), Ok(5));
2734        assert_eq!(stream.send.buffered_bytes(), 10);
2735
2736        let mut buf = [0; 10];
2737        let (written, _) = stream.send.emit(&mut buf).unwrap();
2738        assert_eq!(written, 10);
2739        assert_eq!(stream.send.buffered_bytes(), 0);
2740
2741        // Mark both chunks for retransmission.
2742        let retransmitted = stream.send.retransmit(0, 10);
2743        assert_eq!(retransmitted, 10);
2744        assert_eq!(stream.send.buffered_bytes(), 10);
2745
2746        // Ack and drop only first chunk.
2747        let dropped = stream.send.ack_and_drop(0, 5);
2748        assert_eq!(dropped, 5);
2749        assert_eq!(stream.send.buffered_bytes(), 5);
2750
2751        // Ack and drop second chunk.
2752        let dropped = stream.send.ack_and_drop(5, 5);
2753        assert_eq!(dropped, 5);
2754        assert_eq!(stream.send.buffered_bytes(), 0);
2755    }
2756
2757    #[test]
2758    fn ack_and_drop_returns_zero_when_nothing_dropped() {
2759        let mut stream = <Stream>::new(0, 15, 15, true, 0, 15);
2760
2761        // Write and emit data.
2762        assert_eq!(stream.send.write(b"hello", false), Ok(5));
2763        let mut buf = [0; 10];
2764        let (written, _) = stream.send.emit(&mut buf).unwrap();
2765        assert_eq!(written, 5);
2766
2767        // Ack data that's already been fully emitted and not retransmitted.
2768        // Nothing should be dropped since there's no buffered data.
2769        let dropped = stream.send.ack_and_drop(0, 5);
2770        assert_eq!(dropped, 0);
2771        assert_eq!(stream.send.buffered_bytes(), 0);
2772    }
2773
2774    #[test]
2775    fn cache_consistency_through_full_lifecycle() {
2776        // This test verifies that StreamMap.tx_buffered stays in sync with
2777        // actual buffered data through a full lifecycle: write → emit →
2778        // retransmit → ack.
2779        let mut streams = <StreamMap>::new(5, 5, 15);
2780
2781        // Create a stream using low-level StreamMap interface.
2782        let local_params = crate::TransportParams {
2783            initial_max_data: 30,
2784            initial_max_stream_data_bidi_local: 15,
2785            initial_max_stream_data_bidi_remote: 15,
2786            initial_max_stream_data_uni: 10,
2787            initial_max_streams_bidi: 5,
2788            initial_max_streams_uni: 5,
2789            ..Default::default()
2790        };
2791        let peer_params = local_params.clone();
2792
2793        // Update peer stream limits to allow locally-initiated streams.
2794        streams.update_peer_max_streams_bidi(5);
2795        streams.update_peer_max_streams_uni(5);
2796
2797        let stream_id = 0u64;
2798
2799        // Writing raises `stream.send.buffered_bytes()` and `tx_buffered`.
2800        {
2801            let stream = streams
2802                .get_or_create(
2803                    stream_id,
2804                    &local_params,
2805                    &peer_params,
2806                    true,
2807                    false,
2808                )
2809                .unwrap();
2810            assert_eq!(stream.send.write(b"hello", false), Ok(5));
2811        }
2812        streams.add_tx_buffered(5);
2813        assert_eq!(streams.get(stream_id).unwrap().send.buffered_bytes(), 5);
2814        assert_eq!(streams.tx_buffered(), 5);
2815        assert!(streams.tx_buffered_is_consistent());
2816
2817        // Emitting lowers `stream.send.buffered_bytes()` and `tx_buffered`.
2818        let mut buf = [0; 10];
2819        let written = {
2820            let stream = streams.get_mut(stream_id).unwrap();
2821            let (written, _) = stream.send.emit(&mut buf).unwrap();
2822            written
2823        };
2824        assert_eq!(written, 5);
2825        streams.sub_tx_buffered(5);
2826        assert_eq!(streams.get(stream_id).unwrap().send.buffered_bytes(), 0);
2827        assert_eq!(streams.tx_buffered(), 0);
2828        assert!(streams.tx_buffered_is_consistent());
2829
2830        // Retransmitting raises both values by the amount retransmitted.
2831        let retransmitted = {
2832            let stream = streams.get_mut(stream_id).unwrap();
2833            stream.send.retransmit(0, 5)
2834        };
2835        assert_eq!(retransmitted, 5);
2836        streams.add_tx_buffered(retransmitted);
2837        assert_eq!(streams.get(stream_id).unwrap().send.buffered_bytes(), 5);
2838        assert_eq!(streams.tx_buffered(), 5);
2839        assert!(streams.tx_buffered_is_consistent());
2840
2841        // Ack and drop: both stream.send.buffered_bytes() and tx_buffered
2842        // decrease by actual amount dropped.
2843        let dropped = {
2844            let stream = streams.get_mut(stream_id).unwrap();
2845            stream.send.ack_and_drop(0, 5)
2846        };
2847        assert_eq!(dropped, 5);
2848        streams.sub_tx_buffered(dropped);
2849        assert_eq!(streams.get(stream_id).unwrap().send.buffered_bytes(), 0);
2850        assert_eq!(streams.tx_buffered(), 0);
2851        assert!(streams.tx_buffered_is_consistent());
2852    }
2853
2854    #[test]
2855    fn send_buf_len_reflects_buffered_data() {
2856        let mut stream = <Stream>::new(0, 15, 15, true, 0, 15);
2857
2858        // Initially empty.
2859        assert_eq!(stream.send.buffered_bytes(), 0);
2860
2861        // After write.
2862        assert_eq!(stream.send.write(b"hello", false), Ok(5));
2863        assert_eq!(stream.send.buffered_bytes(), 5);
2864
2865        // After emit.
2866        let mut buf = [0; 10];
2867        let (written, _) = stream.send.emit(&mut buf).unwrap();
2868        assert_eq!(written, 5);
2869        assert_eq!(stream.send.buffered_bytes(), 0);
2870
2871        // After retransmit.
2872        let retransmitted = stream.send.retransmit(0, 5);
2873        assert_eq!(retransmitted, 5);
2874        assert_eq!(stream.send.buffered_bytes(), 5);
2875
2876        // After ack_and_drop.
2877        let dropped = stream.send.ack_and_drop(0, 5);
2878        assert_eq!(dropped, 5);
2879        assert_eq!(stream.send.buffered_bytes(), 0);
2880    }
2881}
2882
2883mod recv_buf;
2884mod send_buf;