From 7f407af3f201fd6bfa55b497e46b73632f39d6fa Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 30 Jan 2018 10:33:56 -0500 Subject: [PATCH] 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 --- .../store/RedisMessageGroupStoreTests.java | 197 +++++++++--------- 1 file changed, 98 insertions(+), 99 deletions(-) diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageGroupStoreTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageGroupStoreTests.java index a84d7354c3..d3f1a625ba 100644 --- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageGroupStoreTests.java +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageGroupStoreTests.java @@ -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("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("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("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("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("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("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("2"); - store.addMessagesToGroup(messageGroup.getGroupId(), new GenericMessage("1"), message); - messageGroup = store.addMessageToGroup(messageGroup.getGroupId(), new GenericMessage("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("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("1")); - store.removeMessagesFromGroup(1, new GenericMessage("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("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("1"); - store1.addMessagesToGroup(1, message); - MessageGroup messageGroup = store2.addMessageToGroup(1, new GenericMessage("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("1")); - store2.addMessagesToGroup(2, new GenericMessage("2")); - store1.addMessagesToGroup(3, new GenericMessage("3"), new GenericMessage("3A")); + store1.addMessagesToGroup(this.groupId, new GenericMessage<>("1")); + store2.addMessagesToGroup(group2, new GenericMessage<>("2")); + store1.addMessagesToGroup(group3, new GenericMessage<>("3"), new GenericMessage<>("3A")); Iterator 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("1"); + store1.removeMessageGroup(this.groupId); + + final Message message = new GenericMessage<>("1"); ExecutorService executor = null; - final List failures = new ArrayList(); + final List 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> messages = new ArrayList>(); for (int i = 0; i < 25; i++) { - Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); - messageStore.addMessagesToGroup(groupId, message); + Message 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> messages = new ArrayList<>(messageGroup.getMessages()); assertEquals(genericMessage.getPayload(), messages.get(0).getPayload()); @@ -481,7 +480,7 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { Message 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()); }