From a89edd612170fd8e1a01a91978aa7f6aa5c5004d Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 11 Jun 2020 16:49:32 +0200 Subject: [PATCH] GH-1973 Fix test binder to ensure it does not duplicate messages added test that reproduced and this validated the issue after it was fixed as a side work some tests were migrated to junit 5 Resolves #1973 --- .../AutoconfigurationDisabledTest.java | 5 +- .../stream/test/example/ExampleTest.java | 7 +- .../test/matcher/MessageQueueMatcherTest.java | 2 +- .../stream/binder/test/OutputDestination.java | 8 ++- .../ImplicitFunctionBindingTests.java | 65 ++++++++++++++++--- 5 files changed, 65 insertions(+), 22 deletions(-) diff --git a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/disable/AutoconfigurationDisabledTest.java b/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/disable/AutoconfigurationDisabledTest.java index 18ad79613..b524ec749 100644 --- a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/disable/AutoconfigurationDisabledTest.java +++ b/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/disable/AutoconfigurationDisabledTest.java @@ -18,8 +18,7 @@ package org.springframework.cloud.stream.test.disable; import java.util.concurrent.TimeUnit; -import org.junit.Test; -import org.junit.runner.RunWith; +import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.SpringBootApplication; @@ -32,14 +31,12 @@ import org.springframework.integration.annotation.Transformer; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; import org.springframework.test.annotation.DirtiesContext; -import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import static org.assertj.core.api.Assertions.assertThat; /** * @author Marius Bogoevici */ -@RunWith(SpringJUnit4ClassRunner.class) @SpringBootTest(classes = AutoconfigurationDisabledTest.MyProcessor.class, properties = { "server.port=-1", "spring.cloud.stream.defaultBinder=test", "--spring.cloud.stream.bindings.input.contentType=text/plain", diff --git a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/example/ExampleTest.java b/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/example/ExampleTest.java index ef14959e7..a95887e28 100644 --- a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/example/ExampleTest.java +++ b/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/example/ExampleTest.java @@ -16,8 +16,7 @@ package org.springframework.cloud.stream.test.example; -import org.junit.Test; -import org.junit.runner.RunWith; +import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.SpringBootApplication; @@ -29,7 +28,6 @@ import org.springframework.integration.annotation.Transformer; import org.springframework.messaging.Message; import org.springframework.messaging.support.GenericMessage; import org.springframework.test.annotation.DirtiesContext; -import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import static org.assertj.core.api.Assertions.assertThat; @@ -38,12 +36,9 @@ import static org.assertj.core.api.Assertions.assertThat; * {@link org.springframework.cloud.stream.test.binder.TestSupportBinder} applies * correctly. */ -@RunWith(SpringJUnit4ClassRunner.class) -// @checkstyle:off @SpringBootTest(classes = ExampleTest.MyProcessor.class, webEnvironment = SpringBootTest.WebEnvironment.NONE, properties = { "--spring.cloud.stream.bindings.input.contentType=text/plain", "--spring.cloud.stream.bindings.output.contentType=text/plain" }) -// @checkstyle:on @DirtiesContext public class ExampleTest { diff --git a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/matcher/MessageQueueMatcherTest.java b/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/matcher/MessageQueueMatcherTest.java index 8af5cd3b1..5202f552a 100644 --- a/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/matcher/MessageQueueMatcherTest.java +++ b/spring-cloud-stream-test-support/src/test/java/org/springframework/cloud/stream/test/matcher/MessageQueueMatcherTest.java @@ -22,7 +22,7 @@ import java.util.concurrent.LinkedBlockingDeque; import java.util.concurrent.TimeUnit; import org.hamcrest.StringDescription; -import org.junit.Test; +import org.junit.jupiter.api.Test; import org.springframework.messaging.Message; import org.springframework.messaging.support.GenericMessage; diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java index 53c7e1cd0..3f2281b25 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/OutputDestination.java @@ -80,9 +80,11 @@ public class OutputDestination extends AbstractDestination { @SuppressWarnings("unchecked") @Override void afterChannelIsSet(int channelIndex, String bindingName) { - BlockingQueue> messageQueue = new LinkedTransferQueue<>(); - this.messageQueues.put(bindingName, messageQueue); - this.getChannelByName(bindingName).subscribe(message -> this.messageQueues.get(bindingName).offer((Message) message)); + if (!this.messageQueues.containsKey(bindingName)) { + BlockingQueue> messageQueue = new LinkedTransferQueue<>(); + this.messageQueues.put(bindingName, messageQueue); + this.getChannelByName(bindingName).subscribe(message -> this.messageQueues.get(bindingName).offer((Message) message)); + } } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 786a14518..e58353ca1 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -27,8 +27,9 @@ import java.util.function.Consumer; import java.util.function.Function; import java.util.function.Supplier; -import org.junit.After; -import org.junit.Test; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import reactor.core.publisher.EmitterProcessor; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -68,7 +69,7 @@ import static org.junit.Assert.fail; */ public class ImplicitFunctionBindingTests { - @After + @AfterEach public void after() { System.clearProperty("spring.cloud.function.definition"); } @@ -382,7 +383,7 @@ public class ImplicitFunctionBindingTests { } } - @Test(expected = Exception.class) + @Test public void testDeclaredTypeVsActualInstance() { System.clearProperty("spring.cloud.function.definition"); try (ConfigurableApplicationContext context = new SpringApplicationBuilder( @@ -394,6 +395,10 @@ public class ImplicitFunctionBindingTests { Message inputMessageOne = MessageBuilder.withPayload("Hello".getBytes()).build(); inputDestination.send(inputMessageOne); + fail(); + } + catch (Exception ex) { + // good } } @@ -759,13 +764,57 @@ public class ImplicitFunctionBindingTests { } } - @Test(expected = BeanCreationException.class) + @Test public void testReactiveConsumerWithConcurrencyFailureConfiguration() { System.clearProperty("spring.cloud.function.definition"); - new SpringApplicationBuilder( - TestChannelBinderConfiguration.getCompleteConfiguration(ReactiveConsumerWithConcurrencyFailureConfiguration.class)) + try { + new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(ReactiveConsumerWithConcurrencyFailureConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false", + "--spring.cloud.stream.bindings.input-in-0.consumer.concurrency=2"); + fail(); + } + catch (BeanCreationException e) { + // good + } + } + + @Test + public void testGh1973() { + System.clearProperty("spring.cloud.function.definition"); + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(SupplierAndProcessorConfiguration.class)) .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false", - "--spring.cloud.stream.bindings.input-in-0.consumer.concurrency=2"); + "--spring.cloud.function.definition=echo;supplier", + "--spring.cloud.stream.bindings.supplier-out-0.destination=output", + "--spring.cloud.stream.bindings.echo-out-0.destination=output")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context.getBean(OutputDestination.class); + + inputDestination.send(MessageBuilder.withPayload("hello").build()); + assertThat(outputDestination.receive(1000, "output")).isNotNull(); + assertThat(outputDestination.receive(1000, "output")).isNull(); + assertThat(outputDestination.receive(1000, "output")).isNull(); + + } + } + + @EnableAutoConfiguration + public static class SupplierAndProcessorConfiguration { + EmitterProcessor> processor = EmitterProcessor.create(); + + @Bean + public Supplier>> supplier() { + return () -> processor.doOnNext(v -> { + System.out.println("Hello " + v); + }); + } + + @Bean + public Function, Message> echo() { + return v -> v; + } } @EnableAutoConfiguration