diff --git a/spring-integration-core/src/main/java/org/springframework/integration/filter/MessageFilter.java b/spring-integration-core/src/main/java/org/springframework/integration/filter/MessageFilter.java
index f88e6e6802..ab21e3fc0f 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/filter/MessageFilter.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/filter/MessageFilter.java
@@ -19,6 +19,7 @@ package org.springframework.integration.filter;
import org.springframework.beans.factory.BeanFactoryAware;
import org.springframework.integration.core.Message;
import org.springframework.integration.core.MessageChannel;
+import org.springframework.integration.core.MessageHeaders;
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
import org.springframework.integration.message.MessageDeliveryException;
import org.springframework.integration.message.MessageRejectedException;
@@ -35,6 +36,7 @@ import org.springframework.util.Assert;
* provided, the rejected Messages will be sent to that channel.
*
* @author Mark Fisher
+ * @author Oleg Zhurakousky
*/
public class MessageFilter extends AbstractReplyProducingMessageHandler {
@@ -115,4 +117,11 @@ public class MessageFilter extends AbstractReplyProducingMessageHandler {
return null;
}
+ protected void handleResult(Object replyMessage, MessageHeaders requestHeaders, MessageChannel replyChannel) {
+ if (!this.sendReplyMessage((Message>) replyMessage, replyChannel)) {
+ throw new MessageDeliveryException((Message>) replyMessage,
+ "failed to send reply Message to channel '" + replyChannel + "'. Consider increasing the " +
+ "send timeout of this endpoint.");
+ }
+ }
}
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java
index 439023b9bc..098f3935b7 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractReplyProducingMessageHandler.java
@@ -32,6 +32,7 @@ import org.springframework.util.Assert;
*
* @author Mark Fisher
* @author Iwein Fuld
+ * @author Oleg Zhurakousky
*/
public abstract class AbstractReplyProducingMessageHandler extends AbstractMessageHandler
implements MessageProducer {
@@ -116,7 +117,7 @@ public abstract class AbstractReplyProducingMessageHandler extends AbstractMessa
@SuppressWarnings("unchecked")
private Message> createReplyMessage(Object reply, MessageHeaders requestHeaders) {
if (reply instanceof Message) {
- return (Message>) reply;
+ return MessageBuilder.fromMessage((Message>) reply).copyHeadersIfAbsent(requestHeaders).build();
}
MessageBuilder> builder = (reply instanceof MessageBuilder)
? (MessageBuilder>) reply : MessageBuilder.withPayload(reply);
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/transformer/MessageTransformingHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/transformer/MessageTransformingHandler.java
index 6e2b1f8599..f9d5c478de 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/transformer/MessageTransformingHandler.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/transformer/MessageTransformingHandler.java
@@ -17,9 +17,11 @@
package org.springframework.integration.transformer;
import org.springframework.beans.factory.BeanFactoryAware;
-
import org.springframework.integration.core.Message;
+import org.springframework.integration.core.MessageChannel;
+import org.springframework.integration.core.MessageHeaders;
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
+import org.springframework.integration.message.MessageDeliveryException;
import org.springframework.integration.message.MessageHandler;
import org.springframework.util.Assert;
@@ -29,6 +31,7 @@ import org.springframework.util.Assert;
* and sends the result to its output channel.
*
* @author Mark Fisher
+ * @author Oleg Zhurakousky
*/
public class MessageTransformingHandler extends AbstractReplyProducingMessageHandler {
@@ -69,5 +72,12 @@ public class MessageTransformingHandler extends AbstractReplyProducingMessageHan
throw new MessageTransformationException(message, e);
}
}
-
+
+ protected void handleResult(Object replyMessage, MessageHeaders requestHeaders, MessageChannel replyChannel) {
+ if (!this.sendReplyMessage((Message>) replyMessage, replyChannel)) {
+ throw new MessageDeliveryException((Message>) replyMessage,
+ "failed to send reply Message to channel '" + replyChannel + "'. Consider increasing the " +
+ "send timeout of this endpoint.");
+ }
+ }
}
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/gateway/MultipleEndpointGatewayTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/gateway/MultipleEndpointGatewayTests-context.xml
new file mode 100644
index 0000000000..36bc540f06
--- /dev/null
+++ b/spring-integration-core/src/test/java/org/springframework/integration/gateway/MultipleEndpointGatewayTests-context.xml
@@ -0,0 +1,35 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/spring-integration-core/src/test/java/org/springframework/integration/gateway/MultipleEndpointGatewayTests.java b/spring-integration-core/src/test/java/org/springframework/integration/gateway/MultipleEndpointGatewayTests.java
new file mode 100644
index 0000000000..696db434e7
--- /dev/null
+++ b/spring-integration-core/src/test/java/org/springframework/integration/gateway/MultipleEndpointGatewayTests.java
@@ -0,0 +1,72 @@
+/*
+ * Copyright 2002-2010 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.gateway;
+
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Qualifier;
+import org.springframework.integration.core.Message;
+import org.springframework.integration.message.MessageBuilder;
+import org.springframework.test.context.ContextConfiguration;
+import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
+
+/**
+ * @author Oleg Zhurakousky
+ *
+ */
+@ContextConfiguration
+@RunWith(SpringJUnit4ClassRunner.class)
+public class MultipleEndpointGatewayTests {
+
+ @Autowired
+ @Qualifier("gatewayA")
+ private SampleGateway gatewayA;
+
+ @Autowired
+ @Qualifier("gatewayB")
+ private SampleGateway gatewayB;
+
+ @Test
+ public void gatewayNoDefaultReplyChannel(){
+ gatewayA.echo("echoAsMessageChannel");
+ // there is nothing to assert. Successful execution of the above is all we care in this test
+ }
+ @Test
+ public void gatewayWithDefaultReplyChannel(){
+ gatewayB.echo("echoAsMessageChannelIgnoreDefOutChannel");
+ // there is nothing to assert. Successful execution of the above is all we care in this test
+ }
+
+ @Test
+ public void gatewayWithReplySentBackToDefaultReplyChannel(){
+ gatewayB.echo("echoAsMessageChannelDefaultOutputChannel");
+ // there is nothing to assert. Successful execution of the above is all we care in this test
+ }
+
+ public static interface SampleGateway{
+ public Object echo(Object value);
+ }
+
+ public static class SampleEchoService {
+ public Object echo(Object value){
+ return "R:" + value;
+ }
+ public Message echoAsMessage(Object value){
+ return MessageBuilder.withPayload("R:" + value).build();
+ }
+ }
+}
diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/ActiveMqTestUtils.java b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/ActiveMqTestUtils.java
new file mode 100644
index 0000000000..06f5e35e0c
--- /dev/null
+++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/ActiveMqTestUtils.java
@@ -0,0 +1,49 @@
+/*
+ * Copyright 2002-2010 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.jms.config;
+
+import java.io.File;
+
+import org.junit.Before;
+
+/**
+ * @author Oleg Zhurakousky
+ *
+ */
+public class ActiveMqTestUtils {
+
+ @Before
+ public static void prepare() {
+ System.out.println("####### Refreshing ActiveMq ########");
+ File activeMqTempDir = new File("activemq-data");
+ deleteDir(activeMqTempDir);
+
+ }
+ /*
+ *
+ */
+ private static void deleteDir(File directory){
+ if (directory.exists()){
+ String[] children = directory.list();
+ if (children != null){
+ for (int i=0; i < children.length; i++) {
+ deleteDir(new File(directory, children[i]));
+ }
+ }
+ }
+ directory.delete();
+ }
+}
diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/ExceptionHandlingSiConsumerTests.java b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/ExceptionHandlingSiConsumerTests.java
index bddb6e90dd..89c75c2f97 100644
--- a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/ExceptionHandlingSiConsumerTests.java
+++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/ExceptionHandlingSiConsumerTests.java
@@ -42,6 +42,7 @@ public class ExceptionHandlingSiConsumerTests {
@Test
public void nonSiProducer_siConsumer_sync_withReturn() throws Exception {
+ ActiveMqTestUtils.prepare();
ConfigurableApplicationContext applicationContext = new ClassPathXmlApplicationContext("Exception-nonSiProducer-siConsumer.xml", ExceptionHandlingSiConsumerTests.class);
JmsTemplate jmsTemplate = new JmsTemplate(applicationContext.getBean("connectionFactory", ConnectionFactory.class));
Destination request = applicationContext.getBean("requestQueue", Destination.class);
@@ -61,6 +62,7 @@ public class ExceptionHandlingSiConsumerTests {
}
@Test
public void nonSiProducer_siConsumer_sync_withReturnNoException() throws Exception {
+ ActiveMqTestUtils.prepare();
ConfigurableApplicationContext applicationContext = new ClassPathXmlApplicationContext("Exception-nonSiProducer-siConsumer.xml", ExceptionHandlingSiConsumerTests.class);
JmsTemplate jmsTemplate = new JmsTemplate(applicationContext.getBean("connectionFactory", ConnectionFactory.class));
Destination request = applicationContext.getBean("requestQueue", Destination.class);
@@ -81,6 +83,7 @@ public class ExceptionHandlingSiConsumerTests {
@Test
public void nonSiProducer_siConsumer_sync_withOutboundGateway() throws Exception{
+ ActiveMqTestUtils.prepare();
final ConfigurableApplicationContext applicationContext = new ClassPathXmlApplicationContext("Exception-nonSiProducer-siConsumer.xml", ExceptionHandlingSiConsumerTests.class);
SampleGateway gateway = applicationContext.getBean("sampleGateway", SampleGateway.class);
String reply = gateway.echo("echoWithExceptionChannel");
@@ -88,30 +91,6 @@ public class ExceptionHandlingSiConsumerTests {
applicationContext.close();
}
-
-
- @Before
- public void prepare() throws Exception {
- System.out.println("####### Refreshing ActiveMq ########");
- File activeMqTempDir = new File("activemq-data");
- this.deleteDir(activeMqTempDir);
-
- }
- /*
- *
- */
- private void deleteDir(File directory){
- if (directory.exists()){
- String[] children = directory.list();
- if (children != null){
- for (int i=0; i < children.length; i++) {
- deleteDir(new File(directory, children[i]));
- }
- }
- }
- directory.delete();
- }
-
public static class SampleService{
public String echoWithException(String value){
throw new SampleException("echoWithException");
diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/JmsOutboundGatewayParserTests.java b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/JmsOutboundGatewayParserTests.java
index c19755b0cf..bc85813b62 100644
--- a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/JmsOutboundGatewayParserTests.java
+++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/JmsOutboundGatewayParserTests.java
@@ -29,6 +29,7 @@ import org.springframework.jms.support.converter.MessageConverter;
/**
* @author Jonas Partner
+ * @author Oleg Zhurakousky
*/
public class JmsOutboundGatewayParserTests {
@@ -54,5 +55,25 @@ public class JmsOutboundGatewayParserTests {
Object order = accessor.getPropertyValue("order");
assertEquals(99, order);
}
-
+
+ @Test
+ public void gatewayMaintainsReplyChannel() {
+ ActiveMqTestUtils.prepare();
+ ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext(
+ "gatewayMaintainsReplyChannel.xml", this.getClass());
+ SampleGateway gateway = context.getBean("gateway", SampleGateway.class);
+ String result = gateway.echo("hello");
+ assertEquals("HELLO", result);
+ }
+
+ public static interface SampleGateway{
+ public String echo(String value);
+ }
+
+ public static class SampleService{
+ public String echo(String value){
+ return value.toUpperCase();
+ }
+ }
+
}
diff --git a/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/gatewayMaintainsReplyChannel.xml b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/gatewayMaintainsReplyChannel.xml
new file mode 100644
index 0000000000..b379e685ba
--- /dev/null
+++ b/spring-integration-jms/src/test/java/org/springframework/integration/jms/config/gatewayMaintainsReplyChannel.xml
@@ -0,0 +1,52 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+