MQTT - Update to Spring Integration 4.0.0.M1

Conflicts:
	spring-integration-mqtt/build.gradle
This commit is contained in:
Gary Russell
2013-11-08 10:28:28 -05:00
committed by Artem Bilan
parent 46cf219b44
commit b22a130a87
6 changed files with 11 additions and 11 deletions

View File

@@ -14,7 +14,7 @@ apply plugin: 'idea'
group = 'org.springframework.integration'
repositories {
maven { url 'http://repo.springsource.org/snapshot' }
maven { url 'http://repo.springsource.org/milestone' }
maven { url 'https://repo.eclipse.org/content/repositories/paho-releases/' }
maven { url 'http://repo.springsource.org/plugins-release' }
}

View File

@@ -140,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);
MqttMessage mqttMessage = this.converter.fromMessage(message, MqttMessage.class);
MqttMessage 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

@@ -15,12 +15,11 @@
*/
package org.springframework.integration.mqtt.support;
import java.lang.reflect.Type;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHeaders;
import org.springframework.messaging.support.converter.MessageConversionException;
import org.springframework.util.Assert;
@@ -54,11 +53,12 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter {
}
@Override
@SuppressWarnings("unchecked")
public Message<?> toMessage(MqttMessage mqttMessage) {
return toMessage(null, mqttMessage);
public Message<?> toMessage(Object mqttMessage, MessageHeaders headers) {
Assert.isInstanceOf(MqttMessage.class, mqttMessage);
return toMessage(null, (MqttMessage) mqttMessage);
}
@Override
public Message<String> toMessage(String topic, MqttMessage mqttMessage) {
try {
MessageBuilder<String> messageBuilder = MessageBuilder.withPayload(new String(mqttMessage.getPayload(), this.charset))
@@ -76,7 +76,7 @@ public class DefaultPahoMessageConverter implements MqttMessageConverter {
}
@Override
public MqttMessage fromMessage(Message<?> message, Type targetClass) {
public MqttMessage fromMessage(Message<?> message, Class<?> targetClass) {
Object payload = message.getPayload();
Assert.isTrue(payload instanceof byte[] || payload instanceof String);
byte[] payloadBytes;

View File

@@ -27,7 +27,7 @@ import org.springframework.messaging.support.converter.MessageConverter;
* @since 1.0
*
*/
public interface MqttMessageConverter extends MessageConverter<MqttMessage> {
public interface MqttMessageConverter extends MessageConverter {
Message<String> toMessage(String topic, MqttMessage mqttMessage);
}

View File

@@ -22,12 +22,12 @@ import org.junit.Rule;
import org.junit.Test;
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.messaging.support.GenericMessage;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
/**

View File

@@ -44,12 +44,12 @@ import org.mockito.invocation.InvocationOnMock;
import org.mockito.stubbing.Answer;
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.messaging.support.GenericMessage;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
/**