diff --git a/src/h3/mod.rs b/src/h3/mod.rs index f5203d129..4eb325c75 100644 --- a/src/h3/mod.rs +++ b/src/h3/mod.rs @@ -1688,6 +1688,10 @@ impl Connection { return Ok((finished, Event::Finished)); } + for stream_id in conn.streams.collected_streams() { + self.streams.remove(&stream_id); + } + // Process queued DATAGRAMs if the poll threshold allows it. match self.process_dgrams(conn) { Ok(v) => return Ok(v), diff --git a/src/stream.rs b/src/stream.rs index 6e978bbc7..1e45c6bd6 100644 --- a/src/stream.rs +++ b/src/stream.rs @@ -105,6 +105,10 @@ pub struct StreamMap { /// created streams, to prevent peers from re-creating them. collected: StreamIdHashSet, + /// Queue of streams to notify a transport stream of their completion + /// via the `.collected_streams()` call. + pending_collected: VecDeque, + /// Peer's maximum bidirectional stream count limit. peer_max_streams_bidi: u64, @@ -548,7 +552,9 @@ impl StreamMap { self.mark_writable(stream_id, false); self.streams.remove(&stream_id); + self.collected.insert(stream_id); + self.pending_collected.push_front(stream_id); } /// Creates an iterator over streams that have outstanding data to read. @@ -638,6 +644,11 @@ impl StreamMap { pub fn len(&self) -> usize { self.streams.len() } + + /// Returns streams which have been removed but not yet processed + pub fn collected_streams<'a>(&'a mut self) -> impl Iterator + 'a { + self.pending_collected.drain(..) + } } /// A QUIC stream.