From 00f43e12d242cd86239ac43ef8a6a37e8d344532 Mon Sep 17 00:00:00 2001 From: BackendArchitectX <96111851+BackendArchitectX@users.noreply.github.com> Date: Wed, 9 Sep 2026 19:43:16 +0530 Subject: [PATCH 1/2] [rpc] Handle endpoint changes for cached server connections --- .../fluss/rpc/netty/client/NettyClient.java | 89 +++++++++--- .../rpc/netty/client/ServerConnection.java | 4 +- .../rpc/netty/client/NettyClientTest.java | 136 +++++++++++++++++- 3 files changed, 207 insertions(+), 22 deletions(-) diff --git a/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/NettyClient.java b/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/NettyClient.java index 50e714d0269..05a00f8b1e8 100644 --- a/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/NettyClient.java +++ b/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/NettyClient.java @@ -41,6 +41,7 @@ import javax.annotation.concurrent.ThreadSafe; import java.util.ArrayList; +import java.util.Collection; import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; @@ -66,10 +67,10 @@ public final class NettyClient implements RpcClient { private final EventLoopGroup eventGroup; /** - * Managed connections to Netty servers. The key is the server uid (e.g., "cs-2", "ts-3"), the - * value is the connection. + * Managed connections to Netty servers. Connections are identified by server uid, host and + * port. */ - private final Map connections; + private final Map connections; /** Metric groups for client. */ private final ClientMetricGroup clientMetricGroup; @@ -130,11 +131,21 @@ public boolean connect(ServerNode node) { public CompletableFuture disconnect(String serverUid) { LOG.debug("Disconnecting from server {}.", serverUid); checkArgument(!isClosed, "Netty client is closed."); - ServerConnection connection = connections.remove(serverUid); - if (connection != null) { - return connection.close(); + + List> shutdownFutures = new ArrayList<>(); + + for (Map.Entry entry : connections.entrySet()) { + if (entry.getKey().belongsTo(serverUid) + && connections.remove(entry.getKey(), entry.getValue())) { + shutdownFutures.add(entry.getValue().close()); + } } - return FutureUtils.completedVoidFuture(); + + if (shutdownFutures.isEmpty()) { + return FutureUtils.completedVoidFuture(); + } + + return CompletableFuture.allOf(shutdownFutures.toArray(new CompletableFuture[0])); } /** @@ -146,11 +157,14 @@ public CompletableFuture disconnect(String serverUid) { @Override public boolean isReady(String serverUid) { checkArgument(!isClosed, "Netty client is closed."); - ServerConnection connection = connections.get(serverUid); - if (connection == null) { - return false; + + for (Map.Entry entry : connections.entrySet()) { + if (entry.getKey().belongsTo(serverUid) && entry.getValue().isReady()) { + return true; + } } - return connection.isReady(); + + return false; } /** Send an RPC request to the given server and return a future for the response. */ @@ -166,7 +180,7 @@ public void close() throws Exception { try { isClosed = true; final List> shutdownFutures = new ArrayList<>(); - for (Map.Entry conn : connections.entrySet()) { + for (Map.Entry conn : connections.entrySet()) { if (connections.remove(conn.getKey(), conn.getValue())) { shutdownFutures.add(conn.getValue().close()); } @@ -181,9 +195,9 @@ public void close() throws Exception { } private ServerConnection getOrCreateConnection(ServerNode node) { - String serverId = node.uid(); + ConnectionKey connectionKey = ConnectionKey.from(node); return connections.computeIfAbsent( - serverId, + connectionKey, ignored -> { LOG.debug("Creating connection to server {}.", node); return new ServerConnection( @@ -191,12 +205,53 @@ private ServerConnection getOrCreateConnection(ServerNode node) { node, clientMetricGroup, authenticatorSupplier.get(), - (con, ignore) -> connections.remove(serverId, con)); + (con, ignore) -> connections.remove(connectionKey, con)); }); } @VisibleForTesting - Map connections() { - return connections; + Collection connections() { + return connections.values(); + } + + private static final class ConnectionKey { + + private final String serverUid; + private final String host; + private final int port; + + private ConnectionKey(String serverUid, String host, int port) { + this.serverUid = serverUid; + this.host = host; + this.port = port; + } + + private static ConnectionKey from(ServerNode node) { + return new ConnectionKey(node.uid(), node.host(), node.port()); + } + + private boolean belongsTo(String serverUid) { + return this.serverUid.equals(serverUid); + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (!(o instanceof ConnectionKey)) { + return false; + } + ConnectionKey that = (ConnectionKey) o; + return port == that.port && serverUid.equals(that.serverUid) && host.equals(that.host); + } + + @Override + public int hashCode() { + int result = serverUid.hashCode(); + result = 31 * result + host.hashCode(); + result = 31 * result + port; + return result; + } } } diff --git a/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/ServerConnection.java b/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/ServerConnection.java index 1bc678452eb..35bb3198495 100644 --- a/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/ServerConnection.java +++ b/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/ServerConnection.java @@ -115,7 +115,9 @@ final class ServerConnection { BiConsumer closeCallback) { this.node = node; this.state = ConnectionState.CONNECTING; - this.connectionMetrics = clientMetricGroup.createConnectionMetricGroup(node.uid()); + this.connectionMetrics = + clientMetricGroup.createConnectionMetricGroup( + node.uid() + "@" + node.host() + ":" + node.port()); this.authenticator = authenticator; this.backoff = new ExponentialBackoff(100L, 2, 5000L, 0.2); whenClose(closeCallback); diff --git a/fluss-rpc/src/test/java/org/apache/fluss/rpc/netty/client/NettyClientTest.java b/fluss-rpc/src/test/java/org/apache/fluss/rpc/netty/client/NettyClientTest.java index 329f5d2f289..60df8c59e47 100644 --- a/fluss-rpc/src/test/java/org/apache/fluss/rpc/netty/client/NettyClientTest.java +++ b/fluss-rpc/src/test/java/org/apache/fluss/rpc/netty/client/NettyClientTest.java @@ -25,6 +25,7 @@ import org.apache.fluss.metrics.groups.MetricGroup; import org.apache.fluss.metrics.util.NOPMetricsGroup; import org.apache.fluss.rpc.TestingGatewayService; +import org.apache.fluss.rpc.TestingTabletGatewayService; import org.apache.fluss.rpc.messages.ApiMessage; import org.apache.fluss.rpc.messages.ApiVersionsRequest; import org.apache.fluss.rpc.messages.GetTableInfoRequest; @@ -147,8 +148,9 @@ void testServerDisconnection() throws Exception { .setClientSoftwareVersion("1.0"); nettyClient.sendRequest(serverNode, ApiKeys.API_VERSIONS, request).get(); assertThat(nettyClient.connections().size()).isEqualTo(1); - assertThat(nettyClient.connections().get(serverNode.uid()).getServerNode()) - .isEqualTo(serverNode); + assertThat(nettyClient.connections()) + .extracting(ServerConnection::getServerNode) + .containsExactly(serverNode); // close the netty server. nettyServer.close(); @@ -166,8 +168,134 @@ void testServerDisconnection() throws Exception { buildNettyServer(1); nettyClient.sendRequest(serverNode, ApiKeys.API_VERSIONS, request).get(); assertThat(nettyClient.connections().size()).isEqualTo(1); - assertThat(nettyClient.connections().get(serverNode.uid()).getServerNode()) - .isEqualTo(serverNode); + assertThat(nettyClient.connections()) + .extracting(ServerConnection::getServerNode) + .containsExactly(serverNode); + } + + @Test + void testSameServerUidWithChangedEndpointUsesNewEndpoint() throws Exception { + try (NetUtils.Port firstPort = getAvailablePort(); + NetUtils.Port secondPort = getAvailablePort()) { + + TestingTabletGatewayService firstService = new TestingTabletGatewayService(); + TestingTabletGatewayService secondService = new TestingTabletGatewayService(); + + MetricGroup firstMetricGroup = NOPMetricsGroup.newInstance(); + MetricGroup secondMetricGroup = NOPMetricsGroup.newInstance(); + + ServerNode firstNode = + new ServerNode(1, "localhost", firstPort.getPort(), ServerType.TABLET_SERVER); + + ServerNode secondNode = + new ServerNode(1, "localhost", secondPort.getPort(), ServerType.TABLET_SERVER); + + try (NettyServer firstServer = + new NettyServer( + conf, + Collections.singleton( + new Endpoint( + firstNode.host(), + firstNode.port(), + "INTERNAL")), + firstService, + firstMetricGroup, + RequestsMetrics.createTabletServerRequestMetrics( + firstMetricGroup)); + NettyServer secondServer = + new NettyServer( + conf, + Collections.singleton( + new Endpoint( + secondNode.host(), + secondNode.port(), + "INTERNAL")), + secondService, + secondMetricGroup, + RequestsMetrics.createTabletServerRequestMetrics( + secondMetricGroup))) { + + firstServer.start(); + secondServer.start(); + + ApiVersionsRequest request = + new ApiVersionsRequest() + .setClientSoftwareName("testing_client") + .setClientSoftwareVersion("1.0"); + + nettyClient.sendRequest(firstNode, ApiKeys.API_VERSIONS, request).get(); + + assertThat(firstService.getProcessorThreadNames()).hasSize(2); + assertThat(secondService.getProcessorThreadNames()).isEmpty(); + + nettyClient.sendRequest(secondNode, ApiKeys.API_VERSIONS, request).get(); + + assertThat(secondService.getProcessorThreadNames()).hasSize(2); + assertThat(firstService.getProcessorThreadNames()).hasSize(2); + } + } + } + + @Test + void testDisconnectClosesAllConnectionsForSameServerUid() throws Exception { + try (NetUtils.Port firstPort = getAvailablePort(); + NetUtils.Port secondPort = getAvailablePort()) { + + TestingTabletGatewayService firstService = new TestingTabletGatewayService(); + TestingTabletGatewayService secondService = new TestingTabletGatewayService(); + + MetricGroup firstMetricGroup = NOPMetricsGroup.newInstance(); + MetricGroup secondMetricGroup = NOPMetricsGroup.newInstance(); + + ServerNode firstNode = + new ServerNode(1, "localhost", firstPort.getPort(), ServerType.TABLET_SERVER); + + ServerNode secondNode = + new ServerNode(1, "localhost", secondPort.getPort(), ServerType.TABLET_SERVER); + + try (NettyServer firstServer = + new NettyServer( + conf, + Collections.singleton( + new Endpoint( + firstNode.host(), + firstNode.port(), + "INTERNAL")), + firstService, + firstMetricGroup, + RequestsMetrics.createTabletServerRequestMetrics( + firstMetricGroup)); + NettyServer secondServer = + new NettyServer( + conf, + Collections.singleton( + new Endpoint( + secondNode.host(), + secondNode.port(), + "INTERNAL")), + secondService, + secondMetricGroup, + RequestsMetrics.createTabletServerRequestMetrics( + secondMetricGroup))) { + + firstServer.start(); + secondServer.start(); + + ApiVersionsRequest request = + new ApiVersionsRequest() + .setClientSoftwareName("testing_client") + .setClientSoftwareVersion("1.0"); + + nettyClient.sendRequest(firstNode, ApiKeys.API_VERSIONS, request).get(); + nettyClient.sendRequest(secondNode, ApiKeys.API_VERSIONS, request).get(); + + assertThat(nettyClient.connections()).hasSize(2); + + nettyClient.disconnect(firstNode.uid()).get(); + + assertThat(nettyClient.connections()).isEmpty(); + } + } } @Test From d9fe4ccc4c45a123e568524b2c2e30f26b958935 Mon Sep 17 00:00:00 2001 From: BackendArchitectX <96111851+BackendArchitectX@users.noreply.github.com> Date: Thu, 10 Sep 2026 20:15:07 +0530 Subject: [PATCH 2/2] Address endpoint-specific RPC lifecycle review --- .../java/org/apache/fluss/rpc/RpcClient.java | 22 ++- .../fluss/rpc/netty/client/NettyClient.java | 46 +++++-- .../rpc/netty/client/NettyClientTest.java | 129 +++++++++++++++++- 3 files changed, 173 insertions(+), 24 deletions(-) diff --git a/fluss-rpc/src/main/java/org/apache/fluss/rpc/RpcClient.java b/fluss-rpc/src/main/java/org/apache/fluss/rpc/RpcClient.java index ce4e6fd2dba..a015ff53f2f 100644 --- a/fluss-rpc/src/main/java/org/apache/fluss/rpc/RpcClient.java +++ b/fluss-rpc/src/main/java/org/apache/fluss/rpc/RpcClient.java @@ -53,8 +53,17 @@ static RpcClient create(Configuration conf, ClientMetricGroup clientMetricGroup) boolean connect(ServerNode node); /** - * Disconnects the connection to the given server node, if there is one. Any in-flight/pending - * requests for this connection will receive disconnections. + * Disconnects the connection to the given server endpoint, if there is one. Any in-flight or + * pending requests for this connection will receive disconnections. + * + * @param node The server node to disconnect + * @return A future that is completed when the disconnection is complete + */ + CompletableFuture disconnect(ServerNode node); + + /** + * Disconnects all connections associated with the given logical server uid. Any in-flight or + * pending requests for these connections will receive disconnections. * * @param serverUid The uid of the server node * @return A future that is completed when the disconnection is complete @@ -62,12 +71,13 @@ static RpcClient create(Configuration conf, ClientMetricGroup clientMetricGroup) CompletableFuture disconnect(String serverUid); /** - * Check if we are currently ready to send another request to the given server but don't attempt - * to connect if we aren't. + * Check if we are currently ready to send another request to the given server endpoint but + * don't attempt to connect if we aren't. * - * @return true if the node is ready + * @param node The server node to check + * @return true if the connection to the node is ready */ - boolean isReady(String serverUid); + boolean isReady(ServerNode node); /** * Send an RPC request to the given server and return a future for the response. If the diff --git a/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/NettyClient.java b/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/NettyClient.java index 05a00f8b1e8..ec7054fb49c 100644 --- a/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/NettyClient.java +++ b/fluss-rpc/src/main/java/org/apache/fluss/rpc/netty/client/NettyClient.java @@ -121,13 +121,34 @@ public boolean connect(ServerNode node) { } /** - * Disconnects the connection to the given server node, if there is one. Any inflight/pending - * requests for this connection will receive disconnections. + * Disconnects the connection to the given server endpoint, if there is one. Any + * inflight/pending requests for this connection will receive disconnections. * - * @param serverUid The uid of the server node + * @param node The server node to disconnect * @return A future that completes when the connection is fully closed */ @Override + public CompletableFuture disconnect(ServerNode node) { + LOG.debug("Disconnecting from server {}.", node); + checkArgument(!isClosed, "Netty client is closed."); + + ServerConnection connection = connections.remove(ConnectionKey.from(node)); + + if (connection == null) { + return FutureUtils.completedVoidFuture(); + } + + return connection.close(); + } + + /** + * Disconnects all connections associated with the given logical server uid. Any + * inflight/pending requests for these connections will receive disconnections. + * + * @param serverUid The uid of the server node + * @return A future that completes when all associated connections are fully closed + */ + @Override public CompletableFuture disconnect(String serverUid) { LOG.debug("Disconnecting from server {}.", serverUid); checkArgument(!isClosed, "Netty client is closed."); @@ -149,22 +170,19 @@ public CompletableFuture disconnect(String serverUid) { } /** - * Check if we are currently ready to send another request to the given server but don't attempt - * to connect if we aren't. + * Check if we are currently ready to send another request to the given server endpoint but + * don't attempt to connect if we aren't. * - * @return true if the node is ready + * @param node The server node to check + * @return true if the connection to the node is ready */ @Override - public boolean isReady(String serverUid) { + public boolean isReady(ServerNode node) { checkArgument(!isClosed, "Netty client is closed."); - for (Map.Entry entry : connections.entrySet()) { - if (entry.getKey().belongsTo(serverUid) && entry.getValue().isReady()) { - return true; - } - } + ServerConnection connection = connections.get(ConnectionKey.from(node)); - return false; + return connection != null && connection.isReady(); } /** Send an RPC request to the given server and return a future for the response. */ @@ -239,9 +257,11 @@ public boolean equals(Object o) { if (this == o) { return true; } + if (!(o instanceof ConnectionKey)) { return false; } + ConnectionKey that = (ConnectionKey) o; return port == that.port && serverUid.equals(that.serverUid) && host.equals(that.host); } diff --git a/fluss-rpc/src/test/java/org/apache/fluss/rpc/netty/client/NettyClientTest.java b/fluss-rpc/src/test/java/org/apache/fluss/rpc/netty/client/NettyClientTest.java index 60df8c59e47..52cc9aca695 100644 --- a/fluss-rpc/src/test/java/org/apache/fluss/rpc/netty/client/NettyClientTest.java +++ b/fluss-rpc/src/test/java/org/apache/fluss/rpc/netty/client/NettyClientTest.java @@ -124,17 +124,23 @@ void testSendRequestToWrongServerType() { void testRequestsProcessedInOrder() throws Exception { int numRequests = 100; List> futures = new ArrayList<>(); + for (int i = 0; i < numRequests; i++) { ApiVersionsRequest request = new ApiVersionsRequest() .setClientSoftwareName("testing_client" + i) .setClientSoftwareVersion("1.0"); + futures.add(nettyClient.sendRequest(serverNode, ApiKeys.API_VERSIONS, request)); } + FutureUtils.waitForAll(futures).get(); + // we have one more api version request for rpc handshake. assertThat(service.getProcessorThreadNames()).hasSize(numRequests + 1); + Set deduplicatedThreadNames = new HashSet<>(service.getProcessorThreadNames()); + // there should only one thread to process the requests // since all requests are from the same client. assertThat(deduplicatedThreadNames).hasSize(1); @@ -146,14 +152,17 @@ void testServerDisconnection() throws Exception { new ApiVersionsRequest() .setClientSoftwareName("testing_client_100") .setClientSoftwareVersion("1.0"); + nettyClient.sendRequest(serverNode, ApiKeys.API_VERSIONS, request).get(); - assertThat(nettyClient.connections().size()).isEqualTo(1); + + assertThat(nettyClient.connections()).hasSize(1); assertThat(nettyClient.connections()) .extracting(ServerConnection::getServerNode) .containsExactly(serverNode); // close the netty server. nettyServer.close(); + assertThatThrownBy( () -> nettyClient @@ -162,12 +171,15 @@ void testServerDisconnection() throws Exception { .rootCause() .isInstanceOf(ConnectException.class) .hasMessageContaining("Connection refused"); - assertThat(nettyClient.connections().size()).isEqualTo(0); + + assertThat(nettyClient.connections()).isEmpty(); // restart the netty server. buildNettyServer(1); + nettyClient.sendRequest(serverNode, ApiKeys.API_VERSIONS, request).get(); - assertThat(nettyClient.connections().size()).isEqualTo(1); + + assertThat(nettyClient.connections()).hasSize(1); assertThat(nettyClient.connections()) .extracting(ServerConnection::getServerNode) .containsExactly(serverNode); @@ -223,15 +235,103 @@ void testSameServerUidWithChangedEndpointUsesNewEndpoint() throws Exception { .setClientSoftwareName("testing_client") .setClientSoftwareVersion("1.0"); + // Establish only the first endpoint connection. nettyClient.sendRequest(firstNode, ApiKeys.API_VERSIONS, request).get(); assertThat(firstService.getProcessorThreadNames()).hasSize(2); assertThat(secondService.getProcessorThreadNames()).isEmpty(); + // Readiness must be endpoint-specific. The second node has the same UID, + // but there is no connection to its host/port yet. + assertThat(nettyClient.isReady(firstNode)).isTrue(); + assertThat(nettyClient.isReady(secondNode)).isFalse(); + + // Send to the changed endpoint. nettyClient.sendRequest(secondNode, ApiKeys.API_VERSIONS, request).get(); assertThat(secondService.getProcessorThreadNames()).hasSize(2); assertThat(firstService.getProcessorThreadNames()).hasSize(2); + + // Both physical endpoint connections now coexist. + assertThat(nettyClient.isReady(firstNode)).isTrue(); + assertThat(nettyClient.isReady(secondNode)).isTrue(); + assertThat(nettyClient.connections()).hasSize(2); + } + } + } + + @Test + void testDisconnectSpecificEndpoint() throws Exception { + try (NetUtils.Port firstPort = getAvailablePort(); + NetUtils.Port secondPort = getAvailablePort()) { + + TestingTabletGatewayService firstService = new TestingTabletGatewayService(); + TestingTabletGatewayService secondService = new TestingTabletGatewayService(); + + MetricGroup firstMetricGroup = NOPMetricsGroup.newInstance(); + MetricGroup secondMetricGroup = NOPMetricsGroup.newInstance(); + + ServerNode firstNode = + new ServerNode(1, "localhost", firstPort.getPort(), ServerType.TABLET_SERVER); + + ServerNode secondNode = + new ServerNode(1, "localhost", secondPort.getPort(), ServerType.TABLET_SERVER); + + try (NettyServer firstServer = + new NettyServer( + conf, + Collections.singleton( + new Endpoint( + firstNode.host(), + firstNode.port(), + "INTERNAL")), + firstService, + firstMetricGroup, + RequestsMetrics.createTabletServerRequestMetrics( + firstMetricGroup)); + NettyServer secondServer = + new NettyServer( + conf, + Collections.singleton( + new Endpoint( + secondNode.host(), + secondNode.port(), + "INTERNAL")), + secondService, + secondMetricGroup, + RequestsMetrics.createTabletServerRequestMetrics( + secondMetricGroup))) { + + firstServer.start(); + secondServer.start(); + + ApiVersionsRequest request = + new ApiVersionsRequest() + .setClientSoftwareName("testing_client") + .setClientSoftwareVersion("1.0"); + + nettyClient.sendRequest(firstNode, ApiKeys.API_VERSIONS, request).get(); + nettyClient.sendRequest(secondNode, ApiKeys.API_VERSIONS, request).get(); + + assertThat(nettyClient.connections()).hasSize(2); + assertThat(nettyClient.isReady(firstNode)).isTrue(); + assertThat(nettyClient.isReady(secondNode)).isTrue(); + + // Disconnect only the first physical endpoint. + nettyClient.disconnect(firstNode).get(); + + assertThat(nettyClient.isReady(firstNode)).isFalse(); + assertThat(nettyClient.isReady(secondNode)).isTrue(); + + assertThat(nettyClient.connections()) + .extracting(ServerConnection::getServerNode) + .containsExactly(secondNode); + + // Verify that the remaining endpoint is still usable. + nettyClient.sendRequest(secondNode, ApiKeys.API_VERSIONS, request).get(); + + // handshake + first application request + request above + assertThat(secondService.getProcessorThreadNames()).hasSize(3); } } } @@ -290,10 +390,15 @@ void testDisconnectClosesAllConnectionsForSameServerUid() throws Exception { nettyClient.sendRequest(secondNode, ApiKeys.API_VERSIONS, request).get(); assertThat(nettyClient.connections()).hasSize(2); + assertThat(nettyClient.isReady(firstNode)).isTrue(); + assertThat(nettyClient.isReady(secondNode)).isTrue(); + // UID-level disconnect must close every physical endpoint for the logical server. nettyClient.disconnect(firstNode.uid()).get(); assertThat(nettyClient.connections()).isEmpty(); + assertThat(nettyClient.isReady(firstNode)).isFalse(); + assertThat(nettyClient.isReady(secondNode)).isFalse(); } } } @@ -316,6 +421,7 @@ void testBindFailureDetection() { @Test void testMultipleEndpoint() throws Exception { MetricGroup metricGroup = NOPMetricsGroup.newInstance(); + try (NetUtils.Port availablePort1 = getAvailablePort(); NetUtils.Port availablePort2 = getAvailablePort(); NettyServer multipleEndpointsServer = @@ -330,11 +436,14 @@ void testMultipleEndpoint() throws Exception { metricGroup, RequestsMetrics.createCoordinatorServerRequestMetrics( metricGroup))) { + multipleEndpointsServer.start(); + ApiVersionsRequest request = new ApiVersionsRequest() .setClientSoftwareName("testing_client_100") .setClientSoftwareVersion("1.0"); + nettyClient .sendRequest( new ServerNode( @@ -345,9 +454,12 @@ void testMultipleEndpoint() throws Exception { ApiKeys.API_VERSIONS, request) .get(); - assertThat(nettyClient.connections().size()).isEqualTo(1); + + assertThat(nettyClient.connections()).hasSize(1); + try (NettyClient client = new NettyClient(conf, TestingClientMetricGroup.newInstance())) { + client.sendRequest( new ServerNode( 2, @@ -357,7 +469,8 @@ void testMultipleEndpoint() throws Exception { ApiKeys.API_VERSIONS, request) .get(); - assertThat(client.connections().size()).isEqualTo(1); + + assertThat(client.connections()).hasSize(1); } } } @@ -368,6 +481,7 @@ void testExceptionWhenInitializeServerConnection() throws Exception { new ApiVersionsRequest() .setClientSoftwareName("testing_client_100") .setClientSoftwareVersion("1.0"); + // close the netty server. nettyServer.close(); @@ -378,6 +492,7 @@ void testExceptionWhenInitializeServerConnection() throws Exception { .sendRequest(serverNode, ApiKeys.API_VERSIONS, request) .get()) .hasMessageContaining("Disconnected from node"); + assertThat(nettyClient.connections()).isEmpty(); } @@ -386,8 +501,11 @@ private void buildNettyServer(int serverId) throws Exception { serverNode = new ServerNode( serverId, "localhost", availablePort.getPort(), ServerType.COORDINATOR); + service = new TestingGatewayService(); + MetricGroup metricGroup = NOPMetricsGroup.newInstance(); + nettyServer = new NettyServer( conf, @@ -396,6 +514,7 @@ private void buildNettyServer(int serverId) throws Exception { service, metricGroup, RequestsMetrics.createCoordinatorServerRequestMetrics(metricGroup)); + nettyServer.start(); } }