Skip to main content

quiche/h3/
stream.rs

1// Copyright (C) 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::mem::MaybeUninit;
28
29use crate::buffers::BufFactory;
30
31use super::Error;
32use super::Result;
33
34use super::frame;
35
36pub const HTTP3_CONTROL_STREAM_TYPE_ID: u64 = 0x0;
37pub const HTTP3_PUSH_STREAM_TYPE_ID: u64 = 0x1;
38pub const QPACK_ENCODER_STREAM_TYPE_ID: u64 = 0x2;
39pub const QPACK_DECODER_STREAM_TYPE_ID: u64 = 0x3;
40
41const MAX_STATE_BUF_SIZE: usize = (1 << 24) - 1;
42const MAX_STATE_BUF_ALLOC_SIZE: usize = 4096;
43#[derive(Clone, Copy, Debug, PartialEq, Eq)]
44pub enum Type {
45    Control,
46    Request,
47    Push,
48    QpackEncoder,
49    QpackDecoder,
50    Unknown,
51}
52
53impl Type {
54    #[cfg(feature = "qlog")]
55    pub fn to_qlog(self) -> qlog::events::http3::StreamType {
56        match self {
57            Type::Control => qlog::events::http3::StreamType::Control,
58            Type::Request => qlog::events::http3::StreamType::Request,
59            Type::Push => qlog::events::http3::StreamType::Push,
60            Type::QpackEncoder => qlog::events::http3::StreamType::QpackEncode,
61            Type::QpackDecoder => qlog::events::http3::StreamType::QpackDecode,
62            Type::Unknown => qlog::events::http3::StreamType::Unknown,
63        }
64    }
65}
66
67#[derive(Clone, Copy, Debug, PartialEq, Eq)]
68pub enum State {
69    /// Reading the stream's type.
70    StreamType,
71
72    /// Reading the stream's current frame's type.
73    FrameType,
74
75    /// Reading the stream's current frame's payload length.
76    FramePayloadLen,
77
78    /// Reading the stream's current frame's payload.
79    FramePayload,
80
81    /// Reading DATA payload.
82    Data,
83
84    /// Reading the push ID.
85    PushId,
86
87    /// Reading a QPACK instruction.
88    QpackInstruction,
89
90    /// Reading the stream's current frame's payload without buffering the data.
91    SkipFramePayload,
92
93    /// Reading and discarding data.
94    Drain,
95
96    /// All data has been read.
97    Finished,
98}
99
100impl Type {
101    pub fn deserialize(v: u64) -> Result<Type> {
102        match v {
103            HTTP3_CONTROL_STREAM_TYPE_ID => Ok(Type::Control),
104            HTTP3_PUSH_STREAM_TYPE_ID => Ok(Type::Push),
105            QPACK_ENCODER_STREAM_TYPE_ID => Ok(Type::QpackEncoder),
106            QPACK_DECODER_STREAM_TYPE_ID => Ok(Type::QpackDecoder),
107
108            _ => Ok(Type::Unknown),
109        }
110    }
111}
112
113/// An HTTP/3 stream.
114///
115/// This maintains the HTTP/3 state for streams of any type (control, request,
116/// QPACK, ...).
117///
118/// A number of bytes, depending on the current stream's state, is read from the
119/// transport stream into the HTTP/3 stream's "state buffer". This intermediate
120/// buffering is required due to the fact that data read from the transport
121/// might not be complete (e.g. a varint might be split across multiple QUIC
122/// packets).
123///
124/// When enough data to complete the current state has been buffered, it is
125/// consumed from the state buffer and the stream is transitioned to the next
126/// state (see `State` for a list of possible states).
127#[derive(Debug)]
128pub struct Stream {
129    /// The corresponding transport stream's ID.
130    id: u64,
131
132    /// The stream's type (if known).
133    ty: Option<Type>,
134
135    /// The current stream state.
136    state: State,
137
138    /// The buffer holding partial data for the current state.
139    state_buf: Vec<u8>,
140
141    /// The expected amount of bytes required to complete the state.
142    state_len: usize,
143
144    /// The write offset in the state buffer, that is, how many bytes have
145    /// already been read from the transport for the current state. When
146    /// it reaches `stream_len` the state can be completed.
147    state_off: usize,
148
149    /// The type of the frame currently being parsed.
150    frame_type: Option<u64>,
151
152    /// Whether the stream was created locally, or by the peer.
153    is_local: bool,
154
155    /// Whether the stream has been remotely initialized.
156    remote_initialized: bool,
157
158    /// Whether the stream has been locally initialized.
159    local_initialized: bool,
160
161    /// Whether the local send-side of the stream has finished.
162    local_finished: bool,
163
164    /// Whether a `Data` event has been triggered for this stream.
165    data_event_triggered: bool,
166
167    /// The last `PRIORITY_UPDATE` frame encoded field value, if any.
168    last_priority_update: Option<Vec<u8>>,
169
170    /// The count of HEADERS frames that have been received.
171    headers_received_count: usize,
172
173    /// Whether a DATA frame has been received.
174    data_received: bool,
175
176    /// Whether a trailing HEADER field has been sent.
177    trailers_sent: bool,
178
179    /// Whether a trailing HEADER field has been received.
180    trailers_received: bool,
181
182    /// Max size of QPACK encoded headers carried in HEADERS or PUSH_PROMISE
183    /// frames. Related to SETTINGS_MAX_FIELD_LIST_SIZE; see
184    /// <https://datatracker.ietf.org/doc/html/rfc9114#section-7.2.4.1>
185    max_encoded_headers_payload_size: u64,
186
187    /// Max PRIORITY_UPDATE frame payload size; see
188    /// <https://datatracker.ietf.org/doc/html/rfc9218#section-7.2>
189    max_priority_update_size: u64,
190}
191
192impl Stream {
193    /// Creates a new HTTP/3 stream.
194    ///
195    /// The `is_local` parameter indicates whether the stream was created by the
196    /// local endpoint, or by the peer.
197    pub fn new(
198        id: u64, is_local: bool, max_field_section_size: u64,
199        max_priority_update_size: u64,
200    ) -> Stream {
201        let (ty, state) = if crate::stream::is_bidi(id) {
202            // All bidirectional streams are "request" streams, so we don't
203            // need to read the stream type.
204            (Some(Type::Request), State::FrameType)
205        } else {
206            // The stream's type is yet to be determined.
207            (None, State::StreamType)
208        };
209
210        // Huffman encoding might inflate the size of encoded headers
211        // when transferred in HEADERS or PUSH_PROMISE frames. Scale up
212        // the limit by 50% to accommodate this.
213        let max_encoded_headers_payload_size =
214            max_field_section_size.saturating_add(max_field_section_size / 2);
215
216        Stream {
217            id,
218            ty,
219
220            state,
221
222            // Pre-allocate a buffer to avoid multiple tiny early allocations.
223            state_buf: Vec::with_capacity(16),
224
225            // Expect one byte for the initial state, to parse the initial
226            // varint length.
227            state_len: 1,
228            state_off: 0,
229
230            frame_type: None,
231
232            is_local,
233
234            remote_initialized: false,
235
236            local_initialized: false,
237            local_finished: false,
238
239            data_event_triggered: false,
240
241            last_priority_update: None,
242
243            headers_received_count: 0,
244
245            data_received: false,
246
247            trailers_sent: false,
248            trailers_received: false,
249
250            max_encoded_headers_payload_size,
251            max_priority_update_size,
252        }
253    }
254
255    pub fn ty(&self) -> Option<Type> {
256        self.ty
257    }
258
259    pub fn state(&self) -> State {
260        self.state
261    }
262
263    /// Sets the stream's type and transitions to the next state.
264    pub fn set_ty(&mut self, ty: Type) -> Result<()> {
265        assert_eq!(self.state, State::StreamType);
266
267        self.ty = Some(ty);
268
269        let state = match ty {
270            Type::Control | Type::Request => State::FrameType,
271
272            Type::Push => State::PushId,
273
274            Type::QpackEncoder | Type::QpackDecoder => {
275                self.remote_initialized = true;
276
277                State::QpackInstruction
278            },
279
280            Type::Unknown => State::Drain,
281        };
282
283        self.state_transition(state, 1, true)?;
284
285        Ok(())
286    }
287
288    /// Sets the push ID and transitions to the next state.
289    pub fn set_push_id(&mut self, _id: u64) -> Result<()> {
290        assert_eq!(self.state, State::PushId);
291
292        // TODO: implement push ID.
293
294        self.state_transition(State::FrameType, 1, true)?;
295
296        Ok(())
297    }
298
299    /// Sets the frame type and transitions to the next state.
300    pub fn set_frame_type(&mut self, ty: u64) -> Result<()> {
301        assert_eq!(self.state, State::FrameType);
302
303        // Only expect frames on Control, Request and Push streams.
304        match self.ty {
305            Some(Type::Control) => {
306                // Control stream starts uninitialized and only SETTINGS is
307                // accepted in that state. Other frames cause an error. Once
308                // initialized, no more SETTINGS are permitted.
309                match (ty, self.remote_initialized) {
310                    // Initialize control stream.
311                    (frame::SETTINGS_FRAME_TYPE_ID, false) =>
312                        self.remote_initialized = true,
313
314                    // Non-SETTINGS frames not allowed on control stream
315                    // before initialization.
316                    (_, false) => return Err(Error::MissingSettings),
317
318                    // Additional SETTINGS frame.
319                    (frame::SETTINGS_FRAME_TYPE_ID, true) =>
320                        return Err(Error::FrameUnexpected),
321
322                    // Frames that can't be received on control stream
323                    // after initialization.
324                    (frame::DATA_FRAME_TYPE_ID, true) =>
325                        return Err(Error::FrameUnexpected),
326
327                    (frame::HEADERS_FRAME_TYPE_ID, true) =>
328                        return Err(Error::FrameUnexpected),
329
330                    (frame::PUSH_PROMISE_FRAME_TYPE_ID, true) =>
331                        return Err(Error::FrameUnexpected),
332
333                    // All other frames are ignored after initialization.
334                    (_, true) => (),
335                }
336            },
337
338            Some(Type::Request) => {
339                self.validate_request_frame_type(ty)?;
340            },
341
342            Some(Type::Push) => {
343                match ty {
344                    // Frames that can never be received on request streams.
345                    frame::CANCEL_PUSH_FRAME_TYPE_ID =>
346                        return Err(Error::FrameUnexpected),
347
348                    frame::SETTINGS_FRAME_TYPE_ID =>
349                        return Err(Error::FrameUnexpected),
350
351                    frame::PUSH_PROMISE_FRAME_TYPE_ID =>
352                        return Err(Error::FrameUnexpected),
353
354                    frame::GOAWAY_FRAME_TYPE_ID =>
355                        return Err(Error::FrameUnexpected),
356
357                    frame::MAX_PUSH_FRAME_TYPE_ID =>
358                        return Err(Error::FrameUnexpected),
359
360                    _ => (),
361                }
362            },
363
364            _ => return Err(Error::FrameUnexpected),
365        }
366
367        self.frame_type = Some(ty);
368
369        self.state_transition(State::FramePayloadLen, 1, true)?;
370
371        Ok(())
372    }
373
374    /// Validates a frame type received on a request stream and advances the
375    /// request-stream HTTP message phase tracking accordingly.
376    ///
377    /// Request streams start uninitialized and only HEADERS is accepted. After
378    /// initialization, DATA and HEADERS frames may be acceptable, depending on
379    /// the HTTP message phase (informational/final headers, body, trailers).
380    /// Receiving any other known frame type on a request stream is always an
381    /// error per RFC 9114, regardless of which endpoint opened the stream.
382    ///
383    /// HTTP message phase bookkeeping (`remote_initialized`, `data_received`,
384    /// `trailers_received`) only applies to peer-initiated streams, since it
385    /// tracks the receive direction.
386    fn validate_request_frame_type(&mut self, ty: u64) -> Result<()> {
387        // Frames that are never valid on a request stream, regardless of which
388        // endpoint opened it.
389        if matches!(
390            ty,
391            frame::CANCEL_PUSH_FRAME_TYPE_ID |
392                frame::SETTINGS_FRAME_TYPE_ID |
393                frame::GOAWAY_FRAME_TYPE_ID |
394                frame::MAX_PUSH_FRAME_TYPE_ID |
395                frame::PRIORITY_UPDATE_FRAME_REQUEST_TYPE_ID |
396                frame::PRIORITY_UPDATE_FRAME_PUSH_TYPE_ID
397        ) {
398            return Err(Error::FrameUnexpected);
399        }
400
401        // HTTP message phase bookkeeping only applies to peer-initiated
402        // streams, since it tracks the receive direction.
403        if self.is_local {
404            return Ok(());
405        }
406
407        match (ty, self.remote_initialized) {
408            (frame::HEADERS_FRAME_TYPE_ID, false) => {
409                self.remote_initialized = true;
410            },
411
412            (frame::DATA_FRAME_TYPE_ID, false) =>
413                return Err(Error::FrameUnexpected),
414
415            (frame::HEADERS_FRAME_TYPE_ID, true) => {
416                if self.trailers_received {
417                    return Err(Error::FrameUnexpected);
418                }
419
420                if self.data_received {
421                    self.trailers_received = true;
422                }
423            },
424
425            (frame::DATA_FRAME_TYPE_ID, true) => {
426                if self.trailers_received {
427                    return Err(Error::FrameUnexpected);
428                }
429
430                self.data_received = true;
431            },
432
433            // All other frames can be ignored regardless of stream state.
434            _ => (),
435        }
436
437        Ok(())
438    }
439
440    // Returns the stream's current frame type, if any
441    pub fn frame_type(&self) -> Option<u64> {
442        self.frame_type
443    }
444
445    /// Sets the frame's payload length and transitions to the next state.
446    pub fn set_frame_payload_len(&mut self, len: u64) -> Result<()> {
447        assert_eq!(self.state, State::FramePayloadLen);
448
449        // Only expect frames on Control, Request and Push streams.
450        if !matches!(self.ty, Some(Type::Control | Type::Request | Type::Push)) {
451            return Err(Error::InternalError);
452        }
453
454        let (state, resize) = match self.frame_type {
455            Some(frame::DATA_FRAME_TYPE_ID) => (State::Data, false),
456
457            Some(frame::HEADERS_FRAME_TYPE_ID) => {
458                if len > self.max_encoded_headers_payload_size {
459                    return Err(Error::ExcessiveLoad);
460                }
461
462                (State::FramePayload, true)
463            },
464
465            // These frames carry a mandatory single varint, so their payload
466            // size has to be at least 1 byte and at most 8 bytes.
467            Some(frame::CANCEL_PUSH_FRAME_TYPE_ID) |
468            Some(frame::GOAWAY_FRAME_TYPE_ID) |
469            Some(frame::MAX_PUSH_FRAME_TYPE_ID) => {
470                if !(1..=8).contains(&len) {
471                    return Err(Error::FrameError);
472                }
473
474                (State::FramePayload, true)
475            },
476
477            Some(frame::SETTINGS_FRAME_TYPE_ID) => {
478                if len > frame::MAX_SETTINGS_PAYLOAD_SIZE as u64 {
479                    return Err(Error::FrameError);
480                }
481
482                (State::FramePayload, true)
483            },
484
485            Some(frame::PUSH_PROMISE_FRAME_TYPE_ID) => {
486                // A push promise payload includes a varint and a field section.
487                let max_push_promise_size =
488                    self.max_encoded_headers_payload_size.saturating_add(8);
489
490                if len == 0 {
491                    return Err(Error::FrameError);
492                }
493
494                if len > max_push_promise_size {
495                    return Err(Error::ExcessiveLoad);
496                }
497
498                (State::FramePayload, true)
499            },
500
501            Some(frame::PRIORITY_UPDATE_FRAME_REQUEST_TYPE_ID) |
502            Some(frame::PRIORITY_UPDATE_FRAME_PUSH_TYPE_ID) => {
503                if len == 0 {
504                    return Err(Error::FrameError);
505                }
506
507                if len > self.max_priority_update_size {
508                    return Err(Error::ExcessiveLoad);
509                }
510
511                (State::FramePayload, true)
512            },
513
514            // Ignore unknown frames' payloads.
515            _ => {
516                if len > MAX_STATE_BUF_SIZE as u64 {
517                    return Err(Error::ExcessiveLoad);
518                }
519
520                (State::SkipFramePayload, false)
521            },
522        };
523
524        self.state_transition(state, len as usize, resize)?;
525
526        Ok(())
527    }
528
529    /// Returns a mutable slice over the state buffer's spare capacity,
530    /// reserving additional space if needed.
531    ///
532    /// Callers must initialize the returned bytes and call
533    /// [`commit_state_buf_read()`] with the number of bytes written.
534    fn spare_state_buf(&mut self) -> &mut [u8] {
535        let need = self.state_len - self.state_off;
536        let spare = self
537            .state_buf
538            .capacity()
539            .saturating_sub(self.state_buf.len());
540
541        if spare == 0 {
542            let additional = std::cmp::min(MAX_STATE_BUF_ALLOC_SIZE, need);
543            self.state_buf.reserve(additional);
544        }
545
546        let buf = self.state_buf.spare_capacity_mut();
547        let usable = std::cmp::min(need, buf.len());
548
549        // SAFETY: MaybeUninit<u8> has the same layout as u8. The caller
550        // contract requires initializing all returned bytes before calling
551        // commit_state_buf_read(), so no uninitialized memory is read.
552        unsafe {
553            std::mem::transmute::<&mut [MaybeUninit<u8>], &mut [u8]>(
554                &mut buf[..usable],
555            )
556        }
557    }
558
559    /// Advances the state buffer's initialized length and offset by `read`
560    /// bytes. The caller must have written `read` bytes into the slice
561    /// previously returned by [`spare_state_buf()`].
562    fn commit_state_buf_read(&mut self, read: usize) {
563        let buf_len = self.state_buf.len();
564        debug_assert!(buf_len + read <= self.state_buf.capacity());
565        // SAFETY: `read` bytes were written to spare capacity by the I/O
566        // layer. read <= spare_buf.len() <= spare_capacity, so
567        // buf_len + read <= capacity.
568        unsafe { self.state_buf.set_len(buf_len + read) };
569
570        self.state_off += read;
571    }
572
573    /// Tries to fill the state buffer by reading data from the corresponding
574    /// transport stream.
575    ///
576    /// When not enough data can be read to complete the state, this returns
577    /// `Error::Done`.
578    pub fn try_fill_buffer<F: BufFactory>(
579        &mut self, conn: &mut crate::Connection<F>,
580    ) -> Result<()> {
581        // If no bytes are required to be read, return early.
582        if self.state_buffer_complete() {
583            return Ok(());
584        }
585
586        loop {
587            let stream_id = self.id;
588
589            let spare_buf = self.spare_state_buf();
590            let spare_len = spare_buf.len();
591
592            match conn.stream_recv(stream_id, spare_buf) {
593                Ok((read, fin)) => {
594                    self.commit_state_buf_read(read);
595
596                    if self.critical_stream_closed(fin) {
597                        super::close_conn_critical_stream(conn)?;
598                    }
599
600                    trace!(
601                        "{} read {} bytes on stream {}",
602                        conn.trace_id(),
603                        read,
604                        self.id,
605                    );
606
607                    if read < spare_len {
608                        break;
609                    }
610
611                    if self.state_buffer_complete() {
612                        return Ok(());
613                    }
614                },
615
616                Err(e @ crate::Error::StreamReset(_)) => {
617                    if self.critical_stream_closed(true) {
618                        super::close_conn_critical_stream(conn)?;
619                    }
620
621                    return Err(e.into());
622                },
623
624                Err(e) => {
625                    // The stream is not readable anymore, so re-arm the Data
626                    // event.
627                    if e == crate::Error::Done {
628                        self.reset_data_event();
629                    }
630
631                    return Err(e.into());
632                },
633            };
634        }
635
636        if !self.state_buffer_complete() {
637            self.reset_data_event();
638
639            return Err(Error::Done);
640        }
641
642        Ok(())
643    }
644
645    /// Tries to read data from the corresponding transport stream up to the
646    /// state's size, without storing the data in the state buffer.
647    ///
648    /// When not enough data can be read to complete the state, this returns
649    /// `Error::Done`.
650    pub fn try_skip_data<F: BufFactory>(
651        &mut self, conn: &mut crate::Connection<F>,
652    ) -> Result<()> {
653        // If no bytes are required to be read, return early.
654        if self.state_buffer_complete() {
655            return Ok(());
656        }
657
658        let len = self.state_len - self.state_off;
659
660        let read = match conn.stream_discard(self.id, len) {
661            Ok((len, fin)) => {
662                if self.critical_stream_closed(fin) {
663                    super::close_conn_critical_stream(conn)?;
664                }
665
666                len
667            },
668
669            Err(e @ crate::Error::StreamReset(_)) => {
670                if self.critical_stream_closed(true) {
671                    super::close_conn_critical_stream(conn)?;
672                }
673
674                return Err(e.into());
675            },
676
677            Err(e) => {
678                // The stream is not readable anymore, so re-arm the Data
679                // event.
680                if e == crate::Error::Done {
681                    self.reset_data_event();
682                }
683
684                return Err(e.into());
685            },
686        };
687
688        trace!(
689            "{} discarded {} bytes on stream {}",
690            conn.trace_id(),
691            read,
692            self.id,
693        );
694
695        self.state_off += read;
696
697        if !self.state_buffer_complete() {
698            self.reset_data_event();
699
700            return Err(Error::Done);
701        }
702
703        Ok(())
704    }
705
706    /// Initialize the local part of the stream.
707    pub fn initialize_local(&mut self) {
708        self.local_initialized = true
709    }
710
711    /// Whether the stream has been locally initialized.
712    pub fn local_initialized(&self) -> bool {
713        self.local_initialized
714    }
715
716    /// Finish the local part of the stream.
717    pub fn finish_local(&mut self) {
718        self.local_finished = true
719    }
720
721    /// Whether the local send-side of the stream has finished.
722    pub fn local_finished(&self) -> bool {
723        self.local_finished
724    }
725
726    pub fn increment_headers_received(&mut self) {
727        self.headers_received_count =
728            self.headers_received_count.saturating_add(1);
729    }
730
731    pub fn headers_received_count(&self) -> usize {
732        self.headers_received_count
733    }
734
735    pub fn mark_trailers_sent(&mut self) {
736        self.trailers_sent = true;
737    }
738
739    pub fn trailers_sent(&self) -> bool {
740        self.trailers_sent
741    }
742
743    /// Tries to fill the state buffer by reading data from the given cursor.
744    ///
745    /// This is intended to replace `try_fill_buffer()` in tests, in order to
746    /// avoid having to setup a transport connection.
747    #[cfg(test)]
748    fn try_fill_buffer_for_tests(
749        &mut self, stream: &mut std::io::Cursor<Vec<u8>>,
750    ) -> Result<()> {
751        // If no bytes are required to be read, return early
752        if self.state_buffer_complete() {
753            return Ok(());
754        }
755
756        loop {
757            let spare_buf = self.spare_state_buf();
758            let spare_len = spare_buf.len();
759
760            let read = match std::io::Read::read(stream, spare_buf) {
761                Ok(0) => {
762                    // end of stream, stop
763                    break;
764                },
765
766                Ok(v) => v,
767
768                Err(_) => {
769                    panic!("Test buffer reading should never fail");
770                },
771            };
772
773            self.commit_state_buf_read(read);
774
775            if read < spare_len {
776                break;
777            }
778
779            if self.state_buffer_complete() {
780                break;
781            }
782        }
783
784        if !self.state_buffer_complete() {
785            return Err(Error::Done);
786        }
787
788        Ok(())
789    }
790
791    /// Tries to parse a varint (including length) from the state buffer.
792    pub fn try_consume_varint(&mut self) -> Result<u64> {
793        if self.state_off == 1 {
794            self.state_len = octets::varint_parse_len(self.state_buf[0]);
795            self.state_buf.reserve(self.state_len);
796        }
797
798        // Return early if we don't have enough data in the state buffer to
799        // parse the whole varint.
800        if !self.state_buffer_complete() {
801            return Err(Error::Done);
802        }
803
804        let varint = octets::Octets::with_slice(&self.state_buf).get_varint()?;
805
806        Ok(varint)
807    }
808
809    /// Tries to parse a frame from the state buffer.
810    ///
811    /// If successful, returns the `frame::Frame` and the payload length.
812    pub fn try_consume_frame(&mut self) -> Result<(frame::Frame, u64)> {
813        debug_assert_eq!(self.state, State::FramePayload);
814        // Processing a frame other than DATA, so re-arm the Data event.
815        self.reset_data_event();
816
817        let payload_len = self.state_len as u64;
818
819        // TODO: properly propagate frame parsing errors.
820        let frame = frame::Frame::from_bytes(
821            self.frame_type.unwrap(),
822            payload_len,
823            &self.state_buf,
824        )?;
825
826        self.state_transition(State::FrameType, 1, true)?;
827
828        Ok((frame, payload_len))
829    }
830
831    /// Tries to skip the current frame's payload.
832    pub fn try_skip_frame<F: BufFactory>(
833        &mut self, conn: &mut crate::Connection<F>,
834    ) -> Result<()> {
835        self.try_skip_data(conn)?;
836
837        // Processing a frame other than DATA, so re-arm the Data event.
838        self.reset_data_event();
839
840        self.state_transition(State::FrameType, 1, true)?;
841
842        Ok(())
843    }
844
845    /// Tries to read DATA payload from the transport stream.
846    pub fn try_consume_data<F: BufFactory, OUT: bytes::BufMut>(
847        &mut self, conn: &mut crate::Connection<F>, out: OUT,
848    ) -> Result<(usize, bool)> {
849        debug_assert_eq!(self.state, State::Data);
850        let out = out.limit(self.state_len - self.state_off);
851
852        let (len, fin) = match conn.stream_recv_buf(self.id, out) {
853            Ok(v) => v,
854
855            Err(e) => {
856                // The stream is not readable anymore, so re-arm the Data event.
857                if e == crate::Error::Done {
858                    self.reset_data_event();
859                }
860
861                return Err(e.into());
862            },
863        };
864
865        self.state_off += len;
866        debug_assert!(self.state_len >= self.state_off);
867
868        // The stream is not readable anymore, so re-arm the Data event.
869        if !conn.stream_readable(self.id) {
870            self.reset_data_event();
871        }
872
873        if self.state_buffer_complete() {
874            self.state_transition(State::FrameType, 1, true)?;
875        }
876
877        Ok((len, fin))
878    }
879
880    /// Marks the stream as finished.
881    pub fn finished(&mut self) {
882        let _ = self.state_transition(State::Finished, 0, false);
883    }
884
885    /// Tries to read DATA payload from the given cursor.
886    ///
887    /// This is intended to replace `try_consume_data()` in tests, in order to
888    /// avoid having to setup a transport connection.
889    #[cfg(test)]
890    fn try_consume_data_for_tests(
891        &mut self, stream: &mut std::io::Cursor<Vec<u8>>, out: &mut [u8],
892    ) -> Result<usize> {
893        let left = std::cmp::min(out.len(), self.state_len - self.state_off);
894
895        let len = std::io::Read::read(stream, &mut out[..left]).unwrap();
896
897        self.state_off += len;
898
899        if self.state_buffer_complete() {
900            self.state_transition(State::FrameType, 1, true)?;
901        }
902
903        Ok(len)
904    }
905
906    /// Tries to update the data triggered state for the stream.
907    ///
908    /// This returns `true` if a Data event was not already triggered before
909    /// the last reset, and updates the state. Returns `false` otherwise.
910    pub fn try_trigger_data_event(&mut self) -> bool {
911        if self.data_event_triggered {
912            return false;
913        }
914
915        self.data_event_triggered = true;
916
917        true
918    }
919
920    /// Resets the data triggered state.
921    fn reset_data_event(&mut self) {
922        self.data_event_triggered = false;
923    }
924
925    /// Set the last priority update for the stream.
926    pub fn set_last_priority_update(&mut self, priority_update: Option<Vec<u8>>) {
927        self.last_priority_update = priority_update;
928    }
929
930    /// Take the last priority update and leave `None` in its place.
931    pub fn take_last_priority_update(&mut self) -> Option<Vec<u8>> {
932        self.last_priority_update.take()
933    }
934
935    /// Returns `true` if there is a priority update.
936    pub fn has_last_priority_update(&self) -> bool {
937        self.last_priority_update.is_some()
938    }
939
940    /// Checks whether one of the critical streams was closed.
941    fn critical_stream_closed(&self, fin: bool) -> bool {
942        fin && matches!(
943            self.ty,
944            Some(Type::Control) |
945                Some(Type::QpackEncoder) |
946                Some(Type::QpackDecoder)
947        )
948    }
949
950    /// Returns true if the state buffer has enough data to complete the state.
951    fn state_buffer_complete(&self) -> bool {
952        self.state_off == self.state_len
953    }
954
955    /// Transitions the stream to a new state, and optionally resets the state
956    /// buffer.
957    fn state_transition(
958        &mut self, new_state: State, expected_len: usize, resize: bool,
959    ) -> Result<()> {
960        self.state_buf.clear();
961
962        // Some states don't need the state buffer, so don't resize it if not
963        // necessary.
964        if resize {
965            // A peer can influence the size of the state buffer (e.g. with the
966            // payload size of a GREASE frame), so we need to limit the maximum
967            // size to avoid DoS.
968            if expected_len > MAX_STATE_BUF_SIZE {
969                return Err(Error::ExcessiveLoad);
970            }
971
972            let reserve_len =
973                std::cmp::min(expected_len, MAX_STATE_BUF_ALLOC_SIZE);
974            self.state_buf.reserve(reserve_len);
975        }
976
977        self.state = new_state;
978        self.state_off = 0;
979        self.state_len = expected_len;
980
981        Ok(())
982    }
983}
984
985#[cfg(test)]
986mod tests {
987    use crate::h3::frame::*;
988    use crate::h3::PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT;
989    use crate::h3::SETTINGS_MAX_FIELD_SECTION_SIZE_DEFAULT;
990
991    use super::*;
992
993    fn open_uni(b: &mut octets::OctetsMut, ty: u64) -> Result<Stream> {
994        let stream = <Stream>::new(
995            2,
996            false,
997            SETTINGS_MAX_FIELD_SECTION_SIZE_DEFAULT,
998            PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT,
999        );
1000        assert_eq!(stream.state, State::StreamType);
1001
1002        b.put_varint(ty)?;
1003
1004        Ok(stream)
1005    }
1006
1007    fn open_remote_request_stream() -> Stream {
1008        Stream::new(
1009            0,
1010            false,
1011            SETTINGS_MAX_FIELD_SECTION_SIZE_DEFAULT,
1012            PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT,
1013        )
1014    }
1015
1016    fn parse_uni(
1017        stream: &mut Stream, ty: u64, cursor: &mut std::io::Cursor<Vec<u8>>,
1018    ) -> Result<()> {
1019        stream.try_fill_buffer_for_tests(cursor)?;
1020
1021        let stream_ty = stream.try_consume_varint()?;
1022        assert_eq!(stream_ty, ty);
1023        stream.set_ty(Type::deserialize(stream_ty).unwrap())?;
1024
1025        Ok(())
1026    }
1027
1028    /// Fill the buffer and parse a multi-byte varint that requires two
1029    /// fills (the first read gets the length prefix, the second gets the
1030    /// remaining bytes).
1031    fn parse_multibyte_varint(
1032        stream: &mut Stream, cursor: &mut std::io::Cursor<Vec<u8>>,
1033    ) -> Result<u64> {
1034        stream.try_fill_buffer_for_tests(cursor)?;
1035        assert_eq!(stream.try_consume_varint(), Err(Error::Done));
1036        stream.try_fill_buffer_for_tests(cursor)?;
1037        stream.try_consume_varint()
1038    }
1039
1040    fn parse_skip_frame(
1041        stream: &mut Stream, cursor: &mut std::io::Cursor<Vec<u8>>,
1042    ) -> Result<()> {
1043        // Parse the frame type.
1044        stream.try_fill_buffer_for_tests(cursor)?;
1045
1046        let frame_ty = stream.try_consume_varint()?;
1047
1048        stream.set_frame_type(frame_ty)?;
1049        assert_eq!(stream.state, State::FramePayloadLen);
1050
1051        // Parse the frame payload length.
1052        stream.try_fill_buffer_for_tests(cursor)?;
1053
1054        let frame_payload_len = stream.try_consume_varint()?;
1055        stream.set_frame_payload_len(frame_payload_len)?;
1056        assert_eq!(stream.state, State::FramePayload);
1057
1058        // Parse the frame payload.
1059        stream.try_fill_buffer_for_tests(cursor)?;
1060
1061        stream.try_consume_frame()?;
1062        assert_eq!(stream.state, State::FrameType);
1063
1064        Ok(())
1065    }
1066
1067    #[test]
1068    /// Process incoming SETTINGS frame on control stream.
1069    fn control_good() {
1070        let mut d = vec![42; 40];
1071        let mut b = octets::OctetsMut::with_slice(&mut d);
1072
1073        let raw_settings = vec![
1074            (SETTINGS_MAX_FIELD_SECTION_SIZE, 0),
1075            (SETTINGS_QPACK_MAX_TABLE_CAPACITY, 0),
1076            (SETTINGS_QPACK_BLOCKED_STREAMS, 0),
1077        ];
1078
1079        let frame = Frame::Settings {
1080            max_field_section_size: Some(0),
1081            qpack_max_table_capacity: Some(0),
1082            qpack_blocked_streams: Some(0),
1083            connect_protocol_enabled: None,
1084            h3_datagram: None,
1085            grease: None,
1086            additional_settings: None,
1087            raw: Some(raw_settings),
1088        };
1089
1090        let mut stream = open_uni(&mut b, HTTP3_CONTROL_STREAM_TYPE_ID).unwrap();
1091        frame.to_bytes(&mut b).unwrap();
1092
1093        let mut cursor = std::io::Cursor::new(d);
1094
1095        parse_uni(&mut stream, HTTP3_CONTROL_STREAM_TYPE_ID, &mut cursor)
1096            .unwrap();
1097        assert_eq!(stream.state, State::FrameType);
1098
1099        // Parse the SETTINGS frame type.
1100        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1101
1102        let frame_ty = stream.try_consume_varint().unwrap();
1103        assert_eq!(frame_ty, SETTINGS_FRAME_TYPE_ID);
1104
1105        stream.set_frame_type(frame_ty).unwrap();
1106        assert_eq!(stream.state, State::FramePayloadLen);
1107
1108        // Parse the SETTINGS frame payload length.
1109        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1110
1111        let frame_payload_len = stream.try_consume_varint().unwrap();
1112        assert_eq!(frame_payload_len, 6);
1113        stream.set_frame_payload_len(frame_payload_len).unwrap();
1114        assert_eq!(stream.state, State::FramePayload);
1115
1116        // Parse the SETTINGS frame payload.
1117        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1118
1119        assert_eq!(stream.try_consume_frame(), Ok((frame, 6)));
1120        assert_eq!(stream.state, State::FrameType);
1121    }
1122
1123    #[test]
1124    /// Process incoming empty SETTINGS frame on control stream.
1125    fn control_empty_settings() {
1126        let mut d = vec![42; 40];
1127        let mut b = octets::OctetsMut::with_slice(&mut d);
1128
1129        let frame = Frame::Settings {
1130            max_field_section_size: None,
1131            qpack_max_table_capacity: None,
1132            qpack_blocked_streams: None,
1133            connect_protocol_enabled: None,
1134            h3_datagram: None,
1135            grease: None,
1136            additional_settings: None,
1137            raw: Some(vec![]),
1138        };
1139
1140        let mut stream = open_uni(&mut b, HTTP3_CONTROL_STREAM_TYPE_ID).unwrap();
1141        frame.to_bytes(&mut b).unwrap();
1142
1143        let mut cursor = std::io::Cursor::new(d);
1144
1145        parse_uni(&mut stream, HTTP3_CONTROL_STREAM_TYPE_ID, &mut cursor)
1146            .unwrap();
1147        assert_eq!(stream.state, State::FrameType);
1148
1149        // Parse the SETTINGS frame type.
1150        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1151
1152        let frame_ty = stream.try_consume_varint().unwrap();
1153        assert_eq!(frame_ty, SETTINGS_FRAME_TYPE_ID);
1154
1155        stream.set_frame_type(frame_ty).unwrap();
1156        assert_eq!(stream.state, State::FramePayloadLen);
1157
1158        // Parse the SETTINGS frame payload length.
1159        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1160
1161        let frame_payload_len = stream.try_consume_varint().unwrap();
1162        assert_eq!(frame_payload_len, 0);
1163        stream.set_frame_payload_len(frame_payload_len).unwrap();
1164        assert_eq!(stream.state, State::FramePayload);
1165
1166        // Parse the SETTINGS frame payload.
1167        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1168
1169        assert_eq!(stream.try_consume_frame(), Ok((frame, 0)));
1170        assert_eq!(stream.state, State::FrameType);
1171    }
1172
1173    #[test]
1174    /// Process duplicate SETTINGS frame on control stream.
1175    fn control_bad_multiple_settings() {
1176        let mut d = vec![42; 40];
1177        let mut b = octets::OctetsMut::with_slice(&mut d);
1178
1179        let raw_settings = vec![
1180            (SETTINGS_MAX_FIELD_SECTION_SIZE, 0),
1181            (SETTINGS_QPACK_MAX_TABLE_CAPACITY, 0),
1182            (SETTINGS_QPACK_BLOCKED_STREAMS, 0),
1183        ];
1184
1185        let frame = Frame::Settings {
1186            max_field_section_size: Some(0),
1187            qpack_max_table_capacity: Some(0),
1188            qpack_blocked_streams: Some(0),
1189            connect_protocol_enabled: None,
1190            h3_datagram: None,
1191            grease: None,
1192            additional_settings: None,
1193            raw: Some(raw_settings),
1194        };
1195
1196        let mut stream = open_uni(&mut b, HTTP3_CONTROL_STREAM_TYPE_ID).unwrap();
1197        frame.to_bytes(&mut b).unwrap();
1198        frame.to_bytes(&mut b).unwrap();
1199
1200        let mut cursor = std::io::Cursor::new(d);
1201
1202        parse_uni(&mut stream, HTTP3_CONTROL_STREAM_TYPE_ID, &mut cursor)
1203            .unwrap();
1204        assert_eq!(stream.state, State::FrameType);
1205
1206        // Parse the SETTINGS frame type.
1207        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1208
1209        let frame_ty = stream.try_consume_varint().unwrap();
1210        assert_eq!(frame_ty, SETTINGS_FRAME_TYPE_ID);
1211
1212        stream.set_frame_type(frame_ty).unwrap();
1213        assert_eq!(stream.state, State::FramePayloadLen);
1214
1215        // Parse the SETTINGS frame payload length.
1216        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1217
1218        let frame_payload_len = stream.try_consume_varint().unwrap();
1219        assert_eq!(frame_payload_len, 6);
1220        stream.set_frame_payload_len(frame_payload_len).unwrap();
1221        assert_eq!(stream.state, State::FramePayload);
1222
1223        // Parse the SETTINGS frame payload.
1224        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1225
1226        assert_eq!(stream.try_consume_frame(), Ok((frame, 6)));
1227        assert_eq!(stream.state, State::FrameType);
1228
1229        // Parse the second SETTINGS frame type.
1230        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1231
1232        let frame_ty = stream.try_consume_varint().unwrap();
1233        assert_eq!(stream.set_frame_type(frame_ty), Err(Error::FrameUnexpected));
1234    }
1235
1236    #[test]
1237    /// Process other frame before SETTINGS frame on control stream.
1238    fn control_bad_late_settings() {
1239        let mut d = vec![42; 40];
1240        let mut b = octets::OctetsMut::with_slice(&mut d);
1241
1242        let goaway = Frame::GoAway { id: 0 };
1243
1244        let raw_settings = vec![
1245            (SETTINGS_MAX_FIELD_SECTION_SIZE, 0),
1246            (SETTINGS_QPACK_MAX_TABLE_CAPACITY, 0),
1247            (SETTINGS_QPACK_BLOCKED_STREAMS, 0),
1248        ];
1249
1250        let settings = Frame::Settings {
1251            max_field_section_size: Some(0),
1252            qpack_max_table_capacity: Some(0),
1253            qpack_blocked_streams: Some(0),
1254            connect_protocol_enabled: None,
1255            h3_datagram: None,
1256            grease: None,
1257            additional_settings: None,
1258            raw: Some(raw_settings),
1259        };
1260
1261        let mut stream = open_uni(&mut b, HTTP3_CONTROL_STREAM_TYPE_ID).unwrap();
1262        goaway.to_bytes(&mut b).unwrap();
1263        settings.to_bytes(&mut b).unwrap();
1264
1265        let mut cursor = std::io::Cursor::new(d);
1266
1267        parse_uni(&mut stream, HTTP3_CONTROL_STREAM_TYPE_ID, &mut cursor)
1268            .unwrap();
1269        assert_eq!(stream.state, State::FrameType);
1270
1271        // Parse GOAWAY.
1272        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1273
1274        let frame_ty = stream.try_consume_varint().unwrap();
1275        assert_eq!(stream.set_frame_type(frame_ty), Err(Error::MissingSettings));
1276    }
1277
1278    #[test]
1279    /// Process not-allowed frame on control stream.
1280    fn control_bad_frame() {
1281        let mut d = vec![42; 40];
1282        let mut b = octets::OctetsMut::with_slice(&mut d);
1283
1284        let header_block = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12];
1285        let hdrs = Frame::Headers { header_block };
1286
1287        let raw_settings = vec![
1288            (SETTINGS_MAX_FIELD_SECTION_SIZE, 0),
1289            (SETTINGS_QPACK_MAX_TABLE_CAPACITY, 0),
1290            (SETTINGS_QPACK_BLOCKED_STREAMS, 0),
1291            (33, 33),
1292        ];
1293
1294        let settings = Frame::Settings {
1295            max_field_section_size: Some(0),
1296            qpack_max_table_capacity: Some(0),
1297            qpack_blocked_streams: Some(0),
1298            connect_protocol_enabled: None,
1299            h3_datagram: None,
1300            grease: None,
1301            additional_settings: None,
1302            raw: Some(raw_settings),
1303        };
1304
1305        let mut stream = open_uni(&mut b, HTTP3_CONTROL_STREAM_TYPE_ID).unwrap();
1306        settings.to_bytes(&mut b).unwrap();
1307        hdrs.to_bytes(&mut b).unwrap();
1308
1309        let mut cursor = std::io::Cursor::new(d);
1310
1311        parse_uni(&mut stream, HTTP3_CONTROL_STREAM_TYPE_ID, &mut cursor)
1312            .unwrap();
1313        assert_eq!(stream.state, State::FrameType);
1314
1315        // Parse first SETTINGS frame.
1316        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1317
1318        let frame_ty = stream.try_consume_varint().unwrap();
1319        stream.set_frame_type(frame_ty).unwrap();
1320
1321        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1322
1323        let frame_payload_len = stream.try_consume_varint().unwrap();
1324        stream.set_frame_payload_len(frame_payload_len).unwrap();
1325
1326        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1327
1328        assert!(stream.try_consume_frame().is_ok());
1329
1330        // Parse HEADERS.
1331        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1332
1333        let frame_ty = stream.try_consume_varint().unwrap();
1334        assert_eq!(stream.set_frame_type(frame_ty), Err(Error::FrameUnexpected));
1335    }
1336
1337    #[test]
1338    fn request_no_data() {
1339        let mut stream = open_remote_request_stream();
1340
1341        assert_eq!(stream.ty, Some(Type::Request));
1342        assert_eq!(stream.state, State::FrameType);
1343
1344        assert_eq!(stream.try_consume_varint(), Err(Error::Done));
1345    }
1346
1347    #[test]
1348    fn request_good() {
1349        let mut stream = open_remote_request_stream();
1350
1351        let mut d = vec![42; 128];
1352        let mut b = octets::OctetsMut::with_slice(&mut d);
1353
1354        let header_block = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12];
1355        let payload = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12];
1356        let hdrs = Frame::Headers { header_block };
1357        let data = Frame::Data {
1358            payload: payload.clone(),
1359        };
1360
1361        hdrs.to_bytes(&mut b).unwrap();
1362        data.to_bytes(&mut b).unwrap();
1363
1364        let mut cursor = std::io::Cursor::new(d);
1365
1366        // Parse the HEADERS frame type.
1367        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1368
1369        let frame_ty = stream.try_consume_varint().unwrap();
1370        assert_eq!(frame_ty, HEADERS_FRAME_TYPE_ID);
1371
1372        stream.set_frame_type(frame_ty).unwrap();
1373        assert_eq!(stream.state, State::FramePayloadLen);
1374
1375        // Parse the HEADERS frame payload length.
1376        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1377
1378        let frame_payload_len = stream.try_consume_varint().unwrap();
1379        assert_eq!(frame_payload_len, 12);
1380
1381        stream.set_frame_payload_len(frame_payload_len).unwrap();
1382        assert_eq!(stream.state, State::FramePayload);
1383
1384        // Parse the HEADERS frame.
1385        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1386
1387        assert_eq!(stream.try_consume_frame(), Ok((hdrs, 12)));
1388        assert_eq!(stream.state, State::FrameType);
1389
1390        // Parse the DATA frame type.
1391        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1392
1393        let frame_ty = stream.try_consume_varint().unwrap();
1394        assert_eq!(frame_ty, DATA_FRAME_TYPE_ID);
1395
1396        stream.set_frame_type(frame_ty).unwrap();
1397        assert_eq!(stream.state, State::FramePayloadLen);
1398
1399        // Parse the DATA frame payload length.
1400        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1401
1402        let frame_payload_len = stream.try_consume_varint().unwrap();
1403        assert_eq!(frame_payload_len, 12);
1404
1405        stream.set_frame_payload_len(frame_payload_len).unwrap();
1406        assert_eq!(stream.state, State::Data);
1407
1408        // Parse the DATA payload.
1409        let mut recv_buf = vec![0; payload.len()];
1410        assert_eq!(
1411            stream.try_consume_data_for_tests(&mut cursor, &mut recv_buf),
1412            Ok(payload.len())
1413        );
1414        assert_eq!(payload, recv_buf);
1415
1416        assert_eq!(stream.state, State::FrameType);
1417    }
1418
1419    #[test]
1420    fn push_good() {
1421        let mut d = vec![42; 128];
1422        let mut b = octets::OctetsMut::with_slice(&mut d);
1423
1424        let header_block = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12];
1425        let payload = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12];
1426        let hdrs = Frame::Headers { header_block };
1427        let data = Frame::Data {
1428            payload: payload.clone(),
1429        };
1430
1431        let mut stream = open_uni(&mut b, HTTP3_PUSH_STREAM_TYPE_ID).unwrap();
1432        b.put_varint(1).unwrap();
1433        hdrs.to_bytes(&mut b).unwrap();
1434        data.to_bytes(&mut b).unwrap();
1435
1436        let mut cursor = std::io::Cursor::new(d);
1437
1438        parse_uni(&mut stream, HTTP3_PUSH_STREAM_TYPE_ID, &mut cursor).unwrap();
1439        assert_eq!(stream.state, State::PushId);
1440
1441        // Parse push ID.
1442        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1443
1444        let push_id = stream.try_consume_varint().unwrap();
1445        assert_eq!(push_id, 1);
1446
1447        stream.set_push_id(push_id).unwrap();
1448        assert_eq!(stream.state, State::FrameType);
1449
1450        // Parse the HEADERS frame type.
1451        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1452
1453        let frame_ty = stream.try_consume_varint().unwrap();
1454        assert_eq!(frame_ty, HEADERS_FRAME_TYPE_ID);
1455
1456        stream.set_frame_type(frame_ty).unwrap();
1457        assert_eq!(stream.state, State::FramePayloadLen);
1458
1459        // Parse the HEADERS frame payload length.
1460        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1461
1462        let frame_payload_len = stream.try_consume_varint().unwrap();
1463        assert_eq!(frame_payload_len, 12);
1464
1465        stream.set_frame_payload_len(frame_payload_len).unwrap();
1466        assert_eq!(stream.state, State::FramePayload);
1467
1468        // Parse the HEADERS frame.
1469        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1470
1471        assert_eq!(stream.try_consume_frame(), Ok((hdrs, 12)));
1472        assert_eq!(stream.state, State::FrameType);
1473
1474        // Parse the DATA frame type.
1475        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1476
1477        let frame_ty = stream.try_consume_varint().unwrap();
1478        assert_eq!(frame_ty, DATA_FRAME_TYPE_ID);
1479
1480        stream.set_frame_type(frame_ty).unwrap();
1481        assert_eq!(stream.state, State::FramePayloadLen);
1482
1483        // Parse the DATA frame payload length.
1484        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1485
1486        let frame_payload_len = stream.try_consume_varint().unwrap();
1487        assert_eq!(frame_payload_len, 12);
1488
1489        stream.set_frame_payload_len(frame_payload_len).unwrap();
1490        assert_eq!(stream.state, State::Data);
1491
1492        // Parse the DATA payload.
1493        let mut recv_buf = vec![0; payload.len()];
1494        assert_eq!(
1495            stream.try_consume_data_for_tests(&mut cursor, &mut recv_buf),
1496            Ok(payload.len())
1497        );
1498        assert_eq!(payload, recv_buf);
1499
1500        assert_eq!(stream.state, State::FrameType);
1501    }
1502
1503    #[test]
1504    fn grease() {
1505        let mut d = vec![42; 20];
1506        let mut b = octets::OctetsMut::with_slice(&mut d);
1507
1508        let mut stream = open_uni(&mut b, 33).unwrap();
1509
1510        let mut cursor = std::io::Cursor::new(d);
1511
1512        // Parse stream type.
1513        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1514
1515        let stream_ty = stream.try_consume_varint().unwrap();
1516        assert_eq!(stream_ty, 33);
1517        stream
1518            .set_ty(Type::deserialize(stream_ty).unwrap())
1519            .unwrap();
1520        assert_eq!(stream.state, State::Drain);
1521    }
1522
1523    #[test]
1524    fn data_before_headers() {
1525        let mut stream = open_remote_request_stream();
1526
1527        let mut d = vec![42; 128];
1528        let mut b = octets::OctetsMut::with_slice(&mut d);
1529
1530        let data = Frame::Data {
1531            payload: vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12],
1532        };
1533
1534        data.to_bytes(&mut b).unwrap();
1535
1536        let mut cursor = std::io::Cursor::new(d);
1537
1538        // Parse the DATA frame type.
1539        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1540
1541        let frame_ty = stream.try_consume_varint().unwrap();
1542        assert_eq!(frame_ty, DATA_FRAME_TYPE_ID);
1543
1544        assert_eq!(stream.set_frame_type(frame_ty), Err(Error::FrameUnexpected));
1545    }
1546
1547    #[test]
1548    fn additional_headers() {
1549        let mut stream = open_remote_request_stream();
1550
1551        let mut d = vec![42; 128];
1552        let mut b = octets::OctetsMut::with_slice(&mut d);
1553
1554        let header_block = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12];
1555        let payload = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12];
1556        let info_hdrs = Frame::Headers {
1557            header_block: header_block.clone(),
1558        };
1559        let non_info_hdrs = Frame::Headers {
1560            header_block: header_block.clone(),
1561        };
1562        let trailers = Frame::Headers { header_block };
1563        let data = Frame::Data {
1564            payload: payload.clone(),
1565        };
1566
1567        info_hdrs.to_bytes(&mut b).unwrap();
1568        non_info_hdrs.to_bytes(&mut b).unwrap();
1569        data.to_bytes(&mut b).unwrap();
1570        trailers.to_bytes(&mut b).unwrap();
1571
1572        let mut cursor = std::io::Cursor::new(d);
1573
1574        // Parse the HEADERS frame type.
1575        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1576
1577        let frame_ty = stream.try_consume_varint().unwrap();
1578        assert_eq!(frame_ty, HEADERS_FRAME_TYPE_ID);
1579
1580        stream.set_frame_type(frame_ty).unwrap();
1581        assert_eq!(stream.state, State::FramePayloadLen);
1582
1583        // Parse the HEADERS frame payload length.
1584        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1585
1586        let frame_payload_len = stream.try_consume_varint().unwrap();
1587        assert_eq!(frame_payload_len, 12);
1588
1589        stream.set_frame_payload_len(frame_payload_len).unwrap();
1590        assert_eq!(stream.state, State::FramePayload);
1591
1592        // Parse the HEADERS frame.
1593        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1594
1595        assert_eq!(stream.try_consume_frame(), Ok((info_hdrs, 12)));
1596        assert_eq!(stream.state, State::FrameType);
1597
1598        // Parse the non-info HEADERS frame type.
1599        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1600
1601        let frame_ty = stream.try_consume_varint().unwrap();
1602        assert_eq!(frame_ty, HEADERS_FRAME_TYPE_ID);
1603
1604        stream.set_frame_type(frame_ty).unwrap();
1605        assert_eq!(stream.state, State::FramePayloadLen);
1606
1607        // Parse the HEADERS frame payload length.
1608        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1609
1610        let frame_payload_len = stream.try_consume_varint().unwrap();
1611        assert_eq!(frame_payload_len, 12);
1612
1613        stream.set_frame_payload_len(frame_payload_len).unwrap();
1614        assert_eq!(stream.state, State::FramePayload);
1615
1616        // Parse the HEADERS frame.
1617        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1618
1619        assert_eq!(stream.try_consume_frame(), Ok((non_info_hdrs, 12)));
1620        assert_eq!(stream.state, State::FrameType);
1621
1622        // Parse the DATA frame type.
1623        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1624
1625        let frame_ty = stream.try_consume_varint().unwrap();
1626        assert_eq!(frame_ty, DATA_FRAME_TYPE_ID);
1627
1628        stream.set_frame_type(frame_ty).unwrap();
1629        assert_eq!(stream.state, State::FramePayloadLen);
1630
1631        // Parse the DATA frame payload length.
1632        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1633
1634        let frame_payload_len = stream.try_consume_varint().unwrap();
1635        assert_eq!(frame_payload_len, 12);
1636
1637        stream.set_frame_payload_len(frame_payload_len).unwrap();
1638        assert_eq!(stream.state, State::Data);
1639
1640        // Parse the DATA payload.
1641        let mut recv_buf = vec![0; payload.len()];
1642        assert_eq!(
1643            stream.try_consume_data_for_tests(&mut cursor, &mut recv_buf),
1644            Ok(payload.len())
1645        );
1646        assert_eq!(payload, recv_buf);
1647
1648        assert_eq!(stream.state, State::FrameType);
1649
1650        // Parse the trailing HEADERS frame type.
1651        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1652
1653        let frame_ty = stream.try_consume_varint().unwrap();
1654        assert_eq!(frame_ty, HEADERS_FRAME_TYPE_ID);
1655
1656        stream.set_frame_type(frame_ty).unwrap();
1657        assert_eq!(stream.state, State::FramePayloadLen);
1658
1659        // Parse the HEADERS frame payload length.
1660        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1661
1662        let frame_payload_len = stream.try_consume_varint().unwrap();
1663        assert_eq!(frame_payload_len, 12);
1664
1665        stream.set_frame_payload_len(frame_payload_len).unwrap();
1666        assert_eq!(stream.state, State::FramePayload);
1667
1668        // Parse the HEADERS frame.
1669        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1670
1671        assert_eq!(stream.try_consume_frame(), Ok((trailers, 12)));
1672        assert_eq!(stream.state, State::FrameType);
1673    }
1674
1675    /// Returns the frame type ID for a given frame variant.
1676    fn frame_type_id(frame: &Frame) -> u64 {
1677        match frame {
1678            Frame::Data { .. } => DATA_FRAME_TYPE_ID,
1679            Frame::Headers { .. } => HEADERS_FRAME_TYPE_ID,
1680            Frame::CancelPush { .. } => CANCEL_PUSH_FRAME_TYPE_ID,
1681            Frame::Settings { .. } => SETTINGS_FRAME_TYPE_ID,
1682            Frame::PushPromise { .. } => PUSH_PROMISE_FRAME_TYPE_ID,
1683            Frame::GoAway { .. } => GOAWAY_FRAME_TYPE_ID,
1684            Frame::MaxPushId { .. } => MAX_PUSH_FRAME_TYPE_ID,
1685            Frame::PriorityUpdateRequest { .. } =>
1686                PRIORITY_UPDATE_FRAME_REQUEST_TYPE_ID,
1687            Frame::PriorityUpdatePush { .. } =>
1688                PRIORITY_UPDATE_FRAME_PUSH_TYPE_ID,
1689            Frame::Unknown { .. } => unreachable!(),
1690        }
1691    }
1692
1693    /// Parse a large frame and check the size limit behavior.
1694    ///
1695    /// Writes the frame to a buffer, parses the frame type and payload
1696    /// length (handling multi-byte varint retry), then checks whether
1697    /// `set_frame_payload_len` accepts or rejects the frame. If accepted,
1698    /// also parses the payload and verifies `try_consume_frame` output.
1699    ///
1700    /// - `stream`: the H3 stream to parse the frame on.
1701    /// - `frame`: the frame to encode and parse back.
1702    /// - `expected_payload_len`: the expected encoded payload length.
1703    /// - `expect_accept`: if true, the frame should be accepted and fully
1704    ///   parsed; if false, `set_frame_payload_len` should reject it with
1705    ///   `Error::ExcessiveLoad`.
1706    fn check_large_frame_size_limit(
1707        stream: &mut Stream, frame: Frame, expected_payload_len: u64,
1708        expect_accept: bool,
1709    ) {
1710        let expected_type_id = frame_type_id(&frame);
1711
1712        let mut d = vec![42; 20000];
1713        let mut b = octets::OctetsMut::with_slice(&mut d);
1714        frame.to_bytes(&mut b).unwrap();
1715        let mut cursor = std::io::Cursor::new(d);
1716
1717        // Parse frame type.
1718        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1719        let frame_ty = stream.try_consume_varint().unwrap();
1720        assert_eq!(frame_ty, expected_type_id);
1721
1722        stream.set_frame_type(frame_ty).unwrap();
1723        assert_eq!(stream.state, State::FramePayloadLen);
1724
1725        // Parse frame payload length.
1726        let frame_payload_len =
1727            parse_multibyte_varint(stream, &mut cursor).unwrap();
1728        assert_eq!(frame_payload_len, expected_payload_len);
1729
1730        if expect_accept {
1731            stream.set_frame_payload_len(frame_payload_len).unwrap();
1732            assert_eq!(stream.state, State::FramePayload);
1733
1734            stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1735
1736            assert_eq!(
1737                stream.try_consume_frame(),
1738                Ok((frame, expected_payload_len))
1739            );
1740            assert_eq!(stream.state, State::FrameType);
1741        } else {
1742            assert_eq!(
1743                stream.set_frame_payload_len(frame_payload_len),
1744                Err(Error::ExcessiveLoad)
1745            );
1746        }
1747    }
1748
1749    #[test]
1750    fn large_headers_default_limit() {
1751        let mut stream = open_remote_request_stream();
1752        let header_block = vec![0; 16384];
1753        let frame = Frame::Headers {
1754            header_block: header_block.clone(),
1755        };
1756
1757        check_large_frame_size_limit(&mut stream, frame, 16384, true);
1758    }
1759
1760    #[test]
1761    fn large_headers_limit_with_huffman() {
1762        // Create stream with a max_field_section_size limit of 4k.
1763        let mut stream = Stream::new(
1764            0,
1765            false,
1766            4196,
1767            PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT,
1768        );
1769
1770        // Size the header block so it will fit within the Huffman margin.
1771        // On this branch the margin is x + x/2 = 4196 + 2098 = 6294.
1772        let header_block = vec![0; 6294];
1773        let frame = Frame::Headers {
1774            header_block: header_block.clone(),
1775        };
1776
1777        check_large_frame_size_limit(&mut stream, frame, 6294, true);
1778    }
1779
1780    #[test]
1781    fn large_headers_small_limit() {
1782        // Create stream with a max_field_section_size limit of 4k.
1783        // Encoded headers at 16k are larger than the 4k limit.
1784        let mut stream = Stream::new(
1785            0,
1786            false,
1787            4196,
1788            PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT,
1789        );
1790        let header_block = vec![0; 16384];
1791        let frame = Frame::Headers {
1792            header_block: header_block.clone(),
1793        };
1794
1795        check_large_frame_size_limit(&mut stream, frame, 16384, false);
1796    }
1797
1798    #[test]
1799    fn large_push_promise_default_limit() {
1800        let mut stream = open_remote_request_stream();
1801        let header_block = vec![0; 16384];
1802        let frame = Frame::PushPromise {
1803            push_id: 0,
1804            header_block: header_block.clone(),
1805        };
1806
1807        check_large_frame_size_limit(&mut stream, frame, 1 + 16384, true);
1808    }
1809
1810    #[test]
1811    fn large_push_promise_limit_with_huffman() {
1812        // Create stream with a max_field_section_size limit of 4k.
1813        let mut stream = Stream::new(
1814            0,
1815            false,
1816            4196,
1817            PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT,
1818        );
1819
1820        // Size the header block so it will fit within the Huffman margin.
1821        // On this branch the margin is x + x/2 = 4196 + 2098 = 6294.
1822        let header_block = vec![0; 6294];
1823        let frame = Frame::PushPromise {
1824            push_id: 0,
1825            header_block: header_block.clone(),
1826        };
1827
1828        check_large_frame_size_limit(&mut stream, frame, 1 + 6294, true);
1829    }
1830
1831    #[test]
1832    fn large_push_promise_small_limit() {
1833        // Create stream with a max_field_section_size limit of 4k.
1834        // Encoded push promise at 16k is larger than the 4k limit.
1835        let mut stream = Stream::new(
1836            0,
1837            false,
1838            4196,
1839            PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT,
1840        );
1841        let header_block = vec![0; 16384];
1842        let frame = Frame::PushPromise {
1843            push_id: 0,
1844            header_block: header_block.clone(),
1845        };
1846
1847        check_large_frame_size_limit(&mut stream, frame, 1 + 16384, false);
1848    }
1849
1850    #[test]
1851    fn large_priority_update_large_limit() {
1852        let settings = Frame::Settings {
1853            max_field_section_size: None,
1854            qpack_max_table_capacity: None,
1855            qpack_blocked_streams: None,
1856            connect_protocol_enabled: None,
1857            h3_datagram: None,
1858            grease: None,
1859            additional_settings: None,
1860            raw: Some(vec![]),
1861        };
1862
1863        let mut d = vec![42; 20000];
1864        let mut b = octets::OctetsMut::with_slice(&mut d);
1865
1866        // Control stream needs a SETTINGS frame to transition it into
1867        // being able to parse other frame types.
1868        let mut stream = <Stream>::new(
1869            2,
1870            false,
1871            SETTINGS_MAX_FIELD_SECTION_SIZE_DEFAULT,
1872            20000,
1873        );
1874        b.put_varint(HTTP3_CONTROL_STREAM_TYPE_ID).unwrap();
1875        settings.to_bytes(&mut b).unwrap();
1876
1877        let priority_field_value = vec![0; 16384];
1878        let pu = Frame::PriorityUpdateRequest {
1879            prioritized_element_id: 0,
1880            priority_field_value,
1881        };
1882
1883        pu.to_bytes(&mut b).unwrap();
1884
1885        let mut cursor = std::io::Cursor::new(d);
1886
1887        parse_uni(&mut stream, HTTP3_CONTROL_STREAM_TYPE_ID, &mut cursor)
1888            .unwrap();
1889
1890        // Skip SETTINGS frame type.
1891        parse_skip_frame(&mut stream, &mut cursor).unwrap();
1892
1893        // Parse the frame type.
1894        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1895
1896        // Parse fails because we need more bytes for the 4-byte encoded length.
1897        // This trial then sets the expected buffer size for us to fill.
1898        assert_eq!(stream.try_consume_varint(), Err(Error::Done));
1899        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1900        let frame_ty = stream.try_consume_varint().unwrap();
1901        assert_eq!(frame_ty, PRIORITY_UPDATE_FRAME_REQUEST_TYPE_ID);
1902
1903        stream.set_frame_type(frame_ty).unwrap();
1904        assert_eq!(stream.state, State::FramePayloadLen);
1905
1906        // Parse the frame payload length.
1907        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1908
1909        // Parse fails because we need more bytes for the 4-byte encoded length.
1910        // This trial then sets the expected buffer size for us to fill.
1911        assert_eq!(stream.try_consume_varint(), Err(Error::Done));
1912        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1913
1914        let frame_payload_len = stream.try_consume_varint().unwrap();
1915        assert_eq!(frame_payload_len, 1 + 16384);
1916
1917        stream.set_frame_payload_len(frame_payload_len).unwrap();
1918        assert_eq!(stream.state, State::FramePayload);
1919
1920        // Parse the frame.
1921        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1922
1923        assert_eq!(stream.try_consume_frame(), Ok((pu, 1 + 16384)));
1924        assert_eq!(stream.state, State::FrameType);
1925    }
1926
1927    #[test]
1928    fn large_priority_update_small_limit() {
1929        let settings = Frame::Settings {
1930            max_field_section_size: None,
1931            qpack_max_table_capacity: None,
1932            qpack_blocked_streams: None,
1933            connect_protocol_enabled: None,
1934            h3_datagram: None,
1935            grease: None,
1936            additional_settings: None,
1937            raw: Some(vec![]),
1938        };
1939
1940        let mut d = vec![42; 20000];
1941        let mut b = octets::OctetsMut::with_slice(&mut d);
1942
1943        // Control stream needs a SETTINGS frame to transition it into
1944        // being able to parse other frame types.
1945        let mut stream =
1946            <Stream>::new(2, false, SETTINGS_MAX_FIELD_SECTION_SIZE_DEFAULT, 123);
1947        b.put_varint(HTTP3_CONTROL_STREAM_TYPE_ID).unwrap();
1948
1949        settings.to_bytes(&mut b).unwrap();
1950
1951        let priority_field_value = vec![0; 16384];
1952        let pu = Frame::PriorityUpdateRequest {
1953            prioritized_element_id: 0,
1954            priority_field_value,
1955        };
1956
1957        pu.to_bytes(&mut b).unwrap();
1958
1959        let mut cursor = std::io::Cursor::new(d);
1960
1961        parse_uni(&mut stream, HTTP3_CONTROL_STREAM_TYPE_ID, &mut cursor)
1962            .unwrap();
1963
1964        // Skip SETTINGS frame type.
1965        parse_skip_frame(&mut stream, &mut cursor).unwrap();
1966
1967        // Parse the frame type.
1968        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1969
1970        // Parse fails because we need more bytes for the 4-byte encoded length.
1971        // This trial then sets the expected buffer size for us to fill.
1972        assert_eq!(stream.try_consume_varint(), Err(Error::Done));
1973        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1974        let frame_ty = stream.try_consume_varint().unwrap();
1975        assert_eq!(frame_ty, PRIORITY_UPDATE_FRAME_REQUEST_TYPE_ID);
1976
1977        stream.set_frame_type(frame_ty).unwrap();
1978        assert_eq!(stream.state, State::FramePayloadLen);
1979
1980        // Parse the frame payload length.
1981        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1982
1983        // Parse fails because we need more bytes for the 4-byte encoded length.
1984        // This trial then sets the expected buffer size for us to fill.
1985        assert_eq!(stream.try_consume_varint(), Err(Error::Done));
1986        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
1987
1988        let frame_payload_len = stream.try_consume_varint().unwrap();
1989        assert_eq!(frame_payload_len, 1 + 16384);
1990
1991        assert_eq!(
1992            stream.set_frame_payload_len(frame_payload_len),
1993            Err(Error::ExcessiveLoad)
1994        );
1995    }
1996
1997    #[test]
1998    fn finite_sized_frame_limits() {
1999        let settings = Frame::Settings {
2000            max_field_section_size: None,
2001            qpack_max_table_capacity: None,
2002            qpack_blocked_streams: None,
2003            connect_protocol_enabled: None,
2004            h3_datagram: None,
2005            grease: None,
2006            additional_settings: None,
2007            raw: Some(vec![]),
2008        };
2009
2010        for ty in [
2011            CANCEL_PUSH_FRAME_TYPE_ID,
2012            GOAWAY_FRAME_TYPE_ID,
2013            MAX_PUSH_FRAME_TYPE_ID,
2014        ] {
2015            // These frames must have a size between 1 and 8 bytes inclusive.
2016            for size in [0, 9] {
2017                let mut d = vec![42; 128];
2018                let mut b = octets::OctetsMut::with_slice(&mut d);
2019
2020                // Control stream needs a SETTINGS frame to transition it into
2021                // being able to parse other frame types.
2022                let mut stream =
2023                    open_uni(&mut b, HTTP3_CONTROL_STREAM_TYPE_ID).unwrap();
2024                settings.to_bytes(&mut b).unwrap();
2025
2026                // Write bytes as far as frame length.
2027                b.put_varint(ty).unwrap();
2028                b.put_varint(size).unwrap();
2029
2030                let mut cursor = std::io::Cursor::new(d);
2031
2032                parse_uni(&mut stream, HTTP3_CONTROL_STREAM_TYPE_ID, &mut cursor)
2033                    .unwrap();
2034
2035                // Skip SETTINGS frame type.
2036                parse_skip_frame(&mut stream, &mut cursor).unwrap();
2037
2038                // Parse frame type.
2039                stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
2040                let frame_ty = stream.try_consume_varint().unwrap();
2041                assert_eq!(frame_ty, ty);
2042
2043                stream.set_frame_type(frame_ty).unwrap();
2044                assert_eq!(stream.state, State::FramePayloadLen);
2045
2046                // Parse frame payload length.
2047                stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
2048                let frame_payload_len = stream.try_consume_varint().unwrap();
2049                assert_eq!(
2050                    Err(Error::FrameError),
2051                    stream.set_frame_payload_len(frame_payload_len)
2052                );
2053            }
2054        }
2055    }
2056
2057    #[test]
2058    fn zero_length_push_promise() {
2059        let mut d = vec![42; 128];
2060        let mut b = octets::OctetsMut::with_slice(&mut d);
2061
2062        let mut stream = open_remote_request_stream();
2063
2064        assert_eq!(stream.ty, Some(Type::Request));
2065        assert_eq!(stream.state, State::FrameType);
2066
2067        // Write a 0-length payload frame.
2068        b.put_varint(PUSH_PROMISE_FRAME_TYPE_ID).unwrap();
2069        b.put_varint(0).unwrap();
2070
2071        let mut cursor = std::io::Cursor::new(d);
2072
2073        // Parse frame type.
2074        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
2075        let frame_ty = stream.try_consume_varint().unwrap();
2076        assert_eq!(frame_ty, PUSH_PROMISE_FRAME_TYPE_ID);
2077
2078        stream.set_frame_type(frame_ty).unwrap();
2079        assert_eq!(stream.state, State::FramePayloadLen);
2080
2081        // Parse frame payload length.
2082        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
2083        let frame_payload_len = stream.try_consume_varint().unwrap();
2084        assert_eq!(
2085            Err(Error::FrameError),
2086            stream.set_frame_payload_len(frame_payload_len)
2087        );
2088    }
2089
2090    #[test]
2091    /// Drip feed data in chunks that exactly match spare capacity, forcing
2092    /// spare to hit 0 on every re-entry to try_fill_buffer_for_tests.
2093    fn large_state_buf_exact_spare_drip_feed() {
2094        const LARGE_HEADER_LEN: usize = 16384;
2095        let mut stream = Stream::new(
2096            0,
2097            false,
2098            LARGE_HEADER_LEN as u64,
2099            PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT,
2100        );
2101
2102        let mut d = vec![42; 20000];
2103        let mut b = octets::OctetsMut::with_slice(&mut d);
2104
2105        // Use nonzero fill to catch off-by-one errors in payload offset.
2106        let header_block = vec![0xAB; LARGE_HEADER_LEN];
2107        let hdrs = Frame::Headers {
2108            header_block: header_block.clone(),
2109        };
2110
2111        hdrs.to_bytes(&mut b).unwrap();
2112
2113        let mut cursor = std::io::Cursor::new(d);
2114
2115        // Parse frame type.
2116        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
2117        let frame_ty = stream.try_consume_varint().unwrap();
2118        assert_eq!(frame_ty, HEADERS_FRAME_TYPE_ID);
2119        stream.set_frame_type(frame_ty).unwrap();
2120
2121        // Parse frame payload length.
2122        let frame_payload_len =
2123            parse_multibyte_varint(&mut stream, &mut cursor).unwrap();
2124        assert_eq!(frame_payload_len, LARGE_HEADER_LEN as u64);
2125
2126        stream.set_frame_payload_len(frame_payload_len).unwrap();
2127        assert_eq!(stream.state, State::FramePayload);
2128
2129        // After state transition, initial reserve gives us
2130        // MAX_STATE_BUF_ALLOC_SIZE of spare capacity.
2131        assert_eq!(stream.state_buf.capacity(), MAX_STATE_BUF_ALLOC_SIZE);
2132        assert_eq!(stream.state_buf.len(), 0);
2133
2134        // Save the full cursor data then replace with an empty one so we
2135        // can drip feed exact amounts.
2136        let full_data = cursor.into_inner();
2137        let pos = 5; // 1 byte frame type + 4 byte varint (16384 >= 2^14)
2138        let payload_data = &full_data[pos..pos + LARGE_HEADER_LEN];
2139
2140        // Drip feed in chunks that exactly match MAX_STATE_BUF_ALLOC_SIZE.
2141        // Each chunk fully consumes spare, so on re-entry spare is exactly 0
2142        // and spare_state_buf must reserve.
2143        let mut fed = 0;
2144        while fed + MAX_STATE_BUF_ALLOC_SIZE <= LARGE_HEADER_LEN {
2145            let chunk = &payload_data[fed..fed + MAX_STATE_BUF_ALLOC_SIZE];
2146            let mut chunk_cursor = std::io::Cursor::new(chunk.to_vec());
2147
2148            let result = stream.try_fill_buffer_for_tests(&mut chunk_cursor);
2149
2150            fed += MAX_STATE_BUF_ALLOC_SIZE;
2151
2152            if fed < LARGE_HEADER_LEN {
2153                assert_eq!(result, Err(Error::Done));
2154                assert_eq!(stream.state_off, fed);
2155                // spare_state_buf should have reserved on each re-entry
2156                // since previous chunk consumed all spare.
2157                assert!(stream.state_buf.capacity() >= fed);
2158            } else {
2159                assert_eq!(result, Ok(()));
2160                assert_eq!(stream.state_off, LARGE_HEADER_LEN);
2161            }
2162        }
2163
2164        assert_eq!(
2165            stream.try_consume_frame(),
2166            Ok((hdrs, LARGE_HEADER_LEN as u64))
2167        );
2168        assert_eq!(stream.state, State::FrameType);
2169    }
2170
2171    #[test]
2172    /// Drip feed data in chunks smaller than spare capacity, so spare is
2173    /// small but nonzero on re-entry. Verifies that the buffer eventually
2174    /// grows when spare is fully consumed across multiple small reads.
2175    fn large_state_buf_small_leftover_spare() {
2176        const LARGE_HEADER_LEN: usize = 260000;
2177        let mut stream = Stream::new(
2178            0,
2179            false,
2180            LARGE_HEADER_LEN as u64,
2181            PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT,
2182        );
2183
2184        let mut d = vec![42; LARGE_HEADER_LEN + 10];
2185        let mut b = octets::OctetsMut::with_slice(&mut d);
2186
2187        let header_block = vec![0; LARGE_HEADER_LEN];
2188        let hdrs = Frame::Headers {
2189            header_block: header_block.clone(),
2190        };
2191
2192        hdrs.to_bytes(&mut b).unwrap();
2193
2194        let mut cursor = std::io::Cursor::new(d);
2195
2196        // Parse frame type.
2197        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
2198        let frame_ty = stream.try_consume_varint().unwrap();
2199        assert_eq!(frame_ty, HEADERS_FRAME_TYPE_ID);
2200        stream.set_frame_type(frame_ty).unwrap();
2201
2202        // Parse frame payload length.
2203        let frame_payload_len =
2204            parse_multibyte_varint(&mut stream, &mut cursor).unwrap();
2205        assert_eq!(frame_payload_len, LARGE_HEADER_LEN as u64);
2206
2207        stream.set_frame_payload_len(frame_payload_len).unwrap();
2208        assert_eq!(stream.state, State::FramePayload);
2209        assert_eq!(stream.state_buf.capacity(), MAX_STATE_BUF_ALLOC_SIZE);
2210
2211        let full_data = cursor.into_inner();
2212        let pos = 5; // 1 byte frame type + 4 byte varint length
2213        let payload_data = &full_data[pos..pos + LARGE_HEADER_LEN];
2214
2215        // Drip feed 1000-byte chunks. These don't align with the 4096 spare
2216        // capacity, so spare will be nonzero but shrinking on each re-entry
2217        // until it hits 0 and forces a reserve.
2218        let chunk_size = 1000;
2219        let mut fed = 0;
2220        while fed < LARGE_HEADER_LEN {
2221            let end = std::cmp::min(fed + chunk_size, LARGE_HEADER_LEN);
2222            let chunk = &payload_data[fed..end];
2223            let mut chunk_cursor = std::io::Cursor::new(chunk.to_vec());
2224
2225            let result = stream.try_fill_buffer_for_tests(&mut chunk_cursor);
2226
2227            fed = end;
2228
2229            if fed < LARGE_HEADER_LEN {
2230                assert_eq!(result, Err(Error::Done));
2231                assert_eq!(stream.state_off, fed);
2232                // Capacity must always be at least as large as what we've
2233                // buffered.
2234                assert!(stream.state_buf.capacity() >= stream.state_buf.len());
2235                // Capacity should grow incrementally. The allocator may
2236                // over-allocate (typically doubling), so we allow up to
2237                // 2x (bytes_read + reserve_size) to account for that.
2238                // The key property: capacity never jumps to the full
2239                // frame size before we've read a proportional amount.
2240                assert!(
2241                    stream.state_buf.capacity() <=
2242                        (fed + MAX_STATE_BUF_ALLOC_SIZE) * 2,
2243                    "capacity {} grew too far ahead of bytes read {} \
2244                     (max alloc size {})",
2245                    stream.state_buf.capacity(),
2246                    fed,
2247                    MAX_STATE_BUF_ALLOC_SIZE,
2248                );
2249            } else {
2250                assert_eq!(result, Ok(()));
2251            }
2252        }
2253
2254        assert_eq!(
2255            stream.try_consume_frame(),
2256            Ok((hdrs, LARGE_HEADER_LEN as u64))
2257        );
2258        assert_eq!(stream.state, State::FrameType);
2259    }
2260
2261    #[test]
2262    fn large_state_buf_allocation() {
2263        const LARGE_HEADER_LEN: usize = 260000;
2264        let mut stream = Stream::new(
2265            0,
2266            false,
2267            LARGE_HEADER_LEN as u64,
2268            PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT,
2269        );
2270        assert_eq!(stream.state_buf.capacity(), 16);
2271
2272        let mut d = vec![42; 5];
2273        let mut b = octets::OctetsMut::with_slice(&mut d);
2274
2275        // Drip feed a large HEADERS frame into the "stream"
2276        b.put_varint(HEADERS_FRAME_TYPE_ID).unwrap();
2277        b.put_varint(LARGE_HEADER_LEN as u64).unwrap();
2278
2279        let mut cursor = std::io::Cursor::new(d);
2280
2281        // Parse the HEADERS frame type.
2282        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
2283        assert_eq!(stream.state_buf.capacity(), 16);
2284
2285        let frame_ty = stream.try_consume_varint().unwrap();
2286        assert_eq!(frame_ty, HEADERS_FRAME_TYPE_ID);
2287        assert_eq!(stream.state_buf.capacity(), 16);
2288
2289        stream.set_frame_type(frame_ty).unwrap();
2290        assert_eq!(stream.state_buf.capacity(), 16);
2291
2292        // Parse the HEADERS frame payload length.
2293        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
2294        assert_eq!(stream.state_buf.capacity(), 16);
2295
2296        // Parse fails because we need more bytes for the 4-byte encoded length.
2297        // This trial then sets the expected buffer size for us to fill.
2298        assert_eq!(stream.try_consume_varint(), Err(Error::Done));
2299        stream.try_fill_buffer_for_tests(&mut cursor).unwrap();
2300        assert_eq!(stream.state_buf.capacity(), 16);
2301
2302        let frame_payload_len = stream.try_consume_varint().unwrap();
2303        assert_eq!(frame_payload_len, LARGE_HEADER_LEN as u64);
2304        assert_eq!(stream.state_buf.capacity(), 16);
2305
2306        stream.set_frame_payload_len(frame_payload_len).unwrap();
2307        assert_eq!(stream.state_buf.capacity(), MAX_STATE_BUF_ALLOC_SIZE);
2308
2309        /// Assert state_len, state_off, and state_buf.capacity() in one
2310        /// call, with labeled messages on failure.
2311        fn assert_state_buf_props(
2312            stream: &Stream, len: usize, off: usize, capacity: usize,
2313        ) {
2314            assert_eq!(stream.state_len, len, "state_len");
2315            assert_eq!(stream.state_off, off, "state_off");
2316            assert_eq!(stream.state_buf.capacity(), capacity, "capacity");
2317        }
2318
2319        // Start consuming HEADERS frame payload. It fails because the cursor
2320        // doesn't have the target size of data in it.
2321        assert_eq!(
2322            stream.try_fill_buffer_for_tests(&mut cursor),
2323            Err(Error::Done)
2324        );
2325        assert_state_buf_props(
2326            &stream,
2327            LARGE_HEADER_LEN,
2328            0,
2329            MAX_STATE_BUF_ALLOC_SIZE,
2330        );
2331
2332        // Drip feed data into the cursor to emulate a series of transport
2333        // reads. After set_frame_payload_len, the initial reserve gives us
2334        // MAX_STATE_BUF_ALLOC_SIZE (4096) bytes of spare capacity.
2335
2336        // Feed 2048 bytes: fits within the 4096 spare, no growth needed.
2337        cursor.get_mut().extend_from_slice(&[123; 2048]);
2338        assert_eq!(
2339            stream.try_fill_buffer_for_tests(&mut cursor),
2340            Err(Error::Done)
2341        );
2342        assert_state_buf_props(
2343            &stream,
2344            LARGE_HEADER_LEN,
2345            2048,
2346            MAX_STATE_BUF_ALLOC_SIZE,
2347        );
2348
2349        // Feed 1024 bytes: still fits in the remaining 2048 spare.
2350        cursor.get_mut().extend_from_slice(&[123; 1024]);
2351        assert_eq!(
2352            stream.try_fill_buffer_for_tests(&mut cursor),
2353            Err(Error::Done)
2354        );
2355        assert_state_buf_props(
2356            &stream,
2357            LARGE_HEADER_LEN,
2358            3072,
2359            MAX_STATE_BUF_ALLOC_SIZE,
2360        );
2361
2362        // Feed 512 bytes: fits in the remaining 1024 spare. Capacity
2363        // between state_off 4096 and 6144 stays stable because spare
2364        // doesn't hit 0.
2365        cursor.get_mut().extend_from_slice(&[123; 512]);
2366        assert_eq!(
2367            stream.try_fill_buffer_for_tests(&mut cursor),
2368            Err(Error::Done)
2369        );
2370        assert_state_buf_props(
2371            &stream,
2372            LARGE_HEADER_LEN,
2373            3584,
2374            MAX_STATE_BUF_ALLOC_SIZE,
2375        );
2376
2377        // Feed 4096 bytes: exceeds the remaining 512 spare. The loop reads
2378        // 512 to fill spare, then spare hits 0, triggers a reserve of
2379        // MAX_STATE_BUF_ALLOC_SIZE (4096), and reads the remaining 3584.
2380        // The allocator doubles capacity from 4096 to 8192.
2381        cursor.get_mut().extend_from_slice(&[123; 4096]);
2382        assert_eq!(
2383            stream.try_fill_buffer_for_tests(&mut cursor),
2384            Err(Error::Done)
2385        );
2386        assert_state_buf_props(
2387            &stream,
2388            LARGE_HEADER_LEN,
2389            7680,
2390            MAX_STATE_BUF_ALLOC_SIZE * 2,
2391        );
2392
2393        // Each subsequent feed exceeds spare (512 after the previous
2394        // reserve+read), so the loop drains the leftover spare, hits 0,
2395        // reserves MAX_STATE_BUF_ALLOC_SIZE, and the allocator doubles.
2396        cursor.get_mut().extend_from_slice(&[123; 8192]);
2397        assert_eq!(
2398            stream.try_fill_buffer_for_tests(&mut cursor),
2399            Err(Error::Done)
2400        );
2401        assert_state_buf_props(
2402            &stream,
2403            LARGE_HEADER_LEN,
2404            15872,
2405            MAX_STATE_BUF_ALLOC_SIZE * 4,
2406        );
2407
2408        cursor.get_mut().extend_from_slice(&[123; 16384]);
2409        assert_eq!(
2410            stream.try_fill_buffer_for_tests(&mut cursor),
2411            Err(Error::Done)
2412        );
2413        assert_state_buf_props(
2414            &stream,
2415            LARGE_HEADER_LEN,
2416            32256,
2417            MAX_STATE_BUF_ALLOC_SIZE * 8,
2418        );
2419
2420        cursor.get_mut().extend_from_slice(&[123; 32768]);
2421        assert_eq!(
2422            stream.try_fill_buffer_for_tests(&mut cursor),
2423            Err(Error::Done)
2424        );
2425        assert_state_buf_props(
2426            &stream,
2427            LARGE_HEADER_LEN,
2428            65024,
2429            MAX_STATE_BUF_ALLOC_SIZE * 16,
2430        );
2431
2432        cursor.get_mut().extend_from_slice(&[123; 65536]);
2433        assert_eq!(
2434            stream.try_fill_buffer_for_tests(&mut cursor),
2435            Err(Error::Done)
2436        );
2437        assert_state_buf_props(
2438            &stream,
2439            LARGE_HEADER_LEN,
2440            130560,
2441            MAX_STATE_BUF_ALLOC_SIZE * 32,
2442        );
2443
2444        // Feed the remaining bytes to complete the frame.
2445        let remaining = LARGE_HEADER_LEN - 130560;
2446        cursor.get_mut().extend_from_slice(&vec![123; remaining]);
2447        assert_eq!(stream.try_fill_buffer_for_tests(&mut cursor), Ok(()));
2448        assert_state_buf_props(
2449            &stream,
2450            LARGE_HEADER_LEN,
2451            LARGE_HEADER_LEN,
2452            MAX_STATE_BUF_ALLOC_SIZE * 64,
2453        );
2454
2455        let header_block = vec![123; LARGE_HEADER_LEN];
2456        let hdrs = Frame::Headers {
2457            header_block: header_block.clone(),
2458        };
2459        assert_eq!(
2460            stream.try_consume_frame(),
2461            Ok((hdrs, LARGE_HEADER_LEN as u64))
2462        );
2463        assert_eq!(stream.state, State::FrameType);
2464
2465        assert_state_buf_props(&stream, 1, 0, MAX_STATE_BUF_ALLOC_SIZE * 64);
2466    }
2467}