Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 7 additions & 5 deletions .github/workflows/multi-language-client.yml
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ on:
- master
- "rc/*"
paths:
- 'pom.xml'
- 'iotdb-client/pom.xml'
- 'iotdb-client/client-py/**'
- 'iotdb-client/client-cpp/**'
Expand All @@ -20,6 +21,7 @@ on:
- "rc/*"
- 'force_ci/**'
paths:
- 'pom.xml'
- 'iotdb-client/pom.xml'
- 'iotdb-client/client-py/**'
- 'iotdb-client/client-cpp/**'
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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: |
Expand All @@ -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

Expand Down
2 changes: 1 addition & 1 deletion LICENSE-binary
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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.
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Comment thread
HTHou marked this conversation as resolved.
RpcStat.readCompressedBytes.addAndGet(size);
try {
int uncompressedLength = uncompressedLength(readBuffer.getBuffer(), 0, size);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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;

Expand Down Expand Up @@ -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);
Comment thread
HTHou marked this conversation as resolved.
Comment thread
Copilot marked this conversation as resolved.
}

protected void validateFrame(int size) throws TTransportException {
Expand Down Expand Up @@ -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);
Comment thread
Copilot marked this conversation as resolved.
}

@Override
public void write(byte[] buf, int off, int len) throws TTransportException {
writeBuffer.write(buf, off, len);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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() {

Expand Down
8 changes: 6 additions & 2 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,10 @@
<id>jitpack.io</id>
<url>https://jitpack.io</url>
</repository>
<repository>
<id>apache.iotdb.thrift.0.25.0.0.staging</id>
<url>https://repository.apache.org/content/repositories/orgapacheiotdb-1203/</url>
</repository>
</repositories>
<parent>
<groupId>org.apache</groupId>
Expand Down Expand Up @@ -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
-->
<iotdb-tools-thrift.version>0.23.0.0</iotdb-tools-thrift.version>
<iotdb-tools-thrift.version>0.25.0.0</iotdb-tools-thrift.version>
<jackson.version>2.18.9</jackson.version>
<jakarta.annotation-api.version>3.0.0</jakarta.annotation-api.version>
<jakarta.servlet-api.version>6.0.0</jakarta.servlet-api.version>
Expand Down Expand Up @@ -142,7 +146,7 @@
<swagger.version>2.2.50</swagger.version>
<thrift.exec-cmd.executable>chmod</thrift.exec-cmd.executable>
<thrift.exec.absolute.path/>
<thrift.version>0.24.0</thrift.version>
<thrift.version>0.25.0</thrift.version>
<xz.version>1.9</xz.version>
<zstd-jni.version>1.5.6-3</zstd-jni.version>
<tsfile.version>2.4.1-260915-SNAPSHOT</tsfile.version>
Expand Down
Loading