INT-3087:Fix MongoDbMessageStore RemoveMessage Bug
Add new query so the same message found in pollMessageFromGroup is the one which gets deleted. See jira for more details: https://jira.springsource.org/browse/INT-3087 INT-3087 Polishing Use an atomic method in pollMessageFromGroup to pop the first message. Add test cases for pollMessageFromGroup and removeMessageFromGroup where the same message exists in multiple groups.
This commit is contained in:
committed by
Gary Russell
parent
c5355a3417
commit
56af118656
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2012 the original author or authors.
|
||||
* Copyright 2002-2013 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.
|
||||
@@ -57,7 +57,6 @@ import org.springframework.integration.store.SimpleMessageGroup;
|
||||
import org.springframework.jmx.export.annotation.ManagedAttribute;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.ClassUtils;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
import com.mongodb.BasicDBList;
|
||||
@@ -71,6 +70,8 @@ import com.mongodb.DBObject;
|
||||
* @author Mark Fisher
|
||||
* @author Oleg Zhurakousky
|
||||
* @author Sean Brandt
|
||||
* @author Jodie StJohn
|
||||
* @author Gary Russell
|
||||
* @since 2.1
|
||||
*/
|
||||
public class MongoDbMessageStore extends AbstractMessageGroupStore implements MessageStore, BeanClassLoaderAware {
|
||||
@@ -205,7 +206,8 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
public MessageGroup removeMessageFromGroup(Object groupId, Message<?> messageToRemove) {
|
||||
Assert.notNull(groupId, "'groupId' must not be null");
|
||||
Assert.notNull(messageToRemove, "'messageToRemove' must not be null");
|
||||
this.removeMessage(messageToRemove.getHeaders().getId());
|
||||
this.template.findAndRemove(whereMessageIdIsAndGroupIdIs(
|
||||
messageToRemove.getHeaders().getId(), groupId), MessageWrapper.class, this.collectionName);
|
||||
this.updateGroup(groupId);
|
||||
return this.getMessageGroup(groupId);
|
||||
}
|
||||
@@ -245,12 +247,10 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
|
||||
public Message<?> pollMessageFromGroup(Object groupId) {
|
||||
Assert.notNull(groupId, "'groupId' must not be null");
|
||||
List<MessageWrapper> messageWrappers = this.template.find(whereGroupIdIsOrdered(groupId), MessageWrapper.class, this.collectionName);
|
||||
MessageWrapper messageWrapper = this.template.findAndRemove(whereGroupIdIsOrdered(groupId), MessageWrapper.class, this.collectionName);
|
||||
Message<?> message = null;
|
||||
|
||||
if (!CollectionUtils.isEmpty(messageWrappers)){
|
||||
message = messageWrappers.get(0).getMessage();
|
||||
this.removeMessageFromGroup(groupId, message);
|
||||
if (messageWrapper != null) {
|
||||
message = messageWrapper.getMessage();
|
||||
}
|
||||
this.updateGroup(groupId);
|
||||
return message;
|
||||
@@ -270,6 +270,10 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
|
||||
return new Query(where("headers.id._value").is(id.toString()));
|
||||
}
|
||||
|
||||
private static Query whereMessageIdIsAndGroupIdIs(UUID id, Object groupId) {
|
||||
return new Query(where("headers.id._value").is(id.toString()).and(GROUP_ID_KEY).is(groupId));
|
||||
}
|
||||
|
||||
private static Query whereGroupIdIs(Object groupId) {
|
||||
Query q = new Query(where(GROUP_ID_KEY).is(groupId));
|
||||
q.with(new Sort(Direction.DESC, GROUP_UPDATE_TIMESTAMP_KEY));
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2011 the original author or authors.
|
||||
* Copyright 2002-2013 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.
|
||||
@@ -18,7 +18,6 @@ package org.springframework.integration.mongodb.rules;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.junit.rules.MethodRule;
|
||||
import org.junit.runners.model.FrameworkMethod;
|
||||
import org.junit.runners.model.Statement;
|
||||
@@ -27,8 +26,9 @@ import com.mongodb.Mongo;
|
||||
|
||||
/**
|
||||
* A {@link MethodRule} implementation that checks for a running MongoDB process.
|
||||
*
|
||||
*
|
||||
* @author Oleg Zhurakousky
|
||||
* @author Gary Russell
|
||||
* @since 2.1
|
||||
*/
|
||||
public final class MongoDbAvailableRule implements MethodRule {
|
||||
@@ -39,10 +39,10 @@ public final class MongoDbAvailableRule implements MethodRule {
|
||||
return new Statement() {
|
||||
@Override
|
||||
public void evaluate() throws Throwable {
|
||||
MongoDbAvailable redisAvailable = method.getAnnotation(MongoDbAvailable.class);
|
||||
if (redisAvailable != null) {
|
||||
MongoDbAvailable mongoAvailable = method.getAnnotation(MongoDbAvailable.class);
|
||||
if (mongoAvailable != null) {
|
||||
try {
|
||||
Mongo mongo = new Mongo();
|
||||
Mongo mongo = new Mongo();
|
||||
mongo.getDatabaseNames();
|
||||
}
|
||||
catch (Exception e) {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2007-2012 the original author or authors
|
||||
* Copyright 2007-2013 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.
|
||||
@@ -15,6 +15,12 @@
|
||||
*/
|
||||
package org.springframework.integration.mongodb.store;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertFalse;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
import java.util.Iterator;
|
||||
import java.util.Properties;
|
||||
import java.util.UUID;
|
||||
@@ -35,14 +41,9 @@ import org.springframework.integration.store.MessageGroup;
|
||||
import org.springframework.integration.store.SimpleMessageGroup;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertFalse;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
/**
|
||||
* @author Oleg Zhurakousky
|
||||
* @author Gary Russell
|
||||
*
|
||||
*/
|
||||
public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
|
||||
@@ -74,7 +75,7 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
|
||||
Message<?> retrievedMessage = store.getMessage(messageA.getHeaders().getId());
|
||||
assertNotNull(retrievedMessage);
|
||||
assertEquals(retrievedMessage.getHeaders().getId(), messageA.getHeaders().getId());
|
||||
// ensure that 'message_group' header that is only used internally is not propagated
|
||||
// ensure that 'message_group' header that is only used internally is not propagated
|
||||
assertNull(retrievedMessage.getHeaders().get("message_group"));
|
||||
}
|
||||
|
||||
@@ -115,6 +116,99 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
|
||||
assertEquals(2, store.messageGroupSize(1));
|
||||
}
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testPollMessages() throws Exception{
|
||||
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
|
||||
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
|
||||
|
||||
Message<?> messageA = new GenericMessage<String>("A");
|
||||
Message<?> messageB = new GenericMessage<String>("B");
|
||||
store.addMessageToGroup(1, messageA);
|
||||
store.addMessageToGroup(1, messageB);
|
||||
assertEquals(2, store.messageGroupSize(1));
|
||||
Message<?> out = store.pollMessageFromGroup(1);
|
||||
assertEquals("A", out.getPayload());
|
||||
assertEquals(1, store.messageGroupSize(1));
|
||||
out = store.pollMessageFromGroup(1);
|
||||
assertEquals("B", out.getPayload());
|
||||
assertEquals(0, store.messageGroupSize(1));
|
||||
}
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testSameMessageMultipleGroupsPoll() throws Exception{
|
||||
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
|
||||
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
|
||||
|
||||
Message<?> messageA = new GenericMessage<String>("A");
|
||||
store.addMessageToGroup(1, messageA);
|
||||
store.addMessageToGroup(2, messageA);
|
||||
store.addMessageToGroup(3, messageA);
|
||||
store.addMessageToGroup(4, messageA);
|
||||
assertEquals(1, store.messageGroupSize(1));
|
||||
assertEquals(1, store.messageGroupSize(2));
|
||||
assertEquals(1, store.messageGroupSize(3));
|
||||
assertEquals(1, store.messageGroupSize(4));
|
||||
store.pollMessageFromGroup(3);
|
||||
assertEquals(1, store.messageGroupSize(1));
|
||||
assertEquals(1, store.messageGroupSize(2));
|
||||
assertEquals(0, store.messageGroupSize(3));
|
||||
assertEquals(1, store.messageGroupSize(4));
|
||||
store.pollMessageFromGroup(4);
|
||||
assertEquals(1, store.messageGroupSize(1));
|
||||
assertEquals(1, store.messageGroupSize(2));
|
||||
assertEquals(0, store.messageGroupSize(3));
|
||||
assertEquals(0, store.messageGroupSize(4));
|
||||
store.pollMessageFromGroup(2);
|
||||
assertEquals(1, store.messageGroupSize(1));
|
||||
assertEquals(0, store.messageGroupSize(2));
|
||||
assertEquals(0, store.messageGroupSize(3));
|
||||
assertEquals(0, store.messageGroupSize(4));
|
||||
store.pollMessageFromGroup(1);
|
||||
assertEquals(0, store.messageGroupSize(1));
|
||||
assertEquals(0, store.messageGroupSize(2));
|
||||
assertEquals(0, store.messageGroupSize(3));
|
||||
assertEquals(0, store.messageGroupSize(4));
|
||||
}
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testSameMessageMultipleGroupsRemove() throws Exception{
|
||||
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
|
||||
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
|
||||
|
||||
Message<?> messageA = new GenericMessage<String>("A");
|
||||
store.addMessageToGroup(1, messageA);
|
||||
store.addMessageToGroup(2, messageA);
|
||||
store.addMessageToGroup(3, messageA);
|
||||
store.addMessageToGroup(4, messageA);
|
||||
assertEquals(1, store.messageGroupSize(1));
|
||||
assertEquals(1, store.messageGroupSize(2));
|
||||
assertEquals(1, store.messageGroupSize(3));
|
||||
assertEquals(1, store.messageGroupSize(4));
|
||||
store.removeMessageFromGroup(3, messageA);
|
||||
assertEquals(1, store.messageGroupSize(1));
|
||||
assertEquals(1, store.messageGroupSize(2));
|
||||
assertEquals(0, store.messageGroupSize(3));
|
||||
assertEquals(1, store.messageGroupSize(4));
|
||||
store.removeMessageFromGroup(4, messageA);
|
||||
assertEquals(1, store.messageGroupSize(1));
|
||||
assertEquals(1, store.messageGroupSize(2));
|
||||
assertEquals(0, store.messageGroupSize(3));
|
||||
assertEquals(0, store.messageGroupSize(4));
|
||||
store.removeMessageFromGroup(2, messageA);
|
||||
assertEquals(1, store.messageGroupSize(1));
|
||||
assertEquals(0, store.messageGroupSize(2));
|
||||
assertEquals(0, store.messageGroupSize(3));
|
||||
assertEquals(0, store.messageGroupSize(4));
|
||||
store.removeMessageFromGroup(1, messageA);
|
||||
assertEquals(0, store.messageGroupSize(1));
|
||||
assertEquals(0, store.messageGroupSize(2));
|
||||
assertEquals(0, store.messageGroupSize(3));
|
||||
assertEquals(0, store.messageGroupSize(4));
|
||||
}
|
||||
|
||||
@Test
|
||||
@MongoDbAvailable
|
||||
public void testMessageGroupUpdatedDateChangesWithEachAddedMessage() throws Exception{
|
||||
@@ -297,7 +391,7 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
|
||||
// final MongoDbMessageStore store1 = new MongoDbMessageStore(mongoDbFactory);
|
||||
// final MongoDbMessageStore store2 = new MongoDbMessageStore(mongoDbFactory);
|
||||
//
|
||||
// final Message<?> message = new GenericMessage<String>("1");
|
||||
// final Message<?> message = new GenericMessage<String>("1");
|
||||
//
|
||||
// ExecutorService executor = null;
|
||||
//
|
||||
|
||||
Reference in New Issue
Block a user