From 875b0cc02acacae8fee0941ad31ae875655fd1d6 Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Thu, 9 Oct 2008 17:40:17 +0000 Subject: [PATCH] Added an 'adviceChain' property to AbstractPollingEndpoint so that any AOP Advice instance(s) may be applied to the Poller (a Runnable within the AbstractPollingEndpoint). Also, added support for configuring this 'adviceChain' via the @Poller annotation (INT-407). --- .../integration/annotation/Poller.java | 2 + ...AbstractMethodAnnotationPostProcessor.java | 34 +++++++++++----- .../AggregatorAnnotationPostProcessor.java | 4 +- .../MessagingAnnotationPostProcessor.java | 10 ++--- .../RouterAnnotationPostProcessor.java | 4 +- ...rviceActivatorAnnotationPostProcessor.java | 4 +- .../SplitterAnnotationPostProcessor.java | 4 +- .../TransformerAnnotationPostProcessor.java | 4 +- .../endpoint/AbstractPollingEndpoint.java | 40 ++++++++++++++++++- .../endpoint/SourcePollingChannelAdapter.java | 5 +-- 10 files changed, 82 insertions(+), 29 deletions(-) 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); }