diff --git a/.changes/unreleased/Changed-20260924-190142.yaml b/.changes/unreleased/Changed-20260924-190142.yaml new file mode 100644 index 0000000..6d5d08e --- /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 ac5e090..d4af361 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, @@ -4972,22 +4978,42 @@ where B: Body, B::Error: Into>, { - 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)) => { @@ -5000,7 +5026,16 @@ where None => break, } } - Ok(buf.freeze()) + + 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, + (None, None) => Bytes::new(), + }) } /// Percent-decode a gRPC message string. @@ -9803,6 +9838,84 @@ 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. `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, None) + .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" + ); + } + + /// 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_leading_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(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, None).await.unwrap(); + assert_eq!(got, base); + assert!( + !base[..].as_ptr_range().contains(&got.as_ptr()), + "a body of several frames must not alias its first frame" + ); + } + #[tokio::test] async fn collect_body_bounded_propagates_body_error() { let (tx, rx) = tokio::sync::mpsc::channel(4);