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
This commit is contained in:
@@ -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);
|
||||
|
||||
@@ -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<String>("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<String>("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<String>("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();
|
||||
|
||||
Reference in New Issue
Block a user