From c60dc0be27fef967a3f3a22a00742d17d7547a4e Mon Sep 17 00:00:00 2001 From: rdhabalia Date: Thu, 5 Mar 2020 19:00:21 -0800 Subject: [PATCH] [pulsar-client] remove duplicate cnx method --- .../apache/pulsar/client/api/ProducerCreationTest.java | 2 +- .../apache/pulsar/client/impl/ConnectionHandler.java | 8 ++------ .../org/apache/pulsar/client/impl/ConsumerImpl.java | 2 +- .../org/apache/pulsar/client/impl/ProducerImpl.java | 10 +++++----- 4 files changed, 9 insertions(+), 13 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java index b13b61d4d74b2..8aed4f978f4be 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerCreationTest.java @@ -86,7 +86,7 @@ public void testGeneratedNameProducerReconnect(TopicDomain domain) throws Pulsar //simulate create producer timeout. Thread.sleep(3000); - producer.getConnectionHandler().connectionClosed(producer.getConnectionHandler().getClientCnx()); + producer.getConnectionHandler().connectionClosed(producer.getConnectionHandler().cnx()); Assert.assertFalse(producer.isConnected()); Thread.sleep(3000); Assert.assertEquals(producer.getConnectionHandler().getEpoch(), 1); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java index ef58da6930a59..8eab9ab469e30 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java @@ -125,7 +125,8 @@ protected void resetBackoff() { backoff.reset(); } - protected ClientCnx cnx() { + @VisibleForTesting + public ClientCnx cnx() { return CLIENT_CNX_UPDATER.get(this); } @@ -133,11 +134,6 @@ protected boolean isRetriableError(PulsarClientException e) { return e instanceof PulsarClientException.LookupException; } - @VisibleForTesting - public ClientCnx getClientCnx() { - return CLIENT_CNX_UPDATER.get(this); - } - protected void setClientCnx(ClientCnx clientCnx) { CLIENT_CNX_UPDATER.set(this, clientCnx); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index ce4ee1ebc8ea8..14d0aeb14daee 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -1842,7 +1842,7 @@ void connectionClosed(ClientCnx cnx) { @VisibleForTesting public ClientCnx getClientCnx() { - return this.connectionHandler.getClientCnx(); + return this.connectionHandler.cnx(); } void setClientCnx(ClientCnx clientCnx) { diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index eb059096e254e..038db6123ce58 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -573,8 +573,8 @@ protected ByteBufPair sendMessage(long producerId, long lowestSequenceId, long h } private ChecksumType getChecksumType() { - if (connectionHandler.getClientCnx() == null - || connectionHandler.getClientCnx().getRemoteEndpointProtocolVersion() >= brokerChecksumSupportedVersion()) { + if (connectionHandler.cnx() == null + || connectionHandler.cnx().getRemoteEndpointProtocolVersion() >= brokerChecksumSupportedVersion()) { return ChecksumType.Crc32c; } else { return ChecksumType.None; @@ -783,11 +783,11 @@ public CompletableFuture closeAsync() { @Override public boolean isConnected() { - return connectionHandler.getClientCnx() != null && (getState() == State.Ready); + return connectionHandler.cnx() != null && (getState() == State.Ready); } public boolean isWritable() { - ClientCnx cnx = connectionHandler.getClientCnx(); + ClientCnx cnx = connectionHandler.cnx(); return cnx != null && cnx.channel().isWritable(); } @@ -1601,7 +1601,7 @@ void connectionClosed(ClientCnx cnx) { } ClientCnx getClientCnx() { - return this.connectionHandler.getClientCnx(); + return this.connectionHandler.cnx(); } void setClientCnx(ClientCnx clientCnx) {