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()); }