From 0ddc1fe6673ff3199ef62bfe993970a81288514f Mon Sep 17 00:00:00 2001 From: Iwein Fuld Date: Fri, 15 Oct 2010 17:20:58 +0200 Subject: [PATCH] INT-1525: Add treatment for releasePartialSequences to ResequencerParser and CorrelatingMessageHandler --- .../aggregator/CorrelatingMessageHandler.java | 10 ++++++++++ .../integration/config/xml/ResequencerParser.java | 2 ++ .../integration/config/ResequencerParserTests.java | 11 ++++++++++- .../integration/config/resequencerParserTests.xml | 2 +- 4 files changed, 23 insertions(+), 2 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/CorrelatingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/CorrelatingMessageHandler.java index b8c200b889..75b56aa947 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/CorrelatingMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/CorrelatingMessageHandler.java @@ -142,6 +142,13 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements this.sendPartialResultOnExpiry = sendPartialResultOnExpiry; } + public void setReleasePartialSequences(boolean releasePartialSequences){ + Assert.isInstanceOf(SequenceSizeReleaseStrategy.class, this.releaseStrategy, + "Release strategy of type [" + this.releaseStrategy.getClass().getSimpleName() + + "] cannot release partial sequences. Use the default SequenceSizeReleaseStrategy instead."); + ((SequenceSizeReleaseStrategy)this.releaseStrategy).setReleasePartialSequences(releasePartialSequences); + } + @Override public String getComponentType() { return "aggregator"; @@ -162,6 +169,9 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements synchronized (lock) { MessageGroup group = messageStore.getMessageGroup(correlationKey); if (group.canAdd(message)) { + if (logger.isTraceEnabled()) { + logger.trace("Adding message to group [ " + group + "]"); + } group = store(correlationKey, message); if (releaseStrategy.canRelease(group)) { Collection completedMessages = null; diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResequencerParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResequencerParser.java index 70f78a7648..c069a88bee 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResequencerParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/ResequencerParser.java @@ -26,6 +26,7 @@ import org.w3c.dom.Element; * * @author Marius Bogoevici * @author Dave Syer + * @author Iwein Fuld */ public class ResequencerParser extends AbstractConsumerEndpointParser { @@ -84,6 +85,7 @@ public class ResequencerParser extends AbstractConsumerEndpointParser { IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, DISCARD_CHANNEL_ATTRIBUTE); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, SEND_TIMEOUT_ATTRIBUTE); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, SEND_PARTIAL_RESULT_ON_EXPIRY_ATTRIBUTE); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, RELEASE_PARTIAL_SEQUENCES_ATTRIBUTE); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "auto-startup"); return builder; } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/ResequencerParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/ResequencerParserTests.java index 82f2f1c196..9c9dea0ba4 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/ResequencerParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/ResequencerParserTests.java @@ -77,7 +77,7 @@ public class ResequencerParserTests { "The ResequencerEndpoint is not configured with the appropriate 'send partial results on timeout' flag", true, getPropertyValue(resequencer, "sendPartialResultOnExpiry")); assertEquals("The ResequencerEndpoint is not configured with the appropriate 'release partial sequences' flag", - false, getPropertyValue(getPropertyValue(resequencer, "releaseStrategy"), "releasePartialSequences")); + true, getPropertyValue(getPropertyValue(resequencer, "releaseStrategy"), "releasePartialSequences")); } @Test @@ -90,6 +90,15 @@ public class ResequencerParserTests { .getBean("testCorrelationStrategy"), getPropertyValue(resequencer, "correlationStrategy")); } + @Test + public void shouldSetReleasePartialSequencesFlag(){ + EventDrivenConsumer endpoint = (EventDrivenConsumer) context.getBean("completelyDefinedResequencer"); + CorrelatingMessageHandler resequencer = TestUtils.getPropertyValue(endpoint, "handler", + CorrelatingMessageHandler.class); + assertEquals("The ResequencerEndpoint is not configured with the appropriate 'release partial sequences' flag", + true, getPropertyValue(getPropertyValue(resequencer, "releaseStrategy"), "releasePartialSequences")); + } + @Test public void testCorrelationStrategyRefAndMethod() throws Exception { EventDrivenConsumer endpoint = (EventDrivenConsumer) context diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/resequencerParserTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/resequencerParserTests.xml index 7a766c4656..f6201a0a3b 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/resequencerParserTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/resequencerParserTests.xml @@ -35,7 +35,7 @@ discard-channel="discardChannel" send-timeout="86420000" send-partial-result-on-expiry="true" - release-partial-sequences="false"/> + release-partial-sequences="true"/>