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
This commit is contained in:
@@ -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",
|
||||
|
||||
@@ -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 {
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -80,9 +80,11 @@ public class OutputDestination extends AbstractDestination {
|
||||
@SuppressWarnings("unchecked")
|
||||
@Override
|
||||
void afterChannelIsSet(int channelIndex, String bindingName) {
|
||||
BlockingQueue<Message<byte[]>> messageQueue = new LinkedTransferQueue<>();
|
||||
this.messageQueues.put(bindingName, messageQueue);
|
||||
this.getChannelByName(bindingName).subscribe(message -> this.messageQueues.get(bindingName).offer((Message<byte[]>) message));
|
||||
if (!this.messageQueues.containsKey(bindingName)) {
|
||||
BlockingQueue<Message<byte[]>> messageQueue = new LinkedTransferQueue<>();
|
||||
this.messageQueues.put(bindingName, messageQueue);
|
||||
this.getChannelByName(bindingName).subscribe(message -> this.messageQueues.get(bindingName).offer((Message<byte[]>) message));
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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<byte[]> 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<Message<String>> processor = EmitterProcessor.create();
|
||||
|
||||
@Bean
|
||||
public Supplier<Flux<Message<String>>> supplier() {
|
||||
return () -> processor.doOnNext(v -> {
|
||||
System.out.println("Hello " + v);
|
||||
});
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Function<Message<String>, Message<String>> echo() {
|
||||
return v -> v;
|
||||
}
|
||||
}
|
||||
|
||||
@EnableAutoConfiguration
|
||||
|
||||
Reference in New Issue
Block a user