Skip to content

Commit 7984abc

Browse files
authored
[Pipe] Fix AirGap receiver retry request body reuse (#18721)
1 parent d6b602e commit 7984abc

2 files changed

Lines changed: 73 additions & 1 deletion

File tree

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/airgap/IoTDBAirGapReceiver.java‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -163,7 +163,7 @@ private void receive() throws IOException {
163163

164164
private void handleReq(final AirGapPseudoTPipeTransferRequest req, final long startTime)
165165
throws IOException {
166-
final TPipeTransferResp resp = agent.receive(req);
166+
final TPipeTransferResp resp = agent.receive(duplicateReq(req));
167167

168168
final TSStatus status = resp.getStatus();
169169
if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
@@ -197,6 +197,15 @@ private void handleReq(final AirGapPseudoTPipeTransferRequest req, final long st
197197
}
198198
}
199199

200+
private AirGapPseudoTPipeTransferRequest duplicateReq(
201+
final AirGapPseudoTPipeTransferRequest req) {
202+
return (AirGapPseudoTPipeTransferRequest)
203+
new AirGapPseudoTPipeTransferRequest()
204+
.setVersion(req.getVersion())
205+
.setType(req.getType())
206+
.setBody(req.body.duplicate());
207+
}
208+
200209
private void ok() throws IOException {
201210
final OutputStream outputStream = socket.getOutputStream();
202211
outputStream.write(AirGapOneByteResponse.OK);

‎iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/receiver/protocol/airgap/IoTDBAirGapReceiverTest.java‎

Lines changed: 63 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,7 @@
5454
import java.net.SocketTimeoutException;
5555
import java.nio.ByteBuffer;
5656
import java.util.concurrent.atomic.AtomicBoolean;
57+
import java.util.concurrent.atomic.AtomicInteger;
5758

5859
public class IoTDBAirGapReceiverTest {
5960

@@ -146,6 +147,68 @@ public void testTemporaryUnavailableRetryTimeoutReturnsFail() throws Exception {
146147
}
147148
}
148149

150+
@Test
151+
public void testTemporaryUnavailableRetryUsesFreshRequestBody() throws Exception {
152+
final CommonConfig commonConfig = CommonDescriptor.getInstance().getConfig();
153+
final long originalRetryLocalIntervalMs = commonConfig.getPipeAirGapRetryLocalIntervalMs();
154+
final long originalRetryMaxMs = commonConfig.getPipeAirGapRetryMaxMs();
155+
156+
try {
157+
commonConfig.setPipeAirGapRetryLocalIntervalMs(0);
158+
commonConfig.setPipeAirGapRetryMaxMs(10_000);
159+
160+
final RecordingSocket socket = new RecordingSocket();
161+
final IoTDBAirGapReceiver receiver = new IoTDBAirGapReceiver(socket, 4L);
162+
final StubIoTDBDataNodeReceiverAgent stubAgent = new StubIoTDBDataNodeReceiverAgent();
163+
final byte[] expectedBody = new byte[] {1, 2, 3};
164+
final AtomicInteger receiveCount = new AtomicInteger();
165+
stubAgent.setStubReceiver(
166+
new IoTDBReceiver() {
167+
@Override
168+
public TPipeTransferResp receive(final TPipeTransferReq req) {
169+
final byte[] actualBody = new byte[req.body.remaining()];
170+
req.body.get(actualBody);
171+
Assert.assertArrayEquals(expectedBody, actualBody);
172+
return new TPipeTransferResp(
173+
new TSStatus(
174+
receiveCount.getAndIncrement() == 0
175+
? TSStatusCode.PIPE_RECEIVER_TEMPORARY_UNAVAILABLE_EXCEPTION
176+
.getStatusCode()
177+
: TSStatusCode.SUCCESS_STATUS.getStatusCode()));
178+
}
179+
180+
@Override
181+
public void handleExit() {
182+
// noop for unit test
183+
}
184+
185+
@Override
186+
public IoTDBSinkRequestVersion getVersion() {
187+
return IoTDBSinkRequestVersion.VERSION_1;
188+
}
189+
});
190+
setField(receiver, "agent", stubAgent);
191+
192+
final AirGapPseudoTPipeTransferRequest req = new AirGapPseudoTPipeTransferRequest();
193+
req.setVersion(IoTDBSinkRequestVersion.VERSION_1.getVersion());
194+
req.setType((short) 0);
195+
req.setBody(ByteBuffer.wrap(expectedBody));
196+
197+
final Method handleReq =
198+
IoTDBAirGapReceiver.class.getDeclaredMethod(
199+
"handleReq", AirGapPseudoTPipeTransferRequest.class, long.class);
200+
handleReq.setAccessible(true);
201+
handleReq.invoke(receiver, req, System.currentTimeMillis());
202+
203+
Assert.assertEquals(2, receiveCount.get());
204+
Assert.assertEquals(0, req.body.position());
205+
Assert.assertArrayEquals(AirGapOneByteResponse.OK, socket.getWrittenBytes());
206+
} finally {
207+
commonConfig.setPipeAirGapRetryLocalIntervalMs(originalRetryLocalIntervalMs);
208+
commonConfig.setPipeAirGapRetryMaxMs(originalRetryMaxMs);
209+
}
210+
}
211+
149212
@Test
150213
public void testAirGapReceiverExitCleansThriftReceiverRuntime() throws Throwable {
151214
final String sessionKey = "DataNode-1-air_gap-4";

0 commit comments

Comments
 (0)