Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
229 changes: 222 additions & 7 deletions h3/src/frame.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,9 @@ pub struct FrameStream<S, B> {
// Already read data from the stream
decoder: FrameDecoder,
remaining_data: usize,
// A connection error observed while frame bytes were still buffered.
// Those bytes are delivered first, then this error is surfaced exactly once.
pending_quic_error: Option<FrameStreamError>,
}

impl<S, B> FrameStream<S, B> {
Expand All @@ -33,6 +36,7 @@ impl<S, B> FrameStream<S, B> {
stream,
decoder: FrameDecoder::default(),
remaining_data: 0,
pending_quic_error: None,
}
}

Expand Down Expand Up @@ -83,6 +87,11 @@ where
None => {}
}

// Every buffered frame has been delivered: surface the deferred connection error.
if let Some(error) = self.pending_quic_error.take() {
return Poll::Ready(Err(error));
}

match self.try_recv(cx)? {
// Received a chunk but the frame is incomplete, poll until we get `Pending`.
Poll::Ready(false) => continue,
Expand Down Expand Up @@ -112,18 +121,34 @@ where
return Poll::Ready(Ok(None));
};

let end = match self.try_recv(cx) {
Poll::Ready(Ok(end)) => end,
Poll::Ready(Err(e)) => return Poll::Ready(Err(e)),
Poll::Pending => false,
// A connection error must not discard body bytes that are already buffered: defer it
// until they are drained. Any other error, notably a peer stream reset, surfaces now.
let end = if self.pending_quic_error.is_some() {
true
} else {
match self.try_recv(cx) {
Poll::Ready(Ok(end)) => end,
Poll::Ready(Err(e)) if e.is_connection_error() => {
self.pending_quic_error = Some(e);
true
}
Poll::Ready(Err(e)) => return Poll::Ready(Err(e)),
Poll::Pending => false,
}
};
let data = self.stream.buf_mut().take_chunk(self.remaining_data);

match (data, end) {
(None, true) => Poll::Ready(Ok(None)),
(None, true) => match self.pending_quic_error.take() {
Some(error) => Poll::Ready(Err(error)),
None => Poll::Ready(Ok(None)),
},
(None, false) => Poll::Pending,
// The stream finished cleanly in the middle of the frame. After a connection
// error, the partial body is delivered and the connection error follows instead.
(Some(d), true)
if d.remaining() < self.remaining_data
if self.pending_quic_error.is_none()
&& d.remaining() < self.remaining_data
&& !self.stream.buf_mut().has_remaining() =>
{
Poll::Ready(Err(FrameStreamError::UnexpectedEnd))
Expand Down Expand Up @@ -202,11 +227,13 @@ where
stream: send,
decoder: FrameDecoder::default(),
remaining_data: 0,
pending_quic_error: None,
},
FrameStream {
stream: recv,
decoder: self.decoder,
remaining_data: self.remaining_data,
pending_quic_error: self.pending_quic_error,
},
)
}
Expand Down Expand Up @@ -306,6 +333,22 @@ pub enum FrameStreamError {
UnexpectedEnd,
}

impl FrameStreamError {
/// Whether this error closed the whole QUIC connection, not just this stream.
///
/// Only a connection error is deferred behind buffered bytes. A peer stream reset is
/// reported once and frees the stream, and RFC 9114 section 7.1 only treats a truncated
/// frame as a connection error when the stream terminates cleanly. Holding a reset back
/// would turn a DATA frame it truncated into `UnexpectedEnd`, a connection-level
/// H3_FRAME_ERROR that tears down every other stream on the connection.
fn is_connection_error(&self) -> bool {
matches!(
self,
FrameStreamError::Quic(StreamErrorIncoming::ConnectionErrorIncoming { .. })
)
}
}

#[derive(Debug, PartialEq)]
/// Protocol specific errors that can occur while decoding frames in a stream
pub enum FrameProtocolError {
Expand All @@ -327,6 +370,7 @@ mod tests {
use std::{cell::Cell, collections::VecDeque, rc::Rc};

use crate::proto::{coding::Encode, frame::FrameType, varint::VarInt};
use crate::quic::ConnectionErrorIncoming;

// Decoder

Expand Down Expand Up @@ -480,6 +524,147 @@ mod tests {
);
}

/// Regression: a connection error read while a complete HEADERS frame is still buffered must
/// not discard it. The `poll_data` that drains the body reads the error with the trailing
/// HEADERS frame in the buffer.
#[tokio::test]
async fn poll_next_drains_buffered_headers_before_quic_close() {
let mut recv = FakeRecv::default();
let mut buf = BytesMut::with_capacity(64);
Frame::Data(&b"body"[..]).encode_with_payload(&mut buf);
Frame::headers(&b"trailer"[..]).encode_with_payload(&mut buf);
recv.chunk_then_error(buf.freeze(), connection_close());
let transport_polls = recv.poll_count.clone();
let mut stream: FrameStream<_, ()> = FrameStream::new(BufRecvStream::new(recv));

assert_poll_matches!(
|cx| stream.poll_next(cx),
Ok(Some(Frame::Data(PayloadLen(4))))
);
// This poll reads the connection error, with the body and the trailers still buffered.
assert_poll_matches!(
|cx| to_bytes(stream.poll_data(cx)),
Ok(Some(b)) if &*b == b"body"
);
assert_eq!(transport_polls.get(), 2);
assert_poll_matches!(|cx| to_bytes(stream.poll_data(cx)), Ok(None));
assert_poll_matches!(|cx| stream.poll_next(cx), Ok(Some(Frame::Headers(_))));
assert_poll_matches!(|cx| stream.poll_next(cx), Err(FrameStreamError::Quic(_)));
// The saved error is surfaced exactly once, without polling the transport again.
assert_eq!(transport_polls.get(), 2);
assert_poll_matches!(|cx| stream.poll_next(cx), Ok(None));
}

/// Regression: a connection error must not strand DATA frame body bytes that were buffered
/// before it arrived, and must still be reported once the body is delivered.
#[tokio::test]
async fn poll_data_drains_buffered_body_before_quic_close() {
let mut recv = FakeRecv::default();
let mut buf = BytesMut::with_capacity(64);
Frame::Data(Bytes::from("body")).encode_with_payload(&mut buf);
recv.chunk_then_error(buf.freeze(), connection_close());
let mut stream: FrameStream<_, ()> = FrameStream::new(BufRecvStream::new(recv));

assert_poll_matches!(
|cx| stream.poll_next(cx),
Ok(Some(Frame::Data(PayloadLen(4))))
);
// This poll reads the connection error but still returns the buffered body.
assert_poll_matches!(
|cx| to_bytes(stream.poll_data(cx)),
Ok(Some(b)) if b.remaining() == 4
);
// Otherwise the close would look like a clean end of stream to the caller.
assert_poll_matches!(|cx| stream.poll_next(cx), Err(FrameStreamError::Quic(_)));
}

/// A connection error that truncates a DATA frame delivers the buffered part of the body,
/// then surfaces as the connection error rather than as `UnexpectedEnd`.
#[tokio::test]
async fn poll_data_surfaces_quic_close_over_a_truncated_buffered_body() {
let mut recv = FakeRecv::default();
let mut buf = BytesMut::with_capacity(64);
FrameType::DATA.encode(&mut buf);
VarInt::from(8u32).encode(&mut buf);
buf.put_slice(&b"bo"[..]);
recv.chunk(buf.freeze());
recv.chunk_then_error(Bytes::from_static(b"dy"), connection_close());
let mut stream: FrameStream<_, ()> = FrameStream::new(BufRecvStream::new(recv));

assert_poll_matches!(
|cx| stream.poll_next(cx),
Ok(Some(Frame::Data(PayloadLen(8))))
);
assert_poll_matches!(
|cx| to_bytes(stream.poll_data(cx)),
Ok(Some(b)) if &*b == b"bo"
);
// This poll reads the connection error with "dy" buffered and the frame 4 bytes short.
assert_poll_matches!(
|cx| to_bytes(stream.poll_data(cx)),
Ok(Some(b)) if &*b == b"dy"
);
assert_poll_matches!(
|cx| to_bytes(stream.poll_data(cx)),
Err(FrameStreamError::Quic(StreamErrorIncoming::ConnectionErrorIncoming { .. }))
);
}

/// Regression: a peer RESET_STREAM that truncates a DATA frame while part of that frame is
/// still buffered must surface as the reset. Holding it back to drain the buffer would end
/// the body in `UnexpectedEnd`, which the request stream escalates to a connection-level
/// H3_FRAME_ERROR ("received incomplete frame").
#[tokio::test]
async fn poll_data_surfaces_stream_reset_over_a_truncated_buffered_body() {
let mut recv = FakeRecv::default();
let mut buf = BytesMut::with_capacity(64);
FrameType::DATA.encode(&mut buf);
VarInt::from(8u32).encode(&mut buf);
buf.put_slice(&b"bo"[..]);
recv.chunk(buf.freeze());
recv.chunk_then_error(Bytes::from_static(b"dy"), stream_reset());
let mut stream: FrameStream<_, ()> = FrameStream::new(BufRecvStream::new(recv));

assert_poll_matches!(
|cx| stream.poll_next(cx),
Ok(Some(Frame::Data(PayloadLen(8))))
);
// This poll pulls "dy" into the buffer behind "bo".
assert_poll_matches!(
|cx| to_bytes(stream.poll_data(cx)),
Ok(Some(b)) if &*b == b"bo"
);
// The reset arrives with "dy" buffered and the frame 4 bytes short.
assert_poll_matches!(
|cx| to_bytes(stream.poll_data(cx)),
Err(FrameStreamError::Quic(StreamErrorIncoming::StreamTerminated { error_code }))
if error_code == RESET_CODE
);
}

/// Regression: a peer RESET_STREAM read while whole frames are still buffered surfaces as the
/// reset instead of those frames. Quinn reports a reset only once, so holding it back would
/// lose its error code.
#[tokio::test]
async fn poll_data_surfaces_stream_reset_over_a_buffered_frame() {
let mut recv = FakeRecv::default();
let mut buf = BytesMut::with_capacity(64);
Frame::Data(&b"body"[..]).encode_with_payload(&mut buf);
Frame::headers(&b"trailer"[..]).encode_with_payload(&mut buf);
recv.chunk_then_error(buf.freeze(), stream_reset());
let mut stream: FrameStream<_, ()> = FrameStream::new(BufRecvStream::new(recv));

assert_poll_matches!(
|cx| stream.poll_next(cx),
Ok(Some(Frame::Data(PayloadLen(4))))
);
assert_poll_matches!(
|cx| to_bytes(stream.poll_data(cx)),
Err(FrameStreamError::Quic(StreamErrorIncoming::StreamTerminated { error_code }))
if error_code == RESET_CODE
);
}

#[tokio::test]
async fn poll_next_incomplete_frame() {
let mut recv = FakeRecv::default();
Expand Down Expand Up @@ -636,9 +821,25 @@ mod tests {

// Helpers

/// `H3_REQUEST_CANCELLED`, the code a peer that cuts a response resets the stream with.
const RESET_CODE: u64 = 0x010c;

fn connection_close() -> StreamErrorIncoming {
StreamErrorIncoming::ConnectionErrorIncoming {
connection_error: ConnectionErrorIncoming::ApplicationClose { error_code: 0x100 },
}
}

fn stream_reset() -> StreamErrorIncoming {
StreamErrorIncoming::StreamTerminated {
error_code: RESET_CODE,
}
}

#[derive(Default)]
struct FakeRecv {
chunks: VecDeque<Bytes>,
pending_error: Option<StreamErrorIncoming>,
poll_count: Rc<Cell<usize>>,
}

Expand All @@ -647,6 +848,14 @@ mod tests {
self.chunks.push_back(buf);
self
}

/// Queue a last chunk, after which the transport reports `err` instead of the end of the
/// stream.
fn chunk_then_error(&mut self, buf: Bytes, err: StreamErrorIncoming) -> &mut Self {
self.chunks.push_back(buf);
self.pending_error = Some(err);
self
}
}

impl RecvStream for FakeRecv {
Expand All @@ -657,7 +866,13 @@ mod tests {
_: &mut Context<'_>,
) -> Poll<Result<Option<Self::Buf>, StreamErrorIncoming>> {
self.poll_count.set(self.poll_count.get() + 1);
Poll::Ready(Ok(self.chunks.pop_front()))
if let Some(chunk) = self.chunks.pop_front() {
return Poll::Ready(Ok(Some(chunk)));
}
match self.pending_error.take() {
Some(err) => Poll::Ready(Err(err)),
None => Poll::Ready(Ok(None)),
}
}

fn stop_sending(&mut self, _: u64) {
Expand Down