Add support for XCLAIM in StreamOperations

Closes #2465
This commit is contained in:
Marcin Zielinski
2022-11-26 21:54:16 +01:00
committed by John Blum
parent 92143f7333
commit 321e3f03f9
6 changed files with 165 additions and 11 deletions

View File

@@ -20,8 +20,10 @@ import static org.junit.Assume.*;
import reactor.test.StepVerifier;
import java.time.Duration;
import java.util.Collection;
import java.util.Collections;
import java.util.Map;
import org.junit.jupiter.api.BeforeEach;
@@ -358,4 +360,29 @@ public class DefaultReactiveStreamOperationsIntegrationTests<K, HK, HV> {
}).verifyComplete();
}
@ParameterizedRedisTest // https://github.com/spring-projects/spring-data-redis/issues/2465
void claimShouldReadMessageDetails() {
K key = keyFactory.instance();
HK hashKey = hashKeyFactory.instance();
HV value = valueFactory.instance();
Map<HK, HV> content = Collections.singletonMap(hashKey, value);
RecordId messageId = streamOperations.add(key, content).block();
streamOperations.createGroup(key, ReadOffset.from("0-0"), "my-group").then().as(StepVerifier::create)
.verifyComplete();
streamOperations.read(Consumer.from("my-group", "my-consumer"), StreamOffset.create(key, ReadOffset.lastConsumed()))
.then().as(StepVerifier::create).verifyComplete();
streamOperations.claim(key, "my-group", "name", Duration.ZERO, messageId).as(StepVerifier::create)
.assertNext(claimed -> {
assertThat(claimed.getStream()).isEqualTo(key);
assertThat(claimed.getValue()).isEqualTo(content);
assertThat(claimed.getId()).isEqualTo(messageId);
}).verifyComplete();
}
}

View File

@@ -18,6 +18,7 @@ package org.springframework.data.redis.core;
import static org.assertj.core.api.Assertions.*;
import static org.assertj.core.api.Assumptions.*;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
@@ -34,16 +35,7 @@ import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.connection.jedis.extension.JedisConnectionFactoryExtension;
import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory;
import org.springframework.data.redis.connection.lettuce.extension.LettuceConnectionFactoryExtension;
import org.springframework.data.redis.connection.stream.Consumer;
import org.springframework.data.redis.connection.stream.MapRecord;
import org.springframework.data.redis.connection.stream.ObjectRecord;
import org.springframework.data.redis.connection.stream.PendingMessages;
import org.springframework.data.redis.connection.stream.PendingMessagesSummary;
import org.springframework.data.redis.connection.stream.ReadOffset;
import org.springframework.data.redis.connection.stream.RecordId;
import org.springframework.data.redis.connection.stream.StreamOffset;
import org.springframework.data.redis.connection.stream.StreamReadOptions;
import org.springframework.data.redis.connection.stream.StreamRecords;
import org.springframework.data.redis.connection.stream.*;
import org.springframework.data.redis.test.condition.EnabledOnCommand;
import org.springframework.data.redis.test.condition.EnabledOnRedisDriver;
import org.springframework.data.redis.test.condition.EnabledOnRedisVersion;
@@ -72,7 +64,7 @@ public class DefaultStreamOperationsIntegrationTests<K, HK, HV> {
private final StreamOperations<K, HK, HV> streamOps;
public DefaultStreamOperationsIntegrationTests(RedisTemplate<K, ?> redisTemplate, ObjectFactory<K> keyFactory,
ObjectFactory<?> objectFactory) {
ObjectFactory<?> objectFactory) {
this.redisTemplate = redisTemplate;
this.connectionFactory = redisTemplate.getRequiredConnectionFactory();
@@ -420,4 +412,29 @@ public class DefaultStreamOperationsIntegrationTests<K, HK, HV> {
assertThat(pending.get(0).getConsumerName()).isEqualTo("my-consumer");
assertThat(pending.get(0).getTotalDeliveryCount()).isOne();
}
@ParameterizedRedisTest // https://github.com/spring-projects/spring-data-redis/issues/2465
void claimShouldReadMessageDetails() {
K key = keyFactory.instance();
HK hashKey = hashKeyFactory.instance();
HV value = hashValueFactory.instance();
RecordId messageId = streamOps.add(key, Collections.singletonMap(hashKey, value));
streamOps.createGroup(key, ReadOffset.from("0-0"), "my-group");
streamOps.read(Consumer.from("my-group", "name"), StreamOffset.create(key, ReadOffset.lastConsumed()));
List<MapRecord<K, HK, HV>> messages = streamOps.claim(key, "my-group", "new-owner", Duration.ZERO, messageId);
assertThat(messages).hasSize(1);
MapRecord<K, HK, HV> message = messages.get(0);
assertThat(message.getId()).isEqualTo(messageId);
assertThat(message.getStream()).isEqualTo(key);
if (!(key instanceof byte[] || value instanceof byte[])) {
assertThat(message.getValue()).containsEntry(hashKey, value);
}
}
}