Skip to content

Commit 7b6f5d7

Browse files
committed
fix: Release flow control capacity when RecvStream is dropped
Dropping `RecvStream` discards buffered DATA, but doesn't release connection-level flow control capacity until all stream references are dropped. That can starve other streams of flow control capacity. With this change, we release capacity immediately on `RecvStream` drop.
1 parent 8fea6b2 commit 7b6f5d7

3 files changed

Lines changed: 121 additions & 5 deletions

File tree

src/proto/streams/recv.rs

Lines changed: 63 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -506,7 +506,7 @@ impl Recv {
506506
self.release_connection_capacity(stream.in_flight_recv_data, task);
507507
stream.in_flight_recv_data = 0;
508508

509-
self.clear_recv_buffer(stream);
509+
self.clear_recv_buffer(stream, task);
510510
}
511511

512512
/// Set the "target" connection window size.
@@ -933,9 +933,22 @@ impl Recv {
933933
stream.notify_push();
934934
}
935935

936-
pub(super) fn clear_recv_buffer(&mut self, stream: &mut Stream) {
937-
while stream.pending_recv.pop_front(&mut self.buffer).is_some() {
938-
// drop it
936+
pub(super) fn clear_recv_buffer(&mut self, stream: &mut Stream, task: &mut Option<Waker>) {
937+
let mut to_release: WindowSize = 0;
938+
while let Some(event) = stream.pending_recv.pop_front(&mut self.buffer) {
939+
if let Event::Data(data) = &event {
940+
to_release = to_release
941+
.saturating_add(data.len() as WindowSize)
942+
.min(stream.in_flight_recv_data);
943+
}
944+
}
945+
// Release flow control capacity. Cases:
946+
// * User read data but hasn't released: buf=0, in_flight>0 -> release 0
947+
// * User released without reading: buf>0, in_flight=0 -> release 0
948+
// * Normal drop without reading: buf=in_flight -> full release
949+
if to_release > 0 {
950+
stream.in_flight_recv_data -= to_release;
951+
self.release_connection_capacity(to_release, task);
939952
}
940953
}
941954

@@ -1253,6 +1266,52 @@ impl Recv {
12531266
}
12541267
}
12551268

1269+
#[cfg(test)]
1270+
mod tests {
1271+
use super::*;
1272+
1273+
#[test]
1274+
fn clear_recv_buffer_caps_capacity_before_overflow() {
1275+
const FRAME_LEN: usize = 1 << 20;
1276+
const FRAME_COUNT: usize = (u32::MAX as usize / FRAME_LEN) + 1;
1277+
1278+
let config = Config {
1279+
initial_max_send_streams: 0,
1280+
local_max_buffer_size: 0,
1281+
local_next_stream_id: 2.into(),
1282+
local_push_enabled: false,
1283+
extended_connect_protocol_enabled: false,
1284+
local_reset_duration: Duration::ZERO,
1285+
local_reset_max: 0,
1286+
remote_reset_max: 0,
1287+
remote_init_window_sz: DEFAULT_INITIAL_WINDOW_SIZE,
1288+
remote_max_initiated: None,
1289+
local_max_error_reset_streams: None,
1290+
};
1291+
let mut recv = Recv::new(peer::Dyn::Server, &config);
1292+
let mut store = Store::new();
1293+
let mut stream = store.insert(
1294+
StreamId::from(1),
1295+
Stream::new(StreamId::from(1), 0, DEFAULT_INITIAL_WINDOW_SIZE),
1296+
);
1297+
let data = Bytes::from(vec![0; FRAME_LEN]);
1298+
1299+
for _ in 0..FRAME_COUNT {
1300+
stream
1301+
.pending_recv
1302+
.push_back(&mut recv.buffer, Event::Data(data.clone()));
1303+
}
1304+
stream.in_flight_recv_data = DEFAULT_INITIAL_WINDOW_SIZE;
1305+
recv.in_flight_data = DEFAULT_INITIAL_WINDOW_SIZE;
1306+
1307+
recv.clear_recv_buffer(&mut stream, &mut None);
1308+
1309+
assert!(stream.pending_recv.is_empty());
1310+
assert_eq!(stream.in_flight_recv_data, 0);
1311+
assert_eq!(recv.in_flight_data, 0);
1312+
}
1313+
}
1314+
12561315
// ===== impl Open =====
12571316

12581317
impl Open {

src/proto/streams/streams.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1551,7 +1551,9 @@ impl OpaqueStreamRef {
15511551

15521552
let mut stream = me.store.resolve(self.key);
15531553
stream.is_recv = false;
1554-
me.actions.recv.clear_recv_buffer(&mut stream);
1554+
me.actions
1555+
.recv
1556+
.clear_recv_buffer(&mut stream, &mut me.actions.task);
15551557
}
15561558

15571559
pub fn stream_id(&self) -> StreamId {

tests/h2-tests/tests/flow_control.rs

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -544,6 +544,61 @@ async fn stream_error_release_connection_capacity() {
544544
join(srv, client).await;
545545
}
546546

547+
#[tokio::test]
548+
async fn recv_stream_drop_releases_only_buffered_connection_capacity() {
549+
h2_support::trace_init!();
550+
551+
const FRAME_LEN: usize = 16_384;
552+
const TOTAL_LEN: usize = FRAME_LEN * 2;
553+
554+
// Exercise all relationships between buffered and in-flight capacity:
555+
// equal, buffered > in-flight, and buffered < in-flight.
556+
for read_frames in 0usize..=1 {
557+
for released_frames in 0usize..=1 {
558+
let expected_used = read_frames.saturating_sub(released_frames) * FRAME_LEN;
559+
let (io, mut peer) = mock::new();
560+
561+
let peer = async move {
562+
let _ = peer.assert_server_handshake().await;
563+
peer.send_frame(frames::headers(1).request("POST", "https://example.com/"))
564+
.await;
565+
for _ in 0..2 {
566+
peer.send_frame(frames::data(1, vec![0; FRAME_LEN])).await;
567+
}
568+
569+
peer.recv_frame(frames::window_update(0, (TOTAL_LEN - expected_used) as u32))
570+
.await;
571+
if released_frames > 0 {
572+
peer.recv_frame(frames::window_update(1, FRAME_LEN as u32))
573+
.await;
574+
}
575+
};
576+
577+
let server = async move {
578+
let mut server = server::handshake(io).await.unwrap();
579+
let (request, _respond) = server.next().await.unwrap().unwrap();
580+
let mut body = request.into_body();
581+
582+
for _ in 0..read_frames {
583+
assert_eq!(body.data().await.unwrap().unwrap().len(), FRAME_LEN);
584+
}
585+
586+
let mut flow = body.flow_control().clone();
587+
for _ in 0..released_frames {
588+
flow.release_capacity(FRAME_LEN).unwrap();
589+
}
590+
drop(body);
591+
592+
assert_eq!(flow.used_capacity(), expected_used);
593+
594+
let _ = server.next().await;
595+
};
596+
597+
join(peer, server).await;
598+
}
599+
}
600+
}
601+
547602
// Regression test for TODO
548603
#[tokio::test]
549604
async fn padded_data_stream_error_releases_connection_capacity() {

0 commit comments

Comments
 (0)