diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java index 1d3cfa9c72..0dd39ae595 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java @@ -100,7 +100,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP private final Map> expireGroupScheduledFutures = new ConcurrentHashMap<>(); - private final Set groupIds = ConcurrentHashMap.newKeySet(); + private final Set groupIds = ConcurrentHashMap.newKeySet(); private MessageGroupProcessor outputProcessor; @@ -144,15 +144,23 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP public AbstractCorrelatingMessageHandler(MessageGroupProcessor processor, MessageGroupStore store, CorrelationStrategy correlationStrategy, ReleaseStrategy releaseStrategy) { + Assert.notNull(processor, "'processor' must not be null"); Assert.notNull(store, "'store' must not be null"); setMessageStore(store); this.outputProcessor = processor; - this.correlationStrategy = (correlationStrategy == null - ? new HeaderAttributeCorrelationStrategy(IntegrationMessageHeaderAccessor.CORRELATION_ID) - : correlationStrategy); - this.releaseStrategy = releaseStrategy == null ? new SimpleSequenceSizeReleaseStrategy() : releaseStrategy; + + this.correlationStrategy = + correlationStrategy == null + ? new HeaderAttributeCorrelationStrategy(IntegrationMessageHeaderAccessor.CORRELATION_ID) + : correlationStrategy; + + this.releaseStrategy = + releaseStrategy == null + ? new SimpleSequenceSizeReleaseStrategy() + : releaseStrategy; + this.releaseStrategySet = releaseStrategy != null; this.sequenceAware = this.releaseStrategy instanceof SequenceSizeReleaseStrategy; } @@ -195,7 +203,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP } public void setForceReleaseAdviceChain(List forceReleaseAdviceChain) { - Assert.notNull(forceReleaseAdviceChain, "forceReleaseAdviceChain must not be null"); + Assert.notNull(forceReleaseAdviceChain, "'forceReleaseAdviceChain' must not be null"); this.forceReleaseAdviceChain = forceReleaseAdviceChain; } @@ -242,7 +250,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP } if (this.releasePartialSequences) { - Assert.isInstanceOf(SequenceSizeReleaseStrategy.class, this.releaseStrategy, + Assert.isInstanceOf(SequenceSizeReleaseStrategy.class, this.releaseStrategy, () -> "Release strategy of type [" + this.releaseStrategy.getClass().getSimpleName() + "] cannot release partial sequences. Use a SequenceSizeReleaseStrategy instead."); ((SequenceSizeReleaseStrategy) this.releaseStrategy) @@ -519,7 +527,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP * When 'groupTimeout' is evaluated to 'null' we do nothing. * The 'MessageGroupStoreReaper' can be used to 'forceComplete' message groups. */ - if (groupTimeout != null && groupTimeout >= 0) { + if (groupTimeout != null) { if (groupTimeout > 0) { final Object groupId = messageGroup.getGroupId(); final long timestamp = messageGroup.getTimestamp(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ConcurrentAggregatorTests.java b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ConcurrentAggregatorTests.java index 5f0ebe9864..ea9b54887a 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ConcurrentAggregatorTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ConcurrentAggregatorTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2018 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. @@ -47,6 +47,7 @@ import org.springframework.messaging.MessageHeaders; * @author Mark Fisher * @author Marius Bogoevici * @author Iwein Fuld + * @author Artem Bilan */ public class ConcurrentAggregatorTests { @@ -90,8 +91,7 @@ public class ConcurrentAggregatorTests { @Test @Ignore // dropped backwards compatibility for duplicate ID's - public void testCompleteGroupWithinTimeoutWithSameId() - throws InterruptedException { + public void testCompleteGroupWithinTimeoutWithSameId() { QueueChannel replyChannel = new QueueChannel(); Message message1 = createMessage(3, "ABC", 3, 1, replyChannel, "ID#1"); @@ -111,8 +111,8 @@ public class ConcurrentAggregatorTests { } @Test - public void testShouldNotSendPartialResultOnTimeoutByDefault() - throws InterruptedException { + public void testShouldNotSendPartialResultOnTimeoutByDefault() throws InterruptedException { + QueueChannel discardChannel = new QueueChannel(); this.aggregator.setDiscardChannel(discardChannel); QueueChannel replyChannel = new QueueChannel(); @@ -126,17 +126,16 @@ public class ConcurrentAggregatorTests { assertEquals("Task should have completed within timeout", 0, latch .getCount()); - Message reply = replyChannel.receive(1000); + Message reply = replyChannel.receive(10); assertNull("No message should have been sent normally", reply); this.store.expireMessageGroups(-10000); - Message discardedMessage = discardChannel.receive(1000); + Message discardedMessage = discardChannel.receive(10000); assertNotNull("A message should have been discarded", discardedMessage); assertEquals(message, discardedMessage); } @Test - public void testShouldSendPartialResultOnTimeoutTrue() - throws InterruptedException { + public void testShouldSendPartialResultOnTimeoutTrue() throws InterruptedException { this.aggregator.setSendPartialResultOnExpiry(true); QueueChannel replyChannel = new QueueChannel(); Message message1 = createMessage(3, "ABC", 3, 1, replyChannel, null); @@ -252,16 +251,14 @@ public class ConcurrentAggregatorTests { } @Test(expected = MessageHandlingException.class) - public void testExceptionThrownIfNoCorrelationId() - throws InterruptedException { + public void testExceptionThrownIfNoCorrelationId() { Message message = createMessage(3, null, 2, 1, new QueueChannel(), null); this.aggregator.handleMessage(message); } @Test - public void testAdditionalMessageAfterCompletion() - throws InterruptedException { + public void testAdditionalMessageAfterCompletion() throws InterruptedException { QueueChannel replyChannel = new QueueChannel(); Message message1 = createMessage(3, "ABC", 3, 1, replyChannel, null); Message message2 = createMessage(5, "ABC", 3, 2, replyChannel, null); @@ -291,6 +288,7 @@ public class ConcurrentAggregatorTests { private static Message createMessage(Object payload, Object correlationId, int sequenceSize, int sequenceNumber, MessageChannel replyChannel, String predefinedId) { + MessageBuilder builder = MessageBuilder.withPayload(payload) .setCorrelationId(correlationId).setSequenceSize(sequenceSize) .setSequenceNumber(sequenceNumber) diff --git a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/integration/AggregatorIntegrationTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/integration/AggregatorIntegrationTests-context.xml index 48b4040ada..d760b273cd 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/integration/AggregatorIntegrationTests-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/integration/AggregatorIntegrationTests-context.xml @@ -42,7 +42,7 @@ headers = stubHeaders(i, 5, 1); - input.send(new GenericMessage(i, headers)); + input.send(new GenericMessage<>(i, headers)); } Message receive = output.receive(10000); assertNotNull(receive); @@ -102,10 +100,10 @@ public class AggregatorIntegrationTests { } @Test - public void testNonExpiringAggregator() throws Exception { + public void testNonExpiringAggregator() { for (int i = 0; i < 5; i++) { Map headers = stubHeaders(i, 5, 1); - nonExpiringAggregatorInput.send(new GenericMessage(i, headers)); + nonExpiringAggregatorInput.send(new GenericMessage<>(i, headers)); } assertNotNull(output.receive(0)); @@ -113,7 +111,7 @@ public class AggregatorIntegrationTests { for (int i = 5; i < 10; i++) { Map headers = stubHeaders(i, 5, 1); - nonExpiringAggregatorInput.send(new GenericMessage(i, headers)); + nonExpiringAggregatorInput.send(new GenericMessage<>(i, headers)); } assertNull(output.receive(0)); @@ -125,10 +123,10 @@ public class AggregatorIntegrationTests { } @Test - public void testExpiringAggregator() throws Exception { + public void testExpiringAggregator() { for (int i = 0; i < 5; i++) { Map headers = stubHeaders(i, 5, 1); - expiringAggregatorInput.send(new GenericMessage(i, headers)); + expiringAggregatorInput.send(new GenericMessage<>(i, headers)); } assertNotNull(output.receive(0)); @@ -136,7 +134,7 @@ public class AggregatorIntegrationTests { for (int i = 5; i < 10; i++) { Map headers = stubHeaders(i, 5, 1); - expiringAggregatorInput.send(new GenericMessage(i, headers)); + expiringAggregatorInput.send(new GenericMessage<>(i, headers)); } assertNotNull(output.receive(0)); @@ -148,7 +146,7 @@ public class AggregatorIntegrationTests { public void testGroupTimeoutScheduling() throws Exception { for (int i = 0; i < 5; i++) { Map headers = stubHeaders(i, 5, 1); - this.groupTimeoutAggregatorInput.send(new GenericMessage(i, headers)); + this.groupTimeoutAggregatorInput.send(new GenericMessage<>(i, headers)); //Wait until 'group-timeout' does its stuff. MessageGroupStore mgs = TestUtils.getPropertyValue(this.context.getBean("gta.handler"), "messageStore", @@ -166,11 +164,11 @@ public class AggregatorIntegrationTests { @Test public void testGroupTimeoutReschedulingOnMessageDeliveryException() throws Exception { for (int i = 0; i < 5; i++) { - this.output.send(new GenericMessage("fake message")); + this.output.send(new GenericMessage<>("fake message")); } Map headers = stubHeaders(1, 2, 1); - this.groupTimeoutAggregatorInput.send(new GenericMessage(1, headers)); + this.groupTimeoutAggregatorInput.send(new GenericMessage<>(1, headers)); //Wait until 'group-timeout' does its stuff. MessageGroupStore mgs = TestUtils.getPropertyValue(this.context.getBean("gta.handler"), "messageStore", @@ -190,14 +188,14 @@ public class AggregatorIntegrationTests { } @Test - public void testGroupTimeoutExpressionScheduling() throws Exception { - // Since group-timeout-expression="size() >= 2 ? 100 : -1". The first message won't be scheduled to 'forceComplete' - this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage(1, stubHeaders(1, 6, 1))); + public void testGroupTimeoutExpressionScheduling() { + // Since group-timeout-expression="size() >= 2 ? 100 : null". The first message won't be scheduled to 'forceComplete' + this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage<>(1, stubHeaders(1, 6, 1))); assertNull(this.output.receive(0)); assertNull(this.discard.receive(0)); // As far as 'group.size() >= 2' it will be scheduled to 'forceComplete' - this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage(2, stubHeaders(2, 6, 1))); + this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage<>(2, stubHeaders(2, 6, 1))); assertNull(this.output.receive(0)); Message receive = this.output.receive(10000); assertNotNull(receive); @@ -205,15 +203,15 @@ public class AggregatorIntegrationTests { assertNull(this.discard.receive(0)); // The same with these three messages - this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage(3, stubHeaders(3, 6, 1))); + this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage<>(3, stubHeaders(3, 6, 1))); assertNull(this.output.receive(0)); assertNull(this.discard.receive(0)); - this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage(4, stubHeaders(4, 6, 1))); + this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage<>(4, stubHeaders(4, 6, 1))); assertNull(this.output.receive(0)); assertNull(this.discard.receive(0)); - this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage(5, stubHeaders(5, 6, 1))); + this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage<>(5, stubHeaders(5, 6, 1))); assertNull(this.output.receive(0)); receive = this.output.receive(10000); assertNotNull(receive); @@ -221,7 +219,7 @@ public class AggregatorIntegrationTests { assertNull(this.discard.receive(0)); // The last message in the sequence - normal release by provided 'ReleaseStrategy' - this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage(6, stubHeaders(6, 6, 1))); + this.groupTimeoutExpressionAggregatorInput.send(new GenericMessage<>(6, stubHeaders(6, 6, 1))); receive = this.output.receive(10000); assertNotNull(receive); assertEquals(1, ((Collection) receive.getPayload()).size()); @@ -229,17 +227,17 @@ public class AggregatorIntegrationTests { } @Test - public void testZeroGroupTimeoutExpressionScheduling() throws Exception { + public void testZeroGroupTimeoutExpressionScheduling() { try { this.output.purge(null); this.errors.purge(null); - GenericMessage message = new GenericMessage("foo"); + GenericMessage message = new GenericMessage<>("foo"); this.output.send(message); this.output.send(message); this.output.send(message); this.output.send(message); this.output.send(message); - this.zeroGroupTimeoutExpressionAggregatorInput.send(new GenericMessage(1, stubHeaders(1, 2, 1))); + this.zeroGroupTimeoutExpressionAggregatorInput.send(new GenericMessage<>(1, stubHeaders(1, 2, 1))); ErrorMessage em = (ErrorMessage) this.errors.receive(10000); assertNotNull(em); assertThat(em.getPayload().getMessage().toLowerCase(), @@ -263,7 +261,7 @@ public class AggregatorIntegrationTests { } private Map stubHeaders(int sequenceNumber, int sequenceSize, int correlationId) { - Map headers = new HashMap(); + Map headers = new HashMap<>(); headers.put(IntegrationMessageHeaderAccessor.SEQUENCE_NUMBER, sequenceNumber); headers.put(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE, sequenceSize); headers.put(IntegrationMessageHeaderAccessor.CORRELATION_ID, correlationId); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/PriorityChannelTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/PriorityChannelTests.java index dcecf3d752..e2f00dd6f6 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/PriorityChannelTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/PriorityChannelTests.java @@ -38,25 +38,26 @@ import org.springframework.messaging.support.GenericMessage; /** * @author Mark Fisher * @author Gary Russell + * @author Artem Bilan */ public class PriorityChannelTests { @Test public void testCapacityEnforced() { PriorityChannel channel = new PriorityChannel(3); - assertTrue(channel.send(new GenericMessage("test1"), 0)); - assertTrue(channel.send(new GenericMessage("test2"), 0)); - assertTrue(channel.send(new GenericMessage("test3"), 0)); - assertFalse(channel.send(new GenericMessage("test4"), 0)); + assertTrue(channel.send(new GenericMessage<>("test1"), 0)); + assertTrue(channel.send(new GenericMessage<>("test2"), 0)); + assertTrue(channel.send(new GenericMessage<>("test3"), 0)); + assertFalse(channel.send(new GenericMessage<>("test4"), 0)); channel.receive(0); - assertTrue(channel.send(new GenericMessage("test5"))); + assertTrue(channel.send(new GenericMessage<>("test5"))); } @Test - public void testDefaultComparatorWithTimestampFallback() throws Exception { + public void testDefaultComparatorWithTimestampFallback() { PriorityChannel channel = new PriorityChannel(); for (int i = 0; i < 1000; i++) { - channel.send(new GenericMessage(i)); + channel.send(new GenericMessage<>(i)); } for (int i = 0; i < 1000; i++) { assertEquals(i, channel.receive().getPayload()); @@ -86,9 +87,9 @@ public class PriorityChannelTests { // although this test has no assertions it results in ConcurrentModificationException // if executed before changes for INT-2508 @Test - public void testPriorityChannelWithConcurrentModification() throws Exception { + public void testPriorityChannelWithConcurrentModification() { final PriorityChannel channel = new PriorityChannel(); - final Message message = new GenericMessage("hello"); + final Message message = new GenericMessage<>("hello"); for (int i = 0; i < 1000; i++) { channel.send(message); new Thread(() -> channel.receive()).start(); @@ -99,11 +100,11 @@ public class PriorityChannelTests { @Test public void testCustomComparator() { PriorityChannel channel = new PriorityChannel(5, new StringPayloadComparator()); - Message messageA = new GenericMessage("A"); - Message messageB = new GenericMessage("B"); - Message messageC = new GenericMessage("C"); - Message messageD = new GenericMessage("D"); - Message messageE = new GenericMessage("E"); + Message messageA = new GenericMessage<>("A"); + Message messageB = new GenericMessage<>("B"); + Message messageC = new GenericMessage<>("C"); + Message messageD = new GenericMessage<>("D"); + Message messageE = new GenericMessage<>("E"); channel.send(messageC); channel.send(messageA); channel.send(messageE); @@ -200,7 +201,7 @@ public class PriorityChannelTests { PriorityChannel channel = new PriorityChannel(5); Message highPriority = createPriorityMessage(5); Message lowPriority = createPriorityMessage(-5); - Message nullPriority = new GenericMessage("test:NULL"); + Message nullPriority = new GenericMessage<>("test:NULL"); channel.send(lowPriority); channel.send(highPriority); channel.send(nullPriority); @@ -214,7 +215,7 @@ public class PriorityChannelTests { PriorityChannel channel = new PriorityChannel(); Message highPriority = createPriorityMessage(5); Message lowPriority = createPriorityMessage(-5); - Message nullPriority = new GenericMessage("test:NULL"); + Message nullPriority = new GenericMessage<>("test:NULL"); channel.send(lowPriority); channel.send(highPriority); channel.send(nullPriority); @@ -228,8 +229,8 @@ public class PriorityChannelTests { final PriorityChannel channel = new PriorityChannel(1); final AtomicBoolean sentSecondMessage = new AtomicBoolean(false); ExecutorService executor = Executors.newSingleThreadScheduledExecutor(); - channel.send(new GenericMessage("test-1")); - executor.execute(() -> sentSecondMessage.set(channel.send(new GenericMessage("test-2"), 10))); + channel.send(new GenericMessage<>("test-1")); + executor.execute(() -> sentSecondMessage.set(channel.send(new GenericMessage<>("test-2"), 10))); assertFalse(sentSecondMessage.get()); executor.shutdown(); @@ -249,15 +250,15 @@ public class PriorityChannelTests { ExecutorService executor = Executors.newSingleThreadScheduledExecutor(); channel.send(new GenericMessage("test-1")); executor.execute(() -> { - sentSecondMessage.set(channel.send(new GenericMessage("test-2"), 3000)); + sentSecondMessage.set(channel.send(new GenericMessage<>("test-2"), 3000)); latch.countDown(); }); assertFalse(sentSecondMessage.get()); - Thread.sleep(500); + Thread.sleep(10); Message message1 = channel.receive(); assertNotNull(message1); assertEquals("test-1", message1.getPayload()); - latch.await(1000, TimeUnit.MILLISECONDS); + assertTrue(latch.await(10000, TimeUnit.MILLISECONDS)); assertTrue(sentSecondMessage.get()); Message message2 = channel.receive(); assertNotNull(message2); @@ -271,10 +272,10 @@ public class PriorityChannelTests { final AtomicBoolean sentSecondMessage = new AtomicBoolean(false); ExecutorService executor = Executors.newSingleThreadScheduledExecutor(); channel.send(new GenericMessage("test-1")); - executor.execute(() -> sentSecondMessage.set(channel.send(new GenericMessage("test-2"), -1))); + executor.execute(() -> sentSecondMessage.set(channel.send(new GenericMessage<>("test-2"), -1))); assertFalse(sentSecondMessage.get()); - Thread.sleep(500); - Message message1 = channel.receive(1000); + Thread.sleep(10); + Message message1 = channel.receive(10000); assertNotNull(message1); assertEquals("test-1", message1.getPayload()); executor.shutdown(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/ControlBusTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/ControlBusTests.java index ca06dc232e..1e5db69abd 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/ControlBusTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/ControlBusTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2015 the original author or authors. + * Copyright 2002-2018 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. @@ -52,6 +52,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; * @author Gunnar Hillert * @author Gary Russell * @author Artem Bilan + * * @since 2.0 */ @ContextConfiguration @@ -92,7 +93,7 @@ public class ControlBusTests { "ControlBusLifecycleTests-context.xml", this.getClass()); MessageChannel inputChannel = context.getBean("inputChannel", MessageChannel.class); PollableChannel outputChannel = context.getBean("outputChannel", PollableChannel.class); - assertNull(outputChannel.receive(1000)); + assertNull(outputChannel.receive(10)); Message message = MessageBuilder.withPayload("@adapter.start()").build(); inputChannel.send(message); assertNotNull(outputChannel.receive(1000)); @@ -109,7 +110,7 @@ public class ControlBusTests { assertEquals(0, result.getPayload()); this.registry.channelToChannelName(new DirectChannel()); // Sleep a bit to be sure that we aren't reaped by registry TTL as 60000 - Thread.sleep(100); + Thread.sleep(10); messagingTemplate.convertAndSend(input, "@integrationHeaderChannelRegistry.size()"); result = this.output.receive(0); assertNotNull(result); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/PollerWithErrorChannelTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/PollerWithErrorChannelTests.java index 9a0730d7ef..c501f4342c 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/PollerWithErrorChannelTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/PollerWithErrorChannelTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2018 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. @@ -17,12 +17,15 @@ package org.springframework.integration.config.xml; import static org.junit.Assert.assertNotNull; -import static org.mockito.Mockito.atLeastOnce; +import static org.junit.Assert.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.BDDMockito.willAnswer; import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.verify; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import org.junit.Test; -import org.mockito.Mockito; import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.integration.channel.QueueChannel; @@ -37,6 +40,7 @@ import org.springframework.messaging.support.GenericMessage; /** * @author Oleg Zhurakousky * @author Gunnar Hillert + * @author Artem Bilan * */ public class PollerWithErrorChannelTests { @@ -54,10 +58,20 @@ public class PollerWithErrorChannelTests { SubscribableChannel errorChannel = ac.getBean("errorChannel", SubscribableChannel.class); MessageHandler handler = mock(MessageHandler.class); + + CountDownLatch handleLatch = new CountDownLatch(1); + + willAnswer(invocation -> { + handleLatch.countDown(); + return null; + }) + .given(handler) + .handleMessage(any(Message.class)); errorChannel.subscribe(handler); adapter.start(); - Thread.sleep(1000); - verify(handler, atLeastOnce()).handleMessage(Mockito.any(Message.class)); + + assertTrue(handleLatch.await(10, TimeUnit.SECONDS)); + adapter.stop(); ac.close(); } @@ -75,7 +89,7 @@ public class PollerWithErrorChannelTests { } @Test - public void testWithErrorChannelAndHeader() throws Exception { + public void testWithErrorChannelAndHeader() { ClassPathXmlApplicationContext ac = new ClassPathXmlApplicationContext("PollerWithErrorChannel-context.xml", this.getClass()); SourcePollingChannelAdapter adapter = ac.getBean("withErrorChannelAndHeader", @@ -89,7 +103,7 @@ public class PollerWithErrorChannelTests { @Test // config the same as above but the error wil come from the send - public void testWithErrorChannelAndHeaderWithSendFailure() throws Exception { + public void testWithErrorChannelAndHeaderWithSendFailure() { ClassPathXmlApplicationContext ac = new ClassPathXmlApplicationContext("PollerWithErrorChannel-context.xml", this.getClass()); SourcePollingChannelAdapter adapter = ac.getBean("withErrorChannelAndHeaderErrorOnSend", @@ -103,12 +117,12 @@ public class PollerWithErrorChannelTests { @Test // INT-1952 - public void testWithErrorChannelAndPollingConsumer() throws Exception { + public void testWithErrorChannelAndPollingConsumer() { ClassPathXmlApplicationContext ac = new ClassPathXmlApplicationContext("PollerWithErrorChannel-context.xml", this.getClass()); MessageChannel serviceWithPollerChannel = ac.getBean("serviceWithPollerChannel", MessageChannel.class); QueueChannel errorChannel = ac.getBean("serviceErrorChannel", QueueChannel.class); - serviceWithPollerChannel.send(new GenericMessage("")); + serviceWithPollerChannel.send(new GenericMessage<>("")); assertNotNull(errorChannel.receive(10000)); ac.close(); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java index e936cc201d..9dcafaddc0 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java @@ -152,7 +152,7 @@ public class PollerAdviceTests { adapter.setAdviceChain(adviceChain); adapter.afterPropertiesSet(); adapter.start(); - assertFalse(latch.await(1, TimeUnit.SECONDS)); + assertFalse(latch.await(10, TimeUnit.MILLISECONDS)); adapter.stop(); skipper.reset(); latch = new CountDownLatch(1); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java index 9a85339226..58c04f17c2 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollingLifecycleTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2014 the original author or authors. + * Copyright 2002-2018 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. @@ -60,7 +60,7 @@ public class PollingLifecycleTests { private final TestErrorHandler errorHandler = new TestErrorHandler(); @Before - public void init() throws Exception { + public void init() { taskScheduler.afterPropertiesSet(); } @@ -68,13 +68,16 @@ public class PollingLifecycleTests { public void ensurePollerTaskStops() throws Exception { final CountDownLatch latch = new CountDownLatch(1); QueueChannel channel = new QueueChannel(); - channel.send(new GenericMessage("foo")); + channel.send(new GenericMessage<>("foo")); + //Has to be an explicit implementation - Mockito cannot mock/spy lambdas MessageHandler handler = Mockito.spy(new MessageHandler() { + @Override public void handleMessage(Message message) throws MessagingException { latch.countDown(); } + }); PollingConsumer consumer = new PollingConsumer(channel, handler); consumer.setTrigger(new PeriodicTrigger(0)); @@ -87,7 +90,7 @@ public class PollingLifecycleTests { Mockito.verify(handler, times(1)).handleMessage(Mockito.any(Message.class)); consumer.stop(); for (int i = 0; i < 10; i++) { - channel.send(new GenericMessage("foo")); + channel.send(new GenericMessage<>("foo")); } Thread.sleep(2000); // give enough time for poller to kick in if it didn't stop properly // we'll still have a natural race condition between call to stop() and poller polling @@ -105,12 +108,16 @@ public class PollingLifecycleTests { PollerMetadata pollerMetadata = new PollerMetadata(); pollerMetadata.setTrigger(new PeriodicTrigger(2000)); adapterFactory.setPollerMetadata(pollerMetadata); + + //Has to be an explicit implementation - Mockito cannot mock/spy lambdas MessageSource source = spy(new MessageSource() { + @Override public Message receive() { latch.countDown(); - return new GenericMessage("hello"); + return new GenericMessage<>("hello"); } + }); adapterFactory.setSource(source); adapterFactory.setOutputChannel(channel); @@ -122,7 +129,7 @@ public class PollingLifecycleTests { assertTrue(latch.await(20, TimeUnit.SECONDS)); assertNotNull(channel.receive(100)); adapter.stop(); - assertNull(channel.receive(1000)); + assertNull(channel.receive(10)); Mockito.verify(source, times(1)).receive(); } @@ -136,20 +143,20 @@ public class PollingLifecycleTests { pollerMetadata.setMaxMessagesPerPoll(-1); pollerMetadata.setTrigger(new PeriodicTrigger(2000)); adapterFactory.setPollerMetadata(pollerMetadata); - final Runnable coughtInterrupted = mock(Runnable.class); + final Runnable caughtInterrupted = mock(Runnable.class); MessageSource source = () -> { try { for (int i = 0; i < 10; i++) { - Thread.sleep(1000); + Thread.sleep(10); latch.countDown(); } } catch (InterruptedException e) { - coughtInterrupted.run(); + caughtInterrupted.run(); } - return new GenericMessage("hello"); + return new GenericMessage<>("hello"); }; adapterFactory.setSource(source); adapterFactory.setOutputChannel(channel); @@ -161,8 +168,8 @@ public class PollingLifecycleTests { assertTrue(latch.await(3000, TimeUnit.SECONDS)); // adapter.stop(); - Thread.sleep(1000); - Mockito.verify(coughtInterrupted, times(1)).run(); + Thread.sleep(10); + Mockito.verify(caughtInterrupted, times(1)).run(); } @Test diff --git a/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayRequiresReplyTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayRequiresReplyTests-context.xml index 6bcf7e80fe..2a16172dc5 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayRequiresReplyTests-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayRequiresReplyTests-context.xml @@ -13,7 +13,7 @@ default-request-channel="requestChannel" default-reply-timeout="3000" service-interface="org.springframework.integration.gateway.GatewayRequiresReplyTests$TestService" /> - + - + - + - + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayRequiresReplyTests.java b/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayRequiresReplyTests.java index 3533708e43..3140eaf75a 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayRequiresReplyTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/gateway/GatewayRequiresReplyTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2018 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. @@ -16,8 +16,8 @@ package org.springframework.integration.gateway; -import static org.junit.Assert.assertNull; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNull; import org.junit.Test; import org.junit.runner.RunWith; @@ -31,6 +31,8 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; /** * @author Mark Fisher * @author Gunnar Hillert + * @author Artem Bilan + * * @since 2.0 */ @ContextConfiguration @@ -43,34 +45,38 @@ public class GatewayRequiresReplyTests { @Test public void replyReceived() { - TestService gateway = (TestService) applicationContext.getBean("gateway"); + TestService gateway = this.applicationContext.getBean("gateway", TestService.class); String result = gateway.test("foo"); assertEquals("bar", result); } @Test(expected = ReplyRequiredException.class) public void noReplyReceived() { - TestService gateway = (TestService) applicationContext.getBean("gateway"); + TestService gateway = this.applicationContext.getBean("gateway", TestService.class); gateway.test("bad"); } @Test public void timedOutGateway() { - TestService gateway = (TestService) applicationContext.getBean("timeoutGateway"); + TestService gateway = this.applicationContext.getBean("timeoutGateway", TestService.class); String result = gateway.test("hello"); assertNull(result); } public interface TestService { + String test(String s); + } public static class LongRunningService { + public String echo(String value) throws Exception { Thread.sleep(5000); return value; } + } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java index 19e459423b..7199850190 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java @@ -425,17 +425,17 @@ public class DelayHandlerTests { It's difficult to test it from real ctx, because any async process from 'inbound-channel-adapter' can't achieve the DelayHandler before the main thread emits 'ContextRefreshedEvent'. */ - public void testRescheduleAndHandleAtTheSameTime() throws Exception { + public void testRescheduleAndHandleAtTheSameTime() { QueueChannel results = new QueueChannel(); delayHandler.setOutputChannel(results); - this.delayHandler.setDefaultDelay(100); + this.delayHandler.setDefaultDelay(10); startDelayerHandler(); this.input.send(new GenericMessage<>("foo")); this.delayHandler.reschedulePersistedMessages(); Message message = results.receive(10000); assertNotNull(message); - message = results.receive(500); + message = results.receive(50); assertNull(message); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/router/RouterConcurrencyTest.java b/spring-integration-core/src/test/java/org/springframework/integration/router/RouterConcurrencyTest.java index abe999610c..88b0a59a42 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/router/RouterConcurrencyTest.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/router/RouterConcurrencyTest.java @@ -42,6 +42,7 @@ import org.springframework.messaging.MessageChannel; /** * @author Gary Russell + * * @since 3.0 * */ @@ -61,11 +62,11 @@ public class RouterConcurrencyTest { protected void setConversionService(ConversionService conversionService) { try { if (count.incrementAndGet() > 1) { - Thread.sleep(2000); + Thread.sleep(20); } super.setConversionService(conversionService); semaphore.release(); - Thread.sleep(1000); + Thread.sleep(10); } catch (InterruptedException e) { Thread.currentThread().interrupt(); @@ -85,8 +86,7 @@ public class RouterConcurrencyTest { router.setBeanFactory(beanFactory); ExecutorService exec = Executors.newFixedThreadPool(2); - final List returns = Collections.synchronizedList( - new ArrayList()); + final List returns = Collections.synchronizedList(new ArrayList<>()); Runnable runnable = () -> { ConversionService requiredConversionService = router.getRequiredConversionService(); returns.add(requiredConversionService); diff --git a/src/reference/asciidoc/aggregator.adoc b/src/reference/asciidoc/aggregator.adoc index c3bb5e20e1..46ad598e89 100644 --- a/src/reference/asciidoc/aggregator.adoc +++ b/src/reference/asciidoc/aggregator.adoc @@ -477,8 +477,8 @@ See <>, <>, <> for This attribute provides a built-in _Time-base Release Strategy_ for the aggregator, when there is a need to emit a partial result (or discard the group), if a new Message does not arrive for the `MessageGroup` within the timeout. When a new Message arrives at the aggregator, any existing `ScheduledFuture` for its `MessageGroup` is canceled. If the `ReleaseStrategy` returns `false` (don't release) and the `groupTimeout > 0` a new task will be scheduled to expire the group. -Setting this attribute to zero is not advised because it will effectively disable the aggregator because every message group will be immediately completed. -It is possible, however to conditionally set it to zero using an expression; see `group-timeout-expression` for information. +Setting this attribute to zero (or negative value) is not advised because it will effectively disable the aggregator because every message group will be immediately completed. +It is possible, however to conditionally set it to zero (or negative value) using an expression; see `group-timeout-expression` for information. The action taken during the completion depends on the `ReleaseStrategy` and the `send-partial-group-on-expiry` attribute. See <> for more information. Mutually exclusive with 'group-timeout-expression' attribute. @@ -486,7 +486,7 @@ Mutually exclusive with 'group-timeout-expression' attribute. <22> The SpEL expression that evaluates to a `groupTimeout` with the `MessageGroup` as the `#root` evaluation context object. Used for scheduling the `MessageGroup` to be forced complete. -If the expression evaluates to null or `< 0`, the completion is not scheduled. +If the expression evaluates to `null`, the completion is not scheduled. If it evaluates to zero, the group is completed immediately on the current thread. In effect, this provides a dynamic `group-timeout` property. See `group-timeout` for more information. diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index d4805b80b3..6a306afef3 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -63,4 +63,11 @@ See <> for more information. A confusing `max-rows-per-poll` property on the JDBC Inbound Channel Adapter and JDBC Outbound Gateway has been deprecated in favor newly introduced `max-rows` property. -See <> for more information. +See <> for more information. + +==== Aggregator Changes + +An aggregator now expires the group immediately, if the `groupTimeout` is evaluated to a negative value. +Only `null` is considered as a signal do nothing for the current message. + +See <> for more information.