Skip to content

Commit dd9f8eb

Browse files
core: Coalesce Contiguous Small Buffers for ReadableBuffer (v1.82.x backport) (#12943)
Backport of #12924 to v1.82.x. --- b/519106357 Co-authored-by: MV Shiva <speakupshiva@gmail.com>
1 parent 96b6849 commit dd9f8eb

2 files changed

Lines changed: 246 additions & 9 deletions

File tree

core/src/main/java/io/grpc/internal/CompositeReadableBuffer.java

Lines changed: 76 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616

1717
package io.grpc.internal;
1818

19+
import com.google.common.annotations.VisibleForTesting;
1920
import java.io.IOException;
2021
import java.io.OutputStream;
2122
import java.nio.ByteBuffer;
@@ -61,21 +62,78 @@ public void addBuffer(ReadableBuffer buffer) {
6162
}
6263
}
6364

65+
private static final int MIN_LARGE_BUFFER_SIZE = 1024;
66+
private static final int MAX_SMALL_BUFFERS = 1000;
67+
68+
// Tracks the number of consecutive small buffers currently at the tail of the queue
69+
private int tailSmallBufferCount = 0;
70+
6471
private void enqueueBuffer(ReadableBuffer buffer) {
65-
if (!(buffer instanceof CompositeReadableBuffer)) {
72+
int bytes = buffer.readableBytes();
73+
74+
if (bytes >= MIN_LARGE_BUFFER_SIZE) {
75+
// A large buffer arrived. Compact any preceding small buffers FIRST.
76+
if (tailSmallBufferCount > 1) {
77+
coalesceTailSmallBuffers(tailSmallBufferCount);
78+
}
79+
// Reset the counter and enqueue the large buffer.
80+
// This strictly excludes the large buffer from any copying.
81+
tailSmallBufferCount = 0;
6682
readableBuffers.add(buffer);
67-
readableBytes += buffer.readableBytes();
83+
readableBytes += bytes;
84+
} else {
85+
readableBuffers.add(buffer);
86+
readableBytes += bytes;
87+
tailSmallBufferCount++;
88+
89+
if (tailSmallBufferCount >= MAX_SMALL_BUFFERS) {
90+
coalesceTailSmallBuffers(tailSmallBufferCount);
91+
92+
// Resetting to 0 ensures this newly coalesced chunk is NOT re-copied
93+
// into the next batch. This restricts our time complexity to strictly O(N).
94+
tailSmallBufferCount = 0;
95+
}
96+
}
97+
}
98+
99+
private void coalesceTailSmallBuffers(int count) {
100+
if (marked) {
68101
return;
69102
}
70103

71-
CompositeReadableBuffer compositeBuffer = (CompositeReadableBuffer) buffer;
72-
while (!compositeBuffer.readableBuffers.isEmpty()) {
73-
ReadableBuffer subBuffer = compositeBuffer.readableBuffers.remove();
74-
readableBuffers.add(subBuffer);
104+
// Extract ONLY the last 'count' buffers from the tail of the queue
105+
ReadableBuffer[] toMerge = new ReadableBuffer[count];
106+
int totalCoalescedBytes = 0;
107+
108+
// ArrayDeque.pollLast() retrieves elements in reverse order, so we populate backwards
109+
for (int i = count - 1; i >= 0; i--) {
110+
ReadableBuffer b = readableBuffers.pollLast();
111+
toMerge[i] = b;
112+
totalCoalescedBytes += b.readableBytes();
113+
}
114+
115+
byte[] coalescedBytes = new byte[totalCoalescedBytes];
116+
int offset = 0;
117+
118+
for (int i = 0; i < count; i++) {
119+
ReadableBuffer b = toMerge[i];
120+
int len = b.readableBytes();
121+
b.readBytes(coalescedBytes, offset, len);
122+
offset += len;
123+
b.close();
75124
}
76-
readableBytes += compositeBuffer.readableBytes;
77-
compositeBuffer.readableBytes = 0;
78-
compositeBuffer.close();
125+
126+
// Wrap and enqueue the single compacted buffer back at the tail
127+
ReadableBuffer singleBuffer = ReadableBuffers.wrap(coalescedBytes);
128+
readableBuffers.add(singleBuffer);
129+
130+
// Note: The global `readableBytes` remains perfectly synced since we
131+
// subtracted and added the exact same amount of bytes.
132+
}
133+
134+
@VisibleForTesting
135+
int getBufferCount() {
136+
return readableBuffers.size();
79137
}
80138

81139
@Override
@@ -162,6 +220,7 @@ public ReadableBuffer readBytes(int length) {
162220
advanceBuffer();
163221
} else {
164222
readBuffer = readableBuffers.poll();
223+
adjustTailSmallBufferCount();
165224
}
166225
length -= readable;
167226
}
@@ -252,6 +311,7 @@ public void close() {
252311
rewindableBuffers.remove().close();
253312
}
254313
}
314+
tailSmallBufferCount = 0;
255315
}
256316

257317
/**
@@ -315,6 +375,13 @@ private void advanceBuffer() {
315375
} else {
316376
readableBuffers.remove().close();
317377
}
378+
adjustTailSmallBufferCount();
379+
}
380+
381+
private void adjustTailSmallBufferCount() {
382+
if (tailSmallBufferCount > readableBuffers.size()) {
383+
tailSmallBufferCount = readableBuffers.size();
384+
}
318385
}
319386

320387
/**

core/src/test/java/io/grpc/internal/CompositeReadableBufferTest.java

Lines changed: 170 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -287,4 +287,174 @@ private void splitAndAdd(String value) {
287287

288288
assertEquals(value.length(), composite.readableBytes());
289289
}
290+
291+
@Test
292+
public void coalesceOnMaxSmallBuffers() {
293+
composite = new CompositeReadableBuffer();
294+
// 1000 1-byte buffers
295+
for (int i = 0; i < 1000; i++) {
296+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
297+
}
298+
assertEquals(1, composite.getBufferCount());
299+
assertEquals(1000, composite.readableBytes());
300+
}
301+
302+
@Test
303+
public void coalesceBeyondMaxSmallBuffers() {
304+
composite = new CompositeReadableBuffer();
305+
// 1001 1-byte buffers
306+
for (int i = 0; i < 1001; i++) {
307+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
308+
}
309+
assertEquals(2, composite.getBufferCount());
310+
assertEquals(1001, composite.readableBytes());
311+
}
312+
313+
@Test
314+
public void coalesceMultipleBatchesOfSmallBuffers() {
315+
composite = new CompositeReadableBuffer();
316+
// 2000 1-byte buffers
317+
for (int i = 0; i < 2000; i++) {
318+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
319+
}
320+
assertEquals(2, composite.getBufferCount());
321+
assertEquals(2000, composite.readableBytes());
322+
}
323+
324+
@Test
325+
public void coalesceBeforeLargeBuffer() {
326+
composite = new CompositeReadableBuffer();
327+
// 500 small frames
328+
for (int i = 0; i < 500; i++) {
329+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
330+
}
331+
// 1 large frame
332+
composite.addBuffer(ReadableBuffers.wrap(new byte[1024]));
333+
334+
// The 500 small frames should be coalesced into 1, followed by the 1 large frame
335+
assertEquals(2, composite.getBufferCount());
336+
assertEquals(1524, composite.readableBytes());
337+
}
338+
339+
@Test
340+
public void largeBufferResetsTailSmallBufferCount() {
341+
composite = new CompositeReadableBuffer();
342+
// Add 1 large buffer
343+
composite.addBuffer(ReadableBuffers.wrap(new byte[1024]));
344+
assertEquals(1, composite.getBufferCount());
345+
346+
// Add 999 small buffers right after the large buffer (leaving total small tail at 999)
347+
for (int i = 0; i < 999; i++) {
348+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
349+
}
350+
351+
// Since only 999 small buffers are at the tail right after the large buffer,
352+
// they must NOT be coalesced yet. Total buffers in queue should be 1 + 999 = 1000.
353+
assertEquals(1000, composite.getBufferCount());
354+
assertEquals(1024 + 999, composite.readableBytes());
355+
}
356+
357+
@Test
358+
public void noCoalesceOnLargeFrames() {
359+
composite = new CompositeReadableBuffer();
360+
// 1001 1024-byte frames
361+
for (int i = 0; i < 1001; i++) {
362+
composite.addBuffer(ReadableBuffers.wrap(new byte[1024]));
363+
}
364+
assertEquals(1001, composite.getBufferCount());
365+
assertEquals(1001 * 1024, composite.readableBytes());
366+
}
367+
368+
@Test
369+
public void skipCoalesceIfMarked() {
370+
composite = new CompositeReadableBuffer();
371+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
372+
composite.mark();
373+
374+
// Add 1000 more 1-byte buffers, reaching 1001 total
375+
for (int i = 0; i < 1000; i++) {
376+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
377+
}
378+
379+
// Should skip coalescing due to marked=true
380+
assertEquals(1001, composite.getBufferCount());
381+
assertEquals(1001, composite.readableBytes());
382+
}
383+
384+
@Test
385+
public void readBytesAdjustsTailSmallBufferCount() {
386+
composite = new CompositeReadableBuffer();
387+
// Add 500 1-byte buffers
388+
for (int i = 0; i < 500; i++) {
389+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
390+
}
391+
// Read 200 buffers via readBytes(int) without mark
392+
ReadableBuffer read = composite.readBytes(200);
393+
read.close();
394+
395+
// Now add 500 more 1-byte buffers
396+
for (int i = 0; i < 500; i++) {
397+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
398+
}
399+
assertEquals(800, composite.getBufferCount());
400+
assertEquals(800, composite.readableBytes());
401+
}
402+
403+
@Test
404+
public void advanceBufferAdjustsTailSmallBufferCount() {
405+
composite = new CompositeReadableBuffer();
406+
// Add 500 1-byte buffers
407+
for (int i = 0; i < 500; i++) {
408+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
409+
}
410+
// Read 200 buffers via skipBytes (which calls advanceBuffer)
411+
composite.skipBytes(200);
412+
413+
// Now add 500 more 1-byte buffers
414+
for (int i = 0; i < 500; i++) {
415+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
416+
}
417+
assertEquals(800, composite.getBufferCount());
418+
assertEquals(800, composite.readableBytes());
419+
}
420+
421+
@Test
422+
public void coalesceClosesCoalescedBuffers() {
423+
composite = new CompositeReadableBuffer();
424+
ReadableBuffer mock1 = mock(ReadableBuffer.class);
425+
when(mock1.readableBytes()).thenReturn(1);
426+
ReadableBuffer mock2 = mock(ReadableBuffer.class);
427+
when(mock2.readableBytes()).thenReturn(1);
428+
composite.addBuffer(mock1);
429+
composite.addBuffer(mock2);
430+
431+
// Large buffer triggers coalesce of mock1 and mock2 around line 76
432+
composite.addBuffer(ReadableBuffers.wrap(new byte[1024]));
433+
434+
verify(mock1).close();
435+
verify(mock2).close();
436+
}
437+
438+
@Test
439+
public void closeResetsTailSmallBufferCount() {
440+
composite = new CompositeReadableBuffer();
441+
// Add 500 small buffers
442+
for (int i = 0; i < 500; i++) {
443+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
444+
}
445+
composite.close();
446+
assertEquals(0, composite.getBufferCount());
447+
448+
// After close(), tailSmallBufferCount should be exactly 0.
449+
// Adding 999 (MAX_SMALL_BUFFERS - 1) small buffers right after close should NOT trigger
450+
// coalescing, resulting in exactly 999 buffers.
451+
for (int i = 0; i < 999; i++) {
452+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
453+
}
454+
assertEquals(999, composite.getBufferCount());
455+
456+
// Adding the 1000th small buffer should trigger coalescing down to 1 buffer.
457+
composite.addBuffer(ReadableBuffers.wrap(new byte[] {1}));
458+
assertEquals(1, composite.getBufferCount());
459+
}
290460
}

0 commit comments

Comments
 (0)