From 7713ee4764285f3886bcc2172d6399b0e0ef6ef6 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 11 Apr 2011 16:09:40 -0400 Subject: [PATCH] INT-1841 fixed the way inbound gateways react to async channel as request-channel when such channel has no subscribers --- .../dispatcher/UnicastingDispatcher.java | 2 +- .../dispatcher/FailOverDispatcherTests.java | 3 +- .../RoundRobinDispatcherConcurrentTests.java | 30 ++++++---- .../dispatcher/UnicastingDispatcherTests.java | 60 +++++++++++++++++++ .../dispatcher/unicasting-with-async.xml | 17 ++++++ 5 files changed, 97 insertions(+), 15 deletions(-) create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/dispatcher/UnicastingDispatcherTests.java create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/dispatcher/unicasting-with-async.xml diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/UnicastingDispatcher.java b/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/UnicastingDispatcher.java index ba9f093799..dd1e5a3c17 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/UnicastingDispatcher.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dispatcher/UnicastingDispatcher.java @@ -101,7 +101,7 @@ public class UnicastingDispatcher extends AbstractDispatcher { boolean success = false; Iterator handlerIterator = this.getHandlerIterator(message); if (!handlerIterator.hasNext()) { - throw new IllegalStateException("Dispatcher has no subscribers."); + throw new MessageDeliveryException(message, "Dispatcher has no subscribers."); } List exceptions = new ArrayList(); while (success == false && handlerIterator.hasNext()) { diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/FailOverDispatcherTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/FailOverDispatcherTests.java index 4dbbf32e0d..54aaeb72e6 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/FailOverDispatcherTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/FailOverDispatcherTests.java @@ -27,6 +27,7 @@ import java.util.concurrent.atomic.AtomicInteger; import org.junit.Test; import org.springframework.integration.Message; +import org.springframework.integration.MessageDeliveryException; import org.springframework.integration.MessageRejectedException; import org.springframework.integration.core.MessageHandler; import org.springframework.integration.handler.ServiceActivatingHandler; @@ -133,7 +134,7 @@ public class FailOverDispatcherTests { assertEquals(6, counter.get()); } - @Test(expected = IllegalStateException.class) + @Test(expected = MessageDeliveryException.class) public void removeConsumerLastTargetCausesDeliveryException() { UnicastingDispatcher dispatcher = new UnicastingDispatcher(); final AtomicInteger counter = new AtomicInteger(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/RoundRobinDispatcherConcurrentTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/RoundRobinDispatcherConcurrentTests.java index b25a2c55db..a09f58a4e1 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/RoundRobinDispatcherConcurrentTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/RoundRobinDispatcherConcurrentTests.java @@ -15,23 +15,27 @@ package org.springframework.integration.dispatcher; -import org.junit.Before; -import org.junit.Test; -import org.junit.runner.RunWith; -import org.mockito.Mock; -import org.mockito.runners.MockitoJUnitRunner; -import org.springframework.integration.Message; -import org.springframework.integration.MessageRejectedException; -import org.springframework.integration.core.MessageHandler; -import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.fail; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; -import static org.junit.Assert.assertFalse; -import static org.junit.Assert.fail; -import static org.mockito.Mockito.*; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.runners.MockitoJUnitRunner; + +import org.springframework.integration.Message; +import org.springframework.integration.MessageRejectedException; +import org.springframework.integration.MessagingException; +import org.springframework.integration.core.MessageHandler; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; /** @@ -123,7 +127,7 @@ public class RoundRobinDispatcherConcurrentTests { dispatcher.dispatch(message); fail("this shouldn't happen"); } - catch (IllegalStateException e) { + catch (MessagingException e) { // expected } allDone.countDown(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/UnicastingDispatcherTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/UnicastingDispatcherTests.java new file mode 100644 index 0000000000..a23ff20ed9 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/UnicastingDispatcherTests.java @@ -0,0 +1,60 @@ +/* + * Copyright 2002-2011 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.dispatcher; + +import static junit.framework.Assert.assertEquals; +import static junit.framework.Assert.assertTrue; + +import org.junit.Test; + +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.integration.Message; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.MessageDeliveryException; +import org.springframework.integration.MessagingException; +import org.springframework.integration.core.MessageHandler; +import org.springframework.integration.core.SubscribableChannel; +import org.springframework.integration.gateway.RequestReplyExchanger; +import org.springframework.integration.message.GenericMessage; + +/** + * @author Oleg Zhurakousky + * + */ +public class UnicastingDispatcherTests { + + @SuppressWarnings("unchecked") + @Test + public void withInboundGatewayAsyncRequestChannelAndExplicitErrorChannel() throws Exception{ + ApplicationContext context = new ClassPathXmlApplicationContext("unicasting-with-async.xml", this.getClass()); + SubscribableChannel errorChannel = context.getBean("errorChannel", SubscribableChannel.class); + MessageHandler errorHandler = new MessageHandler() { + + public void handleMessage(Message message) throws MessagingException { + MessageChannel replyChannel = (MessageChannel) message.getHeaders().getReplyChannel(); + assertTrue(message.getPayload() instanceof MessageDeliveryException); + replyChannel.send(new GenericMessage("reply")); + } + }; + errorChannel.subscribe(errorHandler); + + RequestReplyExchanger exchanger = context.getBean(RequestReplyExchanger.class); + Message reply = (Message) exchanger.exchange(new GenericMessage("Hello")); + assertEquals("reply", reply.getPayload()); + } + +} diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/unicasting-with-async.xml b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/unicasting-with-async.xml new file mode 100644 index 0000000000..822ec08635 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/unicasting-with-async.xml @@ -0,0 +1,17 @@ + + + + + + + + + + +