From 3508361a29436079b91fe7250dccdafddfb11753 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Mon, 25 Feb 2008 14:30:11 +0000 Subject: [PATCH] DefaultMessageEndpoint now accepts the MessageHandler as a constructor arg (the setter has been removed). The MessageBus now has a 'defaultConcurrencyPolicy' rather than passing 'null' when creating an endpoint without an explicitly provided policy (INT-106). The ConcurrentHandler.HandlerTask now logs all errors at DEBUG level and continues to log at WARN level if no 'errorHandler' has been provided to the ConcurrentHandler (INT-120). --- .../integration/bus/MessageBus.java | 4 +- ...essageEndpointAnnotationPostProcessor.java | 5 +- .../endpoint/ConcurrentHandler.java | 5 +- .../integration/adapter/adapterTests.xml | 2 +- .../DefaultMessageDispatcherTests.java | 6 +- .../endpoint/DefaultMessageEndpointTests.java | 60 +++++++------------ 6 files changed, 34 insertions(+), 48 deletions(-) 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 e6839ecf4c..f6a446d126 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 @@ -78,6 +78,8 @@ public class MessageBus implements ChannelRegistry, ApplicationContextAware, Lif private volatile ScheduledExecutorService executor; + private volatile ConcurrencyPolicy defaultConcurrencyPolicy = new ConcurrencyPolicy(1, 10); + private volatile boolean autoCreateChannels; private volatile boolean autoStartup = true; @@ -211,7 +213,7 @@ public class MessageBus implements ChannelRegistry, ApplicationContextAware, Lif } public void registerHandler(String name, MessageHandler handler, Subscription subscription) { - this.registerHandler(name, handler, subscription, null); + this.registerHandler(name, handler, subscription, this.defaultConcurrencyPolicy); } public void registerHandler(String name, MessageHandler handler, Subscription subscription, ConcurrencyPolicy concurrencyPolicy) { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java index e2d508018e..76feb64334 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/MessageEndpointAnnotationPostProcessor.java @@ -111,10 +111,9 @@ public class MessageEndpointAnnotationPostProcessor implements BeanPostProcessor } return bean; } - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); - this.configureInput(bean, beanName, endpointAnnotation, endpoint); MessageHandlerChain handlerChain = this.createHandlerChain(bean); - endpoint.setHandler(handlerChain); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(handlerChain); + this.configureInput(bean, beanName, endpointAnnotation, endpoint); this.configureDefaultOutput(bean, beanName, endpointAnnotation, endpoint); this.messageBus.registerEndpoint(beanName + "-endpoint", endpoint); return bean; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ConcurrentHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ConcurrentHandler.java index 9d53b56224..12a5df4334 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ConcurrentHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ConcurrentHandler.java @@ -101,10 +101,13 @@ public class ConcurrentHandler implements MessageHandler, DisposableBean { } } catch (Throwable t) { + if (logger.isDebugEnabled()) { + logger.debug("error occurred in handler execution", t); + } if (errorHandler != null) { errorHandler.handle(t); } - else if (logger.isWarnEnabled()) { + else if (logger.isWarnEnabled() && !logger.isDebugEnabled()) { logger.warn("error occurred in handler execution", t); } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/adapter/adapterTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/adapter/adapterTests.xml index 4a792cb0cc..96eec93029 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/adapter/adapterTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/adapter/adapterTests.xml @@ -34,7 +34,7 @@ - + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/DefaultMessageDispatcherTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/DefaultMessageDispatcherTests.java index 2748fb31c3..5ce08027e8 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/DefaultMessageDispatcherTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/DefaultMessageDispatcherTests.java @@ -454,11 +454,9 @@ public class DefaultMessageDispatcherTests { SimpleChannel channel = new SimpleChannel(new DispatcherPolicy(true)); channel.send(new StringMessage(1, "test")); DefaultMessageDispatcher dispatcher = new DefaultMessageDispatcher(channel, scheduler); - DefaultMessageEndpoint endpoint1 = new DefaultMessageEndpoint(); - endpoint1.setHandler(handler1); + DefaultMessageEndpoint endpoint1 = new DefaultMessageEndpoint(handler1); endpoint1.setConcurrencyPolicy(new ConcurrencyPolicy(1, 1)); - DefaultMessageEndpoint endpoint2 = new DefaultMessageEndpoint(); - endpoint2.setHandler(handler2); + DefaultMessageEndpoint endpoint2 = new DefaultMessageEndpoint(handler2); endpoint2.setConcurrencyPolicy(new ConcurrencyPolicy(1, 1)); endpoint1.addMessageSelector(new PayloadTypeSelector(Integer.class)); endpoint2.addMessageSelector(new PayloadTypeSelector(String.class)); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/DefaultMessageEndpointTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/DefaultMessageEndpointTests.java index d2dbf71246..32b4c7806b 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/DefaultMessageEndpointTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/DefaultMessageEndpointTests.java @@ -58,9 +58,8 @@ public class DefaultMessageEndpointTests { return new StringMessage("123", "hello " + message.getPayload()); } }; - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(handler); endpoint.setChannelRegistry(channelRegistry); - endpoint.setHandler(handler); endpoint.setDefaultOutputChannelName("replyChannel"); endpoint.start(); endpoint.handle(new StringMessage(1, "test")); @@ -78,8 +77,7 @@ public class DefaultMessageEndpointTests { return new StringMessage("123", "hello " + message.getPayload()); } }; - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); - endpoint.setHandler(handler); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(handler); endpoint.start(); StringMessage testMessage = new StringMessage(1, "test"); testMessage.getHeader().setReturnAddress(replyChannel); @@ -100,9 +98,8 @@ public class DefaultMessageEndpointTests { return new StringMessage("123", "hello " + message.getPayload()); } }; - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(handler); endpoint.setChannelRegistry(channelRegistry); - endpoint.setHandler(handler); endpoint.start(); StringMessage testMessage = new StringMessage(1, "test"); testMessage.getHeader().setReturnAddress("replyChannel"); @@ -124,9 +121,8 @@ public class DefaultMessageEndpointTests { return new StringMessage("123", "hello " + message.getPayload()); } }; - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(handler); endpoint.setChannelRegistry(channelRegistry); - endpoint.setHandler(handler); endpoint.start(); StringMessage testMessage = new StringMessage("test"); testMessage.getHeader().setReturnAddress(replyChannel1); @@ -148,15 +144,14 @@ public class DefaultMessageEndpointTests { @Test public void testCustomErrorHandler() throws InterruptedException { - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); - endpoint.setConcurrencyPolicy(new ConcurrencyPolicy(1, 1)); final CountDownLatch latch = new CountDownLatch(2); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(TestHandlers.rejectingCountDownHandler(latch)); + endpoint.setConcurrencyPolicy(new ConcurrencyPolicy(1, 1)); endpoint.setErrorHandler(new ErrorHandler() { public void handle(Throwable t) { latch.countDown(); } }); - endpoint.setHandler(TestHandlers.rejectingCountDownHandler(latch)); endpoint.start(); endpoint.handle(new StringMessage("test")); latch.await(500, TimeUnit.MILLISECONDS); @@ -175,9 +170,8 @@ public class DefaultMessageEndpointTests { return new StringMessage("123", "hello " + message.getPayload()); } }; - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(new ConcurrentHandler(handler, createExecutor())); endpoint.setChannelRegistry(channelRegistry); - endpoint.setHandler(new ConcurrentHandler(handler, createExecutor())); endpoint.setDefaultOutputChannelName("replyChannel"); endpoint.start(); endpoint.handle(new StringMessage(1, "test")); @@ -201,9 +195,8 @@ public class DefaultMessageEndpointTests { return null; } }; - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(handler); endpoint.setChannelRegistry(channelRegistry); - endpoint.setHandler(handler); endpoint.setDefaultOutputChannelName("replyChannel"); endpoint.start(); endpoint.handle(new StringMessage(1, "test")); @@ -226,9 +219,8 @@ public class DefaultMessageEndpointTests { return null; } }; - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(new ConcurrentHandler(handler, createExecutor())); endpoint.setChannelRegistry(channelRegistry); - endpoint.setHandler(new ConcurrentHandler(handler, createExecutor())); endpoint.setDefaultOutputChannelName("replyChannel"); endpoint.start(); endpoint.handle(new StringMessage(1, "test")); @@ -251,9 +243,8 @@ public class DefaultMessageEndpointTests { return new StringMessage("123", "hello " + message.getPayload()); } }; - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(new ConcurrentHandler(handler, createExecutor())); endpoint.setChannelRegistry(channelRegistry); - endpoint.setHandler(new ConcurrentHandler(handler, createExecutor())); endpoint.start(); StringMessage message = new StringMessage(1, "test"); message.getHeader().setReturnAddress("replyChannel"); @@ -278,9 +269,8 @@ public class DefaultMessageEndpointTests { return new StringMessage("123", "hello " + message.getPayload()); } }; - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(handler); endpoint.setChannelRegistry(channelRegistry); - endpoint.setHandler(handler); endpoint.setConcurrencyPolicy(new ConcurrencyPolicy(3, 14)); endpoint.setDefaultOutputChannelName("replyChannel"); endpoint.start(); @@ -305,9 +295,8 @@ public class DefaultMessageEndpointTests { return new StringMessage("123", "hello " + message.getPayload()); } }; - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(handler); endpoint.setChannelRegistry(channelRegistry); - endpoint.setHandler(handler); endpoint.setConcurrencyPolicy(new ConcurrencyPolicy(3, 14)); endpoint.start(); StringMessage message = new StringMessage(1, "test"); @@ -323,18 +312,17 @@ public class DefaultMessageEndpointTests { @Test(expected=MessageHandlerNotRunningException.class) public void testEndpointDoesNotHandleMessagesWhenNotYetStarted() { - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(TestHandlers.nullHandler()); endpoint.handle(new StringMessage("test")); } @Test public void testEndpointDoesNotHandleMessagesAfterBeingStopped() { - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); AtomicInteger counter = new AtomicInteger(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(TestHandlers.countingHandler(counter)); boolean exceptionThrown = false; try { endpoint.start(); - endpoint.setHandler(TestHandlers.countingHandler(counter)); endpoint.handle(new StringMessage("test1")); endpoint.stop(); endpoint.handle(new StringMessage("test2")); @@ -348,7 +336,7 @@ public class DefaultMessageEndpointTests { @Test(expected=MessageSelectorRejectedException.class) public void testEndpointWithSelectorRejecting() { - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(TestHandlers.nullHandler()); endpoint.addMessageSelector(new MessageSelector() { public boolean accept(Message message) { return false; @@ -360,14 +348,13 @@ public class DefaultMessageEndpointTests { @Test public void testEndpointWithSelectorAccepting() throws InterruptedException { - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); + CountDownLatch latch = new CountDownLatch(1); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(TestHandlers.countDownHandler(latch)); endpoint.addMessageSelector(new MessageSelector() { public boolean accept(Message message) { return true; } }); - CountDownLatch latch = new CountDownLatch(1); - endpoint.setHandler(TestHandlers.countDownHandler(latch)); endpoint.start(); endpoint.handle(new StringMessage("test")); latch.await(100, TimeUnit.MILLISECONDS); @@ -377,9 +364,9 @@ public class DefaultMessageEndpointTests { @Test public void testEndpointWithMultipleSelectorsAndFirstRejects() { - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); - boolean exceptionThrown = false; final AtomicInteger counter = new AtomicInteger(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(TestHandlers.countingHandler(counter)); + boolean exceptionThrown = false; endpoint.addMessageSelector(new MessageSelector() { public boolean accept(Message message) { counter.incrementAndGet(); @@ -392,7 +379,6 @@ public class DefaultMessageEndpointTests { return true; } }); - endpoint.setHandler(TestHandlers.countingHandler(counter)); endpoint.start(); try { endpoint.handle(new StringMessage("test")); @@ -407,9 +393,9 @@ public class DefaultMessageEndpointTests { @Test public void testEndpointWithMultipleSelectorsAndFirstAccepts() { - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); - boolean exceptionThrown = false; final AtomicInteger counter = new AtomicInteger(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(TestHandlers.countingHandler(counter)); + boolean exceptionThrown = false; endpoint.addMessageSelector(new MessageSelector() { public boolean accept(Message message) { counter.incrementAndGet(); @@ -422,7 +408,6 @@ public class DefaultMessageEndpointTests { return false; } }); - endpoint.setHandler(TestHandlers.countingHandler(counter)); endpoint.start(); try { endpoint.handle(new StringMessage("test")); @@ -437,8 +422,8 @@ public class DefaultMessageEndpointTests { @Test public void testEndpointWithMultipleSelectorsAndBothAccept() { - DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(); final AtomicInteger counter = new AtomicInteger(); + DefaultMessageEndpoint endpoint = new DefaultMessageEndpoint(TestHandlers.countingHandler(counter)); endpoint.addMessageSelector(new MessageSelector() { public boolean accept(Message message) { counter.incrementAndGet(); @@ -451,7 +436,6 @@ public class DefaultMessageEndpointTests { return true; } }); - endpoint.setHandler(TestHandlers.countingHandler(counter)); endpoint.start(); endpoint.handle(new StringMessage("test")); assertEquals("both selectors and handler should have been invoked", 3, counter.get());