diff --git a/spring-integration-core/src/main/java/org/springframework/integration/adapter/DefaultTargetAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/adapter/DefaultTargetAdapter.java index 2db89fb6ca..0f081f2c2f 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/adapter/DefaultTargetAdapter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/adapter/DefaultTargetAdapter.java @@ -43,6 +43,7 @@ public class DefaultTargetAdapter implements TargetAdapter { public DefaultTargetAdapter(Target target) { + Assert.notNull(target, "'target' must not be null"); this.target = target; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java index b7f27a26a8..5843764491 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java @@ -182,6 +182,9 @@ public class MessageBus implements ChannelRegistry, ApplicationContextAware, Lif ConsumerPolicy policy = adapter.getConsumerPolicy(); DispatcherTask dispatcherTask = new DispatcherTask((MessageDispatcher) adapter, policy); this.addDispatcherTask(dispatcherTask); + if (logger.isInfoEnabled()) { + logger.info("registered source adapter '" + name + "'"); + } } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/ChannelAdapterParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/ChannelAdapterParser.java index 8c6d776005..2832700001 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/ChannelAdapterParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/ChannelAdapterParser.java @@ -44,6 +44,10 @@ public class ChannelAdapterParser implements BeanDefinitionParser { private static final String METHOD_ATTRIBUTE = "method"; + private static final String CHANNEL_ATTRIBUTE = "channel"; + + private static final String PERIOD_ATTRIBUTE = "period"; + private final boolean isInbound; @@ -56,14 +60,25 @@ public class ChannelAdapterParser implements BeanDefinitionParser { public BeanDefinition parse(Element element, ParserContext parserContext) { String ref = element.getAttribute(REF_ATTRIBUTE); String method = element.getAttribute(METHOD_ATTRIBUTE); - if (!StringUtils.hasText(ref) || !StringUtils.hasText(method)) { - throw new MessagingConfigurationException("'ref' and 'method' are both required"); + String channel = element.getAttribute(CHANNEL_ATTRIBUTE); + if (!StringUtils.hasText(ref)) { + throw new MessagingConfigurationException("'ref' is required"); + } + if (!StringUtils.hasText(method)) { + throw new MessagingConfigurationException("'method' is required"); + } + if (!StringUtils.hasText(channel)) { + throw new MessagingConfigurationException("'channel' is required"); } RootBeanDefinition adapterDef = null; RootBeanDefinition invokerDef = null; if (this.isInbound) { adapterDef = new RootBeanDefinition(PollingSourceAdapter.class); invokerDef = new RootBeanDefinition(MethodInvokingSource.class); + String period = element.getAttribute(PERIOD_ATTRIBUTE); + if (StringUtils.hasText(period)) { + adapterDef.getPropertyValues().addPropertyValue("period", period); + } } else { adapterDef = new RootBeanDefinition(DefaultTargetAdapter.class); @@ -74,6 +89,7 @@ public class ChannelAdapterParser implements BeanDefinitionParser { String invokerBeanName = parserContext.getReaderContext().generateBeanName(invokerDef); parserContext.registerBeanComponent(new BeanComponentDefinition(invokerDef, invokerBeanName)); adapterDef.getConstructorArgumentValues().addGenericArgumentValue(new RuntimeBeanReference(invokerBeanName)); + adapterDef.getPropertyValues().addPropertyValue("channel", new RuntimeBeanReference(channel)); adapterDef.setSource(parserContext.extractSource(element)); String beanName = element.getAttribute(ID_ATTRIBUTE); if (!StringUtils.hasText(beanName)) { 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 a72f6687bc..10cfca1e11 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 @@ -29,8 +29,8 @@ public class IntegrationNamespaceHandler extends NamespaceHandlerSupport { registerBeanDefinitionParser("message-bus", new MessageBusParser()); registerBeanDefinitionParser("annotation-driven", new AnnotationDrivenParser()); registerBeanDefinitionParser("channel", new ChannelParser()); - registerBeanDefinitionParser("inbound-channel-adapter", new ChannelAdapterParser(true)); - registerBeanDefinitionParser("outbound-channel-adapter", new ChannelAdapterParser(false)); + registerBeanDefinitionParser("source-adapter", new ChannelAdapterParser(true)); + registerBeanDefinitionParser("target-adapter", new ChannelAdapterParser(false)); registerBeanDefinitionParser("endpoint", new EndpointParser()); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-1.0.xsd b/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-1.0.xsd index 42ffd58128..be79e8161c 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-1.0.xsd +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-1.0.xsd @@ -54,33 +54,36 @@ - + - Defines an inbound channel adapter. + Defines a source (inbound) channel adapter. + + - + - Defines an outbound channel adapter. + Defines a target (outbound) channel adapter. + diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandlerAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandlerAdapter.java index f71323cf30..9761f80e28 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandlerAdapter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/AbstractMessageHandlerAdapter.java @@ -51,6 +51,10 @@ public abstract class AbstractMessageHandlerAdapter implements MessageHandler private int order = Integer.MAX_VALUE; + private volatile boolean initialized; + + private Object lifecycleMonitor = new Object(); + public void setObject(T object) { Assert.notNull(object, "'object' must not be null"); @@ -88,7 +92,17 @@ public abstract class AbstractMessageHandlerAdapter implements MessageHandler } public final void afterPropertiesSet() { - this.invoker = new SimpleMethodInvoker(this.object, this.methodName); + this.validate(); + synchronized (this.lifecycleMonitor) { + this.invoker = new SimpleMethodInvoker(this.object, this.methodName); + this.initialized = true; + } + } + + public final boolean isInitialized() { + synchronized (this.lifecycleMonitor) { + return this.initialized; + } } public final Message handle(Message message) { @@ -102,6 +116,12 @@ public abstract class AbstractMessageHandlerAdapter implements MessageHandler return null; } + /** + * Subclasses may override this method to provide validation upon initialization. + */ + protected void validate() { + } + /** * Subclasses must implement this method. The invoker has been created for * the provided target object and method. May return an object of type diff --git a/spring-integration-core/src/main/java/org/springframework/integration/router/RouterMessageHandlerAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/router/RouterMessageHandlerAdapter.java index a2bc41a4c1..f545393fb7 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/router/RouterMessageHandlerAdapter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/router/RouterMessageHandlerAdapter.java @@ -28,6 +28,7 @@ import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.handler.AbstractMessageHandlerAdapter; import org.springframework.integration.message.Message; import org.springframework.integration.util.SimpleMethodInvoker; +import org.springframework.util.Assert; import org.springframework.util.StringUtils; /** @@ -50,7 +51,11 @@ public class RouterMessageHandlerAdapter extends AbstractMessageHandlerAdapter i public RouterMessageHandlerAdapter(Object object, Method method, Map attributes) { + Assert.notNull(object, "'object' must not be null"); + Assert.notNull(method, "'method' must not be null"); + Assert.notNull(attributes, "'attributes' must not be null"); this.setObject(object); + this.setMethodName(method.getName()); this.method = method; this.attributes = attributes; } @@ -61,6 +66,9 @@ public class RouterMessageHandlerAdapter extends AbstractMessageHandlerAdapter i @Override protected Object doHandle(Message message, SimpleMethodInvoker invoker) { + if (!this.isInitialized()) { + this.afterPropertiesSet(); + } if (method.getParameterTypes().length != 1) { throw new MessagingConfigurationException( "method must accept exactly one parameter"); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/adapter/AdapterBusTests.java b/spring-integration-core/src/test/java/org/springframework/integration/adapter/AdapterTests.java similarity index 66% rename from spring-integration-core/src/test/java/org/springframework/integration/adapter/AdapterBusTests.java rename to spring-integration-core/src/test/java/org/springframework/integration/adapter/AdapterTests.java index be60e7954e..904a1bcfa4 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/adapter/AdapterBusTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/adapter/AdapterTests.java @@ -29,11 +29,27 @@ import org.springframework.context.support.ClassPathXmlApplicationContext; /** * @author Mark Fisher */ -public class AdapterBusTests { +public class AdapterTests { @Test - public void testAdapters() throws IOException, InterruptedException { - AbstractApplicationContext context = new ClassPathXmlApplicationContext("adapterBusTests.xml", this.getClass()); + public void testAdaptersWithBeanDefinitions() throws IOException, InterruptedException { + AbstractApplicationContext context = new ClassPathXmlApplicationContext("adapterTests.xml", this.getClass()); + TestSink sink = (TestSink) context.getBean("sink"); + assertNull(sink.get()); + context.start(); + String result = null; + int attempts = 0; + while (result == null && attempts++ < 100) { + Thread.sleep(5); + result = sink.get(); + } + assertNotNull(result); + context.close(); + } + + @Test + public void testAdaptersWithNamespace() throws IOException, InterruptedException { + AbstractApplicationContext context = new ClassPathXmlApplicationContext("adapterTestsWithNamespace.xml", this.getClass()); TestSink sink = (TestSink) context.getBean("sink"); assertNull(sink.get()); context.start(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/adapter/adapterBusTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/adapter/adapterTests.xml similarity index 100% rename from spring-integration-core/src/test/java/org/springframework/integration/adapter/adapterBusTests.xml rename to spring-integration-core/src/test/java/org/springframework/integration/adapter/adapterTests.xml diff --git a/spring-integration-core/src/test/java/org/springframework/integration/adapter/adapterTestsWithNamespace.xml b/spring-integration-core/src/test/java/org/springframework/integration/adapter/adapterTestsWithNamespace.xml new file mode 100644 index 0000000000..52322d548d --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/adapter/adapterTestsWithNamespace.xml @@ -0,0 +1,26 @@ + + + + + + + + + + + + + + + + + + + +