From 49985df8580e296b7aa4aed9e32032537c8e0c0b Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 18 May 2016 13:35:30 -0400 Subject: [PATCH] Fix `AbstractKeyValueMessageStore` create time The `RedisMessageGroupStoreTests.testMessageGroupUpdatedDateChangesWithEachAddedMessage()` can fail sporadically because there is no guaranty the several calls to the `System.getCurrentTimeMillis()` return the same value. * Fix `AbstractKeyValueMessageStore` to reuse `MessageGroup.timestamp` for the `lastModified` in case of a new group --- .../store/AbstractKeyValueMessageStore.java | 6 ++- .../store/RedisMessageGroupStoreTests.java | 42 +++++++++++++------ 2 files changed, 35 insertions(+), 13 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java index fcca5ddb1c..ac4fb4a9b6 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/AbstractKeyValueMessageStore.java @@ -144,9 +144,13 @@ public abstract class AbstractKeyValueMessageStore extends AbstractMessageGroupS if (group != null) { metadata = new MessageGroupMetadata(group); + // When the group is new reuse "create time" as a "last modified" + metadata.setLastModified(group.getTimestamp()); + } + else { + metadata.setLastModified(System.currentTimeMillis()); } - metadata.setLastModified(System.currentTimeMillis()); // store MessageGroupMetadata built from enriched MG doStore(MESSAGE_GROUP_KEY_PREFIX + groupId, metadata); 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 29bcf6fba8..d0f1e711a2 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 @@ -84,14 +84,13 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { RedisConnectionFactory jcf = getConnectionFactoryForTest(); RedisMessageStore store = new RedisMessageStore(jcf); - MessageGroup messageGroup = store.getMessageGroup(1); Message message = new GenericMessage("Hello"); - messageGroup = store.addMessageToGroup(1, message); + MessageGroup messageGroup = store.addMessageToGroup(1, message); assertEquals(1, messageGroup.size()); long createdTimestamp = messageGroup.getTimestamp(); long updatedTimestamp = messageGroup.getLastModified(); - assertTrue(updatedTimestamp - createdTimestamp <= 2); - Thread.sleep(1000); + assertEquals(createdTimestamp, updatedTimestamp); + Thread.sleep(10); message = new GenericMessage("Hello"); messageGroup = store.addMessageToGroup(1, message); createdTimestamp = messageGroup.getTimestamp(); @@ -111,9 +110,8 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { RedisConnectionFactory jcf = getConnectionFactoryForTest(); RedisMessageStore store = new RedisMessageStore(jcf); - MessageGroup messageGroup = store.getMessageGroup(1); Message message = new GenericMessage("Hello"); - messageGroup = store.addMessageToGroup(1, message); + MessageGroup messageGroup = store.addMessageToGroup(1, message); assertEquals(1, messageGroup.size()); // make sure the store is properly rebuild from Redis @@ -228,6 +226,7 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { assertEquals("fooChannel", fooChannelHistory.get("name")); assertEquals("channel", fooChannelHistory.get("type")); } + @Test @RedisAvailable public void testRemoveNonExistingMessageFromTheGroup() throws Exception { @@ -248,7 +247,6 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { } - @Test @RedisAvailable public void testMultipleInstancesOfGroupStore() throws Exception { @@ -313,7 +311,8 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { } @Test - @RedisAvailable @Ignore + @RedisAvailable + @Ignore public void testConcurrentModifications() throws Exception { RedisConnectionFactory jcf = getConnectionFactoryForTest(); final RedisMessageStore store1 = new RedisMessageStore(jcf); @@ -329,6 +328,7 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { executor = Executors.newCachedThreadPool(); executor.execute(new Runnable() { + @Override public void run() { MessageGroup group = store1.addMessageToGroup(1, message); @@ -339,6 +339,7 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { } }); executor.execute(new Runnable() { + @Override public void run() { store2.removeMessagesFromGroup(1, message); @@ -360,23 +361,40 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { @Test @RedisAvailable public void testWithAggregatorWithShutdown() { - ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("redis-aggregator-config.xml", getClass()); + ClassPathXmlApplicationContext context = + new ClassPathXmlApplicationContext("redis-aggregator-config.xml", getClass()); MessageChannel input = context.getBean("inputChannel", MessageChannel.class); QueueChannel output = context.getBean("outputChannel", QueueChannel.class); - Message m1 = MessageBuilder.withPayload("1").setSequenceNumber(1).setSequenceSize(3).setCorrelationId(1).build(); - Message m2 = MessageBuilder.withPayload("2").setSequenceNumber(2).setSequenceSize(3).setCorrelationId(1).build(); + Message m1 = MessageBuilder.withPayload("1") + .setSequenceNumber(1) + .setSequenceSize(3) + .setCorrelationId(1) + .build(); + + Message m2 = MessageBuilder.withPayload("2") + .setSequenceNumber(2) + .setSequenceSize(3) + .setCorrelationId(1) + .build(); + input.send(m1); assertNull(output.receive(1000)); input.send(m2); assertNull(output.receive(1000)); + context.close(); context = new ClassPathXmlApplicationContext("redis-aggregator-config.xml", getClass()); input = context.getBean("inputChannel", MessageChannel.class); output = context.getBean("outputChannel", QueueChannel.class); - Message m3 = MessageBuilder.withPayload("3").setSequenceNumber(3).setSequenceSize(3).setCorrelationId(1).build(); + Message m3 = MessageBuilder.withPayload("3") + .setSequenceNumber(3) + .setSequenceSize(3) + .setCorrelationId(1) + .build(); + input.send(m3); assertNotNull(output.receive(1000)); context.close();