INT-904, INT-907 added exception mapper support to AbstractMessagingGateway, Changed ChannelPublishingMessageListener to extend form AbstractMesagingGateway, added namespace support, tests, etc...
This commit is contained in:
@@ -39,7 +39,7 @@ import org.w3c.dom.Element;
|
|||||||
public class GatewayParser extends AbstractSimpleBeanDefinitionParser {
|
public class GatewayParser extends AbstractSimpleBeanDefinitionParser {
|
||||||
|
|
||||||
private static String[] referenceAttributes = new String[] {
|
private static String[] referenceAttributes = new String[] {
|
||||||
"default-request-channel", "default-reply-channel", "message-mapper"
|
"default-request-channel", "default-reply-channel", "message-mapper", "exception-mapper"
|
||||||
};
|
};
|
||||||
|
|
||||||
private static String[] innerAttributes = new String[] {
|
private static String[] innerAttributes = new String[] {
|
||||||
|
|||||||
@@ -27,6 +27,8 @@ import org.springframework.integration.endpoint.PollingConsumer;
|
|||||||
import org.springframework.integration.endpoint.EventDrivenConsumer;
|
import org.springframework.integration.endpoint.EventDrivenConsumer;
|
||||||
import org.springframework.integration.handler.BridgeHandler;
|
import org.springframework.integration.handler.BridgeHandler;
|
||||||
import org.springframework.integration.message.ErrorMessage;
|
import org.springframework.integration.message.ErrorMessage;
|
||||||
|
import org.springframework.integration.message.InboundMessageMapper;
|
||||||
|
import org.springframework.integration.message.MessageBuilder;
|
||||||
import org.springframework.integration.message.MessageHandler;
|
import org.springframework.integration.message.MessageHandler;
|
||||||
import org.springframework.integration.message.MessageDeliveryException;
|
import org.springframework.integration.message.MessageDeliveryException;
|
||||||
import org.springframework.scheduling.support.PeriodicTrigger;
|
import org.springframework.scheduling.support.PeriodicTrigger;
|
||||||
@@ -41,6 +43,8 @@ import org.springframework.util.Assert;
|
|||||||
* @author Mark Fisher
|
* @author Mark Fisher
|
||||||
*/
|
*/
|
||||||
public abstract class AbstractMessagingGateway extends AbstractEndpoint {
|
public abstract class AbstractMessagingGateway extends AbstractEndpoint {
|
||||||
|
|
||||||
|
private volatile InboundMessageMapper<Throwable> exceptionMapper;
|
||||||
|
|
||||||
private volatile MessageChannel requestChannel;
|
private volatile MessageChannel requestChannel;
|
||||||
|
|
||||||
@@ -120,7 +124,7 @@ public abstract class AbstractMessagingGateway extends AbstractEndpoint {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public void send(Object object) {
|
protected void send(Object object) {
|
||||||
this.initializeIfNecessary();
|
this.initializeIfNecessary();
|
||||||
Assert.state(this.requestChannel != null,
|
Assert.state(this.requestChannel != null,
|
||||||
"send is not supported, because no request channel has been configured");
|
"send is not supported, because no request channel has been configured");
|
||||||
@@ -131,7 +135,7 @@ public abstract class AbstractMessagingGateway extends AbstractEndpoint {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public Object receive() {
|
protected Object receive() {
|
||||||
this.initializeIfNecessary();
|
this.initializeIfNecessary();
|
||||||
Assert.state(this.replyChannel != null && (this.replyChannel instanceof PollableChannel),
|
Assert.state(this.replyChannel != null && (this.replyChannel instanceof PollableChannel),
|
||||||
"receive is not supported, because no pollable reply channel has been configured");
|
"receive is not supported, because no pollable reply channel has been configured");
|
||||||
@@ -147,15 +151,15 @@ public abstract class AbstractMessagingGateway extends AbstractEndpoint {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
public Object sendAndReceive(Object object) {
|
protected Object sendAndReceive(Object object) {
|
||||||
return this.sendAndReceive(object, true);
|
return this.sendAndReceive(object, true);
|
||||||
}
|
}
|
||||||
|
|
||||||
public Message<?> sendAndReceiveMessage(Object object) {
|
protected Message<?> sendAndReceiveMessage(Object object) {
|
||||||
return (Message<?>) this.sendAndReceive(object, false);
|
return (Message<?>) this.sendAndReceive(object, false);
|
||||||
}
|
}
|
||||||
|
|
||||||
private Object sendAndReceive(Object object, boolean shouldMapMessage) {
|
Object sendAndReceive(Object object, boolean shouldMapMessage) {
|
||||||
Message<?> request = this.toMessage(object);
|
Message<?> request = this.toMessage(object);
|
||||||
Message<?> reply = this.sendAndReceiveMessage(request);
|
Message<?> reply = this.sendAndReceiveMessage(request);
|
||||||
if (!shouldMapMessage) {
|
if (!shouldMapMessage) {
|
||||||
@@ -174,7 +178,19 @@ public abstract class AbstractMessagingGateway extends AbstractEndpoint {
|
|||||||
if (this.replyChannel != null && this.replyMessageCorrelator == null) {
|
if (this.replyChannel != null && this.replyMessageCorrelator == null) {
|
||||||
this.registerReplyMessageCorrelator();
|
this.registerReplyMessageCorrelator();
|
||||||
}
|
}
|
||||||
Message<?> reply = this.channelTemplate.sendAndReceive(message, this.requestChannel);
|
Message<?> reply = null;
|
||||||
|
try {
|
||||||
|
reply = this.channelTemplate.sendAndReceive(message, this.requestChannel);
|
||||||
|
} catch (Exception e) {
|
||||||
|
logger.warn("Execution of endpoint by the MessageListener resulted in : " + e);
|
||||||
|
reply = this.toMessage(e);
|
||||||
|
if (reply == null){ // if reply wasn't mapped re-throw
|
||||||
|
if (e instanceof RuntimeException){
|
||||||
|
throw (RuntimeException)e;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
if (reply != null && this.shouldThrowErrors && reply instanceof ErrorMessage) {
|
if (reply != null && this.shouldThrowErrors && reply instanceof ErrorMessage) {
|
||||||
Throwable error = ((ErrorMessage) reply).getPayload();
|
Throwable error = ((ErrorMessage) reply).getPayload();
|
||||||
if (error instanceof RuntimeException) {
|
if (error instanceof RuntimeException) {
|
||||||
@@ -226,10 +242,29 @@ public abstract class AbstractMessagingGateway extends AbstractEndpoint {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public InboundMessageMapper<Throwable> getExceptionMapper() {
|
||||||
|
return exceptionMapper;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setExceptionMapper(InboundMessageMapper<Throwable> exceptionMapper) {
|
||||||
|
this.exceptionMapper = exceptionMapper;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Subclasses must implement this to map from an Object to a Message.
|
* Subclasses must implement this to map from an Object to a Message.
|
||||||
*/
|
*/
|
||||||
protected abstract Message<?> toMessage(Object object);
|
protected Message<?> toMessage(Object object){
|
||||||
|
if (object instanceof Throwable){
|
||||||
|
if (this.exceptionMapper != null){
|
||||||
|
try {
|
||||||
|
return exceptionMapper.toMessage((Throwable) object);
|
||||||
|
} catch (Exception e2) {
|
||||||
|
logger.warn("Problem mapping " + object + " to message with: " + exceptionMapper);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Subclasses must implement this to map from a Message to an Object.
|
* Subclasses must implement this to map from a Message to an Object.
|
||||||
|
|||||||
@@ -35,6 +35,7 @@ import org.springframework.integration.core.Message;
|
|||||||
import org.springframework.integration.core.MessageChannel;
|
import org.springframework.integration.core.MessageChannel;
|
||||||
import org.springframework.integration.endpoint.AbstractEndpoint;
|
import org.springframework.integration.endpoint.AbstractEndpoint;
|
||||||
import org.springframework.integration.handler.ArgumentArrayMessageMapper;
|
import org.springframework.integration.handler.ArgumentArrayMessageMapper;
|
||||||
|
import org.springframework.integration.message.InboundMessageMapper;
|
||||||
import org.springframework.util.Assert;
|
import org.springframework.util.Assert;
|
||||||
import org.springframework.util.ClassUtils;
|
import org.springframework.util.ClassUtils;
|
||||||
import org.springframework.util.StringUtils;
|
import org.springframework.util.StringUtils;
|
||||||
@@ -48,6 +49,8 @@ import org.springframework.util.StringUtils;
|
|||||||
*/
|
*/
|
||||||
public class GatewayProxyFactoryBean extends AbstractEndpoint implements FactoryBean<Object>, MethodInterceptor, BeanClassLoaderAware {
|
public class GatewayProxyFactoryBean extends AbstractEndpoint implements FactoryBean<Object>, MethodInterceptor, BeanClassLoaderAware {
|
||||||
|
|
||||||
|
private volatile InboundMessageMapper<Throwable> exceptionMapper;
|
||||||
|
|
||||||
private volatile Class<?> serviceInterface;
|
private volatile Class<?> serviceInterface;
|
||||||
|
|
||||||
private volatile MessageChannel defaultRequestChannel;
|
private volatile MessageChannel defaultRequestChannel;
|
||||||
@@ -287,6 +290,7 @@ public class GatewayProxyFactoryBean extends AbstractEndpoint implements Factory
|
|||||||
}
|
}
|
||||||
ArgumentArrayMessageMapper messageMapper = new ArgumentArrayMessageMapper(method, staticHeaders);
|
ArgumentArrayMessageMapper messageMapper = new ArgumentArrayMessageMapper(method, staticHeaders);
|
||||||
SimpleMessagingGateway gateway = new SimpleMessagingGateway(messageMapper, new SimpleMessageMapper());
|
SimpleMessagingGateway gateway = new SimpleMessagingGateway(messageMapper, new SimpleMessageMapper());
|
||||||
|
gateway.setExceptionMapper(exceptionMapper);
|
||||||
if (this.getTaskScheduler() != null) {
|
if (this.getTaskScheduler() != null) {
|
||||||
gateway.setTaskScheduler(this.getTaskScheduler());
|
gateway.setTaskScheduler(this.getTaskScheduler());
|
||||||
}
|
}
|
||||||
@@ -332,5 +336,13 @@ public class GatewayProxyFactoryBean extends AbstractEndpoint implements Factory
|
|||||||
public void setMethodToChannelMap(Map<String, GatewayMethodDefinition> methodToChannelMap) {
|
public void setMethodToChannelMap(Map<String, GatewayMethodDefinition> methodToChannelMap) {
|
||||||
this.methodToChannelMap = methodToChannelMap;
|
this.methodToChannelMap = methodToChannelMap;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public InboundMessageMapper<Throwable> getExceptionMapper() {
|
||||||
|
return exceptionMapper;
|
||||||
|
}
|
||||||
|
|
||||||
|
public void setExceptionMapper(InboundMessageMapper<Throwable> exceptionMapper) {
|
||||||
|
this.exceptionMapper = exceptionMapper;
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -52,6 +52,14 @@ public class SimpleMessagingGateway extends AbstractMessagingGateway {
|
|||||||
this.inboundMapper = inboundMapper;
|
this.inboundMapper = inboundMapper;
|
||||||
this.outboundMapper = outboundMapper;
|
this.outboundMapper = outboundMapper;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public Message<?> sendAndReceiveMessage(Object object) {
|
||||||
|
return (Message<?>) super.sendAndReceive(object, false);
|
||||||
|
}
|
||||||
|
|
||||||
|
public Object sendAndReceive(Object object) {
|
||||||
|
return super.sendAndReceive(object, true);
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@@ -69,19 +77,24 @@ public class SimpleMessagingGateway extends AbstractMessagingGateway {
|
|||||||
|
|
||||||
@Override
|
@Override
|
||||||
protected Message<?> toMessage(Object object) {
|
protected Message<?> toMessage(Object object) {
|
||||||
try {
|
Message<?> message = null;
|
||||||
Message<?> message = this.inboundMapper.toMessage(object);
|
if (object instanceof Throwable){
|
||||||
if (message != null) {
|
message = super.toMessage(object);
|
||||||
message.getHeaders().getHistory().addEvent(this);
|
} else {
|
||||||
|
try {
|
||||||
|
message = this.inboundMapper.toMessage(object);
|
||||||
|
if (message != null) {
|
||||||
|
message.getHeaders().getHistory().addEvent(this);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
return message;
|
catch (Exception e) {
|
||||||
}
|
if (e instanceof RuntimeException) {
|
||||||
catch (Exception e) {
|
throw (RuntimeException) e;
|
||||||
if (e instanceof RuntimeException) {
|
}
|
||||||
throw (RuntimeException) e;
|
throw new MessagingException("failed to create Message", e);
|
||||||
}
|
}
|
||||||
throw new MessagingException("failed to create Message", e);
|
|
||||||
}
|
}
|
||||||
|
return message;
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -533,6 +533,16 @@
|
|||||||
</xsd:appinfo>
|
</xsd:appinfo>
|
||||||
</xsd:annotation>
|
</xsd:annotation>
|
||||||
</xsd:attribute>
|
</xsd:attribute>
|
||||||
|
<xsd:attribute name="exception-mapper" type="xsd:string">
|
||||||
|
<xsd:annotation>
|
||||||
|
<xsd:documentation>
|
||||||
|
<![CDATA[
|
||||||
|
Allows you to provide implementation of InboundMessageMapper which allows you to map
|
||||||
|
an Exception thrown by the endpoint to a successfull return message.
|
||||||
|
]]>
|
||||||
|
</xsd:documentation>
|
||||||
|
</xsd:annotation>
|
||||||
|
</xsd:attribute>
|
||||||
<xsd:attribute name="default-request-channel" type="xsd:string">
|
<xsd:attribute name="default-request-channel" type="xsd:string">
|
||||||
<xsd:annotation>
|
<xsd:annotation>
|
||||||
<xsd:documentation>
|
<xsd:documentation>
|
||||||
|
|||||||
@@ -13,14 +13,15 @@
|
|||||||
<si:gateway id="gatewayWithError"
|
<si:gateway id="gatewayWithError"
|
||||||
default-request-channel="routingChannel"
|
default-request-channel="routingChannel"
|
||||||
service-interface="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleGateway"/>
|
service-interface="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleGateway"/>
|
||||||
|
|
||||||
|
<si:gateway id="gatewayWithErrorAndMapper"
|
||||||
|
default-request-channel="routingChannel"
|
||||||
|
service-interface="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleGateway"
|
||||||
|
exception-mapper="exceptionMapper"/>
|
||||||
|
|
||||||
|
|
||||||
<si:router input-channel="routingChannel" expression="payload"/>
|
<si:router input-channel="routingChannel" expression="payload"/>
|
||||||
|
|
||||||
|
|
||||||
<si:service-activator input-channel="echoWithErrorMessageChannel" method="echoWithErrorMessage">
|
|
||||||
<bean class="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleService" />
|
|
||||||
</si:service-activator>
|
|
||||||
|
|
||||||
<si:service-activator input-channel="echoWithRuntimeExceptionChannel" method="echoWithRuntimeException">
|
<si:service-activator input-channel="echoWithRuntimeExceptionChannel" method="echoWithRuntimeException">
|
||||||
<bean class="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleService" />
|
<bean class="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleService" />
|
||||||
</si:service-activator>
|
</si:service-activator>
|
||||||
@@ -29,18 +30,6 @@
|
|||||||
<bean class="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleService" />
|
<bean class="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleService" />
|
||||||
</si:service-activator>
|
</si:service-activator>
|
||||||
|
|
||||||
<si:service-activator input-channel="echoWithCheckedExceptionChannel" method="echoWithCheckedException">
|
|
||||||
<bean class="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleService" />
|
|
||||||
</si:service-activator>
|
|
||||||
|
|
||||||
<si:service-activator input-channel="echoWithRuntimeExceptionChannel" method="echoWithRuntimeExceptionThrown">
|
|
||||||
<bean class="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleService" />
|
|
||||||
</si:service-activator>
|
|
||||||
|
|
||||||
<si:service-activator input-channel="echoWithMessagingExceptionChannel" method="echoWithMessagingExceptionThrown">
|
|
||||||
<bean class="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleService" />
|
|
||||||
</si:service-activator>
|
|
||||||
|
|
||||||
<si:publish-subscribe-channel id="inputA" />
|
<si:publish-subscribe-channel id="inputA" />
|
||||||
<si:publish-subscribe-channel id="inputB" />
|
<si:publish-subscribe-channel id="inputB" />
|
||||||
<si:publish-subscribe-channel id="inputC" />
|
<si:publish-subscribe-channel id="inputC" />
|
||||||
@@ -75,5 +64,7 @@
|
|||||||
<bean class="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleService" />
|
<bean class="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SimpleService" />
|
||||||
</si:service-activator>
|
</si:service-activator>
|
||||||
</si:chain>
|
</si:chain>
|
||||||
|
|
||||||
|
<bean id="exceptionMapper" class="org.springframework.integration.gateway.GatewayInvokingMessageHandlerTests$SampleExceptionMapper"/>
|
||||||
|
|
||||||
</beans>
|
</beans>
|
||||||
|
|||||||
@@ -24,6 +24,7 @@ import org.springframework.beans.factory.annotation.Autowired;
|
|||||||
import org.springframework.beans.factory.annotation.Qualifier;
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
import org.springframework.integration.channel.SubscribableChannel;
|
import org.springframework.integration.channel.SubscribableChannel;
|
||||||
import org.springframework.integration.core.Message;
|
import org.springframework.integration.core.Message;
|
||||||
|
import org.springframework.integration.message.InboundMessageMapper;
|
||||||
import org.springframework.integration.message.MessageBuilder;
|
import org.springframework.integration.message.MessageBuilder;
|
||||||
import org.springframework.integration.message.MessageHandler;
|
import org.springframework.integration.message.MessageHandler;
|
||||||
import org.springframework.integration.message.MessageHandlingException;
|
import org.springframework.integration.message.MessageHandlingException;
|
||||||
@@ -51,6 +52,11 @@ public class GatewayInvokingMessageHandlerTests {
|
|||||||
@Qualifier("gatewayWithError")
|
@Qualifier("gatewayWithError")
|
||||||
SimpleGateway gatewayWithError;
|
SimpleGateway gatewayWithError;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
@Qualifier("gatewayWithErrorAndMapper")
|
||||||
|
SimpleGateway gatewayWithErrorAndMapper;
|
||||||
|
|
||||||
|
|
||||||
@Autowired
|
@Autowired
|
||||||
@Qualifier("inputB")
|
@Qualifier("inputB")
|
||||||
SubscribableChannel output;
|
SubscribableChannel output;
|
||||||
@@ -82,39 +88,44 @@ public class GatewayInvokingMessageHandlerTests {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void validateGatewayWithErrorMessageReturned() {
|
public void validateGatewayWithErrorMessageReturned() {
|
||||||
|
|
||||||
try {
|
try {
|
||||||
gatewayWithError.sendRecieve("echoWithErrorMessageChannel");
|
String result = gatewayWithErrorAndMapper.sendRecieve("echoWithRuntimeExceptionChannel");
|
||||||
Assert.fail();
|
Assert.assertNotNull(result);
|
||||||
|
Assert.assertEquals("Error happened in message: echoWithRuntimeExceptionChannel", result);
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
Assert.assertEquals("echoWithErrorMessageChannel", e.getMessage());
|
Assert.fail();
|
||||||
}
|
}
|
||||||
|
|
||||||
try {
|
try {
|
||||||
gatewayWithError.sendRecieve("echoWithRuntimeExceptionChannel");
|
gatewayWithError.sendRecieve("echoWithRuntimeExceptionChannel");
|
||||||
Assert.fail();
|
Assert.fail();
|
||||||
} catch (Exception e) {
|
} catch (MessageHandlingException e) {
|
||||||
Assert.assertEquals("echoWithRuntimeExceptionChannel", e.getMessage());
|
Assert.assertEquals("echoWithRuntimeExceptionChannel", e.getFailedMessage().getPayload());
|
||||||
}
|
}
|
||||||
|
|
||||||
try {
|
try {
|
||||||
gatewayWithError.sendRecieve("echoWithMessagingExceptionChannel");
|
gatewayWithError.sendRecieve("echoWithMessagingExceptionChannel");
|
||||||
Assert.fail();
|
Assert.fail();
|
||||||
} catch (MessageHandlingException e) {
|
} catch (MessageHandlingException e) {
|
||||||
Assert.assertEquals("echoWithMessagingExceptionChannel", e.getFailedMessage().getPayload());
|
Assert.assertEquals("echoWithMessagingExceptionChannel", e.getFailedMessage().getPayload());
|
||||||
}
|
}
|
||||||
|
|
||||||
try {
|
try {
|
||||||
gatewayWithError.sendRecieve("echoWithCheckedExceptionChannel");
|
String result = gatewayWithErrorAndMapper.sendRecieve("echoWithMessagingExceptionChannel");
|
||||||
Assert.fail();
|
Assert.assertNotNull(result);
|
||||||
|
Assert.assertEquals("Error happened in message: echoWithMessagingExceptionChannel", result);
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
Assert.assertEquals("echoWithCheckedExceptionChannel", e.getCause().getMessage());
|
Assert.fail();
|
||||||
}
|
}
|
||||||
|
|
||||||
//String result = gatewayWithError.sendRecieve("echoWithErrorMessageChannel");
|
|
||||||
//System.out.println("Result: " + result);
|
|
||||||
//Assert.assertEquals("echo:echo:echo:hello", result);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
public static class SampleExceptionMapper implements InboundMessageMapper<Throwable>{
|
||||||
|
public Message<?> toMessage(Throwable object) throws Exception {
|
||||||
|
MessageHandlingException ex = (MessageHandlingException) object;
|
||||||
|
return MessageBuilder.withPayload("Error happened in message: " + ex.getFailedMessage().getPayload()).build();
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
public static interface SimpleGateway {
|
public static interface SimpleGateway {
|
||||||
public String sendRecieve(String str);
|
public String sendRecieve(String str);
|
||||||
@@ -124,22 +135,11 @@ public class GatewayInvokingMessageHandlerTests {
|
|||||||
public String echo(String value) {
|
public String echo(String value) {
|
||||||
return "echo:" + value;
|
return "echo:" + value;
|
||||||
}
|
}
|
||||||
public Message echoWithErrorMessage(String value) {
|
|
||||||
return MessageBuilder.withPayload(new RuntimeException(value)).build();
|
|
||||||
}
|
|
||||||
public RuntimeException echoWithRuntimeException(String value) {
|
public RuntimeException echoWithRuntimeException(String value) {
|
||||||
return new RuntimeException(value);
|
|
||||||
}
|
|
||||||
public MessageHandlingException echoWithMessagingException(String value) {
|
|
||||||
return new MessageHandlingException(new StringMessage(value));
|
|
||||||
}
|
|
||||||
public SampleCheckedException echoWithCheckedException(String value) {
|
|
||||||
return new SampleCheckedException(value);
|
|
||||||
}
|
|
||||||
public String echoWithRuntimeExceptionThrown(String value) {
|
|
||||||
throw new RuntimeException(value);
|
throw new RuntimeException(value);
|
||||||
}
|
}
|
||||||
public String echoWithMessagingExceptionThrown(String value) {
|
|
||||||
|
public MessageHandlingException echoWithMessagingException(String value) {
|
||||||
throw new MessageHandlingException(new StringMessage(value));
|
throw new MessageHandlingException(new StringMessage(value));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -23,12 +23,12 @@ import javax.jms.JMSException;
|
|||||||
import javax.jms.MessageProducer;
|
import javax.jms.MessageProducer;
|
||||||
import javax.jms.Session;
|
import javax.jms.Session;
|
||||||
|
|
||||||
|
import org.apache.commons.logging.Log;
|
||||||
|
import org.apache.commons.logging.LogFactory;
|
||||||
import org.springframework.beans.factory.InitializingBean;
|
import org.springframework.beans.factory.InitializingBean;
|
||||||
import org.springframework.integration.channel.MessageChannelTemplate;
|
|
||||||
import org.springframework.integration.core.Message;
|
import org.springframework.integration.core.Message;
|
||||||
import org.springframework.integration.core.MessageChannel;
|
import org.springframework.integration.gateway.AbstractMessagingGateway;
|
||||||
import org.springframework.integration.message.MessageBuilder;
|
import org.springframework.integration.message.MessageBuilder;
|
||||||
import org.springframework.integration.message.MessageDeliveryException;
|
|
||||||
import org.springframework.jms.listener.SessionAwareMessageListener;
|
import org.springframework.jms.listener.SessionAwareMessageListener;
|
||||||
import org.springframework.jms.support.converter.MessageConverter;
|
import org.springframework.jms.support.converter.MessageConverter;
|
||||||
import org.springframework.jms.support.destination.DestinationResolver;
|
import org.springframework.jms.support.destination.DestinationResolver;
|
||||||
@@ -43,8 +43,12 @@ import org.springframework.util.Assert;
|
|||||||
*
|
*
|
||||||
* @author Mark Fisher
|
* @author Mark Fisher
|
||||||
* @author Juergen Hoeller
|
* @author Juergen Hoeller
|
||||||
|
* @author Oleg Zhurakousky
|
||||||
*/
|
*/
|
||||||
public class ChannelPublishingJmsMessageListener implements SessionAwareMessageListener<javax.jms.Message>, InitializingBean {
|
public class ChannelPublishingJmsMessageListener extends AbstractMessagingGateway
|
||||||
|
implements SessionAwareMessageListener<javax.jms.Message>, InitializingBean {
|
||||||
|
|
||||||
|
private final Log logger = LogFactory.getLog(this.getClass());
|
||||||
|
|
||||||
private volatile boolean expectReply;
|
private volatile boolean expectReply;
|
||||||
|
|
||||||
@@ -68,41 +72,12 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL
|
|||||||
|
|
||||||
private volatile JmsHeaderMapper headerMapper;
|
private volatile JmsHeaderMapper headerMapper;
|
||||||
|
|
||||||
private final MessageChannelTemplate channelTemplate = new MessageChannelTemplate();
|
|
||||||
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Specify the channel to which request Messages should be sent.
|
|
||||||
*/
|
|
||||||
public void setRequestChannel(MessageChannel requestChannel) {
|
|
||||||
this.channelTemplate.setDefaultChannel(requestChannel);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Specify whether a JMS reply Message is expected.
|
* Specify whether a JMS reply Message is expected.
|
||||||
*/
|
*/
|
||||||
public void setExpectReply(boolean expectReply) {
|
public void setExpectReply(boolean expectReply) {
|
||||||
this.expectReply = expectReply;
|
this.expectReply = expectReply;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* Specify the maximum time to wait when sending a request message to the
|
|
||||||
* request channel. The default value will be that of
|
|
||||||
* {@link MessageChannelTemplate}.
|
|
||||||
*/
|
|
||||||
public void setRequestTimeout(long requestTimeout) {
|
|
||||||
this.channelTemplate.setSendTimeout(requestTimeout);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Specify the maximum time to wait for reply Messages. This value is only
|
|
||||||
* relevant if {@link #expectReply} is <code>true</code>. The default
|
|
||||||
* value will be that of {@link MessageChannelTemplate}.
|
|
||||||
*/
|
|
||||||
public void setReplyTimeout(long replyTimeout) {
|
|
||||||
this.channelTemplate.setReceiveTimeout(replyTimeout);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Set the default reply destination to send reply messages to. This will
|
* Set the default reply destination to send reply messages to. This will
|
||||||
* be applied in case of a request message that does not carry a
|
* be applied in case of a request message that does not carry a
|
||||||
@@ -226,8 +201,7 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL
|
|||||||
public void setExtractReplyPayload(boolean extractReplyPayload) {
|
public void setExtractReplyPayload(boolean extractReplyPayload) {
|
||||||
this.extractReplyPayload = extractReplyPayload;
|
this.extractReplyPayload = extractReplyPayload;
|
||||||
}
|
}
|
||||||
|
public final void onInit() {
|
||||||
public final void afterPropertiesSet() {
|
|
||||||
if (!(this.messageConverter instanceof HeaderMappingMessageConverter)) {
|
if (!(this.messageConverter instanceof HeaderMappingMessageConverter)) {
|
||||||
HeaderMappingMessageConverter hmmc = new HeaderMappingMessageConverter(this.messageConverter, this.headerMapper);
|
HeaderMappingMessageConverter hmmc = new HeaderMappingMessageConverter(this.messageConverter, this.headerMapper);
|
||||||
hmmc.setExtractJmsMessageBody(this.extractRequestPayload);
|
hmmc.setExtractJmsMessageBody(this.extractRequestPayload);
|
||||||
@@ -241,31 +215,31 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL
|
|||||||
Message<?> requestMessage = (object instanceof Message<?>) ?
|
Message<?> requestMessage = (object instanceof Message<?>) ?
|
||||||
(Message<?>) object : MessageBuilder.withPayload(object).build();
|
(Message<?>) object : MessageBuilder.withPayload(object).build();
|
||||||
if (!this.expectReply) {
|
if (!this.expectReply) {
|
||||||
boolean sent = this.channelTemplate.send(requestMessage);
|
this.send(requestMessage);
|
||||||
if (!sent) {
|
|
||||||
throw new MessageDeliveryException(requestMessage, "failed to send Message to request channel");
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
else {
|
else {
|
||||||
Message<?> replyMessage = this.channelTemplate.sendAndReceive(requestMessage);
|
Message<?> replyMessage = this.sendAndReceiveMessage(requestMessage);
|
||||||
if (replyMessage != null) {
|
|
||||||
Destination destination = this.getReplyDestination(jmsMessage, session);
|
if (replyMessage != null){
|
||||||
javax.jms.Message jmsReply = this.messageConverter.toMessage(replyMessage, session);
|
Destination destination = this.getReplyDestination(jmsMessage, session, false);
|
||||||
if (jmsReply.getJMSCorrelationID() == null) {
|
if (destination != null){
|
||||||
jmsReply.setJMSCorrelationID(jmsMessage.getJMSMessageID());
|
javax.jms.Message jmsReply = this.messageConverter.toMessage(replyMessage, session);
|
||||||
}
|
if (jmsReply.getJMSCorrelationID() == null) {
|
||||||
MessageProducer producer = session.createProducer(destination);
|
jmsReply.setJMSCorrelationID(jmsMessage.getJMSMessageID());
|
||||||
try {
|
|
||||||
if (this.explicitQosEnabledForReplies) {
|
|
||||||
producer.send(jmsReply,
|
|
||||||
this.replyDeliveryMode, this.replyPriority, this.replyTimeToLive);
|
|
||||||
}
|
}
|
||||||
else {
|
MessageProducer producer = session.createProducer(destination);
|
||||||
producer.send(jmsReply);
|
try {
|
||||||
|
if (this.explicitQosEnabledForReplies) {
|
||||||
|
producer.send(jmsReply,
|
||||||
|
this.replyDeliveryMode, this.replyPriority, this.replyTimeToLive);
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
producer.send(jmsReply);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
finally {
|
||||||
|
producer.close();
|
||||||
}
|
}
|
||||||
}
|
|
||||||
finally {
|
|
||||||
producer.close();
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -273,7 +247,9 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL
|
|||||||
|
|
||||||
/**
|
/**
|
||||||
* Determine a reply destination for the given message.
|
* Determine a reply destination for the given message.
|
||||||
* <p>This implementation first checks the JMS Reply-To {@link Destination}
|
* <p>This implementation first checks the boolean 'error' flag which signifies that the reply is an error message.
|
||||||
|
* If it is it will attempt to get the destination value from {@link JmsHeaders#SEND_ERROR_TO}.
|
||||||
|
* If reply is not an error it will first check the JMS Reply-To {@link Destination}
|
||||||
* of the supplied request message; if that is not <code>null</code> it is
|
* of the supplied request message; if that is not <code>null</code> it is
|
||||||
* returned; if it is <code>null</code>, then the configured
|
* returned; if it is <code>null</code>, then the configured
|
||||||
* {@link #resolveDefaultReplyDestination default reply destination}
|
* {@link #resolveDefaultReplyDestination default reply destination}
|
||||||
@@ -287,7 +263,8 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL
|
|||||||
* @see #setDefaultReplyDestination
|
* @see #setDefaultReplyDestination
|
||||||
* @see javax.jms.Message#getJMSReplyTo()
|
* @see javax.jms.Message#getJMSReplyTo()
|
||||||
*/
|
*/
|
||||||
private Destination getReplyDestination(javax.jms.Message request, Session session) throws JMSException {
|
private Destination getReplyDestination(javax.jms.Message request, Session session, boolean error) throws JMSException {
|
||||||
|
|
||||||
Destination replyTo = request.getJMSReplyTo();
|
Destination replyTo = request.getJMSReplyTo();
|
||||||
if (replyTo == null) {
|
if (replyTo == null) {
|
||||||
replyTo = resolveDefaultReplyDestination(session);
|
replyTo = resolveDefaultReplyDestination(session);
|
||||||
@@ -296,6 +273,7 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL
|
|||||||
"Request message does not contain reply-to destination, and no default reply destination set.");
|
"Request message does not contain reply-to destination, and no default reply destination set.");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return replyTo;
|
return replyTo;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -337,4 +315,19 @@ public class ChannelPublishingJmsMessageListener implements SessionAwareMessageL
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
protected Object fromMessage(Message<?> message) {
|
||||||
|
throw new UnsupportedOperationException("'fromMessage' is not supported within this instance");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
protected Message<?> toMessage(Object object) {
|
||||||
|
Message<?> message = null;
|
||||||
|
if (object instanceof Throwable){
|
||||||
|
message = super.toMessage(object);
|
||||||
|
} else if (object instanceof Message<?>) {
|
||||||
|
message = (Message<?>) object;
|
||||||
|
}
|
||||||
|
return message;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -156,6 +156,7 @@ public class JmsMessageDrivenEndpointParser extends AbstractSingleBeanDefinition
|
|||||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "reply-timeout");
|
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "reply-timeout");
|
||||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "extract-request-payload");
|
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "extract-request-payload");
|
||||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "extract-reply-payload");
|
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "extract-reply-payload");
|
||||||
|
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "exception-mapper");
|
||||||
int defaults = 0;
|
int defaults = 0;
|
||||||
if (StringUtils.hasText(element.getAttribute(DEFAULT_REPLY_DESTINATION_ATTRIB))) {
|
if (StringUtils.hasText(element.getAttribute(DEFAULT_REPLY_DESTINATION_ATTRIB))) {
|
||||||
defaults++;
|
defaults++;
|
||||||
|
|||||||
@@ -468,6 +468,7 @@
|
|||||||
</xsd:attribute>
|
</xsd:attribute>
|
||||||
<xsd:attribute name="request-destination-name" type="xsd:string"/>
|
<xsd:attribute name="request-destination-name" type="xsd:string"/>
|
||||||
<xsd:attribute name="request-pub-sub-domain" type="xsd:string"/>
|
<xsd:attribute name="request-pub-sub-domain" type="xsd:string"/>
|
||||||
|
<xsd:attribute name="exception-mapper" type="xsd:string"/>
|
||||||
<xsd:attribute name="default-reply-destination" type="xsd:string">
|
<xsd:attribute name="default-reply-destination" type="xsd:string">
|
||||||
<xsd:annotation>
|
<xsd:annotation>
|
||||||
<xsd:documentation>
|
<xsd:documentation>
|
||||||
|
|||||||
@@ -0,0 +1,66 @@
|
|||||||
|
<?xml version="1.0" encoding="UTF-8"?>
|
||||||
|
<beans xmlns="http://www.springframework.org/schema/beans"
|
||||||
|
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||||
|
xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans-3.0.xsd
|
||||||
|
http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration-2.0.xsd
|
||||||
|
http://www.springframework.org/schema/integration/jms http://www.springframework.org/schema/integration/jms/spring-integration-jms-2.0.xsd
|
||||||
|
http://www.springframework.org/schema/jms http://www.springframework.org/schema/jms/spring-jms-3.0.xsd
|
||||||
|
http://www.springframework.org/schema/task http://www.springframework.org/schema/task/spring-task-3.0.xsd"
|
||||||
|
xmlns:int="http://www.springframework.org/schema/integration"
|
||||||
|
xmlns:int-jms="http://www.springframework.org/schema/integration/jms"
|
||||||
|
xmlns:jms="http://www.springframework.org/schema/jms"
|
||||||
|
xmlns:task="http://www.springframework.org/schema/task">
|
||||||
|
|
||||||
|
<int:gateway id="sampleGateway"
|
||||||
|
service-interface="org.springframework.integration.jms.config.ExceptionHandlingSiConsumerTests$SampleGateway"
|
||||||
|
default-request-channel="outbound-channel">
|
||||||
|
</int:gateway>
|
||||||
|
|
||||||
|
<int:channel id="outbound-channel"/>
|
||||||
|
|
||||||
|
<int-jms:outbound-gateway request-channel="outbound-channel" request-destination="requestQueue"/>
|
||||||
|
|
||||||
|
|
||||||
|
<int-jms:inbound-gateway request-destination="requestQueue"
|
||||||
|
request-channel="jmsinputchannel"
|
||||||
|
exception-mapper="errorMessageMapper"/>
|
||||||
|
|
||||||
|
<bean id="errorMessageMapper" class="org.springframework.integration.jms.config.ExceptionHandlingSiConsumerTests$SampleErrorMessageMapper"/>
|
||||||
|
<int:channel id="jmsinputchannel"/>
|
||||||
|
|
||||||
|
<int:router input-channel="jmsinputchannel" expression="payload"/>
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
<int:service-activator input-channel="echoWithExceptionChannel" method="echoWithException">
|
||||||
|
<bean class="org.springframework.integration.jms.config.ExceptionHandlingSiConsumerTests$SampleService"/>
|
||||||
|
</int:service-activator>
|
||||||
|
|
||||||
|
|
||||||
|
<int:service-activator input-channel="echoChannel" method="echo">
|
||||||
|
<bean class="org.springframework.integration.jms.config.ExceptionHandlingSiConsumerTests$SampleService"/>
|
||||||
|
</int:service-activator>
|
||||||
|
|
||||||
|
<bean id="requestQueue" class="org.apache.activemq.command.ActiveMQQueue">
|
||||||
|
<constructor-arg value="request.queue"/>
|
||||||
|
</bean>
|
||||||
|
|
||||||
|
<bean id="replyQueue" class="org.apache.activemq.command.ActiveMQQueue">
|
||||||
|
<constructor-arg value="reply.queue"/>
|
||||||
|
</bean>
|
||||||
|
|
||||||
|
<bean id="replyTopic" class="org.apache.activemq.command.ActiveMQTopic">
|
||||||
|
<constructor-arg value="reply.topic"/>
|
||||||
|
</bean>
|
||||||
|
|
||||||
|
<bean id="connectionFactory" class="org.springframework.jms.connection.CachingConnectionFactory">
|
||||||
|
<property name="targetConnectionFactory">
|
||||||
|
<bean class="org.apache.activemq.ActiveMQConnectionFactory">
|
||||||
|
<property name="brokerURL" value="vm://localhost"/>
|
||||||
|
</bean>
|
||||||
|
</property>
|
||||||
|
<property name="sessionCacheSize" value="10"/>
|
||||||
|
<property name="cacheProducers" value="false"/>
|
||||||
|
</bean>
|
||||||
|
|
||||||
|
</beans>
|
||||||
@@ -0,0 +1,142 @@
|
|||||||
|
/*
|
||||||
|
* 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 javax.jms.ConnectionFactory;
|
||||||
|
import javax.jms.Destination;
|
||||||
|
import javax.jms.JMSException;
|
||||||
|
import javax.jms.Message;
|
||||||
|
import javax.jms.Session;
|
||||||
|
import javax.jms.TextMessage;
|
||||||
|
|
||||||
|
import org.junit.Assert;
|
||||||
|
import org.junit.Before;
|
||||||
|
import org.junit.Test;
|
||||||
|
import org.springframework.context.ConfigurableApplicationContext;
|
||||||
|
import org.springframework.context.support.ClassPathXmlApplicationContext;
|
||||||
|
import org.springframework.integration.message.InboundMessageMapper;
|
||||||
|
import org.springframework.integration.message.MessageBuilder;
|
||||||
|
import org.springframework.jms.core.JmsTemplate;
|
||||||
|
import org.springframework.jms.core.MessageCreator;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* @author Oleg Zhurakousky
|
||||||
|
*
|
||||||
|
*/
|
||||||
|
public class ExceptionHandlingSiConsumerTests {
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void nonSiProducer_siConsumer_sync_withReturn() throws Exception {
|
||||||
|
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);
|
||||||
|
final Destination reply = applicationContext.getBean("replyQueue", Destination.class);
|
||||||
|
jmsTemplate.send(request, new MessageCreator() {
|
||||||
|
|
||||||
|
public Message createMessage(Session session) throws JMSException {
|
||||||
|
TextMessage message = session.createTextMessage();
|
||||||
|
message.setText("echoChannel");
|
||||||
|
message.setJMSReplyTo(reply);
|
||||||
|
return message;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
Message message = jmsTemplate.receive(reply);
|
||||||
|
Assert.assertNotNull(message);
|
||||||
|
applicationContext.close();
|
||||||
|
}
|
||||||
|
@Test
|
||||||
|
public void nonSiProducer_siConsumer_sync_withReturnNoException() throws Exception {
|
||||||
|
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);
|
||||||
|
final Destination reply = applicationContext.getBean("replyQueue", Destination.class);
|
||||||
|
jmsTemplate.send(request, new MessageCreator() {
|
||||||
|
|
||||||
|
public Message createMessage(Session session) throws JMSException {
|
||||||
|
TextMessage message = session.createTextMessage();
|
||||||
|
message.setText("echoWithExceptionChannel");
|
||||||
|
message.setJMSReplyTo(reply);
|
||||||
|
return message;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
Message message = jmsTemplate.receive(reply);
|
||||||
|
Assert.assertNotNull(message);
|
||||||
|
applicationContext.close();
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void nonSiProducer_siConsumer_sync_withOutboundGateway() throws Exception{
|
||||||
|
final ConfigurableApplicationContext applicationContext = new ClassPathXmlApplicationContext("Exception-nonSiProducer-siConsumer.xml", ExceptionHandlingSiConsumerTests.class);
|
||||||
|
SampleGateway gateway = applicationContext.getBean("sampleGateway", SampleGateway.class);
|
||||||
|
String reply = gateway.echo("echoWithExceptionChannel");
|
||||||
|
System.out.println("Reply: " + reply);
|
||||||
|
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");
|
||||||
|
}
|
||||||
|
public String echo(String value){
|
||||||
|
return value;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
@SuppressWarnings("serial")
|
||||||
|
public static class SampleException extends RuntimeException{
|
||||||
|
public SampleException(String message){
|
||||||
|
super(message);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public static interface SampleGateway{
|
||||||
|
public String echo(String value);
|
||||||
|
}
|
||||||
|
|
||||||
|
public static class SampleErrorMessageMapper implements InboundMessageMapper<Throwable>{
|
||||||
|
public org.springframework.integration.core.Message<?> toMessage(
|
||||||
|
Throwable t) throws Exception {
|
||||||
|
return MessageBuilder.withPayload(t.getCause().getMessage()).build();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user