diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/annotation/Poller.java b/org.springframework.integration/src/main/java/org/springframework/integration/annotation/Poller.java index dd52a8d4d8..e37705841a 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/annotation/Poller.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/annotation/Poller.java @@ -53,4 +53,6 @@ public @interface Poller { String transactionManager() default ""; + String[] adviceChain() default {}; + } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java index 8e5d698836..d07dcd8b77 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java @@ -19,9 +19,13 @@ package org.springframework.integration.config.annotation; import java.lang.annotation.Annotation; import java.lang.reflect.Method; import java.util.ArrayList; +import java.util.List; + +import org.aopalliance.aop.Advice; -import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.ListableBeanFactory; +import org.springframework.beans.factory.generic.GenericBeanFactoryAccessor; import org.springframework.core.annotation.AnnotationUtils; import org.springframework.integration.annotation.Poller; import org.springframework.integration.channel.ChannelRegistry; @@ -58,16 +62,16 @@ public abstract class AbstractMethodAnnotationPostProcessor 0) { + List adviceChain = new ArrayList(); + for (String adviceChainString : adviceChainArray) { + String[] adviceRefs = StringUtils.tokenizeToStringArray(adviceChainString, ","); + for (String adviceRef : adviceRefs) { + adviceChain.add(this.beanFactoryAccessor.getBean(adviceRef, Advice.class)); + } + } + pollingEndpoint.setAdviceChain(adviceChain); + } } endpoint = pollingEndpoint; } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java index ce0d0447d3..560d5ed4e9 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java @@ -19,7 +19,7 @@ package org.springframework.integration.config.annotation; import java.lang.annotation.Annotation; import java.lang.reflect.Method; -import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.core.annotation.AnnotationUtils; import org.springframework.integration.aggregator.AbstractMessageAggregator; import org.springframework.integration.aggregator.CompletionStrategyAdapter; @@ -39,7 +39,7 @@ import org.springframework.util.StringUtils; */ public class AggregatorAnnotationPostProcessor extends AbstractMethodAnnotationPostProcessor { - public AggregatorAnnotationPostProcessor(BeanFactory beanFactory) { + public AggregatorAnnotationPostProcessor(ListableBeanFactory beanFactory) { super(beanFactory); } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java index 4a9933e9dd..f348c5b6c6 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationPostProcessor.java @@ -31,7 +31,7 @@ import org.springframework.beans.factory.BeanFactoryAware; import org.springframework.beans.factory.BeanNameAware; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.config.BeanPostProcessor; -import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.core.annotation.AnnotationUtils; import org.springframework.integration.annotation.Aggregator; import org.springframework.integration.annotation.ChannelAdapter; @@ -60,7 +60,7 @@ public class MessagingAnnotationPostProcessor implements BeanPostProcessor, Bean private volatile MessageBus messageBus; - private volatile ConfigurableBeanFactory beanFactory; + private volatile ConfigurableListableBeanFactory beanFactory; private final Map, MethodAnnotationPostProcessor> postProcessors = @@ -68,9 +68,9 @@ public class MessagingAnnotationPostProcessor implements BeanPostProcessor, Bean public void setBeanFactory(BeanFactory beanFactory) { - Assert.isAssignable(ConfigurableBeanFactory.class, beanFactory.getClass(), - "a ConfigurableBeanFactory is required"); - this.beanFactory = (ConfigurableBeanFactory) beanFactory; + Assert.isAssignable(ConfigurableListableBeanFactory.class, beanFactory.getClass(), + "a ConfigurableListableBeanFactory is required"); + this.beanFactory = (ConfigurableListableBeanFactory) beanFactory; } public void afterPropertiesSet() { diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/RouterAnnotationPostProcessor.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/RouterAnnotationPostProcessor.java index dcb9b17bb0..15b0a76123 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/RouterAnnotationPostProcessor.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/RouterAnnotationPostProcessor.java @@ -18,7 +18,7 @@ package org.springframework.integration.config.annotation; import java.lang.reflect.Method; -import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.integration.annotation.Router; import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.message.MessageConsumer; @@ -34,7 +34,7 @@ import org.springframework.util.StringUtils; */ public class RouterAnnotationPostProcessor extends AbstractMethodAnnotationPostProcessor { - public RouterAnnotationPostProcessor(BeanFactory beanFactory) { + public RouterAnnotationPostProcessor(ListableBeanFactory beanFactory) { super(beanFactory); } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/ServiceActivatorAnnotationPostProcessor.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/ServiceActivatorAnnotationPostProcessor.java index 261fbbd1d2..f0fed611b8 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/ServiceActivatorAnnotationPostProcessor.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/ServiceActivatorAnnotationPostProcessor.java @@ -18,7 +18,7 @@ package org.springframework.integration.config.annotation; import java.lang.reflect.Method; -import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.integration.annotation.ServiceActivator; import org.springframework.integration.endpoint.ServiceActivatorEndpoint; import org.springframework.integration.message.MessageConsumer; @@ -31,7 +31,7 @@ import org.springframework.integration.message.MessageMappingMethodInvoker; */ public class ServiceActivatorAnnotationPostProcessor extends AbstractMethodAnnotationPostProcessor { - public ServiceActivatorAnnotationPostProcessor(BeanFactory beanFactory) { + public ServiceActivatorAnnotationPostProcessor(ListableBeanFactory beanFactory) { super(beanFactory); } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/SplitterAnnotationPostProcessor.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/SplitterAnnotationPostProcessor.java index f52b3cf539..fa717b298e 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/SplitterAnnotationPostProcessor.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/SplitterAnnotationPostProcessor.java @@ -18,7 +18,7 @@ package org.springframework.integration.config.annotation; import java.lang.reflect.Method; -import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.integration.annotation.Splitter; import org.springframework.integration.message.MessageConsumer; import org.springframework.integration.splitter.MethodInvokingSplitter; @@ -30,7 +30,7 @@ import org.springframework.integration.splitter.MethodInvokingSplitter; */ public class SplitterAnnotationPostProcessor extends AbstractMethodAnnotationPostProcessor { - public SplitterAnnotationPostProcessor(BeanFactory beanFactory) { + public SplitterAnnotationPostProcessor(ListableBeanFactory beanFactory) { super(beanFactory); } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/TransformerAnnotationPostProcessor.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/TransformerAnnotationPostProcessor.java index e8132425c8..ec9717525a 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/TransformerAnnotationPostProcessor.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/TransformerAnnotationPostProcessor.java @@ -18,7 +18,7 @@ package org.springframework.integration.config.annotation; import java.lang.reflect.Method; -import org.springframework.beans.factory.BeanFactory; +import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.integration.annotation.Transformer; import org.springframework.integration.message.MessageConsumer; import org.springframework.integration.transformer.MethodInvokingTransformer; @@ -31,7 +31,7 @@ import org.springframework.integration.transformer.TransformerEndpoint; */ public class TransformerAnnotationPostProcessor extends AbstractMethodAnnotationPostProcessor { - public TransformerAnnotationPostProcessor(BeanFactory beanFactory) { + public TransformerAnnotationPostProcessor(ListableBeanFactory beanFactory) { super(beanFactory); } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java index 185af50eff..72c85f7024 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java @@ -16,8 +16,14 @@ package org.springframework.integration.endpoint; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.ScheduledFuture; +import org.aopalliance.aop.Advice; + +import org.springframework.aop.framework.ProxyFactory; +import org.springframework.beans.factory.BeanClassLoaderAware; import org.springframework.beans.factory.InitializingBean; import org.springframework.context.Lifecycle; import org.springframework.core.task.TaskExecutor; @@ -32,11 +38,13 @@ import org.springframework.transaction.support.DefaultTransactionDefinition; import org.springframework.transaction.support.TransactionCallback; import org.springframework.transaction.support.TransactionTemplate; import org.springframework.util.Assert; +import org.springframework.util.ClassUtils; /** * @author Mark Fisher */ -public abstract class AbstractPollingEndpoint implements MessageEndpoint, TaskSchedulerAware, Lifecycle, InitializingBean { +public abstract class AbstractPollingEndpoint implements MessageEndpoint, + TaskSchedulerAware, Lifecycle, InitializingBean, BeanClassLoaderAware { public static final int MAX_MESSAGES_UNBOUNDED = -1; @@ -53,10 +61,16 @@ public abstract class AbstractPollingEndpoint implements MessageEndpoint, TaskSc private volatile TransactionTemplate transactionTemplate; + private final List adviceChain = new CopyOnWriteArrayList(); + + private volatile ClassLoader classLoader = ClassUtils.getDefaultClassLoader(); + private volatile TaskScheduler taskScheduler; private volatile ScheduledFuture runningTask; + private volatile Runnable poller; + private volatile boolean initialized; private final Object lifecycleMonitor = new Object(); @@ -100,6 +114,16 @@ public abstract class AbstractPollingEndpoint implements MessageEndpoint, TaskSc this.transactionDefinition = transactionDefinition; } + public void setBeanClassLoader(ClassLoader classLoader) { + Assert.notNull(classLoader, "ClassLoader must not be null"); + this.classLoader = classLoader; + } + + public void setAdviceChain(List adviceChain) { + this.adviceChain.clear(); + this.adviceChain.addAll(adviceChain); + } + private TransactionTemplate getTransactionTemplate() { if (!this.initialized) { this.afterPropertiesSet(); @@ -122,10 +146,22 @@ public abstract class AbstractPollingEndpoint implements MessageEndpoint, TaskSc this.transactionTemplate = new TransactionTemplate( this.transactionManager, this.transactionDefinition); } + this.poller = this.createPoller(); this.initialized = true; } } + private Runnable createPoller() { + if (this.adviceChain.isEmpty()) { + return new Poller(); + } + ProxyFactory proxyFactory = new ProxyFactory(new Poller()); + for (Advice advice : this.adviceChain) { + proxyFactory.addAdvice(advice); + } + return (Runnable) proxyFactory.getProxy(this.classLoader); + } + // Lifecycle implementation @@ -145,7 +181,7 @@ public abstract class AbstractPollingEndpoint implements MessageEndpoint, TaskSc } Assert.state(this.taskScheduler != null, "unable to start polling, no taskScheduler available"); - this.runningTask = this.taskScheduler.schedule(new Poller(), this.trigger); + this.runningTask = this.taskScheduler.schedule(this.poller, this.trigger); } } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java index ed5791525d..7d9ce96246 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java @@ -21,7 +21,6 @@ import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.MessageChannelTemplate; import org.springframework.integration.message.Message; import org.springframework.integration.message.MessageSource; -import org.springframework.integration.message.MethodInvokingSource; import org.springframework.util.Assert; /** @@ -62,8 +61,8 @@ public class SourcePollingChannelAdapter extends AbstractPollingEndpoint impleme Assert.notNull(this.source, "source must not be null"); Assert.notNull(this.outputChannel, "outputChannel must not be null"); super.afterPropertiesSet(); - if (this.maxMessagesPerPoll < 0 && source instanceof MethodInvokingSource) { - // the default is 1 since a MethodInvokingSource might return + if (this.maxMessagesPerPoll < 0) { + // the default is 1 since a source might return // a non-null value every time it is invoked this.setMaxMessagesPerPoll(1); }