diff --git a/quinn-h3/src/connection.rs b/quinn-h3/src/connection.rs index 62aace861..2818096ed 100644 --- a/quinn-h3/src/connection.rs +++ b/quinn-h3/src/connection.rs @@ -106,8 +106,8 @@ impl ConnectionRef { pub(crate) struct ConnectionInner { pub inner: Connection, - pub requests: VecDeque<(SendStream, RecvStream)>, - pub requests_task: Option, + requests: VecDeque<(SendStream, RecvStream)>, + requests_task: Option, side: Side, driver: Option, incoming_bi: IncomingBiStreams, @@ -135,6 +135,16 @@ impl ConnectionInner { Ok(self.inner.is_closing() && self.inner.requests_in_flight() == 0) } + pub fn next_request(&mut self, cx: &mut Context) -> Option<(SendStream, RecvStream)> { + match self.requests.pop_front() { + Some(x) => Some(x), + None => { + self.requests_task = Some(cx.waker().clone()); + Ok(None) + } + } + } + pub fn wake(&mut self) { if let Some(w) = self.driver.take() { w.wake(); diff --git a/quinn-h3/src/server.rs b/quinn-h3/src/server.rs index 45752d357..c27dfd6d7 100644 --- a/quinn-h3/src/server.rs +++ b/quinn-h3/src/server.rs @@ -168,18 +168,11 @@ pub struct IncomingRequest(ConnectionRef); impl Stream for IncomingRequest { type Item = RecvRequest; - fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context) -> Poll> { - let (send, recv) = { - let conn = &mut self.0.h3.lock().unwrap(); - match conn.requests.pop_front() { - Some(s) => s, - None => { - conn.requests_task = Some(cx.waker().clone()); - return Poll::Pending; - } - } - }; - Poll::Ready(Some(RecvRequest::new(recv, send, self.0.clone()))) + fn poll_next(self: Pin<&mut Self>, cx: &mut Context) -> Poll> { + match self.0.h3.lock().unwrap().next_request(cx) { + Some((s, r)) => Poll::Ready(Some(RecvRequest::new(r, s, self.0.clone()))), + None => Poll::Pending + } } }