Merge pull request #380 from olegz/INT-2479

This commit is contained in:
Gary Russell
2012-03-29 15:50:45 -04:00
3 changed files with 201 additions and 185 deletions

View File

@@ -554,6 +554,7 @@ project('spring-integration-mongodb') {
'org.springframework.jmx.*;version="[3.0.5, 4.0.0)"',
'org.springframework.data.mongodb.*;version="[1.0.0, 2.0.0)"',
'org.springframework.data.mapping.*;version="[1.0.0, 2.0.0)"',
'org.springframework.data.annotation.*;version="[1.0.0, 2.0.0)"',
'com.mongodb.*;version="[0.0.0, 2.5.0]"',
'javax.*;version="0"',
'org.w3c.dom.*;version="0"'

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2011 the original author or authors.
* Copyright 2002-2012 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.
@@ -27,6 +27,7 @@ import java.util.UUID;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.BeanClassLoaderAware;
import org.springframework.core.convert.converter.Converter;
import org.springframework.data.annotation.Transient;
import org.springframework.data.mapping.context.MappingContext;
import org.springframework.data.mongodb.MongoDbFactory;
import org.springframework.data.mongodb.core.MongoTemplate;
@@ -78,15 +79,15 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
private final static String GROUP_ID_KEY = "_groupId";
private final static String GROUP_COMPLETE_KEY = "_group_complete";
private final static String LAST_RELEASED_SEQUENCE_NUMBER = "_last_released_sequence";
private final static String GROUP_TIMESTAMP_KEY = "_group_timestamp";
private final static String GROUP_UPDATE_TIMESTAMP_KEY = "_group_update_timestamp";
private final static String PAYLOAD_TYPE_KEY = "_payloadType";
private final static String CREATED_DATE = "_createdDate";
@@ -154,12 +155,12 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
boolean completeGroup = false;
if (messageWrappers.size() > 0){
MessageWrapper messageWrapper = messageWrappers.get(0);
timestamp = messageWrapper.getGroupTimestamp();
lastmodified = messageWrapper.getLastModified();
completeGroup = messageWrapper.isCompletedGroup();
lastReleasedSequenceNumber = messageWrapper.getLastReleasedSequenceNumber();
timestamp = messageWrapper.get_Group_timestamp();
lastmodified = messageWrapper.get_Group_update_timestamp();
completeGroup = messageWrapper.get_Group_complete();
lastReleasedSequenceNumber = messageWrapper.get_LastReleasedSequenceNumber();
}
for (MessageWrapper messageWrapper : messageWrappers) {
messages.add(messageWrapper.getMessage());
}
@@ -169,7 +170,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
if (lastReleasedSequenceNumber > 0){
messageGroup.setLastReleasedMessageSequenceNumber(lastReleasedSequenceNumber);
}
return messageGroup;
}
@@ -180,7 +181,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
long messageGroupTimestamp = messageGroup.getTimestamp();
long lastModified = messageGroup.getLastModified();
if (messageGroupTimestamp == 0){
messageGroupTimestamp = System.currentTimeMillis();
lastModified = messageGroupTimestamp;
@@ -188,14 +189,14 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
else {
lastModified = System.currentTimeMillis();
}
MessageWrapper wrapper = new MessageWrapper(message);
wrapper.setGroupId(groupId);
wrapper.setGroupTimestamp(messageGroupTimestamp);
wrapper.setLastModified(lastModified);
wrapper.setCompletedGroup(messageGroup.isComplete());
wrapper.setLastReleasedSequenceNumber(messageGroup.getLastReleasedMessageSequenceNumber());
wrapper.set_GroupId(groupId);
wrapper.set_Group_timestamp(messageGroupTimestamp);
wrapper.set_Group_update_timestamp(lastModified);
wrapper.set_Group_complete(messageGroup.isComplete());
wrapper.set_LastReleasedSequenceNumber(messageGroup.getLastReleasedMessageSequenceNumber());
this.template.insert(wrapper, this.collectionName);
return this.getMessageGroup(groupId);
}
@@ -219,14 +220,14 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
List<MessageWrapper> groupedMessages = this.template.find(whereGroupIdExists(), MessageWrapper.class, this.collectionName);
Map<Object, MessageGroup> messageGroups = new HashMap<Object, MessageGroup>();
for (MessageWrapper groupedMessage : groupedMessages) {
Object groupId = groupedMessage.getGroupId();
Object groupId = groupedMessage.get_GroupId();
if (!messageGroups.containsKey(groupId)) {
messageGroups.put(groupId, this.getMessageGroup(groupId));
}
}
return messageGroups.values().iterator();
}
public void completeGroup(Object groupId) {
Update update = Update.update(GROUP_COMPLETE_KEY, true);
Query q = whereGroupIdIs(groupId);
@@ -240,12 +241,12 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
this.template.updateFirst(q, update, this.collectionName);
this.updateGroup(groupId);
}
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);
Message<?> message = null;
if (!CollectionUtils.isEmpty(messageWrappers)){
message = messageWrappers.get(0).getMessage();
this.removeMessageFromGroup(groupId, message);
@@ -253,7 +254,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
this.updateGroup(groupId);
return message;
}
public int messageGroupSize(Object groupId) {
long lCount = this.template.count(new Query(where(GROUP_ID_KEY).is(groupId)), this.collectionName);
Assert.isTrue(lCount <= Integer.MAX_VALUE, "Message count is out of Integer's range");
@@ -265,7 +266,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
*/
private static Query whereMessageIdIs(UUID id) {
return new Query(where("headers.id").is(id.toString()));
return new Query(where("headers.id._value").is(id.toString()));
}
private static Query whereGroupIdIs(Object groupId) {
@@ -277,13 +278,13 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
private static Query whereGroupIdExists() {
return new Query(where(GROUP_ID_KEY).exists(true));
}
private static Query whereGroupIdIsOrdered(Object groupId) {
Query q = new Query(where(GROUP_ID_KEY).is(groupId)).limit(1);
q.sort().on(CREATED_DATE, Order.ASCENDING);
return q;
}
private void updateGroup(Object groupId) {
Update update = Update.update(GROUP_UPDATE_TIMESTAMP_KEY, System.currentTimeMillis());
Query q = whereGroupIdIs(groupId);
@@ -304,8 +305,8 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
@Override
public void afterPropertiesSet() {
List<Converter<?, ?>> customConverters = new ArrayList<Converter<?,?>>();
customConverters.add(new UuidToStringConverter());
customConverters.add(new StringToUuidConverter());
customConverters.add(new UuidToDBObjectConverter());
customConverters.add(new DBObjectToUUIDConverter());
customConverters.add(new MessageHistoryToDBObjectConverter());
this.setCustomConversions(new CustomConversions(customConverters));
super.afterPropertiesSet();
@@ -313,36 +314,11 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
@Override
public void write(Object source, DBObject target) {
Message<?> message = null;
Object groupId = null;
Assert.isInstanceOf(MessageWrapper.class, source);
boolean groupComplete = false;
long groupTimestamp = 0;
long lastModified = 0;
int lastReleasedSequenceNumber = 0;
if (source instanceof MessageWrapper) {
MessageWrapper wrapper = (MessageWrapper) source;
message = wrapper.getMessage();
groupId = wrapper.getGroupId();
groupComplete = wrapper.isCompletedGroup();
lastReleasedSequenceNumber = wrapper.getLastReleasedSequenceNumber();
groupTimestamp = wrapper.getGroupTimestamp();
lastModified = wrapper.getLastModified();
}
else {
Class<?> sourceType = (source != null) ? source.getClass() : null;
throw new IllegalArgumentException("Unexpected source type [" + sourceType + "]. Should be a MessageWrapper.");
}
target.put(CREATED_DATE, System.currentTimeMillis());
target.put(PAYLOAD_TYPE_KEY, message.getPayload().getClass().getName());
if (groupId != null) {
target.put(GROUP_ID_KEY, groupId);
target.put(GROUP_COMPLETE_KEY, groupComplete);
target.put(LAST_RELEASED_SEQUENCE_NUMBER, lastReleasedSequenceNumber);
target.put(GROUP_TIMESTAMP_KEY, groupTimestamp);
target.put(GROUP_UPDATE_TIMESTAMP_KEY, lastModified);
}
super.write(message, target);
super.write(source, target);
}
@Override
@@ -353,7 +329,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
}
if (source != null) {
Map<String, Object> headers = this.normalizeHeaders((Map<String, Object>) source.get("headers"));
Object payload = source.get("payload");
Object payloadType = source.get(PAYLOAD_TYPE_KEY);
if (payloadType != null && payload instanceof DBObject) {
@@ -368,32 +344,32 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
GenericMessage message = new GenericMessage(payload, headers);
Map innerMap = (Map) new DirectFieldAccessor(message.getHeaders()).getPropertyValue("headers");
// using reflection to set ID and TIMESTAMP since they are immutable through MessageHeaders
innerMap.put(MessageHeaders.ID, UUID.fromString((String) headers.get(MessageHeaders.ID)));
innerMap.put(MessageHeaders.ID, headers.get(MessageHeaders.ID));
innerMap.put(MessageHeaders.TIMESTAMP, headers.get(MessageHeaders.TIMESTAMP));
Long groupTimestamp = (Long)source.get(GROUP_TIMESTAMP_KEY);
Long lastModified = (Long)source.get(GROUP_UPDATE_TIMESTAMP_KEY);
Integer lastReleasedSequenceNumber = (Integer)source.get(LAST_RELEASED_SEQUENCE_NUMBER);
Boolean completeGroup = (Boolean)source.get(GROUP_COMPLETE_KEY);
MessageWrapper wrapper = new MessageWrapper(message);
if (source.containsField(GROUP_ID_KEY)){
wrapper.setGroupId(source.get(GROUP_ID_KEY));
wrapper.set_GroupId(source.get(GROUP_ID_KEY));
}
if (groupTimestamp != null){
wrapper.setGroupTimestamp(groupTimestamp);
wrapper.set_Group_timestamp(groupTimestamp);
}
if (lastModified != null){
wrapper.setLastModified(lastModified);
wrapper.set_Group_update_timestamp(lastModified);
}
if (lastReleasedSequenceNumber != null){
wrapper.setLastReleasedSequenceNumber(lastReleasedSequenceNumber);
wrapper.set_LastReleasedSequenceNumber(lastReleasedSequenceNumber);
}
if (completeGroup != null){
wrapper.setCompletedGroup(completeGroup.booleanValue());
wrapper.set_Group_complete(completeGroup.booleanValue());
}
return (S) wrapper;
}
return null;
@@ -422,17 +398,19 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
}
}
private static class UuidToStringConverter implements Converter<UUID, String> {
public String convert(UUID source) {
return source.toString();
private static class UuidToDBObjectConverter implements Converter<UUID, DBObject> {
public DBObject convert(UUID source) {
BasicDBObject dbObject = new BasicDBObject();
dbObject.put("_value", source.toString());
dbObject.put("_class", source.getClass().getName());
return dbObject;
}
}
private static class StringToUuidConverter implements Converter<String, UUID> {
public UUID convert(String source) {
return UUID.fromString(source);
private static class DBObjectToUUIDConverter implements Converter<DBObject, UUID> {
public UUID convert(DBObject source) {
UUID id = UUID.fromString((String) source.get("_value"));
return id;
}
}
@@ -460,64 +438,77 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore implements Me
*/
private static final class MessageWrapper {
private volatile Object groupId;
private volatile Object _groupId;
@Transient
private final Message<?> message;
private volatile long groupTimestamp;
private volatile long lastModified;
private volatile int lastReleasedSequenceNumber;
private final Object payload;
private volatile boolean completedGroup;
@SuppressWarnings("unused")
private final Map<String, ?> headers;
@SuppressWarnings("unused")
private final String _payloadType;
private volatile long _group_timestamp;
private volatile long _group_update_timestamp;
private volatile int _last_released_sequence;
private volatile boolean _group_complete;
public MessageWrapper(Message<?> message) {
Assert.notNull(message, "'message' must not be null");
this.message = message;
}
public int getLastReleasedSequenceNumber() {
return lastReleasedSequenceNumber;
}
public long getGroupTimestamp() {
return groupTimestamp;
this.payload = message.getPayload();
this.headers = message.getHeaders();
this._payloadType = this.payload.getClass().getName();
}
public boolean isCompletedGroup() {
return completedGroup;
public int get_LastReleasedSequenceNumber() {
return _last_released_sequence;
}
public Object getGroupId() {
return groupId;
public long get_Group_timestamp() {
return _group_timestamp;
}
public boolean get_Group_complete() {
return _group_complete;
}
public Object get_GroupId() {
return _groupId;
}
public Message<?> getMessage() {
return message;
}
public void setGroupId(Object groupId) {
this.groupId = groupId;
public void set_GroupId(Object groupId) {
this._groupId = groupId;
}
public void setGroupTimestamp(long groupTimestamp) {
this.groupTimestamp = groupTimestamp;
}
public long getLastModified() {
return lastModified;
public void set_Group_timestamp(long groupTimestamp) {
this._group_timestamp = groupTimestamp;
}
public void setLastModified(long lastModified) {
this.lastModified = lastModified;
public long get_Group_update_timestamp() {
return _group_update_timestamp;
}
public void setLastReleasedSequenceNumber(int lastReleasedSequenceNumber) {
this.lastReleasedSequenceNumber = lastReleasedSequenceNumber;
public void set_Group_update_timestamp(long lastModified) {
this._group_update_timestamp = lastModified;
}
public void setCompletedGroup(boolean completedGroup) {
this.completedGroup = completedGroup;
public void set_LastReleasedSequenceNumber(int lastReleasedSequenceNumber) {
this._last_released_sequence = lastReleasedSequenceNumber;
}
public void set_Group_complete(boolean completedGroup) {
this._group_complete = completedGroup;
}
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2007-2011 the original author or authors
* Copyright 2007-2012 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.
@@ -49,22 +49,22 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
@Test
@MongoDbAvailable
public void testNonExistingEmptyMessageGroup() throws Exception{
public void testNonExistingEmptyMessageGroup() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
MessageGroup messageGroup = store.getMessageGroup(1);
assertNotNull(messageGroup);
assertTrue(messageGroup instanceof SimpleMessageGroup);
assertEquals(0, messageGroup.size());
}
@Test
@MongoDbAvailable
public void testMessageGroupWithAddedMessage() throws Exception{
public void testMessageGroupWithAddedMessagePrimitiveGroupId() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
MessageGroup messageGroup = store.getMessageGroup(1);
Message<?> messageA = new GenericMessage<String>("A");
Message<?> messageB = new GenericMessage<String>("B");
@@ -77,23 +77,47 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
// ensure that 'message_group' header that is only used internally is not propagated
assertNull(retrievedMessage.getHeaders().get("message_group"));
}
@Test
@MongoDbAvailable
public void testCountMessagesInGroup() throws Exception{
public void testMessageGroupWithAddedMessageUUIDGroupIdAndUUIDHeader() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
Object id = UUID.randomUUID();
MessageGroup messageGroup = store.getMessageGroup(id);
UUID uuidA = UUID.randomUUID();
Message<?> messageA = MessageBuilder.withPayload("A").setHeader("foo", uuidA).build();
UUID uuidB = UUID.randomUUID();
Message<?> messageB = MessageBuilder.withPayload("B").setHeader("foo", uuidB).build();
store.addMessageToGroup(id, messageA);
messageGroup = store.addMessageToGroup(id, messageB);
assertEquals(2, messageGroup.size());
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
assertNull(retrievedMessage.getHeaders().get("message_group"));
Object fooHeader = retrievedMessage.getHeaders().get("foo");
assertTrue(fooHeader instanceof UUID);
assertEquals(uuidA, fooHeader);
}
@Test
@MongoDbAvailable
public void testCountMessagesInGroup() 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));
}
@Test
@MongoDbAvailable
public void testMessageGroupUpdatedDateChangesWithEachAddedMessage() throws Exception{
public void testMessageGroupUpdatedDateChangesWithEachAddedMessage() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
@@ -111,43 +135,43 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
updatedTimestamp = messageGroup.getLastModified();
assertTrue(updatedTimestamp > createdTimestamp);
assertEquals(2, messageGroup.size());
// make sure the store is properly rebuild from MongoDB
store = new MongoDbMessageStore(mongoDbFactory);
messageGroup = store.getMessageGroup(1);
assertEquals(2, messageGroup.size());
}
@Test
@MongoDbAvailable
public void testMessageGroupMarkingMessage() throws Exception{
public void testMessageGroupMarkingMessage() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
MessageGroup messageGroup = store.getMessageGroup(1);
Message<?> messageA = new GenericMessage<String>("A");
Message<?> messageB = new GenericMessage<String>("B");
store.addMessageToGroup(1, messageA);
messageGroup = store.addMessageToGroup(1, messageB);
assertEquals(2, messageGroup.size());
messageGroup = store.removeMessageFromGroup(1, messageA);
assertEquals(1, messageGroup.size());
// validate that the updates were propagated to Mongo as well
store = new MongoDbMessageStore(mongoDbFactory);
messageGroup = store.getMessageGroup(1);
assertEquals(1, messageGroup.size());
}
@Test
@MongoDbAvailable
public void testRemoveMessageGroup() throws Exception{
public void testRemoveMessageGroup() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
MessageGroup messageGroup = store.getMessageGroup(1);
Message<?> message = new GenericMessage<String>("Hello");
UUID id = message.getHeaders().getId();
@@ -155,20 +179,20 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
assertEquals(1, messageGroup.size());
message = store.getMessage(id);
assertNotNull(message);
store.removeMessageGroup(1);
MessageGroup messageGroupA = store.getMessageGroup(1);
assertEquals(0, messageGroupA.size());
assertFalse(messageGroupA.equals(messageGroup));
}
@Test
@MongoDbAvailable
public void testCompleteMessageGroup() throws Exception{
public void testCompleteMessageGroup() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
MessageGroup messageGroup = store.getMessageGroup(1);
Message<?> message = new GenericMessage<String>("Hello");
store.addMessageToGroup(messageGroup.getGroupId(), message);
@@ -176,13 +200,13 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
messageGroup = store.getMessageGroup(1);
assertTrue(messageGroup.isComplete());
}
@Test
@MongoDbAvailable
public void testLastReleasedSequenceNumber() throws Exception{
public void testLastReleasedSequenceNumber() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
MessageGroup messageGroup = store.getMessageGroup(1);
Message<?> message = new GenericMessage<String>("Hello");
store.addMessageToGroup(messageGroup.getGroupId(), message);
@@ -190,28 +214,28 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
messageGroup = store.getMessageGroup(1);
assertEquals(5, messageGroup.getLastReleasedMessageSequenceNumber());
}
@Test
@MongoDbAvailable
public void testRemoveMessageFromTheGroup() throws Exception{
public void testRemoveMessageFromTheGroup() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
MessageGroup messageGroup = store.getMessageGroup(1);
Message<?> message = new GenericMessage<String>("2");
store.addMessageToGroup(1, new GenericMessage<String>("1"));
store.addMessageToGroup(1, message);
messageGroup = store.addMessageToGroup(1, new GenericMessage<String>("3"));
assertEquals(3, messageGroup.size());
messageGroup = store.removeMessageFromGroup(1, message);
assertEquals(2, messageGroup.size());
}
@Test
@MongoDbAvailable
public void testMultipleMessageStores() throws Exception{
public void testMultipleMessageStores() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store1 = new MongoDbMessageStore(mongoDbFactory);
MongoDbMessageStore store2 = new MongoDbMessageStore(mongoDbFactory);
@@ -220,33 +244,33 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
store1.addMessageToGroup(1, message);
store2.addMessageToGroup(1, new GenericMessage<String>("2"));
store1.addMessageToGroup(1, new GenericMessage<String>("3"));
MongoDbMessageStore store3 = new MongoDbMessageStore(mongoDbFactory);
MessageGroup messageGroup = store3.getMessageGroup(1);
assertEquals(3, messageGroup.size());
store3.removeMessageFromGroup(1, message);
messageGroup = store2.getMessageGroup(1);
assertEquals(2, messageGroup.size());
}
@Test
@MongoDbAvailable
public void testMessageGroupIterator() throws Exception{
public void testMessageGroupIterator() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store1 = new MongoDbMessageStore(mongoDbFactory);
MongoDbMessageStore store2 = new MongoDbMessageStore(mongoDbFactory);
Message<?> message = new GenericMessage<String>("1");
store2.addMessageToGroup(1, message);
store1.addMessageToGroup(2, new GenericMessage<String>("2"));
store2.addMessageToGroup(3, new GenericMessage<String>("3"));
MongoDbMessageStore store3 = new MongoDbMessageStore(mongoDbFactory);
Iterator<MessageGroup> iterator = store3.iterator();
int counter = 0;
while (iterator.hasNext()) {
@@ -254,9 +278,9 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
counter++;
}
assertEquals(3, counter);
store2.removeMessageFromGroup(1, message);
iterator = store3.iterator();
counter = 0;
while (iterator.hasNext()) {
@@ -265,59 +289,59 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
}
assertEquals(2, counter);
}
// @Test
// @MongoDbAvailable
// public void testConcurrentModifications() throws Exception{
// public void testConcurrentModifications() throws Exception{
// MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
// final MongoDbMessageStore store1 = new MongoDbMessageStore(mongoDbFactory);
// final MongoDbMessageStore store2 = new MongoDbMessageStore(mongoDbFactory);
//
//
// final Message<?> message = new GenericMessage<String>("1");
//
// ExecutorService executor = null;
//
//
// final List<Object> failures = new ArrayList<Object>();
//
//
// for (int i = 0; i < 100; i++) {
// executor = Executors.newCachedThreadPool();
//
// executor.execute(new Runnable() {
// public void run() {
//
// executor.execute(new Runnable() {
// public void run() {
// MessageGroup group = store1.addMessageToGroup(1, message);
// if (group.getUnmarked().size() != 1){
// failures.add("ADD");
// throw new AssertionFailedError("Failed on ADD");
// }
// }
// }
// });
// executor.execute(new Runnable() {
// executor.execute(new Runnable() {
// public void run() {
// MessageGroup group = store2.removeMessageFromGroup(1, message);
// if (group.getUnmarked().size() != 0){
// failures.add("REMOVE");
// throw new AssertionFailedError("Failed on Remove");
// }
// }
// }
// });
//
//
// executor.shutdown();
// executor.awaitTermination(10, TimeUnit.SECONDS);
// store2.removeMessageFromGroup(1, message); // ensures that if ADD thread executed after REMOVE, the store is empty for the next cycle
// }
// assertTrue(failures.size() == 0);
// }
@Test
@MongoDbAvailable
public void testWithAggregatorWithShutdown() throws Exception{
public void testWithAggregatorWithShutdown() throws Exception{
this.prepareMongoFactory(); // for this test it only ensures that DB was flushed before test
ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("mongo-aggregator-config.xml", this.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();
input.send(m1);
@@ -325,36 +349,36 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
input.send(m2);
assertNull(output.receive(1000));
context.close();
context = new ClassPathXmlApplicationContext("mongo-aggregator-config.xml", this.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();
input.send(m3);
assertNotNull(output.receive(2000));
}
@Test
@MongoDbAvailable
public void testWithMessageHistory() throws Exception{
public void testWithMessageHistory() throws Exception{
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoDbMessageStore store = new MongoDbMessageStore(mongoDbFactory);
store.getMessageGroup(1);
Message<?> message = new GenericMessage<String>("Hello");
DirectChannel fooChannel = new DirectChannel();
fooChannel.setBeanName("fooChannel");
DirectChannel barChannel = new DirectChannel();
barChannel.setBeanName("barChannel");
message = MessageHistory.write(message, fooChannel);
message = MessageHistory.write(message, barChannel);
store.addMessageToGroup(1, message);
message = store.getMessageGroup(1).getMessages().iterator().next();
MessageHistory messageHistory = MessageHistory.read(message);
assertNotNull(messageHistory);
assertEquals(2, messageHistory.size());
@@ -362,5 +386,5 @@ public class MongoDbMessageGroupStoreTests extends MongoDbAvailableTests {
assertEquals("fooChannel", fooChannelHistory.get("name"));
assertEquals("channel", fooChannelHistory.get("type"));
}
}