• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

grpc / grpc-java / #20374

29 Jul 2026 10:45AM UTC coverage: 88.837% (+0.003%) from 88.834%
#20374

push

github

web-flow
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>

36288 of 40848 relevant lines covered (88.84%)

0.89 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

97.01
/../core/src/main/java/io/grpc/internal/CompositeReadableBuffer.java
1
/*
2
 * Copyright 2014 The gRPC Authors
3
 *
4
 * Licensed under the Apache License, Version 2.0 (the "License");
5
 * you may not use this file except in compliance with the License.
6
 * You may obtain a copy of the License at
7
 *
8
 *     http://www.apache.org/licenses/LICENSE-2.0
9
 *
10
 * Unless required by applicable law or agreed to in writing, software
11
 * distributed under the License is distributed on an "AS IS" BASIS,
12
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13
 * See the License for the specific language governing permissions and
14
 * limitations under the License.
15
 */
16

17
package io.grpc.internal;
18

19
import com.google.common.annotations.VisibleForTesting;
20
import java.io.IOException;
21
import java.io.OutputStream;
22
import java.nio.ByteBuffer;
23
import java.nio.InvalidMarkException;
24
import java.util.ArrayDeque;
25
import java.util.Deque;
26
import javax.annotation.Nullable;
27

28
/**
29
 * A {@link ReadableBuffer} that is composed of 0 or more {@link ReadableBuffer}s. This provides a
30
 * facade that allows multiple buffers to be treated as one.
31
 *
32
 * <p>When a buffer is added to a composite, its life cycle is controlled by the composite. Once
33
 * the composite has read past the end of a given buffer, that buffer is automatically closed and
34
 * removed from the composite.
35
 */
36
public class CompositeReadableBuffer extends AbstractReadableBuffer {
37

38
  private final Deque<ReadableBuffer> readableBuffers;
39
  private Deque<ReadableBuffer> rewindableBuffers;
40
  private int readableBytes;
41
  private boolean marked;
42

43
  public CompositeReadableBuffer(int initialCapacity) {
1✔
44
    readableBuffers = new ArrayDeque<>(initialCapacity);
1✔
45
  }
1✔
46

47
  public CompositeReadableBuffer() {
1✔
48
    readableBuffers = new ArrayDeque<>();
1✔
49
  }
1✔
50

51
  /**
52
   * Adds a new {@link ReadableBuffer} at the end of the buffer list. After a buffer is added, it is
53
   * expected that this {@code CompositeBuffer} has complete ownership. Any attempt to modify the
54
   * buffer (i.e. modifying the readable bytes) may result in corruption of the internal state of
55
   * this {@code CompositeBuffer}.
56
   */
57
  public void addBuffer(ReadableBuffer buffer) {
58
    boolean markHead = marked && readableBuffers.isEmpty();
1✔
59
    enqueueBuffer(buffer);
1✔
60
    if (markHead) {
1✔
61
      readableBuffers.peek().mark();
1✔
62
    }
63
  }
1✔
64

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;
1✔
70

71
  private void enqueueBuffer(ReadableBuffer buffer) {
72
    int bytes = buffer.readableBytes();
1✔
73

74
    if (bytes >= MIN_LARGE_BUFFER_SIZE) {
1✔
75
      // A large buffer arrived. Compact any preceding small buffers FIRST.
76
      if (tailSmallBufferCount > 1) {
1✔
77
        coalesceTailSmallBuffers(tailSmallBufferCount);
1✔
78
      }
79
      // Reset the counter and enqueue the large buffer.
80
      // This strictly excludes the large buffer from any copying.
81
      tailSmallBufferCount = 0;
1✔
82
      readableBuffers.add(buffer);
1✔
83
      readableBytes += bytes;
1✔
84
    } else {
85
      readableBuffers.add(buffer);
1✔
86
      readableBytes += bytes;
1✔
87
      tailSmallBufferCount++;
1✔
88

89
      if (tailSmallBufferCount >= MAX_SMALL_BUFFERS) {
1✔
90
        coalesceTailSmallBuffers(tailSmallBufferCount);
1✔
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;
1✔
95
      }
96
    }
97
  }
1✔
98

99
  private void coalesceTailSmallBuffers(int count) {
100
    if (marked) {
1✔
101
      return;
1✔
102
    }
103

104
    // Extract ONLY the last 'count' buffers from the tail of the queue
105
    ReadableBuffer[] toMerge = new ReadableBuffer[count];
1✔
106
    int totalCoalescedBytes = 0;
1✔
107

108
    // ArrayDeque.pollLast() retrieves elements in reverse order, so we populate backwards
109
    for (int i = count - 1; i >= 0; i--) {
1✔
110
      ReadableBuffer b = readableBuffers.pollLast();
1✔
111
      toMerge[i] = b;
1✔
112
      totalCoalescedBytes += b.readableBytes();
1✔
113
    }
114

115
    byte[] coalescedBytes = new byte[totalCoalescedBytes];
1✔
116
    int offset = 0;
1✔
117

118
    for (int i = 0; i < count; i++) {
1✔
119
      ReadableBuffer b = toMerge[i];
1✔
120
      int len = b.readableBytes();
1✔
121
      b.readBytes(coalescedBytes, offset, len);
1✔
122
      offset += len;
1✔
123
      b.close();
1✔
124
    }
125

126
    // Wrap and enqueue the single compacted buffer back at the tail
127
    ReadableBuffer singleBuffer = ReadableBuffers.wrap(coalescedBytes);
1✔
128
    readableBuffers.add(singleBuffer);
1✔
129

130
    // Note: The global `readableBytes` remains perfectly synced since we
131
    // subtracted and added the exact same amount of bytes.
132
  }
1✔
133

134
  @VisibleForTesting
135
  int getBufferCount() {
136
    return readableBuffers.size();
1✔
137
  }
138

139
  @Override
140
  public int readableBytes() {
141
    return readableBytes;
1✔
142
  }
143

144
  private static final NoThrowReadOperation<Void> UBYTE_OP =
1✔
145
      new NoThrowReadOperation<Void>() {
1✔
146
        @Override
147
        public int read(ReadableBuffer buffer, int length, Void unused, int value) {
148
          return buffer.readUnsignedByte();
1✔
149
        }
150
      };
151

152
  @Override
153
  public int readUnsignedByte() {
154
    return executeNoThrow(UBYTE_OP, 1, null, 0);
1✔
155
  }
156

157
  private static final NoThrowReadOperation<Void> SKIP_OP =
1✔
158
      new NoThrowReadOperation<Void>() {
1✔
159
        @Override
160
        public int read(ReadableBuffer buffer, int length, Void unused, int unused2) {
161
          buffer.skipBytes(length);
1✔
162
          return 0;
1✔
163
        }
164
      };
165

166
  @Override
167
  public void skipBytes(int length) {
168
    executeNoThrow(SKIP_OP, length, null, 0);
1✔
169
  }
1✔
170

171
  private static final NoThrowReadOperation<byte[]> BYTE_ARRAY_OP =
1✔
172
      new NoThrowReadOperation<byte[]>() {
1✔
173
        @Override
174
        public int read(ReadableBuffer buffer, int length, byte[] dest, int offset) {
175
          buffer.readBytes(dest, offset, length);
1✔
176
          return offset + length;
1✔
177
        }
178
      };
179

180
  private static final ReadOperation<OutputStream> STREAM_OP =
1✔
181
      new ReadOperation<OutputStream>() {
1✔
182
        @Override
183
        public int read(ReadableBuffer buffer, int length, OutputStream dest, int unused)
184
            throws IOException {
185
          buffer.readBytes(dest, length);
1✔
186
          return 0;
1✔
187
        }
188
      };
189

190
  @Override
191
  public void readBytes(byte[] dest, int destOffset, int length) {
192
    executeNoThrow(BYTE_ARRAY_OP, length, dest, destOffset);
1✔
193
  }
1✔
194

195
  @Override
196
  public void readBytes(OutputStream dest, int length) throws IOException {
197
    execute(STREAM_OP, length, dest, 0);
1✔
198
  }
1✔
199

200
  @Override
201
  public ReadableBuffer readBytes(int length) {
202
    if (length <= 0) {
1✔
203
      return ReadableBuffers.empty();
×
204
    }
205
    checkReadable(length);
1✔
206
    readableBytes -= length;
1✔
207

208
    ReadableBuffer newBuffer = null;
1✔
209
    CompositeReadableBuffer newComposite = null;
1✔
210
    do {
211
      ReadableBuffer buffer = readableBuffers.peek();
1✔
212
      int readable = buffer.readableBytes();
1✔
213
      ReadableBuffer readBuffer;
214
      if (readable > length) {
1✔
215
        readBuffer = buffer.readBytes(length);
1✔
216
        length = 0;
1✔
217
      } else {
218
        if (marked) {
1✔
219
          readBuffer = buffer.readBytes(readable);
1✔
220
          advanceBuffer();
1✔
221
        } else {
222
          readBuffer = readableBuffers.poll();
1✔
223
          adjustTailSmallBufferCount();
1✔
224
        }
225
        length -= readable;
1✔
226
      }
227
      if (newBuffer == null) {
1✔
228
        newBuffer = readBuffer;
1✔
229
      } else {
230
        if (newComposite == null) {
1✔
231
          newComposite = new CompositeReadableBuffer(
1✔
232
              length == 0 ? 2 : Math.min(readableBuffers.size() + 2, 16));
1✔
233
          newComposite.addBuffer(newBuffer);
1✔
234
          newBuffer = newComposite;
1✔
235
        }
236
        newComposite.addBuffer(readBuffer);
1✔
237
      }
238
    } while (length > 0);
1✔
239
    return newBuffer;
1✔
240
  }
241

242
  @Override
243
  public boolean markSupported() {
244
    for (ReadableBuffer buffer : readableBuffers) {
1✔
245
      if (!buffer.markSupported()) {
1✔
246
        return false;
1✔
247
      }
248
    }
1✔
249
    return true;
1✔
250
  }
251

252
  @Override
253
  public void mark() {
254
    if (rewindableBuffers == null) {
1✔
255
      rewindableBuffers = new ArrayDeque<>(Math.min(readableBuffers.size(), 16));
1✔
256
    }
257
    while (!rewindableBuffers.isEmpty()) {
1✔
258
      rewindableBuffers.remove().close();
1✔
259
    }
260
    marked = true;
1✔
261
    ReadableBuffer buffer = readableBuffers.peek();
1✔
262
    if (buffer != null) {
1✔
263
      buffer.mark();
1✔
264
    }
265
  }
1✔
266

267
  @Override
268
  public void reset() {
269
    if (!marked) {
1✔
270
      throw new InvalidMarkException();
1✔
271
    }
272
    ReadableBuffer buffer;
273
    if ((buffer = readableBuffers.peek()) != null) {
1✔
274
      int currentRemain = buffer.readableBytes();
1✔
275
      buffer.reset();
1✔
276
      readableBytes += (buffer.readableBytes() - currentRemain);
1✔
277
    }
278
    while ((buffer = rewindableBuffers.pollLast()) != null) {
1✔
279
      buffer.reset();
1✔
280
      readableBuffers.addFirst(buffer);
1✔
281
      readableBytes += buffer.readableBytes();
1✔
282
    }
283
  }
1✔
284

285
  @Override
286
  public boolean byteBufferSupported() {
287
    for (ReadableBuffer buffer : readableBuffers) {
1✔
288
      if (!buffer.byteBufferSupported()) {
1✔
289
        return false;
1✔
290
      }
291
    }
1✔
292
    return true;
1✔
293
  }
294

295
  @Nullable
296
  @Override
297
  public ByteBuffer getByteBuffer() {
298
    if (readableBuffers.isEmpty()) {
1✔
299
      return null;
×
300
    }
301
    return readableBuffers.peek().getByteBuffer();
1✔
302
  }
303

304
  @Override
305
  public void close() {
306
    while (!readableBuffers.isEmpty()) {
1✔
307
      readableBuffers.remove().close();
1✔
308
    }
309
    if (rewindableBuffers != null) {
1✔
310
      while (!rewindableBuffers.isEmpty()) {
1✔
311
        rewindableBuffers.remove().close();
1✔
312
      }
313
    }
314
    tailSmallBufferCount = 0;
1✔
315
  }
1✔
316

317
  /**
318
   * Executes the given {@link ReadOperation} against the {@link ReadableBuffer}s required to
319
   * satisfy the requested {@code length}.
320
   */
321
  private <T> int execute(ReadOperation<T> op, int length, T dest, int value) throws IOException {
322
    checkReadable(length);
1✔
323

324
    if (!readableBuffers.isEmpty()) {
1✔
325
      advanceBufferIfNecessary();
1✔
326
    }
327

328
    for (; length > 0 && !readableBuffers.isEmpty(); advanceBufferIfNecessary()) {
1✔
329
      ReadableBuffer buffer = readableBuffers.peek();
1✔
330
      int lengthToCopy = Math.min(length, buffer.readableBytes());
1✔
331

332
      // Perform the read operation for this buffer.
333
      value = op.read(buffer, lengthToCopy, dest, value);
1✔
334

335
      length -= lengthToCopy;
1✔
336
      readableBytes -= lengthToCopy;
1✔
337
    }
338

339
    if (length > 0) {
1✔
340
      // Should never get here.
341
      throw new AssertionError("Failed executing read operation");
×
342
    }
343

344
    return value;
1✔
345
  }
346

347
  private <T> int executeNoThrow(NoThrowReadOperation<T> op, int length, T dest, int value) {
348
    try {
349
      return execute(op, length, dest, value);
1✔
350
    } catch (IOException e) {
×
351
      throw new AssertionError(e); // shouldn't happen
×
352
    }
353
  }
354

355
  /**
356
   * If the current buffer is exhausted, removes and closes it.
357
   */
358
  private void advanceBufferIfNecessary() {
359
    ReadableBuffer buffer = readableBuffers.peek();
1✔
360
    if (buffer.readableBytes() == 0) {
1✔
361
      advanceBuffer();
1✔
362
    }
363
  }
1✔
364

365
  /**
366
   * Removes one buffer from the front and closes it.
367
   */
368
  private void advanceBuffer() {
369
    if (marked) {
1✔
370
      rewindableBuffers.add(readableBuffers.remove());
1✔
371
      ReadableBuffer next = readableBuffers.peek();
1✔
372
      if (next != null) {
1✔
373
        next.mark();
1✔
374
      }
375
    } else {
1✔
376
      readableBuffers.remove().close();
1✔
377
    }
378
    adjustTailSmallBufferCount();
1✔
379
  }
1✔
380

381
  private void adjustTailSmallBufferCount() {
382
    if (tailSmallBufferCount > readableBuffers.size()) {
1✔
383
      tailSmallBufferCount = readableBuffers.size();
1✔
384
    }
385
  }
1✔
386

387
  /**
388
   * A simple read operation to perform on a single {@link ReadableBuffer}.
389
   * All state management for the buffers is done by
390
   * {@link CompositeReadableBuffer#execute(ReadOperation, int, Object, int)}.
391
   */
392
  private interface ReadOperation<T> {
393
    /**
394
     * This method can also be used to simultaneously perform operation-specific int-valued
395
     * aggregation over the sequence of buffers in a {@link CompositeReadableBuffer}.
396
     * {@code value} is the return value from the prior buffer, or the "initial" value passed
397
     * to {@code execute()} in the case of the first buffer. {@code execute()} returns the value
398
     * returned by the operation called on the last buffer.
399
     */
400
    int read(ReadableBuffer buffer, int length, T dest, int value) throws IOException;
401
  }
402

403
  private interface NoThrowReadOperation<T> extends ReadOperation<T> {
404
    @Override
405
    int read(ReadableBuffer buffer, int length, T dest, int value);
406
  }
407
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc