From 0cc9273a2e05e975e541022badfe4c9a37cbafb6 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 6 Nov 2014 19:42:57 +0200 Subject: [PATCH] INT-3551: Idempotent Receiver: Add value-strategy JIRA: https://jira.spring.io/browse/INT-3551 * Rename `MetadataKeyStrategy` -> `MetadataEntryStrategy` * Add `valueStrategy` to the `MetadataStoreSelector` * Add `value-strategy` and `value-expression` to the `` INT-3551: Add `@IR` support on service methods Rework `MetadataKeyStrategy` just to the `MessageProcessor` Fix Docs Minor Doc Polishing. --- .../annotation/IdempotentReceiver.java | 2 +- ...AbstractMethodAnnotationPostProcessor.java | 26 +++++++ .../IdempotentReceiverInterceptorParser.java | 69 ++++++++++++++----- .../ExpressionEvaluatingMessageProcessor.java | 11 +-- .../ExpressionMetadataKeyStrategy.java | 63 ----------------- .../metadata/MetadataKeyStrategy.java | 32 --------- .../selector/MetadataStoreSelector.java | 37 +++++++--- .../config/xml/spring-integration-4.1.xsd | 34 ++++++++- .../IdempotentReceiverParserTests-context.xml | 18 +++-- .../xml/IdempotentReceiverParserTests.java | 68 ++++++++++++++---- .../idempotent-receiver-configs.properties | 15 +++- .../advice/IdempotentReceiverTests.java | 9 ++- .../IdempotentReceiverIntegrationTests.java | 66 ++++++++++++++++-- src/reference/docbook/handler-advice.xml | 50 +++++++++++--- 14 files changed, 327 insertions(+), 173 deletions(-) delete mode 100644 spring-integration-core/src/main/java/org/springframework/integration/metadata/ExpressionMetadataKeyStrategy.java delete mode 100644 spring-integration-core/src/main/java/org/springframework/integration/metadata/MetadataKeyStrategy.java diff --git a/spring-integration-core/src/main/java/org/springframework/integration/annotation/IdempotentReceiver.java b/spring-integration-core/src/main/java/org/springframework/integration/annotation/IdempotentReceiver.java index 9cf5ed22f4..cded71b469 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/annotation/IdempotentReceiver.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/annotation/IdempotentReceiver.java @@ -23,7 +23,7 @@ import java.lang.annotation.RetentionPolicy; import java.lang.annotation.Target; /** - * A {@code @Bean} that has a MessagingAnnotation (@code @ServiceActivator, @Router etc.) + * A {@code method} that has a MessagingAnnotation (@code @ServiceActivator, @Router etc.) * that also has this annotation, has an * {@link org.springframework.integration.handler.advice.IdempotentReceiverInterceptor} applied * to the associated {@link org.springframework.messaging.MessageHandler#handleMessage} method. 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 c848ac8593..e275f2673a 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 @@ -27,6 +27,9 @@ import org.aopalliance.aop.Advice; import org.springframework.aop.TargetSource; import org.springframework.aop.framework.Advised; +import org.springframework.aop.framework.ProxyFactory; +import org.springframework.aop.support.DefaultBeanFactoryPointcutAdvisor; +import org.springframework.aop.support.NameMatchMethodPointcut; import org.springframework.beans.factory.ListableBeanFactory; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.context.annotation.Bean; @@ -38,6 +41,7 @@ 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.IdempotentReceiver; import org.springframework.integration.annotation.Poller; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.config.IntegrationConfigUtils; @@ -131,6 +135,28 @@ public abstract class AbstractMethodAnnotationPostProcessor extends AbstractMessageProcessor { @@ -37,18 +38,15 @@ public class ExpressionEvaluatingMessageProcessor extends AbstractMessageProc /** * Create an {@link ExpressionEvaluatingMessageProcessor} for the given expression. - * * @param expression The expression. */ public ExpressionEvaluatingMessageProcessor(Expression expression) { this(expression, null); } - /** * Create an {@link ExpressionEvaluatingMessageProcessor} for the given expression * and expected type for its evaluation result. - * * @param expression The expression. * @param expectedType The expected type. */ @@ -63,11 +61,9 @@ public class ExpressionEvaluatingMessageProcessor extends AbstractMessageProc } } - /** * Processes the Message by evaluating the expression with that Message as the * root object. The expression evaluation result Object will be returned. - * * @param message The message. * @return The result of processing the message. */ @@ -76,4 +72,9 @@ public class ExpressionEvaluatingMessageProcessor extends AbstractMessageProc return this.evaluateExpression(this.expression, message, this.expectedType); } + @Override + public String toString() { + return "ExpressionEvaluatingMessageProcessor for: [" + this.expression.getExpressionString() + "]"; + } + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/metadata/ExpressionMetadataKeyStrategy.java b/spring-integration-core/src/main/java/org/springframework/integration/metadata/ExpressionMetadataKeyStrategy.java deleted file mode 100644 index cab93683ad..0000000000 --- a/spring-integration-core/src/main/java/org/springframework/integration/metadata/ExpressionMetadataKeyStrategy.java +++ /dev/null @@ -1,63 +0,0 @@ -/* - * 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.metadata; - -import org.springframework.beans.BeansException; -import org.springframework.beans.factory.BeanFactory; -import org.springframework.beans.factory.BeanFactoryAware; -import org.springframework.expression.ExpressionParser; -import org.springframework.expression.spel.standard.SpelExpressionParser; -import org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor; -import org.springframework.integration.handler.MessageProcessor; -import org.springframework.messaging.Message; - -/** - * The expression based {@link MetadataKeyStrategy} implementation. - * The provided {@link Message} is used as the evaluation context root object. - * - * @author Artem Bilan - * @since 4.1 - */ -public class ExpressionMetadataKeyStrategy implements MetadataKeyStrategy, BeanFactoryAware { - - private static final ExpressionParser PARSER = new SpelExpressionParser(); - - private final MessageProcessor processor; - - private final String expressionString; - - public ExpressionMetadataKeyStrategy(String expressionString) { - this.processor = new ExpressionEvaluatingMessageProcessor(PARSER.parseExpression(expressionString)); - this.expressionString = expressionString; - } - - @Override - public void setBeanFactory(BeanFactory beanFactory) throws BeansException { - ((BeanFactoryAware) this.processor).setBeanFactory(beanFactory); - } - - @Override - public String getKey(Message message) { - return this.processor.processMessage(message); - } - - @Override - public String toString() { - return "ExpressionEvaluatingSelector for: [" + this.expressionString + "]"; - } - -} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/metadata/MetadataKeyStrategy.java b/spring-integration-core/src/main/java/org/springframework/integration/metadata/MetadataKeyStrategy.java deleted file mode 100644 index c1b5a0a94a..0000000000 --- a/spring-integration-core/src/main/java/org/springframework/integration/metadata/MetadataKeyStrategy.java +++ /dev/null @@ -1,32 +0,0 @@ -/* - * 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.metadata; - -import org.springframework.messaging.Message; - -/** - * The strategy to extract a {@code key} for the {@code MetadataStore} - * from the provided {@link Message}. - * - * @author Artem Bilan - * @since 4.1 - */ -public interface MetadataKeyStrategy { - - String getKey(Message message); - -} diff --git a/spring-integration-core/src/main/java/org/springframework/integration/selector/MetadataStoreSelector.java b/spring-integration-core/src/main/java/org/springframework/integration/selector/MetadataStoreSelector.java index a9b37b88fc..ff55a88cc6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/selector/MetadataStoreSelector.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/selector/MetadataStoreSelector.java @@ -17,19 +17,19 @@ package org.springframework.integration.selector; import org.springframework.integration.core.MessageSelector; +import org.springframework.integration.handler.MessageProcessor; import org.springframework.integration.metadata.ConcurrentMetadataStore; -import org.springframework.integration.metadata.MetadataKeyStrategy; import org.springframework.integration.metadata.SimpleMetadataStore; import org.springframework.messaging.Message; import org.springframework.util.Assert; /** * The {@link MessageSelector} implementation using a {@link ConcurrentMetadataStore} - * and {@link MetadataKeyStrategy}. + * and {@link MessageProcessor}. *

* The {@link #accept} method extracts {@code metadataKey} from the provided {@code message} - * using {@link MetadataKeyStrategy} and uses the {@code timestamp} header as the {@code value} - * (hex). + * using {@link MessageProcessor} and uses the {@code timestamp} header as the {@code value} + * (hex) by default. The {@link #valueStrategy} can be provided to override the default behaviour. *

* The successful result of the {@link #accept} method is based on the * {@link ConcurrentMetadataStore#putIfAbsent} return value. {@code true} is returned @@ -52,23 +52,38 @@ public class MetadataStoreSelector implements MessageSelector { private final ConcurrentMetadataStore metadataStore; - private final MetadataKeyStrategy keyStrategy; + private final MessageProcessor keyStrategy; - public MetadataStoreSelector(MetadataKeyStrategy keyStrategy) { - this(keyStrategy, new SimpleMetadataStore()); + private final MessageProcessor valueStrategy; + + public MetadataStoreSelector(MessageProcessor keyStrategy) { + this(keyStrategy, (MessageProcessor) null); } - public MetadataStoreSelector(MetadataKeyStrategy keyStrategy, ConcurrentMetadataStore metadataStore) { - Assert.notNull(metadataStore); + public MetadataStoreSelector(MessageProcessor keyStrategy, MessageProcessor valueStrategy) { + this(keyStrategy, valueStrategy, new SimpleMetadataStore()); + } + + public MetadataStoreSelector(MessageProcessor keyStrategy, ConcurrentMetadataStore metadataStore) { + this(keyStrategy, null, metadataStore); + } + + public MetadataStoreSelector(MessageProcessor keyStrategy, MessageProcessor valueStrategy, + ConcurrentMetadataStore metadataStore) { Assert.notNull(keyStrategy); + Assert.notNull(metadataStore); this.metadataStore = metadataStore; this.keyStrategy = keyStrategy; + this.valueStrategy = valueStrategy; } + @Override public boolean accept(Message message) { - String key = this.keyStrategy.getKey(message); - String value = Long.toString(message.getHeaders().getTimestamp()); + String key = this.keyStrategy.processMessage(message); + String value = (this.valueStrategy != null) + ? this.valueStrategy.processMessage(message) + : Long.toString(message.getHeaders().getTimestamp()); return this.metadataStore.putIfAbsent(key, value) == null; } diff --git a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-4.1.xsd b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-4.1.xsd index bde78fa5f9..5948b66f66 100644 --- a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-4.1.xsd +++ b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-4.1.xsd @@ -4364,11 +4364,11 @@ The list of component name patterns you want to track (e.g., tracked-components - + - The 'MetadataKeyStrategy' reference. Used by the underlying + A 'MessageProcessor' reference. Used by the underlying 'org.springframework.integration.selector.MetadataStoreSelector'. Evaluates an 'idempotentKey' from the request Message. Mutually exclusive with 'selector' and 'key-expression'. @@ -4378,13 +4378,41 @@ The list of component name patterns you want to track (e.g., tracked-components + + + + + + + + + A 'MessageProcessor' reference. Used by the underlying + 'org.springframework.integration.selector.MetadataStoreSelector'. + Evaluates a 'value' for the 'idempotentKey' from the request Message. + Mutually exclusive with 'selector' and 'value-expression'. + By default, the 'MetadataStoreSelector' uses the 'timestamp' message header as the Metadata 'value'. + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IdempotentReceiverParserTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IdempotentReceiverParserTests-context.xml index 9a909d1b8e..8800c9532d 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IdempotentReceiverParserTests-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IdempotentReceiverParserTests-context.xml @@ -16,11 +16,19 @@ - + - + + + + + @@ -31,7 +39,7 @@ + metadata-store="store" + key-expression="headers.foo"/> diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IdempotentReceiverParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IdempotentReceiverParserTests.java index 9bc3560dfc..9592b706b4 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IdempotentReceiverParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/IdempotentReceiverParserTests.java @@ -46,9 +46,9 @@ import org.springframework.context.support.GenericApplicationContext; import org.springframework.core.io.ClassPathResource; import org.springframework.core.io.InputStreamResource; import org.springframework.integration.core.MessageSelector; +import org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor; +import org.springframework.integration.handler.MessageProcessor; import org.springframework.integration.handler.advice.IdempotentReceiverInterceptor; -import org.springframework.integration.metadata.ExpressionMetadataKeyStrategy; -import org.springframework.integration.metadata.MetadataKeyStrategy; import org.springframework.integration.metadata.MetadataStore; import org.springframework.integration.selector.MetadataStoreSelector; import org.springframework.messaging.MessageChannel; @@ -79,7 +79,10 @@ public class IdempotentReceiverParserTests { private IdempotentReceiverInterceptor strategyInterceptor; @Autowired - private MetadataKeyStrategy keyStrategy; + private MessageProcessor keyStrategy; + + @Autowired + private MessageProcessor valueStrategy; @Autowired @Qualifier("nullChannel") @@ -113,6 +116,7 @@ public class IdempotentReceiverParserTests { Object messageSelector = getPropertyValue(this.strategyInterceptor, "messageSelector"); assertThat(messageSelector, instanceOf(MetadataStoreSelector.class)); assertSame(this.keyStrategy, getPropertyValue(messageSelector, "keyStrategy")); + assertSame(this.valueStrategy, getPropertyValue(messageSelector, "valueStrategy")); @SuppressWarnings("unchecked") Map> idempotentEndpoints = (Map>) getPropertyValue(this.idempotentReceiverAutoProxyCreator, @@ -129,7 +133,7 @@ public class IdempotentReceiverParserTests { assertThat(messageSelector, instanceOf(MetadataStoreSelector.class)); assertSame(this.store, getPropertyValue(messageSelector, "metadataStore")); Object keyStrategy = getPropertyValue(messageSelector, "keyStrategy"); - assertThat(keyStrategy, instanceOf(ExpressionMetadataKeyStrategy.class)); + assertThat(keyStrategy, instanceOf(ExpressionEvaluatingMessageProcessor.class)); assertThat(keyStrategy.toString(), containsString("headers.foo")); @SuppressWarnings("unchecked") Map> idempotentEndpoints = @@ -176,40 +180,66 @@ public class IdempotentReceiverParserTests { catch (BeanDefinitionParsingException e) { assertThat(e.getMessage(), containsString("The 'selector' attribute is mutually exclusive with 'metadata-store', " + - "'key-strategy' or 'key-expression'")); + "'key-strategy', 'key-expression', 'value-strategy' or 'value-expression'")); } } @Test - public void testSelectorAndStrategy() throws Exception { + public void testSelectorAndKeyStrategy() throws Exception { try { - bootStrap("selector-and-strategy"); + bootStrap("selector-and-key-strategy"); fail("BeanDefinitionParsingException expected"); } catch (BeanDefinitionParsingException e) { assertThat(e.getMessage(), containsString("The 'selector' attribute is mutually exclusive with 'metadata-store', " + - "'key-strategy' or 'key-expression'")); + "'key-strategy', 'key-expression', 'value-strategy' or 'value-expression'")); } } @Test - public void testSelectorAndExpression() throws Exception { + public void testSelectorAndKeyExpression() throws Exception { try { - bootStrap("selector-and-expression"); + bootStrap("selector-and-key-expression"); fail("BeanDefinitionParsingException expected"); } catch (BeanDefinitionParsingException e) { assertThat(e.getMessage(), containsString("The 'selector' attribute is mutually exclusive with 'metadata-store', " + - "'key-strategy' or 'key-expression'")); + "'key-strategy', 'key-expression', 'value-strategy' or 'value-expression'")); } } @Test - public void testStrategyAndExpression() throws Exception { + public void testSelectorAndValueStrategy() throws Exception { try { - bootStrap("strategy-and-expression"); + bootStrap("selector-and-value-strategy"); + fail("BeanDefinitionParsingException expected"); + } + catch (BeanDefinitionParsingException e) { + assertThat(e.getMessage(), + containsString("The 'selector' attribute is mutually exclusive with 'metadata-store', " + + "'key-strategy', 'key-expression', 'value-strategy' or 'value-expression'")); + } + } + + @Test + public void testSelectorAndValueExpression() throws Exception { + try { + bootStrap("selector-and-value-expression"); + fail("BeanDefinitionParsingException expected"); + } + catch (BeanDefinitionParsingException e) { + assertThat(e.getMessage(), + containsString("The 'selector' attribute is mutually exclusive with 'metadata-store', " + + "'key-strategy', 'key-expression', 'value-strategy' or 'value-expression'")); + } + } + + @Test + public void testKeyStrategyAndKeyExpression() throws Exception { + try { + bootStrap("key-strategy-and-key-expression"); fail("BeanDefinitionParsingException expected"); } catch (BeanDefinitionParsingException e) { @@ -218,6 +248,18 @@ public class IdempotentReceiverParserTests { } } + @Test + public void testValueStrategyAndValueExpression() throws Exception { + try { + bootStrap("value-strategy-and-value-expression"); + fail("BeanDefinitionParsingException expected"); + } + catch (BeanDefinitionParsingException e) { + assertThat(e.getMessage(), + containsString("The 'value-strategy' and 'value-expression' attributes are mutually exclusive")); + } + } + private ApplicationContext bootStrap(String configProperty) throws Exception { PropertiesFactoryBean pfb = new PropertiesFactoryBean(); pfb.setLocation(new ClassPathResource( diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/idempotent-receiver-configs.properties b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/idempotent-receiver-configs.properties index 0b05605ba6..bef7d6af06 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/idempotent-receiver-configs.properties +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/idempotent-receiver-configs.properties @@ -13,8 +13,17 @@ without-endpoint= selector-and-store= -selector-and-strategy= +selector-and-key-strategy= -selector-and-expression= +selector-and-key-expression= -strategy-and-expression= +selector-and-value-strategy= + +selector-and-value-expression= + +key-strategy-and-key-expression= + +value-strategy-and-value-expression= diff --git a/spring-integration-core/src/test/java/org/springframework/integration/handler/advice/IdempotentReceiverTests.java b/spring-integration-core/src/test/java/org/springframework/integration/handler/advice/IdempotentReceiverTests.java index 00f9baa3e5..1d7267c71c 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/handler/advice/IdempotentReceiverTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/handler/advice/IdempotentReceiverTests.java @@ -33,16 +33,14 @@ import org.mockito.Mockito; import org.springframework.aop.framework.ProxyFactory; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.MessageRejectedException; +import org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor; import org.springframework.integration.metadata.ConcurrentMetadataStore; -import org.springframework.integration.metadata.ExpressionMetadataKeyStrategy; import org.springframework.integration.metadata.MetadataStore; import org.springframework.integration.metadata.SimpleMetadataStore; import org.springframework.integration.selector.MetadataStoreSelector; -import org.springframework.integration.support.DefaultMessageBuilderFactory; -import org.springframework.integration.support.MessageBuilderFactory; -import org.springframework.integration.support.utils.IntegrationUtils; import org.springframework.integration.test.util.TestUtils; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; @@ -86,7 +84,8 @@ public class IdempotentReceiverTests { @Test public void testIdempotentReceiverInterceptor() { ConcurrentMetadataStore store = new SimpleMetadataStore(); - ExpressionMetadataKeyStrategy idempotentKeyStrategy = new ExpressionMetadataKeyStrategy("payload"); + ExpressionEvaluatingMessageProcessor idempotentKeyStrategy = + new ExpressionEvaluatingMessageProcessor<>(new SpelExpressionParser().parseExpression("payload")); BeanFactory beanFactory = Mockito.mock(BeanFactory.class); idempotentKeyStrategy.setBeanFactory(beanFactory); IdempotentReceiverInterceptor idempotentReceiverInterceptor = diff --git a/spring-integration-jmx/src/test/java/org/springframework/integration/monitor/IdempotentReceiverIntegrationTests.java b/spring-integration-jmx/src/test/java/org/springframework/integration/monitor/IdempotentReceiverIntegrationTests.java index 5a2b298b25..8c8ce38311 100644 --- a/spring-integration-jmx/src/test/java/org/springframework/integration/monitor/IdempotentReceiverIntegrationTests.java +++ b/spring-integration-jmx/src/test/java/org/springframework/integration/monitor/IdempotentReceiverIntegrationTests.java @@ -23,6 +23,8 @@ import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; +import java.util.ArrayList; +import java.util.List; import java.util.Map; import java.util.concurrent.atomic.AtomicInteger; @@ -36,14 +38,15 @@ import org.springframework.context.annotation.Configuration; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.MessageRejectedException; import org.springframework.integration.annotation.IdempotentReceiver; +import org.springframework.integration.annotation.ServiceActivator; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.config.EnableIntegration; +import org.springframework.integration.handler.MessageProcessor; import org.springframework.integration.handler.advice.AbstractRequestHandlerAdvice; import org.springframework.integration.handler.advice.IdempotentReceiverInterceptor; import org.springframework.integration.jmx.config.EnableIntegrationMBeanExport; import org.springframework.integration.metadata.ConcurrentMetadataStore; -import org.springframework.integration.metadata.MetadataKeyStrategy; import org.springframework.integration.metadata.MetadataStore; import org.springframework.integration.metadata.SimpleMetadataStore; import org.springframework.integration.selector.MetadataStoreSelector; @@ -54,6 +57,7 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.PollableChannel; import org.springframework.messaging.support.GenericMessage; +import org.springframework.stereotype.Component; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @@ -82,15 +86,24 @@ public class IdempotentReceiverIntegrationTests { @Autowired private AtomicInteger adviceCalled; + @Autowired + private MessageChannel annotatedMethodChannel; + + @Autowired + private FooService fooService; + @Test public void testIdempotentReceiver() { + this.idempotentReceiverInterceptor.setThrowExceptionOnRejection(true); + TestUtils.getPropertyValue(this.store, "metadata", Map.class).clear(); Message message = new GenericMessage("foo"); this.input.send(message); Message receive = this.output.receive(10000); assertNotNull(receive); assertEquals(1, this.adviceCalled.get()); assertEquals(1, TestUtils.getPropertyValue(this.store, "metadata", Map.class).size()); - assertNotNull(this.store.get("foo")); + String foo = this.store.get("foo"); + assertEquals("FOO", foo); try { this.input.send(message); @@ -108,6 +121,18 @@ public class IdempotentReceiverIntegrationTests { assertEquals(1, TestUtils.getPropertyValue(store, "metadata", Map.class).size()); } + @Test + public void testIdempotentReceiverOnMethod() { + TestUtils.getPropertyValue(this.store, "metadata", Map.class).clear(); + Message message = new GenericMessage("foo"); + this.annotatedMethodChannel.send(message); + this.annotatedMethodChannel.send(message); + + assertEquals(2, this.fooService.messages.size()); + assertTrue(this.fooService.messages.get(1).getHeaders().get(IntegrationMessageHeaderAccessor.DUPLICATE_MESSAGE, + Boolean.class)); + } + @Configuration @EnableIntegration @EnableIntegrationMBeanExport(server = "mBeanServer") @@ -125,17 +150,21 @@ public class IdempotentReceiverIntegrationTests { @Bean public IdempotentReceiverInterceptor idempotentReceiverInterceptor() { - IdempotentReceiverInterceptor idempotentReceiverInterceptor = - new IdempotentReceiverInterceptor(new MetadataStoreSelector(new MetadataKeyStrategy() { + return new IdempotentReceiverInterceptor(new MetadataStoreSelector(new MessageProcessor() { @Override - public String getKey(Message message) { + public String processMessage(Message message) { return message.getPayload().toString(); } + }, new MessageProcessor() { + + @Override + public String processMessage(Message message) { + return message.getPayload().toString().toUpperCase(); + } + }, store())); - idempotentReceiverInterceptor.setThrowExceptionOnRejection(true); - return idempotentReceiverInterceptor; } @Bean @@ -182,6 +211,29 @@ public class IdempotentReceiverIntegrationTests { }; } + @Bean + public MessageChannel annotatedMethodChannel() { + return new DirectChannel(); + } + + @Bean + public FooService fooService() { + return new FooService(); + } + + } + + @Component + private static class FooService { + + private List> messages = new ArrayList>(); + + @ServiceActivator(inputChannel = "annotatedMethodChannel") + @IdempotentReceiver("idempotentReceiverInterceptor") + public void handle(Message message) { + this.messages.add(message); + } + } } diff --git a/src/reference/docbook/handler-advice.xml b/src/reference/docbook/handler-advice.xml index 108d489104..66af59ee3e 100644 --- a/src/reference/docbook/handler-advice.xml +++ b/src/reference/docbook/handler-advice.xml @@ -616,12 +616,13 @@ public class MyAdvisedFilter { To maintain state between messages and provide the ability to compare messages for the idempotency, the MetadataStoreSelector is provided. It accepts a - MetadataKeyStrategy implementation (which creates a lookup key + MessageProcessor implementation (which creates a lookup key based on the Message) and an optional ConcurrentMetadataStore (). - See the MetadataStoreSelector JavaDocs for more information. An - ExpressionMetadataKeyStrategy implementation is provided, allowing - simple SpEL expressions to be used to determine the key from the message. + See the MetadataStoreSelector JavaDocs for more information. + The value for ConcurrentMetadataStore also can be customized + using additional MessageProcessor. By default + MetadataStoreSelector uses timestamp message header. For convenience, the MetadataStoreSelector options are configurable directly on @@ -635,7 +636,9 @@ public class MyAdvisedFilter { metadata-store="" ]]> ]]> + value-strategy="" ]]> ]]> @@ -658,7 +661,9 @@ public class MyAdvisedFilter { A MessageSelector bean reference. Mutually exclusive with metadata-store and - key-strategy (key-expression). + key-strategy (key-expression). When selector + is not provided, one of key-strategy or key-strategy-expression + is required. @@ -681,23 +686,50 @@ public class MyAdvisedFilter { - A MetadataKeyStrategy reference. Used by the underlying + A MessageProcessor reference. Used by the underlying MetadataStoreSelector. Evaluates an idempotentKey from the request Message. Mutually exclusive with selector and key-expression. + When a selector + is not provided, one of key-strategy or key-strategy-expression + is required. - A SpEL expression to populate an ExpressionMetadataKeyStrategy. + A SpEL expression to populate an ExpressionEvaluatingMessageProcessor. Used by the underlying MetadataStoreSelector. Evaluates an idempotentKey using the request Message as the evaluation context root object. Mutually exclusive with selector and key-strategy. + When a selector + is not provided, one of key-strategy or key-strategy-expression + is required. + + A MessageProcessor reference. Used by the underlying + MetadataStoreSelector. + Evaluates a value for the idempotentKey from the request Message. + Mutually exclusive with selector and value-expression. + By default, the 'MetadataStoreSelector' uses the 'timestamp' message header as the Metadata 'value'. + + + + + + A SpEL expression to populate an ExpressionEvaluatingMessageProcessor. + Used by the underlying MetadataStoreSelector. + Evaluates a value for the idempotentKey using the request Message + as the evaluation context root object. + Mutually exclusive with selector and value-strategy. + By default, the 'MetadataStoreSelector' uses the 'timestamp' message header as the Metadata 'value'. + + + + Throw an exception if the IdempotentReceiverInterceptor rejects the message defaults to false. @@ -707,7 +739,7 @@ public class MyAdvisedFilter { For Java configuration, the method level IdempotentReceiver annotation is provided. It - is used to mark a @Bean that has a Messaging annotation (@ServiceActivator, + is used to mark a method that has a Messaging annotation (@ServiceActivator, @Router etc.) to specify which IdempotentReceiverInterceptors will be applied to this endpoint: