From 3664d02e8d48cd816af17e8867f328cf20a59760 Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Thu, 4 Jun 2015 06:56:13 +0100 Subject: [PATCH] Lazily initialize the MessageBus in BeanPostProcessor Otherwise bad things happen when the MessageBus is not decorated by other BPPs (especially the @EnableIntegration one). --- roadmap.md | 6 ++- .../MessageBusAdapterConfiguration.java | 47 +++++++++++++++++-- 2 files changed, 48 insertions(+), 5 deletions(-) diff --git a/roadmap.md b/roadmap.md index aa8493d29..6535746df 100644 --- a/roadmap.md +++ b/roadmap.md @@ -77,7 +77,11 @@ To be deployable as an XD module in a "traditional" way you need `/config/*.prop - [x] Support for pubsub as "primary" input/output (in addition to the existing queue semantics) -- [ ] Endpoint "/messages" for module configuration metadata ("/bus" is taken by Spring Cloud) +- [x] Endpoint "/channels" for module configuration metadata ("/bus" is taken by Spring Cloud) + +- [ ] Discover channel names through Spring Cloud service discovery (via "/channels" endpoint on remote components) + +- [ ] Optional validation of application as XD module - [ ] Support for more than one `MessageBus` (e.g. local and redis) in the same app diff --git a/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusAdapterConfiguration.java b/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusAdapterConfiguration.java index 6c30a0852..cdd69a87c 100644 --- a/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusAdapterConfiguration.java +++ b/spring-bus-core/src/main/java/org/springframework/bus/runner/config/MessageBusAdapterConfiguration.java @@ -15,11 +15,16 @@ */ package org.springframework.bus.runner.config; +import java.util.Arrays; import java.util.Collection; import java.util.LinkedHashSet; import java.util.Properties; import java.util.Set; +import org.aopalliance.intercept.MethodInterceptor; +import org.aopalliance.intercept.MethodInvocation; +import org.springframework.aop.framework.ProxyFactory; +import org.springframework.aop.target.LazyInitTargetSource; import org.springframework.beans.factory.BeanFactoryUtils; import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.beans.factory.annotation.Autowired; @@ -31,6 +36,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.ImportResource; import org.springframework.messaging.MessageChannel; +import org.springframework.util.Assert; import org.springframework.xd.dirt.integration.bus.MessageBus; import org.springframework.xd.dirt.integration.bus.MessageBusAwareRouterBeanPostProcessor; @@ -149,11 +155,44 @@ public class MessageBusAdapterConfiguration { @Configuration protected static class MessageBusAwareRouterConfiguration { + @Autowired + private ListableBeanFactory beanFactory; + @Bean - public MessageBusAwareRouterBeanPostProcessor messageBusAwareRouterBeanPostProcessor( - MessageBus messageBus) { - return new MessageBusAwareRouterBeanPostProcessor(messageBus, - new Properties()); + public MessageBusAwareRouterBeanPostProcessor messageBusAwareRouterBeanPostProcessor() { + + return new MessageBusAwareRouterBeanPostProcessor( + createLazyProxy(beanFactory, MessageBus.class), new Properties()); + } + + private T createLazyProxy(ListableBeanFactory beanFactory, Class type) { + ProxyFactory factory = new ProxyFactory(); + LazyInitTargetSource source = new LazyInitTargetSource(); + source.setTargetClass(type); + source.setTargetBeanName(getBeanNameFor(beanFactory, MessageBus.class)); + source.setBeanFactory(beanFactory); + factory.setTargetSource(source); + factory.addAdvice(new PassthruAdvice()); + factory.setInterfaces(new Class[] { type }); + @SuppressWarnings("unchecked") + T proxy = (T) factory.getProxy(); + return proxy; + } + + private String getBeanNameFor(ListableBeanFactory beanFactory, Class type) { + String[] names = BeanFactoryUtils.beanNamesForTypeIncludingAncestors(beanFactory, type, false, false); + Assert.state(names.length==1, "No unique MessageBus (found " + names.length + ": " + + Arrays.asList(names) + ")"); + return names[0]; + } + + private class PassthruAdvice implements MethodInterceptor { + + @Override + public Object invoke(MethodInvocation invocation) throws Throwable { + return invocation.proceed(); + } + } }