MQTT - Use INT/SPR 4.0.0

Conflicts:
	spring-integration-mqtt/build.gradle
This commit is contained in:
Gary Russell
2013-10-15 10:47:20 -04:00
committed by Artem Bilan
parent 2dd27c5df9
commit 46cf219b44
11 changed files with 41 additions and 32 deletions

View File

@@ -26,9 +26,8 @@ ext {
junitVersion = '4.11'
log4jVersion = '1.2.17'
mockitoVersion = '1.9.5'
springVersion = '3.2.8.RELEASE'
springIntegrationVersion = '3.0.1.RELEASE'
springVersion = '4.0.5.RELEASE'
springIntegrationVersion = '4.0.1.RELEASE'
idPrefix = 'mqtt'
linkHomepage = 'https://github.com/SpringSource/spring-integration-extensions'

View File

@@ -1 +1 @@
version=1.0.0.BUILD-SNAPSHOT
version=4.0.0.BUILD-SNAPSHOT

View File

@@ -24,9 +24,9 @@ import org.eclipse.paho.client.mqttv3.MqttClient;
import org.eclipse.paho.client.mqttv3.MqttException;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.integration.Message;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
import org.springframework.messaging.Message;
/**
* Eclipse Paho Implementation.

View File

@@ -16,13 +16,15 @@
package org.springframework.integration.mqtt.outbound;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.context.SmartLifecycle;
import org.springframework.integration.Message;
import org.springframework.integration.MessageHandlingException;
import org.springframework.integration.handler.AbstractMessageHandler;
import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter;
import org.springframework.integration.mqtt.support.MqttHeaders;
import org.springframework.integration.support.converter.MessageConverter;
import org.springframework.integration.mqtt.support.MqttMessageConverter;
import org.springframework.messaging.Message;
import org.springframework.util.Assert;
/**
@@ -43,7 +45,7 @@ public abstract class AbstractMqttMessageHandler extends AbstractMessageHandler
private volatile boolean defaultRetained = false;
private volatile MessageConverter converter;
private volatile MqttMessageConverter converter;
private boolean running;
@@ -70,7 +72,7 @@ public abstract class AbstractMqttMessageHandler extends AbstractMessageHandler
this.defaultRetained = defaultRetain;
}
public void setConverter(MessageConverter converter) {
public void setConverter(MqttMessageConverter converter) {
Assert.notNull(converter, "'converter' cannot be null");
this.converter = converter;
}
@@ -138,7 +140,7 @@ public abstract class AbstractMqttMessageHandler extends AbstractMessageHandler
protected void handleMessageInternal(Message<?> message) throws Exception {
this.connectIfNeeded();
String topic = (String) message.getHeaders().get(MqttHeaders.TOPIC);
Object mqttMessage = this.converter.fromMessage(message);
MqttMessage mqttMessage = this.converter.fromMessage(message, MqttMessage.class);
if (topic == null && this.defaultTopic == null) {
throw new MessageHandlingException(message,
"No '" + MqttHeaders.TOPIC + "' header and no default topic defined");

View File

@@ -20,9 +20,10 @@ import org.eclipse.paho.client.mqttv3.MqttCallback;
import org.eclipse.paho.client.mqttv3.MqttClient;
import org.eclipse.paho.client.mqttv3.MqttException;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.integration.MessagingException;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.core.MqttPahoClientFactory;
import org.springframework.messaging.MessagingException;
import org.springframework.util.Assert;
/**

View File

@@ -15,10 +15,13 @@
*/
package org.springframework.integration.mqtt.support;
import java.lang.reflect.Type;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.integration.Message;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.integration.support.converter.MessageConversionException;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.converter.MessageConversionException;
import org.springframework.util.Assert;
@@ -52,18 +55,16 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter {
@Override
@SuppressWarnings("unchecked")
public <P> Message<P> toMessage(Object object) {
return (Message<P>) toMessage(null, object);
public Message<?> toMessage(MqttMessage mqttMessage) {
return toMessage(null, mqttMessage);
}
public Message<String> toMessage(String topic, Object object) {
Assert.isInstanceOf(MqttMessage.class, object);
MqttMessage message = (MqttMessage) object;
public Message<String> toMessage(String topic, MqttMessage mqttMessage) {
try {
MessageBuilder<String> messageBuilder = MessageBuilder.withPayload(new String(message.getPayload(), this.charset))
.setHeader(MqttHeaders.QOS, message.getQos())
.setHeader(MqttHeaders.DUPLICATE, message.isDuplicate())
.setHeader(MqttHeaders.RETAINED, message.isRetained());
MessageBuilder<String> messageBuilder = MessageBuilder.withPayload(new String(mqttMessage.getPayload(), this.charset))
.setHeader(MqttHeaders.QOS, mqttMessage.getQos())
.setHeader(MqttHeaders.DUPLICATE, mqttMessage.isDuplicate())
.setHeader(MqttHeaders.RETAINED, mqttMessage.isRetained());
if (topic != null) {
messageBuilder.setHeader(MqttHeaders.TOPIC, topic);
}
@@ -75,7 +76,7 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter {
}
@Override
public <P> Object fromMessage(Message<P> message) {
public MqttMessage fromMessage(Message<?> message, Type targetClass) {
Object payload = message.getPayload();
Assert.isTrue(payload instanceof byte[] || payload instanceof String);
byte[] payloadBytes;

View File

@@ -15,8 +15,10 @@
*/
package org.springframework.integration.mqtt.support;
import org.springframework.integration.Message;
import org.springframework.integration.support.converter.MessageConverter;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.converter.MessageConverter;
/**
* Extension of {@link MessageConverter} allowing the topic to be added as
@@ -25,7 +27,7 @@ import org.springframework.integration.support.converter.MessageConverter;
* @since 1.0
*
*/
public interface MqttMessageConverter extends MessageConverter {
public interface MqttMessageConverter extends MessageConverter<MqttMessage> {
Message<String> toMessage(String topic, Object object);
Message<String> toMessage(String topic, MqttMessage mqttMessage);
}

View File

@@ -20,13 +20,14 @@ import static org.junit.Assert.assertNotNull;
import org.junit.Rule;
import org.junit.Test;
import org.springframework.integration.Message;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.message.GenericMessage;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import org.springframework.integration.mqtt.outbound.MqttPahoMessageHandler;
import org.springframework.integration.mqtt.support.MqttHeaders;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.messaging.Message;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
/**

View File

@@ -42,13 +42,14 @@ import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.junit.Test;
import org.mockito.invocation.InvocationOnMock;
import org.mockito.stubbing.Answer;
import org.springframework.integration.Message;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.message.GenericMessage;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory.Will;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import org.springframework.integration.mqtt.outbound.MqttPahoMessageHandler;
import org.springframework.messaging.Message;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
/**

View File

@@ -21,12 +21,13 @@ import static org.junit.Assert.assertSame;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.integration.MessageChannel;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.inbound.MqttPahoMessageDrivenChannelAdapter;
import org.springframework.integration.mqtt.support.MqttMessageConverter;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.messaging.MessageChannel;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;

View File

@@ -22,13 +22,13 @@ import static org.junit.Assert.assertTrue;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.integration.mqtt.core.DefaultMqttPahoClientFactory;
import org.springframework.integration.mqtt.outbound.MqttPahoMessageHandler;
import org.springframework.integration.mqtt.support.DefaultPahoMessageConverter;
import org.springframework.integration.mqtt.support.MqttMessageConverter;
import org.springframework.integration.support.converter.MessageConverter;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
@@ -74,7 +74,8 @@ public class MqttOutboundChannelAdapterParserTests {
assertEquals("bar", TestUtils.getPropertyValue(withDefaultConverterHandler, "defaultTopic"));
assertEquals(1, TestUtils.getPropertyValue(withDefaultConverterHandler, "defaultQos"));
assertTrue(TestUtils.getPropertyValue(withDefaultConverterHandler, "defaultRetained", Boolean.class));
MessageConverter defaultConverter = TestUtils.getPropertyValue(withDefaultConverterHandler, "converter", MessageConverter.class);
MqttMessageConverter defaultConverter = TestUtils.getPropertyValue(withDefaultConverterHandler, "converter",
MqttMessageConverter.class);
assertTrue(defaultConverter instanceof DefaultPahoMessageConverter);
assertEquals(1, TestUtils.getPropertyValue(defaultConverter, "defaultQos"));
assertTrue(TestUtils.getPropertyValue(defaultConverter, "defaultRetained", Boolean.class));