From a4d542a179563ff7c89e590ba6e308ca21272c08 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Wed, 2 Jan 2008 20:20:23 +0000 Subject: [PATCH] Added an implementation of MessageHandler for routing messages. Supports either ChannelResolver or ChannelNameResolver strategy. --- .../router/ChannelNameResolver.java | 6 +- .../integration/router/ChannelResolver.java | 5 +- .../router/RoutingMessageHandler.java | 102 ++++++++++++++ .../router/RoutingMessageHandlerTests.java | 127 ++++++++++++++++++ 4 files changed, 236 insertions(+), 4 deletions(-) create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/router/RoutingMessageHandler.java create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/router/RoutingMessageHandlerTests.java diff --git a/spring-integration-core/src/main/java/org/springframework/integration/router/ChannelNameResolver.java b/spring-integration-core/src/main/java/org/springframework/integration/router/ChannelNameResolver.java index 2cf076b421..09a9f00505 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/router/ChannelNameResolver.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/router/ChannelNameResolver.java @@ -16,13 +16,15 @@ package org.springframework.integration.router; +import org.springframework.integration.message.Message; + /** * Strategy interface for content-based routing to a channel name. * * @author Mark Fisher */ -public interface ChannelNameResolver { +public interface ChannelNameResolver { - String resolve(T t); + String resolve(Message message); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/router/ChannelResolver.java b/spring-integration-core/src/main/java/org/springframework/integration/router/ChannelResolver.java index 47d5938bb1..89881c2809 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/router/ChannelResolver.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/router/ChannelResolver.java @@ -17,14 +17,15 @@ package org.springframework.integration.router; import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.message.Message; /** * Strategy interface for content-based routing to a channel instance. * * @author Mark Fisher */ -public interface ChannelResolver { +public interface ChannelResolver { - MessageChannel resolve(T t); + MessageChannel resolve(Message message); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/router/RoutingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/router/RoutingMessageHandler.java new file mode 100644 index 0000000000..6ca3591a21 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/router/RoutingMessageHandler.java @@ -0,0 +1,102 @@ +/* + * Copyright 2002-2007 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 org.springframework.beans.factory.InitializingBean; +import org.springframework.integration.MessageHandlingException; +import org.springframework.integration.MessagingConfigurationException; +import org.springframework.integration.channel.ChannelRegistry; +import org.springframework.integration.channel.ChannelRegistryAware; +import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.handler.MessageHandler; +import org.springframework.integration.message.Message; + +/** + * @author Mark Fisher + */ +public class RoutingMessageHandler implements MessageHandler, ChannelRegistryAware, InitializingBean { + + private ChannelResolver channelResolver; + + private ChannelNameResolver channelNameResolver; + + private ChannelRegistry channelRegistry; + + private long timeout = -1; + + + /** + * Set the timeout for sending a message to the resolved channel. + * Default is no timeout, meaning the send will block indefinitely. + */ + public void setTimeout(long timeout) { + this.timeout = timeout; + } + + public void setChannelResolver(ChannelResolver channelResolver) { + this.channelResolver = channelResolver; + } + + public void setChannelNameResolver(ChannelNameResolver channelNameResolver) { + this.channelNameResolver = channelNameResolver; + } + + public void setChannelRegistry(ChannelRegistry channelRegistry) { + this.channelRegistry = channelRegistry; + } + + public void afterPropertiesSet() { + if(!(this.channelResolver != null ^ this.channelNameResolver != null)) { + throw new MessagingConfigurationException("exactly one of 'channelResolver' or 'channelNameResolver' must be provided"); + } + if (this.channelNameResolver != null && this.channelRegistry == null) { + throw new MessagingConfigurationException("'channelRegistry' is required when resolving by channel name"); + } + } + + public Message handle(Message message) { + MessageChannel channel = this.resolveChannel(message); + if (channel == null) { + throw new MessageHandlingException("failed to resolve channel"); + } + boolean sent = false; + if (timeout < 0) { + sent = channel.send(message); + } + else { + sent = channel.send(message, timeout); + } + if (!sent) { + throw new MessageHandlingException( + "failed to send message to channel '" + channel.getName() + "'"); + } + return null; + } + + private MessageChannel resolveChannel(Message message) { + if (this.channelResolver != null) { + return this.channelResolver.resolve(message); + } + else if (this.channelNameResolver != null && this.channelRegistry != null) { + String channelName = this.channelNameResolver.resolve(message); + return this.channelRegistry.lookupChannel(channelName); + } + throw new MessagingConfigurationException("router configuration requires either " + + "a 'channelResolver' or both 'channelNameResolver' and 'channelRegistry'"); + } + +} diff --git a/spring-integration-core/src/test/java/org/springframework/integration/router/RoutingMessageHandlerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/router/RoutingMessageHandlerTests.java new file mode 100644 index 0000000000..b71c6b4f90 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/router/RoutingMessageHandlerTests.java @@ -0,0 +1,127 @@ +/* + * Copyright 2002-2007 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 static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; + +import org.junit.Test; + +import org.springframework.integration.MessageHandlingException; +import org.springframework.integration.MessagingConfigurationException; +import org.springframework.integration.channel.ChannelRegistry; +import org.springframework.integration.channel.DefaultChannelRegistry; +import org.springframework.integration.channel.MessageChannel; +import org.springframework.integration.channel.PointToPointChannel; +import org.springframework.integration.message.Message; +import org.springframework.integration.message.StringMessage; + +/** + * @author Mark Fisher + */ +public class RoutingMessageHandlerTests { + + @Test + public void testRoutingWithChannelResolver() { + final PointToPointChannel channel = new PointToPointChannel(); + ChannelResolver channelResolver = new ChannelResolver() { + public MessageChannel resolve(Message message) { + return channel; + } + }; + RoutingMessageHandler router = new RoutingMessageHandler(); + router.setChannelResolver(channelResolver); + router.afterPropertiesSet(); + Message message = new StringMessage("123", "test"); + router.handle(message); + Message result = channel.receive(25); + assertNotNull(result); + assertEquals("test", result.getPayload()); + } + + @Test + public void testRoutingWithChannelNameResolver() { + ChannelNameResolver channelNameResolver = new ChannelNameResolver() { + public String resolve(Message message) { + return "testChannel"; + } + }; + PointToPointChannel channel = new PointToPointChannel(); + ChannelRegistry channelRegistry = new DefaultChannelRegistry(); + channelRegistry.registerChannel("testChannel", channel); + RoutingMessageHandler router = new RoutingMessageHandler(); + router.setChannelNameResolver(channelNameResolver); + router.setChannelRegistry(channelRegistry); + router.afterPropertiesSet(); + Message message = new StringMessage("123", "test"); + router.handle(message); + Message result = channel.receive(25); + assertNotNull(result); + assertEquals("test", result.getPayload()); + } + + @Test(expected=MessagingConfigurationException.class) + public void testConfiguringBothChannelResolverAndChannelNameResolverIsNotAllowed() { + ChannelResolver channelResolver = new ChannelResolver() { + public MessageChannel resolve(Message message) { + return new PointToPointChannel(); + } + }; + ChannelNameResolver channelNameResolver = new ChannelNameResolver() { + public String resolve(Message message) { + return ""; + } + }; + RoutingMessageHandler router = new RoutingMessageHandler(); + router.setChannelResolver(channelResolver); + router.setChannelNameResolver(channelNameResolver); + router.afterPropertiesSet(); + } + + @Test(expected=MessageHandlingException.class) + public void testChannelResolutionFailure() { + ChannelResolver channelResolver = new ChannelResolver() { + public MessageChannel resolve(Message message) { + return null; + } + }; + RoutingMessageHandler router = new RoutingMessageHandler(); + router.setChannelResolver(channelResolver); + router.afterPropertiesSet(); + Message message = new StringMessage("123", "test"); + router.handle(message); + } + + @Test(expected=MessageHandlingException.class) + public void testChannelNameResolutionFailure() { + ChannelNameResolver channelNameResolver = new ChannelNameResolver() { + public String resolve(Message message) { + return "noSuchChannel"; + } + }; + PointToPointChannel channel = new PointToPointChannel(); + ChannelRegistry channelRegistry = new DefaultChannelRegistry(); + channelRegistry.registerChannel("testChannel", channel); + RoutingMessageHandler router = new RoutingMessageHandler(); + router.setChannelNameResolver(channelNameResolver); + router.setChannelRegistry(channelRegistry); + router.afterPropertiesSet(); + Message message = new StringMessage("123", "test"); + router.handle(message); + } + +}