Skip to content

Commit af8788e

Browse files
authored
stream: track FIN acknowledgements separately
SendBuf::is_complete() used acknowledged byte ranges as a proxy for FIN acknowledgement. A zero-length FIN contributes no byte range, so an ACK for an earlier data-only frame is indistinguishable from an ACK for the FIN. Fixes #2525.
1 parent ce57b25 commit af8788e

4 files changed

Lines changed: 110 additions & 4 deletions

File tree

‎quiche/src/lib.rs‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3616,7 +3616,7 @@ impl<F: BufFactory> Connection<F> {
36163616
stream_id,
36173617
offset,
36183618
length,
3619-
..
3619+
fin,
36203620
} => {
36213621
// Emit qlog before checking if the stream still exists.
36223622
// The client does need to ACK frames that were received
@@ -3647,6 +3647,11 @@ impl<F: BufFactory> Connection<F> {
36473647
};
36483648

36493649
let dropped = stream.send.ack_and_drop(offset, length);
3650+
3651+
if fin {
3652+
stream.send.ack_fin();
3653+
}
3654+
36503655
let priority_key = Arc::clone(&stream.priority_key);
36513656

36523657
// Only collect the stream if it is complete and not

‎quiche/src/stream/mod.rs‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1750,6 +1750,9 @@ mod tests {
17501750
assert!(!stream.send.is_complete());
17511751

17521752
stream.send.ack(0, 1);
1753+
assert!(!stream.send.is_complete());
1754+
1755+
stream.send.ack_fin();
17531756
assert!(stream.send.is_complete());
17541757

17551758
assert!(!stream.is_complete());

‎quiche/src/stream/send_buf.rs‎

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -122,6 +122,9 @@ where
122122
/// The final stream offset written to the stream, if any.
123123
fin_off: Option<u64>,
124124

125+
/// Whether a STREAM frame carrying FIN has been acknowledged.
126+
fin_acked: bool,
127+
125128
/// Whether the stream's send-side has been shut down.
126129
shutdown: bool,
127130

@@ -324,6 +327,11 @@ impl<F: BufFactory> SendBuf<F> {
324327
self.acked.insert(off..off + len as u64);
325328
}
326329

330+
/// Marks the stream's final size as acknowledged by the peer.
331+
pub(crate) fn ack_fin(&mut self) {
332+
self.fin_acked = true;
333+
}
334+
327335
pub fn ack_and_drop(&mut self, off: u64, len: usize) -> usize {
328336
self.ack(off, len);
329337

@@ -528,11 +536,13 @@ impl<F: BufFactory> SendBuf<F> {
528536

529537
/// Returns true if the send-side of the stream is complete.
530538
///
531-
/// This happens when the stream's send final size is known, and the peer
532-
/// has already acked all stream data up to that point.
539+
/// This happens when the peer has acknowledged the stream's final size and
540+
/// all stream data up to that point, or the stream has been reset.
533541
pub fn is_complete(&self) -> bool {
534542
if let Some(fin_off) = self.fin_off {
535-
if self.acked == (0..fin_off) {
543+
if (self.fin_acked || self.shutdown || self.error.is_some()) &&
544+
self.acked == (0..fin_off)
545+
{
536546
return true;
537547
}
538548
}

‎quiche/src/tests.rs‎

Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13155,6 +13155,94 @@ fn stop_sending_no_retransmit_after_fin(
1315513155
assert_eq!(pipe.server.send(&mut buf), Err(Error::Done));
1315613156
}
1315713157

13158+
#[rstest]
13159+
/// Tests that a zero-length FIN is sent and acknowledged before the stream is
13160+
/// collected.
13161+
fn stream_empty_fin_is_not_complete_until_fin_acked(
13162+
#[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
13163+
) {
13164+
let mut buf = [0; 65535];
13165+
13166+
let mut pipe = test_utils::Pipe::new(cc_algorithm_name).unwrap();
13167+
assert_eq!(pipe.handshake(), Ok(()));
13168+
13169+
assert_eq!(pipe.client.stream_send(4, b"hello", true), Ok(5));
13170+
assert_eq!(pipe.advance(), Ok(()));
13171+
assert_eq!(pipe.server.stream_recv(4, &mut buf), Ok((5, true)));
13172+
13173+
assert_eq!(pipe.server.stream_send(4, b"response", false), Ok(8));
13174+
13175+
let flight = test_utils::emit_flight(&mut pipe.server).unwrap();
13176+
test_utils::process_flight(&mut pipe.client, flight).unwrap();
13177+
13178+
assert_eq!(pipe.client.stream_recv(4, &mut buf), Ok((8, false)));
13179+
13180+
assert_eq!(pipe.server.stream_send(4, b"", true), Ok(0));
13181+
13182+
let flight = test_utils::emit_flight(&mut pipe.client).unwrap();
13183+
test_utils::process_flight(&mut pipe.server, flight).unwrap();
13184+
13185+
assert!(!pipe.server.streams.is_collected(4));
13186+
13187+
let flight = test_utils::emit_flight(&mut pipe.server).unwrap();
13188+
let mut frames = Vec::new();
13189+
13190+
for (mut pkt, _) in flight {
13191+
frames
13192+
.extend(test_utils::decode_pkt(&mut pipe.client, &mut pkt).unwrap());
13193+
}
13194+
13195+
assert!(frames.iter().any(|f| matches!(
13196+
f,
13197+
frame::Frame::Stream { stream_id: 4, data }
13198+
if data.off() == 8 && data.is_empty() && data.fin()
13199+
)));
13200+
}
13201+
13202+
#[rstest]
13203+
/// Tests that a lost zero-length FIN is retransmitted after all stream data has
13204+
/// already been acknowledged.
13205+
fn stream_empty_fin_is_retransmitted(
13206+
#[values("cubic", "bbr2_gcongestion")] cc_algorithm_name: &str,
13207+
) {
13208+
let mut buf = [0; 65535];
13209+
13210+
let mut pipe = test_utils::Pipe::new(cc_algorithm_name).unwrap();
13211+
assert_eq!(pipe.handshake(), Ok(()));
13212+
13213+
assert_eq!(pipe.client.stream_send(4, b"hello", true), Ok(5));
13214+
assert_eq!(pipe.advance(), Ok(()));
13215+
assert_eq!(pipe.server.stream_recv(4, &mut buf), Ok((5, true)));
13216+
13217+
assert_eq!(pipe.server.stream_send(4, b"response", false), Ok(8));
13218+
assert_eq!(pipe.advance(), Ok(()));
13219+
assert_eq!(pipe.client.stream_recv(4, &mut buf), Ok((8, false)));
13220+
13221+
assert_eq!(pipe.server.stream_send(4, b"", true), Ok(0));
13222+
13223+
let (len, _) = pipe.server.send(&mut buf).unwrap();
13224+
let frames =
13225+
test_utils::decode_pkt(&mut pipe.client, &mut buf[..len]).unwrap();
13226+
13227+
assert!(frames.iter().any(|f| matches!(
13228+
f,
13229+
frame::Frame::Stream { stream_id: 4, data }
13230+
if data.off() == 8 && data.is_empty() && data.fin()
13231+
)));
13232+
13233+
test_utils::trigger_ack_based_loss(&mut pipe.server, &mut pipe.client);
13234+
13235+
let (len, _) = pipe.server.send(&mut buf).unwrap();
13236+
let frames =
13237+
test_utils::decode_pkt(&mut pipe.client, &mut buf[..len]).unwrap();
13238+
13239+
assert!(frames.iter().any(|f| matches!(
13240+
f,
13241+
frame::Frame::Stream { stream_id: 4, data }
13242+
if data.off() == 8 && data.is_empty() && data.fin()
13243+
)));
13244+
}
13245+
1315813246
#[rstest]
1315913247
/// Tests that MAX_STREAMS_BIDI frames are retransmitted if lost
1316013248
fn max_streams_bidi_frame_retransmit(

0 commit comments

Comments
 (0)