MQTT - Update to Spring Integration 4.0.0.M1

This commit is contained in:
Gary Russell
2013-11-08 10:28:28 -05:00
parent f316ba86dd
commit 89cf737ad4
5 changed files with 10 additions and 10 deletions

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;
/**