diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageProducerSupport.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageProducerSupport.java index 1b7e9ff212..7eb6eb3340 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageProducerSupport.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/MessageProducerSupport.java @@ -18,9 +18,11 @@ package org.springframework.integration.endpoint; import org.springframework.integration.Message; import org.springframework.integration.MessageChannel; +import org.springframework.integration.MessagingException; import org.springframework.integration.core.MessageProducer; import org.springframework.integration.core.MessagingTemplate; import org.springframework.integration.history.MessageHistory; +import org.springframework.integration.history.TrackableComponent; import org.springframework.util.Assert; /** @@ -29,10 +31,12 @@ import org.springframework.util.Assert; * * @author Mark Fisher */ -public abstract class MessageProducerSupport extends AbstractEndpoint implements MessageProducer { +public abstract class MessageProducerSupport extends AbstractEndpoint implements MessageProducer, TrackableComponent { private volatile MessageChannel outputChannel; + private volatile boolean shouldTrack = false; + private final MessagingTemplate messagingTemplate = new MessagingTemplate(); @@ -44,13 +48,20 @@ public abstract class MessageProducerSupport extends AbstractEndpoint implements this.messagingTemplate.setSendTimeout(sendTimeout); } + public void setShouldTrack(boolean shouldTrack) { + this.shouldTrack = shouldTrack; + } + @Override protected void onInit() { Assert.notNull(this.outputChannel, "outputChannel is required"); } protected void sendMessage(Message message) { - if (message != null) { + if (message == null) { + throw new MessagingException("cannot send a null message"); + } + if (this.shouldTrack) { message = MessageHistory.write(message, this); } this.messagingTemplate.send(this.outputChannel, message);