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);
+        }
+    }
 }

Reply via email to