INT-1095 Aggregator now properly handles the resolution of replyChannel by name.

This commit is contained in:
Mark Fisher
2010-04-25 16:06:41 +00:00
parent a6aff09329
commit a70ad276db
5 changed files with 111 additions and 9 deletions

View File

@@ -84,8 +84,6 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements
private volatile MessageChannel discardChannel = new NullChannel();
private volatile ChannelResolver channelResolver;
private final IdTracker tracker = new IdTracker();
private final BlockingQueue<DelayedKey> keysInBuffer = new DelayQueue<DelayedKey>();
@@ -153,7 +151,7 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements
}
public void setChannelResolver(ChannelResolver channelResolver) {
this.channelResolver = channelResolver;
super.setChannelResolver(channelResolver);
}
public void setDiscardChannel(MessageChannel discardChannel) {
@@ -192,7 +190,7 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements
logger.debug("Completing group with correlationKey [" + correlationKey + "]");
}
outputProcessor.processAndSend(group, channelTemplate,
this.resolveReplyChannel(message, this.outputChannel, this.channelResolver));
this.resolveReplyChannel(message, this.outputChannel));
}
}
else {
@@ -292,8 +290,7 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements
MessageGroup group = new MessageGroup(all, completionStrategy, key, deleteOrTrackCallback());
if (all.size() > 0) {
// last chance for normal completion
MessageChannel outputChannel = resolveReplyChannel(
all.get(0), this.outputChannel, this.channelResolver);
MessageChannel outputChannel = resolveReplyChannel(all.get(0), this.outputChannel);
boolean processed = false;
if (group.isComplete()) {
outputProcessor.processAndSend(group, channelTemplate, outputChannel);

View File

@@ -74,8 +74,8 @@ public abstract class AbstractMessageHandler extends IntegrationObjectSupport im
protected abstract void handleMessageInternal(Message<?> message) throws Exception;
protected final MessageChannel resolveReplyChannel(Message<?> requestMessage,
MessageChannel defaultOutputChannel,
ChannelResolver channelResolver) {
MessageChannel defaultOutputChannel) {
ChannelResolver channelResolver = this.getChannelResolver();
MessageChannel replyChannel = defaultOutputChannel;
if (replyChannel == null) {
Object replyChannelHeader = requestMessage.getHeaders().getReplyChannel();

View File

@@ -99,7 +99,7 @@ public abstract class AbstractReplyProducingMessageHandler extends AbstractMessa
}
return;
}
MessageChannel replyChannel = resolveReplyChannel(message, this.outputChannel, this.getChannelResolver());
MessageChannel replyChannel = resolveReplyChannel(message, this.outputChannel);
MessageHeaders requestHeaders = message.getHeaders();
this.handleResult(result, requestHeaders, replyChannel);
}