From 8e7c1067f351fa5c9199b3cdfe5be4b12372ee76 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Mon, 2 Dec 2019 14:03:12 -0500 Subject: [PATCH] Fix pub/sub race conditions in Reactive tests --- .../reactivestreams/ReactiveStreamsTests.java | 11 +-- .../splitter/DefaultSplitterTests.java | 69 +++++++++++-------- .../file/splitter/FileSplitterTests.java | 52 +++++++------- 3 files changed, 71 insertions(+), 61 deletions(-) diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java index d0ff48e767..b916d3e7cc 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java @@ -152,11 +152,12 @@ public class ReactiveStreamsTests { @Test void testFromPublisher() { - Flux> messageFlux = Flux.just("1,2,3,4") - .map(v -> v.split(",")) - .flatMapIterable(Arrays::asList) - .map(Integer::parseInt) - .map(GenericMessage::new); + Flux> messageFlux = + Flux.just("1,2,3,4") + .map(v -> v.split(",")) + .flatMapIterable(Arrays::asList) + .map(Integer::parseInt) + .map(GenericMessage::new); QueueChannel resultChannel = new QueueChannel(); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/splitter/DefaultSplitterTests.java b/spring-integration-core/src/test/java/org/springframework/integration/splitter/DefaultSplitterTests.java index 452cc9dfaa..4977d7295b 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/splitter/DefaultSplitterTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/splitter/DefaultSplitterTests.java @@ -161,63 +161,74 @@ class DefaultSplitterTests { void splitArrayPayloadReactive() { Message message = new GenericMessage<>(new String[] { "x", "y", "z" }); FluxMessageChannel replyChannel = new FluxMessageChannel(); - DefaultMessageSplitter splitter = new DefaultMessageSplitter(); - splitter.setOutputChannel(replyChannel); - - splitter.handleMessage(message); Flux testFlux = Flux.from(replyChannel) .map(Message::getPayload) .cast(String.class); - StepVerifier.create(testFlux) - .expectNext("x", "y", "z") - .expectNoEvent(Duration.ofMillis(100)) - .thenCancel() - .verify(Duration.ofSeconds(1)); + StepVerifier verifier = + StepVerifier.create(testFlux) + .expectNext("x", "y", "z") + .expectNoEvent(Duration.ofMillis(100)) + .thenCancel() + .verifyLater(); + + DefaultMessageSplitter splitter = new DefaultMessageSplitter(); + splitter.setOutputChannel(replyChannel); + + splitter.handleMessage(message); + + verifier.verify(Duration.ofSeconds(1)); } @Test void splitStreamReactive() { Message message = new GenericMessage<>(Stream.of("x", "y", "z")); FluxMessageChannel replyChannel = new FluxMessageChannel(); - DefaultMessageSplitter splitter = new DefaultMessageSplitter(); - splitter.setOutputChannel(replyChannel); - - splitter.handleMessage(message); - Flux testFlux = Flux.from(replyChannel) .map(Message::getPayload) .cast(String.class); - StepVerifier.create(testFlux) - .expectNext("x", "y", "z") - .expectNoEvent(Duration.ofMillis(100)) - .thenCancel() - .verify(Duration.ofSeconds(1)); + StepVerifier verifier = + StepVerifier.create(testFlux) + .expectNext("x", "y", "z") + .expectNoEvent(Duration.ofMillis(100)) + .thenCancel() + .verifyLater(); + + DefaultMessageSplitter splitter = new DefaultMessageSplitter(); + splitter.setOutputChannel(replyChannel); + + splitter.handleMessage(message); + + verifier.verify(Duration.ofSeconds(1)); } @Test void splitFluxReactive() { Message message = new GenericMessage<>(Flux.just("x", "y", "z")); FluxMessageChannel replyChannel = new FluxMessageChannel(); - DefaultMessageSplitter splitter = new DefaultMessageSplitter(); - splitter.setOutputChannel(replyChannel); - - splitter.handleMessage(message); - Flux testFlux = Flux.from(replyChannel) .map(Message::getPayload) .cast(String.class); - StepVerifier.create(testFlux) - .expectNext("x", "y", "z") - .expectNoEvent(Duration.ofMillis(100)) - .thenCancel() - .verify(Duration.ofSeconds(1)); + StepVerifier verifier = + StepVerifier.create(testFlux) + .expectNext("x", "y", "z") + .expectNoEvent(Duration.ofMillis(100)) + .thenCancel() + .verifyLater(); + + DefaultMessageSplitter splitter = new DefaultMessageSplitter(); + splitter.setOutputChannel(replyChannel); + + splitter.handleMessage(message); + + + verifier.verify(Duration.ofSeconds(1)); } } diff --git a/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java b/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java index 7690448c96..cdceabfa43 100644 --- a/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java +++ b/spring-integration-file/src/test/java/org/springframework/integration/file/splitter/FileSplitterTests.java @@ -58,9 +58,7 @@ import org.springframework.messaging.MessagingException; import org.springframework.messaging.PollableChannel; import org.springframework.messaging.support.GenericMessage; import org.springframework.test.annotation.DirtiesContext; -import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; -import org.springframework.test.context.support.AnnotationConfigContextLoader; import org.springframework.util.FileCopyUtils; import reactor.test.StepVerifier; @@ -72,8 +70,7 @@ import reactor.test.StepVerifier; * * @since 4.1.2 */ -@ContextConfiguration(loader = AnnotationConfigContextLoader.class) -@SpringJUnitConfig +@SpringJUnitConfig(FileSplitterTests.ContextConfiguration.class) @DirtiesContext public class FileSplitterTests { @@ -265,34 +262,35 @@ public class FileSplitterTests { @Test void testFileSplitterReactive() { FluxMessageChannel outputChannel = new FluxMessageChannel(); + StepVerifier verifier = + StepVerifier.create(outputChannel) + .assertNext(m -> { + assertThat(m.getHeaders()) + .containsKey(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE) + .containsEntry(FileHeaders.MARKER, "START"); + assertThat(m.getPayload()).isInstanceOf(FileMarker.class); + FileMarker fileMarker = (FileMarker) m.getPayload(); + assertThat(fileMarker.getMark()).isEqualTo(FileMarker.Mark.START); + assertThat(fileMarker.getFilePath()).isEqualTo(file.getAbsolutePath()); + }) + .expectNextCount(2) + .assertNext(m -> { + assertThat(m.getHeaders()).containsEntry(FileHeaders.MARKER, "END"); + assertThat(m.getPayload()).isInstanceOf(FileMarker.class); + FileMarker fileMarker = (FileMarker) m.getPayload(); + assertThat(fileMarker.getMark()).isEqualTo(FileMarker.Mark.END); + assertThat(fileMarker.getFilePath()).isEqualTo(file.getAbsolutePath()); + assertThat(fileMarker.getLineCount()).isEqualTo(2); + }) + .expectNoEvent(Duration.ofMillis(100)) + .thenCancel() + .verifyLater(); FileSplitter splitter = new FileSplitter(true, true); splitter.setApplySequence(true); splitter.setOutputChannel(outputChannel); splitter.handleMessage(new GenericMessage<>(file)); - - StepVerifier.create(outputChannel) - .assertNext(m -> { - assertThat(m.getHeaders()) - .containsKey(IntegrationMessageHeaderAccessor.SEQUENCE_SIZE) - .containsEntry(FileHeaders.MARKER, "START"); - assertThat(m.getPayload()).isInstanceOf(FileSplitter.FileMarker.class); - FileMarker fileMarker = (FileSplitter.FileMarker) m.getPayload(); - assertThat(fileMarker.getMark()).isEqualTo(FileMarker.Mark.START); - assertThat(fileMarker.getFilePath()).isEqualTo(file.getAbsolutePath()); - }) - .expectNextCount(2) - .assertNext(m -> { - assertThat(m.getHeaders()).containsEntry(FileHeaders.MARKER, "END"); - assertThat(m.getPayload()).isInstanceOf(FileSplitter.FileMarker.class); - FileMarker fileMarker = (FileSplitter.FileMarker) m.getPayload(); - assertThat(fileMarker.getMark()).isEqualTo(FileMarker.Mark.END); - assertThat(fileMarker.getFilePath()).isEqualTo(file.getAbsolutePath()); - assertThat(fileMarker.getLineCount()).isEqualTo(2); - }) - .expectNoEvent(Duration.ofMillis(100)) - .thenCancel() - .verify(Duration.ofSeconds(1)); + verifier.verify(Duration.ofSeconds(1)); } @Test