INT-2479 UUID Header Conversion
Refactored MongoDbMessageStore to ensure that the Message metadata is persisted consistently using already provided converters. Major refactoring went into MessageWrapper and MongoDbMessageStore.write method which was greatly simplified based on the fact that MessageWrapper now represents the structure that needs to be persisted INT-2479 fixed default converters to make sure that UUID is stored in Mongo with type information, so the actuall UUID type could be restored
This commit is contained in:
committed by
Gary Russell
parent
2df1240220
commit
a37c7a0821
@@ -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"'
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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"));
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user