From 3cb13af813128ccf811c639854f4376f930af1d6 Mon Sep 17 00:00:00 2001 From: Attoumane Date: Tue, 11 Aug 2020 15:30:09 +0200 Subject: [PATCH] More tests for Redis Stream support Related to https://github.com/spring-projects/spring-integration/issues/3226 --- .../ReactiveRedisStreamMessageProducer.java | 4 +-- ...activeRedisStreamMessageProducerTests.java | 29 +++++++++++++++++-- 2 files changed, 29 insertions(+), 4 deletions(-) diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java index b1281b9369..6ab5a15e02 100644 --- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/inbound/ReactiveRedisStreamMessageProducer.java @@ -44,7 +44,7 @@ import reactor.core.publisher.Mono; * output channel. * By default this adapter reads message as a standalone client {@code XREAD} (Redis command) but can be switched to a * Consumer Group feature {@code XREADGROUP} by setting {@link #consumerName} field. - * By default the Consumer Group name is an id of this bean {@link #getBeanName()}. + * By default the Consumer Group name is the id of this bean {@link #getBeanName()}. * * @author Attoumane Ahamadi * @author Artem Bilan @@ -195,7 +195,7 @@ public class ReactiveRedisStreamMessageProducer extends MessageProducerSupport { Consumer consumer = Consumer.from(this.consumerGroup, this.consumerName); // NOSONAR if (offset.getOffset().equals(ReadOffset.latest())) { - // for consumer group offset id should be equal '>' + // for consumer group offset id should be equal to '>' offset = StreamOffset.create(this.streamKey, ReadOffset.lastConsumed()); } 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 f6b1cbfe9b..2e8c208749 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 @@ -38,6 +38,7 @@ import org.springframework.integration.redis.outbound.ReactiveRedisStreamMessage import org.springframework.integration.redis.rules.RedisAvailable; import org.springframework.integration.redis.rules.RedisAvailableRule; import org.springframework.integration.redis.rules.RedisAvailableTests; +import org.springframework.integration.redis.support.RedisHeaders; import org.springframework.integration.redis.util.Address; import org.springframework.integration.redis.util.Person; import org.springframework.messaging.support.GenericMessage; @@ -118,7 +119,11 @@ public class ReactiveRedisStreamMessageProducerTests extends RedisAvailableTests Flux.from(this.fluxMessageChannel) .as(StepVerifier::create) - .assertNext(message -> assertThat(message.getPayload()).isEqualTo(person)) + .assertNext(message -> { + assertThat(message.getPayload()).isEqualTo(person); + assertThat(message.getHeaders()).containsKeys(RedisHeaders.STREAM_KEY, + RedisHeaders.STREAM_MESSAGE_ID); + }) .thenCancel() .verify(Duration.ofSeconds(10)); } @@ -126,7 +131,27 @@ public class ReactiveRedisStreamMessageProducerTests extends RedisAvailableTests @Test @RedisAvailable public void testReadingMessageAsConsumerInConsumerGroup() { - //TODO find why the test above does not execute before implementing this one + Address address = new Address("Winterfell, Westeros"); + Person person = new Person(address, "John Snow"); + + this.template.opsForStream().createGroup(STREAM_KEY, this.redisStreamMessageProducer.getBeanName()) + .subscribe(); + + this.redisStreamMessageProducer.setCreateConsumerGroup(false); + this.redisStreamMessageProducer.setConsumerName(CONSUMER); + this.redisStreamMessageProducer.afterPropertiesSet(); + this.redisStreamMessageProducer.start(); + + this.messageHandler.handleMessage(new GenericMessage<>(person)); + + Flux.from(this.fluxMessageChannel) + .as(StepVerifier::create) + .assertNext(message -> { + assertThat(message.getPayload()).isEqualTo(person); + assertThat(message.getHeaders()).containsKeys(RedisHeaders.CONSUMER_GROUP, RedisHeaders.CONSUMER); + }) + .thenCancel() + .verify(Duration.ofSeconds(10)); } @Configuration