From d634f38d10015df68cf9b3528436723e746b1e9f Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Fri, 11 Apr 2008 02:37:51 +0000 Subject: [PATCH] A rethrowing Exception handler is provided for any Endpoint subscribed to a SynchronousChannel. --- .../integration/bus/MessageBus.java | 12 ++++++ .../SynchronousChannelSubscriptionTests.java | 40 +++++++++++++++++++ 2 files changed, 52 insertions(+) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java index 64647d0fca..f1c9f69265 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java @@ -48,12 +48,14 @@ import org.springframework.integration.endpoint.DefaultMessageEndpoint; import org.springframework.integration.endpoint.EndpointRegistry; import org.springframework.integration.endpoint.MessageEndpoint; import org.springframework.integration.handler.MessageHandler; +import org.springframework.integration.message.MessagingException; import org.springframework.integration.scheduling.MessagePublishingErrorHandler; import org.springframework.integration.scheduling.MessagingTaskScheduler; import org.springframework.integration.scheduling.MessagingTaskSchedulerAware; import org.springframework.integration.scheduling.Schedule; import org.springframework.integration.scheduling.SimpleMessagingTaskScheduler; import org.springframework.integration.scheduling.Subscription; +import org.springframework.integration.util.ErrorHandler; import org.springframework.util.Assert; /** @@ -373,6 +375,16 @@ public class MessageBus implements ChannelRegistry, EndpointRegistry, Applicatio if (handler instanceof Lifecycle) { ((Lifecycle) handler).start(); } + if (handler instanceof DefaultMessageEndpoint) { + ((DefaultMessageEndpoint) handler).setErrorHandler(new ErrorHandler() { + public void handle(Throwable t) { + if (t instanceof MessagingException) { + throw (MessagingException) t; + } + throw new MessagingException("error occurred in handler", t); + } + }); + } return; } SchedulingMessageDispatcher dispatcher = dispatchers.get(channel); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/bus/SynchronousChannelSubscriptionTests.java b/spring-integration-core/src/test/java/org/springframework/integration/bus/SynchronousChannelSubscriptionTests.java index 9511c6bd4c..ce060603d3 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/bus/SynchronousChannelSubscriptionTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/bus/SynchronousChannelSubscriptionTests.java @@ -24,11 +24,13 @@ import org.junit.Test; import org.springframework.integration.annotation.Handler; import org.springframework.integration.annotation.MessageEndpoint; import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.channel.SimpleChannel; import org.springframework.integration.config.MessageEndpointAnnotationPostProcessor; import org.springframework.integration.dispatcher.SynchronousChannel; import org.springframework.integration.endpoint.DefaultMessageEndpoint; import org.springframework.integration.handler.MessageHandler; import org.springframework.integration.message.Message; +import org.springframework.integration.message.MessagingException; import org.springframework.integration.message.StringMessage; import org.springframework.integration.scheduling.Subscription; @@ -77,6 +79,34 @@ public class SynchronousChannelSubscriptionTests { bus.stop(); } + @Test(expected=MessagingException.class) + public void testExceptionThrownFromRegisteredEndpoint() { + SimpleChannel errorChannel = new SimpleChannel(); + bus.setErrorChannel(errorChannel); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(new MessageHandler() { + public Message handle(Message message) { + throw new RuntimeException("intentional test failure"); + } + }); + endpoint.setSubscription(new Subscription("sourceChannel")); + endpoint.setDefaultOutputChannelName("targetChannel"); + bus.registerEndpoint("testEndpoint", endpoint); + bus.start(); + this.sourceChannel.send(new StringMessage("foo")); + } + + @Test(expected=MessagingException.class) + public void testExceptionThrownFromAnnotatedEndpoint() { + SimpleChannel errorChannel = new SimpleChannel(); + bus.setErrorChannel(errorChannel); + MessageEndpointAnnotationPostProcessor postProcessor = new MessageEndpointAnnotationPostProcessor(bus); + postProcessor.afterPropertiesSet(); + FailingTestEndpoint endpoint = new FailingTestEndpoint(); + postProcessor.postProcessAfterInitialization(endpoint, "testEndpoint"); + bus.start(); + this.sourceChannel.send(new StringMessage("foo")); + } + private static class TestHandler implements MessageHandler { @@ -95,4 +125,14 @@ public class SynchronousChannelSubscriptionTests { } } + + @MessageEndpoint(input="sourceChannel", defaultOutput="targetChannel") + public static class FailingTestEndpoint { + + @Handler + public Message handle(Message message) { + throw new RuntimeException("intentional test failure"); + } + } + }