Skip to content

Commit 1acc5ba

Browse files
committed
h3: ignore priority updates for closed streams
Prevent late PRIORITY_UPDATE frames from recreating H3 stream state after the request stream has completed but before transport collection. Add a regression test covering H3 collection before QUIC stream collection.
1 parent cbc8173 commit 1acc5ba

3 files changed

Lines changed: 148 additions & 7 deletions

File tree

quiche/src/h3/mod.rs

Lines changed: 63 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -3202,13 +3202,11 @@ impl Connection {
32023202
return Err(Error::IdError);
32033203
}
32043204

3205-
// If the PRIORITY_UPDATE is valid, consider storing the latest
3206-
// contents. Due to reordering, it is possible that we might
3207-
// receive frames that reference streams that have not yet to
3208-
// been opened and that's OK because it's within our concurrency
3209-
// limit. However, we discard PRIORITY_UPDATE that refers to
3210-
// streams that we know have been collected.
3211-
if conn.streams.is_collected(prioritized_element_id) {
3205+
// PRIORITY_UPDATE can arrive before the request stream exists,
3206+
// so a missing transport stream is allowed. Ignore updates only
3207+
// once the transport stream was collected or both transport
3208+
// directions are finished.
3209+
if conn.stream_closed(prioritized_element_id) {
32123210
return Err(Error::Done);
32133211
}
32143212

@@ -5236,6 +5234,64 @@ mod tests {
52365234
assert_eq!(s.poll_server(), Err(Error::Done));
52375235
}
52385236

5237+
#[test]
5238+
/// Send a PRIORITY_UPDATE for a request stream after H3 has collected it,
5239+
/// but before the transport stream has been collected.
5240+
fn priority_update_request_after_h3_collection() {
5241+
let mut s = Session::new().unwrap();
5242+
s.handshake().unwrap();
5243+
5244+
let init_streams_server = s.server.streams.len();
5245+
5246+
let (stream, req) = s.send_request(true).unwrap();
5247+
let ev_headers = Event::Headers {
5248+
list: req,
5249+
more_frames: false,
5250+
};
5251+
5252+
assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
5253+
assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
5254+
assert_eq!(s.poll_server(), Err(Error::Done));
5255+
5256+
let resp = vec![
5257+
Header::new(b":status", b"200"),
5258+
Header::new(b"server", b"quiche-test"),
5259+
];
5260+
5261+
s.server
5262+
.send_response(&mut s.pipe.server, stream, &resp, true)
5263+
.unwrap();
5264+
5265+
// H3 no longer needs its stream state once it has sent the response
5266+
// FIN and consumed the request FIN. The QUIC stream remains until the
5267+
// response FIN is acknowledged by the peer.
5268+
assert_eq!(s.server.streams.len(), init_streams_server);
5269+
assert!(s.pipe.server.stream_finished(stream));
5270+
assert!(s.pipe.server.stream_closed(stream));
5271+
5272+
let stream_state = s.pipe.server.streams.get(stream).unwrap();
5273+
assert!(stream_state.recv.is_fin());
5274+
assert!(stream_state.send.is_fin());
5275+
assert!(!s.pipe.server.streams.is_collected(stream));
5276+
5277+
s.client
5278+
.send_priority_update_for_request(
5279+
&mut s.pipe.client,
5280+
stream,
5281+
&Priority {
5282+
urgency: 3,
5283+
incremental: false,
5284+
},
5285+
)
5286+
.unwrap();
5287+
5288+
let flight = crate::test_utils::emit_flight(&mut s.pipe.client).unwrap();
5289+
crate::test_utils::process_flight(&mut s.pipe.server, flight).unwrap();
5290+
5291+
assert_eq!(s.poll_server(), Err(Error::Done));
5292+
assert_eq!(s.server.streams.len(), init_streams_server);
5293+
}
5294+
52395295
#[test]
52405296
/// Send a PRIORITY_UPDATE for a request stream, before and after the stream
52415297
/// has been stopped.

quiche/src/lib.rs

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6540,6 +6540,36 @@ impl<F: BufFactory> Connection<F> {
65406540
stream.recv.is_fin()
65416541
}
65426542

6543+
/// Returns true if the specified stream is closed.
6544+
///
6545+
/// For bidirectional streams this happens when both the receive and send
6546+
/// sides have signaled `fin`. For unidirectional streams only the
6547+
/// relevant direction is checked, depending on whether the stream was
6548+
/// created locally or not.
6549+
///
6550+
/// This also returns true if the stream has already been collected, but
6551+
/// returns false if the stream was never opened.
6552+
#[inline]
6553+
pub fn stream_closed(&self, stream_id: u64) -> bool {
6554+
let Some(stream) = self.streams.get(stream_id) else {
6555+
return self.streams.is_collected(stream_id);
6556+
};
6557+
6558+
match (stream.bidi, stream.local) {
6559+
// For bidirectional streams both directions must have signaled
6560+
// FIN.
6561+
(true, _) => stream.recv.is_fin() && stream.send.is_fin(),
6562+
6563+
// For unidirectional streams created locally, only the send side
6564+
// is checked.
6565+
(false, true) => stream.send.is_fin(),
6566+
6567+
// For unidirectional streams created by the peer, only the
6568+
// receive side is checked.
6569+
(false, false) => stream.recv.is_fin(),
6570+
}
6571+
}
6572+
65436573
/// Returns the number of bidirectional streams that can be created
65446574
/// before the peer's stream count limit is reached.
65456575
///

quiche/src/tests.rs

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1013,6 +1013,61 @@ fn streamio(#[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str) {
10131013
assert!(pipe.server.stream_finished(4));
10141014
}
10151015

1016+
#[rstest]
1017+
fn stream_closed_bidi(
1018+
#[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
1019+
) {
1020+
let mut pipe = test_utils::Pipe::new(cc_algorithm_name).unwrap();
1021+
assert_eq!(pipe.handshake(), Ok(()));
1022+
1023+
assert!(!pipe.client.stream_closed(0));
1024+
assert!(!pipe.server.stream_closed(0));
1025+
1026+
assert_eq!(pipe.client.stream_send(0, b"hello", true), Ok(5));
1027+
assert_eq!(pipe.advance(), Ok(()));
1028+
1029+
assert!(!pipe.client.stream_closed(0));
1030+
assert!(!pipe.server.stream_closed(0));
1031+
1032+
let mut buf = [0; 5];
1033+
assert_eq!(pipe.server.stream_recv(0, &mut buf), Ok((5, true)));
1034+
1035+
assert!(!pipe.server.stream_closed(0));
1036+
1037+
assert_eq!(pipe.server.stream_send(0, b"world", true), Ok(5));
1038+
assert!(pipe.server.stream_closed(0));
1039+
1040+
assert_eq!(pipe.advance(), Ok(()));
1041+
1042+
assert!(!pipe.client.stream_closed(0));
1043+
assert_eq!(pipe.client.stream_recv(0, &mut buf), Ok((5, true)));
1044+
assert!(pipe.client.stream_closed(0));
1045+
}
1046+
1047+
#[rstest]
1048+
fn stream_closed_uni(
1049+
#[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
1050+
) {
1051+
let mut pipe = test_utils::Pipe::new(cc_algorithm_name).unwrap();
1052+
assert_eq!(pipe.handshake(), Ok(()));
1053+
1054+
assert!(!pipe.client.stream_closed(2));
1055+
assert!(!pipe.server.stream_closed(2));
1056+
1057+
assert_eq!(pipe.client.stream_send(2, b"hello", true), Ok(5));
1058+
1059+
assert!(pipe.client.stream_closed(2));
1060+
1061+
assert_eq!(pipe.advance(), Ok(()));
1062+
1063+
assert!(!pipe.server.stream_closed(2));
1064+
1065+
let mut buf = [0; 5];
1066+
assert_eq!(pipe.server.stream_recv(2, &mut buf), Ok((5, true)));
1067+
1068+
assert!(pipe.server.stream_closed(2));
1069+
}
1070+
10161071
/// Test receiving into `BufMut`
10171072
#[rstest]
10181073
fn stream_recv_buf(

0 commit comments

Comments
 (0)