diff --git a/build.gradle b/build.gradle index 49ad7e896f..77c9d2b6e1 100644 --- a/build.gradle +++ b/build.gradle @@ -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"' diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java index 56ef77716d..ce0aa3012d 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java @@ -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 groupedMessages = this.template.find(whereGroupIdExists(), MessageWrapper.class, this.collectionName); Map messageGroups = new HashMap(); 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 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> customConverters = new ArrayList>(); - 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 headers = this.normalizeHeaders((Map) 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 { - public String convert(UUID source) { - return source.toString(); + private static class UuidToDBObjectConverter implements Converter { + 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 { - public UUID convert(String source) { - return UUID.fromString(source); + private static class DBObjectToUUIDConverter implements Converter { + 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 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; } } } diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoDbMessageGroupStoreTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoDbMessageGroupStoreTests.java index e3d07a56af..85d1494265 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoDbMessageGroupStoreTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/MongoDbMessageGroupStoreTests.java @@ -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("A"); Message messageB = new GenericMessage("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("A"); Message messageB = new GenericMessage("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("A"); Message messageB = new GenericMessage("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("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("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("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("2"); store.addMessageToGroup(1, new GenericMessage("1")); store.addMessageToGroup(1, message); messageGroup = store.addMessageToGroup(1, new GenericMessage("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("2")); store1.addMessageToGroup(1, new GenericMessage("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("1"); store2.addMessageToGroup(1, message); store1.addMessageToGroup(2, new GenericMessage("2")); store2.addMessageToGroup(3, new GenericMessage("3")); - + MongoDbMessageStore store3 = new MongoDbMessageStore(mongoDbFactory); - + Iterator 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("1"); // // ExecutorService executor = null; -// +// // final List failures = new ArrayList(); -// +// // 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("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")); } - + }