Skip to content

Commit f5b0081

Browse files
committed
Reset message budget for compressed frames
1 parent ee8852d commit f5b0081

2 files changed

Lines changed: 43 additions & 0 deletions

File tree

‎iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TCompressedElasticFramedTransport.java‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,10 +59,14 @@ protected void closeAllocatedBuffers() {
5959

6060
@Override
6161
protected void readFrame() throws TTransportException {
62+
// Discard the previous frame's budget before reading the next frame header and payload.
63+
resetMessageSizeAndConsumedBytes();
6264
underlying.readAll(i32buf, 0, 4);
6365
int size = TFramedTransport.decodeFrameSize(i32buf);
6466
validateFrame(size);
6567
readBuffer.fill(underlying, size);
68+
// Bind subsequent protocol reads to the current compressed frame size.
69+
resetMessageSizeAndConsumedBytes(size);
6670
RpcStat.readCompressedBytes.addAndGet(size);
6771
try {
6872
int uncompressedLength = uncompressedLength(readBuffer.getBuffer(), 0, size);

‎iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java‎

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727

2828
import java.nio.ByteBuffer;
2929
import java.nio.charset.StandardCharsets;
30+
import java.util.Arrays;
3031

3132
import static org.junit.Assert.assertArrayEquals;
3233
import static org.junit.Assert.assertEquals;
@@ -61,6 +62,44 @@ public void testReadFramesResetMessageSize() throws TTransportException {
6162
assertArrayEquals(secondFrame, actualSecondFrame);
6263
}
6364

65+
@Test
66+
public void testReadCompressedFramesResetMessageSize() throws TTransportException {
67+
byte[] firstFrame = {1, 2};
68+
byte[] secondFrame = {3, 4, 5, 6, 7, 8};
69+
TMemoryBuffer wire = new TMemoryBuffer(128);
70+
TSnappyElasticFramedTransport output =
71+
new TSnappyElasticFramedTransport(wire, 4, 128, false);
72+
output.write(firstFrame);
73+
output.flush();
74+
output.write(secondFrame);
75+
output.flush();
76+
77+
byte[] framedData = Arrays.copyOf(wire.getArray(), wire.length());
78+
ByteBuffer frameSizes = ByteBuffer.wrap(framedData);
79+
int firstCompressedSize = frameSizes.getInt();
80+
frameSizes.position(Integer.BYTES + firstCompressedSize);
81+
int secondCompressedSize = frameSizes.getInt();
82+
int maxMessageSize = Integer.BYTES + Math.max(firstCompressedSize, secondCompressedSize);
83+
TConfiguration configuration =
84+
TConfiguration.custom()
85+
.setMaxMessageSize(maxMessageSize)
86+
.setMaxFrameSize(maxMessageSize)
87+
.build();
88+
TMemoryBuffer underlying = new TMemoryBuffer(configuration, maxMessageSize);
89+
underlying.write(framedData);
90+
TSnappyElasticFramedTransport input =
91+
new TSnappyElasticFramedTransport(
92+
underlying, 4, configuration.getMaxFrameSize(), false);
93+
94+
byte[] actualFirstFrame = new byte[firstFrame.length];
95+
input.readAll(actualFirstFrame, 0, actualFirstFrame.length);
96+
assertArrayEquals(firstFrame, actualFirstFrame);
97+
98+
byte[] actualSecondFrame = new byte[secondFrame.length];
99+
input.readAll(actualSecondFrame, 0, actualSecondFrame.length);
100+
assertArrayEquals(secondFrame, actualSecondFrame);
101+
}
102+
64103
@Test
65104
public void testSingularSize() {
66105

0 commit comments

Comments
 (0)