From 9fcc71645fcf76874f3c662d1ae81b1c771b1e1f Mon Sep 17 00:00:00 2001 From: dongjiang Date: Wed, 15 Jul 2026 15:27:02 +0800 Subject: [PATCH 1/2] =?UTF-8?q?Avoid=20O(n=C2=B2)=20buffer=20copies=20in?= =?UTF-8?q?=20EnvelopeReader=20streaming=20path?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Use in-place `del buffer[:n]` instead of `buffer = buffer[n:]` to avoid allocating a new bytearray and copying all remaining bytes on every message. For long-running streaming RPCs with N messages, this reduces total bytes copied from O(N²) to O(N), yielding up to 66x throughput improvement in single-feed scenarios. Signed-off-by: dongjiang --- src/connectrpc/_envelope.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/src/connectrpc/_envelope.py b/src/connectrpc/_envelope.py index 5ee8a353..2f1db6e0 100644 --- a/src/connectrpc/_envelope.py +++ b/src/connectrpc/_envelope.py @@ -53,7 +53,10 @@ def _read_messages(self) -> Iterator[_RES]: compressed = prefix_byte & 0b01 != 0 message_data = self._buffer[5 : 5 + self._next_message_length] - self._buffer = self._buffer[5 + self._next_message_length :] + # Use in-place deletion to avoid O(n²) memory copies in streaming RPCs. + # `del buf[:n]` uses memmove within the existing buffer, while + # `buf = buf[n:]` allocates a new bytearray and copies all remaining bytes. + del self._buffer[: 5 + self._next_message_length] self._next_message_length = None if compressed: if isinstance(self._compression, IdentityCompression): From b1a864bb9f6f0292dd71f5cd877ffa8055f29aef Mon Sep 17 00:00:00 2001 From: "Anuraag (Rag) Agrawal" Date: Wed, 15 Jul 2026 16:58:16 +0900 Subject: [PATCH 2/2] Remove comments on buffer memory optimization Remove comments explaining in-place deletion for buffer. Signed-off-by: Anuraag (Rag) Agrawal --- src/connectrpc/_envelope.py | 3 --- 1 file changed, 3 deletions(-) diff --git a/src/connectrpc/_envelope.py b/src/connectrpc/_envelope.py index 2f1db6e0..7867b201 100644 --- a/src/connectrpc/_envelope.py +++ b/src/connectrpc/_envelope.py @@ -53,9 +53,6 @@ def _read_messages(self) -> Iterator[_RES]: compressed = prefix_byte & 0b01 != 0 message_data = self._buffer[5 : 5 + self._next_message_length] - # Use in-place deletion to avoid O(n²) memory copies in streaming RPCs. - # `del buf[:n]` uses memmove within the existing buffer, while - # `buf = buf[n:]` allocates a new bytearray and copies all remaining bytes. del self._buffer[: 5 + self._next_message_length] self._next_message_length = None if compressed: