diff --git a/.github/workflows/multi-language-client.yml b/.github/workflows/multi-language-client.yml
index 61a947c3e08a4..b49b63b2e5c77 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/**'
@@ -80,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
@@ -108,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:
@@ -182,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: |
@@ -209,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
diff --git a/LICENSE-binary b/LICENSE-binary
index 84d24b8bdbcbf..5d1de65c54ee7 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/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 33778bf8fb23d..9adb0ccf0e0bd 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,9 @@ 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
+ 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!";
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 81b519fd4779f..8cc7d9e1adfd9 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,9 @@ 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
+ EXCEPTION_FRAME_SIZE_ARG_EXCEEDS_THE_MAXIMUM_SUPPORTED_MESSAGE_SIZE_E83C0952 =
+ "帧大小 (%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/NonOpenTransport.java b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/NonOpenTransport.java
index 8348fef373ec6..09de65d68300b 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/TCompressedElasticFramedTransport.java b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TCompressedElasticFramedTransport.java
index 88c9683db970b..b4c6c0994f556 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/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TElasticFramedTransport.java
index d0c6a3d24cadd..339dee595e7b7 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,28 @@ 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,
+ String.format(
+ RpcMessages
+ .EXCEPTION_FRAME_SIZE_ARG_EXCEEDS_THE_MAXIMUM_SUPPORTED_MESSAGE_SIZE_E83C0952,
+ thriftMaxFrameSize));
+ }
+ if (configuration.getMaxFrameSize() < thriftMaxFrameSize) {
+ configuration.setMaxFrameSize(thriftMaxFrameSize);
+ }
+ if (configuration.getMaxMessageSize() < maxMessageSize) {
+ configuration.setMaxMessageSize((int) maxMessageSize);
+ }
+ }
+
protected final int thriftDefaultBufferSize;
protected final int thriftMaxFrameSize;
@@ -194,10 +217,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 +367,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 d9e99ec8232df..0cbab4409e585 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,18 +19,119 @@
package org.apache.iotdb.rpc;
+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;
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;
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(14).setMaxFrameSize(10).build();
+ TMemoryBuffer underlying = new BudgetConsumingMemoryBuffer(configuration, 14);
+ underlying.write(framedData.array());
+ TElasticFramedTransport transport =
+ new TElasticFramedTransport(underlying, 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 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 maxFrameSize = Math.max(firstCompressedSize, secondCompressedSize);
+ TConfiguration configuration =
+ TConfiguration.custom()
+ .setMaxMessageSize(maxFrameSize)
+ .setMaxFrameSize(maxFrameSize)
+ .build();
+ 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);
+ assertArrayEquals(firstFrame, actualFirstFrame);
+
+ byte[] actualSecondFrame = new byte[secondFrame.length];
+ input.readAll(actualSecondFrame, 0, actualSecondFrame.length);
+ 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() {
diff --git a/pom.xml b/pom.xml
index 762db1d9ef73d..abf184a78442e 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