From 2324f549ec2827ea4e516bc1826f620c20ab8734 Mon Sep 17 00:00:00 2001 From: HTHou Date: Tue, 6 Oct 2026 18:32:46 +0800 Subject: [PATCH 01/11] Bump Thrift to 0.25.0 --- LICENSE-binary | 2 +- pom.xml | 8 ++++++-- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/LICENSE-binary b/LICENSE-binary index 84d24b8bdbcb..5d1de65c54ee 100644 --- a/LICENSE-binary +++ b/LICENSE-binary @@ -243,7 +243,7 @@ org.eclipse.jetty.ee10:jetty-ee10-servlet:12.0.36 org.eclipse.jetty:jetty-util:12.0.36 com.google.code.findbugs:jsr305:3.0.2 com.librato.metrics:librato-java:2.1.0 -org.apache.thrift:libthrift:0.24.0 +org.apache.thrift:libthrift:0.25.0 io.dropwizard.metrics:metrics-core:4.2.19 io.dropwizard.metrics:metrics-jvm:3.2.2 com.librato.metrics:metrics-librato:5.1.0 diff --git a/pom.xml b/pom.xml index 762db1d9ef73..abf184a78442 100644 --- a/pom.xml +++ b/pom.xml @@ -27,6 +27,10 @@ jitpack.io https://jitpack.io + + apache.iotdb.thrift.0.25.0.0.staging + https://repository.apache.org/content/repositories/orgapacheiotdb-1203/ + org.apache @@ -91,7 +95,7 @@ This is the version of the thrift binary, that we release separately from here: https://github.com/apache/iotdb-bin-resources/tree/main/iotdb-tools-thrift --> - 0.23.0.0 + 0.25.0.0 2.18.9 3.0.0 6.0.0 @@ -142,7 +146,7 @@ 2.2.50 chmod - 0.24.0 + 0.25.0 1.9 1.5.6-3 2.4.1-260915-SNAPSHOT From d5d21555d49a938bffbe08d305f638a383fdfb0d Mon Sep 17 00:00:00 2001 From: HTHou Date: Tue, 6 Oct 2026 19:47:35 +0800 Subject: [PATCH 02/11] Adapt custom transports to Thrift 0.25 --- .../apache/iotdb/rpc/NonOpenTransport.java | 6 ++++ .../iotdb/rpc/TElasticFramedTransport.java | 9 ++++++ .../rpc/TElasticFramedTransportTest.java | 28 +++++++++++++++++++ 3 files changed, 43 insertions(+) diff --git a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/NonOpenTransport.java b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/NonOpenTransport.java index 8348fef373ec..09de65d68300 100644 --- a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/NonOpenTransport.java +++ b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/NonOpenTransport.java @@ -20,6 +20,7 @@ package org.apache.iotdb.rpc; import org.apache.thrift.transport.TTransport; +import org.apache.thrift.transport.TTransportException; /** A TTransport that does not require open/close. */ public abstract class NonOpenTransport extends TTransport { @@ -40,4 +41,9 @@ public boolean isOpen() { public void open() { isOpen = true; } + + @Override + public void resetMessageSizeAndConsumedBytes(long newSize) throws TTransportException { + // This in-memory transport does not track a message-size budget. + } } diff --git a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java index d0c6a3d24cad..da00259b148c 100644 --- a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java +++ b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java @@ -194,10 +194,14 @@ public int read(byte[] buf, int off, int len) throws TTransportException { } protected void readFrame() throws TTransportException { + // Discard the previous frame's budget before reading the next frame header and payload. + resetMessageSizeAndConsumedBytes(); underlying.readAll(i32buf, 0, 4); int size = TFramedTransport.decodeFrameSize(i32buf); validateFrame(size); readBuffer.fill(underlying, size); + // Bind subsequent protocol reads to the current frame size. + resetMessageSizeAndConsumedBytes(size); } protected void validateFrame(int size) throws TTransportException { @@ -340,6 +344,11 @@ public void checkReadBytesAvailable(long numBytes) throws TTransportException { error.throwException(numBytes, remoteInfo, limit); } + @Override + public void resetMessageSizeAndConsumedBytes(long newSize) throws TTransportException { + underlying.resetMessageSizeAndConsumedBytes(newSize); + } + @Override public void write(byte[] buf, int off, int len) throws TTransportException { writeBuffer.write(buf, off, len); diff --git a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java index d9e99ec8232d..7f9d7fa7de04 100644 --- a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java +++ b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java @@ -19,6 +19,7 @@ package org.apache.iotdb.rpc; +import org.apache.thrift.TConfiguration; import org.apache.thrift.transport.TByteBuffer; import org.apache.thrift.transport.TTransportException; import org.junit.Test; @@ -26,11 +27,38 @@ import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; +import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; import static org.junit.Assert.fail; public class TElasticFramedTransportTest { + @Test + public void testReadFramesResetMessageSize() throws TTransportException { + byte[] firstFrame = {1, 2}; + byte[] secondFrame = {3, 4, 5, 6, 7, 8}; + ByteBuffer framedData = + ByteBuffer.allocate(8 + firstFrame.length + secondFrame.length) + .putInt(firstFrame.length) + .put(firstFrame) + .putInt(secondFrame.length) + .put(secondFrame) + .flip(); + TConfiguration configuration = + TConfiguration.custom().setMaxMessageSize(10).setMaxFrameSize(10).build(); + TElasticFramedTransport transport = + new TElasticFramedTransport( + new TByteBuffer(configuration, framedData), 4, configuration.getMaxFrameSize(), false); + + byte[] actualFirstFrame = new byte[firstFrame.length]; + transport.readAll(actualFirstFrame, 0, actualFirstFrame.length); + assertArrayEquals(firstFrame, actualFirstFrame); + + byte[] actualSecondFrame = new byte[secondFrame.length]; + transport.readAll(actualSecondFrame, 0, actualSecondFrame.length); + assertArrayEquals(secondFrame, actualSecondFrame); + } + @Test public void testSingularSize() { From adfda68fef1af75b3014b45dd2c992a1f23ba440 Mon Sep 17 00:00:00 2001 From: HTHou Date: Tue, 6 Oct 2026 19:56:24 +0800 Subject: [PATCH 03/11] Fix framed transport message budget test --- .../org/apache/iotdb/rpc/TElasticFramedTransportTest.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java index 7f9d7fa7de04..0e12bdbf2132 100644 --- a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java +++ b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java @@ -21,6 +21,7 @@ import org.apache.thrift.TConfiguration; import org.apache.thrift.transport.TByteBuffer; +import org.apache.thrift.transport.TMemoryBuffer; import org.apache.thrift.transport.TTransportException; import org.junit.Test; @@ -46,9 +47,11 @@ public void testReadFramesResetMessageSize() throws TTransportException { .flip(); TConfiguration configuration = TConfiguration.custom().setMaxMessageSize(10).setMaxFrameSize(10).build(); + TMemoryBuffer underlying = new TMemoryBuffer(configuration, 10); + underlying.write(framedData.array()); TElasticFramedTransport transport = new TElasticFramedTransport( - new TByteBuffer(configuration, framedData), 4, configuration.getMaxFrameSize(), false); + underlying, 4, configuration.getMaxFrameSize(), false); byte[] actualFirstFrame = new byte[firstFrame.length]; transport.readAll(actualFirstFrame, 0, actualFirstFrame.length); From e4b141ca07ad2ffc65d981db64ebc969b268c2f9 Mon Sep 17 00:00:00 2001 From: HTHou Date: Tue, 6 Oct 2026 20:01:47 +0800 Subject: [PATCH 04/11] Format framed transport test --- .../java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java index 0e12bdbf2132..492f023e5dda 100644 --- a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java +++ b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java @@ -50,8 +50,7 @@ public void testReadFramesResetMessageSize() throws TTransportException { TMemoryBuffer underlying = new TMemoryBuffer(configuration, 10); underlying.write(framedData.array()); TElasticFramedTransport transport = - new TElasticFramedTransport( - underlying, 4, configuration.getMaxFrameSize(), false); + new TElasticFramedTransport(underlying, 4, configuration.getMaxFrameSize(), false); byte[] actualFirstFrame = new byte[firstFrame.length]; transport.readAll(actualFirstFrame, 0, actualFirstFrame.length); From 9e3767239e868afdcbb303c55d843967e4e48aba Mon Sep 17 00:00:00 2001 From: HTHou Date: Tue, 6 Oct 2026 21:11:49 +0800 Subject: [PATCH 05/11] Trigger client CI for root POM changes --- .github/workflows/multi-language-client.yml | 2 ++ 1 file changed, 2 insertions(+) diff --git a/.github/workflows/multi-language-client.yml b/.github/workflows/multi-language-client.yml index 61a947c3e08a..f7e4cc8c315a 100644 --- a/.github/workflows/multi-language-client.yml +++ b/.github/workflows/multi-language-client.yml @@ -6,6 +6,7 @@ on: - master - "rc/*" paths: + - 'pom.xml' - 'iotdb-client/pom.xml' - 'iotdb-client/client-py/**' - 'iotdb-client/client-cpp/**' @@ -20,6 +21,7 @@ on: - "rc/*" - 'force_ci/**' paths: + - 'pom.xml' - 'iotdb-client/pom.xml' - 'iotdb-client/client-py/**' - 'iotdb-client/client-cpp/**' From ee8852d26c9e8e433df193e9f8f504f3b2e7d570 Mon Sep 17 00:00:00 2001 From: HTHou Date: Tue, 6 Oct 2026 21:33:13 +0800 Subject: [PATCH 06/11] Add Linux AArch64 C++ client CI --- .github/workflows/multi-language-client.yml | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/.github/workflows/multi-language-client.yml b/.github/workflows/multi-language-client.yml index f7e4cc8c315a..81b5084c6ea6 100644 --- a/.github/workflows/multi-language-client.yml +++ b/.github/workflows/multi-language-client.yml @@ -110,7 +110,7 @@ jobs: fail-fast: false max-parallel: 15 matrix: - os: [ubuntu-22.04, ubuntu-24.04, windows-2022, windows-2025-vs2026, macos-latest] + os: [ubuntu-22.04, ubuntu-24.04, ubuntu-22.04-arm, windows-2022, windows-2025-vs2026, macos-latest] runs-on: ${{ matrix.os}} steps: @@ -184,8 +184,8 @@ jobs: uses: actions/cache@v5 with: path: ~/.m2 - key: ${{ runner.os }}-m2-${{ hashFiles('**/pom.xml') }} - restore-keys: ${{ runner.os }}-m2- + key: ${{ runner.os }}-${{ runner.arch }}-m2-${{ hashFiles('**/pom.xml') }} + restore-keys: ${{ runner.os }}-${{ runner.arch }}-m2- - name: Check C++ format (Spotless) shell: bash run: | @@ -211,7 +211,7 @@ jobs: if: failure() uses: actions/upload-artifact@v6 with: - name: cpp-IT-${{ runner.os }} + name: cpp-IT-${{ matrix.os }} path: distribution/target/apache-iotdb-*-all-bin/apache-iotdb-*-all-bin/logs retention-days: 1 From f5b00810e6516ec445a6579a2718d7d5757cfe70 Mon Sep 17 00:00:00 2001 From: HTHou Date: Tue, 6 Oct 2026 22:31:40 +0800 Subject: [PATCH 07/11] Reset message budget for compressed frames --- .../TCompressedElasticFramedTransport.java | 4 ++ .../rpc/TElasticFramedTransportTest.java | 39 +++++++++++++++++++ 2 files changed, 43 insertions(+) diff --git a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TCompressedElasticFramedTransport.java b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TCompressedElasticFramedTransport.java index 88c9683db970..b4c6c0994f55 100644 --- a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TCompressedElasticFramedTransport.java +++ b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TCompressedElasticFramedTransport.java @@ -59,10 +59,14 @@ protected void closeAllocatedBuffers() { @Override protected void readFrame() throws TTransportException { + // Discard the previous frame's budget before reading the next frame header and payload. + resetMessageSizeAndConsumedBytes(); underlying.readAll(i32buf, 0, 4); int size = TFramedTransport.decodeFrameSize(i32buf); validateFrame(size); readBuffer.fill(underlying, size); + // Bind subsequent protocol reads to the current compressed frame size. + resetMessageSizeAndConsumedBytes(size); RpcStat.readCompressedBytes.addAndGet(size); try { int uncompressedLength = uncompressedLength(readBuffer.getBuffer(), 0, size); diff --git a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java index 492f023e5dda..76c923d751be 100644 --- a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java +++ b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java @@ -27,6 +27,7 @@ import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; +import java.util.Arrays; import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; @@ -61,6 +62,44 @@ public void testReadFramesResetMessageSize() throws TTransportException { assertArrayEquals(secondFrame, actualSecondFrame); } + @Test + public void testReadCompressedFramesResetMessageSize() throws TTransportException { + byte[] firstFrame = {1, 2}; + byte[] secondFrame = {3, 4, 5, 6, 7, 8}; + TMemoryBuffer wire = new TMemoryBuffer(128); + TSnappyElasticFramedTransport output = + new TSnappyElasticFramedTransport(wire, 4, 128, false); + output.write(firstFrame); + output.flush(); + output.write(secondFrame); + output.flush(); + + byte[] framedData = Arrays.copyOf(wire.getArray(), wire.length()); + ByteBuffer frameSizes = ByteBuffer.wrap(framedData); + int firstCompressedSize = frameSizes.getInt(); + frameSizes.position(Integer.BYTES + firstCompressedSize); + int secondCompressedSize = frameSizes.getInt(); + int maxMessageSize = Integer.BYTES + Math.max(firstCompressedSize, secondCompressedSize); + TConfiguration configuration = + TConfiguration.custom() + .setMaxMessageSize(maxMessageSize) + .setMaxFrameSize(maxMessageSize) + .build(); + TMemoryBuffer underlying = new TMemoryBuffer(configuration, maxMessageSize); + underlying.write(framedData); + TSnappyElasticFramedTransport input = + new TSnappyElasticFramedTransport( + underlying, 4, configuration.getMaxFrameSize(), false); + + byte[] actualFirstFrame = new byte[firstFrame.length]; + input.readAll(actualFirstFrame, 0, actualFirstFrame.length); + assertArrayEquals(firstFrame, actualFirstFrame); + + byte[] actualSecondFrame = new byte[secondFrame.length]; + input.readAll(actualSecondFrame, 0, actualSecondFrame.length); + assertArrayEquals(secondFrame, actualSecondFrame); + } + @Test public void testSingularSize() { From 3209ba21dc6d634340f7cbc0eeec84a779f975a3 Mon Sep 17 00:00:00 2001 From: HTHou Date: Tue, 6 Oct 2026 22:34:54 +0800 Subject: [PATCH 08/11] Format compressed frame reset test --- .../org/apache/iotdb/rpc/TElasticFramedTransportTest.java | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java index 76c923d751be..4fb86c388de1 100644 --- a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java +++ b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java @@ -67,8 +67,7 @@ public void testReadCompressedFramesResetMessageSize() throws TTransportExceptio byte[] firstFrame = {1, 2}; byte[] secondFrame = {3, 4, 5, 6, 7, 8}; TMemoryBuffer wire = new TMemoryBuffer(128); - TSnappyElasticFramedTransport output = - new TSnappyElasticFramedTransport(wire, 4, 128, false); + TSnappyElasticFramedTransport output = new TSnappyElasticFramedTransport(wire, 4, 128, false); output.write(firstFrame); output.flush(); output.write(secondFrame); @@ -88,8 +87,7 @@ public void testReadCompressedFramesResetMessageSize() throws TTransportExceptio TMemoryBuffer underlying = new TMemoryBuffer(configuration, maxMessageSize); underlying.write(framedData); TSnappyElasticFramedTransport input = - new TSnappyElasticFramedTransport( - underlying, 4, configuration.getMaxFrameSize(), false); + new TSnappyElasticFramedTransport(underlying, 4, configuration.getMaxFrameSize(), false); byte[] actualFirstFrame = new byte[firstFrame.length]; input.readAll(actualFirstFrame, 0, actualFirstFrame.length); From b3e2f8f2dd22886f748da043a6bac22a2daf8b03 Mon Sep 17 00:00:00 2001 From: HTHou Date: Wed, 7 Oct 2026 14:26:18 +0800 Subject: [PATCH 09/11] Align framed transport message budgets --- .../iotdb/rpc/TElasticFramedTransport.java | 20 ++++++++ .../rpc/TElasticFramedTransportTest.java | 46 ++++++++++++++++--- 2 files changed, 60 insertions(+), 6 deletions(-) diff --git a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java index da00259b148c..09df552a1c1e 100644 --- a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java +++ b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java @@ -91,6 +91,7 @@ public TElasticFramedTransport( this.thriftDefaultBufferSize = thriftDefaultBufferSize; this.thriftMaxFrameSize = thriftMaxFrameSize; this.copyBinary = copyBinary; + alignUnderlyingConfiguration(); try { readBuffer = new AutoScalingBufferReadTransport(thriftDefaultBufferSize); writeBuffer = new AutoScalingBufferWriteTransport(thriftDefaultBufferSize); @@ -100,6 +101,25 @@ public TElasticFramedTransport( } } + private void alignUnderlyingConfiguration() throws TTransportException { + TConfiguration configuration = underlying.getConfiguration(); + if (configuration == null) { + return; + } + long maxMessageSize = (long) thriftMaxFrameSize + Integer.BYTES; + if (maxMessageSize > Integer.MAX_VALUE) { + throw new TTransportException( + TTransportException.MESSAGE_SIZE_LIMIT, + "Frame size exceeds the maximum supported message size: " + thriftMaxFrameSize); + } + if (configuration.getMaxFrameSize() < thriftMaxFrameSize) { + configuration.setMaxFrameSize(thriftMaxFrameSize); + } + if (configuration.getMaxMessageSize() < maxMessageSize) { + configuration.setMaxMessageSize((int) maxMessageSize); + } + } + protected final int thriftDefaultBufferSize; protected final int thriftMaxFrameSize; diff --git a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java index 4fb86c388de1..0cbab4409e58 100644 --- a/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java +++ b/iotdb-client/service-rpc/src/test/java/org/apache/iotdb/rpc/TElasticFramedTransportTest.java @@ -47,8 +47,8 @@ public void testReadFramesResetMessageSize() throws TTransportException { .put(secondFrame) .flip(); TConfiguration configuration = - TConfiguration.custom().setMaxMessageSize(10).setMaxFrameSize(10).build(); - TMemoryBuffer underlying = new TMemoryBuffer(configuration, 10); + TConfiguration.custom().setMaxMessageSize(14).setMaxFrameSize(10).build(); + TMemoryBuffer underlying = new BudgetConsumingMemoryBuffer(configuration, 14); underlying.write(framedData.array()); TElasticFramedTransport transport = new TElasticFramedTransport(underlying, 4, configuration.getMaxFrameSize(), false); @@ -78,16 +78,17 @@ public void testReadCompressedFramesResetMessageSize() throws TTransportExceptio int firstCompressedSize = frameSizes.getInt(); frameSizes.position(Integer.BYTES + firstCompressedSize); int secondCompressedSize = frameSizes.getInt(); - int maxMessageSize = Integer.BYTES + Math.max(firstCompressedSize, secondCompressedSize); + int maxFrameSize = Math.max(firstCompressedSize, secondCompressedSize); TConfiguration configuration = TConfiguration.custom() - .setMaxMessageSize(maxMessageSize) - .setMaxFrameSize(maxMessageSize) + .setMaxMessageSize(maxFrameSize) + .setMaxFrameSize(maxFrameSize) .build(); - TMemoryBuffer underlying = new TMemoryBuffer(configuration, maxMessageSize); + TMemoryBuffer underlying = new BudgetConsumingMemoryBuffer(configuration, maxFrameSize); underlying.write(framedData); TSnappyElasticFramedTransport input = new TSnappyElasticFramedTransport(underlying, 4, configuration.getMaxFrameSize(), false); + assertEquals(Integer.BYTES + maxFrameSize, configuration.getMaxMessageSize()); byte[] actualFirstFrame = new byte[firstFrame.length]; input.readAll(actualFirstFrame, 0, actualFirstFrame.length); @@ -98,6 +99,39 @@ public void testReadCompressedFramesResetMessageSize() throws TTransportExceptio assertArrayEquals(secondFrame, actualSecondFrame); } + @Test + public void testFrameSizeAlignsUnderlyingMessageSize() throws TTransportException { + byte[] frame = {1, 2, 3, 4, 5, 6, 7, 8}; + ByteBuffer framedData = + ByteBuffer.allocate(Integer.BYTES + frame.length).putInt(frame.length).put(frame).flip(); + TConfiguration configuration = + TConfiguration.custom().setMaxMessageSize(6).setMaxFrameSize(frame.length).build(); + TMemoryBuffer underlying = new BudgetConsumingMemoryBuffer(configuration, 6); + underlying.write(framedData.array()); + TElasticFramedTransport transport = + new TElasticFramedTransport(underlying, 4, frame.length, false); + + assertEquals(Integer.BYTES + frame.length, configuration.getMaxMessageSize()); + byte[] actualFrame = new byte[frame.length]; + transport.readAll(actualFrame, 0, actualFrame.length); + assertArrayEquals(frame, actualFrame); + } + + private static class BudgetConsumingMemoryBuffer extends TMemoryBuffer { + + private BudgetConsumingMemoryBuffer(TConfiguration configuration, int size) + throws TTransportException { + super(configuration, size); + } + + @Override + public int read(byte[] buffer, int offset, int length) throws TTransportException { + int bytesRead = super.read(buffer, offset, length); + countConsumedMessageBytes(bytesRead); + return bytesRead; + } + } + @Test public void testSingularSize() { From 01904af0e1e1f91df61d11f718fe637d626f8552 Mon Sep 17 00:00:00 2001 From: HTHou Date: Wed, 7 Oct 2026 17:28:35 +0800 Subject: [PATCH 10/11] Localize frame size limit error --- .../main/i18n/en/org/apache/iotdb/rpc/i18n/RpcMessages.java | 2 ++ .../main/i18n/zh/org/apache/iotdb/rpc/i18n/RpcMessages.java | 2 ++ .../java/org/apache/iotdb/rpc/TElasticFramedTransport.java | 4 +++- 3 files changed, 7 insertions(+), 1 deletion(-) diff --git a/iotdb-client/service-rpc/src/main/i18n/en/org/apache/iotdb/rpc/i18n/RpcMessages.java b/iotdb-client/service-rpc/src/main/i18n/en/org/apache/iotdb/rpc/i18n/RpcMessages.java index 33778bf8fb23..1c563776effa 100644 --- a/iotdb-client/service-rpc/src/main/i18n/en/org/apache/iotdb/rpc/i18n/RpcMessages.java +++ b/iotdb-client/service-rpc/src/main/i18n/en/org/apache/iotdb/rpc/i18n/RpcMessages.java @@ -33,6 +33,8 @@ public final class RpcMessages { "Read a negative frame size (%d)%s!"; public static final String FRAME_ERROR_FRAME_SIZE_EXCEEDED = "Frame size (%d) larger than protect max size (%d)%s!"; + public static final String FRAME_ERROR_FRAME_SIZE_EXCEEDS_MAX_MESSAGE_SIZE = + "Frame size (%d) exceeds the maximum supported message size!"; public static final String FRAME_ERROR_STRING_LENGTH_EXCEEDED = "String length (%d) larger than protect max size (%d)%s!"; public static final String diff --git a/iotdb-client/service-rpc/src/main/i18n/zh/org/apache/iotdb/rpc/i18n/RpcMessages.java b/iotdb-client/service-rpc/src/main/i18n/zh/org/apache/iotdb/rpc/i18n/RpcMessages.java index 81b519fd4779..d85ae82f6603 100644 --- a/iotdb-client/service-rpc/src/main/i18n/zh/org/apache/iotdb/rpc/i18n/RpcMessages.java +++ b/iotdb-client/service-rpc/src/main/i18n/zh/org/apache/iotdb/rpc/i18n/RpcMessages.java @@ -30,6 +30,8 @@ public final class RpcMessages { public static final String FRAME_ERROR_NEGATIVE_FRAME_SIZE = "读取到负数帧大小 (%d)%s!"; public static final String FRAME_ERROR_FRAME_SIZE_EXCEEDED = "帧大小 (%d) 超过保护最大值 (%d)%s!"; + public static final String FRAME_ERROR_FRAME_SIZE_EXCEEDS_MAX_MESSAGE_SIZE = + "帧大小 (%d) 超过支持的最大消息大小!"; public static final String FRAME_ERROR_STRING_LENGTH_EXCEEDED = "字符串长度 (%d) 超过保护最大值 (%d)%s!"; public static final String diff --git a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java index 09df552a1c1e..e7ba0fcf0cb6 100644 --- a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java +++ b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java @@ -110,7 +110,9 @@ private void alignUnderlyingConfiguration() throws TTransportException { if (maxMessageSize > Integer.MAX_VALUE) { throw new TTransportException( TTransportException.MESSAGE_SIZE_LIMIT, - "Frame size exceeds the maximum supported message size: " + thriftMaxFrameSize); + String.format( + RpcMessages.FRAME_ERROR_FRAME_SIZE_EXCEEDS_MAX_MESSAGE_SIZE, + thriftMaxFrameSize)); } if (configuration.getMaxFrameSize() < thriftMaxFrameSize) { configuration.setMaxFrameSize(thriftMaxFrameSize); From 10ba8ea1fa252825ef098ffa498ebf217ffe7b41 Mon Sep 17 00:00:00 2001 From: HTHou Date: Wed, 7 Oct 2026 17:48:35 +0800 Subject: [PATCH 11/11] Fix Thrift upgrade validation issues --- .github/workflows/multi-language-client.yml | 2 +- .../main/i18n/en/org/apache/iotdb/rpc/i18n/RpcMessages.java | 3 ++- .../main/i18n/zh/org/apache/iotdb/rpc/i18n/RpcMessages.java | 3 ++- .../java/org/apache/iotdb/rpc/TElasticFramedTransport.java | 3 ++- 4 files changed, 7 insertions(+), 4 deletions(-) diff --git a/.github/workflows/multi-language-client.yml b/.github/workflows/multi-language-client.yml index 81b5084c6ea6..b49b63b2e5c7 100644 --- a/.github/workflows/multi-language-client.yml +++ b/.github/workflows/multi-language-client.yml @@ -82,7 +82,7 @@ jobs: go=false while IFS= read -r file; do case "$file" in - iotdb-client/pom.xml|iotdb-client/client-cpp/*|iotdb-protocol/thrift-datanode/src/main/thrift/client.thrift|iotdb-protocol/thrift-commons/src/main/thrift/common.thrift|.github/workflows/multi-language-client.yml|.github/workflows/client-cpp-package.yml|.github/scripts/package-client-cpp-*.sh) + pom.xml|iotdb-client/pom.xml|iotdb-client/client-cpp/*|iotdb-protocol/thrift-datanode/src/main/thrift/client.thrift|iotdb-protocol/thrift-commons/src/main/thrift/common.thrift|.github/workflows/multi-language-client.yml|.github/workflows/client-cpp-package.yml|.github/scripts/package-client-cpp-*.sh) cpp=true ;; esac diff --git a/iotdb-client/service-rpc/src/main/i18n/en/org/apache/iotdb/rpc/i18n/RpcMessages.java b/iotdb-client/service-rpc/src/main/i18n/en/org/apache/iotdb/rpc/i18n/RpcMessages.java index 1c563776effa..9adb0ccf0e0b 100644 --- a/iotdb-client/service-rpc/src/main/i18n/en/org/apache/iotdb/rpc/i18n/RpcMessages.java +++ b/iotdb-client/service-rpc/src/main/i18n/en/org/apache/iotdb/rpc/i18n/RpcMessages.java @@ -33,7 +33,8 @@ public final class RpcMessages { "Read a negative frame size (%d)%s!"; public static final String FRAME_ERROR_FRAME_SIZE_EXCEEDED = "Frame size (%d) larger than protect max size (%d)%s!"; - public static final String FRAME_ERROR_FRAME_SIZE_EXCEEDS_MAX_MESSAGE_SIZE = + public static final String + EXCEPTION_FRAME_SIZE_ARG_EXCEEDS_THE_MAXIMUM_SUPPORTED_MESSAGE_SIZE_E83C0952 = "Frame size (%d) exceeds the maximum supported message size!"; public static final String FRAME_ERROR_STRING_LENGTH_EXCEEDED = "String length (%d) larger than protect max size (%d)%s!"; diff --git a/iotdb-client/service-rpc/src/main/i18n/zh/org/apache/iotdb/rpc/i18n/RpcMessages.java b/iotdb-client/service-rpc/src/main/i18n/zh/org/apache/iotdb/rpc/i18n/RpcMessages.java index d85ae82f6603..8cc7d9e1adfd 100644 --- a/iotdb-client/service-rpc/src/main/i18n/zh/org/apache/iotdb/rpc/i18n/RpcMessages.java +++ b/iotdb-client/service-rpc/src/main/i18n/zh/org/apache/iotdb/rpc/i18n/RpcMessages.java @@ -30,7 +30,8 @@ public final class RpcMessages { public static final String FRAME_ERROR_NEGATIVE_FRAME_SIZE = "读取到负数帧大小 (%d)%s!"; public static final String FRAME_ERROR_FRAME_SIZE_EXCEEDED = "帧大小 (%d) 超过保护最大值 (%d)%s!"; - public static final String FRAME_ERROR_FRAME_SIZE_EXCEEDS_MAX_MESSAGE_SIZE = + public static final String + EXCEPTION_FRAME_SIZE_ARG_EXCEEDS_THE_MAXIMUM_SUPPORTED_MESSAGE_SIZE_E83C0952 = "帧大小 (%d) 超过支持的最大消息大小!"; public static final String FRAME_ERROR_STRING_LENGTH_EXCEEDED = "字符串长度 (%d) 超过保护最大值 (%d)%s!"; diff --git a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java index e7ba0fcf0cb6..339dee595e7b 100644 --- a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java +++ b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java @@ -111,7 +111,8 @@ private void alignUnderlyingConfiguration() throws TTransportException { throw new TTransportException( TTransportException.MESSAGE_SIZE_LIMIT, String.format( - RpcMessages.FRAME_ERROR_FRAME_SIZE_EXCEEDS_MAX_MESSAGE_SIZE, + RpcMessages + .EXCEPTION_FRAME_SIZE_ARG_EXCEEDS_THE_MAXIMUM_SUPPORTED_MESSAGE_SIZE_E83C0952, thriftMaxFrameSize)); } if (configuration.getMaxFrameSize() < thriftMaxFrameSize) {