diff --git a/org.springframework.integration.adapter/src/main/java/org/springframework/integration/adapter/AbstractRemotingOutboundGateway.java b/org.springframework.integration.adapter/src/main/java/org/springframework/integration/adapter/AbstractRemotingOutboundGateway.java index b4db69d314..7affe87234 100644 --- a/org.springframework.integration.adapter/src/main/java/org/springframework/integration/adapter/AbstractRemotingOutboundGateway.java +++ b/org.springframework.integration.adapter/src/main/java/org/springframework/integration/adapter/AbstractRemotingOutboundGateway.java @@ -19,7 +19,7 @@ package org.springframework.integration.adapter; import java.io.Serializable; import org.springframework.integration.channel.MessageChannel; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.message.Message; import org.springframework.integration.message.MessageHandlingException; import org.springframework.remoting.RemoteAccessException; @@ -29,7 +29,7 @@ import org.springframework.remoting.RemoteAccessException; * * @author Mark Fisher */ -public abstract class AbstractRemotingOutboundGateway extends AbstractMessageHandlingEndpoint { +public abstract class AbstractRemotingOutboundGateway extends AbstractReplyProducingMessageConsumer { private final MessageHandler handlerProxy; diff --git a/org.springframework.integration.ws/src/main/java/org/springframework/integration/ws/handler/AbstractWebServiceHandler.java b/org.springframework.integration.ws/src/main/java/org/springframework/integration/ws/handler/AbstractWebServiceHandler.java index 5ae9dd4917..62ff4ba401 100644 --- a/org.springframework.integration.ws/src/main/java/org/springframework/integration/ws/handler/AbstractWebServiceHandler.java +++ b/org.springframework.integration.ws/src/main/java/org/springframework/integration/ws/handler/AbstractWebServiceHandler.java @@ -19,7 +19,7 @@ package org.springframework.integration.ws.handler; import java.io.IOException; import java.net.URI; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.message.GenericMessage; import org.springframework.integration.message.Message; import org.springframework.util.Assert; @@ -37,7 +37,7 @@ import org.springframework.ws.transport.WebServiceMessageSender; * * @author Mark Fisher */ -public abstract class AbstractWebServiceHandler extends AbstractMessageHandlingEndpoint { +public abstract class AbstractWebServiceHandler extends AbstractReplyProducingMessageConsumer { public static final String SOAP_ACTION_PROPERTY_KEY = "_ws.soapAction"; diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/AbstractMessageBarrierEndpoint.java b/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/AbstractMessageBarrierEndpoint.java index 1eaf4bdcd6..749aad4265 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/AbstractMessageBarrierEndpoint.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/AbstractMessageBarrierEndpoint.java @@ -31,7 +31,7 @@ import org.apache.commons.logging.LogFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.integration.channel.BlockingChannel; import org.springframework.integration.channel.MessageChannel; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.message.Message; import org.springframework.integration.message.MessageConsumer; import org.springframework.integration.message.MessageHandlingException; @@ -63,7 +63,7 @@ import org.springframework.util.ObjectUtils; * @author Mark Fisher * @author Marius Bogoevici */ -public abstract class AbstractMessageBarrierEndpoint extends AbstractMessageHandlingEndpoint implements TaskSchedulerAware, InitializingBean { +public abstract class AbstractMessageBarrierEndpoint extends AbstractReplyProducingMessageConsumer implements TaskSchedulerAware, InitializingBean { public final static long DEFAULT_SEND_TIMEOUT = 1000; diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java index b181dae985..e91259b143 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java @@ -29,7 +29,7 @@ import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.PollableChannel; import org.springframework.integration.channel.SubscribableChannel; import org.springframework.integration.endpoint.AbstractMessageConsumer; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.endpoint.MessageEndpoint; import org.springframework.integration.endpoint.PollingConsumerEndpoint; import org.springframework.integration.endpoint.SubscribingConsumerEndpoint; @@ -119,12 +119,12 @@ public abstract class AbstractMethodAnnotationPostProcessor message) { + protected final void onMessageInternal(Message message) { if (!this.supports(message)) { throw new MessageRejectedException(message, "unsupported message"); } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/ServiceActivatorEndpoint.java b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/ServiceActivatorEndpoint.java index 29f9268232..d845c28921 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/ServiceActivatorEndpoint.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/ServiceActivatorEndpoint.java @@ -31,7 +31,7 @@ import org.springframework.util.Assert; /** * @author Mark Fisher */ -public class ServiceActivatorEndpoint extends AbstractMessageHandlingEndpoint implements InitializingBean { +public class ServiceActivatorEndpoint extends AbstractReplyProducingMessageConsumer implements InitializingBean { private final MethodResolver methodResolver = new DefaultMethodResolver(ServiceActivator.class); diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/filter/FilterEndpoint.java b/org.springframework.integration/src/main/java/org/springframework/integration/filter/FilterEndpoint.java index e388f4b88f..c938a45bfa 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/filter/FilterEndpoint.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/filter/FilterEndpoint.java @@ -16,7 +16,7 @@ package org.springframework.integration.filter; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.message.Message; import org.springframework.integration.message.selector.MessageSelector; import org.springframework.util.Assert; @@ -27,7 +27,7 @@ import org.springframework.util.Assert; * * @author Mark Fisher */ -public class FilterEndpoint extends AbstractMessageHandlingEndpoint { +public class FilterEndpoint extends AbstractReplyProducingMessageConsumer { private MessageSelector selector; diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/gateway/AbstractMessagingGateway.java b/org.springframework.integration/src/main/java/org/springframework/integration/gateway/AbstractMessagingGateway.java index 73f7a6522b..c5a94749ed 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/gateway/AbstractMessagingGateway.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/gateway/AbstractMessagingGateway.java @@ -22,7 +22,7 @@ import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.MessageChannelTemplate; import org.springframework.integration.channel.PollableChannel; import org.springframework.integration.channel.SubscribableChannel; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.endpoint.MessageEndpoint; import org.springframework.integration.endpoint.MessagingGateway; import org.springframework.integration.endpoint.PollingConsumerEndpoint; @@ -150,9 +150,9 @@ public abstract class AbstractMessagingGateway implements MessagingGateway, Mess if (this.replyMessageCorrelator != null) { return; } - Assert.state(this.messageBus != null, "No MessageBus available. Cannot register ReplyMessageCorrelator."); + Assert.state(this.messageBus != null, "No MessageBus available. Cannot register reply correlator."); MessageEndpoint correlator = null; - MessageConsumer consumer = new AbstractMessageHandlingEndpoint() { + MessageConsumer consumer = new AbstractReplyProducingMessageConsumer() { @Override protected Object handle(Message message) { return message; diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/splitter/SplitterEndpoint.java b/org.springframework.integration/src/main/java/org/springframework/integration/splitter/SplitterEndpoint.java index f386f4e7db..097557fe26 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/splitter/SplitterEndpoint.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/splitter/SplitterEndpoint.java @@ -18,7 +18,7 @@ package org.springframework.integration.splitter; import java.util.List; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.message.CompositeMessage; import org.springframework.integration.message.Message; import org.springframework.util.Assert; @@ -26,7 +26,7 @@ import org.springframework.util.Assert; /** * @author Mark Fisher */ -public class SplitterEndpoint extends AbstractMessageHandlingEndpoint { +public class SplitterEndpoint extends AbstractReplyProducingMessageConsumer { private final Splitter splitter; diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/transformer/TransformerEndpoint.java b/org.springframework.integration/src/main/java/org/springframework/integration/transformer/TransformerEndpoint.java index 9a966f6273..789bb5a4b5 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/transformer/TransformerEndpoint.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/transformer/TransformerEndpoint.java @@ -16,14 +16,14 @@ package org.springframework.integration.transformer; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.message.Message; import org.springframework.util.Assert; /** * @author Mark Fisher */ -public class TransformerEndpoint extends AbstractMessageHandlingEndpoint { +public class TransformerEndpoint extends AbstractReplyProducingMessageConsumer { private final Transformer transformer; diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/bus/DefaultMessageBusTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/bus/DefaultMessageBusTests.java index 0e065674b8..27f92e715c 100644 --- a/org.springframework.integration/src/test/java/org/springframework/integration/bus/DefaultMessageBusTests.java +++ b/org.springframework.integration/src/test/java/org/springframework/integration/bus/DefaultMessageBusTests.java @@ -35,7 +35,7 @@ import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.PollableChannel; import org.springframework.integration.channel.PublishSubscribeChannel; import org.springframework.integration.channel.QueueChannel; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.endpoint.PollingConsumerEndpoint; import org.springframework.integration.endpoint.SourcePollingChannelAdapter; import org.springframework.integration.endpoint.SubscribingConsumerEndpoint; @@ -64,7 +64,7 @@ public class DefaultMessageBusTests { Message message = MessageBuilder.withPayload("test") .setReturnAddress("targetChannel").build(); sourceChannel.send(message); - AbstractMessageHandlingEndpoint consumer = new AbstractMessageHandlingEndpoint() { + AbstractReplyProducingMessageConsumer consumer = new AbstractReplyProducingMessageConsumer() { public Message handle(Message message) { return message; } @@ -119,12 +119,12 @@ public class DefaultMessageBusTests { QueueChannel inputChannel = new QueueChannel(); QueueChannel outputChannel1 = new QueueChannel(); QueueChannel outputChannel2 = new QueueChannel(); - AbstractMessageHandlingEndpoint consumer1 = new AbstractMessageHandlingEndpoint() { + AbstractReplyProducingMessageConsumer consumer1 = new AbstractReplyProducingMessageConsumer() { public Message handle(Message message) { return MessageBuilder.fromMessage(message).build(); } }; - AbstractMessageHandlingEndpoint consumer2 = new AbstractMessageHandlingEndpoint() { + AbstractReplyProducingMessageConsumer consumer2 = new AbstractReplyProducingMessageConsumer() { public Message handle(Message message) { return MessageBuilder.fromMessage(message).build(); } @@ -160,14 +160,14 @@ public class DefaultMessageBusTests { QueueChannel outputChannel1 = new QueueChannel(); QueueChannel outputChannel2 = new QueueChannel(); final CountDownLatch latch = new CountDownLatch(2); - AbstractMessageHandlingEndpoint consumer1 = new AbstractMessageHandlingEndpoint() { + AbstractReplyProducingMessageConsumer consumer1 = new AbstractReplyProducingMessageConsumer() { public Message handle(Message message) { Message reply = MessageBuilder.fromMessage(message).build(); latch.countDown(); return reply; } }; - AbstractMessageHandlingEndpoint consumer2 = new AbstractMessageHandlingEndpoint() { + AbstractReplyProducingMessageConsumer consumer2 = new AbstractReplyProducingMessageConsumer() { public Message handle(Message message) { Message reply = MessageBuilder.fromMessage(message).build(); latch.countDown(); @@ -247,7 +247,7 @@ public class DefaultMessageBusTests { errorChannel.setBeanName(ChannelRegistry.ERROR_CHANNEL_NAME); context.getBeanFactory().registerSingleton(ChannelRegistry.ERROR_CHANNEL_NAME, errorChannel); final CountDownLatch latch = new CountDownLatch(1); - AbstractMessageHandlingEndpoint consumer = new AbstractMessageHandlingEndpoint() { + AbstractReplyProducingMessageConsumer consumer = new AbstractReplyProducingMessageConsumer() { public Message handle(Message message) { latch.countDown(); return null; diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/bus/DirectChannelSubscriptionTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/bus/DirectChannelSubscriptionTests.java index f793e53567..b436671091 100644 --- a/org.springframework.integration/src/test/java/org/springframework/integration/bus/DirectChannelSubscriptionTests.java +++ b/org.springframework.integration/src/test/java/org/springframework/integration/bus/DirectChannelSubscriptionTests.java @@ -29,7 +29,7 @@ import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.channel.ThreadLocalChannel; import org.springframework.integration.config.annotation.MessagingAnnotationPostProcessor; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.endpoint.ServiceActivatorEndpoint; import org.springframework.integration.endpoint.SubscribingConsumerEndpoint; import org.springframework.integration.message.Message; @@ -93,7 +93,7 @@ public class DirectChannelSubscriptionTests { QueueChannel errorChannel = new QueueChannel(); errorChannel.setBeanName(ChannelRegistry.ERROR_CHANNEL_NAME); bus.registerChannel(errorChannel); - AbstractMessageHandlingEndpoint consumer = new AbstractMessageHandlingEndpoint() { + AbstractReplyProducingMessageConsumer consumer = new AbstractReplyProducingMessageConsumer() { public Message handle(Message message) { throw new RuntimeException("intentional test failure"); } diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/channel/MessageChannelTemplateTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/channel/MessageChannelTemplateTests.java index 1987d78303..8a90fde55c 100644 --- a/org.springframework.integration/src/test/java/org/springframework/integration/channel/MessageChannelTemplateTests.java +++ b/org.springframework.integration/src/test/java/org/springframework/integration/channel/MessageChannelTemplateTests.java @@ -29,7 +29,7 @@ import org.junit.Test; import org.springframework.context.support.GenericApplicationContext; import org.springframework.integration.bus.DefaultMessageBus; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.endpoint.PollingConsumerEndpoint; import org.springframework.integration.message.Message; import org.springframework.integration.message.MessageBuilder; @@ -47,7 +47,7 @@ public class MessageChannelTemplateTests { public void setUp() { this.requestChannel = new QueueChannel(); this.requestChannel.setBeanName("requestChannel"); - AbstractMessageHandlingEndpoint consumer = new AbstractMessageHandlingEndpoint() { + AbstractReplyProducingMessageConsumer consumer = new AbstractReplyProducingMessageConsumer() { public Message handle(Message message) { return new StringMessage(message.getPayload().toString().toUpperCase()); } diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/dispatcher/SimpleDispatcherTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/dispatcher/SimpleDispatcherTests.java index 34b7865f4d..2f3a7990d1 100644 --- a/org.springframework.integration/src/test/java/org/springframework/integration/dispatcher/SimpleDispatcherTests.java +++ b/org.springframework.integration/src/test/java/org/springframework/integration/dispatcher/SimpleDispatcherTests.java @@ -26,7 +26,7 @@ import java.util.concurrent.atomic.AtomicInteger; import org.junit.Test; -import org.springframework.integration.endpoint.AbstractMessageHandlingEndpoint; +import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer; import org.springframework.integration.endpoint.ServiceActivatorEndpoint; import org.springframework.integration.message.Message; import org.springframework.integration.message.MessageConsumer; @@ -161,9 +161,9 @@ public class SimpleDispatcherTests { final AtomicInteger counter2 = new AtomicInteger(); final AtomicInteger counter3 = new AtomicInteger(); final AtomicInteger selectorCounter = new AtomicInteger(); - AbstractMessageHandlingEndpoint endpoint1 = createEndpoint(TestHandlers.countingCountDownHandler(counter1, latch)); - AbstractMessageHandlingEndpoint endpoint2 = createEndpoint(TestHandlers.countingCountDownHandler(counter2, latch)); - AbstractMessageHandlingEndpoint endpoint3 = createEndpoint(TestHandlers.countingCountDownHandler(counter3, latch)); + AbstractReplyProducingMessageConsumer endpoint1 = createEndpoint(TestHandlers.countingCountDownHandler(counter1, latch)); + AbstractReplyProducingMessageConsumer endpoint2 = createEndpoint(TestHandlers.countingCountDownHandler(counter2, latch)); + AbstractReplyProducingMessageConsumer endpoint3 = createEndpoint(TestHandlers.countingCountDownHandler(counter3, latch)); endpoint1.setSelector(new TestMessageSelector(selectorCounter, false)); endpoint2.setSelector(new TestMessageSelector(selectorCounter, false)); endpoint3.setSelector(new TestMessageSelector(selectorCounter, true)); @@ -186,9 +186,9 @@ public class SimpleDispatcherTests { final AtomicInteger counter2 = new AtomicInteger(); final AtomicInteger counter3 = new AtomicInteger(); final AtomicInteger selectorCounter = new AtomicInteger(); - AbstractMessageHandlingEndpoint endpoint1 = createEndpoint(TestHandlers.countingCountDownHandler(counter1, latch)); - AbstractMessageHandlingEndpoint endpoint2 = createEndpoint(TestHandlers.countingCountDownHandler(counter2, latch)); - AbstractMessageHandlingEndpoint endpoint3 = createEndpoint(TestHandlers.countingCountDownHandler(counter3, latch)); + AbstractReplyProducingMessageConsumer endpoint1 = createEndpoint(TestHandlers.countingCountDownHandler(counter1, latch)); + AbstractReplyProducingMessageConsumer endpoint2 = createEndpoint(TestHandlers.countingCountDownHandler(counter2, latch)); + AbstractReplyProducingMessageConsumer endpoint3 = createEndpoint(TestHandlers.countingCountDownHandler(counter3, latch)); endpoint1.setSelector(new TestMessageSelector(selectorCounter, false)); endpoint2.setSelector(new TestMessageSelector(selectorCounter, false)); endpoint3.setSelector(new TestMessageSelector(selectorCounter, false));