From 490ae1b49befc3c393e9a69272535d64602ed256 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 10 Oct 2018 14:48:13 -0400 Subject: [PATCH] Some various polishing * Add `RotatingServerAdvice.StandardRotationPolicy.getCurrent()` method for better end-user experience when this class is extended * Add NPE protection into the MQTT Channel Adapters. Some code style polishing for them **Cherry-pick to 5.0.x** --- .../file/remote/aop/RotatingServerAdvice.java | 7 ++++- .../MqttPahoMessageDrivenChannelAdapter.java | 24 ++++++++++------- .../mqtt/outbound/MqttPahoMessageHandler.java | 26 ++++++++++--------- 3 files changed, 34 insertions(+), 23 deletions(-) diff --git a/spring-integration-file/src/main/java/org/springframework/integration/file/remote/aop/RotatingServerAdvice.java b/spring-integration-file/src/main/java/org/springframework/integration/file/remote/aop/RotatingServerAdvice.java index 1b3f0d6aa1..a2893914d4 100644 --- a/spring-integration-file/src/main/java/org/springframework/integration/file/remote/aop/RotatingServerAdvice.java +++ b/spring-integration-file/src/main/java/org/springframework/integration/file/remote/aop/RotatingServerAdvice.java @@ -36,6 +36,7 @@ import org.springframework.util.Assert; * * @author Gary Russell * @author Michael Forstner + * @author Artem Bilan * * @since 5.0.7 * @@ -118,7 +119,7 @@ public class RotatingServerAdvice extends AbstractMessageSourceAdvice { protected final Log logger = LogFactory.getLog(getClass()); - private final DelegatingSessionFactory factory; + protected final DelegatingSessionFactory factory; private final List keyDirectories = new ArrayList<>(); @@ -170,6 +171,10 @@ public class RotatingServerAdvice extends AbstractMessageSourceAdvice { return this.fair; } + protected KeyDirectory getCurrent() { + return this.current; + } + @Override public void beforeReceive(MessageSource source) { if (this.fair || !this.initialized) { 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 71288b9411..b2c834a065 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 @@ -162,6 +162,7 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv if (this.consumerStopAction.equals(ConsumerStopAction.UNSUBSCRIBE_ALWAYS) || (this.consumerStopAction.equals(ConsumerStopAction.UNSUBSCRIBE_CLEAN) && this.cleanSession)) { + this.client.unsubscribe(getTopic()); } } @@ -249,8 +250,8 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv if (grantedQos[i] != requestedQos[i]) { if (logger.isWarnEnabled()) { logger.warn("Granted QOS different to Requested QOS; topics: " + Arrays.toString(topics) - + " requested: " + Arrays.toString(requestedQos) - + " granted: " + Arrays.toString(grantedQos)); + + " requested: " + Arrays.toString(requestedQos) + + " granted: " + Arrays.toString(grantedQos)); } break; } @@ -294,7 +295,8 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv } } - private void scheduleReconnect() { + private synchronized void scheduleReconnect() { + cancelReconnect(); try { this.reconnectFuture = getTaskScheduler().schedule(() -> { try { @@ -324,12 +326,14 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv if (isRunning()) { this.logger.error("Lost connection: " + cause.getMessage() + "; retrying..."); this.connected = false; - try { - this.client.setCallback(null); - this.client.close(); - } - catch (MqttException e) { - // NOSONAR + if (this.client != null) { + try { + this.client.setCallback(null); + this.client.close(); + } + catch (MqttException e) { + // NOSONAR + } } this.client = null; scheduleReconnect(); @@ -340,7 +344,7 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv } @Override - public void messageArrived(String topic, MqttMessage mqttMessage) throws Exception { + public void messageArrived(String topic, MqttMessage mqttMessage) { Message message = this.getConverter().toMessage(topic, mqttMessage); try { sendMessage(message); diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java index a1e15f5ede..81971b42ee 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/MqttPahoMessageHandler.java @@ -39,6 +39,7 @@ import org.springframework.util.Assert; * * @author Gary Russell * @author Artem Bilan + * * @since 4.0 * */ @@ -51,13 +52,13 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler private final MqttPahoClientFactory clientFactory; - private volatile IMqttAsyncClient client; - private volatile boolean async; private volatile boolean asyncEvents; - private volatile ApplicationEventPublisher applicationEventPublisher; + private ApplicationEventPublisher applicationEventPublisher; + + private volatile IMqttAsyncClient client; /** * Use this constructor for a single url (although it may be overridden @@ -181,7 +182,6 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler catch (MqttException e) { if (client != null) { client.close(); - client = null; } throw new MessagingException("Failed to connect", e); } @@ -215,18 +215,20 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler @Override public synchronized void connectionLost(Throwable cause) { logger.error("Lost connection; will attempt reconnect on next request"); - try { - this.client.setCallback(null); - this.client.close(); + if (this.client != null) { + try { + this.client.setCallback(null); + this.client.close(); + } + catch (MqttException e) { + // NOSONAR + } + this.client = null; } - catch (MqttException e) { - // NOSONAR - } - this.client = null; } @Override - public void messageArrived(String topic, MqttMessage message) throws Exception { + public void messageArrived(String topic, MqttMessage message) { }