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