More MQTT polishing according Sonar report
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2018 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -33,18 +33,18 @@ import org.springframework.util.StringUtils;
|
||||
* respective {@link BeanDefinition}s.
|
||||
*
|
||||
* @author Gary Russell
|
||||
* @author Artem Bilan
|
||||
*
|
||||
* @since 4.0
|
||||
*
|
||||
*/
|
||||
public final class MqttParserUtils {
|
||||
final class MqttParserUtils {
|
||||
|
||||
/** Prevent instantiation. */
|
||||
private MqttParserUtils() {
|
||||
throw new AssertionError();
|
||||
|
||||
}
|
||||
|
||||
public static void parseCommon(Element element, BeanDefinitionBuilder builder, ParserContext parserContext) {
|
||||
|
||||
static void parseCommon(Element element, BeanDefinitionBuilder builder, ParserContext parserContext) {
|
||||
ValueHolder holder;
|
||||
int n = 0;
|
||||
String url = element.getAttribute("url");
|
||||
@@ -54,7 +54,7 @@ public final class MqttParserUtils {
|
||||
holder.setType("java.lang.String");
|
||||
}
|
||||
builder.addConstructorArgValue(element.getAttribute("client-id"));
|
||||
holder = builder.getRawBeanDefinition().getConstructorArgumentValues().getIndexedArgumentValues().get(n++);
|
||||
holder = builder.getRawBeanDefinition().getConstructorArgumentValues().getIndexedArgumentValues().get(n);
|
||||
holder.setType("java.lang.String");
|
||||
String clientFactory = element.getAttribute("client-factory");
|
||||
if (StringUtils.hasText(clientFactory)) {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2018 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -35,6 +35,8 @@ import org.springframework.util.Assert;
|
||||
* Abstract class for MQTT Message-Driven Channel Adapters.
|
||||
*
|
||||
* @author Gary Russell
|
||||
* @author Artem Bilan
|
||||
*
|
||||
* @since 4.0
|
||||
*
|
||||
*/
|
||||
@@ -50,7 +52,7 @@ public abstract class AbstractMqttMessageDrivenChannelAdapter extends MessagePro
|
||||
|
||||
private volatile MqttMessageConverter converter;
|
||||
|
||||
protected final Lock topicLock = new ReentrantLock();
|
||||
protected final Lock topicLock = new ReentrantLock(); // NOSONAR
|
||||
|
||||
public AbstractMqttMessageDrivenChannelAdapter(String url, String clientId, String... topic) {
|
||||
Assert.hasText(clientId, "'clientId' cannot be null or empty");
|
||||
@@ -58,7 +60,7 @@ public abstract class AbstractMqttMessageDrivenChannelAdapter extends MessagePro
|
||||
Assert.noNullElements(topic, "'topics' cannot have null elements");
|
||||
this.url = url;
|
||||
this.clientId = clientId;
|
||||
this.topics = new LinkedHashSet<Topic>();
|
||||
this.topics = new LinkedHashSet<>();
|
||||
for (String t : topic) {
|
||||
this.topics.add(new Topic(t, 1));
|
||||
}
|
||||
@@ -225,10 +227,8 @@ public abstract class AbstractMqttMessageDrivenChannelAdapter extends MessagePro
|
||||
this.topicLock.lock();
|
||||
try {
|
||||
for (String t : topic) {
|
||||
if (this.topics.remove(new Topic(t, 0))) {
|
||||
if (this.logger.isDebugEnabled()) {
|
||||
logger.debug("Removed '" + t + "' from subscriptions.");
|
||||
}
|
||||
if (this.topics.remove(new Topic(t, 0)) && this.logger.isDebugEnabled()) {
|
||||
logger.debug("Removed '" + t + "' from subscriptions.");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -45,7 +45,7 @@ import org.springframework.util.Assert;
|
||||
* @author Gary Russell
|
||||
* @author Artem Bilan
|
||||
*
|
||||
* @since 1.0
|
||||
* @since 4.0
|
||||
*
|
||||
*/
|
||||
public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDrivenChannelAdapter
|
||||
@@ -57,10 +57,10 @@ public class MqttPahoMessageDrivenChannelAdapter extends AbstractMqttMessageDriv
|
||||
|
||||
private final MqttPahoClientFactory clientFactory;
|
||||
|
||||
private int recoveryInterval = DEFAULT_RECOVERY_INTERVAL;
|
||||
|
||||
private volatile long completionTimeout = DEFAULT_COMPLETION_TIMEOUT;
|
||||
|
||||
private volatile int recoveryInterval = DEFAULT_RECOVERY_INTERVAL;
|
||||
|
||||
private volatile IMqttClient client;
|
||||
|
||||
private volatile ScheduledFuture<?> reconnectFuture;
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
* Copyright 2002-2018 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -37,6 +37,7 @@ import org.springframework.util.Assert;
|
||||
*
|
||||
* @author Gary Russell
|
||||
* @author Artem Bilan
|
||||
*
|
||||
* @since 4.0
|
||||
*
|
||||
*/
|
||||
|
||||
@@ -165,23 +165,22 @@ public class MqttPahoMessageHandler extends AbstractMqttMessageHandler
|
||||
this.client = null;
|
||||
}
|
||||
if (this.client == null) {
|
||||
IMqttAsyncClient client = null;
|
||||
try {
|
||||
MqttConnectOptions connectionOptions = this.clientFactory.getConnectionOptions();
|
||||
Assert.state(this.getUrl() != null || connectionOptions.getServerURIs() != null,
|
||||
"If no 'url' provided, connectionOptions.getServerURIs() must not be null");
|
||||
client = this.clientFactory.getAsyncClientInstance(this.getUrl(), this.getClientId());
|
||||
this.client = this.clientFactory.getAsyncClientInstance(this.getUrl(), this.getClientId());
|
||||
incrementClientInstance();
|
||||
client.setCallback(this);
|
||||
client.connect(connectionOptions).waitForCompletion(this.completionTimeout);
|
||||
this.client = client;
|
||||
this.client.setCallback(this);
|
||||
this.client.connect(connectionOptions).waitForCompletion(this.completionTimeout);
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Client connected");
|
||||
}
|
||||
}
|
||||
catch (MqttException e) {
|
||||
if (client != null) {
|
||||
client.close();
|
||||
if (this.client != null) {
|
||||
this.client.close();
|
||||
this.client = null;
|
||||
}
|
||||
throw new MessagingException("Failed to connect", e);
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2017 the original author or authors.
|
||||
* Copyright 2002-2018 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -16,6 +16,8 @@
|
||||
|
||||
package org.springframework.integration.mqtt.support;
|
||||
|
||||
import java.nio.charset.Charset;
|
||||
|
||||
import org.eclipse.paho.client.mqttv3.MqttMessage;
|
||||
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
@@ -38,12 +40,13 @@ import org.springframework.util.Assert;
|
||||
*
|
||||
* @author Gary Russell
|
||||
* @author Artem Bilan
|
||||
*
|
||||
* @since 4.0
|
||||
*
|
||||
*/
|
||||
public class DefaultPahoMessageConverter implements MqttMessageConverter, BeanFactoryAware {
|
||||
|
||||
private final String charset;
|
||||
private final Charset charset;
|
||||
|
||||
private final int defaultQos;
|
||||
|
||||
@@ -68,7 +71,7 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter, BeanFa
|
||||
* Construct a converter with default options (qos=0, retain=false, charset=UTF-8).
|
||||
*/
|
||||
public DefaultPahoMessageConverter() {
|
||||
this (0, false);
|
||||
this(0, false);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -86,7 +89,7 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter, BeanFa
|
||||
/**
|
||||
* Construct a converter with default options (qos=0, retain=false) and
|
||||
* the supplied charset.
|
||||
* @param charset the charset used to convert outbound String paylaods to {@code byte[]} and inbound
|
||||
* @param charset the charset used to convert outbound String payloads to {@code byte[]} and inbound
|
||||
* {@code byte[]} to String (unless {@link #setPayloadAsBytes(boolean) payloadAdBytes} is true).
|
||||
* @since 4.1.2
|
||||
*/
|
||||
@@ -99,7 +102,7 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter, BeanFa
|
||||
* retain settings and the supplied charset.
|
||||
* @param defaultQos the default qos.
|
||||
* @param defaultRetained the default retained.
|
||||
* @param charset the charset used to convert outbound String paylaods to
|
||||
* @param charset the charset used to convert outbound String payloads to
|
||||
* {@code byte[]} and inbound {@code byte[]} to String (unless
|
||||
* {@link #setPayloadAsBytes(boolean) payloadAdBytes} is true).
|
||||
*/
|
||||
@@ -121,6 +124,7 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter, BeanFa
|
||||
*/
|
||||
public DefaultPahoMessageConverter(int defaultQos, MessageProcessor<Integer> qosProcessor, boolean defaultRetained,
|
||||
MessageProcessor<Boolean> retainedProcessor) {
|
||||
|
||||
this(defaultQos, qosProcessor, defaultRetained, retainedProcessor, "UTF-8");
|
||||
}
|
||||
|
||||
@@ -131,20 +135,21 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter, BeanFa
|
||||
* @param qosProcessor a message processor to determine the qos.
|
||||
* @param defaultRetained the default retained.
|
||||
* @param retainedProcessor a message processor to determine the retained flag.
|
||||
* @param charset the charset used to convert outbound String paylaods to
|
||||
* @param charset the charset used to convert outbound String payloads to
|
||||
* {@code byte[]} and inbound {@code byte[]} to String (unless
|
||||
* {@link #setPayloadAsBytes(boolean) payloadAdBytes} is true).
|
||||
* @since 5.0
|
||||
*/
|
||||
public DefaultPahoMessageConverter(int defaultQos, MessageProcessor<Integer> qosProcessor, boolean defaultRetained,
|
||||
MessageProcessor<Boolean> retainedProcessor, String charset) {
|
||||
|
||||
Assert.notNull(qosProcessor, "'qosProcessor' cannot be null");
|
||||
Assert.notNull(retainedProcessor, "'retainedProcessor' cannot be null");
|
||||
this.defaultQos = defaultQos;
|
||||
this.qosProcessor = qosProcessor;
|
||||
this.defaultRetained = defaultRetained;
|
||||
this.retainedProcessor = retainedProcessor;
|
||||
this.charset = charset;
|
||||
this.charset = Charset.forName(charset);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -201,18 +206,19 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter, BeanFa
|
||||
return toMessage(null, (MqttMessage) mqttMessage);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Override
|
||||
public Message<?> toMessage(String topic, MqttMessage mqttMessage) {
|
||||
try {
|
||||
AbstractIntegrationMessageBuilder<Object> messageBuilder;
|
||||
AbstractIntegrationMessageBuilder<?> messageBuilder;
|
||||
if (this.bytesMessageMapper != null) {
|
||||
messageBuilder = (AbstractIntegrationMessageBuilder<Object>) getMessageBuilderFactory()
|
||||
.fromMessage(this.bytesMessageMapper.toMessage(mqttMessage.getPayload()));
|
||||
messageBuilder =
|
||||
getMessageBuilderFactory()
|
||||
.fromMessage(this.bytesMessageMapper.toMessage(mqttMessage.getPayload()));
|
||||
}
|
||||
else {
|
||||
messageBuilder = getMessageBuilderFactory()
|
||||
.withPayload(mqttBytesToPayload(mqttMessage));
|
||||
messageBuilder =
|
||||
getMessageBuilderFactory()
|
||||
.withPayload(mqttBytesToPayload(mqttMessage));
|
||||
}
|
||||
messageBuilder
|
||||
.setHeader(MqttHeaders.RECEIVED_QOS, mqttMessage.getQos())
|
||||
@@ -242,12 +248,10 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter, BeanFa
|
||||
/**
|
||||
* Subclasses can override this method to convert the byte[] to a payload.
|
||||
* The default implementation creates a String (default) or byte[].
|
||||
*
|
||||
* @param mqttMessage The inbound message.
|
||||
* @return The payload for the Spring integration message
|
||||
* @throws Exception Any.
|
||||
*/
|
||||
protected Object mqttBytesToPayload(MqttMessage mqttMessage) throws Exception {
|
||||
protected Object mqttBytesToPayload(MqttMessage mqttMessage) {
|
||||
if (this.payloadAsBytes) {
|
||||
return mqttMessage.getPayload();
|
||||
}
|
||||
@@ -261,7 +265,6 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter, BeanFa
|
||||
* The default implementation accepts a byte[] or String payload.
|
||||
* If a {@link BytesMessageMapper} is provided, conversion to byte[]
|
||||
* is delegated to it, so any payload that it can handle is supported.
|
||||
*
|
||||
* @param message The outbound Message.
|
||||
* @return The byte[] which will become the payload of the MQTT Message.
|
||||
*/
|
||||
@@ -283,12 +286,7 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter, BeanFa
|
||||
+ payload.getClass().getName() + " payloads");
|
||||
byte[] payloadBytes;
|
||||
if (payload instanceof String) {
|
||||
try {
|
||||
payloadBytes = ((String) payload).getBytes(this.charset);
|
||||
}
|
||||
catch (Exception e) {
|
||||
throw new MessageConversionException("failed to convert Message to object", e);
|
||||
}
|
||||
payloadBytes = ((String) payload).getBytes(this.charset);
|
||||
}
|
||||
else {
|
||||
payloadBytes = (byte[]) payload;
|
||||
|
||||
@@ -1,33 +0,0 @@
|
||||
/*
|
||||
* Copyright 2002-2016 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
package org.springframework.integration.mqtt.support;
|
||||
|
||||
|
||||
/**
|
||||
* Contains utility methods used by the MqttAdapter components.
|
||||
*
|
||||
* @author Gary Russell
|
||||
* @since 4.0
|
||||
*
|
||||
*/
|
||||
public final class MqttUtils {
|
||||
|
||||
private MqttUtils() {
|
||||
throw new AssertionError();
|
||||
}
|
||||
|
||||
}
|
||||
Reference in New Issue
Block a user