From c38b0e18325895679f14aec7f27e6d49e7d870e5 Mon Sep 17 00:00:00 2001 From: Iwein Fuld Date: Sat, 10 Oct 2009 18:14:08 +0000 Subject: [PATCH] https://jira.springsource.org/browse/INT-838 Added channel to the message, added testcase for this branch. --- .../AbstractReplyProducingMessageHandler.java | 143 ++++++++++-------- ...tractReplyProducingMessageHandlerTest.java | 58 +++++++ 2 files changed, 135 insertions(+), 66 deletions(-) create mode 100644 org.springframework.integration/src/test/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandlerTest.java diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java b/org.springframework.integration/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java index 91e739df27..45b5345f15 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java @@ -19,7 +19,6 @@ package org.springframework.integration.handler; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.BeanFactoryAware; import org.springframework.integration.channel.BeanFactoryChannelResolver; -import org.springframework.integration.channel.ChannelResolutionException; import org.springframework.integration.channel.ChannelResolver; import org.springframework.integration.channel.MessageChannelTemplate; import org.springframework.integration.core.Message; @@ -32,91 +31,103 @@ import org.springframework.util.Assert; /** * Base class for MessageHandlers that are capable of producing replies. - * + * * @author Mark Fisher + * @author Iwein Fuld */ public abstract class AbstractReplyProducingMessageHandler extends AbstractMessageHandler implements BeanFactoryAware { - public static final long DEFAULT_SEND_TIMEOUT = 1000; + public static final long DEFAULT_SEND_TIMEOUT = 1000; - private MessageChannel outputChannel; + private MessageChannel outputChannel; - private volatile ChannelResolver channelResolver; + private volatile ChannelResolver channelResolver; - private volatile boolean requiresReply = false; + private volatile boolean requiresReply = false; - private final MessageChannelTemplate channelTemplate; + private final MessageChannelTemplate channelTemplate; - public AbstractReplyProducingMessageHandler() { - this.channelTemplate = new MessageChannelTemplate(); - this.channelTemplate.setSendTimeout(DEFAULT_SEND_TIMEOUT); - } + public AbstractReplyProducingMessageHandler() { + this.channelTemplate = new MessageChannelTemplate(); + this.channelTemplate.setSendTimeout(DEFAULT_SEND_TIMEOUT); + } - public void setOutputChannel(MessageChannel outputChannel) { - this.outputChannel = outputChannel; - } + public void setOutputChannel(MessageChannel outputChannel) { + this.outputChannel = outputChannel; + } - protected MessageChannel getOutputChannel() { - return this.outputChannel; - } + protected MessageChannel getOutputChannel() { + return this.outputChannel; + } - /** - * Set the timeout for sending reply Messages. - */ - public void setSendTimeout(long sendTimeout) { - this.channelTemplate.setSendTimeout(sendTimeout); - } + /** + * Set the timeout for sending reply Messages. + */ + public void setSendTimeout(long sendTimeout) { + this.channelTemplate.setSendTimeout(sendTimeout); + } - public void setChannelResolver(ChannelResolver channelResolver) { - Assert.notNull(channelResolver, "channelResolver must not be null"); - this.channelResolver = channelResolver; - } + /** + * Set the ChannelResolver to be used when there is no default output channel. + */ + public void setChannelResolver(ChannelResolver channelResolver) { + Assert.notNull(channelResolver, "'channelResolver' must not be null"); + this.channelResolver = channelResolver; + } - public void setRequiresReply(boolean requiresReply) { - this.requiresReply = requiresReply; - } + /** + * Flag wether reply is required. If true an incoming message MUST result in a reply message being sent. + * If false an incoming message MAY result in a reply message being sent + */ + public void setRequiresReply(boolean requiresReply) { + this.requiresReply = requiresReply; + } - public void setBeanFactory(BeanFactory beanFactory) { - if (this.channelResolver == null) { - this.channelResolver = new BeanFactoryChannelResolver(beanFactory); - } - } + public void setBeanFactory(BeanFactory beanFactory) { + if (this.channelResolver == null) { + this.channelResolver = new BeanFactoryChannelResolver(beanFactory); + } + } - @Override - protected final void handleMessageInternal(Message message) { - ReplyMessageHolder replyMessageHolder = new ReplyMessageHolder(); - this.handleRequestMessage(message, replyMessageHolder); - if (replyMessageHolder.isEmpty()) { - if (this.requiresReply) { - throw new MessageHandlingException(message, "handler '" + this - + "' requires a reply, but no reply was received"); - } - if (logger.isDebugEnabled()) { - logger.debug("handler '" + this + "' produced no reply for request Message: " + message); - } - return; - } - MessageChannel replyChannel = resolveReplyChannel(message, this.outputChannel, this.channelResolver); - MessageHeaders requestHeaders = message.getHeaders(); - for (MessageBuilder builder : replyMessageHolder.builders()) { - builder.copyHeadersIfAbsent(requestHeaders); - Message replyMessage = builder.build(); - if (!this.sendReplyMessage(replyMessage, replyChannel)) { - throw new MessageDeliveryException(replyMessage, "failed to send reply Message"); - } - } - } + /** + * {@inheritDoc} + */ + @Override + protected final void handleMessageInternal(Message message) { + ReplyMessageHolder replyMessageHolder = new ReplyMessageHolder(); + this.handleRequestMessage(message, replyMessageHolder); + if (replyMessageHolder.isEmpty()) { + if (this.requiresReply) { + throw new MessageHandlingException(message, "handler '" + this + + "' requires a reply, but no reply was received"); + } + if (logger.isDebugEnabled()) { + logger.debug("handler '" + this + "' produced no reply for request Message: " + message); + } + return; + } + MessageChannel replyChannel = resolveReplyChannel(message, this.outputChannel, this.channelResolver); + MessageHeaders requestHeaders = message.getHeaders(); + for (MessageBuilder builder : replyMessageHolder.builders()) { + builder.copyHeadersIfAbsent(requestHeaders); + Message replyMessage = builder.build(); + if (!this.sendReplyMessage(replyMessage, replyChannel)) { + throw new MessageDeliveryException(replyMessage, + "failed to send reply Message to channel '" + replyChannel + "'"); + } + } + } - protected abstract void handleRequestMessage(Message requestMessage, ReplyMessageHolder replyMessageHolder); + protected abstract void handleRequestMessage(Message requestMessage, ReplyMessageHolder replyMessageHolder); - protected boolean sendReplyMessage(Message replyMessage, MessageChannel replyChannel) { - if (logger.isDebugEnabled()) { - logger.debug("handler '" + this + "' sending reply Message: " + replyMessage); - } - return this.channelTemplate.send(replyMessage, replyChannel); - } + protected boolean sendReplyMessage(Message replyMessage, MessageChannel replyChannel) { + if (logger.isDebugEnabled()) { + logger.debug("handler '" + this + "' sending reply Message: " + replyMessage); + } + return this.channelTemplate.send(replyMessage, replyChannel); + } } diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandlerTest.java b/org.springframework.integration/src/test/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandlerTest.java new file mode 100644 index 0000000000..f9c83e401c --- /dev/null +++ b/org.springframework.integration/src/test/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandlerTest.java @@ -0,0 +1,58 @@ +/* + * Copyright 2002-2008 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.integration.handler; + +import org.junit.Test; +import org.junit.matchers.JUnitMatchers; +import static org.junit.Assert.assertThat; +import static org.junit.Assert.fail; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import static org.mockito.Mockito.when; +import org.springframework.integration.core.Message; +import org.springframework.integration.core.MessageChannel; +import org.springframework.integration.core.MessagingException; +import org.springframework.integration.message.MessageBuilder; + +/** + * @author Iwein Fuld + */ +@RunWith(org.mockito.runners.MockitoJUnitRunner.class) +public class AbstractReplyProducingMessageHandlerTest { + + private AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() { + @Override + protected void handleRequestMessage(Message requestMessage, ReplyMessageHolder replyMessageHolder) { + replyMessageHolder.add(requestMessage); + } + }; + private Message message = MessageBuilder.withPayload("test").build(); + @Mock + private MessageChannel channel=null; + + @Test + public void errorMessageShouldContainChannelName() { + handler.setOutputChannel(channel); + when(channel.send(message)).thenReturn(false); + when(channel.toString()).thenReturn("testChannel"); + try { + handler.handleMessage(message); + fail("Expected a MessagingException"); + } catch (MessagingException e) { + assertThat(e.getMessage(), JUnitMatchers.containsString("'testChannel'")); + } + } +}