From 34cc492942f41c1998be6db32a5f752bf57246e8 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 28 Apr 2014 15:56:01 +0300 Subject: [PATCH] INT-3381: Messaging Annotations on `@Bean` Method JIRA: https://jira.spring.io/browse/INT-3381 * Add support to mark `@Bean` methods with `Messaging Annotations` * Addition attributes for `Messaging Annotations` * Some refactoring and fixing * Fix `AbstractMappingMessageRouter#onInit()` propagation * Fix `GemfireGroupStoreTests` (https://build.spring.io/browse/INT-MASTERSPRING40-JOB1-246/test/case/135235772) There is still some side effect: we need define `MessageChannel` beans explicitly. Since methods are processed within `BPP`, not all beans with `Messaging Annotations` might be processed. INT-3381: PR comments and other fixes INT-3381 Polishing INT-3381: Revert `RouterFB` & `SplitterFB` --- .../integration/annotation/Filter.java | 4 + .../integration/annotation/Router.java | 20 +- .../annotation/ServiceActivator.java | 2 + .../integration/annotation/Splitter.java | 2 + .../config/SplitterFactoryBean.java | 2 +- ...AbstractMethodAnnotationPostProcessor.java | 98 ++++++-- .../AggregatorAnnotationPostProcessor.java | 4 + .../BridgeFromAnnotationPostProcessor.java | 13 +- .../BridgeToAnnotationPostProcessor.java | 18 +- .../FilterAnnotationPostProcessor.java | 58 ++++- ...ChannelAdapterAnnotationPostProcessor.java | 39 +-- .../annotation/MessagingAnnotationUtils.java | 1 + .../RouterAnnotationPostProcessor.java | 110 ++++++++- ...rviceActivatorAnnotationPostProcessor.java | 44 +++- .../SplitterAnnotationPostProcessor.java | 38 ++- .../TransformerAnnotationPostProcessor.java | 20 +- .../router/AbstractMappingMessageRouter.java | 6 + .../FilterAnnotationPostProcessorTests.java | 5 +- ...ingAnnotationsWithBeanAnnotationTests.java | 223 ++++++++++++++++++ .../gemfire/store/GemfireGroupStoreTests.java | 26 +- 20 files changed, 638 insertions(+), 95 deletions(-) create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/config/annotation/MessagingAnnotationsWithBeanAnnotationTests.java diff --git a/spring-integration-core/src/main/java/org/springframework/integration/annotation/Filter.java b/spring-integration-core/src/main/java/org/springframework/integration/annotation/Filter.java index 45697b1cc0..c4c6cc3ab1 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/annotation/Filter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/annotation/Filter.java @@ -48,6 +48,10 @@ public @interface Filter { String outputChannel() default ""; + String discardChannel() default ""; + + String throwExceptionOnRejection() default ""; + String[] adviceChain() default {}; boolean discardWithinAdvice() default true; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/annotation/Router.java b/spring-integration-core/src/main/java/org/springframework/integration/annotation/Router.java index 1dd3f9de51..abcea0cc8b 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/annotation/Router.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/annotation/Router.java @@ -38,7 +38,8 @@ import java.lang.annotation.Target; * whose elements are either * {@link org.springframework.messaging.MessageChannel channels} or * Strings. In the latter case, the endpoint hosting this router will attempt - * to resolve each channel name with the Channel Registry. + * to resolve each channel name with the Channel Registry or with + * {@link #channelMappings()}, if provided. * * @author Mark Fisher * @author Artem Bilan @@ -53,6 +54,23 @@ public @interface Router { String defaultOutputChannel() default ""; + /** + * The 'key=value' pairs to represent channelMapping entries + * @return the channelMappings + * @see org.springframework.integration.router.AbstractMappingMessageRouter#setChannelMapping(String, String) + */ + String[] channelMappings() default {}; + + String prefix() default ""; + + String suffix() default ""; + + String resolutionRequired() default ""; + + String applySequence() default ""; + + String ignoreSendFailures() default ""; + /* {@code SmartLifecycle} options. Can be specified as 'property placeholder', e.g. {@code ${foo.autoStartup}}. diff --git a/spring-integration-core/src/main/java/org/springframework/integration/annotation/ServiceActivator.java b/spring-integration-core/src/main/java/org/springframework/integration/annotation/ServiceActivator.java index bf669599ea..f67cbb92d6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/annotation/ServiceActivator.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/annotation/ServiceActivator.java @@ -51,6 +51,8 @@ public @interface ServiceActivator { String outputChannel() default ""; + String requiresReply() default ""; + String[] adviceChain() default {}; /* diff --git a/spring-integration-core/src/main/java/org/springframework/integration/annotation/Splitter.java b/spring-integration-core/src/main/java/org/springframework/integration/annotation/Splitter.java index fed04887e1..af773dff28 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/annotation/Splitter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/annotation/Splitter.java @@ -50,6 +50,8 @@ public @interface Splitter { String outputChannel() default ""; + String applySequence() default ""; + String[] adviceChain() default {}; /* diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java b/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java index cffe8f9438..e6762db259 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2013 the original author or authors. + * Copyright 2002-2014 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. diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java index 7cc30d06fb..271e807ccb 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java @@ -25,12 +25,17 @@ import java.util.List; import org.aopalliance.aop.Advice; -import org.springframework.beans.factory.BeanFactory; +import org.springframework.aop.TargetSource; +import org.springframework.aop.framework.Advised; import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.context.annotation.Bean; import org.springframework.core.GenericTypeResolver; +import org.springframework.core.annotation.AnnotatedElementUtils; import org.springframework.core.annotation.AnnotationUtils; import org.springframework.core.annotation.Order; +import org.springframework.core.convert.ConversionService; +import org.springframework.core.convert.support.DefaultConversionService; import org.springframework.core.env.Environment; import org.springframework.core.task.TaskExecutor; import org.springframework.integration.annotation.Poller; @@ -72,7 +77,9 @@ public abstract class AbstractMethodAnnotationPostProcessor) GenericTypeResolver.resolveTypeArgument(this.getClass(), @@ -101,12 +117,19 @@ public abstract class AbstractMethodAnnotationPostProcessor annotations) { String inputChannel = MessagingAnnotationUtils.resolveAttribute(annotations, getInputChannelAttribute(), String.class); - return StringUtils.hasText(inputChannel); + boolean createEndpoint = StringUtils.hasText(inputChannel); + if (!createEndpoint && beanAnnotationAware()) { + boolean isBean = AnnotatedElementUtils.isAnnotated(method, Bean.class.getName()); + Assert.isTrue(!isBean, "A channel name in '" + getInputChannelAttribute() + "' is required when " + this.annotationType + + " is used on '@Bean' methods."); + } + return createEndpoint; } protected String getInputChannelAttribute() { return INPUT_CHANNEL_ATTRIBUTE; } + protected boolean beanAnnotationAware() { + return true; + } + protected final void setAdviceChainIfPresent(String beanName, List annotations, MessageHandler handler) { String[] adviceChainNames = MessagingAnnotationUtils.resolveAttribute(annotations, ADVICE_CHAIN_ATTRIBUTE, String[].class); @@ -173,11 +206,8 @@ public abstract class AbstractMethodAnnotationPostProcessor annotations) { + protected AbstractEndpoint doCreateEndpoint(MessageHandler handler, MessageChannel inputChannel,List annotations) { AbstractEndpoint endpoint; if (inputChannel instanceof PollableChannel) { PollingConsumer pollingConsumer = new PollingConsumer((PollableChannel) inputChannel, handler); @@ -294,11 +323,46 @@ public abstract class AbstractMethodAnnotationPostProcessor H extractTypeIfPossible(Object targetObject, Class expectedType) { + if (targetObject == null) { + return null; + } + if (expectedType.isAssignableFrom(targetObject.getClass())) { + return (H) targetObject; + } + if (targetObject instanceof Advised) { + TargetSource targetSource = ((Advised) targetObject).getTargetSource(); + if (targetSource == null) { + return null; + } + try { + return extractTypeIfPossible(targetSource.getTarget(), expectedType); + } catch (Exception e) { + throw new IllegalStateException(e); + } + } + return null; + } + /** * Subclasses must implement this method to create the MessageHandler. * * @param bean The bean. * @param method The method. + * @param annotations The messaging annotation (or meta-annotation hierarchy) on the method. * @return The MessageHandler. */ protected abstract MessageHandler createHandler(Object bean, Method method, List annotations); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java index 604e70e88c..fefdbf2a2c 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java @@ -94,4 +94,8 @@ public class AggregatorAnnotationPostProcessor extends AbstractMethodAnnotationP return handler; } + protected boolean beanAnnotationAware() { + return false; + } + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/BridgeFromAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/BridgeFromAnnotationPostProcessor.java index dc5d346995..ba980769c6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/BridgeFromAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/BridgeFromAnnotationPostProcessor.java @@ -31,7 +31,6 @@ import org.springframework.integration.handler.BridgeHandler; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; import org.springframework.util.Assert; -import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; /** @@ -74,15 +73,9 @@ public class BridgeFromAnnotationPostProcessor extends AbstractMethodAnnotationP @Override protected MessageHandler createHandler(Object bean, Method method, List annotations) { BridgeHandler handler = new BridgeHandler(); - String outputChannelName = null; - String[] names = AnnotationUtils.getAnnotation(method, Bean.class).name(); - if (!ObjectUtils.isEmpty(names)) { - outputChannelName = names[0]; - } - if (!StringUtils.hasText(outputChannelName)) { - outputChannelName = method.getName(); - } - handler.setOutputChannelName(outputChannelName); + Object outputChannel = resolveTargetBeanFromMethodWithBeanAnnotation(method); + Assert.isInstanceOf(MessageChannel.class, outputChannel); + handler.setOutputChannel((MessageChannel) outputChannel); return handler; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/BridgeToAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/BridgeToAnnotationPostProcessor.java index 2b44046377..898381e4b9 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/BridgeToAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/BridgeToAnnotationPostProcessor.java @@ -23,7 +23,6 @@ import java.util.List; import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.context.annotation.Bean; import org.springframework.core.annotation.AnnotatedElementUtils; -import org.springframework.core.annotation.AnnotationUtils; import org.springframework.core.env.Environment; import org.springframework.integration.annotation.BridgeFrom; import org.springframework.integration.annotation.BridgeTo; @@ -32,8 +31,6 @@ import org.springframework.integration.handler.BridgeHandler; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; import org.springframework.util.Assert; -import org.springframework.util.ObjectUtils; -import org.springframework.util.StringUtils; /** * Post-processor for the {@link BridgeTo @BridgeTo} annotation. @@ -73,18 +70,9 @@ public class BridgeToAnnotationPostProcessor extends AbstractMethodAnnotationPos @Override protected AbstractEndpoint createEndpoint(MessageHandler handler, Method method, List annotations) { - String inputChannelName = null; - String[] names = AnnotationUtils.getAnnotation(method, Bean.class).name(); - if (!ObjectUtils.isEmpty(names)) { - inputChannelName = names[0]; - } - if (!StringUtils.hasText(inputChannelName)) { - inputChannelName = method.getName(); - } - - MessageChannel inputChannel = this.beanFactory.getBean(inputChannelName, MessageChannel.class); - - return doCreateEndpoint(handler, inputChannel, annotations); + Object inputChannel = this.resolveTargetBeanFromMethodWithBeanAnnotation(method); + Assert.isInstanceOf(MessageChannel.class, inputChannel); + return doCreateEndpoint(handler, (MessageChannel) inputChannel, annotations); } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/FilterAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/FilterAnnotationPostProcessor.java index 72bdd14c94..d039a9ac9d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/FilterAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/FilterAnnotationPostProcessor.java @@ -21,18 +21,23 @@ import java.lang.reflect.Method; import java.util.List; import org.springframework.beans.factory.ListableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.core.annotation.AnnotatedElementUtils; import org.springframework.core.env.Environment; import org.springframework.integration.annotation.Filter; +import org.springframework.integration.core.MessageSelector; import org.springframework.integration.filter.MessageFilter; import org.springframework.integration.filter.MethodInvokingSelector; import org.springframework.messaging.MessageHandler; import org.springframework.util.Assert; +import org.springframework.util.StringUtils; /** * Post-processor for Methods annotated with {@link Filter @Filter}. * * @author Mark Fisher * @author Gary Russell + * @author Artem Bilan * @since 2.0 */ public class FilterAnnotationPostProcessor extends AbstractMethodAnnotationPostProcessor { @@ -44,16 +49,51 @@ public class FilterAnnotationPostProcessor extends AbstractMethodAnnotationPostP @Override protected MessageHandler createHandler(Object bean, Method method, List annotations) { - Assert.isTrue(boolean.class.equals(method.getReturnType()) || Boolean.class.equals(method.getReturnType()), - "The Filter annotation may only be applied to methods with a boolean return type."); - MethodInvokingSelector selector = new MethodInvokingSelector(bean, method); - MessageFilter filter = new MessageFilter(selector); - this.setOutputChannelIfPresent(annotations, filter); - Boolean discardWithinAdvice = MessagingAnnotationUtils.resolveAttribute(annotations, "discardWithinAdvice", - Boolean.class); - if (discardWithinAdvice != null) { - filter.setDiscardWithinAdvice(discardWithinAdvice); + MessageSelector selector; + if (AnnotatedElementUtils.isAnnotated(method, Bean.class.getName())) { + Object target = this.resolveTargetBeanFromMethodWithBeanAnnotation(method); + if (target instanceof MessageSelector) { + selector = (MessageSelector) target; + } + else if (this.extractTypeIfPossible(target, MessageFilter.class) != null) { + return (MessageHandler) target; + } + else { + selector = new MethodInvokingSelector(target); + } } + else { + Assert.isTrue(boolean.class.equals(method.getReturnType()) || Boolean.class.equals(method.getReturnType()), + "The Filter annotation may only be applied to methods with a boolean return type."); + selector = new MethodInvokingSelector(bean, method); + } + + MessageFilter filter = new MessageFilter(selector); + + /* TODO will be revised in the future + String discardWithinAdvice = MessagingAnnotationUtils.resolveAttribute(annotations, "discardWithinAdvice", + String.class); + if (StringUtils.hasText(discardWithinAdvice)) { + String discardWithinAdviceValue = this.environment.resolvePlaceholders(discardWithinAdvice); + if (StringUtils.hasText(discardWithinAdviceValue)) { + filter.setDiscardWithinAdvice(Boolean.parseBoolean(discardWithinAdviceValue)); + } + }*/ + + filter.setDiscardWithinAdvice(MessagingAnnotationUtils.resolveAttribute(annotations, "discardWithinAdvice", + Boolean.class)); + + String throwExceptionOnRejection = MessagingAnnotationUtils.resolveAttribute(annotations, + "throwExceptionOnRejection", String.class); + if (StringUtils.hasText(throwExceptionOnRejection)) { + String throwExceptionOnRejectionValue = this.environment.resolvePlaceholders(throwExceptionOnRejection); + filter.setThrowExceptionOnRejection(Boolean.parseBoolean(throwExceptionOnRejectionValue)); + } + + String discardChannelName = MessagingAnnotationUtils.resolveAttribute(annotations, "discardChannel", String.class); + filter.setDiscardChannelName(discardChannelName); + + this.setOutputChannelIfPresent(annotations, filter); return filter; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/InboundChannelAdapterAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/InboundChannelAdapterAnnotationPostProcessor.java index f010ae3e8e..84e0f374fc 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/InboundChannelAdapterAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/InboundChannelAdapterAnnotationPostProcessor.java @@ -21,11 +21,13 @@ import java.lang.reflect.Method; import java.util.List; import org.springframework.beans.factory.ListableBeanFactory; -import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.core.annotation.AnnotatedElementUtils; import org.springframework.core.annotation.AnnotationUtils; import org.springframework.core.env.Environment; import org.springframework.integration.annotation.InboundChannelAdapter; import org.springframework.integration.config.IntegrationConfigUtils; +import org.springframework.integration.core.MessageSource; import org.springframework.integration.endpoint.MethodInvokingMessageSource; import org.springframework.integration.endpoint.SourcePollingChannelAdapter; import org.springframework.messaging.MessageChannel; @@ -53,26 +55,12 @@ public class InboundChannelAdapterAnnotationPostProcessor extends @Override public Object postProcess(Object bean, String beanName, Method method, List annotations) { - Assert.isTrue(!Void.class.isAssignableFrom(method.getReturnType()), "The method '" + method - + "' for 'SourcePollingChannelAdapter' must not have 'void' return type."); - Assert.isTrue(method.getParameterTypes().length == 0, "The method '" + method - + "' for 'SourcePollingChannelAdapter' must not have any parameters."); - String channelName = MessagingAnnotationUtils.resolveAttribute(annotations, AnnotationUtils.VALUE, String.class); Assert.hasText(channelName, "The channel ('value' attribute of @InboundChannelAdapter) can't be empty."); - MessageChannel channel = this.channelResolver.resolveDestination(channelName); + MessageSource messageSource = this.createMessageSource(bean, beanName, method); - MethodInvokingMessageSource messageSource = new MethodInvokingMessageSource(); - messageSource.setObject(bean); - messageSource.setMethod(method); - if (beanFactory instanceof ConfigurableListableBeanFactory) { - String handlerBeanName = this.generateHandlerBeanName(beanName, method); - ConfigurableListableBeanFactory listableBeanFactory = (ConfigurableListableBeanFactory) beanFactory; - listableBeanFactory.registerSingleton(handlerBeanName, messageSource); - messageSource = (MethodInvokingMessageSource) listableBeanFactory - .initializeBean(messageSource, handlerBeanName); - } + MessageChannel channel = this.channelResolver.resolveDestination(channelName); SourcePollingChannelAdapter adapter = new SourcePollingChannelAdapter(); adapter.setOutputChannel(channel); @@ -82,6 +70,23 @@ public class InboundChannelAdapterAnnotationPostProcessor extends return adapter; } + private MessageSource createMessageSource(Object bean, String beanName, Method method) { + if (AnnotatedElementUtils.isAnnotated(method, Bean.class.getName())) { + Object target = this.resolveTargetBeanFromMethodWithBeanAnnotation(method); + Assert.isInstanceOf(MessageSource.class, target, "The '" + this.annotationType + "' on @Bean method " + + "level is allowed only for: " + MessageSource.class.getName() + "beans"); + return (MessageSource) target; + } + else { + MethodInvokingMessageSource messageSource = new MethodInvokingMessageSource(); + messageSource.setObject(bean); + messageSource.setMethod(method); + String messageSourceBeanName = this.generateHandlerBeanName(beanName, method); + this.beanFactory.registerSingleton(messageSourceBeanName, messageSource); + return (MessageSource) this.beanFactory.initializeBean(messageSource, messageSourceBeanName); + } + } + @Override protected String generateHandlerBeanName(String originalBeanName, Method method) { return super.generateHandlerBeanName(originalBeanName, method) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationUtils.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationUtils.java index 1b108b697e..500cf9b3a7 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationUtils.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/MessagingAnnotationUtils.java @@ -48,6 +48,7 @@ public final class MessagingAnnotationUtils { * @param annotations The meta-annotations in order (closest first). * @param name The attribute name. * @param requiredType The expected type. + * @param The type. * @return The value. */ @SuppressWarnings("unchecked") diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/RouterAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/RouterAnnotationPostProcessor.java index ed701fa18b..4a962bb35b 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/RouterAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/RouterAnnotationPostProcessor.java @@ -19,14 +19,19 @@ package org.springframework.integration.config.annotation; import java.lang.annotation.Annotation; import java.lang.reflect.Method; import java.util.List; +import java.util.Properties; import org.springframework.beans.factory.ListableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.core.annotation.AnnotatedElementUtils; +import org.springframework.core.convert.TypeDescriptor; import org.springframework.core.env.Environment; import org.springframework.integration.annotation.Router; +import org.springframework.integration.router.AbstractMessageRouter; import org.springframework.integration.router.MethodInvokingRouter; -import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; import org.springframework.util.Assert; +import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; /** @@ -34,6 +39,7 @@ import org.springframework.util.StringUtils; * * @author Mark Fisher * @author Gary Russell + * @author Artem Bilan */ public class RouterAnnotationPostProcessor extends AbstractMethodAnnotationPostProcessor { @@ -43,18 +49,100 @@ public class RouterAnnotationPostProcessor extends AbstractMethodAnnotationPostP @Override - protected MessageHandler createHandler(Object bean, Method method, - List annotations) { - MethodInvokingRouter router = new MethodInvokingRouter(bean, method); - router.setBeanFactory(this.beanFactory); - String defaultOutputChannelName = MessagingAnnotationUtils.resolveAttribute(annotations, "defaultOutputChannel", - String.class); - if (StringUtils.hasText(defaultOutputChannelName)) { - MessageChannel defaultOutputChannel = this.channelResolver.resolveDestination(defaultOutputChannelName); - Assert.notNull(defaultOutputChannel, "unable to resolve defaultOutputChannel '" + defaultOutputChannelName + "'"); - router.setDefaultOutputChannel(defaultOutputChannel); + protected MessageHandler createHandler(Object bean, Method method, List annotations) { + AbstractMessageRouter router; + if (AnnotatedElementUtils.isAnnotated(method, Bean.class.getName())) { + Object target = this.resolveTargetBeanFromMethodWithBeanAnnotation(method); + router = this.extractTypeIfPossible(target, AbstractMessageRouter.class); + if (router == null) { + if (target instanceof MessageHandler) { + Assert.isTrue(this.routerAttributesProvided(annotations), "'defaultOutputChannel', " + + "'applySequence', 'ignoreSendFailures', 'resolutionRequired' and 'channelMappings' " + + "can be applied to 'AbstractMessageRouter' implementations, but target handler is: " + + target.getClass()); + return (MessageHandler) target; + } + else { + router = new MethodInvokingRouter(target); + } + } + else { + return router; + } } + else { + router = new MethodInvokingRouter(bean, method); + } + String defaultOutputChannelName = MessagingAnnotationUtils.resolveAttribute(annotations, + "defaultOutputChannel", String.class); + router.setDefaultOutputChannelName(defaultOutputChannelName); + + String applySequence = MessagingAnnotationUtils.resolveAttribute(annotations, "applySequence", String.class); + if (StringUtils.hasText(applySequence)) { + router.setApplySequence(Boolean.parseBoolean(this.environment.resolvePlaceholders(applySequence))); + } + + String ignoreSendFailures = MessagingAnnotationUtils.resolveAttribute(annotations, "ignoreSendFailures", + String.class); + if (StringUtils.hasText(ignoreSendFailures)) { + router.setIgnoreSendFailures(Boolean.parseBoolean(this.environment.resolvePlaceholders(ignoreSendFailures))); + } + + if (this.routerAttributesProvided(annotations)) { + + MethodInvokingRouter methodInvokingRouter = (MethodInvokingRouter) router; + + String resolutionRequired = MessagingAnnotationUtils.resolveAttribute(annotations, "resolutionRequired", + String.class); + if (StringUtils.hasText(resolutionRequired)) { + String resolutionRequiredValue = this.environment.resolvePlaceholders(resolutionRequired); + if (StringUtils.hasText(resolutionRequiredValue)) { + methodInvokingRouter.setResolutionRequired(Boolean.parseBoolean(resolutionRequiredValue)); + } + } + + String prefix = MessagingAnnotationUtils.resolveAttribute(annotations, "prefix", String.class); + if (StringUtils.hasText(prefix)) { + methodInvokingRouter.setPrefix(this.environment.resolvePlaceholders(prefix)); + } + + String suffix = MessagingAnnotationUtils.resolveAttribute(annotations, "suffix", String.class); + if (StringUtils.hasText(suffix)) { + methodInvokingRouter.setSuffix(this.environment.resolvePlaceholders(suffix)); + } + + String[] channelMappings = MessagingAnnotationUtils.resolveAttribute(annotations, "channelMappings", + String[].class); + if (!ObjectUtils.isEmpty(channelMappings)) { + StringBuilder mappings = new StringBuilder(); + for (String channelMapping : channelMappings) { + mappings.append(channelMapping).append("\n"); + } + Properties properties = (Properties) this.conversionService.convert(mappings.toString(), + TypeDescriptor.valueOf(String.class), TypeDescriptor.valueOf(Properties.class)); + methodInvokingRouter.replaceChannelMappings(properties); + } + + } + return router; } + private boolean routerAttributesProvided(List annotations) { + String defaultOutputChannel = MessagingAnnotationUtils.resolveAttribute(annotations, "defaultOutputChannel", + String.class); + String[] channelMappings = MessagingAnnotationUtils.resolveAttribute(annotations, "channelMappings", + String[].class); + String prefix = MessagingAnnotationUtils.resolveAttribute(annotations, "prefix", String.class); + String suffix = MessagingAnnotationUtils.resolveAttribute(annotations, "suffix", String.class); + String resolutionRequired = MessagingAnnotationUtils.resolveAttribute(annotations, "resolutionRequired", + String.class); + String applySequence = MessagingAnnotationUtils.resolveAttribute(annotations, "applySequence", String.class); + String ignoreSendFailures = MessagingAnnotationUtils.resolveAttribute(annotations, "ignoreSendFailures", + String.class); + return StringUtils.hasText(defaultOutputChannel) || !ObjectUtils.isEmpty(channelMappings) + || StringUtils.hasText(prefix) || StringUtils.hasText(suffix) || StringUtils.hasText(resolutionRequired) + || StringUtils.hasText(applySequence) || StringUtils.hasText(ignoreSendFailures); + } + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/ServiceActivatorAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/ServiceActivatorAnnotationPostProcessor.java index 3752919892..fcc7a20a78 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/ServiceActivatorAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/ServiceActivatorAnnotationPostProcessor.java @@ -21,16 +21,22 @@ import java.lang.reflect.Method; import java.util.List; import org.springframework.beans.factory.ListableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.core.annotation.AnnotatedElementUtils; import org.springframework.core.env.Environment; import org.springframework.integration.annotation.ServiceActivator; +import org.springframework.integration.handler.AbstractReplyProducingMessageHandler; import org.springframework.integration.handler.ServiceActivatingHandler; +import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandler; +import org.springframework.util.StringUtils; /** * Post-processor for Methods annotated with {@link ServiceActivator @ServiceActivator}. * * @author Mark Fisher * @author Gary Russell + * @author Artem Bilan */ public class ServiceActivatorAnnotationPostProcessor extends AbstractMethodAnnotationPostProcessor { @@ -41,7 +47,43 @@ public class ServiceActivatorAnnotationPostProcessor extends AbstractMethodAnnot @Override protected MessageHandler createHandler(Object bean, Method method, List annotations) { - ServiceActivatingHandler serviceActivator = new ServiceActivatingHandler(bean, method); + AbstractReplyProducingMessageHandler serviceActivator; + if (AnnotatedElementUtils.isAnnotated(method, Bean.class.getName())) { + final Object target = this.resolveTargetBeanFromMethodWithBeanAnnotation(method); + serviceActivator = this.extractTypeIfPossible(target, AbstractReplyProducingMessageHandler.class); + if (serviceActivator == null) { + if (target instanceof MessageHandler) { + /* + * Return a reply-producing message handler so that we still get 'produced no reply' messages + * and the super class will inject the advice chain to advise the handler method if needed. + */ + return new AbstractReplyProducingMessageHandler() { + + @Override + protected Object handleRequestMessage(Message requestMessage) { + + ((MessageHandler) target).handleMessage(requestMessage); + return null; + } + }; + } + else { + serviceActivator = new ServiceActivatingHandler(target); + } + } + else { + return (MessageHandler) target; + } + } + else { + serviceActivator = new ServiceActivatingHandler(bean, method); + } + + String requiresReply = MessagingAnnotationUtils.resolveAttribute(annotations, "requiresReply", String.class); + if (StringUtils.hasText(requiresReply)) { + serviceActivator.setRequiresReply(Boolean.parseBoolean(this.environment.resolvePlaceholders(requiresReply))); + } + this.setOutputChannelIfPresent(annotations, serviceActivator); return serviceActivator; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/SplitterAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/SplitterAnnotationPostProcessor.java index 21ca6d2da2..3872dbb44d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/SplitterAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/SplitterAnnotationPostProcessor.java @@ -21,16 +21,22 @@ import java.lang.reflect.Method; import java.util.List; import org.springframework.beans.factory.ListableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.core.annotation.AnnotatedElementUtils; import org.springframework.core.env.Environment; import org.springframework.integration.annotation.Splitter; +import org.springframework.integration.splitter.AbstractMessageSplitter; import org.springframework.integration.splitter.MethodInvokingSplitter; import org.springframework.messaging.MessageHandler; +import org.springframework.util.Assert; +import org.springframework.util.StringUtils; /** * Post-processor for Methods annotated with {@link Splitter @Splitter}. * * @author Mark Fisher * @author Gary Russell + * @author Artem Bilan */ public class SplitterAnnotationPostProcessor extends AbstractMethodAnnotationPostProcessor { @@ -41,7 +47,37 @@ public class SplitterAnnotationPostProcessor extends AbstractMethodAnnotationPos @Override protected MessageHandler createHandler(Object bean, Method method, List annotations) { - MethodInvokingSplitter splitter = new MethodInvokingSplitter(bean, method); + String applySequence = MessagingAnnotationUtils.resolveAttribute(annotations, "applySequence", String.class); + + AbstractMessageSplitter splitter; + if (AnnotatedElementUtils.isAnnotated(method, Bean.class.getName())) { + Object target = this.resolveTargetBeanFromMethodWithBeanAnnotation(method); + splitter = this.extractTypeIfPossible(target, AbstractMessageSplitter.class); + if (splitter == null) { + if (target instanceof MessageHandler) { + Assert.hasText(applySequence, "'applySequence' can be applied to 'AbstractMessageSplitter', but " + + "target handler is: " + target.getClass()); + return (MessageHandler) target; + } + else { + splitter = new MethodInvokingSplitter(target); + } + } + else { + return splitter; + } + } + else { + splitter = new MethodInvokingSplitter(bean, method); + } + + if (StringUtils.hasText(applySequence)) { + String applySequenceValue = this.environment.resolvePlaceholders(applySequence); + if (StringUtils.hasText(applySequenceValue)) { + splitter.setApplySequence(Boolean.parseBoolean(applySequenceValue)); + } + } + this.setOutputChannelIfPresent(annotations, splitter); return splitter; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/TransformerAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/TransformerAnnotationPostProcessor.java index da0db7e9fd..d1f431595c 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/TransformerAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/TransformerAnnotationPostProcessor.java @@ -21,8 +21,11 @@ import java.lang.reflect.Method; import java.util.List; import org.springframework.beans.factory.ListableBeanFactory; +import org.springframework.context.annotation.Bean; +import org.springframework.core.annotation.AnnotatedElementUtils; import org.springframework.core.env.Environment; import org.springframework.integration.annotation.Transformer; +import org.springframework.integration.handler.AbstractReplyProducingMessageHandler; import org.springframework.integration.transformer.MessageTransformingHandler; import org.springframework.integration.transformer.MethodInvokingTransformer; import org.springframework.messaging.MessageHandler; @@ -32,6 +35,7 @@ import org.springframework.messaging.MessageHandler; * * @author Mark Fisher * @author Gary Russell + * @author Artem Bilan */ public class TransformerAnnotationPostProcessor extends AbstractMethodAnnotationPostProcessor { @@ -42,7 +46,21 @@ public class TransformerAnnotationPostProcessor extends AbstractMethodAnnotation @Override protected MessageHandler createHandler(Object bean, Method method, List annotations) { - MethodInvokingTransformer transformer = new MethodInvokingTransformer(bean, method); + org.springframework.integration.transformer.Transformer transformer; + if (AnnotatedElementUtils.isAnnotated(method, Bean.class.getName())) { + Object target = this.resolveTargetBeanFromMethodWithBeanAnnotation(method); + transformer = this.extractTypeIfPossible(target, org.springframework.integration.transformer.Transformer.class); + if (transformer == null) { + if (this.extractTypeIfPossible(target, AbstractReplyProducingMessageHandler.class) != null) { + return (MessageHandler) target; + } + transformer = new MethodInvokingTransformer(target); + } + } + else { + transformer = new MethodInvokingTransformer(bean, method); + } + MessageTransformingHandler handler = new MessageTransformingHandler(transformer); this.setOutputChannelIfPresent(annotations, handler); return handler; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/router/AbstractMappingMessageRouter.java b/spring-integration-core/src/main/java/org/springframework/integration/router/AbstractMappingMessageRouter.java index bb48c7c7f8..5e3667c2d8 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/router/AbstractMappingMessageRouter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/router/AbstractMappingMessageRouter.java @@ -154,6 +154,12 @@ public abstract class AbstractMappingMessageRouter extends AbstractMessageRouter @Override public void onInit() { + try { + super.onInit(); + } + catch (Exception e) { + throw new IllegalStateException(e); + } BeanFactory beanFactory = this.getBeanFactory(); if (this.channelResolver == null && beanFactory != null) { this.channelResolver = new BeanFactoryChannelResolver(beanFactory); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/annotation/FilterAnnotationPostProcessorTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/annotation/FilterAnnotationPostProcessorTests.java index 15bffb8a33..88053c0941 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/annotation/FilterAnnotationPostProcessorTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/annotation/FilterAnnotationPostProcessorTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2013 the original author or authors. + * Copyright 2002-2014 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. @@ -39,10 +39,12 @@ import org.springframework.integration.handler.advice.AbstractRequestHandlerAdvi import org.springframework.messaging.support.GenericMessage; import org.springframework.integration.test.util.TestUtils; import org.springframework.integration.test.util.TestUtils.TestApplicationContext; +import org.springframework.mock.env.MockEnvironment; /** * @author Mark Fisher * @author Gary Russell + * @author Artem Bilan * @since 2.0 */ public class FilterAnnotationPostProcessorTests { @@ -60,6 +62,7 @@ public class FilterAnnotationPostProcessorTests { context.registerChannel("input", inputChannel); context.registerChannel("output", outputChannel); postProcessor.setBeanFactory(context.getBeanFactory()); + postProcessor.setEnvironment(new MockEnvironment()); postProcessor.afterPropertiesSet(); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/annotation/MessagingAnnotationsWithBeanAnnotationTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/annotation/MessagingAnnotationsWithBeanAnnotationTests.java new file mode 100644 index 0000000000..4045cbb0ae --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/annotation/MessagingAnnotationsWithBeanAnnotationTests.java @@ -0,0 +1,223 @@ +/* + * Copyright 2014 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.config.annotation; + +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertThat; +import static org.junit.Assert.assertTrue; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; + +import javax.annotation.Resource; + +import org.hamcrest.Matchers; +import org.junit.Test; +import org.junit.runner.RunWith; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.expression.ExpressionParser; +import org.springframework.expression.spel.standard.SpelExpressionParser; +import org.springframework.integration.aggregator.AggregatingMessageHandler; +import org.springframework.integration.aggregator.ExpressionEvaluatingCorrelationStrategy; +import org.springframework.integration.aggregator.ExpressionEvaluatingReleaseStrategy; +import org.springframework.integration.aggregator.PassThroughMessageGroupProcessor; +import org.springframework.integration.annotation.Filter; +import org.springframework.integration.annotation.InboundChannelAdapter; +import org.springframework.integration.annotation.Poller; +import org.springframework.integration.annotation.Router; +import org.springframework.integration.annotation.ServiceActivator; +import org.springframework.integration.annotation.Splitter; +import org.springframework.integration.annotation.Transformer; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.config.EnableIntegration; +import org.springframework.integration.config.EnableMessageHistory; +import org.springframework.integration.core.MessageSelector; +import org.springframework.integration.core.MessageSource; +import org.springframework.integration.endpoint.SourcePollingChannelAdapter; +import org.springframework.integration.filter.ExpressionEvaluatingSelector; +import org.springframework.integration.history.MessageHistory; +import org.springframework.integration.splitter.DefaultMessageSplitter; +import org.springframework.integration.transformer.ExpressionEvaluatingTransformer; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.MessageHandler; +import org.springframework.messaging.MessagingException; +import org.springframework.messaging.PollableChannel; +import org.springframework.messaging.support.GenericMessage; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.test.context.support.AnnotationConfigContextLoader; + +/** + * @author Artem Bilan + * @since 4.0 + */ +@ContextConfiguration(loader = AnnotationConfigContextLoader.class) +@RunWith(SpringJUnit4ClassRunner.class) +@DirtiesContext +public class MessagingAnnotationsWithBeanAnnotationTests { + + @Autowired + private SourcePollingChannelAdapter sourcePollingChannelAdapter; + + @Autowired + private PollableChannel discardChannel; + + @Resource(name="collector") + private List> collector; + + @Test + public void testMessagingAnnotationsFlow() { + this.sourcePollingChannelAdapter.start(); + for (int i = 0; i < 10; i++) { + Message receive = this.discardChannel.receive(1000); + assertNotNull(receive); + assertTrue(((Integer) receive.getPayload()) % 2 == 0); + } + for (Message message : collector) { + assertFalse(((Integer) message.getPayload()) % 2 == 0); + MessageHistory messageHistory = MessageHistory.read(message); + assertNotNull(messageHistory); + String messageHistoryString = messageHistory.toString(); + assertThat(messageHistoryString, Matchers.containsString("routerChannel")); + assertThat(messageHistoryString, Matchers.containsString("filterChannel")); + assertThat(messageHistoryString, Matchers.containsString("aggregatorChannel")); + assertThat(messageHistoryString, Matchers.containsString("splitterChannel")); + assertThat(messageHistoryString, Matchers.containsString("serviceChannel")); + assertThat(messageHistoryString, Matchers.not(Matchers.containsString("discardChannel"))); + } + } + + @Configuration + @EnableIntegration + @EnableMessageHistory + public static class ContextConfiguration { + + private static final ExpressionParser PARSER = new SpelExpressionParser(); + + @Bean + public AtomicInteger counter() { + return new AtomicInteger(); + } + + @Bean + @InboundChannelAdapter(value = "routerChannel", autoStartup = "false", + poller = @Poller(fixedRate = "10", maxMessagesPerPoll = "1")) + public MessageSource counterMessageSource(final AtomicInteger counter) { + return new MessageSource() { + + @Override + public Message receive() { + return new GenericMessage(counter.incrementAndGet()); + } + }; + } + + @Bean + public MessageChannel routerChannel() { + return new DirectChannel(); + } + + @Bean + @Router(inputChannel = "routerChannel", channelMappings = {"true=odd", "false=filter"}, suffix = "Channel") + public MessageSelector router() { + return new ExpressionEvaluatingSelector("payload % 2 == 0"); + } + + @Bean + @Transformer(inputChannel = "oddChannel", outputChannel = "filterChannel") + public ExpressionEvaluatingTransformer oddTransformer() { + return new ExpressionEvaluatingTransformer(PARSER.parseExpression("payload / 2")); + } + + @Bean + public MessageChannel filterChannel() { + return new DirectChannel(); + } + + @Bean + @Filter(inputChannel = "filterChannel", outputChannel = "aggregatorChannel", discardChannel = "discardChannel") + public MessageSelector filter() { + return new ExpressionEvaluatingSelector("payload % 2 != 0"); + } + + @Bean + public MessageChannel aggregatorChannel() { + return new DirectChannel(); + } + + @Bean + @ServiceActivator(inputChannel = "aggregatorChannel") + public MessageHandler aggregator() { + AggregatingMessageHandler handler = new AggregatingMessageHandler(new PassThroughMessageGroupProcessor()); + handler.setCorrelationStrategy(new ExpressionEvaluatingCorrelationStrategy("1")); + handler.setReleaseStrategy(new ExpressionEvaluatingReleaseStrategy("size() == 10")); + handler.setOutputChannelName("splitterChannel"); + return handler; + } + + @Bean + public MessageChannel splitterChannel() { + return new DirectChannel(); + } + + @Bean + @Splitter(inputChannel = "splitterChannel") + public MessageHandler splitter() { + DefaultMessageSplitter defaultMessageSplitter = new DefaultMessageSplitter(); + defaultMessageSplitter.setOutputChannelName("serviceChannel"); + return defaultMessageSplitter; + } + + @Bean + public PollableChannel discardChannel() { + return new QueueChannel(); + } + + @Bean + public List> collector() { + return new ArrayList>(); + } + + @Bean + public MessageChannel serviceChannel() { + return new DirectChannel(); + } + + @Bean + @ServiceActivator(inputChannel = "serviceChannel") + public MessageHandler service() { + final List> collector = this.collector(); + return new MessageHandler() { + + @Override + public void handleMessage(Message message) throws MessagingException { + collector.add(message); + } + }; + } + + } + +} diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireGroupStoreTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireGroupStoreTests.java index 45a128f370..83d17b6e09 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireGroupStoreTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireGroupStoreTests.java @@ -39,16 +39,16 @@ import org.junit.Test; import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.data.gemfire.CacheFactoryBean; -import org.springframework.messaging.Message; -import org.springframework.messaging.MessageChannel; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.history.MessageHistory; -import org.springframework.messaging.support.GenericMessage; import org.springframework.integration.store.MessageGroup; import org.springframework.integration.store.SimpleMessageGroup; import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.test.support.LongRunningIntegrationTest; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.support.GenericMessage; import org.springframework.util.Assert; import com.gemstone.gemfire.cache.Cache; @@ -291,6 +291,7 @@ public class GemfireGroupStoreTests { executor = Executors.newCachedThreadPool(); executor.execute(new Runnable() { + @Override public void run() { MessageGroup group = store1.addMessageToGroup(1, message); if (group.getMessages().size() != 1) { @@ -300,6 +301,7 @@ public class GemfireGroupStoreTests { } }); executor.execute(new Runnable() { + @Override public void run() { MessageGroup group = store2.removeMessageFromGroup(1, message); if (group.getMessages().size() != 0) { @@ -319,10 +321,10 @@ public class GemfireGroupStoreTests { @Test public void testWithAggregatorWithShutdown() { - ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("gemfire-aggregator-config.xml", + ClassPathXmlApplicationContext context1 = new ClassPathXmlApplicationContext("gemfire-aggregator-config.xml", this.getClass()); - MessageChannel input = context.getBean("inputChannel", MessageChannel.class); - QueueChannel output = context.getBean("outputChannel", QueueChannel.class); + MessageChannel input = context1.getBean("inputChannel", MessageChannel.class); + QueueChannel output = context1.getBean("outputChannel", QueueChannel.class); Message m1 = MessageBuilder.withPayload("1").setSequenceNumber(1).setSequenceSize(3).setCorrelationId(1) .build(); @@ -333,14 +335,17 @@ public class GemfireGroupStoreTests { input.send(m2); assertNull(output.receive(1000)); - context = new ClassPathXmlApplicationContext("gemfire-aggregator-config-a.xml", this.getClass()); - MessageChannel inputA = context.getBean("inputChannel", MessageChannel.class); - QueueChannel outputA = context.getBean("outputChannel", QueueChannel.class); + ClassPathXmlApplicationContext context2 = new ClassPathXmlApplicationContext("gemfire-aggregator-config-a.xml", + this.getClass()); + MessageChannel inputA = context2.getBean("inputChannel", MessageChannel.class); + QueueChannel outputA = context2.getBean("outputChannel", QueueChannel.class); Message m3 = MessageBuilder.withPayload("3").setSequenceNumber(3).setSequenceSize(3).setCorrelationId(1) .build(); inputA.send(m3); assertNotNull(outputA.receive(1000)); + context1.close(); + context2.close(); } @Test @@ -356,9 +361,10 @@ public class GemfireGroupStoreTests { Thread.sleep(1); } for (int i = 0; i < 20; i++) { - assertNotNull(outputQueue.receive(1)); + assertNotNull(outputQueue.receive(5000)); } assertNull(outputQueue.receive(1)); + context.close(); } @Before