Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions .changes/unreleased/Changed-20260924-190142.yaml
Original file line number Diff line number Diff line change
@@ -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
133 changes: 123 additions & 10 deletions connectrpc/src/client/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4963,6 +4963,12 @@ fn parse_grpc_error_from_trailers(trailers: &http::HeaderMap) -> Option<ConnectE
/// Returns `ResourceExhausted` if the accumulated data exceeds `max_size`,
/// `DeadlineExceeded` if the body fails after the deadline has passed, and
/// `Internal` if it fails before.
///
/// A body with one non-empty data frame is returned as that frame, without
/// a copy, so the result may be a slice of the transport's read buffer and
/// keeps that whole buffer alive while it is held. A body with several
/// non-empty frames is concatenated into a new buffer that does not alias
/// the transport.
async fn collect_body_bounded<B>(
body: B,
max_size: usize,
Expand All @@ -4972,22 +4978,42 @@ where
B: Body<Data = Bytes>,
B::Error: Into<Box<dyn std::error::Error + Send + Sync>>,
{
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<Bytes> = None;
let mut joined: Option<BytesMut> = 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)) => {
Expand All @@ -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.
Expand Down Expand Up @@ -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);
Expand Down
Loading