diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageGroup.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageGroup.java index 2b23d81c99..c3a491a007 100644 --- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageGroup.java +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageGroup.java @@ -12,8 +12,8 @@ */ package org.springframework.integration.redis.store; -import java.util.ArrayList; import java.util.Collection; +import java.util.LinkedList; import java.util.List; import org.springframework.data.redis.core.BoundSetOperations; @@ -43,7 +43,10 @@ class RedisMessageGroup implements MessageGroup { // TODO Auto-generated method stub return false; } - + /** + * + * @param message + */ public void add(Message message) { BoundSetOperations mGroupOps = this.redisTemplate.boundSetOps(UNMARKED_PREFIX + this.groupId.toString()); String key = message.getHeaders().getId().toString(); @@ -51,7 +54,9 @@ class RedisMessageGroup implements MessageGroup { BoundValueOperations bvOps = this.redisTemplate.boundValueOps(key); bvOps.set(message); } - + /** + * + */ public void destroy(){ // delete unmarked BoundSetOperations mGroupOps = this.redisTemplate.boundSetOps(UNMARKED_PREFIX + this.groupId.toString()); @@ -68,7 +73,16 @@ class RedisMessageGroup implements MessageGroup { } public void markAll() { - + BoundSetOperations unmarkedGroupOps = this.redisTemplate.boundSetOps(UNMARKED_PREFIX + this.groupId.toString()); + long uSize = unmarkedGroupOps.size(); + String destinationKey = UNMARKED_PREFIX + this.groupId.toString(); + + for (Object messageId : unmarkedGroupOps.members()) { + unmarkedGroupOps.move(destinationKey, messageId); + } + BoundSetOperations markedGroupOps = this.redisTemplate.boundSetOps(UNMARKED_PREFIX + this.groupId.toString()); + long mSize = markedGroupOps.size(); + Assert.isTrue(mSize == uSize, "Failed to mark ALl messages in the message group"); } @@ -110,11 +124,18 @@ class RedisMessageGroup implements MessageGroup { return mGroupOps.members().size(); } - /* (non-Javadoc) - * @see org.springframework.integration.store.MessageGroup#getOne() - */ + public Message getOne() { - // TODO Auto-generated method stub + List> messages = (List>) this.getUnmarked(); + if (messages.isEmpty()){ + messages = (List>) this.getMarked(); + if (!messages.isEmpty()){ + return messages.get(0); + } + } + else { + return messages.get(0); + } return null; } @@ -127,7 +148,7 @@ class RedisMessageGroup implements MessageGroup { } private Collection> getMessages(BoundSetOperations mGroupOps){ - List> messages = new ArrayList>(); + List> messages = new LinkedList>(); for (Object messageId : mGroupOps.members()) { BoundValueOperations bvOps = this.redisTemplate.boundValueOps(messageId.toString()); Object message = bvOps.get(); 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 215e6d7514..45e53951ea 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 @@ -85,4 +85,33 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { assertEquals(0, messageGroup.size()); } + @Test + @RedisAvailable + public void testMarkAllMessagesInMessageGroup(){ + JedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisMessageStore store = new RedisMessageStore(jcf); + + MessageGroup messageGroup = store.getMessageGroup(1); + ((RedisMessageGroup)messageGroup).add(new GenericMessage("1")); + ((RedisMessageGroup)messageGroup).add(new GenericMessage("2")); + ((RedisMessageGroup)messageGroup).add(new GenericMessage("3")); + ((RedisMessageGroup)messageGroup).markAll(); + // No need to assert for anything here. + // Assertion is done within the method itself and the exception would be thrown if move did not succeed + } + + @Test + @RedisAvailable + public void testGetOneFromMessageGroup(){ + JedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisMessageStore store = new RedisMessageStore(jcf); + + MessageGroup messageGroup = store.getMessageGroup(1); + ((RedisMessageGroup)messageGroup).add(new GenericMessage("1")); + ((RedisMessageGroup)messageGroup).add(new GenericMessage("2")); + ((RedisMessageGroup)messageGroup).add(new GenericMessage("3")); + Message message = messageGroup.getOne(); + assertEquals("1", message.getPayload()); + } + }