Clean Redis keys before testing
https://build.spring.io/browse/INT-FATS5IC-400 When we run several concurrent builds (e.g. CI) and use the same shared Redis server we may end up with the case when one process reads data populated by another because we use the same key (groupId in our case) * Fix `RedisMessageGroupStoreTests` to use `UUID.randomUUID()` for the `groupId` **Cherry-pick to 4.3.x** # Conflicts: # spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageGroupStoreTests.java
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2007-2017 the original author or authors.
|
||||
* Copyright 2007-2018 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -76,6 +76,8 @@ import junit.framework.AssertionFailedError;
|
||||
*/
|
||||
public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
|
||||
private final UUID groupId = UUID.randomUUID();
|
||||
|
||||
@Bean
|
||||
public RedisConnectionFactory redisConnectionFactory() {
|
||||
return getConnectionFactoryForTest();
|
||||
@@ -90,11 +92,11 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testNonExistingEmptyMessageGroup() throws Exception {
|
||||
public void testNonExistingEmptyMessageGroup() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store = new RedisMessageStore(jcf);
|
||||
|
||||
MessageGroup messageGroup = store.getMessageGroup(1);
|
||||
MessageGroup messageGroup = store.getMessageGroup(this.groupId);
|
||||
assertNotNull(messageGroup);
|
||||
assertTrue(messageGroup instanceof SimpleMessageGroup);
|
||||
assertEquals(0, messageGroup.size());
|
||||
@@ -106,15 +108,15 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store = new RedisMessageStore(jcf);
|
||||
|
||||
Message<?> message = new GenericMessage<String>("Hello");
|
||||
MessageGroup messageGroup = store.addMessageToGroup(1, message);
|
||||
Message<?> message = new GenericMessage<>("Hello");
|
||||
MessageGroup messageGroup = store.addMessageToGroup(this.groupId, message);
|
||||
assertEquals(1, messageGroup.size());
|
||||
long createdTimestamp = messageGroup.getTimestamp();
|
||||
long updatedTimestamp = messageGroup.getLastModified();
|
||||
assertEquals(createdTimestamp, updatedTimestamp);
|
||||
Thread.sleep(10);
|
||||
message = new GenericMessage<String>("Hello");
|
||||
messageGroup = store.addMessageToGroup(1, message);
|
||||
message = new GenericMessage<>("Hello");
|
||||
messageGroup = store.addMessageToGroup(this.groupId, message);
|
||||
createdTimestamp = messageGroup.getTimestamp();
|
||||
updatedTimestamp = messageGroup.getLastModified();
|
||||
assertTrue(updatedTimestamp > createdTimestamp);
|
||||
@@ -122,40 +124,40 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
// make sure the store is properly rebuild from Redis
|
||||
store = new RedisMessageStore(jcf);
|
||||
|
||||
messageGroup = store.getMessageGroup(1);
|
||||
messageGroup = store.getMessageGroup(this.groupId);
|
||||
assertEquals(2, messageGroup.size());
|
||||
}
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testMessageGroupWithAddedMessage() throws Exception {
|
||||
public void testMessageGroupWithAddedMessage() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store = new RedisMessageStore(jcf);
|
||||
|
||||
Message<?> message = new GenericMessage<String>("Hello");
|
||||
MessageGroup messageGroup = store.addMessageToGroup(1, message);
|
||||
Message<?> message = new GenericMessage<>("Hello");
|
||||
MessageGroup messageGroup = store.addMessageToGroup(this.groupId, message);
|
||||
assertEquals(1, messageGroup.size());
|
||||
|
||||
// make sure the store is properly rebuild from Redis
|
||||
store = new RedisMessageStore(jcf);
|
||||
|
||||
messageGroup = store.getMessageGroup(1);
|
||||
messageGroup = store.getMessageGroup(this.groupId);
|
||||
assertEquals(1, messageGroup.size());
|
||||
}
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testRemoveMessageGroup() throws Exception {
|
||||
public void testRemoveMessageGroup() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store = new RedisMessageStore(jcf);
|
||||
|
||||
MessageGroup messageGroup = store.getMessageGroup(1);
|
||||
Message<?> message = new GenericMessage<String>("Hello");
|
||||
MessageGroup messageGroup = store.getMessageGroup(this.groupId);
|
||||
Message<?> message = new GenericMessage<>("Hello");
|
||||
messageGroup = store.addMessageToGroup(messageGroup.getGroupId(), message);
|
||||
assertEquals(1, messageGroup.size());
|
||||
|
||||
store.removeMessageGroup(1);
|
||||
MessageGroup messageGroupA = store.getMessageGroup(1);
|
||||
store.removeMessageGroup(this.groupId);
|
||||
MessageGroup messageGroupA = store.getMessageGroup(this.groupId);
|
||||
assertNotSame(messageGroup, messageGroupA);
|
||||
// assertEquals(0, messageGroupA.getMarked().size());
|
||||
assertEquals(0, messageGroupA.getMessages().size());
|
||||
@@ -164,7 +166,7 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
// make sure the store is properly rebuild from Redis
|
||||
store = new RedisMessageStore(jcf);
|
||||
|
||||
messageGroup = store.getMessageGroup(1);
|
||||
messageGroup = store.getMessageGroup(this.groupId);
|
||||
|
||||
assertEquals(0, messageGroup.getMessages().size());
|
||||
assertEquals(0, messageGroup.size());
|
||||
@@ -172,64 +174,62 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testCompleteMessageGroup() throws Exception {
|
||||
public void testCompleteMessageGroup() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store = new RedisMessageStore(jcf);
|
||||
|
||||
MessageGroup messageGroup = store.getMessageGroup(1);
|
||||
Message<?> message = new GenericMessage<String>("Hello");
|
||||
MessageGroup messageGroup = store.getMessageGroup(this.groupId);
|
||||
Message<?> message = new GenericMessage<>("Hello");
|
||||
messageGroup = store.addMessageToGroup(messageGroup.getGroupId(), message);
|
||||
store.completeGroup(messageGroup.getGroupId());
|
||||
messageGroup = store.getMessageGroup(1);
|
||||
messageGroup = store.getMessageGroup(this.groupId);
|
||||
assertTrue(messageGroup.isComplete());
|
||||
}
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testLastReleasedSequenceNumber() throws Exception {
|
||||
public void testLastReleasedSequenceNumber() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store = new RedisMessageStore(jcf);
|
||||
|
||||
MessageGroup messageGroup = store.getMessageGroup(1);
|
||||
Message<?> message = new GenericMessage<String>("Hello");
|
||||
MessageGroup messageGroup = store.getMessageGroup(this.groupId);
|
||||
Message<?> message = new GenericMessage<>("Hello");
|
||||
messageGroup = store.addMessageToGroup(messageGroup.getGroupId(), message);
|
||||
store.setLastReleasedSequenceNumberForGroup(messageGroup.getGroupId(), 5);
|
||||
messageGroup = store.getMessageGroup(1);
|
||||
messageGroup = store.getMessageGroup(this.groupId);
|
||||
assertEquals(5, messageGroup.getLastReleasedMessageSequenceNumber());
|
||||
}
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testRemoveMessageFromTheGroup() throws Exception {
|
||||
public void testRemoveMessageFromTheGroup() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store = new RedisMessageStore(jcf);
|
||||
|
||||
MessageGroup messageGroup = store.getMessageGroup(1);
|
||||
Message<?> message = new GenericMessage<String>("2");
|
||||
store.addMessagesToGroup(messageGroup.getGroupId(), new GenericMessage<String>("1"), message);
|
||||
messageGroup = store.addMessageToGroup(messageGroup.getGroupId(), new GenericMessage<String>("3"));
|
||||
MessageGroup messageGroup = store.getMessageGroup(this.groupId);
|
||||
Message<?> message = new GenericMessage<>("2");
|
||||
store.addMessagesToGroup(messageGroup.getGroupId(), new GenericMessage<>("1"), message);
|
||||
messageGroup = store.addMessageToGroup(messageGroup.getGroupId(), new GenericMessage<>("3"));
|
||||
assertEquals(3, messageGroup.size());
|
||||
|
||||
store.removeMessagesFromGroup(1, message);
|
||||
messageGroup = store.getMessageGroup(1);
|
||||
store.removeMessagesFromGroup(this.groupId, message);
|
||||
messageGroup = store.getMessageGroup(this.groupId);
|
||||
assertEquals(2, messageGroup.size());
|
||||
|
||||
// make sure the store is properly rebuild from Redis
|
||||
store = new RedisMessageStore(jcf);
|
||||
|
||||
messageGroup = store.getMessageGroup(1);
|
||||
messageGroup = store.getMessageGroup(this.groupId);
|
||||
assertEquals(2, messageGroup.size());
|
||||
}
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testWithMessageHistory() throws Exception {
|
||||
public void testWithMessageHistory() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store = new RedisMessageStore(jcf);
|
||||
|
||||
store.getMessageGroup(1);
|
||||
|
||||
Message<?> message = new GenericMessage<String>("Hello");
|
||||
Message<?> message = new GenericMessage<>("Hello");
|
||||
DirectChannel fooChannel = new DirectChannel();
|
||||
fooChannel.setBeanName("fooChannel");
|
||||
DirectChannel barChannel = new DirectChannel();
|
||||
@@ -237,9 +237,9 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
|
||||
message = MessageHistory.write(message, fooChannel);
|
||||
message = MessageHistory.write(message, barChannel);
|
||||
store.addMessagesToGroup(1, message);
|
||||
store.addMessagesToGroup(this.groupId, message);
|
||||
|
||||
message = store.getMessageGroup(1).getMessages().iterator().next();
|
||||
message = store.getMessageGroup(this.groupId).getMessages().iterator().next();
|
||||
|
||||
MessageHistory messageHistory = MessageHistory.read(message);
|
||||
assertNotNull(messageHistory);
|
||||
@@ -251,57 +251,64 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testRemoveNonExistingMessageFromTheGroup() throws Exception {
|
||||
public void testRemoveNonExistingMessageFromTheGroup() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store = new RedisMessageStore(jcf);
|
||||
|
||||
MessageGroup messageGroup = store.getMessageGroup(1);
|
||||
store.addMessagesToGroup(messageGroup.getGroupId(), new GenericMessage<String>("1"));
|
||||
store.removeMessagesFromGroup(1, new GenericMessage<String>("2"));
|
||||
MessageGroup messageGroup = store.getMessageGroup(this.groupId);
|
||||
store.addMessagesToGroup(messageGroup.getGroupId(), new GenericMessage<>("1"));
|
||||
store.removeMessagesFromGroup(this.groupId, new GenericMessage<>("2"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testRemoveNonExistingMessageFromNonExistingTheGroup() throws Exception {
|
||||
public void testRemoveNonExistingMessageFromNonExistingTheGroup() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store = new RedisMessageStore(jcf);
|
||||
store.removeMessagesFromGroup(1, new GenericMessage<String>("2"));
|
||||
store.removeMessagesFromGroup(this.groupId, new GenericMessage<>("2"));
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testMultipleInstancesOfGroupStore() throws Exception {
|
||||
public void testMultipleInstancesOfGroupStore() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store1 = new RedisMessageStore(jcf);
|
||||
|
||||
RedisMessageStore store2 = new RedisMessageStore(jcf);
|
||||
|
||||
Message<?> message = new GenericMessage<String>("1");
|
||||
store1.addMessagesToGroup(1, message);
|
||||
MessageGroup messageGroup = store2.addMessageToGroup(1, new GenericMessage<String>("2"));
|
||||
store1.removeMessageGroup(this.groupId);
|
||||
|
||||
Message<?> message = new GenericMessage<>("1");
|
||||
store1.addMessagesToGroup(this.groupId, message);
|
||||
MessageGroup messageGroup = store2.addMessageToGroup(this.groupId, new GenericMessage<>("2"));
|
||||
|
||||
assertEquals(2, messageGroup.getMessages().size());
|
||||
|
||||
RedisMessageStore store3 = new RedisMessageStore(jcf);
|
||||
|
||||
store3.removeMessagesFromGroup(1, message);
|
||||
messageGroup = store3.getMessageGroup(1);
|
||||
store3.removeMessagesFromGroup(this.groupId, message);
|
||||
messageGroup = store3.getMessageGroup(this.groupId);
|
||||
|
||||
assertEquals(1, messageGroup.getMessages().size());
|
||||
}
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testIteratorOfMessageGroups() throws Exception {
|
||||
public void testIteratorOfMessageGroups() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store1 = new RedisMessageStore(jcf);
|
||||
RedisMessageStore store2 = new RedisMessageStore(jcf);
|
||||
|
||||
store1.removeMessageGroup(this.groupId);
|
||||
UUID group2 = UUID.randomUUID();
|
||||
store1.removeMessageGroup(group2);
|
||||
UUID group3 = UUID.randomUUID();
|
||||
store1.removeMessageGroup(group3);
|
||||
|
||||
store1.addMessagesToGroup(1, new GenericMessage<String>("1"));
|
||||
store2.addMessagesToGroup(2, new GenericMessage<String>("2"));
|
||||
store1.addMessagesToGroup(3, new GenericMessage<String>("3"), new GenericMessage<String>("3A"));
|
||||
store1.addMessagesToGroup(this.groupId, new GenericMessage<>("1"));
|
||||
store2.addMessagesToGroup(group2, new GenericMessage<>("2"));
|
||||
store1.addMessagesToGroup(group3, new GenericMessage<>("3"), new GenericMessage<>("3A"));
|
||||
|
||||
Iterator<MessageGroup> messageGroups = store1.iterator();
|
||||
int counter = 0;
|
||||
@@ -321,7 +328,7 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
}
|
||||
assertEquals(3, counter);
|
||||
|
||||
store2.removeMessageGroup(3);
|
||||
store2.removeMessageGroup(group3);
|
||||
|
||||
messageGroups = store1.iterator();
|
||||
counter = 0;
|
||||
@@ -340,36 +347,30 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
final RedisMessageStore store1 = new RedisMessageStore(jcf);
|
||||
final RedisMessageStore store2 = new RedisMessageStore(jcf);
|
||||
|
||||
final Message<?> message = new GenericMessage<String>("1");
|
||||
store1.removeMessageGroup(this.groupId);
|
||||
|
||||
final Message<?> message = new GenericMessage<>("1");
|
||||
|
||||
ExecutorService executor = null;
|
||||
|
||||
final List<Object> failures = new ArrayList<Object>();
|
||||
final List<Object> failures = new ArrayList<>();
|
||||
|
||||
for (int i = 0; i < 100; i++) {
|
||||
executor = Executors.newCachedThreadPool();
|
||||
|
||||
executor.execute(new Runnable() {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
MessageGroup group = store1.addMessageToGroup(1, message);
|
||||
if (group.getMessages().size() != 1) {
|
||||
failures.add("ADD");
|
||||
throw new AssertionFailedError("Failed on ADD");
|
||||
}
|
||||
executor.execute(() -> {
|
||||
MessageGroup group = store1.addMessageToGroup(this.groupId, message);
|
||||
if (group.getMessages().size() != 1) {
|
||||
failures.add("ADD");
|
||||
throw new AssertionFailedError("Failed on ADD");
|
||||
}
|
||||
});
|
||||
executor.execute(new Runnable() {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
store2.removeMessagesFromGroup(1, message);
|
||||
MessageGroup group = store2.getMessageGroup(1);
|
||||
if (group.getMessages().size() != 0) {
|
||||
failures.add("REMOVE");
|
||||
throw new AssertionFailedError("Failed on Remove");
|
||||
}
|
||||
executor.execute(() -> {
|
||||
store2.removeMessagesFromGroup(this.groupId, message);
|
||||
MessageGroup group = store2.getMessageGroup(this.groupId);
|
||||
if (group.getMessages().size() != 0) {
|
||||
failures.add("REMOVE");
|
||||
throw new AssertionFailedError("Failed on Remove");
|
||||
}
|
||||
});
|
||||
|
||||
@@ -391,13 +392,13 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
Message<?> m1 = MessageBuilder.withPayload("1")
|
||||
.setSequenceNumber(1)
|
||||
.setSequenceSize(3)
|
||||
.setCorrelationId(1)
|
||||
.setCorrelationId(this.groupId)
|
||||
.build();
|
||||
|
||||
Message<?> m2 = MessageBuilder.withPayload("2")
|
||||
.setSequenceNumber(2)
|
||||
.setSequenceSize(3)
|
||||
.setCorrelationId(1)
|
||||
.setCorrelationId(this.groupId)
|
||||
.build();
|
||||
|
||||
input.send(m1);
|
||||
@@ -414,11 +415,11 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
Message<?> m3 = MessageBuilder.withPayload("3")
|
||||
.setSequenceNumber(3)
|
||||
.setSequenceSize(3)
|
||||
.setCorrelationId(1)
|
||||
.setCorrelationId(this.groupId)
|
||||
.build();
|
||||
|
||||
input.send(m3);
|
||||
assertNotNull(output.receive(1000));
|
||||
assertNotNull(output.receive(10000));
|
||||
|
||||
MessageGroupStoreReaper messageGroupStoreReaper = context.getBean(MessageGroupStoreReaper.class);
|
||||
messageGroupStoreReaper.run();
|
||||
@@ -431,28 +432,26 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testAddAndRemoveMessagesFromMessageGroup() throws Exception {
|
||||
public void testAddAndRemoveMessagesFromMessageGroup() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore messageStore = new RedisMessageStore(jcf);
|
||||
String groupId = "X";
|
||||
messageStore.removeMessageGroup("X");
|
||||
List<Message<?>> messages = new ArrayList<Message<?>>();
|
||||
for (int i = 0; i < 25; i++) {
|
||||
Message<String> message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build();
|
||||
messageStore.addMessagesToGroup(groupId, message);
|
||||
Message<String> message = MessageBuilder.withPayload("foo").setCorrelationId(this.groupId).build();
|
||||
messageStore.addMessagesToGroup(this.groupId, message);
|
||||
messages.add(message);
|
||||
}
|
||||
MessageGroup group = messageStore.getMessageGroup(groupId);
|
||||
MessageGroup group = messageStore.getMessageGroup(this.groupId);
|
||||
assertEquals(25, group.size());
|
||||
messageStore.removeMessagesFromGroup(groupId, messages);
|
||||
group = messageStore.getMessageGroup(groupId);
|
||||
messageStore.removeMessagesFromGroup(this.groupId, messages);
|
||||
group = messageStore.getMessageGroup(this.groupId);
|
||||
assertEquals(0, group.size());
|
||||
messageStore.removeMessageGroup("X");
|
||||
messageStore.removeMessageGroup(this.groupId);
|
||||
}
|
||||
|
||||
@Test
|
||||
@RedisAvailable
|
||||
public void testJsonSerialization() throws Exception {
|
||||
public void testJsonSerialization() {
|
||||
RedisConnectionFactory jcf = getConnectionFactoryForTest();
|
||||
RedisMessageStore store = new RedisMessageStore(jcf);
|
||||
|
||||
@@ -466,9 +465,9 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
Message<?> adviceMessage = new AdviceMessage<>("foo", genericMessage);
|
||||
ErrorMessage errorMessage = new ErrorMessage(new RuntimeException("test exception"));
|
||||
|
||||
store.addMessagesToGroup(1, genericMessage, mutableMessage, adviceMessage, errorMessage);
|
||||
store.addMessagesToGroup(this.groupId, genericMessage, mutableMessage, adviceMessage, errorMessage);
|
||||
|
||||
MessageGroup messageGroup = store.getMessageGroup(1);
|
||||
MessageGroup messageGroup = store.getMessageGroup(this.groupId);
|
||||
assertEquals(4, messageGroup.size());
|
||||
List<Message<?>> messages = new ArrayList<>(messageGroup.getMessages());
|
||||
assertEquals(genericMessage.getPayload(), messages.get(0).getPayload());
|
||||
@@ -481,7 +480,7 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
|
||||
Message<Foo> fooMessage = new GenericMessage<>(new Foo("foo"));
|
||||
try {
|
||||
store.addMessageToGroup(1, fooMessage)
|
||||
store.addMessageToGroup(this.groupId, fooMessage)
|
||||
.getMessages()
|
||||
.iterator()
|
||||
.next();
|
||||
@@ -501,8 +500,8 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
serializer = new GenericJackson2JsonRedisSerializer(mapper);
|
||||
store.setValueSerializer(serializer);
|
||||
|
||||
store.removeMessageGroup(1);
|
||||
messageGroup = store.addMessageToGroup(1, fooMessage);
|
||||
store.removeMessageGroup(this.groupId);
|
||||
messageGroup = store.addMessageToGroup(this.groupId, fooMessage);
|
||||
assertEquals(1, messageGroup.size());
|
||||
assertEquals(fooMessage.getPayload(), messageGroup.getMessages().iterator().next().getPayload());
|
||||
|
||||
@@ -511,8 +510,8 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests {
|
||||
serializer = new GenericJackson2JsonRedisSerializer(mapper);
|
||||
store.setValueSerializer(serializer);
|
||||
|
||||
store.removeMessageGroup(1);
|
||||
messageGroup = store.addMessageToGroup(1, fooMessage);
|
||||
store.removeMessageGroup(this.groupId);
|
||||
messageGroup = store.addMessageToGroup(this.groupId, fooMessage);
|
||||
assertEquals(1, messageGroup.size());
|
||||
assertEquals(fooMessage.getPayload(), messageGroup.getMessages().iterator().next().getPayload());
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user