From c2a7377fa4b10807bc5d503dc88df046046de38f Mon Sep 17 00:00:00 2001 From: Daniel Boros Date: Thu, 1 Oct 2026 08:31:16 +0200 Subject: [PATCH] perf: avoid reallocating the read buffer to watch for EOF --- src/proto/h1/conn.rs | 2 +- src/proto/h1/io.rs | 63 +++++++++++++++++++++++++++++++++++++++++++- 2 files changed, 63 insertions(+), 2 deletions(-) diff --git a/src/proto/h1/conn.rs b/src/proto/h1/conn.rs index 69ca50b1a4..3a30849040 100644 --- a/src/proto/h1/conn.rs +++ b/src/proto/h1/conn.rs @@ -506,7 +506,7 @@ where fn force_io_read(&mut self, cx: &mut Context<'_>) -> Poll> { debug_assert!(!self.state.is_read_closed()); - let result = ready!(self.io.poll_read_from_io(cx)); + let result = ready!(self.io.poll_read_from_io_spare(cx)); Poll::Ready(result.map_err(|e| { trace!(error = %e, "force_io_read; io error"); self.state.close(); diff --git a/src/proto/h1/io.rs b/src/proto/h1/io.rs index 5833004768..20f88c138b 100644 --- a/src/proto/h1/io.rs +++ b/src/proto/h1/io.rs @@ -222,6 +222,21 @@ where } pub(crate) fn poll_read_from_io(&mut self, cx: &mut Context<'_>) -> Poll> { + self.poll_read_from_io_inner(cx, ReadSize::Next) + } + + pub(crate) fn poll_read_from_io_spare( + &mut self, + cx: &mut Context<'_>, + ) -> Poll> { + self.poll_read_from_io_inner(cx, ReadSize::Spare) + } + + fn poll_read_from_io_inner( + &mut self, + cx: &mut Context<'_>, + size: ReadSize, + ) -> Poll> { self.read_blocked = false; // Get the next amount to allocate, but make sure we don't go over // the max read buf size configured. @@ -231,7 +246,8 @@ where .max() .saturating_sub(self.read_buf.len()), ); - if self.read_buf_remaining_mut() < next { + let remaining = self.read_buf_remaining_mut(); + if remaining < next && (matches!(size, ReadSize::Next) || remaining == 0) { self.read_buf.reserve(next); } @@ -362,6 +378,12 @@ where } } +#[derive(Clone, Copy, Debug)] +enum ReadSize { + Next, + Spare, +} + #[derive(Clone, Copy, Debug)] enum ReadStrategy { Adaptive { @@ -717,6 +739,45 @@ mod tests { ); } + #[cfg(all(feature = "server", not(miri)))] + #[tokio::test] + async fn spare_read_does_not_reallocate_a_shared_buffer() { + use crate::proto::h1::ServerTransaction; + + let mock = Mock::new() + .read(b"GET / HTTP/1.1\r\nHost: example.com\r\n\r\n") + .wait(Duration::from_secs(1)) + .build(); + let mut buffered = Buffered::<_, Cursor>>::new(Compat::new(mock)); + + let msg = futures_util::future::poll_fn(|cx| { + let parse_ctx = ParseContext { + cached_headers: &mut None, + req_method: &mut None, + h1_parser_config: Default::default(), + h1_max_headers: None, + preserve_header_case: false, + #[cfg(feature = "ffi")] + preserve_header_order: false, + h09_responses: false, + #[cfg(feature = "client")] + on_informational: &mut None, + }; + buffered.parse::(cx, parse_ctx) + }) + .await + .expect("parse"); + + let before = buffered.read_buf.as_ptr(); + futures_util::future::poll_fn(|cx| { + assert!(buffered.poll_read_from_io_spare(cx).is_pending()); + Poll::Ready(()) + }) + .await; + assert_eq!(buffered.read_buf.as_ptr(), before); + drop(msg); + } + #[test] fn read_strategy_adaptive_increments() { let mut strategy = ReadStrategy::default();