INT-1126: rename timeout to expiry
This commit is contained in:
@@ -72,7 +72,7 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements
|
||||
|
||||
private volatile MessageChannel discardChannel = new NullChannel();
|
||||
|
||||
private boolean sendPartialResultOnTimeout = false;
|
||||
private boolean sendPartialResultOnExpiry = false;
|
||||
|
||||
private final ConcurrentMap<Object, Object> locks = new ConcurrentHashMap<Object, Object>();
|
||||
|
||||
@@ -132,8 +132,8 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements
|
||||
this.channelTemplate.setSendTimeout(sendTimeout);
|
||||
}
|
||||
|
||||
public void setSendPartialResultOnTimeout(boolean sendPartialResultOnTimeout) {
|
||||
this.sendPartialResultOnTimeout = sendPartialResultOnTimeout;
|
||||
public void setSendPartialResultOnExpiry(boolean sendPartialResultOnExpiry) {
|
||||
this.sendPartialResultOnExpiry = sendPartialResultOnExpiry;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -209,7 +209,7 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements
|
||||
remove(group);
|
||||
}
|
||||
else {
|
||||
if (sendPartialResultOnTimeout) {
|
||||
if (sendPartialResultOnExpiry) {
|
||||
if (logger.isInfoEnabled()) {
|
||||
logger.info("Processing partially complete messages for key [" + correlationKey + "] to: "
|
||||
+ outputChannel);
|
||||
|
||||
@@ -59,8 +59,8 @@ public @interface Aggregator {
|
||||
long sendTimeout() default CorrelatingMessageHandler.DEFAULT_SEND_TIMEOUT;
|
||||
|
||||
/**
|
||||
* indicates whether to send an incomplete aggregate on timeout
|
||||
* indicates whether to send an incomplete aggregate on expiry of the message group
|
||||
*/
|
||||
boolean sendPartialResultsOnTimeout() default false;
|
||||
boolean sendPartialResultsOnExpiry() default false;
|
||||
|
||||
}
|
||||
|
||||
@@ -65,7 +65,7 @@ public class AggregatorAnnotationPostProcessor extends AbstractMethodAnnotationP
|
||||
handler.setOutputChannel(this.channelResolver.resolveChannelName(outputChannelName));
|
||||
}
|
||||
handler.setSendTimeout(annotation.sendTimeout());
|
||||
handler.setSendPartialResultOnTimeout(annotation.sendPartialResultsOnTimeout());
|
||||
handler.setSendPartialResultOnExpiry(annotation.sendPartialResultsOnExpiry());
|
||||
handler.setBeanFactory(this.beanFactory);
|
||||
handler.afterPropertiesSet();
|
||||
return handler;
|
||||
|
||||
@@ -50,7 +50,7 @@ public class AggregatorParser extends AbstractConsumerEndpointParser {
|
||||
|
||||
private static final String SEND_TIMEOUT_ATTRIBUTE = "send-timeout";
|
||||
|
||||
private static final String SEND_PARTIAL_RESULT_ON_TIMEOUT_ATTRIBUTE = "send-partial-result-on-timeout";
|
||||
private static final String SEND_PARTIAL_RESULT_ON_TIMEOUT_ATTRIBUTE = "send-partial-result-on-expiry";
|
||||
|
||||
private static final String RELEASE_STRATEGY_PROPERTY = "releaseStrategy";
|
||||
|
||||
|
||||
@@ -64,7 +64,7 @@ public class ResequencerParser extends AbstractConsumerEndpointParser {
|
||||
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "message-store");
|
||||
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "discard-channel");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "send-timeout");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "send-partial-result-on-timeout");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "send-partial-result-on-expiry");
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "auto-startup");
|
||||
return builder;
|
||||
}
|
||||
|
||||
@@ -1764,7 +1764,7 @@
|
||||
</xsd:appinfo>
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
<xsd:attribute name="send-partial-result-on-timeout" type="xsd:string" />
|
||||
<xsd:attribute name="send-partial-result-on-expiry" type="xsd:string" />
|
||||
</xsd:extension>
|
||||
</xsd:complexContent>
|
||||
</xsd:complexType>
|
||||
@@ -1832,7 +1832,7 @@
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
<xsd:attribute name="release-partial-sequences" type="xsd:string" />
|
||||
<xsd:attribute name="send-partial-result-on-timeout" type="xsd:string" />
|
||||
<xsd:attribute name="send-partial-result-on-expiry" type="xsd:string" />
|
||||
</xsd:extension>
|
||||
</xsd:complexContent>
|
||||
</xsd:complexType>
|
||||
|
||||
@@ -83,7 +83,7 @@ public class AggregatorTests {
|
||||
|
||||
@Test
|
||||
public void testShouldSendPartialResultOnTimeoutTrue() throws InterruptedException {
|
||||
this.aggregator.setSendPartialResultOnTimeout(true);
|
||||
this.aggregator.setSendPartialResultOnExpiry(true);
|
||||
QueueChannel replyChannel = new QueueChannel();
|
||||
Message<?> message1 = createMessage(3, "ABC", 3, 1, replyChannel, null);
|
||||
Message<?> message2 = createMessage(5, "ABC", 3, 2, replyChannel, null);
|
||||
|
||||
@@ -129,7 +129,7 @@ public class ConcurrentAggregatorTests {
|
||||
@Test
|
||||
public void testShouldSendPartialResultOnTimeoutTrue()
|
||||
throws InterruptedException {
|
||||
this.aggregator.setSendPartialResultOnTimeout(true);
|
||||
this.aggregator.setSendPartialResultOnExpiry(true);
|
||||
QueueChannel replyChannel = new QueueChannel();
|
||||
Message<?> message1 = createMessage(3, "ABC", 3, 1, replyChannel, null);
|
||||
Message<?> message2 = createMessage(5, "ABC", 3, 2, replyChannel, null);
|
||||
|
||||
@@ -158,7 +158,7 @@ public class ResequencerTests {
|
||||
Message<?> message1 = createMessage("123", "ABC", 4, 2, null);
|
||||
Message<?> message2 = createMessage("456", "ABC", 4, 1, null);
|
||||
Message<?> message3 = createMessage("789", "ABC", 4, 4, null);
|
||||
this.resequencer.setSendPartialResultOnTimeout(false);
|
||||
this.resequencer.setSendPartialResultOnExpiry(false);
|
||||
this.processor.setReleasePartialSequences(false);
|
||||
this.resequencer.setDiscardChannel(discardChannel);
|
||||
this.resequencer.handleMessage(message1);
|
||||
@@ -186,7 +186,7 @@ public class ResequencerTests {
|
||||
QueueChannel discardChannel = new QueueChannel();
|
||||
Message<?> message1 = createMessage("123", "ABC", 4, 2, null);
|
||||
Message<?> message2 = createMessage("456", "ABC", 5, 1, null);
|
||||
this.resequencer.setSendPartialResultOnTimeout(false);
|
||||
this.resequencer.setSendPartialResultOnExpiry(false);
|
||||
this.processor.setReleasePartialSequences(false);
|
||||
this.resequencer.setDiscardChannel(discardChannel);
|
||||
this.resequencer.handleMessage(message1);
|
||||
@@ -204,7 +204,7 @@ public class ResequencerTests {
|
||||
public void testResequencingWithWrongSequenceSizeAndNumber() throws InterruptedException {
|
||||
QueueChannel discardChannel = new QueueChannel();
|
||||
Message<?> message1 = createMessage("123", "ABC", 2, 4, null);
|
||||
this.resequencer.setSendPartialResultOnTimeout(false);
|
||||
this.resequencer.setSendPartialResultOnExpiry(false);
|
||||
this.processor.setReleasePartialSequences(false);
|
||||
this.resequencer.setDiscardChannel(discardChannel);
|
||||
this.resequencer.handleMessage(message1);
|
||||
|
||||
@@ -100,7 +100,7 @@ public class AggregatorParserTests {
|
||||
86420000l, TestUtils.getPropertyValue(consumer, "channelTemplate.sendTimeout"));
|
||||
Assert.assertEquals(
|
||||
"The AggregatorEndpoint is not configured with the appropriate 'send partial results on timeout' flag",
|
||||
true, accessor.getPropertyValue("sendPartialResultOnTimeout"));
|
||||
true, accessor.getPropertyValue("sendPartialResultOnExpiry"));
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -86,7 +86,7 @@ public class ResequencerParserTests {
|
||||
resequencer, "channelTemplate.sendTimeout"));
|
||||
assertEquals(
|
||||
"The ResequencerEndpoint is not configured with the appropriate 'send partial results on timeout' flag",
|
||||
false, getPropertyValue(resequencer, "sendPartialResultOnTimeout"));
|
||||
false, getPropertyValue(resequencer, "sendPartialResultOnExpiry"));
|
||||
assertEquals("The ResequencerEndpoint is not configured with the appropriate 'release partial sequences' flag",
|
||||
false, getPropertyValue(getPropertyValue(resequencer, "outputProcessor"), "releasePartialSequences"));
|
||||
}
|
||||
@@ -106,7 +106,7 @@ public class ResequencerParserTests {
|
||||
getPropertyValue(resequencer, "channelTemplate.sendTimeout"));
|
||||
assertEquals(
|
||||
"The ResequencerEndpoint is not configured with the appropriate 'send partial results on timeout' flag",
|
||||
true, getPropertyValue(resequencer, "sendPartialResultOnTimeout"));
|
||||
true, getPropertyValue(resequencer, "sendPartialResultOnExpiry"));
|
||||
assertEquals("The ResequencerEndpoint is not configured with the appropriate 'release partial sequences' flag",
|
||||
false, getPropertyValue(getPropertyValue(resequencer, "outputProcessor"), "releasePartialSequences"));
|
||||
}
|
||||
|
||||
@@ -27,7 +27,7 @@
|
||||
release-strategy="releaseStrategy"
|
||||
correlation-strategy="correlationStrategy"
|
||||
send-timeout="86420000"
|
||||
send-partial-result-on-timeout="true"/>
|
||||
send-partial-result-on-expiry="true"/>
|
||||
|
||||
<channel id="aggregatorWithReferenceAndMethodInput"/>
|
||||
<aggregator id="aggregatorWithReferenceAndMethod"
|
||||
|
||||
@@ -58,7 +58,7 @@ public class AggregatorAnnotationTests {
|
||||
assertTrue(getPropertyValue(aggregator, "discardChannel") instanceof NullChannel);
|
||||
assertEquals(CorrelatingMessageHandler.DEFAULT_SEND_TIMEOUT,
|
||||
getPropertyValue(aggregator, "channelTemplate.sendTimeout"));
|
||||
assertEquals(false, getPropertyValue(aggregator, "sendPartialResultOnTimeout"));
|
||||
assertEquals(false, getPropertyValue(aggregator, "sendPartialResultOnExpiry"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -75,7 +75,7 @@ public class AggregatorAnnotationTests {
|
||||
assertEquals(channelResolver.resolveChannelName("discardChannel"),
|
||||
getPropertyValue(aggregator, "discardChannel"));
|
||||
assertEquals(98765432l, getPropertyValue(aggregator, "channelTemplate.sendTimeout"));
|
||||
assertEquals(true, getPropertyValue(aggregator, "sendPartialResultOnTimeout"));
|
||||
assertEquals(true, getPropertyValue(aggregator, "sendPartialResultOnExpiry"));
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -40,7 +40,7 @@ public class TestAnnotatedEndpointWithCustomizedAggregator {
|
||||
inputChannel = "inputChannel",
|
||||
outputChannel = "outputChannel",
|
||||
discardChannel = "discardChannel",
|
||||
sendPartialResultsOnTimeout = true,
|
||||
sendPartialResultsOnExpiry = true,
|
||||
sendTimeout = 98765432)
|
||||
public Message<?> aggregatingMethod(List<Message<?>> messages) {
|
||||
List<Message<?>> sortableList = new ArrayList<Message<?>>(messages);
|
||||
|
||||
@@ -32,7 +32,7 @@
|
||||
output-channel="outputChannel"
|
||||
discard-channel="discardChannel"
|
||||
send-timeout="86420000"
|
||||
send-partial-result-on-timeout="true"
|
||||
send-partial-result-on-expiry="true"
|
||||
release-partial-sequences="false"/>
|
||||
|
||||
<resequencer id="resequencerWithCorrelationStrategyRefOnly"
|
||||
|
||||
Reference in New Issue
Block a user