From d4e135b13d8888f9876c33efb7d9296664a58e40 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 10 Jul 2012 21:38:47 +0300 Subject: [PATCH] INT-2649: DelayHandler: tx & adviceChain support * 'delayer-type' XSD: add `` & `` * move `PollerParser#configureAdviceChain` into `IntegrationNamespaceUtils` * DelayerParser: parsing `` & `` via `IntegrationNamespaceUtils#configureAdviceChain` * DelayHandler: add `adviceChain` property * DelayHandler: introduce `ReleaseMessageHandler` to apply `adviceChain` capabilities around `DelayHandler#doReleaseMessage` * DelayerParserTests: tests for introduced sub-elements * DelayerHandlerRescheduleIntegrationTests: test for transaction boundaries with intentional rollback in the message-flow after message release JIRA: https://jira.springsource.org/browse/INT-2649 Resolve conflicts; polishing. --- .../integration/config/xml/DelayerParser.java | 13 +++- .../config/xml/IntegrationNamespaceUtils.java | 30 ++++++--- .../integration/config/xml/PollerParser.java | 4 +- .../integration/handler/DelayHandler.java | 64 ++++++++++++++++++- .../config/xml/spring-integration-2.2.xsd | 10 +++ .../config/xml/DelayerParserTests-context.xml | 31 ++++++++- .../config/xml/DelayerParserTests.java | 46 ++++++++++++- .../config/xml/DelayerUsageTests-context.xml | 8 ++- ...dlerRescheduleIntegrationTests-context.xml | 20 ++++++ ...ayerHandlerRescheduleIntegrationTests.java | 60 ++++++++++++++++- 10 files changed, 265 insertions(+), 21 deletions(-) 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 192e4231d1..4e124399f8 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 @@ -16,12 +16,12 @@ package org.springframework.integration.config.xml; -import org.springframework.integration.handler.DelayHandler; -import org.w3c.dom.Element; - import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.xml.ParserContext; +import org.springframework.integration.handler.DelayHandler; import org.springframework.util.StringUtils; +import org.springframework.util.xml.DomUtils; +import org.w3c.dom.Element; /** * Parser for the <delayer> element. @@ -68,6 +68,13 @@ public class DelayerParser extends AbstractConsumerEndpointParser { IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "message-store"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "send-timeout"); + + Element txElement = DomUtils.getChildElementByTagName(element, "transactional"); + Element adviceChainElement = DomUtils.getChildElementByTagName(element, "advice-chain"); + + IntegrationNamespaceUtils.configureAndSetAdviceChainIfPresent(adviceChainElement, txElement, builder, + parserContext, "delayedAdviceChain"); + 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 2c2fe14fe9..5c5b81f00a 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 @@ -292,17 +292,26 @@ public abstract class IntegrationNamespaceUtils { * @see AbstractPollingEndpoint */ public static BeanDefinition configureTransactionAttributes(Element txElement) { + BeanDefinition txDefinition = configureTransactionDefinition(txElement); + BeanDefinitionBuilder attributeSourceBuilder = BeanDefinitionBuilder.genericBeanDefinition(MatchAlwaysTransactionAttributeSource.class); + attributeSourceBuilder.addPropertyValue("transactionAttribute", txDefinition); + BeanDefinitionBuilder txInterceptorBuilder = BeanDefinitionBuilder.genericBeanDefinition(TransactionInterceptor.class); + txInterceptorBuilder.addPropertyReference("transactionManager", txElement.getAttribute("transaction-manager")); + txInterceptorBuilder.addPropertyValue("transactionAttributeSource", attributeSourceBuilder.getBeanDefinition()); + return txInterceptorBuilder.getBeanDefinition(); + } + + /** + * Parse attributes of "transactional" element and configure a {@link DefaultTransactionAttribute} + * with provided "transactionDefinition" properties. + */ + public static BeanDefinition configureTransactionDefinition(Element txElement) { BeanDefinitionBuilder txDefinitionBuilder = BeanDefinitionBuilder.genericBeanDefinition(DefaultTransactionAttribute.class); txDefinitionBuilder.addPropertyValue("propagationBehaviorName", "PROPAGATION_" + txElement.getAttribute("propagation")); txDefinitionBuilder.addPropertyValue("isolationLevelName", "ISOLATION_" + txElement.getAttribute("isolation")); txDefinitionBuilder.addPropertyValue("timeout", txElement.getAttribute("timeout")); txDefinitionBuilder.addPropertyValue("readOnly", txElement.getAttribute("read-only")); - BeanDefinitionBuilder attributeSourceBuilder = BeanDefinitionBuilder.genericBeanDefinition(MatchAlwaysTransactionAttributeSource.class); - attributeSourceBuilder.addPropertyValue("transactionAttribute", txDefinitionBuilder.getBeanDefinition()); - BeanDefinitionBuilder txInterceptorBuilder = BeanDefinitionBuilder.genericBeanDefinition(TransactionInterceptor.class); - txInterceptorBuilder.addPropertyReference("transactionManager", txElement.getAttribute("transaction-manager")); - txInterceptorBuilder.addPropertyValue("transactionAttributeSource", attributeSourceBuilder.getBeanDefinition()); - return txInterceptorBuilder.getBeanDefinition(); + return txDefinitionBuilder.getBeanDefinition(); } public static String[] generateAlias(Element element) { @@ -314,12 +323,17 @@ public abstract class IntegrationNamespaceUtils { return handlerAlias; } - @SuppressWarnings({ "rawtypes" }) public static void configureAndSetAdviceChainIfPresent(Element adviceChainElement, Element txElement, BeanDefinitionBuilder parentBuilder, ParserContext parserContext) { + configureAndSetAdviceChainIfPresent(adviceChainElement, txElement, parentBuilder, parserContext, "adviceChain"); + } + + @SuppressWarnings({ "rawtypes" }) + public static void configureAndSetAdviceChainIfPresent(Element adviceChainElement, Element txElement, + BeanDefinitionBuilder parentBuilder, ParserContext parserContext, String propertyName) { ManagedList adviceChain = configureAdviceChain(adviceChainElement, txElement, parentBuilder, parserContext); if (adviceChain != null) { - parentBuilder.addPropertyValue("adviceChain", adviceChain); + parentBuilder.addPropertyValue(propertyName, adviceChain); } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/PollerParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/PollerParser.java index 9d609eada3..9549c540a2 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/PollerParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/PollerParser.java @@ -81,9 +81,12 @@ public class PollerParser extends AbstractBeanDefinitionParser { parserContext.getReaderContext().error( "the 'ref' attribute must not be present on the top-level 'poller' element", element); } + configureTrigger(element, metadataBuilder, parserContext); + IntegrationNamespaceUtils.setValueIfAttributeDefined(metadataBuilder, element, "max-messages-per-poll"); IntegrationNamespaceUtils.setValueIfAttributeDefined(metadataBuilder, element, "receive-timeout"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(metadataBuilder, element, "task-executor"); Element txElement = DomUtils.getChildElementByTagName(element, "transactional"); Element adviceChainElement = DomUtils.getChildElementByTagName(element, "advice-chain"); @@ -103,7 +106,6 @@ public class PollerParser extends AbstractBeanDefinitionParser { pseudoTxElement = pseudoTxElement == null ? txSyncElement : pseudoTxElement; configureTransactionSync(pseudoTxElement, metadataBuilder, parserContext); - IntegrationNamespaceUtils.setReferenceIfAttributeDefined(metadataBuilder, element, "task-executor"); String errorChannel = element.getAttribute("error-channel"); if (StringUtils.hasText(errorChannel)) { BeanDefinitionBuilder errorHandler = BeanDefinitionBuilder.genericBeanDefinition(MessagePublishingErrorHandler.class); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java index ffb5febce5..6eab5a6bff 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java @@ -18,11 +18,15 @@ package org.springframework.integration.handler; import java.io.Serializable; import java.util.Date; +import java.util.List; import java.util.concurrent.atomic.AtomicBoolean; +import org.aopalliance.aop.Advice; +import org.springframework.aop.framework.ProxyFactory; import org.springframework.context.ApplicationListener; import org.springframework.context.event.ContextRefreshedEvent; import org.springframework.integration.Message; +import org.springframework.integration.MessagingException; import org.springframework.integration.context.IntegrationObjectSupport; import org.springframework.integration.core.MessageHandler; import org.springframework.integration.store.MessageGroup; @@ -34,6 +38,8 @@ import org.springframework.jmx.export.annotation.ManagedResource; import org.springframework.scheduling.TaskScheduler; import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; import org.springframework.util.Assert; +import org.springframework.util.ClassUtils; +import org.springframework.util.CollectionUtils; /** * A {@link MessageHandler} that is capable of delaying the continuation of a @@ -75,8 +81,12 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement private volatile MessageGroupStore messageStore; + private volatile List delayedAdviceChain; + private final AtomicBoolean initialized = new AtomicBoolean(); + private volatile MessageHandler releaseHandler = new ReleaseMessageHandler(); + /** * Create a DelayHandler with the given 'messageGroupId' that is used as 'key' for {@link MessageGroup} * to store delayed Messages in the {@link MessageGroupStore}. The sending of Messages after @@ -126,6 +136,17 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement this.messageStore = messageStore; } + /** + * Specify the List to advise {@link DelayHandler.ReleaseMessageHandler} proxy. + * Usually used to add transactions to delayed messages retrieved from a transactional message store. + * + * @see #createReleaseMessageTask + */ + public void setDelayedAdviceChain(List delayedAdviceChain) { + Assert.notNull(delayedAdviceChain, "delayedAdviceChain must not be null"); + this.delayedAdviceChain = delayedAdviceChain; + } + @Override public String getComponentType() { return "delayer"; @@ -140,6 +161,21 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement else { Assert.isInstanceOf(MessageStore.class, this.messageStore); } + + this.releaseHandler = this.createReleaseMessageTask(); + } + + private MessageHandler createReleaseMessageTask() { + ReleaseMessageHandler releaseHandler = new ReleaseMessageHandler(); + + if (!CollectionUtils.isEmpty(this.delayedAdviceChain)) { + ProxyFactory proxyFactory = new ProxyFactory(releaseHandler); + for (Advice advice : delayedAdviceChain) { + proxyFactory.addAdvice(advice); + } + return (MessageHandler) proxyFactory.getProxy(ClassUtils.getDefaultClassLoader()); + } + return releaseHandler; } /** @@ -150,7 +186,8 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement * * @param requestMessage - the Message which may be delayed. * @return - null if 'requestMessage' is delayed, - * otherwise - 'payload' from 'requestMessage'. + * otherwise - 'payload' from 'requestMessage'. + * * @see #releaseMessage */ @@ -216,6 +253,10 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement } private void releaseMessage(Message message) { + this.releaseHandler.handleMessage(message); + } + + private void doReleaseMessage(Message message) { if (this.messageStore instanceof SimpleMessageStore || ((MessageStore) this.messageStore).removeMessage(message.getHeaders().getId()) != null) { this.messageStore.removeMessageFromGroup(this.messageGroupId, message); @@ -265,7 +306,8 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement * in the 'parent-child' contexts, e.g. in the Spring-MVC applications. * * @param event - {@link ContextRefreshedEvent} which occurs - * after Application context is completely initialized. + * after Application context is completely initialized. + * * @see #reschedulePersistedMessages */ public void onApplicationEvent(ContextRefreshedEvent event) { @@ -274,6 +316,23 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement } } + + /** + * Delegate {@link MessageHandler} implementation for 'release Message task'. + * Used as 'pointcut' to wrap 'release Message task' with adviceChain. + * + * @see @createReleaseMessageTask + * @see @releaseMessage + */ + private class ReleaseMessageHandler implements MessageHandler { + + public void handleMessage(Message message) throws MessagingException { + DelayHandler.this.doReleaseMessage(message); + } + + } + + private static final class DelayedMessageWrapper implements Serializable { private static final long serialVersionUID = -4739802369074947045L; @@ -308,6 +367,7 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement public int hashCode() { return this.original.hashCode(); } + } } diff --git a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.2.xsd b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.2.xsd index 951bc7a012..983085256b 100644 --- a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.2.xsd +++ b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.2.xsd @@ -1298,6 +1298,16 @@ + + + + 'transactional' and 'advice-chain' elements specify the configuration List of AOP Advice + to proxying DelayHandler's 'release Message task'. + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerParserTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerParserTests-context.xml index 1784af583a..c1e6f0c928 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerParserTests-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerParserTests-context.xml @@ -2,11 +2,13 @@ + http://www.springframework.org/schema/integration/spring-integration.xsd + http://www.springframework.org/schema/tx http://www.springframework.org/schema/tx/spring-tx.xsd"> @@ -34,10 +36,37 @@ default-delay="0" message-store="testMessageStore"/> + + + + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerParserTests.java index e9d0b4f856..2bebb0a2f0 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerParserTests.java @@ -17,12 +17,17 @@ package org.springframework.integration.config.xml; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertTrue; + +import java.util.HashMap; +import java.util.List; import org.junit.Test; import org.junit.runner.RunWith; - import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationContext; @@ -32,6 +37,11 @@ import org.springframework.integration.handler.DelayHandler; import org.springframework.integration.test.util.TestUtils; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.transaction.TransactionDefinition; +import org.springframework.transaction.interceptor.MatchAlwaysTransactionAttributeSource; +import org.springframework.transaction.interceptor.NameMatchTransactionAttributeSource; +import org.springframework.transaction.interceptor.TransactionAttributeSource; +import org.springframework.transaction.interceptor.TransactionInterceptor; /** * @author Mark Fisher @@ -45,7 +55,6 @@ public class DelayerParserTests { @Autowired private ApplicationContext context; - @Test public void defaultScheduler() { Object endpoint = context.getBean("delayerWithDefaultScheduler"); @@ -91,4 +100,37 @@ public class DelayerParserTests { assertEquals(context.getBean("testMessageStore"), accessor.getPropertyValue("messageStore")); } + @Test //INT-2649 + public void transactionalSubElement() { + Object endpoint = context.getBean("delayerWithTransactional"); + DelayHandler delayHandler = TestUtils.getPropertyValue(endpoint, "handler", DelayHandler.class); + List adviceChain = TestUtils.getPropertyValue(delayHandler, "delayedAdviceChain", List.class); + assertEquals(1, adviceChain.size()); + Object advice = adviceChain.get(0); + assertTrue(advice instanceof TransactionInterceptor); + TransactionAttributeSource transactionAttributeSource = ((TransactionInterceptor) advice).getTransactionAttributeSource(); + assertTrue(transactionAttributeSource instanceof MatchAlwaysTransactionAttributeSource); + TransactionDefinition definition = transactionAttributeSource.getTransactionAttribute(null, null); + assertEquals(TransactionDefinition.PROPAGATION_REQUIRED, definition.getPropagationBehavior()); + assertEquals(TransactionDefinition.ISOLATION_DEFAULT, definition.getIsolationLevel()); + assertEquals(TransactionDefinition.TIMEOUT_DEFAULT, definition.getTimeout()); + assertFalse(definition.isReadOnly()); + } + + @Test //INT-2649 + public void adviceChainSubElement() { + Object endpoint = context.getBean("delayerWithAdviceChain"); + DelayHandler delayHandler = TestUtils.getPropertyValue(endpoint, "handler", DelayHandler.class); + List adviceChain = TestUtils.getPropertyValue(delayHandler, "delayedAdviceChain", List.class); + assertEquals(2, adviceChain.size()); + assertSame(context.getBean("testAdviceBean"), adviceChain.get(0)); + + Object txAdvice = adviceChain.get(1); + assertEquals(TransactionInterceptor.class, txAdvice.getClass()); + TransactionAttributeSource transactionAttributeSource = ((TransactionInterceptor) txAdvice).getTransactionAttributeSource(); + assertEquals(NameMatchTransactionAttributeSource.class, transactionAttributeSource.getClass()); + HashMap nameMap = TestUtils.getPropertyValue(transactionAttributeSource, "nameMap", HashMap.class); + assertEquals("{*=PROPAGATION_REQUIRES_NEW,ISOLATION_DEFAULT,readOnly}", nameMap.toString()); + } + } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerUsageTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerUsageTests-context.xml index a24a2826cb..b9e4ed48de 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerUsageTests-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DelayerUsageTests-context.xml @@ -37,7 +37,13 @@ - + + + + + + + diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/DelayerHandlerRescheduleIntegrationTests-context.xml b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/DelayerHandlerRescheduleIntegrationTests-context.xml index 78b3f2f9ed..656bf70f91 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/DelayerHandlerRescheduleIntegrationTests-context.xml +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/DelayerHandlerRescheduleIntegrationTests-context.xml @@ -2,12 +2,17 @@ + + @@ -18,4 +23,19 @@ default-delay="200" message-store="messageStore"/> + + + + + + + + + + diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/DelayerHandlerRescheduleIntegrationTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/DelayerHandlerRescheduleIntegrationTests.java index 2e035cb051..388afda528 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/DelayerHandlerRescheduleIntegrationTests.java +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/DelayerHandlerRescheduleIntegrationTests.java @@ -19,6 +19,9 @@ import static org.junit.Assert.assertNotSame; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; + import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.Test; @@ -26,6 +29,8 @@ import org.springframework.context.support.AbstractApplicationContext; import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.integration.Message; import org.springframework.integration.MessageChannel; +import org.springframework.integration.MessagingException; +import org.springframework.integration.core.MessageHandler; import org.springframework.integration.core.PollableChannel; import org.springframework.integration.store.MessageGroup; import org.springframework.integration.store.MessageGroupStore; @@ -35,17 +40,19 @@ import org.springframework.integration.util.UUIDConverter; import org.springframework.jdbc.datasource.embedded.EmbeddedDatabase; import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseBuilder; import org.springframework.jdbc.datasource.embedded.EmbeddedDatabaseType; +import org.springframework.transaction.support.TransactionSynchronization; +import org.springframework.transaction.support.TransactionSynchronizationAdapter; +import org.springframework.transaction.support.TransactionSynchronizationManager; /** * @author Artem Bilan * @author Gary Russell */ -//INT-1132 public class DelayerHandlerRescheduleIntegrationTests { public static final String DELAYER_ID = "delayerWithJdbcMS"; - private static EmbeddedDatabase dataSource; + public static EmbeddedDatabase dataSource; @BeforeClass public static void init() { @@ -60,7 +67,7 @@ public class DelayerHandlerRescheduleIntegrationTests { dataSource.shutdown(); } - @Test + @Test //INT-1132 public void testDelayerHandlerRescheduleWithJdbcMessageStore() throws Exception { AbstractApplicationContext context = new ClassPathXmlApplicationContext("DelayerHandlerRescheduleIntegrationTests-context.xml", this.getClass()); MessageChannel input = context.getBean("input", MessageChannel.class); @@ -115,6 +122,31 @@ public class DelayerHandlerRescheduleIntegrationTests { } + @Test //INT-2649 + public void testRollbackOnDelayerHandlerReleaseTask() throws Exception { + AbstractApplicationContext context = new ClassPathXmlApplicationContext("DelayerHandlerRescheduleIntegrationTests-context.xml", this.getClass()); + MessageChannel input = context.getBean("transactionalDelayerInput", MessageChannel.class); + + MessageGroupStore messageStore = context.getBean("messageStore", MessageGroupStore.class); + String delayerMessageGroupId = UUIDConverter.getUUID("transactionalDelayer.messageGroupId").toString(); + assertEquals(0, messageStore.messageGroupSize(delayerMessageGroupId)); + + input.send(MessageBuilder.withPayload("test").build()); + + Thread.sleep(30); + + assertEquals(1, messageStore.messageGroupSize(delayerMessageGroupId)); + + //To check that 'rescheduling' works in the transaction boundaries too + context.destroy(); + context.refresh(); + + assertTrue(RollbackTxSync.latch.await(2, TimeUnit.SECONDS)); + + //On transaction rollback the delayed Message should remain in the persistent MessageStore + assertEquals(1, messageStore.messageGroupSize(delayerMessageGroupId)); + } + private static class TestJdbcMessageStore extends JdbcMessageStore { private TestJdbcMessageStore() { @@ -124,4 +156,26 @@ public class DelayerHandlerRescheduleIntegrationTests { } + private static class ExceptionMessageHandler implements MessageHandler { + + public void handleMessage(Message message) throws MessagingException { + TransactionSynchronizationManager.registerSynchronization(new RollbackTxSync()); + throw new RuntimeException("intentional"); + } + + } + + private static class RollbackTxSync extends TransactionSynchronizationAdapter { + + public static CountDownLatch latch = new CountDownLatch(2); + + @Override + public void afterCompletion(int status) { + if (TransactionSynchronization.STATUS_ROLLED_BACK == status) { + latch.countDown(); + } + } + + } + }