From a31e72fe420ef51cc86798f16e0bf2a1edd7ba06 Mon Sep 17 00:00:00 2001 From: Blaise Bruer Date: Sat, 1 Aug 2026 18:51:34 -0500 Subject: [PATCH 1/2] client: return one-frame response bodies without copying collect_body_bounded copied every frame into a BytesMut, including the common case of a body that arrives whole in one frame. Hold the first frame by reference count and return it directly when no second frame follows; promote to a buffer only when one does. Using http_body_util::Limited here instead would need B::Error to be std::error::Error rather than Display, which propagates to ServerStream and BidiStream and the transport trait's body. Not worth a public API break for this. Signed-off-by: Blaise Bruer --- .../unreleased/Changed-20260801-190000.yaml | 9 ++ connectrpc/src/client/mod.rs | 82 ++++++++++++++++--- 2 files changed, 81 insertions(+), 10 deletions(-) create mode 100644 .changes/unreleased/Changed-20260801-190000.yaml diff --git a/.changes/unreleased/Changed-20260801-190000.yaml b/.changes/unreleased/Changed-20260801-190000.yaml new file mode 100644 index 00000000..52f05bd8 --- /dev/null +++ b/.changes/unreleased/Changed-20260801-190000.yaml @@ -0,0 +1,9 @@ +kind: Changed +body: |- + **Client response bodies that arrive in one frame are no longer copied.** + `collect_body_bounded` holds the first data frame by reference count and + returns it directly when no second frame follows, which is the common case + for unary responses and error bodies. Multi-frame bodies concatenate as + before, now starting from an exactly-sized buffer. The size limit is still + enforced before any frame is retained. +time: 2026-08-01T19:00:00.000000000Z diff --git a/connectrpc/src/client/mod.rs b/connectrpc/src/client/mod.rs index f6c67e07..3a191fee 100644 --- a/connectrpc/src/client/mod.rs +++ b/connectrpc/src/client/mod.rs @@ -4368,22 +4368,42 @@ where B: Body, B::Error: std::fmt::Display, { - let mut buf = BytesMut::new(); let mut stream = std::pin::pin!(body); + // Hold the first frame by refcount so a one-frame body is never copied. + let mut head: Option = None; + let mut joined: Option = None; + let mut len: usize = 0; + loop { match std::future::poll_fn(|cx| stream.as_mut().poll_frame(cx)).await { Some(Ok(frame)) => { // Trailer frames are skipped: Connect unary/error bodies don't // use HTTP trailers (those come via `trailer-` prefixed headers // or the JSON body). - if let Ok(data) = frame.into_data() { - if buf.len().saturating_add(data.len()) > max_size { - return Err(ConnectError::new( - ErrorCode::ResourceExhausted, - format!("response body size exceeds limit {max_size}"), - )); - } - buf.extend_from_slice(&data); + let Ok(data) = frame.into_data() else { + continue; + }; + len = len.saturating_add(data.len()); + if len > max_size { + return Err(ConnectError::new( + ErrorCode::ResourceExhausted, + format!("response body size exceeds limit {max_size}"), + )); + } + if data.is_empty() { + continue; + } + match &mut joined { + Some(buf) => buf.extend_from_slice(&data), + None => match head.take() { + None => head = Some(data), + Some(first) => { + let mut buf = BytesMut::with_capacity(len); + buf.extend_from_slice(&first); + buf.extend_from_slice(&data); + joined = Some(buf); + } + }, } } Some(Err(e)) => { @@ -4394,7 +4414,12 @@ where None => break, } } - Ok(buf.freeze()) + + Ok(match (head, joined) { + (_, Some(buf)) => buf.freeze(), + (Some(first), None) => first, + (None, None) => Bytes::new(), + }) } /// Percent-decode a gRPC message string. @@ -7215,6 +7240,43 @@ mod tests { assert_eq!(&got[..], b"foobar"); } + /// A one-frame body must be handed back by reference count, not copied. + /// This is the point of the single-frame path, and nothing else pins it. + #[tokio::test] + async fn collect_body_bounded_single_frame_does_not_copy() { + let src = Bytes::from(vec![9u8; 4096]); + + let got = collect_body_bounded(Full::new(src.clone()), 8192) + .await + .unwrap(); + + assert_eq!(got.len(), 4096); + assert!( + std::ptr::addr_eq(got.as_ptr(), src.as_ptr()), + "single-frame body must not be copied" + ); + } + + /// The promotion path still concatenates correctly once a second frame + /// arrives, and the result no longer aliases the first frame. + #[tokio::test] + async fn collect_body_bounded_multi_frame_promotes() { + let (tx, rx) = tokio::sync::mpsc::channel(4); + let body = ChannelBody { rx }; + let first = Bytes::from_static(b"foo"); + tx.send(Ok(first.clone())).await.unwrap(); + tx.send(Ok(Bytes::from_static(b"bar"))).await.unwrap(); + tx.send(Ok(Bytes::from_static(b"baz"))).await.unwrap(); + drop(tx); + + let got = collect_body_bounded(body, 100).await.unwrap(); + assert_eq!(&got[..], b"foobarbaz"); + assert!( + !std::ptr::addr_eq(got.as_ptr(), first.as_ptr()), + "a promoted body must be a fresh buffer" + ); + } + #[tokio::test] async fn collect_body_bounded_propagates_body_error() { let (tx, rx) = tokio::sync::mpsc::channel(4); From 4711503c2248462d329735395bcfbd8ee9386c32 Mon Sep 17 00:00:00 2001 From: Iain McGinniss <309153+iainmcgin@users.noreply.github.com> Date: Thu, 24 Sep 2026 12:06:02 -0700 Subject: [PATCH 2/2] client: fix collect_body_bounded tests and document buffer aliasing The merge with main gave collect_body_bounded a deadline argument; pass None in the tests that predate it. Document that a one-frame body is returned without a copy and so may be a slice of the transport's read buffer, say so in the changelog fragment instead of naming the private function, and drop the claim of an exactly-sized buffer. Cover an empty frame before and after the only data frame, and pin that a multi-frame body is joined into a new buffer. Assert that head and joined are never both set. Signed-off-by: Iain McGinniss <309153+iainmcgin@users.noreply.github.com> --- .../unreleased/Changed-20260801-190000.yaml | 9 --- .../unreleased/Changed-20260924-190142.yaml | 14 ++++ connectrpc/src/client/mod.rs | 77 +++++++++++++++---- 3 files changed, 78 insertions(+), 22 deletions(-) delete mode 100644 .changes/unreleased/Changed-20260801-190000.yaml create mode 100644 .changes/unreleased/Changed-20260924-190142.yaml diff --git a/.changes/unreleased/Changed-20260801-190000.yaml b/.changes/unreleased/Changed-20260801-190000.yaml deleted file mode 100644 index 52f05bd8..00000000 --- a/.changes/unreleased/Changed-20260801-190000.yaml +++ /dev/null @@ -1,9 +0,0 @@ -kind: Changed -body: |- - **Client response bodies that arrive in one frame are no longer copied.** - `collect_body_bounded` holds the first data frame by reference count and - returns it directly when no second frame follows, which is the common case - for unary responses and error bodies. Multi-frame bodies concatenate as - before, now starting from an exactly-sized buffer. The size limit is still - enforced before any frame is retained. -time: 2026-08-01T19:00:00.000000000Z diff --git a/.changes/unreleased/Changed-20260924-190142.yaml b/.changes/unreleased/Changed-20260924-190142.yaml new file mode 100644 index 00000000..6d5d08e9 --- /dev/null +++ b/.changes/unreleased/Changed-20260924-190142.yaml @@ -0,0 +1,14 @@ +kind: Changed +body: |- + **Connect unary and client-streaming responses that arrive in one frame + are no longer copied** ([#265]). A response body that arrives whole, which + is the common case, is now used as received instead of being copied into + a new buffer. With the proto codec, the response `OwnedView` can + therefore share memory with the transport's read buffer, and keeps that + buffer alive while it is held. `to_owned_message()` does not detach + `bytes` fields, which still point into that buffer, so copy those with + `Bytes::copy_from_slice` if you keep them long-term. Bodies of several + frames are concatenated as before. + + [#265]: https://github.com/connectrpc/connect-rust/pull/265 +time: 2026-09-24T19:01:42.699416714+00:00 diff --git a/connectrpc/src/client/mod.rs b/connectrpc/src/client/mod.rs index 526d6b5b..d4af361d 100644 --- a/connectrpc/src/client/mod.rs +++ b/connectrpc/src/client/mod.rs @@ -4963,6 +4963,12 @@ fn parse_grpc_error_from_trailers(trailers: &http::HeaderMap) -> Option( body: B, max_size: usize, @@ -5021,6 +5027,10 @@ where } } + debug_assert!( + head.is_none() || joined.is_none(), + "head and joined are never both set" + ); Ok(match (head, joined) { (_, Some(buf)) => buf.freeze(), (Some(first), None) => first, @@ -9829,12 +9839,13 @@ mod tests { } /// A one-frame body must be handed back by reference count, not copied. - /// This is the point of the single-frame path, and nothing else pins it. + /// This is the point of the single-frame path. `src` holds the second + /// reference, so the pointer comparison cannot pass by coincidence. #[tokio::test] async fn collect_body_bounded_single_frame_does_not_copy() { let src = Bytes::from(vec![9u8; 4096]); - let got = collect_body_bounded(Full::new(src.clone()), 8192) + let got = collect_body_bounded(Full::new(src.clone()), 8192, None) .await .unwrap(); @@ -9845,23 +9856,63 @@ mod tests { ); } - /// The promotion path still concatenates correctly once a second frame - /// arrives, and the result no longer aliases the first frame. + /// An empty frame ahead of the only data frame adds nothing, so the data + /// frame is still returned by reference count. #[tokio::test] - async fn collect_body_bounded_multi_frame_promotes() { + async fn collect_body_bounded_leading_empty_frame_does_not_copy() { let (tx, rx) = tokio::sync::mpsc::channel(4); let body = ChannelBody { rx }; - let first = Bytes::from_static(b"foo"); - tx.send(Ok(first.clone())).await.unwrap(); - tx.send(Ok(Bytes::from_static(b"bar"))).await.unwrap(); - tx.send(Ok(Bytes::from_static(b"baz"))).await.unwrap(); + let src = Bytes::from(vec![7u8; 64]); + tx.send(Ok(Bytes::new())).await.unwrap(); + tx.send(Ok(src.clone())).await.unwrap(); + drop(tx); + + let got = collect_body_bounded(body, 128, None).await.unwrap(); + assert_eq!(got, src); + assert!( + std::ptr::addr_eq(got.as_ptr(), src.as_ptr()), + "an empty frame must not force a copy" + ); + } + + /// An empty frame after the only data frame must not promote the body to + /// a buffer either. + #[tokio::test] + async fn collect_body_bounded_trailing_empty_frame_does_not_copy() { + let (tx, rx) = tokio::sync::mpsc::channel(4); + let body = ChannelBody { rx }; + let src = Bytes::from(vec![7u8; 64]); + tx.send(Ok(src.clone())).await.unwrap(); + tx.send(Ok(Bytes::new())).await.unwrap(); + drop(tx); + + let got = collect_body_bounded(body, 128, None).await.unwrap(); + assert_eq!(got, src); + assert!( + std::ptr::addr_eq(got.as_ptr(), src.as_ptr()), + "an empty frame must not force a copy" + ); + } + + /// A body of several frames is joined into a new buffer, even when the + /// frames are contiguous slices of one allocation, and a frame after the + /// second is appended to it. `base` stays alive, so the address + /// comparison cannot pass by coincidence. + #[tokio::test] + async fn collect_body_bounded_multi_frame_joins_into_a_new_buffer() { + let (tx, rx) = tokio::sync::mpsc::channel(4); + let body = ChannelBody { rx }; + let base = Bytes::from(b"foobarbaz".to_vec()); + for range in [0..3, 3..6, 6..9] { + tx.send(Ok(base.slice(range))).await.unwrap(); + } drop(tx); - let got = collect_body_bounded(body, 100).await.unwrap(); - assert_eq!(&got[..], b"foobarbaz"); + let got = collect_body_bounded(body, 100, None).await.unwrap(); + assert_eq!(got, base); assert!( - !std::ptr::addr_eq(got.as_ptr(), first.as_ptr()), - "a promoted body must be a fresh buffer" + !base[..].as_ptr_range().contains(&got.as_ptr()), + "a body of several frames must not alias its first frame" ); }