diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java index b7a6787873..620bf57600 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java @@ -146,6 +146,7 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS protected void doAddMessage(Message message) { Assert.notNull(message, "'message' must not be null"); UUID messageId = message.getHeaders().getId(); + Assert.notNull(messageId, "Cannot store messages without an ID header"); doStoreIfAbsent(this.messagePrefix + messageId, new MessageHolder(message)); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java b/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java index ee84507709..7317a58789 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2018 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -182,7 +182,9 @@ public class SimpleMessageStore extends AbstractMessageGroupStore + this.individualCapacity + "), try constructing it with a larger capacity."); } - this.idToMessage.put(message.getHeaders().getId(), message); + UUID id = message.getHeaders().getId(); + Assert.notNull(id, "ID header must not be null"); + this.idToMessage.put(id, message); return message; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/transformer/ClaimCheckInTransformer.java b/spring-integration-core/src/main/java/org/springframework/integration/transformer/ClaimCheckInTransformer.java index 17367d6ef5..ca3ff8820d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/transformer/ClaimCheckInTransformer.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/transformer/ClaimCheckInTransformer.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2014 the original author or authors. + * Copyright 2002-2018 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -16,6 +16,8 @@ package org.springframework.integration.transformer; +import java.util.UUID; + import org.springframework.integration.store.MessageStore; import org.springframework.integration.support.AbstractIntegrationMessageBuilder; import org.springframework.messaging.Message; @@ -51,10 +53,10 @@ public class ClaimCheckInTransformer extends AbstractTransformer { @Override protected Object doTransform(Message message) throws Exception { Assert.notNull(message, "message must not be null"); - Object payload = message.getPayload(); - Assert.notNull(payload, "payload must not be null"); - Message storedMessage = this.messageStore.addMessage(message); - AbstractIntegrationMessageBuilder responseBuilder = this.getMessageBuilderFactory().withPayload(storedMessage.getHeaders().getId()); + UUID id = message.getHeaders().getId(); + Assert.notNull(id, "ID header must not be null"); + this.messageStore.addMessage(message); + AbstractIntegrationMessageBuilder responseBuilder = getMessageBuilderFactory().withPayload(id); // headers on the 'current' message take precedence responseBuilder.copyHeaders(message.getHeaders()); return responseBuilder.build(); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/transformer/ContentEnricher.java b/spring-integration-core/src/main/java/org/springframework/integration/transformer/ContentEnricher.java index 6556cdf307..472f133372 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/transformer/ContentEnricher.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/transformer/ContentEnricher.java @@ -375,7 +375,9 @@ public class ContentEnricher extends AbstractReplyProducingMessageHandler implem else { final Object requestMessagePayload = this.requestPayloadExpression.getValue(this.sourceEvaluationContext, requestMessage); - actualRequestMessage = this.getMessageBuilderFactory().withPayload(requestMessagePayload) + Assert.state(requestMessagePayload != null, + () -> "Request payload expression produced null for " + requestMessage); + actualRequestMessage = getMessageBuilderFactory().withPayload(requestMessagePayload) .copyHeaders(requestMessage.getHeaders()).build(); } final Message replyMessage; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/util/DynamicPeriodicTrigger.java b/spring-integration-core/src/main/java/org/springframework/integration/util/DynamicPeriodicTrigger.java index 93d87d8e46..922b9407e7 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/util/DynamicPeriodicTrigger.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/util/DynamicPeriodicTrigger.java @@ -217,13 +217,14 @@ public class DynamicPeriodicTrigger implements Trigger { */ @Override public Date nextExecutionTime(TriggerContext triggerContext) { - if (triggerContext.lastScheduledExecutionTime() == null) { + Date lastScheduled = triggerContext.lastScheduledExecutionTime(); + if (lastScheduled == null) { return new Date(System.currentTimeMillis() + this.initialDuration.toMillis()); } else if (this.fixedRate) { - return new Date(triggerContext.lastScheduledExecutionTime().getTime() + this.duration.toMillis()); + return new Date(lastScheduled.getTime() + this.duration.toMillis()); } - return new Date(triggerContext.lastCompletionTime().getTime() + this.duration.toMillis()); + return new Date(triggerContext.lastCompletionTime().getTime() + this.duration.toMillis()); // NOSONAR never null here } @Override diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java index b64e581387..2db6b50850 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcMessageStore.java @@ -327,6 +327,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa @SuppressWarnings("unchecked") public Message addMessage(final Message message) { UUID id = message.getHeaders().getId(); + Assert.notNull(id, "Cannot store messages without an ID header"); final String messageId = getKey(id); final byte[] messageBytes = this.serializer.convert(message); diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MessageDocument.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MessageDocument.java index f076ea6dba..ebcf980b9d 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MessageDocument.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MessageDocument.java @@ -77,6 +77,7 @@ public class MessageDocument { @PersistenceConstructor MessageDocument(Message message, UUID messageId) { Assert.notNull(message, "'message' must not be null"); + Assert.notNull(messageId, "'message' ID header must not be null"); this.message = message; this.messageId = messageId; } 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 fcc0e8a11c..47f9b35210 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 @@ -212,6 +212,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore private void addMessageDocument(MessageWrapper document) { UUID messageId = (UUID) document.headers.get(MessageHeaders.ID); + Assert.notNull(messageId, "ID header must not be null"); Query query = whereMessageIdIsAndGroupIdIs(messageId, document.get_GroupId()); if (!this.template.exists(query, MessageWrapper.class, this.collectionName)) { if (document.get_Group_timestamp() == 0) {