From a795478263fea897599b500725875fd05f9c9240 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 3 Mar 2020 11:11:19 -0500 Subject: [PATCH] Revert NPE check in the FailoverClientConnFactory * Add `synchronized` to `MqttPahoMessageDrivenChannelAdapter` setters to fix sync inconsistency --- .../ip/tcp/connection/FailoverClientConnectionFactory.java | 3 +-- .../mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java | 4 ++-- 2 files changed, 3 insertions(+), 4 deletions(-) diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactory.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactory.java index 9e699e6946..e3d27f804a 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactory.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/FailoverClientConnectionFactory.java @@ -171,7 +171,7 @@ public class FailoverClientConnectionFactory extends AbstractClientConnectionFac return failoverTcpConnection; } - private void closeRefreshedIfNecessary(@Nullable FailoverTcpConnection sharedConnection, boolean refreshShared, + private void closeRefreshedIfNecessary(FailoverTcpConnection sharedConnection, boolean refreshShared, FailoverTcpConnection failoverTcpConnection) { this.creationTime = System.currentTimeMillis(); @@ -179,7 +179,6 @@ public class FailoverClientConnectionFactory extends AbstractClientConnectionFac * We may have simply wrapped the same connection in a new wrapper; don't close. */ if (refreshShared && this.closeOnRefresh - && sharedConnection != null && !sharedConnection.delegate.equals(failoverTcpConnection.delegate) && sharedConnection.isOpen()) { diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java index 7b86d1bddd..a229966f8c 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/inbound/MqttPahoMessageDrivenChannelAdapter.java @@ -128,7 +128,7 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv * @param completionTimeout The timeout. * @since 4.1 */ - public void setCompletionTimeout(long completionTimeout) { + public synchronized void setCompletionTimeout(long completionTimeout) { this.completionTimeout = completionTimeout; } @@ -138,7 +138,7 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv * @param completionTimeout The timeout. * @since 5.1.10 */ - public void setDisconnectCompletionTimeout(long completionTimeout) { + public synchronized void setDisconnectCompletionTimeout(long completionTimeout) { this.disconnectCompletionTimeout = completionTimeout; }