Renamed ReplyHolder to ReplyMessageHolder. Modified signature of AbstractReplyMessageProducingConsumer's abstract method to be named 'onMessage'.

This commit is contained in:
Mark Fisher
2008-10-13 19:50:04 +00:00
parent 6f7bd01c2a
commit c86b7faef4
13 changed files with 45 additions and 44 deletions

View File

@@ -31,7 +31,7 @@ import org.apache.commons.logging.LogFactory;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.integration.channel.MessageChannel;
import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer;
import org.springframework.integration.endpoint.ReplyHolder;
import org.springframework.integration.endpoint.ReplyMessageHolder;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageConsumer;
import org.springframework.integration.message.MessageHandlingException;
@@ -167,7 +167,7 @@ public abstract class AbstractMessageBarrierConsumer extends AbstractReplyProduc
}
@Override
protected final void handle(Message<?> message, ReplyHolder replyHolder) {
protected final void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
if (!this.initialized) {
this.afterPropertiesSet();
}

View File

@@ -97,16 +97,16 @@ public abstract class AbstractReplyProducingMessageConsumer extends AbstractMess
if (!this.supports(message)) {
throw new MessageRejectedException(message, "unsupported message");
}
ReplyHolder replyHolder = new ReplyHolder();
this.handle(message, replyHolder);
if (replyHolder.isEmpty()) {
ReplyMessageHolder replyMessageHolder = new ReplyMessageHolder();
this.onMessage(message, replyMessageHolder);
if (replyMessageHolder.isEmpty()) {
if (this.requiresReply) {
throw new MessageHandlingException(message, "consumer '" + this
+ "' requires a reply, but no reply was received");
}
return;
}
Object targetChannelValue = replyHolder.getTargetChannel();
Object targetChannelValue = replyMessageHolder.getTargetChannel();
MessageChannel replyChannel = null;
if (targetChannelValue == null) {
replyChannel = this.resolveReplyChannel(message);
@@ -115,13 +115,13 @@ public abstract class AbstractReplyProducingMessageConsumer extends AbstractMess
replyChannel = this.channelResolver.resolveChannelName((String) targetChannelValue);
}
MessageHeaders requestHeaders = message.getHeaders();
for (MessageBuilder<?> builder : replyHolder.builders()) {
for (MessageBuilder<?> builder : replyMessageHolder.builders()) {
builder.copyHeadersIfAbsent(requestHeaders);
this.sendReplyMessage(builder.build(), replyChannel);
}
}
protected abstract void handle(Message<?> message, ReplyHolder replyHolder);
protected abstract void onMessage(Message<?> requestMessage, ReplyMessageHolder replyMessageHolder);
protected boolean supports(Message<?> message) {
if (this.selector != null && !this.selector.accept(message)) {

View File

@@ -27,19 +27,19 @@ import org.springframework.integration.message.MessageBuilder;
/**
* @author Mark Fisher
*/
public class ReplyHolder {
public class ReplyMessageHolder {
private final List<MessageBuilder<?>> builders = new ArrayList<MessageBuilder<?>>();
private volatile Object targetChannel;
public MessageBuilder<?> set(Object replyObject) {
return this.createAndAddBuilder(replyObject, true);
public MessageBuilder<?> set(Object messageOrPayload) {
return this.createAndAddBuilder(messageOrPayload, true);
}
public MessageBuilder<?> add(Object replyObject) {
return this.createAndAddBuilder(replyObject, false);
public MessageBuilder<?> add(Object messageOrPayload) {
return this.createAndAddBuilder(messageOrPayload, false);
}
public void setTargetChannel(MessageChannel targetChannel) {
@@ -62,16 +62,16 @@ public class ReplyHolder {
return Collections.unmodifiableList(this.builders);
}
private MessageBuilder<?> createAndAddBuilder(Object replyObject, boolean clearExistingValues) {
private MessageBuilder<?> createAndAddBuilder(Object messageOrPayload, boolean clearExistingValues) {
MessageBuilder<?> builder = null;
if (replyObject instanceof MessageBuilder) {
builder = (MessageBuilder<?>) replyObject;
if (messageOrPayload instanceof MessageBuilder) {
builder = (MessageBuilder<?>) messageOrPayload;
}
else if (replyObject instanceof Message) {
builder = MessageBuilder.fromMessage((Message<?>) replyObject);
else if (messageOrPayload instanceof Message) {
builder = MessageBuilder.fromMessage((Message<?>) messageOrPayload);
}
else {
builder = MessageBuilder.withPayload(replyObject);
builder = MessageBuilder.withPayload(messageOrPayload);
}
synchronized (this.builders) {
if (clearExistingValues) {

View File

@@ -63,7 +63,7 @@ public class ServiceActivatorEndpoint extends AbstractReplyProducingMessageConsu
}
@Override
protected void handle(Message<?> message, ReplyHolder replyHolder) {
protected void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
try {
Object result = this.invoker.invokeMethod(message);
if (result != null) {

View File

@@ -17,7 +17,7 @@
package org.springframework.integration.filter;
import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer;
import org.springframework.integration.endpoint.ReplyHolder;
import org.springframework.integration.endpoint.ReplyMessageHolder;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.selector.MessageSelector;
import org.springframework.util.Assert;
@@ -41,7 +41,7 @@ public class MessageFilter extends AbstractReplyProducingMessageConsumer {
@Override
protected void handle(Message<?> message, ReplyHolder replyHolder) {
protected void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
if (this.selector.accept(message)) {
replyHolder.set(message);
}

View File

@@ -26,7 +26,7 @@ import org.springframework.integration.endpoint.AbstractReplyProducingMessageCon
import org.springframework.integration.endpoint.MessageEndpoint;
import org.springframework.integration.endpoint.MessagingGateway;
import org.springframework.integration.endpoint.PollingConsumerEndpoint;
import org.springframework.integration.endpoint.ReplyHolder;
import org.springframework.integration.endpoint.ReplyMessageHolder;
import org.springframework.integration.endpoint.SubscribingConsumerEndpoint;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageConsumer;
@@ -155,7 +155,7 @@ public abstract class AbstractMessagingGateway implements MessagingGateway, Mess
MessageEndpoint correlator = null;
MessageConsumer consumer = new AbstractReplyProducingMessageConsumer() {
@Override
protected void handle(Message<?> message, ReplyHolder replyHolder) {
protected void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
replyHolder.set(message);
}
};

View File

@@ -19,7 +19,7 @@ package org.springframework.integration.splitter;
import java.util.Collection;
import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer;
import org.springframework.integration.endpoint.ReplyHolder;
import org.springframework.integration.endpoint.ReplyMessageHolder;
import org.springframework.integration.message.Message;
/**
@@ -30,7 +30,7 @@ import org.springframework.integration.message.Message;
public abstract class AbstractMessageSplitter extends AbstractReplyProducingMessageConsumer {
@Override
protected final void handle(Message<?> message, ReplyHolder replyHolder) {
protected final void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
Object result = this.splitMessage(message);
if (result == null) {
return;
@@ -57,7 +57,7 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess
}
}
private void addReply(ReplyHolder replyHolder, Object item, Object correlationId, int sequenceNumber, int sequenceSize) {
private void addReply(ReplyMessageHolder replyHolder, Object item, Object correlationId, int sequenceNumber, int sequenceSize) {
replyHolder.add(item).setCorrelationId(correlationId)
.setSequenceNumber(sequenceNumber)
.setSequenceSize(sequenceSize);

View File

@@ -17,7 +17,7 @@
package org.springframework.integration.transformer;
import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer;
import org.springframework.integration.endpoint.ReplyHolder;
import org.springframework.integration.endpoint.ReplyMessageHolder;
import org.springframework.integration.message.Message;
import org.springframework.util.Assert;
@@ -44,7 +44,7 @@ public class MessageTransformingConsumer extends AbstractReplyProducingMessageCo
@Override
protected void handle(Message<?> message, ReplyHolder replyHolder) {
protected void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
Message<?> result = transformer.transform(message);
if (result != null) {
replyHolder.set(result);

View File

@@ -37,7 +37,7 @@ import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.config.xml.MessageBusParser;
import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer;
import org.springframework.integration.endpoint.PollingConsumerEndpoint;
import org.springframework.integration.endpoint.ReplyHolder;
import org.springframework.integration.endpoint.ReplyMessageHolder;
import org.springframework.integration.endpoint.SourcePollingChannelAdapter;
import org.springframework.integration.endpoint.SubscribingConsumerEndpoint;
import org.springframework.integration.message.ErrorMessage;
@@ -67,7 +67,7 @@ public class DefaultMessageBusTests {
.setReturnAddress("targetChannel").build();
sourceChannel.send(message);
AbstractReplyProducingMessageConsumer consumer = new AbstractReplyProducingMessageConsumer() {
public void handle(Message<?> message, ReplyHolder replyHolder) {
public void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
replyHolder.set(message);
}
};
@@ -125,12 +125,12 @@ public class DefaultMessageBusTests {
QueueChannel outputChannel1 = new QueueChannel();
QueueChannel outputChannel2 = new QueueChannel();
AbstractReplyProducingMessageConsumer consumer1 = new AbstractReplyProducingMessageConsumer() {
public void handle(Message<?> message, ReplyHolder replyHolder) {
public void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
replyHolder.set(message);
}
};
AbstractReplyProducingMessageConsumer consumer2 = new AbstractReplyProducingMessageConsumer() {
public void handle(Message<?> message, ReplyHolder replyHolder) {
public void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
replyHolder.set(message);
}
};
@@ -167,13 +167,13 @@ public class DefaultMessageBusTests {
QueueChannel outputChannel2 = new QueueChannel();
final CountDownLatch latch = new CountDownLatch(2);
AbstractReplyProducingMessageConsumer consumer1 = new AbstractReplyProducingMessageConsumer() {
public void handle(Message<?> message, ReplyHolder replyHolder) {
public void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
replyHolder.set(message);
latch.countDown();
}
};
AbstractReplyProducingMessageConsumer consumer2 = new AbstractReplyProducingMessageConsumer() {
public void handle(Message<?> message, ReplyHolder replyHolder) {
public void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
replyHolder.set(message);
latch.countDown();
}
@@ -245,7 +245,7 @@ public class DefaultMessageBusTests {
context.getBeanFactory().registerSingleton(DefaultMessageBus.ERROR_CHANNEL_BEAN_NAME, errorChannel);
final CountDownLatch latch = new CountDownLatch(1);
AbstractReplyProducingMessageConsumer consumer = new AbstractReplyProducingMessageConsumer() {
public void handle(Message<?> message, ReplyHolder replyHolder) {
public void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
latch.countDown();
}
};

View File

@@ -30,7 +30,7 @@ import org.springframework.integration.channel.ThreadLocalChannel;
import org.springframework.integration.config.annotation.MessagingAnnotationPostProcessor;
import org.springframework.integration.config.xml.MessageBusParser;
import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer;
import org.springframework.integration.endpoint.ReplyHolder;
import org.springframework.integration.endpoint.ReplyMessageHolder;
import org.springframework.integration.endpoint.ServiceActivatorEndpoint;
import org.springframework.integration.endpoint.SubscribingConsumerEndpoint;
import org.springframework.integration.message.Message;
@@ -96,7 +96,7 @@ public class DirectChannelSubscriptionTests {
@Test(expected = MessagingException.class)
public void exceptionThrownFromRegisteredEndpoint() {
AbstractReplyProducingMessageConsumer consumer = new AbstractReplyProducingMessageConsumer() {
public void handle(Message<?> message, ReplyHolder replyHolder) {
public void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
throw new RuntimeException("intentional test failure");
}
};

View File

@@ -33,7 +33,7 @@ import org.springframework.context.support.GenericApplicationContext;
import org.springframework.integration.bus.DefaultMessageBus;
import org.springframework.integration.endpoint.AbstractReplyProducingMessageConsumer;
import org.springframework.integration.endpoint.PollingConsumerEndpoint;
import org.springframework.integration.endpoint.ReplyHolder;
import org.springframework.integration.endpoint.ReplyMessageHolder;
import org.springframework.integration.message.Message;
import org.springframework.integration.message.MessageBuilder;
import org.springframework.integration.message.StringMessage;
@@ -52,7 +52,7 @@ public class MessageChannelTemplateTests {
this.requestChannel = new QueueChannel();
this.requestChannel.setBeanName("requestChannel");
AbstractReplyProducingMessageConsumer consumer = new AbstractReplyProducingMessageConsumer() {
public void handle(Message<?> message, ReplyHolder replyHolder) {
public void onMessage(Message<?> message, ReplyMessageHolder replyHolder) {
replyHolder.set(message.getPayload().toString().toUpperCase());
}
};