core: Coalesce Contiguous Small Buffers for ReadableBuffer (v1.83.x backport) (#12944)
Backport of #12924 to v1.83.x.
---
b/519106357
Co-authored-by: MV Shiva <speakupshiva@gmail.com>
diff --git a/core/src/main/java/io/grpc/internal/CompositeReadableBuffer.java b/core/src/main/java/io/grpc/internal/CompositeReadableBuffer.java
index 6cedb2c..6578db5 100644
--- a/core/src/main/java/io/grpc/internal/CompositeReadableBuffer.java
+++ b/core/src/main/java/io/grpc/internal/CompositeReadableBuffer.java
@@ -16,6 +16,7 @@
package io.grpc.internal;
+import com.google.common.annotations.VisibleForTesting;
import java.io.IOException;
import java.io.OutputStream;
import java.nio.ByteBuffer;
@@ -61,21 +62,78 @@
}
}
+ private static final int MIN_LARGE_BUFFER_SIZE = 1024;
+ private static final int MAX_SMALL_BUFFERS = 1000;
+
+ // Tracks the number of consecutive small buffers currently at the tail of the queue
+ private int tailSmallBufferCount = 0;
+
private void enqueueBuffer(ReadableBuffer buffer) {
- if (!(buffer instanceof CompositeReadableBuffer)) {
+ int bytes = buffer.readableBytes();
+
+ if (bytes >= MIN_LARGE_BUFFER_SIZE) {
+ // A large buffer arrived. Compact any preceding small buffers FIRST.
+ if (tailSmallBufferCount > 1) {
+ coalesceTailSmallBuffers(tailSmallBufferCount);
+ }
+ // Reset the counter and enqueue the large buffer.
+ // This strictly excludes the large buffer from any copying.
+ tailSmallBufferCount = 0;
readableBuffers.add(buffer);
- readableBytes += buffer.readableBytes();
+ readableBytes += bytes;
+ } else {
+ readableBuffers.add(buffer);
+ readableBytes += bytes;
+ tailSmallBufferCount++;
+
+ if (tailSmallBufferCount >= MAX_SMALL_BUFFERS) {
+ coalesceTailSmallBuffers(tailSmallBufferCount);
+
+ // Resetting to 0 ensures this newly coalesced chunk is NOT re-copied
+ // into the next batch. This restricts our time complexity to strictly O(N).
+ tailSmallBufferCount = 0;
+ }
+ }
+ }
+
+ private void coalesceTailSmallBuffers(int count) {
+ if (marked) {
return;
}
- CompositeReadableBuffer compositeBuffer = (CompositeReadableBuffer) buffer;
- while (!compositeBuffer.readableBuffers.isEmpty()) {
- ReadableBuffer subBuffer = compositeBuffer.readableBuffers.remove();
- readableBuffers.add(subBuffer);
+ // Extract ONLY the last 'count' buffers from the tail of the queue
+ ReadableBuffer[] toMerge = new ReadableBuffer[count];
+ int totalCoalescedBytes = 0;
+
+ // ArrayDeque.pollLast() retrieves elements in reverse order, so we populate backwards
+ for (int i = count - 1; i >= 0; i--) {
+ ReadableBuffer b = readableBuffers.pollLast();
+ toMerge[i] = b;
+ totalCoalescedBytes += b.readableBytes();
}
- readableBytes += compositeBuffer.readableBytes;
- compositeBuffer.readableBytes = 0;
- compositeBuffer.close();
+
+ byte[] coalescedBytes = new byte[totalCoalescedBytes];
+ int offset = 0;
+
+ for (int i = 0; i < count; i++) {
+ ReadableBuffer b = toMerge[i];
+ int len = b.readableBytes();
+ b.readBytes(coalescedBytes, offset, len);
+ offset += len;
+ b.close();
+ }
+
+ // Wrap and enqueue the single compacted buffer back at the tail
+ ReadableBuffer singleBuffer = ReadableBuffers.wrap(coalescedBytes);
+ readableBuffers.add(singleBuffer);
+
+ // Note: The global `readableBytes` remains perfectly synced since we
+ // subtracted and added the exact same amount of bytes.
+ }
+
+ @VisibleForTesting
+ int getBufferCount() {
+ return readableBuffers.size();
}
@Override
@@ -162,6 +220,7 @@
advanceBuffer();
} else {
readBuffer = readableBuffers.poll();
+ adjustTailSmallBufferCount();
}
length -= readable;
}
@@ -252,6 +311,7 @@
rewindableBuffers.remove().close();
}
}
+ tailSmallBufferCount = 0;
}
/**
@@ -315,6 +375,13 @@
} else {
readableBuffers.remove().close();
}
+ adjustTailSmallBufferCount();
+ }
+
+ private void adjustTailSmallBufferCount() {
+ if (tailSmallBufferCount > readableBuffers.size()) {
+ tailSmallBufferCount = readableBuffers.size();
+ }
}
/**
diff --git a/core/src/test/java/io/grpc/internal/CompositeReadableBufferTest.java b/core/src/test/java/io/grpc/internal/CompositeReadableBufferTest.java
index 749b71d..6d9ac3b 100644
--- a/core/src/test/java/io/grpc/internal/CompositeReadableBufferTest.java
+++ b/core/src/test/java/io/grpc/internal/CompositeReadableBufferTest.java
@@ -287,4 +287,174 @@
assertEquals(value.length(), composite.readableBytes());
}
+
+ @Test
+ public void coalesceOnMaxSmallBuffers() {
+ composite = new CompositeReadableBuffer();
+ // 1000 1-byte buffers
+ for (int i = 0; i < 1000; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+ assertEquals(1, composite.getBufferCount());
+ assertEquals(1000, composite.readableBytes());
+ }
+
+ @Test
+ public void coalesceBeyondMaxSmallBuffers() {
+ composite = new CompositeReadableBuffer();
+ // 1001 1-byte buffers
+ for (int i = 0; i < 1001; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+ assertEquals(2, composite.getBufferCount());
+ assertEquals(1001, composite.readableBytes());
+ }
+
+ @Test
+ public void coalesceMultipleBatchesOfSmallBuffers() {
+ composite = new CompositeReadableBuffer();
+ // 2000 1-byte buffers
+ for (int i = 0; i < 2000; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+ assertEquals(2, composite.getBufferCount());
+ assertEquals(2000, composite.readableBytes());
+ }
+
+ @Test
+ public void coalesceBeforeLargeBuffer() {
+ composite = new CompositeReadableBuffer();
+ // 500 small frames
+ for (int i = 0; i < 500; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+ // 1 large frame
+ composite.addBuffer(ReadableBuffers.wrap(new byte[1024]));
+
+ // The 500 small frames should be coalesced into 1, followed by the 1 large frame
+ assertEquals(2, composite.getBufferCount());
+ assertEquals(1524, composite.readableBytes());
+ }
+
+ @Test
+ public void largeBufferResetsTailSmallBufferCount() {
+ composite = new CompositeReadableBuffer();
+ // Add 1 large buffer
+ composite.addBuffer(ReadableBuffers.wrap(new byte[1024]));
+ assertEquals(1, composite.getBufferCount());
+
+ // Add 999 small buffers right after the large buffer (leaving total small tail at 999)
+ for (int i = 0; i < 999; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+
+ // Since only 999 small buffers are at the tail right after the large buffer,
+ // they must NOT be coalesced yet. Total buffers in queue should be 1 + 999 = 1000.
+ assertEquals(1000, composite.getBufferCount());
+ assertEquals(1024 + 999, composite.readableBytes());
+ }
+
+ @Test
+ public void noCoalesceOnLargeFrames() {
+ composite = new CompositeReadableBuffer();
+ // 1001 1024-byte frames
+ for (int i = 0; i < 1001; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[1024]));
+ }
+ assertEquals(1001, composite.getBufferCount());
+ assertEquals(1001 * 1024, composite.readableBytes());
+ }
+
+ @Test
+ public void skipCoalesceIfMarked() {
+ composite = new CompositeReadableBuffer();
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ composite.mark();
+
+ // Add 1000 more 1-byte buffers, reaching 1001 total
+ for (int i = 0; i < 1000; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+
+ // Should skip coalescing due to marked=true
+ assertEquals(1001, composite.getBufferCount());
+ assertEquals(1001, composite.readableBytes());
+ }
+
+ @Test
+ public void readBytesAdjustsTailSmallBufferCount() {
+ composite = new CompositeReadableBuffer();
+ // Add 500 1-byte buffers
+ for (int i = 0; i < 500; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+ // Read 200 buffers via readBytes(int) without mark
+ ReadableBuffer read = composite.readBytes(200);
+ read.close();
+
+ // Now add 500 more 1-byte buffers
+ for (int i = 0; i < 500; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+ assertEquals(800, composite.getBufferCount());
+ assertEquals(800, composite.readableBytes());
+ }
+
+ @Test
+ public void advanceBufferAdjustsTailSmallBufferCount() {
+ composite = new CompositeReadableBuffer();
+ // Add 500 1-byte buffers
+ for (int i = 0; i < 500; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+ // Read 200 buffers via skipBytes (which calls advanceBuffer)
+ composite.skipBytes(200);
+
+ // Now add 500 more 1-byte buffers
+ for (int i = 0; i < 500; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+ assertEquals(800, composite.getBufferCount());
+ assertEquals(800, composite.readableBytes());
+ }
+
+ @Test
+ public void coalesceClosesCoalescedBuffers() {
+ composite = new CompositeReadableBuffer();
+ ReadableBuffer mock1 = mock(ReadableBuffer.class);
+ when(mock1.readableBytes()).thenReturn(1);
+ ReadableBuffer mock2 = mock(ReadableBuffer.class);
+ when(mock2.readableBytes()).thenReturn(1);
+ composite.addBuffer(mock1);
+ composite.addBuffer(mock2);
+
+ // Large buffer triggers coalesce of mock1 and mock2 around line 76
+ composite.addBuffer(ReadableBuffers.wrap(new byte[1024]));
+
+ verify(mock1).close();
+ verify(mock2).close();
+ }
+
+ @Test
+ public void closeResetsTailSmallBufferCount() {
+ composite = new CompositeReadableBuffer();
+ // Add 500 small buffers
+ for (int i = 0; i < 500; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+ composite.close();
+ assertEquals(0, composite.getBufferCount());
+
+ // After close(), tailSmallBufferCount should be exactly 0.
+ // Adding 999 (MAX_SMALL_BUFFERS - 1) small buffers right after close should NOT trigger
+ // coalescing, resulting in exactly 999 buffers.
+ for (int i = 0; i < 999; i++) {
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ }
+ assertEquals(999, composite.getBufferCount());
+
+ // Adding the 1000th small buffer should trigger coalescing down to 1 buffer.
+ composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
+ assertEquals(1, composite.getBufferCount());
+ }
}