INT-1377, merged AbstractSingleChannelNameRouter with AbstractMessageRouter, fixed corresponding tests, removed dependencies on it across the workspace
This commit is contained in:
@@ -20,7 +20,6 @@ import org.springframework.aop.framework.Advised;
|
||||
import org.springframework.expression.Expression;
|
||||
import org.springframework.integration.MessageChannel;
|
||||
import org.springframework.integration.core.MessageHandler;
|
||||
import org.springframework.integration.router.AbstractChannelNameResolvingMessageRouter;
|
||||
import org.springframework.integration.router.AbstractMessageRouter;
|
||||
import org.springframework.integration.router.ExpressionEvaluatingRouter;
|
||||
import org.springframework.integration.router.MethodInvokingRouter;
|
||||
@@ -144,11 +143,11 @@ public class RouterFactoryBean extends AbstractMessageHandlerFactoryBean {
|
||||
}
|
||||
|
||||
private AbstractMessageRouter configureRouter(AbstractMessageRouter router) {
|
||||
if (this.channelResolver != null && router instanceof AbstractChannelNameResolvingMessageRouter) {
|
||||
((AbstractChannelNameResolvingMessageRouter) router).setChannelResolver(this.channelResolver);
|
||||
if (this.channelResolver != null && router instanceof AbstractMessageRouter) {
|
||||
((AbstractMessageRouter) router).setChannelResolver(this.channelResolver);
|
||||
}
|
||||
if (this.channelIdentifierMap != null && router instanceof AbstractChannelNameResolvingMessageRouter) {
|
||||
((AbstractChannelNameResolvingMessageRouter) router).setChannelIdentifierMap(this.channelIdentifierMap);
|
||||
if (this.channelIdentifierMap != null && router instanceof AbstractMessageRouter) {
|
||||
((AbstractMessageRouter) router).setChannelIdentifierMap(this.channelIdentifierMap);
|
||||
}
|
||||
if (this.defaultOutputChannel != null) {
|
||||
router.setDefaultOutputChannel(this.defaultOutputChannel);
|
||||
@@ -157,10 +156,10 @@ public class RouterFactoryBean extends AbstractMessageHandlerFactoryBean {
|
||||
router.setTimeout(timeout.longValue());
|
||||
}
|
||||
if (this.ignoreChannelNameResolutionFailures != null) {
|
||||
Assert.isTrue(router instanceof AbstractChannelNameResolvingMessageRouter,
|
||||
Assert.isTrue(router instanceof AbstractMessageRouter,
|
||||
"The 'ignoreChannelNameResolutionFailures' property can only be set on routers that extend "
|
||||
+ AbstractChannelNameResolvingMessageRouter.class.getName());
|
||||
((AbstractChannelNameResolvingMessageRouter) router)
|
||||
+ AbstractMessageRouter.class.getName());
|
||||
((AbstractMessageRouter) router)
|
||||
.setIgnoreChannelNameResolutionFailures(ignoreChannelNameResolutionFailures);
|
||||
}
|
||||
if (this.applySequence != null) {
|
||||
|
||||
@@ -1,219 +0,0 @@
|
||||
/*
|
||||
* Copyright 2002-2010 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.router;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
import org.springframework.core.convert.ConversionService;
|
||||
import org.springframework.core.convert.support.ConversionServiceFactory;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.MessageChannel;
|
||||
import org.springframework.integration.MessagingException;
|
||||
import org.springframework.integration.support.channel.BeanFactoryChannelResolver;
|
||||
import org.springframework.integration.support.channel.ChannelResolutionException;
|
||||
import org.springframework.integration.support.channel.ChannelResolver;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
* A base class for router implementations that return only the channel name(s)
|
||||
* rather than {@link MessageChannel} instances.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
* @author Jonas Partner
|
||||
* @author Oleg Zhurakousky
|
||||
*/
|
||||
public abstract class AbstractChannelNameResolvingMessageRouter extends AbstractMessageRouter {
|
||||
|
||||
private volatile String prefix;
|
||||
|
||||
private volatile String suffix;
|
||||
|
||||
private volatile ChannelResolver channelResolver;
|
||||
|
||||
private volatile boolean ignoreChannelNameResolutionFailures;
|
||||
|
||||
protected volatile Map<String, String> channelIdentifierMap;
|
||||
|
||||
|
||||
/**
|
||||
* Specify the {@link ChannelResolver} strategy to use.
|
||||
* The default is a BeanFactoryChannelResolver.
|
||||
*/
|
||||
public void setChannelResolver(ChannelResolver channelResolver) {
|
||||
Assert.notNull(channelResolver, "'channelResolver' must not be null");
|
||||
this.channelResolver = channelResolver;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify a prefix to be added to each channel name prior to resolution.
|
||||
*/
|
||||
public void setPrefix(String prefix) {
|
||||
this.prefix = prefix;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify a suffix to be added to each channel name prior to resolution.
|
||||
*/
|
||||
public void setSuffix(String suffix) {
|
||||
this.suffix = suffix;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify whether this router should ignore any failure to resolve a channel name to
|
||||
* an actual MessageChannel instance when delegating to the ChannelResolver strategy.
|
||||
*/
|
||||
public void setIgnoreChannelNameResolutionFailures(boolean ignoreChannelNameResolutionFailures) {
|
||||
this.ignoreChannelNameResolutionFailures = ignoreChannelNameResolutionFailures;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onInit() {
|
||||
BeanFactory beanFactory = this.getBeanFactory();
|
||||
if (this.channelResolver == null && beanFactory != null) {
|
||||
this.channelResolver = new BeanFactoryChannelResolver(beanFactory);
|
||||
}
|
||||
}
|
||||
|
||||
private MessageChannel resolveChannelForName(String channelName, Message<?> message) {
|
||||
Assert.state(this.channelResolver != null,
|
||||
"unable to resolve channel names, no ChannelResolver available");
|
||||
MessageChannel channel = null;
|
||||
try {
|
||||
channel = this.channelResolver.resolveChannelName(channelName);
|
||||
}
|
||||
catch (ChannelResolutionException e) {
|
||||
if (!this.ignoreChannelNameResolutionFailures) {
|
||||
throw new MessagingException(message,
|
||||
"failed to resolve channel name '" + channelName + "'", e);
|
||||
}
|
||||
}
|
||||
if (channel == null && !this.ignoreChannelNameResolutionFailures) {
|
||||
throw new MessagingException(message,
|
||||
"failed to resolve channel name '" + channelName + "'");
|
||||
}
|
||||
return channel;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected Collection<MessageChannel> determineTargetChannels(Message<?> message) {
|
||||
this.afterPropertiesSet();
|
||||
Collection<MessageChannel> channels = new ArrayList<MessageChannel>();
|
||||
Collection<Object> channelsReturned = this.getChannelIndicatorList(message);
|
||||
addToCollection(channels, channelsReturned, message);
|
||||
return channels;
|
||||
}
|
||||
|
||||
private void addToCollection(Collection<MessageChannel> channels, Collection<?> channelIndicators, Message<?> message) {
|
||||
if (channelIndicators == null) {
|
||||
return;
|
||||
}
|
||||
for (Object channelIndicator : channelIndicators) {
|
||||
if (channelIndicator == null) {
|
||||
continue;
|
||||
}
|
||||
else if (channelIndicator instanceof MessageChannel) {
|
||||
channels.add((MessageChannel) channelIndicator);
|
||||
}
|
||||
else if (channelIndicator instanceof MessageChannel[]) {
|
||||
channels.addAll(Arrays.asList((MessageChannel[]) channelIndicator));
|
||||
}
|
||||
else if (channelIndicator instanceof String) {
|
||||
addChannelFromString(channels, (String) channelIndicator, message);
|
||||
}
|
||||
else if (channelIndicator instanceof String[]) {
|
||||
for (String indicatorName : (String[]) channelIndicator) {
|
||||
addChannelFromString(channels, indicatorName, message);
|
||||
}
|
||||
}
|
||||
else if (channelIndicator instanceof Collection) {
|
||||
addToCollection(channels, (Collection<?>) channelIndicator, message);
|
||||
}
|
||||
else if (this.getRequiredConversionService().canConvert(channelIndicator.getClass(), String.class)) {
|
||||
addChannelFromString(channels,
|
||||
this.getConversionService().convert(channelIndicator, String.class), message);
|
||||
}
|
||||
else {
|
||||
throw new MessagingException(
|
||||
"unsupported return type for router [" + channelIndicator.getClass() + "]");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void addChannelFromString(Collection<MessageChannel> channels, String channelIdentifier, Message<?> message) {
|
||||
if (channelIdentifier.indexOf(',') != -1) {
|
||||
for (String name : StringUtils.commaDelimitedListToStringArray(channelIdentifier)) {
|
||||
addChannelFromString(channels, name, message);
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (this.prefix != null) {
|
||||
channelIdentifier = this.prefix + channelIdentifier;
|
||||
}
|
||||
if (this.suffix != null) {
|
||||
channelIdentifier = channelIdentifier + suffix;
|
||||
}
|
||||
/*
|
||||
* Some routers due to their complex nature will already resolve 'channelIdentifier'
|
||||
* to 'channelName' (e.g., PTR, EMETR)
|
||||
*/
|
||||
String channelName = channelIdentifier;
|
||||
if (channelIdentifierMap != null && channelIdentifierMap.containsKey(channelIdentifier)){
|
||||
channelName = channelIdentifierMap.get(channelIdentifier);
|
||||
}
|
||||
|
||||
if (this.channelResolver != null){
|
||||
MessageChannel channel = resolveChannelForName(channelName, message);
|
||||
if (channel != null) {
|
||||
channels.add(channel);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
protected ConversionService getRequiredConversionService() {
|
||||
if (this.getConversionService() == null) {
|
||||
this.setConversionService(ConversionServiceFactory.createDefaultConversionService());
|
||||
}
|
||||
return this.getConversionService();
|
||||
}
|
||||
|
||||
public Map<String, String> getChannelIdentifierMap() {
|
||||
return channelIdentifierMap;
|
||||
}
|
||||
|
||||
public void setChannelIdentifierMap(Map<String, String> channelIdentifierMap) {
|
||||
this.channelIdentifierMap = channelIdentifierMap;
|
||||
}
|
||||
|
||||
public void setChannelMapping(String channelIdentifier, String channelName){
|
||||
this.channelIdentifierMap.put(channelIdentifier, channelName);
|
||||
}
|
||||
|
||||
public void removeChannelMapping(String channelIdentifier){
|
||||
this.channelIdentifierMap.remove(channelIdentifier);
|
||||
}
|
||||
/**
|
||||
* Subclasses must implement this method to return the channel indicators.
|
||||
*/
|
||||
protected abstract List<Object> getChannelIndicatorList(Message<?> message);
|
||||
|
||||
}
|
||||
@@ -32,7 +32,7 @@ import org.springframework.util.Assert;
|
||||
* @author Mark Fisher
|
||||
* @since 2.0
|
||||
*/
|
||||
class AbstractMessageProcessingRouter extends AbstractChannelNameResolvingMessageRouter {
|
||||
class AbstractMessageProcessingRouter extends AbstractMessageRouter {
|
||||
|
||||
private final MessageProcessor<Object> messageProcessor;
|
||||
|
||||
|
||||
@@ -16,8 +16,15 @@
|
||||
|
||||
package org.springframework.integration.router;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.Collection;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
import org.springframework.core.convert.ConversionService;
|
||||
import org.springframework.core.convert.support.ConversionServiceFactory;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.MessageChannel;
|
||||
import org.springframework.integration.MessageDeliveryException;
|
||||
@@ -25,11 +32,17 @@ import org.springframework.integration.MessagingException;
|
||||
import org.springframework.integration.core.MessagingTemplate;
|
||||
import org.springframework.integration.handler.AbstractMessageHandler;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.integration.support.channel.BeanFactoryChannelResolver;
|
||||
import org.springframework.integration.support.channel.ChannelResolutionException;
|
||||
import org.springframework.integration.support.channel.ChannelResolver;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
* Base class for Message Routers.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
* @author Oleg Zhurakousky
|
||||
*/
|
||||
public abstract class AbstractMessageRouter extends AbstractMessageHandler {
|
||||
|
||||
@@ -42,6 +55,177 @@ public abstract class AbstractMessageRouter extends AbstractMessageHandler {
|
||||
private volatile boolean applySequence;
|
||||
|
||||
private final MessagingTemplate messagingTemplate = new MessagingTemplate();
|
||||
|
||||
private volatile String prefix;
|
||||
|
||||
private volatile String suffix;
|
||||
|
||||
private volatile ChannelResolver channelResolver;
|
||||
|
||||
private volatile boolean ignoreChannelNameResolutionFailures;
|
||||
|
||||
protected volatile Map<String, String> channelIdentifierMap;
|
||||
|
||||
/**
|
||||
* Specify the {@link ChannelResolver} strategy to use.
|
||||
* The default is a BeanFactoryChannelResolver.
|
||||
*/
|
||||
public void setChannelResolver(ChannelResolver channelResolver) {
|
||||
Assert.notNull(channelResolver, "'channelResolver' must not be null");
|
||||
this.channelResolver = channelResolver;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify a prefix to be added to each channel name prior to resolution.
|
||||
*/
|
||||
public void setPrefix(String prefix) {
|
||||
this.prefix = prefix;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify a suffix to be added to each channel name prior to resolution.
|
||||
*/
|
||||
public void setSuffix(String suffix) {
|
||||
this.suffix = suffix;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify whether this router should ignore any failure to resolve a channel name to
|
||||
* an actual MessageChannel instance when delegating to the ChannelResolver strategy.
|
||||
*/
|
||||
public void setIgnoreChannelNameResolutionFailures(boolean ignoreChannelNameResolutionFailures) {
|
||||
this.ignoreChannelNameResolutionFailures = ignoreChannelNameResolutionFailures;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onInit() {
|
||||
BeanFactory beanFactory = this.getBeanFactory();
|
||||
if (this.channelResolver == null && beanFactory != null) {
|
||||
this.channelResolver = new BeanFactoryChannelResolver(beanFactory);
|
||||
}
|
||||
}
|
||||
|
||||
private MessageChannel resolveChannelForName(String channelName, Message<?> message) {
|
||||
Assert.state(this.channelResolver != null,
|
||||
"unable to resolve channel names, no ChannelResolver available");
|
||||
MessageChannel channel = null;
|
||||
try {
|
||||
channel = this.channelResolver.resolveChannelName(channelName);
|
||||
}
|
||||
catch (ChannelResolutionException e) {
|
||||
if (!this.ignoreChannelNameResolutionFailures) {
|
||||
throw new MessagingException(message,
|
||||
"failed to resolve channel name '" + channelName + "'", e);
|
||||
}
|
||||
}
|
||||
if (channel == null && !this.ignoreChannelNameResolutionFailures) {
|
||||
throw new MessagingException(message,
|
||||
"failed to resolve channel name '" + channelName + "'");
|
||||
}
|
||||
return channel;
|
||||
}
|
||||
|
||||
private void addChannelFromString(Collection<MessageChannel> channels, String channelIdentifier, Message<?> message) {
|
||||
if (channelIdentifier.indexOf(',') != -1) {
|
||||
for (String name : StringUtils.commaDelimitedListToStringArray(channelIdentifier)) {
|
||||
addChannelFromString(channels, name, message);
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (this.prefix != null) {
|
||||
channelIdentifier = this.prefix + channelIdentifier;
|
||||
}
|
||||
if (this.suffix != null) {
|
||||
channelIdentifier = channelIdentifier + suffix;
|
||||
}
|
||||
/*
|
||||
* Some routers due to their complex nature will already resolve 'channelIdentifier'
|
||||
* to 'channelName' (e.g., PTR, EMETR)
|
||||
*/
|
||||
String channelName = channelIdentifier;
|
||||
if (channelIdentifierMap != null && channelIdentifierMap.containsKey(channelIdentifier)){
|
||||
channelName = channelIdentifierMap.get(channelIdentifier);
|
||||
}
|
||||
|
||||
if (this.channelResolver != null){
|
||||
MessageChannel channel = resolveChannelForName(channelName, message);
|
||||
if (channel != null) {
|
||||
channels.add(channel);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void addToCollection(Collection<MessageChannel> channels, Collection<?> channelIndicators, Message<?> message) {
|
||||
if (channelIndicators == null) {
|
||||
return;
|
||||
}
|
||||
for (Object channelIndicator : channelIndicators) {
|
||||
if (channelIndicator == null) {
|
||||
continue;
|
||||
}
|
||||
else if (channelIndicator instanceof MessageChannel) {
|
||||
channels.add((MessageChannel) channelIndicator);
|
||||
}
|
||||
else if (channelIndicator instanceof MessageChannel[]) {
|
||||
channels.addAll(Arrays.asList((MessageChannel[]) channelIndicator));
|
||||
}
|
||||
else if (channelIndicator instanceof String) {
|
||||
addChannelFromString(channels, (String) channelIndicator, message);
|
||||
}
|
||||
else if (channelIndicator instanceof String[]) {
|
||||
for (String indicatorName : (String[]) channelIndicator) {
|
||||
addChannelFromString(channels, indicatorName, message);
|
||||
}
|
||||
}
|
||||
else if (channelIndicator instanceof Collection) {
|
||||
addToCollection(channels, (Collection<?>) channelIndicator, message);
|
||||
}
|
||||
else if (this.getRequiredConversionService().canConvert(channelIndicator.getClass(), String.class)) {
|
||||
addChannelFromString(channels,
|
||||
this.getConversionService().convert(channelIndicator, String.class), message);
|
||||
}
|
||||
else {
|
||||
throw new MessagingException(
|
||||
"unsupported return type for router [" + channelIndicator.getClass() + "]");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
protected Collection<MessageChannel> determineTargetChannels(Message<?> message) {
|
||||
this.afterPropertiesSet();
|
||||
Collection<MessageChannel> channels = new ArrayList<MessageChannel>();
|
||||
Collection<Object> channelsReturned = this.getChannelIndicatorList(message);
|
||||
addToCollection(channels, channelsReturned, message);
|
||||
return channels;
|
||||
}
|
||||
|
||||
protected ConversionService getRequiredConversionService() {
|
||||
if (this.getConversionService() == null) {
|
||||
this.setConversionService(ConversionServiceFactory.createDefaultConversionService());
|
||||
}
|
||||
return this.getConversionService();
|
||||
}
|
||||
|
||||
public Map<String, String> getChannelIdentifierMap() {
|
||||
return channelIdentifierMap;
|
||||
}
|
||||
|
||||
public void setChannelIdentifierMap(Map<String, String> channelIdentifierMap) {
|
||||
this.channelIdentifierMap = channelIdentifierMap;
|
||||
}
|
||||
|
||||
public void setChannelMapping(String channelIdentifier, String channelName){
|
||||
this.channelIdentifierMap.put(channelIdentifier, channelName);
|
||||
}
|
||||
|
||||
public void removeChannelMapping(String channelIdentifier){
|
||||
this.channelIdentifierMap.remove(channelIdentifier);
|
||||
}
|
||||
/**
|
||||
* Subclasses must implement this method to return the channel indicators.
|
||||
*/
|
||||
protected abstract List<Object> getChannelIndicatorList(Message<?> message);
|
||||
|
||||
/**
|
||||
* Set the default channel where Messages should be sent if channel resolution fails to return any channels. If no
|
||||
@@ -137,9 +321,9 @@ public abstract class AbstractMessageRouter extends AbstractMessageHandler {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Subclasses must implement this method to return the target channels for a given Message.
|
||||
*/
|
||||
protected abstract Collection<MessageChannel> determineTargetChannels(Message<?> message);
|
||||
// /**
|
||||
// * Subclasses must implement this method to return the target channels for a given Message.
|
||||
// */
|
||||
// protected abstract Collection<MessageChannel> determineTargetChannels(Message<?> message);
|
||||
|
||||
}
|
||||
|
||||
@@ -27,7 +27,7 @@ import org.springframework.integration.Message;
|
||||
*
|
||||
* @author Mark Fisher
|
||||
*/
|
||||
public abstract class AbstractSingleChannelNameRouter extends AbstractChannelNameResolvingMessageRouter {
|
||||
public abstract class AbstractSingleChannelNameRouter extends AbstractMessageRouter {
|
||||
|
||||
@Override
|
||||
protected final List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
|
||||
@@ -30,7 +30,7 @@ import org.springframework.integration.MessageChannel;
|
||||
* @author Mark Fisher
|
||||
* @author Oleg Zhurakousky
|
||||
*/
|
||||
public class ErrorMessageExceptionTypeRouter extends AbstractChannelNameResolvingMessageRouter {
|
||||
public class ErrorMessageExceptionTypeRouter extends AbstractMessageRouter {
|
||||
|
||||
@Override
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
|
||||
@@ -30,7 +30,7 @@ import org.springframework.util.StringUtils;
|
||||
* @author Mark Fisher
|
||||
* @since 1.0.3
|
||||
*/
|
||||
public class HeaderValueRouter extends AbstractChannelNameResolvingMessageRouter {
|
||||
public class HeaderValueRouter extends AbstractMessageRouter {
|
||||
|
||||
private final String headerName;
|
||||
|
||||
|
||||
@@ -30,7 +30,7 @@ import org.springframework.util.StringUtils;
|
||||
* @author Mark Fisher
|
||||
* @author Oleg Zhurakousky
|
||||
*/
|
||||
public class PayloadTypeRouter extends AbstractChannelNameResolvingMessageRouter {
|
||||
public class PayloadTypeRouter extends AbstractMessageRouter {
|
||||
|
||||
@Override
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
|
||||
@@ -89,17 +89,17 @@ public class RecipientListRouter extends AbstractMessageRouter implements Initia
|
||||
Assert.notEmpty(this.recipients, "a non-empty recipient list is required");
|
||||
}
|
||||
|
||||
@Override
|
||||
protected Collection<MessageChannel> determineTargetChannels(Message<?> message) {
|
||||
List<MessageChannel> channels = new ArrayList<MessageChannel>();
|
||||
List<Recipient> recipientList = this.recipients;
|
||||
for (Recipient recipient : recipientList) {
|
||||
if (recipient.accept(message)) {
|
||||
channels.add(recipient.getChannel());
|
||||
}
|
||||
}
|
||||
return channels;
|
||||
}
|
||||
// @Override
|
||||
// protected Collection<MessageChannel> determineTargetChannels(Message<?> message) {
|
||||
// List<MessageChannel> channels = new ArrayList<MessageChannel>();
|
||||
// List<Recipient> recipientList = this.recipients;
|
||||
// for (Recipient recipient : recipientList) {
|
||||
// if (recipient.accept(message)) {
|
||||
// channels.add(recipient.getChannel());
|
||||
// }
|
||||
// }
|
||||
// return channels;
|
||||
// }
|
||||
|
||||
|
||||
public static class Recipient {
|
||||
@@ -126,4 +126,17 @@ public class RecipientListRouter extends AbstractMessageRouter implements Initia
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
List<Object> channels = new ArrayList<Object>();
|
||||
List<Recipient> recipientList = this.recipients;
|
||||
for (Recipient recipient : recipientList) {
|
||||
if (recipient.accept(message)) {
|
||||
channels.add(recipient.getChannel());
|
||||
}
|
||||
}
|
||||
return channels;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -40,7 +40,7 @@ public class MultiChannelRouterTests {
|
||||
|
||||
@Test
|
||||
public void routeWithChannelMapping() {
|
||||
AbstractChannelNameResolvingMessageRouter router = new AbstractChannelNameResolvingMessageRouter() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
@SuppressWarnings("unchecked")
|
||||
public List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return CollectionUtils.arrayToList(new String[] {"channel1", "channel2"});
|
||||
@@ -64,7 +64,7 @@ public class MultiChannelRouterTests {
|
||||
|
||||
@Test(expected = MessagingException.class)
|
||||
public void channelNameLookupFailure() {
|
||||
AbstractChannelNameResolvingMessageRouter router = new AbstractChannelNameResolvingMessageRouter() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
@SuppressWarnings("unchecked")
|
||||
public List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return CollectionUtils.arrayToList(new String[] {"noSuchChannel"} );
|
||||
@@ -78,7 +78,7 @@ public class MultiChannelRouterTests {
|
||||
|
||||
@Test(expected = MessagingException.class)
|
||||
public void channelMappingNotAvailable() {
|
||||
AbstractChannelNameResolvingMessageRouter router = new AbstractChannelNameResolvingMessageRouter() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
@SuppressWarnings("unchecked")
|
||||
public List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return CollectionUtils.arrayToList(new String[] {"noSuchChannel"});
|
||||
|
||||
@@ -45,9 +45,11 @@ public class RouterTests {
|
||||
@Test
|
||||
public void nullChannelIgnoredByDefault() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
public List<MessageChannel> determineTargetChannels(Message<?> message) {
|
||||
@Override
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return null;
|
||||
}
|
||||
|
||||
};
|
||||
Message<String> message = new GenericMessage<String>("test");
|
||||
router.handleMessage(message);
|
||||
@@ -56,7 +58,8 @@ public class RouterTests {
|
||||
@Test(expected = MessageDeliveryException.class)
|
||||
public void nullChannelThrowsExceptionWhenResolutionRequired() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
public List<MessageChannel> determineTargetChannels(Message<?> message) {
|
||||
@Override
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return null;
|
||||
}
|
||||
};
|
||||
@@ -68,8 +71,9 @@ public class RouterTests {
|
||||
@Test
|
||||
public void emptyChannelListIgnoredByDefault() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
public List<MessageChannel> determineTargetChannels(Message<?> message) {
|
||||
return Collections.emptyList();
|
||||
@Override
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return null;
|
||||
}
|
||||
};
|
||||
Message<String> message = new GenericMessage<String>("test");
|
||||
@@ -79,8 +83,9 @@ public class RouterTests {
|
||||
@Test(expected = MessageDeliveryException.class)
|
||||
public void emptyChannelListThrowsExceptionWhenResolutionRequired() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
public List<MessageChannel> determineTargetChannels(Message<?> message) {
|
||||
return Collections.emptyList();
|
||||
@Override
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return null;
|
||||
}
|
||||
};
|
||||
router.setResolutionRequired(true);
|
||||
@@ -90,8 +95,9 @@ public class RouterTests {
|
||||
|
||||
@Test
|
||||
public void nullChannelNameArrayIgnoredByDefault() {
|
||||
AbstractChannelNameResolvingMessageRouter router = new AbstractChannelNameResolvingMessageRouter() {
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
@Override
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return null;
|
||||
}
|
||||
};
|
||||
@@ -103,7 +109,7 @@ public class RouterTests {
|
||||
|
||||
@Test(expected = MessageDeliveryException.class)
|
||||
public void nullChannelNameArrayThrowsExceptionWhenResolutionRequired() {
|
||||
AbstractChannelNameResolvingMessageRouter router = new AbstractChannelNameResolvingMessageRouter() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return null;
|
||||
}
|
||||
@@ -118,7 +124,7 @@ public class RouterTests {
|
||||
|
||||
@Test
|
||||
public void emptyChannelNameArrayIgnoredByDefault() {
|
||||
AbstractChannelNameResolvingMessageRouter router = new AbstractChannelNameResolvingMessageRouter() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return new ArrayList<Object>();
|
||||
}
|
||||
@@ -131,7 +137,7 @@ public class RouterTests {
|
||||
|
||||
@Test(expected = MessageDeliveryException.class)
|
||||
public void emptyChannelNameArrayThrowsExceptionWhenResolutionRequired() {
|
||||
AbstractChannelNameResolvingMessageRouter router = new AbstractChannelNameResolvingMessageRouter() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
@SuppressWarnings("unchecked")
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return CollectionUtils.arrayToList(new String[] {});
|
||||
@@ -157,7 +163,7 @@ public class RouterTests {
|
||||
|
||||
@Test(expected = MessagingException.class)
|
||||
public void channelMappingIsRequiredWhenResolvingChannelNamesWithMultiChannelRouter() {
|
||||
AbstractChannelNameResolvingMessageRouter router = new AbstractChannelNameResolvingMessageRouter() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
@SuppressWarnings("unchecked")
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message){
|
||||
return CollectionUtils.arrayToList(new String[] { "notImportant" });
|
||||
@@ -185,7 +191,7 @@ public class RouterTests {
|
||||
|
||||
@Test
|
||||
public void beanFactoryWithMultiChannelRouter() {
|
||||
AbstractChannelNameResolvingMessageRouter router = new AbstractChannelNameResolvingMessageRouter() {
|
||||
AbstractMessageRouter router = new AbstractMessageRouter() {
|
||||
@SuppressWarnings("unchecked")
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return CollectionUtils.arrayToList(new String[] { "testChannel" });
|
||||
|
||||
@@ -26,6 +26,7 @@ import static org.mockito.Mockito.verify;
|
||||
|
||||
import java.util.Collection;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.mockito.Mockito;
|
||||
@@ -205,9 +206,10 @@ public class RouterParserTests {
|
||||
this.channel = channel;
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
protected Collection<MessageChannel> determineTargetChannels(Message<?> message) {
|
||||
return Collections.singletonList(this.channel);
|
||||
protected List<Object> getChannelIndicatorList(Message<?> message) {
|
||||
return Collections.singletonList((Object)this.channel);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -19,7 +19,7 @@ package org.springframework.integration.xml.router;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
import org.springframework.integration.router.AbstractChannelNameResolvingMessageRouter;
|
||||
import org.springframework.integration.router.AbstractMessageRouter;
|
||||
import org.springframework.integration.xml.DefaultXmlPayloadConverter;
|
||||
import org.springframework.integration.xml.XmlPayloadConverter;
|
||||
import org.springframework.xml.xpath.XPathExpression;
|
||||
@@ -31,7 +31,7 @@ import org.springframework.xml.xpath.XPathExpressionFactory;
|
||||
*
|
||||
* @author Jonas Partner
|
||||
*/
|
||||
public abstract class AbstractXPathRouter extends AbstractChannelNameResolvingMessageRouter {
|
||||
public abstract class AbstractXPathRouter extends AbstractMessageRouter {
|
||||
|
||||
private final XPathExpression xPathExpression;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user