This is an automated email from the ASF dual-hosted git repository.
jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 7984abcabb6 [Pipe] Fix AirGap receiver retry request body reuse
(#18721)
7984abcabb6 is described below
commit 7984abcabb628da2eae9dc81a51037738340c032
Author: Caideyipi <[email protected]>
AuthorDate: Thu Sep 24 21:07:58 2026 +0800
[Pipe] Fix AirGap receiver retry request body reuse (#18721)
---
.../protocol/airgap/IoTDBAirGapReceiver.java | 11 +++-
.../protocol/airgap/IoTDBAirGapReceiverTest.java | 63 ++++++++++++++++++++++
2 files changed, 73 insertions(+), 1 deletion(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/airgap/IoTDBAirGapReceiver.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/airgap/IoTDBAirGapReceiver.java
index 9784efee5b2..48fb42d27c2 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/airgap/IoTDBAirGapReceiver.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/airgap/IoTDBAirGapReceiver.java
@@ -163,7 +163,7 @@ public class IoTDBAirGapReceiver extends WrappedRunnable {
private void handleReq(final AirGapPseudoTPipeTransferRequest req, final
long startTime)
throws IOException {
- final TPipeTransferResp resp = agent.receive(req);
+ final TPipeTransferResp resp = agent.receive(duplicateReq(req));
final TSStatus status = resp.getStatus();
if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
@@ -197,6 +197,15 @@ public class IoTDBAirGapReceiver extends WrappedRunnable {
}
}
+ private AirGapPseudoTPipeTransferRequest duplicateReq(
+ final AirGapPseudoTPipeTransferRequest req) {
+ return (AirGapPseudoTPipeTransferRequest)
+ new AirGapPseudoTPipeTransferRequest()
+ .setVersion(req.getVersion())
+ .setType(req.getType())
+ .setBody(req.body.duplicate());
+ }
+
private void ok() throws IOException {
final OutputStream outputStream = socket.getOutputStream();
outputStream.write(AirGapOneByteResponse.OK);
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/protocol/airgap/IoTDBAirGapReceiverTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/protocol/airgap/IoTDBAirGapReceiverTest.java
index 0ce9f1cd19e..d7e69b51cfd 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/protocol/airgap/IoTDBAirGapReceiverTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/protocol/airgap/IoTDBAirGapReceiverTest.java
@@ -54,6 +54,7 @@ import java.net.SocketException;
import java.net.SocketTimeoutException;
import java.nio.ByteBuffer;
import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
public class IoTDBAirGapReceiverTest {
@@ -146,6 +147,68 @@ public class IoTDBAirGapReceiverTest {
}
}
+ @Test
+ public void testTemporaryUnavailableRetryUsesFreshRequestBody() throws
Exception {
+ final CommonConfig commonConfig =
CommonDescriptor.getInstance().getConfig();
+ final long originalRetryLocalIntervalMs =
commonConfig.getPipeAirGapRetryLocalIntervalMs();
+ final long originalRetryMaxMs = commonConfig.getPipeAirGapRetryMaxMs();
+
+ try {
+ commonConfig.setPipeAirGapRetryLocalIntervalMs(0);
+ commonConfig.setPipeAirGapRetryMaxMs(10_000);
+
+ final RecordingSocket socket = new RecordingSocket();
+ final IoTDBAirGapReceiver receiver = new IoTDBAirGapReceiver(socket, 4L);
+ final StubIoTDBDataNodeReceiverAgent stubAgent = new
StubIoTDBDataNodeReceiverAgent();
+ final byte[] expectedBody = new byte[] {1, 2, 3};
+ final AtomicInteger receiveCount = new AtomicInteger();
+ stubAgent.setStubReceiver(
+ new IoTDBReceiver() {
+ @Override
+ public TPipeTransferResp receive(final TPipeTransferReq req) {
+ final byte[] actualBody = new byte[req.body.remaining()];
+ req.body.get(actualBody);
+ Assert.assertArrayEquals(expectedBody, actualBody);
+ return new TPipeTransferResp(
+ new TSStatus(
+ receiveCount.getAndIncrement() == 0
+ ?
TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION
+ .getStatusCode()
+ : TSStatusCode.SUCCESS_STATUS.getStatusCode()));
+ }
+
+ @Override
+ public void handleExit() {
+ // noop for unit test
+ }
+
+ @Override
+ public IoTDBSinkRequestVersion getVersion() {
+ return IoTDBSinkRequestVersion.VERSION_1;
+ }
+ });
+ setField(receiver, "agent", stubAgent);
+
+ final AirGapPseudoTPipeTransferRequest req = new
AirGapPseudoTPipeTransferRequest();
+ req.setVersion(IoTDBSinkRequestVersion.VERSION_1.getVersion());
+ req.setType((short) 0);
+ req.setBody(ByteBuffer.wrap(expectedBody));
+
+ final Method handleReq =
+ IoTDBAirGapReceiver.class.getDeclaredMethod(
+ "handleReq", AirGapPseudoTPipeTransferRequest.class, long.class);
+ handleReq.setAccessible(true);
+ handleReq.invoke(receiver, req, System.currentTimeMillis());
+
+ Assert.assertEquals(2, receiveCount.get());
+ Assert.assertEquals(0, req.body.position());
+ Assert.assertArrayEquals(AirGapOneByteResponse.OK,
socket.getWrittenBytes());
+ } finally {
+
commonConfig.setPipeAirGapRetryLocalIntervalMs(originalRetryLocalIntervalMs);
+ commonConfig.setPipeAirGapRetryMaxMs(originalRetryMaxMs);
+ }
+ }
+
@Test
public void testAirGapReceiverExitCleansThriftReceiverRuntime() throws
Throwable {
final String sessionKey = "DataNode-1-air_gap-4";