From a352f1951408845ed53c11e6d194636bd0878ffe Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Fri, 12 Aug 2011 12:54:24 -0400 Subject: [PATCH] Revert "Merge pull request #6 from olegz/INT-2030" Backing out this change in order to use the merge commit approach. This reverts commit ab47ca704b9e3cf495fe31f4db9c86a757952efa. --- .../redis/store/RedisMessageGroup.java | 63 +++++++++---------- .../store/RedisMessageGroupStoreTests.java | 3 - 2 files changed, 30 insertions(+), 36 deletions(-) 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 5d69139ea4..e00baac1bc 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 @@ -32,9 +32,6 @@ class RedisMessageGroup implements MessageGroup { private final String MARKED_PREFIX = "MARKED_"; private final String UNMARKED_PREFIX = "UNMARKED_"; - private final String unmarkedId; - private final String markedId; - private final List> unmarked = new LinkedList>(); private final List> marked = new LinkedList>(); @@ -48,10 +45,8 @@ class RedisMessageGroup implements MessageGroup { this.groupId = groupId; this.redisTemplate = redisTemplate; this.messageStore = messageStore; - this.unmarkedId = UNMARKED_PREFIX + groupId; - this.markedId = MARKED_PREFIX + groupId; - BoundListOperations unmarkedOps = this.redisTemplate.boundListOps(unmarkedId); - BoundListOperations markedOps = this.redisTemplate.boundListOps(markedId); + BoundListOperations unmarkedOps = this.redisTemplate.boundListOps(UNMARKED_PREFIX + this.groupId.toString()); + BoundListOperations markedOps = this.redisTemplate.boundListOps(MARKED_PREFIX + this.groupId.toString()); this.rebuildLocalCache(unmarkedOps, markedOps); } @@ -90,19 +85,19 @@ class RedisMessageGroup implements MessageGroup { } public int size() { - synchronized (lock) { - return marked.size() + unmarked.size(); - } + BoundListOperations mGroupOps = this.redisTemplate.boundListOps(UNMARKED_PREFIX + this.groupId.toString()); + long usize = mGroupOps.size(); + mGroupOps = this.redisTemplate.boundListOps(MARKED_PREFIX + this.groupId.toString()); + long msize = mGroupOps.size(); + return (int)(usize + msize); } public Message getOne() { - synchronized (lock) { - Message one = unmarked.get(0); - if (one == null) { - one = marked.get(0); - } - return one; + Message one = unmarked.get(0); + if (one == null) { + one = marked.get(0); } + return one; } public long getTimestamp() { @@ -110,12 +105,13 @@ class RedisMessageGroup implements MessageGroup { } protected void markAll() { - BoundListOperations unmarkedOps = this.redisTemplate.boundListOps(unmarkedId); - BoundListOperations markedOps = this.redisTemplate.boundListOps(markedId); + BoundListOperations unmarkedOps = this.redisTemplate.boundListOps(UNMARKED_PREFIX + this.groupId.toString()); + BoundListOperations markedOps = this.redisTemplate.boundListOps(MARKED_PREFIX + this.groupId.toString()); long uSize = unmarkedOps.size(); + String destiinationKey = MARKED_PREFIX + this.groupId.toString(); synchronized (lock) { - unmarkedOps.rename(markedId); + unmarkedOps.rename(destiinationKey); this.unmarked.clear(); this.rebuildLocalCache(null, markedOps); } @@ -125,15 +121,17 @@ class RedisMessageGroup implements MessageGroup { } protected void markMessage(String messageId) { - BoundListOperations unmarkedOps = this.redisTemplate.boundListOps(unmarkedId); + BoundListOperations unmarkedOps = this.redisTemplate.boundListOps(UNMARKED_PREFIX + this.groupId.toString()); List messageIds = unmarkedOps.range(0, unmarkedOps.size()-1); synchronized (lock) { for (Object id : messageIds) { if (messageId.equals(id)){ - BoundListOperations markedOps = this.redisTemplate.boundListOps(markedId); + BoundListOperations markedOps = this.redisTemplate.boundListOps(MARKED_PREFIX + this.groupId.toString()); markedOps.rightPush(id); + System.out.println(unmarkedOps.size()); unmarkedOps.remove(0, id); + System.out.println(unmarkedOps.size()); this.rebuildLocalCache(unmarkedOps, markedOps); return; } @@ -147,7 +145,7 @@ class RedisMessageGroup implements MessageGroup { */ protected void add(Message message) { String messageId = message.getHeaders().getId().toString(); - BoundListOperations unmarkedOps = this.redisTemplate.boundListOps(unmarkedId); + BoundListOperations unmarkedOps = this.redisTemplate.boundListOps(UNMARKED_PREFIX + this.groupId.toString()); synchronized (lock) { unmarkedOps.rightPush(messageId); this.messageStore.addMessage(message); @@ -157,8 +155,8 @@ class RedisMessageGroup implements MessageGroup { protected void remove(Message message) { UUID messageId = message.getHeaders().getId(); - BoundListOperations unmarkedOps = this.redisTemplate.boundListOps(unmarkedId); - BoundListOperations markedOps = this.redisTemplate.boundListOps(markedId); + BoundListOperations unmarkedOps = this.redisTemplate.boundListOps(UNMARKED_PREFIX + this.groupId.toString()); + BoundListOperations markedOps = this.redisTemplate.boundListOps(MARKED_PREFIX + this.groupId.toString()); synchronized (lock) { unmarkedOps.remove(0, messageId.toString()); @@ -172,20 +170,19 @@ class RedisMessageGroup implements MessageGroup { * */ protected void destroy(){ - BoundListOperations unmarkedOps = this.redisTemplate.boundListOps(unmarkedId); - List messageIds = unmarkedOps.range(0, unmarkedOps.size()-1); + BoundListOperations mGroupListOps = this.redisTemplate.boundListOps(UNMARKED_PREFIX + this.groupId.toString()); + List messageIds = mGroupListOps.range(0, mGroupListOps.size()-1); for (Object messageId : messageIds) { this.messageStore.removeMessage(UUID.fromString(messageId.toString())); } - this.redisTemplate.delete(unmarkedId); + this.redisTemplate.delete(UNMARKED_PREFIX + this.groupId.toString()); - BoundListOperations markedOps = this.redisTemplate.boundListOps(markedId); - messageIds = markedOps.range(0, markedOps.size()-1); + mGroupListOps = this.redisTemplate.boundListOps(MARKED_PREFIX + this.groupId.toString()); + messageIds = mGroupListOps.range(0, mGroupListOps.size()-1); for (Object messageId : messageIds) { this.messageStore.removeMessage(UUID.fromString(messageId.toString())); } - this.redisTemplate.delete(markedId); - this.rebuildLocalCache(unmarkedOps, markedOps); + this.redisTemplate.delete(MARKED_PREFIX + this.groupId.toString()); } private Collection> buildMessageList(BoundListOperations mGroupOps){ @@ -211,9 +208,9 @@ class RedisMessageGroup implements MessageGroup { else { synchronized (lock) { - BoundListOperations mGroupOps = this.redisTemplate.boundListOps(unmarkedId); + BoundListOperations mGroupOps = this.redisTemplate.boundListOps(UNMARKED_PREFIX + this.groupId.toString()); Collection> unmarked = this.buildMessageList(mGroupOps); - mGroupOps = this.redisTemplate.boundListOps(markedId); + mGroupOps = this.redisTemplate.boundListOps(MARKED_PREFIX + this.groupId.toString()); Collection> marked = this.buildMessageList(mGroupOps); if (containsSequenceNumber(unmarked, messageSequenceNumber) || containsSequenceNumber(marked, messageSequenceNumber)) { 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 cd13613228..480ad2e2ce 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 @@ -88,8 +88,6 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { store.removeMessageGroup(1); messageGroup = store.getMessageGroup(1); - assertEquals(0, messageGroup.getMarked().size()); - assertEquals(0, messageGroup.getUnmarked().size()); assertEquals(0, messageGroup.size()); } @@ -224,7 +222,6 @@ public class RedisMessageGroupStoreTests extends RedisAvailableTests { assertNotSame(mg1, mg2); } - @Test @RedisAvailable public void testWithAggregatorWithShutdown(){