From e5d859f1c7b33d603310630a30be447e4c0573c8 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 26 Aug 2020 14:07:38 -0400 Subject: [PATCH] Refinement for Redis Stream tests --- ...activeRedisStreamMessageProducerTests.java | 30 ++++++++++++------- 1 file changed, 19 insertions(+), 11 deletions(-) diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java index 2e8c208749..950b2be6f6 100644 --- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducerTests.java @@ -115,17 +115,21 @@ public class ReactiveRedisStreamMessageProducerTests extends RedisAvailableTests this.redisStreamMessageProducer.setConsumerName(null); this.redisStreamMessageProducer.setReadOffset(ReadOffset.from("0")); this.redisStreamMessageProducer.afterPropertiesSet(); + + StepVerifier stepVerifier = + Flux.from(this.fluxMessageChannel) + .as(StepVerifier::create) + .assertNext(message -> { + assertThat(message.getPayload()).isEqualTo(person); + assertThat(message.getHeaders()).containsKeys(RedisHeaders.STREAM_KEY, + RedisHeaders.STREAM_MESSAGE_ID); + }) + .thenCancel() + .verifyLater(); + this.redisStreamMessageProducer.start(); - Flux.from(this.fluxMessageChannel) - .as(StepVerifier::create) - .assertNext(message -> { - assertThat(message.getPayload()).isEqualTo(person); - assertThat(message.getHeaders()).containsKeys(RedisHeaders.STREAM_KEY, - RedisHeaders.STREAM_MESSAGE_ID); - }) - .thenCancel() - .verify(Duration.ofSeconds(10)); + stepVerifier.verify(Duration.ofSeconds(10)); } @Test @@ -134,8 +138,12 @@ public class ReactiveRedisStreamMessageProducerTests extends RedisAvailableTests Address address = new Address("Winterfell, Westeros"); Person person = new Person(address, "John Snow"); - this.template.opsForStream().createGroup(STREAM_KEY, this.redisStreamMessageProducer.getBeanName()) - .subscribe(); + this.template.opsForStream() + .createGroup(STREAM_KEY, this.redisStreamMessageProducer.getBeanName()) + .as(StepVerifier::create) + .assertNext(message -> assertThat(message).isEqualTo("OK")) + .thenCancel() + .verify(Duration.ofSeconds(10)); this.redisStreamMessageProducer.setCreateConsumerGroup(false); this.redisStreamMessageProducer.setConsumerName(CONSUMER);