From 46cf219b448c8c56ff4209d23f766664962a4a17 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 15 Oct 2013 10:47:20 -0400 Subject: [PATCH] MQTT - Use INT/SPR 4.0.0 Conflicts: spring-integration-mqtt/build.gradle --- spring-integration-mqtt/build.gradle | 5 ++-- spring-integration-mqtt/gradle.properties | 2 +- .../MqttPahoMessageDrivenChannelAdapter.java | 2 +- .../outbound/AbstractMqttMessageHandler.java | 12 +++++---- .../mqtt/outbound/MqttPahoMessageHandler.java | 3 ++- .../support/DefaultPahoMessageConverter.java | 25 ++++++++++--------- .../mqtt/support/MqttMessageConverter.java | 10 +++++--- .../mqtt/BackTobackAdapterTests.java | 3 ++- .../integration/mqtt/MqttAdapterTests.java | 3 ++- ...essageDrivenChannelAdapterParserTests.java | 3 ++- ...MqttOutboundChannelAdapterParserTests.java | 5 ++-- 11 files changed, 41 insertions(+), 32 deletions(-) diff --git a/spring-integration-mqtt/build.gradle b/spring-integration-mqtt/build.gradle index 957372d..a966b94 100644 --- a/spring-integration-mqtt/build.gradle +++ b/spring-integration-mqtt/build.gradle @@ -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' diff --git a/spring-integration-mqtt/gradle.properties b/spring-integration-mqtt/gradle.properties index bebfcbc..2db7ae7 100644 --- a/spring-integration-mqtt/gradle.properties +++ b/spring-integration-mqtt/gradle.properties @@ -1 +1 @@ -version=1.0.0.BUILD-SNAPSHOT +version=4.0.0.BUILD-SNAPSHOT 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 d4305c7..99be20b 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 @@ -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. diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java index 12a4ef2..177a0ac 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/outbound/AbstractMqttMessageHandler.java @@ -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"); 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 810aae0..e0570be 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 @@ -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; /** diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/DefaultPahoMessageConverter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/DefaultPahoMessageConverter.java index 20f8137..1d87534 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/DefaultPahoMessageConverter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/DefaultPahoMessageConverter.java @@ -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

Message

toMessage(Object object) { - return (Message

) toMessage(null, object); + public Message toMessage(MqttMessage mqttMessage) { + return toMessage(null, mqttMessage); } - public Message toMessage(String topic, Object object) { - Assert.isInstanceOf(MqttMessage.class, object); - MqttMessage message = (MqttMessage) object; + public Message toMessage(String topic, MqttMessage mqttMessage) { try { - MessageBuilder 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 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

Object fromMessage(Message

message) { + public MqttMessage fromMessage(Message message, Type targetClass) { Object payload = message.getPayload(); Assert.isTrue(payload instanceof byte[] || payload instanceof String); byte[] payloadBytes; diff --git a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/MqttMessageConverter.java b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/MqttMessageConverter.java index 3d13741..61c586a 100644 --- a/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/MqttMessageConverter.java +++ b/spring-integration-mqtt/src/main/java/org/springframework/integration/mqtt/support/MqttMessageConverter.java @@ -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 { - Message toMessage(String topic, Object object); + Message toMessage(String topic, MqttMessage mqttMessage); } diff --git a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackTobackAdapterTests.java b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackTobackAdapterTests.java index 66e2ba3..0c4f31e 100644 --- a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackTobackAdapterTests.java +++ b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/BackTobackAdapterTests.java @@ -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; /** diff --git a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/MqttAdapterTests.java b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/MqttAdapterTests.java index 92c44c4..c3cd912 100644 --- a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/MqttAdapterTests.java +++ b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/MqttAdapterTests.java @@ -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; /** diff --git a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/config/xml/MqttMessageDrivenChannelAdapterParserTests.java b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/config/xml/MqttMessageDrivenChannelAdapterParserTests.java index 5b3452f..00d4fe5 100644 --- a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/config/xml/MqttMessageDrivenChannelAdapterParserTests.java +++ b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/config/xml/MqttMessageDrivenChannelAdapterParserTests.java @@ -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; diff --git a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/config/xml/MqttOutboundChannelAdapterParserTests.java b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/config/xml/MqttOutboundChannelAdapterParserTests.java index 71a752b..bea91f5 100644 --- a/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/config/xml/MqttOutboundChannelAdapterParserTests.java +++ b/spring-integration-mqtt/src/test/java/org/springframework/integration/mqtt/config/xml/MqttOutboundChannelAdapterParserTests.java @@ -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));