From 17294c44ec11ef477ac5d4d083fc955084e43e99 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Thu, 15 May 2008 05:54:12 +0000 Subject: [PATCH] Adding support for defining a ChannelFactory on the MessageBus in the form: The channel factory will be used for auto-created channels, as well as for channels created using the syntax. For specifying a QueueChannel (with capacity), the newly added element must be used. --- .../channel/config/DefaultChannelParser.java | 41 +++++++++++++++++++ .../channel/config/QueueChannelParser.java | 2 +- .../factory/DefaultChannelFactoryBean.java | 36 ++++++++++------ .../config/AbstractTargetEndpointParser.java | 2 +- .../config/IntegrationNamespaceHandler.java | 5 ++- .../integration/config/MessageBusParser.java | 35 ++++++++++++---- .../config/spring-integration-core-1.0.xsd | 21 +++++++++- .../channel/config/channelParserTests.xml | 2 +- .../channel/factory/TestChannelFactory.java | 17 +++++--- .../config/MessageBusParserTests.java | 13 ++++++ .../config/aggregatorParserTests.xml | 2 + .../endpointWithHandlerChainElement.xml | 2 +- .../config/endpointWithSelectors.xml | 2 +- .../config/handlerAdapterEndpointTests.xml | 2 +- .../config/messageBusWithChannelFactory.xml | 20 +++++++++ .../config/messageBusWithDefaults.xml | 2 +- .../config/simpleEndpointTests.xml | 2 +- 17 files changed, 170 insertions(+), 36 deletions(-) create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/channel/config/DefaultChannelParser.java create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithChannelFactory.xml diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/config/DefaultChannelParser.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/config/DefaultChannelParser.java new file mode 100644 index 0000000000..57cc1e02f4 --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/config/DefaultChannelParser.java @@ -0,0 +1,41 @@ +/* + * Copyright 2002-2008 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.channel.config; + +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.integration.channel.DispatcherPolicy; +import org.springframework.integration.channel.factory.DefaultChannelFactoryBean; +import org.w3c.dom.Element; + +/** + * Parser for the <channel> element. + * + * @author Marius Bogoevici + */ +public class DefaultChannelParser extends AbstractChannelParser { + + @Override + protected Class getBeanClass(Element element) { + return DefaultChannelFactoryBean.class; + } + + @Override + protected void configureConstructorArgs(BeanDefinitionBuilder builder, Element element, DispatcherPolicy dispatcherPolicy) { + builder.addConstructorArgValue(dispatcherPolicy); + } + +} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/config/QueueChannelParser.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/config/QueueChannelParser.java index 0ed43f5ffc..32d0fc6a51 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/config/QueueChannelParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/config/QueueChannelParser.java @@ -24,7 +24,7 @@ import org.springframework.integration.channel.QueueChannel; import org.springframework.util.StringUtils; /** - * Parser for the <channel> element. + * Parser for the <queue-channel> element. * * @author Mark Fisher */ diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/factory/DefaultChannelFactoryBean.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/factory/DefaultChannelFactoryBean.java index 77fdbe7755..4b71d2d935 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/factory/DefaultChannelFactoryBean.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/factory/DefaultChannelFactoryBean.java @@ -17,48 +17,58 @@ package org.springframework.integration.channel.factory; import java.util.List; +import java.util.Map; import org.springframework.beans.factory.FactoryBean; -import org.springframework.beans.factory.InitializingBean; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; import org.springframework.integration.bus.MessageBus; -import org.springframework.integration.bus.MessageBusAware; import org.springframework.integration.channel.ChannelInterceptor; import org.springframework.integration.channel.DispatcherPolicy; import org.springframework.integration.channel.MessageChannel; +import org.springframework.util.Assert; /** * Creates a channel by delegating to the current message bus-configured - * ChannelFactory. + * ChannelFactory. Tries to retrieve the {@link ChannelFactory} from the + * single {@link MessageBus} defined in the {@link ApplicationContext}. + * As a {@link FactoryBean}, this class is solely intended to be used within + * an ApplicationContext. * @author Marius Bogoevici */ -public class DefaultChannelFactoryBean implements FactoryBean, MessageBusAware, InitializingBean { +public class DefaultChannelFactoryBean implements ApplicationContextAware, FactoryBean{ private volatile ChannelFactory channelFactory; private volatile List interceptors; private volatile DispatcherPolicy dispatcherPolicy; + + private volatile boolean publisherSubscriber; - public void setMessageBus(MessageBus messageBus) { - this.channelFactory = messageBus.getChannelFactory(); + + public DefaultChannelFactoryBean(DispatcherPolicy dispatcherPolicy) { + this.dispatcherPolicy = dispatcherPolicy; } - public void afterPropertiesSet() throws Exception { - if (null == this.channelFactory) { + + public void setApplicationContext(ApplicationContext applicationContext){ + Map map = applicationContext.getBeansOfType(MessageBus.class); + Assert.state(map.size() <= 1, "There is more than one MessageBus in the ApplicationContext"); + if (map.isEmpty()) { this.channelFactory = new QueueChannelFactory(); } - + else { + this.channelFactory = ((MessageBus)map.values().iterator().next()).getChannelFactory(); + } } public void setInterceptors(List interceptors) { this.interceptors = interceptors; } - public void setDispatcherPolicy(DispatcherPolicy dispatcherPolicy) { - this.dispatcherPolicy = dispatcherPolicy; - } - public Object getObject() throws Exception { + Assert.notNull(channelFactory, "ChannelFactory not set on this instance. Is this used within an ApplicationContext?"); return channelFactory.getChannel(dispatcherPolicy, interceptors); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/AbstractTargetEndpointParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/AbstractTargetEndpointParser.java index ff77dc46a6..a1fcc76f66 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/AbstractTargetEndpointParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/AbstractTargetEndpointParser.java @@ -111,7 +111,7 @@ public abstract class AbstractTargetEndpointParser extends AbstractSingleBeanDef } if (StringUtils.hasText(inputChannel)) { RootBeanDefinition subscriptionDef = new RootBeanDefinition(Subscription.class); - subscriptionDef.getConstructorArgumentValues().addGenericArgumentValue(inputChannel); + subscriptionDef.getConstructorArgumentValues().addGenericArgumentValue(new RuntimeBeanReference(inputChannel)); if (schedule != null) { subscriptionDef.getConstructorArgumentValues().addGenericArgumentValue(schedule); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationNamespaceHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationNamespaceHandler.java index 93ef1741e3..1f559efc5d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationNamespaceHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationNamespaceHandler.java @@ -24,10 +24,10 @@ import java.util.Properties; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; - import org.springframework.beans.factory.xml.BeanDefinitionParser; import org.springframework.beans.factory.xml.NamespaceHandlerSupport; import org.springframework.core.io.support.PropertiesLoaderUtils; +import org.springframework.integration.channel.config.DefaultChannelParser; import org.springframework.integration.channel.config.DirectChannelParser; import org.springframework.integration.channel.config.PriorityChannelParser; import org.springframework.integration.channel.config.QueueChannelParser; @@ -51,7 +51,8 @@ public class IntegrationNamespaceHandler extends NamespaceHandlerSupport { public void init() { registerBeanDefinitionParser("message-bus", new MessageBusParser()); registerBeanDefinitionParser("annotation-driven", new AnnotationDrivenParser()); - registerBeanDefinitionParser("channel", new QueueChannelParser()); + registerBeanDefinitionParser("channel", new DefaultChannelParser()); + registerBeanDefinitionParser("queue-channel", new QueueChannelParser()); registerBeanDefinitionParser("direct-channel", new DirectChannelParser()); registerBeanDefinitionParser("priority-channel", new PriorityChannelParser()); registerBeanDefinitionParser("rendezvous-channel", new RendezvousChannelParser()); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/MessageBusParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/MessageBusParser.java index fcd098b02e..ee5edc88a1 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/MessageBusParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/MessageBusParser.java @@ -42,6 +42,8 @@ import org.w3c.dom.NodeList; */ public class MessageBusParser extends AbstractSimpleBeanDefinitionParser { + private static final String REFERENCE_ATTRIBUTE = "ref"; + public static final String MESSAGE_BUS_BEAN_NAME = "internal.MessageBus"; public static final String MESSAGE_BUS_AWARE_POST_PROCESSOR_BEAN_NAME = "internal.MessageBusAwareBeanPostProcessor"; @@ -53,6 +55,10 @@ public class MessageBusParser extends AbstractSimpleBeanDefinitionParser { private static final String DEFAULT_CONCURRENCY_ELEMENT = "default-concurrency"; private static final String DEFAULT_CONCURRENCY_PROPERTY = "defaultConcurrencyPolicy"; + + private static final String CHANNEL_FACTORY_ELEMENT = "channel-factory"; + + private static final String CHANNEL_FACTORY_PROPERTY = "channelFactory"; @Override protected String resolveId(Element element, AbstractBeanDefinition definition, ParserContext parserContext) @@ -81,23 +87,38 @@ public class MessageBusParser extends AbstractSimpleBeanDefinitionParser { beanDefinition.addPropertyReference(Conventions.attributeNameToPropertyName( ERROR_CHANNEL_ATTRIBUTE), errorChannelRef); } - this.registerDefaultConcurrencyIfAvailable(beanDefinition, element); + this.processAdditionalChildElements(beanDefinition, element); } - private void registerDefaultConcurrencyIfAvailable(BeanDefinitionBuilder beanDefinition, Element element) { + private void processAdditionalChildElements(BeanDefinitionBuilder beanDefinition, Element element) { NodeList childNodes = element.getChildNodes(); for (int i = 0; i < childNodes.getLength(); i++) { Node child = childNodes.item(i); if (child.getNodeType() == Node.ELEMENT_NODE) { - String localName = child.getLocalName(); - if (DEFAULT_CONCURRENCY_ELEMENT.equals(localName)) { - ConcurrencyPolicy policy = IntegrationNamespaceUtils.parseConcurrencyPolicy((Element) child); - beanDefinition.addPropertyValue(DEFAULT_CONCURRENCY_PROPERTY, policy); - } + processIfConcurrencyElement(beanDefinition, child); + processIfChannelFactoryElement(beanDefinition, child); } } } + + private void processIfConcurrencyElement(BeanDefinitionBuilder beanDefinition, Node node) { + String localName = node.getLocalName(); + if (DEFAULT_CONCURRENCY_ELEMENT.equals(localName)) { + ConcurrencyPolicy policy = IntegrationNamespaceUtils.parseConcurrencyPolicy((Element) node); + beanDefinition.addPropertyValue(DEFAULT_CONCURRENCY_PROPERTY, policy); + } + + } + + private void processIfChannelFactoryElement(BeanDefinitionBuilder beanDefinition, Node node) { + String localName = node.getLocalName(); + if (CHANNEL_FACTORY_ELEMENT.equals(localName)) { + beanDefinition.addPropertyReference(CHANNEL_FACTORY_PROPERTY, ((Element) node) + .getAttribute(REFERENCE_ATTRIBUTE)); + } + } + @Override protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { super.doParse(element, parserContext, builder); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-core-1.0.xsd b/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-core-1.0.xsd index 78b35e17fc..90891a2710 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-core-1.0.xsd +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-core-1.0.xsd @@ -26,6 +26,7 @@ + @@ -43,8 +44,17 @@ + + + + + Defines a generic channel type. The actual channel type + will be determined by the channel factory set on the message bus. + + + - + Defines a channel that buffers messages in a queue. @@ -256,6 +266,15 @@ + + + + + Defines a channel factory. + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/config/channelParserTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/channel/config/channelParserTests.xml index f9c8753992..d1cbbf5b64 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/config/channelParserTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/config/channelParserTests.xml @@ -7,7 +7,7 @@ http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration-core-1.0.xsd"> - + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/factory/TestChannelFactory.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/factory/TestChannelFactory.java index bc8f83e3fe..86472416a3 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/factory/TestChannelFactory.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/factory/TestChannelFactory.java @@ -25,9 +25,12 @@ import java.util.List; import org.junit.Before; import org.junit.Test; -import org.omg.CORBA.REBIND; import org.springframework.beans.DirectFieldAccessor; -import org.springframework.beans.factory.FactoryBean; +import org.springframework.beans.factory.config.BeanDefinition; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.context.ApplicationContext; +import org.springframework.context.support.GenericApplicationContext; +import org.springframework.context.support.StaticApplicationContext; import org.springframework.integration.bus.MessageBus; import org.springframework.integration.channel.AbstractMessageChannel; import org.springframework.integration.channel.ChannelInterceptor; @@ -35,6 +38,7 @@ import org.springframework.integration.channel.DispatcherPolicy; import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.channel.RendezvousChannel; +import org.springframework.integration.config.MessageBusParser; import org.springframework.integration.dispatcher.SynchronousChannel; import org.springframework.integration.message.Message; import org.springframework.integration.message.selector.MessageSelector; @@ -97,10 +101,13 @@ public class TestChannelFactory { MessageBus messageBus = new MessageBus(); ChannelFactory channelFactory = new StubChannelFactory(); messageBus.setChannelFactory(channelFactory); - DefaultChannelFactoryBean channelFactoryBean = new DefaultChannelFactoryBean(); - channelFactoryBean.setDispatcherPolicy(dispatcherPolicy); + StaticApplicationContext applicationContext = new StaticApplicationContext(); + BeanDefinitionBuilder messageBusDefinitionBuilder = BeanDefinitionBuilder.rootBeanDefinition(MessageBus.class); + messageBusDefinitionBuilder.getBeanDefinition().getPropertyValues().addPropertyValue("channelFactory", channelFactory); + applicationContext.registerBeanDefinition("messageBus", messageBusDefinitionBuilder.getBeanDefinition()); + DefaultChannelFactoryBean channelFactoryBean = new DefaultChannelFactoryBean(dispatcherPolicy); + channelFactoryBean.setApplicationContext(applicationContext); channelFactoryBean.setInterceptors(interceptors); - channelFactoryBean.setMessageBus(messageBus); StubChannel channel = (StubChannel)channelFactoryBean.getObject(); assertTrue(dispatcherPolicy == channel.getDispatcherPolicy()); assertInterceptors(channel); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/MessageBusParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/MessageBusParserTests.java index f510f46e1c..a47160aad7 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/MessageBusParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/MessageBusParserTests.java @@ -29,6 +29,8 @@ import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.integration.ConfigurationException; import org.springframework.integration.bus.MessageBus; import org.springframework.integration.bus.TestMessageBusAwareImpl; +import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.dispatcher.SynchronousChannel; import org.springframework.integration.endpoint.TargetEndpoint; import org.springframework.integration.handler.TestHandlers; import org.springframework.integration.scheduling.Subscription; @@ -143,5 +145,16 @@ public class MessageBusParserTests { TestMessageBusAwareImpl messageBusAware = (TestMessageBusAwareImpl) context.getBean("messageBusAwareBean"); assertTrue(messageBusAware.getMessageBus() == context.getBean(MessageBusParser.MESSAGE_BUS_BEAN_NAME)); } + + @Test + public void testMessageBusWithChannelFactory() { + ApplicationContext context = new ClassPathXmlApplicationContext("messageBusWithChannelFactory.xml", + this.getClass()); + MessageBus bus = (MessageBus) context.getBean(MessageBusParser.MESSAGE_BUS_BEAN_NAME); + assertTrue (context.getBean("defaultTypeChannel") instanceof SynchronousChannel); + assertTrue (context.getBean("specifiedTypeChannel") instanceof QueueChannel); + + + } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml index e88f817983..c29aeb04e1 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml @@ -7,6 +7,8 @@ http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration-core-1.0.xsd"> + + - + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/endpointWithSelectors.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/endpointWithSelectors.xml index 82e6ceaa27..1a94e5658b 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/endpointWithSelectors.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/endpointWithSelectors.xml @@ -9,7 +9,7 @@ - + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/handlerAdapterEndpointTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/handlerAdapterEndpointTests.xml index 67bacde949..8c7b102982 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/handlerAdapterEndpointTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/handlerAdapterEndpointTests.xml @@ -9,7 +9,7 @@ - + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithChannelFactory.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithChannelFactory.xml new file mode 100644 index 0000000000..e02f9dda0b --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithChannelFactory.xml @@ -0,0 +1,20 @@ + + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithDefaults.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithDefaults.xml index 794bf871ca..35558fd505 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithDefaults.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithDefaults.xml @@ -8,5 +8,5 @@ http://www.springframework.org/schema/integration/spring-integration-core-1.0.xsd"> - + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/simpleEndpointTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/simpleEndpointTests.xml index b05bf1361c..a23a020750 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/simpleEndpointTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/simpleEndpointTests.xml @@ -9,7 +9,7 @@ - +