From 5cca8e8e018c3a5dac4315b81b4ffce6cd01bfc3 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 2 Nov 2016 19:32:16 -0400 Subject: [PATCH] INT-3770: Add TX Support from Mid-flow JIRA: https://jira.spring.io/browse/INT-3770, https://jira.spring.io/browse/INT-4107 Having `TransactionHandleMessageAdvice` we can start TX from any `MessageHandler.handleMessage()` * Add `` alongside with the `` for those components which produce reply * Merge `` and `` configuration to a single `ManagedList` * Rework JPA `` in favor of common solution * Some polishing and refactoring AbstractPollingEndpoint: avoid `new ArrayList` if we don't have `receiveOnlyAdvice`s --- .../xml/AbstractConsumerEndpointParser.java | 14 +++- .../AbstractOutboundChannelAdapterParser.java | 7 +- .../integration/config/xml/DelayerParser.java | 4 + .../config/xml/IntegrationNamespaceUtils.java | 63 ++++++++++---- .../endpoint/AbstractPollingEndpoint.java | 24 +++--- .../endpoint/SourcePollingChannelAdapter.java | 5 +- .../config/spring-integration-5.0.xsd | 4 + ....xml => AggregatorParserTests-context.xml} | 6 +- .../config/AggregatorParserTests.java | 31 +++---- .../config/spring-integration-file-5.0.xsd | 2 + .../ftp/config/spring-integration-ftp-5.0.xsd | 1 + .../config/spring-integration-groovy-5.0.xsd | 1 + .../config/spring-integration-http-5.0.xsd | 1 + .../ip/config/spring-integration-ip-5.0.xsd | 1 + .../jms/config/spring-integration-jms-5.0.xsd | 1 + .../jmx/config/spring-integration-jmx-5.0.xsd | 1 + .../IdempotentReceiverIntegrationTests.java | 79 ++++++++++++++---- .../src/test/resources/log4j.properties | 6 +- .../xml/AbstractJpaOutboundGatewayParser.java | 13 --- .../xml/JpaOutboundChannelAdapterParser.java | 12 --- .../JpaOutboundGatewayFactoryBean.java | 32 ------- .../xml/JpaMessageHandlerParserTests.java | 20 +++-- .../xml/JpaOutboundGatewayParserTests.java | 3 +- .../xml/JpaOutboundGatewayParserTests.xml | 4 +- ...JpaOutboundChannelAdapterTests-context.xml | 22 ++--- .../config/spring-integration-redis-5.0.xsd | 2 + .../rmi/config/spring-integration-rmi-5.0.xsd | 1 + .../rmi/BackToBackTests-context.xml | 9 +- .../integration/rmi/BackToBackTests.java | 25 ++++-- .../src/test/resources/.svnignore | 0 .../src/test/resources/log4j.properties | 8 ++ .../config/spring-integration-sftp-5.0.xsd | 1 + .../src/test/resources/log4j.properties | 6 +- .../config/spring-integration-twitter-5.0.xsd | 1 + .../ws/config/spring-integration-ws-5.0.xsd | 1 + src/reference/asciidoc/handler-advice.adoc | 83 +++++++++++++++++-- src/reference/asciidoc/whats-new.adoc | 3 + 37 files changed, 327 insertions(+), 170 deletions(-) rename spring-integration-core/src/test/java/org/springframework/integration/config/{aggregatorParserTests.xml => AggregatorParserTests-context.xml} (97%) delete mode 100644 spring-integration-rmi/src/test/resources/.svnignore create mode 100644 spring-integration-rmi/src/test/resources/log4j.properties diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractConsumerEndpointParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractConsumerEndpointParser.java index 36cde97540..2d1b564fdd 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractConsumerEndpointParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractConsumerEndpointParser.java @@ -29,6 +29,7 @@ import org.springframework.beans.factory.parsing.BeanComponentDefinition; import org.springframework.beans.factory.support.AbstractBeanDefinition; import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.support.BeanDefinitionReaderUtils; +import org.springframework.beans.factory.support.ManagedList; import org.springframework.beans.factory.support.ManagedSet; import org.springframework.beans.factory.xml.AbstractBeanDefinitionParser; import org.springframework.beans.factory.xml.ParserContext; @@ -86,11 +87,18 @@ public abstract class AbstractConsumerEndpointParser extends AbstractBeanDefinit IntegrationNamespaceUtils.setReferenceIfAttributeDefined(handlerBuilder, element, "output-channel"); IntegrationNamespaceUtils.setValueIfAttributeDefined(handlerBuilder, element, "order"); + Element txElement = DomUtils.getChildElementByTagName(element, "transactional"); Element adviceChainElement = DomUtils.getChildElementByTagName(element, IntegrationNamespaceUtils.REQUEST_HANDLER_ADVICE_CHAIN); - IntegrationNamespaceUtils.configureAndSetAdviceChainIfPresent(adviceChainElement, null, + + @SuppressWarnings("rawtypes") + ManagedList adviceChain = IntegrationNamespaceUtils.configureAdviceChain(adviceChainElement, txElement, true, handlerBuilder.getRawBeanDefinition(), parserContext); + if (!CollectionUtils.isEmpty(adviceChain)) { + handlerBuilder.addPropertyValue("adviceChain", adviceChain); + } + AbstractBeanDefinition handlerBeanDefinition = handlerBuilder.getBeanDefinition(); String inputChannelAttributeName = this.getInputChannelAttributeName(); boolean hasInputChannelAttribute = element.hasAttribute(inputChannelAttributeName); @@ -121,6 +129,10 @@ public abstract class AbstractConsumerEndpointParser extends AbstractBeanDefinit BeanDefinitionBuilder builder = BeanDefinitionBuilder.genericBeanDefinition(ConsumerEndpointFactoryBean.class); + if (!CollectionUtils.isEmpty(adviceChain)) { + builder.addPropertyValue("adviceChain", adviceChain); + } + String handlerBeanName = BeanDefinitionReaderUtils.generateBeanName(handlerBeanDefinition, parserContext.getRegistry()); String[] handlerAlias = IntegrationNamespaceUtils.generateAlias(element); parserContext.registerBeanComponent(new BeanComponentDefinition(handlerBeanDefinition, handlerBeanName, handlerAlias)); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractOutboundChannelAdapterParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractOutboundChannelAdapterParser.java index 053b6c5f2c..7c6eeea2dd 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractOutboundChannelAdapterParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/AbstractOutboundChannelAdapterParser.java @@ -27,6 +27,7 @@ import org.springframework.beans.factory.support.ManagedList; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.ConsumerEndpointFactoryBean; import org.springframework.integration.handler.AbstractReplyProducingMessageHandler; +import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; import org.springframework.util.xml.DomUtils; @@ -81,12 +82,14 @@ public abstract class AbstractOutboundChannelAdapterParser extends AbstractChann private void configureRequestHandlerAdviceChain(Element element, ParserContext parserContext, BeanDefinition handlerBeanDefinition, BeanDefinitionBuilder consumerBuilder) { + Element txElement = DomUtils.getChildElementByTagName(element, "transactional"); Element adviceChainElement = DomUtils.getChildElementByTagName(element, IntegrationNamespaceUtils.REQUEST_HANDLER_ADVICE_CHAIN); @SuppressWarnings("rawtypes") ManagedList adviceChain = - IntegrationNamespaceUtils.configureAdviceChain(adviceChainElement, null, handlerBeanDefinition, parserContext); - if (adviceChain != null) { + IntegrationNamespaceUtils.configureAdviceChain(adviceChainElement, txElement, handlerBeanDefinition, + parserContext); + if (!CollectionUtils.isEmpty(adviceChain)) { /* * For ARPMH, the advice chain is injected so just the handleRequestMessage method is advised. * Sometime ARPMHs do double duty as a gateway and a channel adapter. The parser subclass diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/DelayerParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/DelayerParser.java index 7dc1650a16..4cf4f8e87d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/DelayerParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/DelayerParser.java @@ -101,6 +101,10 @@ public class DelayerParser extends AbstractConsumerEndpointParser { IntegrationNamespaceUtils.configureAndSetAdviceChainIfPresent(adviceChainElement, txElement, builder.getRawBeanDefinition(), parserContext, "delayedAdviceChain"); + if (txElement != null) { + element.removeChild(txElement); + } + return builder; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceUtils.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceUtils.java index 1eb167d7b3..3d6b08f049 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceUtils.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/IntegrationNamespaceUtils.java @@ -48,10 +48,12 @@ import org.springframework.integration.config.FixedSubscriberChannelBeanFactoryP import org.springframework.integration.config.IntegrationConfigUtils; import org.springframework.integration.context.IntegrationContextUtils; import org.springframework.integration.endpoint.AbstractPollingEndpoint; +import org.springframework.integration.transaction.TransactionHandleMessageAdvice; import org.springframework.transaction.interceptor.DefaultTransactionAttribute; import org.springframework.transaction.interceptor.MatchAlwaysTransactionAttributeSource; import org.springframework.transaction.interceptor.TransactionInterceptor; import org.springframework.util.Assert; +import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; import org.springframework.util.xml.DomUtils; @@ -381,19 +383,34 @@ public abstract class IntegrationNamespaceUtils { * Parse a "transactional" element and configure a {@link TransactionInterceptor} * with "transactionManager" and other "transactionDefinition" properties. * For example, this advisor will be applied on the Polling Task proxy. - * * @param txElement The transactional element. * @return The bean definition. - * * @see AbstractPollingEndpoint */ public static BeanDefinition configureTransactionAttributes(Element txElement) { + return configureTransactionAttributes(txElement, false); + } + + /** + * Parse a "transactional" element and configure a {@link TransactionInterceptor} + * or {@link TransactionHandleMessageAdvice} + * with "transactionManager" and other "transactionDefinition" properties. + * For example, this advisor will be applied on the Polling Task proxy. + * @param txElement The transactional element. + * @param handleMessageAdvice flag if to use {@link TransactionHandleMessageAdvice} + * or regular {@link TransactionInterceptor} + * @return The bean definition. + * @see AbstractPollingEndpoint + */ + public static BeanDefinition configureTransactionAttributes(Element txElement, boolean handleMessageAdvice) { BeanDefinition txDefinition = configureTransactionDefinition(txElement); BeanDefinitionBuilder attributeSourceBuilder = BeanDefinitionBuilder.genericBeanDefinition(MatchAlwaysTransactionAttributeSource.class); attributeSourceBuilder.addPropertyValue("transactionAttribute", txDefinition); BeanDefinitionBuilder txInterceptorBuilder = - BeanDefinitionBuilder.genericBeanDefinition(TransactionInterceptor.class); + BeanDefinitionBuilder.genericBeanDefinition(handleMessageAdvice + ? TransactionHandleMessageAdvice.class + : TransactionInterceptor.class); txInterceptorBuilder.addPropertyReference("transactionManager", txElement.getAttribute("transaction-manager")); txInterceptorBuilder.addPropertyValue("transactionAttributeSource", attributeSourceBuilder.getBeanDefinition()); return txInterceptorBuilder.getBeanDefinition(); @@ -426,31 +443,47 @@ public abstract class IntegrationNamespaceUtils { public static void configureAndSetAdviceChainIfPresent(Element adviceChainElement, Element txElement, BeanDefinition parentBeanDefinition, ParserContext parserContext) { - configureAndSetAdviceChainIfPresent(adviceChainElement, txElement, parentBeanDefinition, parserContext, - "adviceChain"); + configureAndSetAdviceChainIfPresent(adviceChainElement, txElement, false, parentBeanDefinition, parserContext); + } + + public static void configureAndSetAdviceChainIfPresent(Element adviceChainElement, + Element txElement, boolean handleMessageAdvice, BeanDefinition parentBeanDefinition, + ParserContext parserContext) { + configureAndSetAdviceChainIfPresent(adviceChainElement, txElement, handleMessageAdvice, + parentBeanDefinition, parserContext, "adviceChain"); + } + + public static void configureAndSetAdviceChainIfPresent(Element adviceChainElement, Element txElement, + BeanDefinition parentBeanDefinition, ParserContext parserContext, String propertyName) { + configureAndSetAdviceChainIfPresent(adviceChainElement, txElement, false, parentBeanDefinition, + parserContext, propertyName); } @SuppressWarnings({ "rawtypes" }) public static void configureAndSetAdviceChainIfPresent(Element adviceChainElement, Element txElement, - BeanDefinition parentBeanDefinition, ParserContext parserContext, String propertyName) { - ManagedList adviceChain = configureAdviceChain(adviceChainElement, txElement, parentBeanDefinition, - parserContext); - if (adviceChain != null) { + boolean handleMessageAdvice, BeanDefinition parentBeanDefinition, ParserContext parserContext, + String propertyName) { + ManagedList adviceChain = configureAdviceChain(adviceChainElement, txElement, handleMessageAdvice, + parentBeanDefinition, parserContext); + if (!CollectionUtils.isEmpty(adviceChain)) { parentBeanDefinition.getPropertyValues().add(propertyName, adviceChain); } } - @SuppressWarnings({ "rawtypes", "unchecked" }) + @SuppressWarnings("rawtypes") public static ManagedList configureAdviceChain(Element adviceChainElement, Element txElement, BeanDefinition parentBeanDefinition, ParserContext parserContext) { - ManagedList adviceChain = null; - // Schema validation ensures txElement and adviceChainElement are mutually exclusive + return configureAdviceChain(adviceChainElement, txElement, false, parentBeanDefinition, parserContext); + } + + @SuppressWarnings({ "rawtypes", "unchecked" }) + public static ManagedList configureAdviceChain(Element adviceChainElement, Element txElement, + boolean handleMessageAdvice, BeanDefinition parentBeanDefinition, ParserContext parserContext) { + ManagedList adviceChain = new ManagedList(); if (txElement != null) { - adviceChain = new ManagedList(); - adviceChain.add(IntegrationNamespaceUtils.configureTransactionAttributes(txElement)); + adviceChain.add(configureTransactionAttributes(txElement, handleMessageAdvice)); } if (adviceChainElement != null) { - adviceChain = new ManagedList(); NodeList childNodes = adviceChainElement.getChildNodes(); for (int i = 0; i < childNodes.getLength(); i++) { Node child = childNodes.item(i); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java index 1b2fb85d79..6ca8e29e5e 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/AbstractPollingEndpoint.java @@ -16,12 +16,12 @@ package org.springframework.integration.endpoint; -import java.util.ArrayList; import java.util.Collection; import java.util.List; import java.util.concurrent.Callable; import java.util.concurrent.Executor; import java.util.concurrent.ScheduledFuture; +import java.util.stream.Collectors; import org.aopalliance.aop.Advice; @@ -174,30 +174,26 @@ public abstract class AbstractPollingEndpoint extends AbstractEndpoint implement @SuppressWarnings("unchecked") private Runnable createPoller() throws Exception { - List receiveOnlyAdviceChain = new ArrayList(); + List receiveOnlyAdviceChain = null; if (!CollectionUtils.isEmpty(this.adviceChain)) { - for (Advice advice : this.adviceChain) { - if (isReceiveOnlyAdvice(advice)) { - receiveOnlyAdviceChain.add(advice); - } - } + receiveOnlyAdviceChain = this.adviceChain.stream() + .filter(this::isReceiveOnlyAdvice) + .collect(Collectors.toList()); } - Callable pollingTask = () -> doPoll(); + Callable pollingTask = this::doPoll; List adviceChain = this.adviceChain; if (!CollectionUtils.isEmpty(adviceChain)) { ProxyFactory proxyFactory = new ProxyFactory(pollingTask); if (!CollectionUtils.isEmpty(adviceChain)) { - for (Advice advice : adviceChain) { - if (!isReceiveOnlyAdvice(advice)) { - proxyFactory.addAdvice(advice); - } - } + adviceChain.stream() + .filter(advice -> !isReceiveOnlyAdvice(advice)) + .forEach(proxyFactory::addAdvice); } pollingTask = (Callable) proxyFactory.getProxy(this.beanClassLoader); } - if (receiveOnlyAdviceChain.size() > 0) { + if (receiveOnlyAdviceChain != null) { applyReceiveOnlyAdviceChain(receiveOnlyAdviceChain); } return new Poller(pollingTask); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java index 22887b405f..65b93a8d58 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/SourcePollingChannelAdapter.java @@ -132,9 +132,10 @@ public class SourcePollingChannelAdapter extends AbstractPollingEndpoint @Override protected void applyReceiveOnlyAdviceChain(Collection chain) { if (AopUtils.isAopProxy(this.source)) { - this.appliedAdvices.forEach(((Advised) this.source)::removeAdvice); + Advised source = (Advised) this.source; + this.appliedAdvices.forEach(source::removeAdvice); for (Advice advice : chain) { - ((Advised) this.source).addAdvisor(adviceToReceiveAdvisor(advice)); + source.addAdvisor(adviceToReceiveAdvisor(advice)); } } else { diff --git a/spring-integration-core/src/main/resources/org/springframework/integration/config/spring-integration-5.0.xsd b/spring-integration-core/src/main/resources/org/springframework/integration/config/spring-integration-5.0.xsd index 3241cbe0b3..8c12e78fb0 100644 --- a/spring-integration-core/src/main/resources/org/springframework/integration/config/spring-integration-5.0.xsd +++ b/spring-integration-core/src/main/resources/org/springframework/integration/config/spring-integration-5.0.xsd @@ -1367,6 +1367,7 @@ + @@ -1665,6 +1666,7 @@ + @@ -2871,6 +2873,7 @@ + @@ -4125,6 +4128,7 @@ + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests-context.xml similarity index 97% rename from spring-integration-core/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml rename to spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests-context.xml index 5029a63130..0f4a3d717a 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests-context.xml @@ -27,6 +27,9 @@ input-channel="aggregatorWithCustomMGPReferenceInput" output-channel="outputChannel"/> + + + - + @@ -116,4 +119,5 @@ class="org.springframework.integration.config.MaxValueReleaseStrategy"> + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java index a7e95c6697..6b29d601c1 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java @@ -34,11 +34,12 @@ import java.util.List; import java.util.concurrent.atomic.AtomicReference; import org.junit.Assert; -import org.junit.Before; import org.junit.Test; +import org.junit.runner.RunWith; import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.BeanCreationException; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.parsing.BeanDefinitionParsingException; import org.springframework.context.ApplicationContext; import org.springframework.context.support.ClassPathXmlApplicationContext; @@ -63,6 +64,7 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.PollableChannel; import org.springframework.messaging.SubscribableChannel; +import org.springframework.test.context.junit4.SpringRunner; /** * @author Marius Bogoevici @@ -73,15 +75,12 @@ import org.springframework.messaging.SubscribableChannel; * @author Gunnar Hillert * @author Gary Russell */ +@RunWith(SpringRunner.class) public class AggregatorParserTests { + @Autowired private ApplicationContext context; - @Before - public void setUp() { - this.context = new ClassPathXmlApplicationContext("aggregatorParserTests.xml", this.getClass()); - } - @Test public void testAggregation() { MessageChannel input = (MessageChannel) context.getBean("aggregatorWithReferenceInput"); @@ -90,11 +89,11 @@ public class AggregatorParserTests { outboundMessages.add(createMessage("123", "id1", 3, 1, null)); outboundMessages.add(createMessage("789", "id1", 3, 3, null)); outboundMessages.add(createMessage("456", "id1", 3, 2, null)); - for (Message message : outboundMessages) { - input.send(message); - } - assertEquals("One and only one message must have been aggregated", 1, aggregatorBean.getAggregatedMessages() - .size()); + + outboundMessages.forEach(input::send); + + assertEquals("One and only one message must have been aggregated", 1, + aggregatorBean.getAggregatedMessages().size()); Message aggregatedMessage = aggregatorBean.getAggregatedMessages().get("id1"); assertEquals("The aggregated message payload is not correct", "123456789", aggregatedMessage.getPayload()); Object mbf = context.getBean(IntegrationUtils.INTEGRATION_MESSAGE_BUILDER_FACTORY_BEAN_NAME); @@ -111,9 +110,9 @@ public class AggregatorParserTests { outboundMessages.add(createMessage("123", "id1", 3, 1, null)); outboundMessages.add(createMessage("789", "id1", 3, 3, null)); outboundMessages.add(createMessage("456", "id1", 3, 2, null)); - for (Message message : outboundMessages) { - input.send(message); - } + + outboundMessages.forEach(input::send); + assertEquals(3, output.getQueueSize()); output.purge(null); } @@ -127,7 +126,9 @@ public class AggregatorParserTests { outboundMessages.add(createMessage("123", "id1", 3, 1, null)); outboundMessages.add(createMessage("789", "id1", 3, 3, null)); outboundMessages.add(createMessage("456", "id1", 3, 2, null)); + outboundMessages.forEach(input::send); + assertEquals(3, output.getQueueSize()); output.purge(null); } @@ -142,7 +143,9 @@ public class AggregatorParserTests { outboundMessages.add(MessageBuilder.withPayload("123").setHeader("foo", "1").build()); outboundMessages.add(MessageBuilder.withPayload("456").setHeader("foo", "1").build()); outboundMessages.add(MessageBuilder.withPayload("789").setHeader("foo", "1").build()); + outboundMessages.forEach(input::send); + assertEquals("The aggregated message payload is not correct", "[123]", aggregatedMessage.get().getPayload() .toString()); Object mbf = context.getBean(IntegrationUtils.INTEGRATION_MESSAGE_BUILDER_FACTORY_BEAN_NAME); diff --git a/spring-integration-file/src/main/resources/org/springframework/integration/file/config/spring-integration-file-5.0.xsd b/spring-integration-file/src/main/resources/org/springframework/integration/file/config/spring-integration-file-5.0.xsd index dabda9e370..1d90625b6d 100644 --- a/spring-integration-file/src/main/resources/org/springframework/integration/file/config/spring-integration-file-5.0.xsd +++ b/spring-integration-file/src/main/resources/org/springframework/integration/file/config/spring-integration-file-5.0.xsd @@ -407,6 +407,7 @@ Only files matching this regular expression will be picked up by this adapter. + @@ -651,6 +652,7 @@ Only files matching this regular expression will be picked up by this adapter. + diff --git a/spring-integration-ftp/src/main/resources/org/springframework/integration/ftp/config/spring-integration-ftp-5.0.xsd b/spring-integration-ftp/src/main/resources/org/springframework/integration/ftp/config/spring-integration-ftp-5.0.xsd index ca817d438f..d7f45240f8 100644 --- a/spring-integration-ftp/src/main/resources/org/springframework/integration/ftp/config/spring-integration-ftp-5.0.xsd +++ b/spring-integration-ftp/src/main/resources/org/springframework/integration/ftp/config/spring-integration-ftp-5.0.xsd @@ -206,6 +206,7 @@ + diff --git a/spring-integration-groovy/src/main/resources/org/springframework/integration/groovy/config/spring-integration-groovy-5.0.xsd b/spring-integration-groovy/src/main/resources/org/springframework/integration/groovy/config/spring-integration-groovy-5.0.xsd index 4ab8dec99e..66861bfd4e 100644 --- a/spring-integration-groovy/src/main/resources/org/springframework/integration/groovy/config/spring-integration-groovy-5.0.xsd +++ b/spring-integration-groovy/src/main/resources/org/springframework/integration/groovy/config/spring-integration-groovy-5.0.xsd @@ -79,6 +79,7 @@ + diff --git a/spring-integration-http/src/main/resources/org/springframework/integration/http/config/spring-integration-http-5.0.xsd b/spring-integration-http/src/main/resources/org/springframework/integration/http/config/spring-integration-http-5.0.xsd index d8c49e3fc3..d636bb27d9 100644 --- a/spring-integration-http/src/main/resources/org/springframework/integration/http/config/spring-integration-http-5.0.xsd +++ b/spring-integration-http/src/main/resources/org/springframework/integration/http/config/spring-integration-http-5.0.xsd @@ -419,6 +419,7 @@ + diff --git a/spring-integration-ip/src/main/resources/org/springframework/integration/ip/config/spring-integration-ip-5.0.xsd b/spring-integration-ip/src/main/resources/org/springframework/integration/ip/config/spring-integration-ip-5.0.xsd index 33b00a26c6..80228410b9 100644 --- a/spring-integration-ip/src/main/resources/org/springframework/integration/ip/config/spring-integration-ip-5.0.xsd +++ b/spring-integration-ip/src/main/resources/org/springframework/integration/ip/config/spring-integration-ip-5.0.xsd @@ -374,6 +374,7 @@ + diff --git a/spring-integration-jms/src/main/resources/org/springframework/integration/jms/config/spring-integration-jms-5.0.xsd b/spring-integration-jms/src/main/resources/org/springframework/integration/jms/config/spring-integration-jms-5.0.xsd index b70340ddad..d5e13e2e8c 100644 --- a/spring-integration-jms/src/main/resources/org/springframework/integration/jms/config/spring-integration-jms-5.0.xsd +++ b/spring-integration-jms/src/main/resources/org/springframework/integration/jms/config/spring-integration-jms-5.0.xsd @@ -774,6 +774,7 @@ + diff --git a/spring-integration-jmx/src/main/resources/org/springframework/integration/jmx/config/spring-integration-jmx-5.0.xsd b/spring-integration-jmx/src/main/resources/org/springframework/integration/jmx/config/spring-integration-jmx-5.0.xsd index 7920ecdc1c..245dd6995d 100644 --- a/spring-integration-jmx/src/main/resources/org/springframework/integration/jmx/config/spring-integration-jmx-5.0.xsd +++ b/spring-integration-jmx/src/main/resources/org/springframework/integration/jmx/config/spring-integration-jmx-5.0.xsd @@ -80,6 +80,7 @@ + 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 54d4573824..a88c2c2eeb 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 @@ -24,10 +24,12 @@ import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; +import static org.mockito.Mockito.spy; import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import org.aopalliance.aop.Advice; @@ -44,6 +46,7 @@ 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.config.GlobalChannelInterceptor; import org.springframework.integration.handler.MessageProcessor; import org.springframework.integration.handler.ServiceActivatingHandler; import org.springframework.integration.handler.advice.AbstractRequestHandlerAdvice; @@ -55,6 +58,8 @@ import org.springframework.integration.metadata.SimpleMetadataStore; import org.springframework.integration.selector.MetadataStoreSelector; import org.springframework.integration.support.MessageBuilder; import org.springframework.integration.test.util.TestUtils; +import org.springframework.integration.transaction.PseudoTransactionManager; +import org.springframework.integration.transaction.TransactionInterceptorBuilder; import org.springframework.integration.transformer.Transformer; import org.springframework.jmx.support.MBeanServerFactoryBean; import org.springframework.messaging.Message; @@ -62,11 +67,15 @@ import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; import org.springframework.messaging.MessageHandlingException; import org.springframework.messaging.PollableChannel; +import org.springframework.messaging.support.ChannelInterceptor; +import org.springframework.messaging.support.ChannelInterceptorAdapter; 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; +import org.springframework.test.context.junit4.SpringRunner; +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.interceptor.TransactionInterceptor; +import org.springframework.transaction.support.TransactionSynchronizationManager; import com.hazelcast.config.Config; import com.hazelcast.core.Hazelcast; @@ -77,8 +86,7 @@ import com.hazelcast.core.HazelcastInstance; * @author Gary Russell * @since 4.1 */ -@ContextConfiguration -@RunWith(SpringJUnit4ClassRunner.class) +@RunWith(SpringRunner.class) @DirtiesContext public class IdempotentReceiverIntegrationTests { @@ -109,11 +117,14 @@ public class IdempotentReceiverIntegrationTests { @Autowired private MessageChannel annotatedBeanMessageHandlerChannel2; + @Autowired + private AtomicBoolean txSupplied; + @Test public void testIdempotentReceiver() { this.idempotentReceiverInterceptor.setThrowExceptionOnRejection(true); TestUtils.getPropertyValue(this.store, "metadata", Map.class).clear(); - Message message = new GenericMessage("foo"); + Message message = new GenericMessage<>("foo"); this.input.send(message); Message receive = this.output.receive(10000); assertNotNull(receive); @@ -136,18 +147,22 @@ public class IdempotentReceiverIntegrationTests { assertEquals(2, this.adviceCalled.get()); assertTrue(receive.getHeaders().get(IntegrationMessageHeaderAccessor.DUPLICATE_MESSAGE, Boolean.class)); assertEquals(1, TestUtils.getPropertyValue(store, "metadata", Map.class).size()); + + assertTrue(this.txSupplied.get()); } @Test public void testIdempotentReceiverOnMethod() { TestUtils.getPropertyValue(this.store, "metadata", Map.class).clear(); - Message message = new GenericMessage("foo"); + 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)); + assertTrue( + this.fooService.messages.get(1) + .getHeaders() + .get(IntegrationMessageHeaderAccessor.DUPLICATE_MESSAGE, Boolean.class)); } @Test @@ -197,7 +212,9 @@ public class IdempotentReceiverIntegrationTests { @Bean public ConcurrentMetadataStore store() { - return new SimpleMetadataStore(hazelcastInstance().getMap("idempotentReceiverMetadataStore")); + return new SimpleMetadataStore( + hazelcastInstance() + .getMap("idempotentReceiverMetadataStore")); } @Bean @@ -208,6 +225,17 @@ public class IdempotentReceiverIntegrationTests { message -> message.getPayload().toString().toUpperCase(), store())); } + @Bean + public PlatformTransactionManager transactionManager() { + return spy(new PseudoTransactionManager()); + } + + @Bean + public TransactionInterceptor transactionInterceptor() { + return new TransactionInterceptorBuilder(true) + .build(); + } + @Bean public MessageChannel input() { return new DirectChannel(); @@ -218,9 +246,31 @@ public class IdempotentReceiverIntegrationTests { return new QueueChannel(); } + @Bean + public AtomicBoolean txSupplied() { + return new AtomicBoolean(); + } + + @Bean + @GlobalChannelInterceptor(patterns = "output") + public ChannelInterceptor txSuppliedChannelInterceptor(final AtomicBoolean txSupplied) { + return new ChannelInterceptorAdapter() { + + @Override + public void postSend(Message message, MessageChannel channel, boolean sent) { + super.postSend(message, channel, sent); + txSupplied.set(TransactionSynchronizationManager.isActualTransactionActive()); + } + + }; + } + @Bean @org.springframework.integration.annotation.Transformer(inputChannel = "input", - outputChannel = "output", adviceChain = {"fooAdvice", "idempotentReceiverInterceptor"}) + outputChannel = "output", + adviceChain = { "fooAdvice", + "idempotentReceiverInterceptor", + "transactionInterceptor" }) public Transformer transformer() { return message -> message; } @@ -258,14 +308,7 @@ public class IdempotentReceiverIntegrationTests { @ServiceActivator(inputChannel = "annotatedBeanMessageHandlerChannel") @IdempotentReceiver("idempotentReceiverInterceptor") public MessageHandler messageHandler() { - return new ServiceActivatingHandler(new MessageProcessor() { - - @Override - public Object processMessage(Message message) { - return message; - } - - }); + return new ServiceActivatingHandler((MessageProcessor) message -> message); } @Bean diff --git a/spring-integration-jmx/src/test/resources/log4j.properties b/spring-integration-jmx/src/test/resources/log4j.properties index 0c503ac468..71c73a07f2 100644 --- a/spring-integration-jmx/src/test/resources/log4j.properties +++ b/spring-integration-jmx/src/test/resources/log4j.properties @@ -6,6 +6,6 @@ log4j.appender.stdout.layout.ConversionPattern=%d{ABSOLUTE} %5p %t %c{2}:%L - %m log4j.category.org.springframework=WARN -#log4j.category.org.springframework.integration=DEBUG -#log4j.category.org.springframework.beans.factory=DEBUG -#log4j.category.org.springframework.integration.monitor=TRACE +log4j.category.org.springframework.integration=WARN +log4j.category.org.springframework.integration.jmx=WARN +log4j.category.org.springframework.integration.monitor=WARN diff --git a/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/config/xml/AbstractJpaOutboundGatewayParser.java b/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/config/xml/AbstractJpaOutboundGatewayParser.java index d1d38437f1..bd970da042 100644 --- a/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/config/xml/AbstractJpaOutboundGatewayParser.java +++ b/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/config/xml/AbstractJpaOutboundGatewayParser.java @@ -18,15 +18,12 @@ package org.springframework.integration.jpa.config.xml; import org.w3c.dom.Element; -import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.support.BeanDefinitionBuilder; -import org.springframework.beans.factory.support.ManagedList; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.xml.AbstractConsumerEndpointParser; import org.springframework.integration.config.xml.IntegrationNamespaceUtils; import org.springframework.integration.jpa.outbound.JpaOutboundGatewayFactoryBean; import org.springframework.util.StringUtils; -import org.springframework.util.xml.DomUtils; /** * The Abstract Parser for the JPA Outbound Gateways. @@ -56,16 +53,6 @@ public abstract class AbstractJpaOutboundGatewayParser extends AbstractConsumerE jpaOutboundGatewayBuilder.addPropertyReference("outputChannel", replyChannel); } - final Element transactionalElement = DomUtils.getChildElementByTagName(gatewayElement, "transactional"); - - if (transactionalElement != null) { - BeanDefinition txAdviceDefinition = - IntegrationNamespaceUtils.configureTransactionAttributes(transactionalElement); - ManagedList adviceChain = new ManagedList(); - adviceChain.add(txAdviceDefinition); - jpaOutboundGatewayBuilder.addPropertyValue("txAdviceChain", adviceChain); - } - return jpaOutboundGatewayBuilder; } diff --git a/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/config/xml/JpaOutboundChannelAdapterParser.java b/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/config/xml/JpaOutboundChannelAdapterParser.java index 4714d945a8..675b6cd0ea 100644 --- a/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/config/xml/JpaOutboundChannelAdapterParser.java +++ b/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/config/xml/JpaOutboundChannelAdapterParser.java @@ -22,12 +22,10 @@ import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.parsing.BeanComponentDefinition; import org.springframework.beans.factory.support.AbstractBeanDefinition; import org.springframework.beans.factory.support.BeanDefinitionBuilder; -import org.springframework.beans.factory.support.ManagedList; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.xml.AbstractOutboundChannelAdapterParser; import org.springframework.integration.config.xml.IntegrationNamespaceUtils; import org.springframework.integration.jpa.outbound.JpaOutboundGatewayFactoryBean; -import org.springframework.util.xml.DomUtils; /** * The parser for JPA outbound channel adapter @@ -77,16 +75,6 @@ public class JpaOutboundChannelAdapterParser extends AbstractOutboundChannelAdap jpaOutboundChannelAdapterBuilder.addPropertyReference("jpaExecutor", jpaExecutorBeanName) .addPropertyValue("producesReply", Boolean.FALSE); - final Element transactionalElement = DomUtils.getChildElementByTagName(element, "transactional"); - - if (transactionalElement != null) { - BeanDefinition txAdviceDefinition = - IntegrationNamespaceUtils.configureTransactionAttributes(transactionalElement); - ManagedList adviceChain = new ManagedList(); - adviceChain.add(txAdviceDefinition); - jpaOutboundChannelAdapterBuilder.addPropertyValue("txAdviceChain", adviceChain); - } - return jpaOutboundChannelAdapterBuilder.getBeanDefinition(); } diff --git a/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/outbound/JpaOutboundGatewayFactoryBean.java b/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/outbound/JpaOutboundGatewayFactoryBean.java index f9b0489da6..32577276ac 100644 --- a/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/outbound/JpaOutboundGatewayFactoryBean.java +++ b/spring-integration-jpa/src/main/java/org/springframework/integration/jpa/outbound/JpaOutboundGatewayFactoryBean.java @@ -20,7 +20,6 @@ import java.util.List; import org.aopalliance.aop.Advice; -import org.springframework.aop.framework.ProxyFactory; import org.springframework.beans.factory.FactoryBean; import org.springframework.beans.factory.config.AbstractFactoryBean; import org.springframework.integration.jpa.core.JpaExecutor; @@ -28,8 +27,6 @@ import org.springframework.integration.jpa.support.OutboundGatewayType; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; import org.springframework.transaction.interceptor.TransactionInterceptor; -import org.springframework.util.ClassUtils; -import org.springframework.util.CollectionUtils; /** * The {@link JpaOutboundGatewayFactoryBean} creates instances of the @@ -50,18 +47,11 @@ public class JpaOutboundGatewayFactoryBean extends AbstractFactoryBean txAdviceChain; - /** * <request-handler-advice-chain /> only applies to the handleRequestMessage. */ private List adviceChain; - private ClassLoader beanClassLoader = ClassUtils.getDefaultClassLoader(); - private boolean producesReply = true; private MessageChannel outputChannel; @@ -85,10 +75,6 @@ public class JpaOutboundGatewayFactoryBean extends AbstractFactoryBean txAdviceChain) { - this.txAdviceChain = txAdviceChain; - } - public void setAdviceChain(List adviceChain) { this.adviceChain = adviceChain; } @@ -128,12 +114,6 @@ public class JpaOutboundGatewayFactoryBean extends AbstractFactoryBean getObjectType() { return MessageHandler.class; @@ -154,18 +134,6 @@ public class JpaOutboundGatewayFactoryBean extends AbstractFactoryBean diff --git a/spring-integration-jpa/src/test/java/org/springframework/integration/jpa/outbound/JpaOutboundChannelAdapterTests-context.xml b/spring-integration-jpa/src/test/java/org/springframework/integration/jpa/outbound/JpaOutboundChannelAdapterTests-context.xml index f94c280cc0..41867f1317 100644 --- a/spring-integration-jpa/src/test/java/org/springframework/integration/jpa/outbound/JpaOutboundChannelAdapterTests-context.xml +++ b/spring-integration-jpa/src/test/java/org/springframework/integration/jpa/outbound/JpaOutboundChannelAdapterTests-context.xml @@ -1,23 +1,17 @@ + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xmlns:jpa="http://www.springframework.org/schema/integration/jpa" + xmlns:int="http://www.springframework.org/schema/integration" + xsi:schemaLocation="http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd + http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd + http://www.springframework.org/schema/integration/jpa http://www.springframework.org/schema/integration/jpa/spring-integration-jpa.xsd"> - + - + diff --git a/spring-integration-redis/src/main/resources/org/springframework/integration/redis/config/spring-integration-redis-5.0.xsd b/spring-integration-redis/src/main/resources/org/springframework/integration/redis/config/spring-integration-redis-5.0.xsd index 98bcf90a59..df2274ea09 100644 --- a/spring-integration-redis/src/main/resources/org/springframework/integration/redis/config/spring-integration-redis-5.0.xsd +++ b/spring-integration-redis/src/main/resources/org/springframework/integration/redis/config/spring-integration-redis-5.0.xsd @@ -658,6 +658,7 @@ + @@ -775,6 +776,7 @@ + diff --git a/spring-integration-rmi/src/main/resources/org/springframework/integration/rmi/config/spring-integration-rmi-5.0.xsd b/spring-integration-rmi/src/main/resources/org/springframework/integration/rmi/config/spring-integration-rmi-5.0.xsd index 7d821fd740..6f3807689c 100644 --- a/spring-integration-rmi/src/main/resources/org/springframework/integration/rmi/config/spring-integration-rmi-5.0.xsd +++ b/spring-integration-rmi/src/main/resources/org/springframework/integration/rmi/config/spring-integration-rmi-5.0.xsd @@ -82,6 +82,7 @@ + diff --git a/spring-integration-rmi/src/test/java/org/springframework/integration/rmi/BackToBackTests-context.xml b/spring-integration-rmi/src/test/java/org/springframework/integration/rmi/BackToBackTests-context.xml index 43ff846688..d5a8d97ced 100644 --- a/spring-integration-rmi/src/test/java/org/springframework/integration/rmi/BackToBackTests-context.xml +++ b/spring-integration-rmi/src/test/java/org/springframework/integration/rmi/BackToBackTests-context.xml @@ -12,8 +12,13 @@ + request-channel="good" reply-channel="reply" port="#{@port}"> + + + + + + diff --git a/spring-integration-rmi/src/test/java/org/springframework/integration/rmi/BackToBackTests.java b/spring-integration-rmi/src/test/java/org/springframework/integration/rmi/BackToBackTests.java index d5c513b8b1..46e4d3643a 100644 --- a/spring-integration-rmi/src/test/java/org/springframework/integration/rmi/BackToBackTests.java +++ b/spring-integration-rmi/src/test/java/org/springframework/integration/rmi/BackToBackTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2013-2014 the original author or authors. + * Copyright 2013-2016 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. @@ -21,25 +21,31 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertThat; import static org.junit.Assert.fail; +import static org.mockito.Matchers.any; +import static org.mockito.Mockito.verify; import org.junit.Test; import org.junit.runner.RunWith; + import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.support.AbstractApplicationContext; -import org.springframework.messaging.support.GenericMessage; import org.springframework.messaging.Message; import org.springframework.messaging.PollableChannel; import org.springframework.messaging.SubscribableChannel; -import org.springframework.test.context.ContextConfiguration; -import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.messaging.support.GenericMessage; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.junit4.SpringRunner; +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.TransactionDefinition; /** * @author Gary Russell + * @author Artem Bilan * @since 3.0 * */ -@ContextConfiguration -@RunWith(SpringJUnit4ClassRunner.class) +@RunWith(SpringRunner.class) +@DirtiesContext public class BackToBackTests { @Autowired @@ -57,12 +63,17 @@ public class BackToBackTests { @Autowired private AbstractApplicationContext context; + @Autowired + private PlatformTransactionManager transactionManager; + @Test public void testGood() { - good.send(new GenericMessage("foo")); + good.send(new GenericMessage<>("foo")); Message reply = this.reply.receive(0); assertNotNull(reply); assertEquals("reply:foo", reply.getPayload()); + + verify(this.transactionManager).getTransaction(any(TransactionDefinition.class)); } @Test diff --git a/spring-integration-rmi/src/test/resources/.svnignore b/spring-integration-rmi/src/test/resources/.svnignore deleted file mode 100644 index e69de29bb2..0000000000 diff --git a/spring-integration-rmi/src/test/resources/log4j.properties b/spring-integration-rmi/src/test/resources/log4j.properties new file mode 100644 index 0000000000..b603c7d86e --- /dev/null +++ b/spring-integration-rmi/src/test/resources/log4j.properties @@ -0,0 +1,8 @@ +log4j.rootCategory=WARN, stdout + +log4j.appender.stdout=org.apache.log4j.ConsoleAppender +log4j.appender.stdout.layout=org.apache.log4j.PatternLayout +log4j.appender.stdout.layout.ConversionPattern=%d{ABSOLUTE} %5p %t %c{2}:%L - %m%n + +log4j.category.org.springframework.integration=WARN +log4j.category.org.springframework.integration.rmi=WARN diff --git a/spring-integration-sftp/src/main/resources/org/springframework/integration/sftp/config/spring-integration-sftp-5.0.xsd b/spring-integration-sftp/src/main/resources/org/springframework/integration/sftp/config/spring-integration-sftp-5.0.xsd index 02b4254682..d02ec6cc53 100644 --- a/spring-integration-sftp/src/main/resources/org/springframework/integration/sftp/config/spring-integration-sftp-5.0.xsd +++ b/spring-integration-sftp/src/main/resources/org/springframework/integration/sftp/config/spring-integration-sftp-5.0.xsd @@ -209,6 +209,7 @@ + diff --git a/spring-integration-sftp/src/test/resources/log4j.properties b/spring-integration-sftp/src/test/resources/log4j.properties index e46776bc79..e3f87b00ee 100644 --- a/spring-integration-sftp/src/test/resources/log4j.properties +++ b/spring-integration-sftp/src/test/resources/log4j.properties @@ -4,6 +4,6 @@ log4j.appender.stdout=org.apache.log4j.ConsoleAppender log4j.appender.stdout.layout=org.apache.log4j.PatternLayout log4j.appender.stdout.layout.ConversionPattern=%d{ABSOLUTE} %5p %t %c{2}:%L - %m%n -log4j.category.com.jcraft.jsch=DEBUG -log4j.category.org.springframework.integration=DEBUG -log4j.category.org.springframework.integration.sftp=DEBUG +log4j.category.com.jcraft.jsch=WARN +log4j.category.org.springframework.integration=WARN +log4j.category.org.springframework.integration.sftp=WARN diff --git a/spring-integration-twitter/src/main/resources/org/springframework/integration/twitter/config/spring-integration-twitter-5.0.xsd b/spring-integration-twitter/src/main/resources/org/springframework/integration/twitter/config/spring-integration-twitter-5.0.xsd index 6de449d9a1..c214499f96 100644 --- a/spring-integration-twitter/src/main/resources/org/springframework/integration/twitter/config/spring-integration-twitter-5.0.xsd +++ b/spring-integration-twitter/src/main/resources/org/springframework/integration/twitter/config/spring-integration-twitter-5.0.xsd @@ -292,6 +292,7 @@ + diff --git a/spring-integration-ws/src/main/resources/org/springframework/integration/ws/config/spring-integration-ws-5.0.xsd b/spring-integration-ws/src/main/resources/org/springframework/integration/ws/config/spring-integration-ws-5.0.xsd index 20ce37ef56..a477023afc 100644 --- a/spring-integration-ws/src/main/resources/org/springframework/integration/ws/config/spring-integration-ws-5.0.xsd +++ b/spring-integration-ws/src/main/resources/org/springframework/integration/ws/config/spring-integration-ws-5.0.xsd @@ -31,6 +31,7 @@ + diff --git a/src/reference/asciidoc/handler-advice.adoc b/src/reference/asciidoc/handler-advice.adoc index 685e2b62d6..26244921a5 100644 --- a/src/reference/asciidoc/handler-advice.adoc +++ b/src/reference/asciidoc/handler-advice.adoc @@ -487,6 +487,79 @@ Note, however, that in that case, the entire downstream flow would be within the In the case of a `MessageHandler` that does **not** return a response, the advice chain order is retained. +[[tx-handle-message-advice]] +==== Transaction Support + +Starting with _version 5.0_ a new `TransactionHandleMessageAdvice` has been introduced to make the whole downstream flow transactional, thanks to the `HandleMessageAdvice` implementation. +When regular `TransactionInterceptor` is used in the ``, for example via `` configuration, a started transaction is only applied only for an internal `AbstractReplyProducingMessageHandler.handleRequestMessage()` and isn't propagated to the downstream flow. + +To simplify XML configuration, alongside with the ``, a `` sub-element has been added to all `` and `` & family components: + +[source,xml] +---- + + + + + + + +---- + +For whom is familiar with <> such a configuration isn't new, but now we can start transaction from any point in our flow, not only from the `` or Message Driven Channel Adapter like in <>. + +Java & Annotation configuration can be simplified via newly introduced `TransactionInterceptorBuilder` and the result bean name can be used in the <> `adviceChain` attribute: + +[source,java] +---- +@Bean +public ConcurrentMetadataStore store() { + return new SimpleMetadataStore(hazelcastInstance() + .getMap("idempotentReceiverMetadataStore")); +} + +@Bean +public IdempotentReceiverInterceptor idempotentReceiverInterceptor() { + return new IdempotentReceiverInterceptor( + new MetadataStoreSelector( + message -> message.getPayload().toString(), + message -> message.getPayload().toString().toUpperCase(), store())); +} + +@Bean +public TransactionInterceptor transactionInterceptor() { + return new TransactionInterceptorBuilder(true) + .transactionManager(this.transactionManager) + .isolation(Isolation.READ_COMMITTED) + .propagation(Propagation.REQUIRES_NEW) + .build(); +} + +@Bean +@org.springframework.integration.annotation.Transformer(inputChannel = "input", + outputChannel = "output", + adviceChain = { "idempotentReceiverInterceptor", + "transactionInterceptor" }) +public Transformer transformer() { + return message -> message; +} +---- + +Note the `true` for the `TransactionInterceptorBuilder` constructor, which means produce `TransactionHandleMessageAdvice`, not regular `TransactionInterceptor`. + +Java DSL supports such an `Advice` via `.transactional()` options on the endpoint configuration: +[source,java] +---- +@Bean +public IntegrationFlow updatingGatewayFlow() { + return f -> f + .handle(Jpa.updatingGateway(this.entityManagerFactory), + e -> e.transactional(true)) + .channel(c -> c.queue("persistResults")); +} +---- + [[advising-filters]] ==== Advising Filters @@ -510,11 +583,11 @@ An example with the discard being performed after the advice is shown below. @MessageEndpoint public class MyAdvisedFilter { - @Filter(inputChannel="input", outputChannel="output", - adviceChain="adviceChain", discardWithinAdvice="false") - public boolean filter(String s) { - return s.contains("good"); - } + @Filter(inputChannel="input", outputChannel="output", + adviceChain="adviceChain", discardWithinAdvice="false") + public boolean filter(String s) { + return s.contains("good"); + } } ---- diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index 9e778bac0c..ad849dbaa9 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -18,6 +18,9 @@ development process. The `@Poller` annotation now has the `errorChannel` attribute for easier configuration of the underlying `MessagePublishingErrorHandler`. See <> for more information. +All the request-reply endpoints (based on `AbstractReplyProducingMessageHandler`) can now start transaction and, therefore, make the whole downstream flow transactional. +See <> for more information. + ==== JMS Changes Previously, Spring Integration JMS XML configuration used a default bean name `connectionFactory` for the JMS Connection Factory, allowing the property to be omitted from component definitions.