This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new dfc2cd7919 [service] Fix connection race leak in
NetworkClient.sendRequest (#8684)
dfc2cd7919 is described below
commit dfc2cd791937a8f614a62db0fceb5e121edcfb32
Author: Eunbin Son <[email protected]>
AuthorDate: Thu Jul 16 14:11:07 2026 +0900
[service] Fix connection race leak in NetworkClient.sendRequest (#8684)
---
.../paimon/service/network/NetworkClient.java | 38 ++++---
.../paimon/service/network/NetworkClientTest.java | 121 +++++++++++++++++++++
2 files changed, 143 insertions(+), 16 deletions(-)
diff --git
a/paimon-service/paimon-service-client/src/main/java/org/apache/paimon/service/network/NetworkClient.java
b/paimon-service/paimon-service-client/src/main/java/org/apache/paimon/service/network/NetworkClient.java
index 748d1d2f05..a14971aa0a 100644
---
a/paimon-service/paimon-service-client/src/main/java/org/apache/paimon/service/network/NetworkClient.java
+++
b/paimon-service/paimon-service-client/src/main/java/org/apache/paimon/service/network/NetworkClient.java
@@ -143,22 +143,28 @@ public class NetworkClient<REQ extends MessageBody, RESP
extends MessageBody> {
new IllegalStateException(clientName + " is already shut
down."));
}
- ServerConnection<REQ, RESP> serverConnection =
connections.get(serverAddress);
- if (serverConnection == null) {
- final ServerConnection<REQ, RESP> newConnection =
- ServerConnection.createPendingConnection(clientName,
messageSerializer, stats);
- serverConnection = newConnection;
- connections.put(serverAddress, newConnection);
- bootstrap
- .connect(serverAddress.getAddress(),
serverAddress.getPort())
- .addListener((ChannelFutureListener)
newConnection::establishConnection);
-
- newConnection
- .getCloseFuture()
- .handle(
- (ignoredA, ignoredB) ->
- connections.remove(serverAddress,
newConnection));
- }
+ final ServerConnection<REQ, RESP> serverConnection =
+ connections.computeIfAbsent(
+ serverAddress,
+ ignored -> {
+ final ServerConnection<REQ, RESP> newConnection =
+ ServerConnection.createPendingConnection(
+ clientName, messageSerializer,
stats);
+ bootstrap
+ .connect(serverAddress.getAddress(),
serverAddress.getPort())
+ .addListener(
+ (ChannelFutureListener)
+
newConnection::establishConnection);
+
+ newConnection
+ .getCloseFuture()
+ .handle(
+ (ignoredA, ignoredB) ->
+ connections.remove(
+ serverAddress,
newConnection));
+
+ return newConnection;
+ });
return serverConnection.sendRequest(request);
}
diff --git
a/paimon-service/paimon-service-runtime/src/test/java/org/apache/paimon/service/network/NetworkClientTest.java
b/paimon-service/paimon-service-runtime/src/test/java/org/apache/paimon/service/network/NetworkClientTest.java
index 39e070fbee..efc08781bc 100644
---
a/paimon-service/paimon-service-runtime/src/test/java/org/apache/paimon/service/network/NetworkClientTest.java
+++
b/paimon-service/paimon-service-runtime/src/test/java/org/apache/paimon/service/network/NetworkClientTest.java
@@ -55,11 +55,13 @@ import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import static org.assertj.core.api.Assertions.assertThat;
@@ -332,6 +334,88 @@ class NetworkClientTest {
}
}
+ /**
+ * Tests that concurrent first requests racing on the initial connect to
the same server address
+ * create exactly one connection. The atomic {@code computeIfAbsent}
guarantees a single {@link
+ * ServerConnection}, hence a single accepted server channel and no
channel leak after shutdown.
+ */
+ @Test
+ void testConcurrentFirstRequestsCreateSingleConnection() throws Exception {
+ AtomicServiceRequestStats stats = new AtomicServiceRequestStats();
+
+ final MessageSerializer<KvRequest, KvResponse> serializer =
+ new MessageSerializer<>(
+ new KvRequest.KvRequestDeserializer(),
+ new KvResponse.KvResponseDeserializer());
+
+ final KvResponse expected = KvResponseTest.random();
+
+ ExecutorService executor = null;
+ NetworkClient<KvRequest, KvResponse> client = null;
+ Channel serverChannel = null;
+
+ try {
+ final int numThreads = 8;
+ executor = Executors.newFixedThreadPool(numThreads);
+
+ client = new NetworkClient<>("Test Client", 1, serializer, stats);
+
+ final AtomicInteger acceptedChannels = new AtomicInteger(0);
+ serverChannel =
+ createServerChannel(
+ new CountingRespondingChannelHandler(
+ serializer, expected, acceptedChannels));
+
+ final InetSocketAddress serverAddress =
getServerAddress(serverChannel);
+
+ // Release all threads simultaneously to maximize contention on
the initial connect.
+ final CountDownLatch startLatch = new CountDownLatch(1);
+ final NetworkClient<KvRequest, KvResponse> finalClient = client;
+ Callable<CompletableFuture<KvResponse>> queryTask =
+ () -> {
+ startLatch.await();
+ KvRequest request = KvRequestTest.random();
+ return finalClient.sendRequest(serverAddress, request);
+ };
+
+ List<Future<CompletableFuture<KvResponse>>> submitted = new
ArrayList<>();
+ for (int i = 0; i < numThreads; i++) {
+ submitted.add(executor.submit(queryTask));
+ }
+
+ startLatch.countDown();
+
+ // Wait for all request futures to complete successfully.
+ for (Future<CompletableFuture<KvResponse>> future : submitted) {
+ KvResponse actual = future.get().get();
+ assertThat(actual).isEqualTo(expected);
+ }
+
+ // computeIfAbsent guarantees a single ServerConnection, hence a
single accepted
+ // channel.
+ assertThat(acceptedChannels.get()).isEqualTo(1);
+ } finally {
+ if (executor != null) {
+ executor.shutdown();
+ }
+
+ if (serverChannel != null) {
+ serverChannel.close();
+ }
+
+ if (client != null) {
+ try {
+ client.shutdown().get();
+ } catch (Exception e) {
+ e.printStackTrace();
+ }
+ assertThat(client.isEventGroupShutdown()).isTrue();
+ }
+
+ assertThat(stats.getNumConnections()).withFailMessage("Channel
leak").isZero();
+ }
+ }
+
/**
* Tests that a server failure closes the connection and removes it from
the established
* connections.
@@ -572,4 +656,41 @@ class NetworkClientTest {
ctx.channel().writeAndFlush(serResponse);
}
}
+
+ @ChannelHandler.Sharable
+ private static final class CountingRespondingChannelHandler
+ extends ChannelInboundHandlerAdapter {
+ private final MessageSerializer<KvRequest, KvResponse> serializer;
+ private final KvResponse response;
+ private final AtomicInteger acceptedChannels;
+
+ private CountingRespondingChannelHandler(
+ MessageSerializer<KvRequest, KvResponse> serializer,
+ KvResponse response,
+ AtomicInteger acceptedChannels) {
+ this.serializer = serializer;
+ this.response = response;
+ this.acceptedChannels = acceptedChannels;
+ }
+
+ @Override
+ public void channelActive(ChannelHandlerContext ctx) {
+ acceptedChannels.incrementAndGet();
+ }
+
+ @Override
+ public void channelRead(ChannelHandlerContext ctx, Object msg) {
+ ByteBuf buf = (ByteBuf) msg;
+
assertThat(MessageSerializer.deserializeHeader(buf)).isEqualTo(MessageType.REQUEST);
+ long requestId = MessageSerializer.getRequestId(buf);
+ serializer.deserializeRequest(buf);
+
+ buf.release();
+
+ ByteBuf serResponse =
+ MessageSerializer.serializeResponse(ctx.alloc(),
requestId, response);
+
+ ctx.channel().writeAndFlush(serResponse);
+ }
+ }
}