diff --git a/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java b/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java index fc84e797aa..bb55ece8f7 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/gateway/MessagingGatewaySupport.java @@ -507,6 +507,7 @@ public abstract class MessagingGatewaySupport extends AbstractEndpoint else { requestMessage = (object instanceof Message) ? (Message) object : this.requestMapper.toMessage(object); + Assert.state(requestMessage != null, () -> "request mapper resulted in no message for " + object); requestMessage = this.historyWritingPostProcessor.postProcessMessage(requestMessage); reply = this.messagingTemplate.sendAndReceive(requestChannel, requestMessage); if (reply instanceof ErrorMessage) { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/mapping/BytesMessageMapper.java b/spring-integration-core/src/main/java/org/springframework/integration/mapping/BytesMessageMapper.java index 88a6f6604b..b6657d46cc 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/mapping/BytesMessageMapper.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/mapping/BytesMessageMapper.java @@ -16,6 +16,12 @@ package org.springframework.integration.mapping; +import java.util.Map; + +import org.springframework.lang.NonNull; +import org.springframework.lang.Nullable; +import org.springframework.messaging.Message; + /** * An {@link OutboundMessageMapper} and {@link InboundMessageMapper} that * maps to/from {@code byte[]}. @@ -26,4 +32,14 @@ package org.springframework.integration.mapping; */ public interface BytesMessageMapper extends InboundMessageMapper, OutboundMessageMapper { + @Override + @NonNull // override + default Message toMessage(byte[] object) { + return toMessage(object, null); + } + + @Override + @NonNull // override + Message toMessage(byte[] bytes, @Nullable Map headers); + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/mapping/InboundMessageMapper.java b/spring-integration-core/src/main/java/org/springframework/integration/mapping/InboundMessageMapper.java index 8e18861745..58875b78b3 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/mapping/InboundMessageMapper.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/mapping/InboundMessageMapper.java @@ -26,6 +26,7 @@ import org.springframework.messaging.Message; * * @author Mark Fisher * @author Artem Bilan + * @author Gary Russell */ @FunctionalInterface public interface InboundMessageMapper { @@ -36,7 +37,8 @@ public interface InboundMessageMapper { * @return the message as a result of mapping * @throws Exception the exception thrown by the underlying mapper implementation */ - default Message toMessage(T object) throws Exception { + @Nullable + default Message toMessage(T object) throws Exception { // NOSONAR - TODO remove Exception in 5.2 return toMessage(object, null); } @@ -50,6 +52,6 @@ public interface InboundMessageMapper { * @since 5.0 */ @Nullable - Message toMessage(T object, @Nullable Map headers) throws Exception; + Message toMessage(T object, @Nullable Map headers) throws Exception; // NOSONAR } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/json/EmbeddedJsonHeadersMessageMapper.java b/spring-integration-core/src/main/java/org/springframework/integration/support/json/EmbeddedJsonHeadersMessageMapper.java index 2325d5f001..e956be9f0e 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/json/EmbeddedJsonHeadersMessageMapper.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/json/EmbeddedJsonHeadersMessageMapper.java @@ -16,6 +16,7 @@ package org.springframework.integration.support.json; +import java.io.IOException; import java.nio.ByteBuffer; import java.util.Arrays; import java.util.Collection; @@ -229,11 +230,12 @@ public class EmbeddedJsonHeadersMessageMapper implements BytesMessageMapper { return message; } else { - return new GenericMessage<>(bytes, headers); + return headers == null ? new GenericMessage<>(bytes) : new GenericMessage<>(bytes, headers); } } - private Message decodeNativeFormat(byte[] bytes, Map headersToAdd) throws Exception { + @Nullable + private Message decodeNativeFormat(byte[] bytes, @Nullable Map headersToAdd) throws IOException { ByteBuffer buffer = ByteBuffer.wrap(bytes); if (buffer.remaining() > 4) { int headersLen = buffer.getInt(); diff --git a/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/IntegrationMBeanExporter.java b/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/IntegrationMBeanExporter.java index 90036a2616..cf34c5ea37 100644 --- a/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/IntegrationMBeanExporter.java +++ b/spring-integration-jmx/src/main/java/org/springframework/integration/monitor/IntegrationMBeanExporter.java @@ -846,9 +846,6 @@ public class IntegrationMBeanExporter extends MBeanExporter implements Applicati return bean; } Advised advised = (Advised) bean; - if (advised.getTargetSource() == null) { - return null; - } try { return extractTarget(advised.getTargetSource().getTarget()); } @@ -1075,7 +1072,7 @@ public class IntegrationMBeanExporter extends MBeanExporter implements Applicati if (target instanceof MessagingGatewaySupport) { outputChannel = ((MessagingGatewaySupport) target).getRequestChannel(); } - else { + else if (target instanceof SourcePollingChannelAdapter) { outputChannel = ((SourcePollingChannelAdapter) target).getOutputChannel(); } diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java index ae05cd8cba..06ae2a5a45 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java @@ -201,7 +201,7 @@ public abstract class AbstractConfigurableMongoDbMessageStore extends AbstractMe new Update().inc(MessageDocumentFields.SEQUENCE, 1), FindAndModifyOptions.options().returnNew(true).upsert(true), Map.class, this.collectionName) - .get(MessageDocumentFields.SEQUENCE); + .get(MessageDocumentFields.SEQUENCE); // NOSONAR - never returns null } protected void addMessageDocument(final MessageDocument document) { diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java index 87b2c34b2c..fcc0e8a11c 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java @@ -473,7 +473,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore new Update().inc(SEQUENCE, 1), FindAndModifyOptions.options().returnNew(true).upsert(true), Map.class, - this.collectionName).get(SEQUENCE); + this.collectionName).get(SEQUENCE); // NOSONAR - never returns null } @SuppressWarnings("unchecked") @@ -481,8 +481,14 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore Map innerMap = (Map) new DirectFieldAccessor(messageHeaders).getPropertyValue("headers"); // using reflection to set ID and TIMESTAMP since they are immutable through MessageHeaders - innerMap.put(MessageHeaders.ID, headers.get(MessageHeaders.ID)); - innerMap.put(MessageHeaders.TIMESTAMP, headers.get(MessageHeaders.TIMESTAMP)); + Object idHeader = headers.get(MessageHeaders.ID); + if (idHeader != null) { + innerMap.put(MessageHeaders.ID, idHeader); + } + Object tsHeader = headers.get(MessageHeaders.TIMESTAMP); + if (tsHeader != null) { + innerMap.put(MessageHeaders.TIMESTAMP, tsHeader); + } } @SuppressWarnings("unchecked") @@ -760,7 +766,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore MongoDbMessageStore.this.converter.normalizeHeaders((Map) source.get("headers")); Object payload = this.deserializingConverter.convert(((Binary) source.get("payload")).getData()); - ErrorMessage message = new ErrorMessage((Throwable) payload, headers); + ErrorMessage message = new ErrorMessage((Throwable) payload, headers); // NOSONAR not null enhanceHeaders(message.getHeaders(), headers); return message; diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MongoDbMessageBytesConverter.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MongoDbMessageBytesConverter.java index 9ce5550535..e2f253f3b1 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MongoDbMessageBytesConverter.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MongoDbMessageBytesConverter.java @@ -59,6 +59,9 @@ public class MongoDbMessageBytesConverter implements GenericConverter { @Override public Object convert(Object source, TypeDescriptor sourceType, TypeDescriptor targetType) { + if (source == null) { + return null; + } if (Message.class.isAssignableFrom(sourceType.getObjectType())) { return new Binary(this.serializingConverter.convert(source)); }