AbstractMessageHandlingEndpoint is now AbstractReplyProducingMessageConsumer.

This commit is contained in:
Mark Fisher
2008-10-06 20:59:55 +00:00
parent b567e9570d
commit f12c6b3748
14 changed files with 41 additions and 39 deletions

View File

@@ -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;

View File

@@ -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<T extends Annotation
+ inputChannel.getClass() + "]");
}
}
if (consumer instanceof AbstractMessageHandlingEndpoint) {
if (consumer instanceof AbstractReplyProducingMessageConsumer) {
String outputChannelName = (String) AnnotationUtils.getValue(annotation, OUTPUT_CHANNEL_ATTRIBUTE);
if (StringUtils.hasText(outputChannelName)) {
MessageChannel outputChannel = this.messageBus.lookupChannel(outputChannelName);
Assert.notNull(outputChannel, "unable to resolve outputChannel '" + outputChannelName + "'");
((AbstractMessageHandlingEndpoint) consumer).setOutputChannel(outputChannel);
((AbstractReplyProducingMessageConsumer) consumer).setOutputChannel(outputChannel);
}
}
}

View File

@@ -34,9 +34,11 @@ import org.springframework.integration.message.selector.MessageSelector;
import org.springframework.util.Assert;
/**
* Base class for MessageConsumers that are capable of producing replies.
*
* @author Mark Fisher
*/
public abstract class AbstractMessageHandlingEndpoint extends AbstractMessageConsumer implements ChannelRegistryAware {
public abstract class AbstractReplyProducingMessageConsumer extends AbstractMessageConsumer implements ChannelRegistryAware {
private MessageChannel outputChannel;
@@ -71,7 +73,7 @@ public abstract class AbstractMessageHandlingEndpoint extends AbstractMessageCon
@Override
protected void onMessageInternal(Message<?> message) {
protected final void onMessageInternal(Message<?> message) {
if (!this.supports(message)) {
throw new MessageRejectedException(message, "unsupported message");
}

View File

@@ -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);

View File

@@ -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;

View File

@@ -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;

View File

@@ -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;

View File

@@ -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;

View File

@@ -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<String> 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;

View File

@@ -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");
}

View File

@@ -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());
}

View File

@@ -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));