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";

Reply via email to