Fix pub/sub race conditions in Reactive tests

This commit is contained in:
Artem Bilan
2019-12-02 14:03:12 -05:00
parent daa89bf091
commit 8e7c1067f3
3 changed files with 71 additions and 61 deletions

View File

@@ -152,11 +152,12 @@ public class ReactiveStreamsTests {
@Test
void testFromPublisher() {
Flux<Message<?>> messageFlux = Flux.just("1,2,3,4")
.map(v -> v.split(","))
.flatMapIterable(Arrays::asList)
.map(Integer::parseInt)
.map(GenericMessage::new);
Flux<Message<?>> messageFlux =
Flux.just("1,2,3,4")
.map(v -> v.split(","))
.flatMapIterable(Arrays::asList)
.map(Integer::parseInt)
.map(GenericMessage::new);
QueueChannel resultChannel = new QueueChannel();

View File

@@ -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<String> 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<String> 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<String> 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));
}
}