Prepare for Milestone Release

Change all the `BUILD-SNAPSHOT`s to their latest Milestones
Fix compatibility with those Milestones

Upgrade to Spring AMQP 2.0.0.M1

Upgrade to Spring Data Kay

Gemfire now is based on the Apache Geode, so changed all the imports to proper new packages
MongoDB now is based on the Mongo 3 Driver, therefore many breaking changes. Mostly to the mapping part

Fix Checkstyle violation

Fix race condition in the  `MongoDbInboundChannelAdapterIntegrationTests` around `QueueChannel` and tx commit

Add `-s` to Travis Gradle command to see stack trace about `MongoDbMetadataStoreTests` problem

Test only MongoDB module on Travis with -d

Add addon to Travis to pull MongoDB-3.0

Fix `MongoDbAvailableRule` for MongoDB 3.0 Driver style

Increase `serverSelectionTimeout` to `100` in the `MongoDbAvailableRule`.
Looks like `0` isn't good value to get immediate answer
This commit is contained in:
Artem Bilan
2016-11-30 07:32:51 -05:00
committed by Gary Russell
parent 9945d04edd
commit c0a507c36c
45 changed files with 346 additions and 263 deletions

View File

@@ -5,6 +5,12 @@ services:
- mongodb
- rabbitmq
- redis-server
addons:
apt:
sources:
- mongodb-3.0-precise
packages:
- mongodb-org-server
before_cache:
- rm -f $HOME/.gradle/caches/modules-2/modules-2.lock
cache:

View File

@@ -122,16 +122,16 @@ subprojects { subproject ->
slf4jVersion = "1.7.21"
tomcatVersion = "8.0.33"
smackVersion = '4.1.7'
springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.0.0.BUILD-SNAPSHOT'
springDataJpaVersion = '1.11.0.BUILD-SNAPSHOT'
springDataMongoVersion = '1.10.0.BUILD-SNAPSHOT'
springDataRedisVersion = '1.8.0.BUILD-SNAPSHOT'
springGemfireVersion = '1.9.0.BUILD-SNAPSHOT'
springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.0.0.M1'
springDataJpaVersion = '2.0.0.M1'
springDataMongoVersion = '2.0.0.M1'
springDataRedisVersion = '2.0.0.M1'
springGemfireVersion = '2.0.0.M1'
springSecurityVersion = '4.2.0.RELEASE'
springSocialTwitterVersion = '1.1.2.RELEASE'
springRetryVersion = '1.2.0.RC1'
springVersion = project.hasProperty('springVersion') ? project.springVersion : '5.0.0.BUILD-SNAPSHOT'
springWsVersion = '2.3.0.RELEASE'
springVersion = project.hasProperty('springVersion') ? project.springVersion : '5.0.0.M3'
springWsVersion = '2.4.0.RELEASE'
xmlUnitVersion = '1.6'
xstreamVersion = '1.4.7'
}

View File

@@ -86,7 +86,7 @@ public class LambdaMessageProcessor implements MessageProcessor<Object>, BeanFac
public void setBeanFactory(BeanFactory beanFactory) throws BeansException {
ConversionService conversionService = IntegrationUtils.getConversionService(beanFactory);
if (conversionService == null) {
conversionService = DefaultConversionService.getSharedInstance();
conversionService = new DefaultConversionService();
}
this.conversionService = conversionService;
}

View File

@@ -17,16 +17,13 @@
package org.springframework.integration.expression;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNotSame;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertSame;
import org.junit.Test;
import org.springframework.beans.factory.support.RootBeanDefinition;
import org.springframework.context.support.ConversionServiceFactoryBean;
import org.springframework.context.support.GenericApplicationContext;
import org.springframework.core.convert.support.DefaultConversionService;
import org.springframework.expression.TypeConverter;
import org.springframework.expression.spel.support.StandardEvaluationContext;
import org.springframework.integration.config.IntegrationEvaluationContextFactoryBean;
@@ -36,8 +33,8 @@ import org.springframework.integration.test.util.TestUtils;
/**
* @author Gary Russell
* @since 3.0
*
* @since 3.0
*/
public class ExpressionUtilsTests {
@@ -52,8 +49,9 @@ public class ExpressionUtilsTests {
StandardEvaluationContext evalContext = ExpressionUtils.createStandardEvaluationContext(context);
assertNotNull(evalContext.getBeanResolver());
assertNotNull(evalContext.getTypeConverter());
IntegrationEvaluationContextFactoryBean factory = context.getBean("&" + IntegrationContextUtils.INTEGRATION_EVALUATION_CONTEXT_BEAN_NAME,
IntegrationEvaluationContextFactoryBean.class);
IntegrationEvaluationContextFactoryBean factory =
context.getBean("&" + IntegrationContextUtils.INTEGRATION_EVALUATION_CONTEXT_BEAN_NAME,
IntegrationEvaluationContextFactoryBean.class);
assertSame(evalContext.getTypeConverter(), TestUtils.getPropertyValue(factory, "typeConverter"));
}
@@ -68,8 +66,6 @@ public class ExpressionUtilsTests {
assertNotNull(evalContext.getBeanResolver());
TypeConverter typeConverter = evalContext.getTypeConverter();
assertNotNull(typeConverter);
assertSame(DefaultConversionService.getSharedInstance(),
TestUtils.getPropertyValue(typeConverter, "conversionService"));
}
@Test
@@ -82,8 +78,6 @@ public class ExpressionUtilsTests {
assertNotNull(evalContext.getBeanResolver());
TypeConverter typeConverter = evalContext.getTypeConverter();
assertNotNull(typeConverter);
assertNotSame(DefaultConversionService.getSharedInstance(),
TestUtils.getPropertyValue(typeConverter, "conversionService"));
assertSame(context.getBean(IntegrationUtils.INTEGRATION_CONVERSION_SERVICE_BEAN_NAME),
TestUtils.getPropertyValue(typeConverter, "conversionService"));
}
@@ -94,7 +88,5 @@ public class ExpressionUtilsTests {
assertNull(evalContext.getBeanResolver());
TypeConverter typeConverter = evalContext.getTypeConverter();
assertNotNull(typeConverter);
assertSame(DefaultConversionService.getSharedInstance(),
TestUtils.getPropertyValue(typeConverter, "conversionService"));
}
}

View File

@@ -37,7 +37,7 @@ import org.springframework.integration.redis.metadata.RedisMetadataStore;
import org.springframework.integration.redis.rules.RedisAvailable;
import org.springframework.integration.redis.rules.RedisAvailableTests;
import com.gemstone.gemfire.cache.CacheFactory;
import org.apache.geode.cache.CacheFactory;
/**
* @author Gary Russell

View File

@@ -27,11 +27,11 @@ import org.springframework.integration.endpoint.ExpressionMessageProducerSupport
import org.springframework.messaging.Message;
import org.springframework.util.Assert;
import com.gemstone.gemfire.cache.CacheClosedException;
import com.gemstone.gemfire.cache.CacheListener;
import com.gemstone.gemfire.cache.EntryEvent;
import com.gemstone.gemfire.cache.Region;
import com.gemstone.gemfire.cache.util.CacheListenerAdapter;
import org.apache.geode.cache.CacheClosedException;
import org.apache.geode.cache.CacheListener;
import org.apache.geode.cache.EntryEvent;
import org.apache.geode.cache.Region;
import org.apache.geode.cache.util.CacheListenerAdapter;
/**
* An inbound endpoint that listens to a GemFire region for events and then publishes Messages to

View File

@@ -30,12 +30,12 @@ import org.springframework.integration.endpoint.ExpressionMessageProducerSupport
import org.springframework.messaging.Message;
import org.springframework.util.Assert;
import com.gemstone.gemfire.cache.query.CqEvent;
import org.apache.geode.cache.query.CqEvent;
/**
* Responds to a Gemfire continuous query (set using the #query field) that is
* constantly evaluated against a cache
* {@link com.gemstone.gemfire.cache.Region}. This is much faster than
* {@link org.apache.geode.cache.Region}. This is much faster than
* re-querying the cache manually.
*
* @author Josh Long

View File

@@ -19,9 +19,9 @@ package org.springframework.integration.gemfire.metadata;
import org.springframework.integration.metadata.ConcurrentMetadataStore;
import org.springframework.util.Assert;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.Region;
import com.gemstone.gemfire.cache.Scope;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.Region;
import org.apache.geode.cache.Scope;
/**
* Gemfire implementation of {@link ConcurrentMetadataStore}.

View File

@@ -32,9 +32,9 @@ import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHandler;
import org.springframework.util.Assert;
import com.gemstone.gemfire.GemFireCheckedException;
import com.gemstone.gemfire.GemFireException;
import com.gemstone.gemfire.cache.Region;
import org.apache.geode.GemFireCheckedException;
import org.apache.geode.GemFireException;
import org.apache.geode.cache.Region;
/**
* A {@link MessageHandler} implementation that writes to a GemFire Region. The

View File

@@ -29,8 +29,8 @@ import org.springframework.integration.store.MessageStore;
import org.springframework.util.Assert;
import org.springframework.util.PatternMatchUtils;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.Region;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.Region;
/**
* Gemfire implementation of the key/value style {@link MessageStore} and

View File

@@ -21,9 +21,9 @@ import java.util.concurrent.locks.Lock;
import org.springframework.integration.support.locks.LockRegistry;
import org.springframework.util.Assert;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.Region;
import com.gemstone.gemfire.cache.Scope;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.Region;
import org.apache.geode.cache.Scope;
/**
* Implementation of {@link LockRegistry} providing a distributed lock using Gemfire.

View File

@@ -16,8 +16,8 @@
package org.springframework.integration.gemfire;
import com.gemstone.gemfire.cache.EntryEvent;
import com.gemstone.gemfire.cache.util.CacheListenerAdapter;
import org.apache.geode.cache.EntryEvent;
import org.apache.geode.cache.util.CacheListenerAdapter;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;

View File

@@ -8,7 +8,7 @@
http://www.springframework.org/schema/beans/spring-beans.xsd">
<bean id="region" class="org.mockito.Mockito" factory-method="mock">
<constructor-arg value="com.gemstone.gemfire.cache.Region"/>
<constructor-arg value="org.apache.geode.cache.Region"/>
</bean>
<int-gfe:inbound-channel-adapter id="channel1"

View File

@@ -22,7 +22,7 @@
</int-gfe:outbound-channel-adapter>
<bean id="region" class="org.mockito.Mockito" factory-method="mock">
<constructor-arg value="com.gemstone.gemfire.cache.Region"/>
<constructor-arg value="org.apache.geode.cache.Region"/>
</bean>
</beans>

View File

@@ -24,12 +24,12 @@ import java.util.Properties;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.CacheFactory;
import com.gemstone.gemfire.cache.Region;
import com.gemstone.gemfire.cache.RegionShortcut;
import com.gemstone.gemfire.cache.Scope;
import com.gemstone.gemfire.cache.server.CacheServer;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.CacheFactory;
import org.apache.geode.cache.Region;
import org.apache.geode.cache.RegionShortcut;
import org.apache.geode.cache.Scope;
import org.apache.geode.cache.server.CacheServer;
/**
* @author Costin Leau

View File

@@ -33,7 +33,7 @@ import org.springframework.expression.spel.standard.SpelExpressionParser;
import org.springframework.integration.channel.QueueChannel;
import org.springframework.messaging.Message;
import com.gemstone.gemfire.cache.Region;
import org.apache.geode.cache.Region;
/**
* @author Mark Fisher

View File

@@ -19,9 +19,12 @@ package org.springframework.integration.gemfire.inbound;
import static org.junit.Assert.assertEquals;
import static org.mockito.Mockito.mock;
import org.apache.geode.cache.Operation;
import org.apache.geode.cache.query.CqEvent;
import org.apache.geode.cache.query.CqQuery;
import org.apache.geode.cache.query.internal.cq.ServerCQImpl;
import org.junit.Before;
import org.junit.Test;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.data.gemfire.listener.ContinuousQueryListenerContainer;
import org.springframework.expression.spel.standard.SpelExpressionParser;
@@ -30,11 +33,6 @@ import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessagingException;
import com.gemstone.gemfire.cache.Operation;
import com.gemstone.gemfire.cache.query.CqEvent;
import com.gemstone.gemfire.cache.query.CqQuery;
import com.gemstone.gemfire.cache.query.internal.CqQueryImpl;
/**
* @author David Turanski
* @author Artem Bilan
@@ -107,7 +105,7 @@ public class ContinuousQueryMessageProducerTests {
return new CqEvent() {
final CqQuery cq = new CqQueryImpl();
final CqQuery cq = new ServerCQImpl();
final byte[] ba = new byte[0];

View File

@@ -38,8 +38,8 @@ import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import com.gemstone.gemfire.cache.query.CqEvent;
import com.gemstone.gemfire.internal.cache.LocalRegion;
import org.apache.geode.cache.query.CqEvent;
import org.apache.geode.internal.cache.LocalRegion;
/**
* @author David Turanski

View File

@@ -32,8 +32,8 @@ import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import com.gemstone.gemfire.cache.EntryEvent;
import com.gemstone.gemfire.internal.cache.LocalRegion;
import org.apache.geode.cache.EntryEvent;
import org.apache.geode.internal.cache.LocalRegion;
/**
* @author David Turanski

View File

@@ -30,9 +30,9 @@ import org.springframework.data.gemfire.GemfireTemplate;
import org.springframework.integration.metadata.ConcurrentMetadataStore;
import org.springframework.util.Assert;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.CacheFactory;
import com.gemstone.gemfire.cache.Region;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.CacheFactory;
import org.apache.geode.cache.Region;
/**
* @author Artem Bilan

View File

@@ -36,9 +36,9 @@ import org.springframework.integration.support.MessageBuilder;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.GenericMessage;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.Region;
import com.gemstone.gemfire.cache.Scope;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.Region;
import org.apache.geode.cache.Scope;
/**
* @author Mark Fisher

View File

@@ -33,7 +33,7 @@ import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import com.gemstone.gemfire.internal.cache.DistributedRegion;
import org.apache.geode.internal.cache.DistributedRegion;
/**
* @author David Turanski

View File

@@ -43,9 +43,9 @@ import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.PollableChannel;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.Region;
import com.gemstone.gemfire.cache.Scope;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.Region;
import org.apache.geode.cache.Scope;
/**

View File

@@ -48,9 +48,9 @@ import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.support.GenericMessage;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.Region;
import com.gemstone.gemfire.cache.Scope;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.Region;
import org.apache.geode.cache.Scope;
import junit.framework.AssertionFailedError;

View File

@@ -39,9 +39,9 @@ import org.springframework.integration.test.util.TestUtils;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.GenericMessage;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.Region;
import com.gemstone.gemfire.cache.Scope;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.Region;
import org.apache.geode.cache.Scope;
/**
* @author Mark Fisher

View File

@@ -19,6 +19,7 @@ package org.springframework.integration.mongodb.metadata;
import java.util.HashMap;
import java.util.Map;
import org.bson.Document;
import org.springframework.data.mongodb.MongoDbFactory;
import org.springframework.data.mongodb.core.FindAndModifyOptions;
import org.springframework.data.mongodb.core.MongoTemplate;
@@ -28,7 +29,6 @@ import org.springframework.data.mongodb.core.query.Update;
import org.springframework.integration.metadata.ConcurrentMetadataStore;
import org.springframework.util.Assert;
import com.mongodb.BasicDBObject;
import com.mongodb.DBCollection;
/**
@@ -109,10 +109,10 @@ public class MongoDbMetadataStore implements ConcurrentMetadataStore {
public void put(String key, String value) {
Assert.hasText(key, "'key' must not be empty.");
Assert.hasText(value, "'value' must not be empty.");
final Map<String, String> entry = new HashMap<>();
final Map<String, Object> entry = new HashMap<>();
entry.put(ID_FIELD, key);
entry.put(VALUE, value);
this.template.execute(this.collectionName, collection -> collection.save(new BasicDBObject(entry)));
this.template.save(new Document(entry), this.collectionName);
}
/**
@@ -195,7 +195,7 @@ public class MongoDbMetadataStore implements ConcurrentMetadataStore {
Assert.hasText(newValue, "'newValue' must not be empty.");
Query query = new Query(Criteria.where(ID_FIELD).is(key).and(VALUE).is(oldValue));
return this.template.updateFirst(query, Update.update(VALUE, newValue), this.collectionName)
.isUpdateOfExisting();
.getModifiedCount() > 0;
}
}

View File

@@ -25,7 +25,6 @@ import java.util.UUID;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.beans.BeansException;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.context.ApplicationContext;
@@ -44,7 +43,8 @@ import org.springframework.data.mongodb.core.mapping.MongoMappingContext;
import org.springframework.data.mongodb.core.query.Criteria;
import org.springframework.data.mongodb.core.query.Query;
import org.springframework.data.mongodb.core.query.Update;
import org.springframework.integration.mongodb.support.MongoDbMessageBytesConverter;
import org.springframework.integration.mongodb.support.BinaryToMessageConverter;
import org.springframework.integration.mongodb.support.MessageToBinaryConverter;
import org.springframework.integration.store.AbstractMessageGroupStore;
import org.springframework.integration.store.BasicMessageGroupStore;
import org.springframework.integration.store.MessageGroup;
@@ -131,7 +131,8 @@ public abstract class AbstractConfigurableMongoDbMessageStore extends AbstractMe
new MongoMappingContext());
this.mappingMongoConverter.setApplicationContext(this.applicationContext);
List<Object> customConverters = new ArrayList<Object>();
customConverters.add(new MongoDbMessageBytesConverter());
customConverters.add(new MessageToBinaryConverter());
customConverters.add(new BinaryToMessageConverter());
this.mappingMongoConverter.setCustomConversions(new CustomConversions(customConverters));
this.mappingMongoConverter.afterPropertiesSet();
}

View File

@@ -259,9 +259,8 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb
List<MessageGroup> messageGroups = new ArrayList<MessageGroup>();
Query query = Query.query(Criteria.where(MessageDocumentFields.GROUP_ID).exists(true));
@SuppressWarnings("rawtypes")
List groupIds = mongoTemplate.getCollection(collectionName)
.distinct(MessageDocumentFields.GROUP_ID, query.getQueryObject());
Iterable<String> groupIds = mongoTemplate.getCollection(collectionName)
.distinct(MessageDocumentFields.GROUP_ID, query.getQueryObject(), String.class);
for (Object groupId : groupIds) {
messageGroups.add(getMessageGroup(groupId));
@@ -310,7 +309,8 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb
public int getMessageGroupCount() {
Query query = Query.query(Criteria.where(MessageDocumentFields.GROUP_ID).exists(true));
return this.mongoTemplate.getCollection(this.collectionName)
.distinct(MessageDocumentFields.GROUP_ID, query.getQueryObject())
.distinct(MessageDocumentFields.GROUP_ID, query.getQueryObject(), Object.class)
.into(new ArrayList<>())
.size();
}

View File

@@ -25,7 +25,11 @@ import java.util.Map;
import java.util.Map.Entry;
import java.util.Properties;
import java.util.UUID;
import java.util.stream.Collectors;
import org.bson.Document;
import org.bson.conversions.Bson;
import org.bson.types.Binary;
import org.springframework.beans.BeansException;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.BeanClassLoaderAware;
@@ -37,10 +41,12 @@ import org.springframework.core.serializer.support.DeserializingConverter;
import org.springframework.core.serializer.support.SerializingConverter;
import org.springframework.data.annotation.Id;
import org.springframework.data.annotation.Transient;
import org.springframework.data.convert.ReadingConverter;
import org.springframework.data.convert.WritingConverter;
import org.springframework.data.domain.Sort;
import org.springframework.data.mapping.context.MappingContext;
import org.springframework.data.mongodb.MongoDbFactory;
import org.springframework.data.mongodb.core.BulkOperations;
import org.springframework.data.mongodb.core.FindAndModifyOptions;
import org.springframework.data.mongodb.core.IndexOperations;
import org.springframework.data.mongodb.core.MongoTemplate;
@@ -74,8 +80,6 @@ import org.springframework.util.ClassUtils;
import org.springframework.util.StringUtils;
import com.mongodb.BasicDBList;
import com.mongodb.BasicDBObject;
import com.mongodb.BulkWriteOperation;
import com.mongodb.DBObject;
@@ -230,7 +234,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
@Override
@ManagedAttribute
public long getMessageCount() {
return this.template.getCollection(this.collectionName).getCount();
return this.template.getCollection(this.collectionName).count();
}
@Override
@@ -300,7 +304,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
Assert.notNull(groupId, "'groupId' must not be null");
Assert.notNull(messages, "'messageToRemove' must not be null");
Collection<UUID> ids = new ArrayList<UUID>();
Collection<UUID> ids = new ArrayList<>();
for (Message<?> messageToRemove : messages) {
ids.add(messageToRemove.getHeaders().getId());
if (ids.size() >= getRemoveBatchSize()) {
@@ -315,13 +319,12 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
}
private void bulkRemove(Object groupId, Collection<UUID> ids) {
BulkWriteOperation bulkOp = this.template.getCollection(this.collectionName)
.initializeOrderedBulkOperation();
BulkOperations bulkOperations = this.template.bulkOps(BulkOperations.BulkMode.ORDERED, this.collectionName);
for (UUID id : ids) {
bulkOp.find(whereMessageIdIsAndGroupIdIs(id, groupId).getQueryObject())
.remove();
bulkOperations.remove(whereMessageIdIsAndGroupIdIs(id, groupId));
}
bulkOp.execute();
bulkOperations.execute();
}
@Override
@@ -331,13 +334,13 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
@Override
public Iterator<MessageGroup> iterator() {
List<MessageGroup> messageGroups = new ArrayList<MessageGroup>();
List<MessageGroup> messageGroups = new ArrayList<>();
Query query = Query.query(Criteria.where(GROUP_ID_KEY).exists(true));
@SuppressWarnings("rawtypes")
List groupIds = this.template.getCollection(this.collectionName)
.distinct(GROUP_ID_KEY, query.getQueryObject());
Iterable<String> groupIds = this.template.getCollection(this.collectionName)
.distinct(GROUP_ID_KEY, query.getQueryObject(), String.class);
for (Object groupId : groupIds) {
messageGroups.add(getMessageGroup(groupId));
@@ -394,12 +397,10 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
Assert.notNull(groupId, "'groupId' must not be null");
Query query = whereGroupIdOrder(groupId);
List<MessageWrapper> messageWrappers = this.template.find(query, MessageWrapper.class, this.collectionName);
List<Message<?>> messages = new ArrayList<Message<?>>();
for (MessageWrapper messageWrapper : messageWrappers) {
messages.add(messageWrapper.getMessage());
}
return messages;
return messageWrappers.stream()
.map(MessageWrapper::getMessage)
.collect(Collectors.toList());
}
@Override
@@ -417,7 +418,8 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
public int getMessageGroupCount() {
Query query = Query.query(Criteria.where(MessageDocumentFields.GROUP_ID).exists(true));
return this.template.getCollection(this.collectionName)
.distinct(MessageDocumentFields.GROUP_ID, query.getQueryObject())
.distinct(MessageDocumentFields.GROUP_ID, query.getQueryObject(), Object.class)
.into(new ArrayList<>())
.size();
}
@@ -430,11 +432,11 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
*/
private static Query whereMessageIdIs(UUID id) {
return new Query(Criteria.where("headers.id._value").is(id.toString()));
return new Query(Criteria.where("headers.id").is(id));
}
private static Query whereMessageIdIsAndGroupIdIs(UUID id, Object groupId) {
return new Query(Criteria.where("headers.id._value").is(id.toString()).and(GROUP_ID_KEY).is(groupId));
return new Query(Criteria.where("headers.id").is(id).and(GROUP_ID_KEY).is(groupId));
}
private static Query whereGroupIdOrder(Object groupId) {
@@ -469,6 +471,20 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
innerMap.put(MessageHeaders.TIMESTAMP, headers.get(MessageHeaders.TIMESTAMP));
}
@SuppressWarnings("unchecked")
private static Map<String, Object> asMap(Bson bson) {
if (bson instanceof Document) {
return (Document) bson;
}
if (bson instanceof DBObject) {
return ((DBObject) bson).toMap();
}
throw new IllegalArgumentException(
String.format("Cannot read %s. as map. Given Bson must be a Document or DBObject!", bson.getClass()));
}
/**
* Custom implementation of the {@link MappingMongoConverter} strategy.
@@ -482,56 +498,56 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
@Override
public void afterPropertiesSet() {
List<Object> customConverters = new ArrayList<Object>();
customConverters.add(new UuidToDBObjectConverter());
customConverters.add(new DBObjectToUUIDConverter());
customConverters.add(new MessageHistoryToDBObjectConverter());
customConverters.add(new DBObjectToGenericMessageConverter());
customConverters.add(new DBObjectToMutableMessageConverter());
customConverters.add(new DBObjectToErrorMessageConverter());
customConverters.add(new DBObjectToAdviceMessageConverter());
List<Object> customConverters = new ArrayList<>();
customConverters.add(new MessageHistoryToDocumentConverter());
customConverters.add(new DocumentToGenericMessageConverter());
customConverters.add(new DocumentToMutableMessageConverter());
customConverters.add(new DocumentToErrorMessageConverter());
customConverters.add(new DocumentToAdviceMessageConverter());
customConverters.add(new ThrowableToBytesConverter());
this.setCustomConversions(new CustomConversions(customConverters));
super.afterPropertiesSet();
}
@Override
public void write(Object source, DBObject target) {
public void write(Object source, Bson target) {
Assert.isInstanceOf(MessageWrapper.class, source);
target.put(CREATED_DATE, System.currentTimeMillis());
asMap(target).put(CREATED_DATE, System.currentTimeMillis());
super.write(source, target);
}
@Override
@SuppressWarnings({ "unchecked" })
public <S> S read(Class<S> clazz, DBObject source) {
@SuppressWarnings({"unchecked"})
public <S> S read(Class<S> clazz, Bson source) {
if (!MessageWrapper.class.equals(clazz)) {
return super.read(clazz, source);
}
if (source != null) {
Map<String, Object> sourceMap = asMap(source);
Message<?> message = null;
Object messageType = source.get("_messageType");
Object messageType = sourceMap.get("_messageType");
if (messageType == null) {
messageType = GenericMessage.class.getName();
}
try {
message = (Message<?>) this.read(ClassUtils.forName(messageType.toString(), MongoDbMessageStore.this.classLoader), source);
message = (Message<?>) read(ClassUtils.forName(messageType.toString(),
MongoDbMessageStore.this.classLoader), source);
}
catch (ClassNotFoundException e) {
throw new IllegalStateException("failed to load class: " + messageType, e);
}
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);
Long groupTimestamp = (Long) sourceMap.get(GROUP_TIMESTAMP_KEY);
Long lastModified = (Long) sourceMap.get(GROUP_UPDATE_TIMESTAMP_KEY);
Integer lastReleasedSequenceNumber = (Integer) sourceMap.get(LAST_RELEASED_SEQUENCE_NUMBER);
Boolean completeGroup = (Boolean) sourceMap.get(GROUP_COMPLETE_KEY);
MessageWrapper wrapper = new MessageWrapper(message);
if (source.containsField(GROUP_ID_KEY)) {
wrapper.set_GroupId(source.get(GROUP_ID_KEY));
if (sourceMap.containsKey(GROUP_ID_KEY)) {
wrapper.set_GroupId(sourceMap.get(GROUP_ID_KEY));
}
if (groupTimestamp != null) {
wrapper.set_Group_timestamp(groupTimestamp);
@@ -557,19 +573,20 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
for (Entry<String, Object> entry : headers.entrySet()) {
String headerName = entry.getKey();
Object headerValue = entry.getValue();
if (headerValue instanceof DBObject) {
DBObject source = (DBObject) headerValue;
if (headerValue instanceof Bson) {
Bson source = (Bson) headerValue;
Map<String, Object> document = asMap(source);
try {
Class<?> typeClass = null;
if (source.containsField("_class")) {
Object type = source.get("_class");
if (document.containsKey("_class")) {
Object type = document.get("_class");
typeClass = ClassUtils.forName(type.toString(), MongoDbMessageStore.this.classLoader);
}
else if (source instanceof BasicDBList) {
typeClass = List.class;
}
else {
throw new IllegalStateException("Unsupported 'DBObject' type: " + source.getClass());
throw new IllegalStateException("Unsupported 'Bson' type: " + source.getClass());
}
normalizedHeaders.put(headerName, super.read(typeClass, source));
}
@@ -584,14 +601,15 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
return normalizedHeaders;
}
private Object extractPayload(DBObject source) {
Object payload = source.get("payload");
if (payload instanceof DBObject) {
DBObject payloadObject = (DBObject) payload;
Object payloadType = payloadObject.get("_class");
private Object extractPayload(Bson source) {
Object payload = asMap(source).get("payload");
if (payload instanceof Bson) {
Bson payloadObject = (Bson) payload;
Object payloadType = asMap(payloadObject).get("_class");
try {
Class<?> payloadClass = ClassUtils.forName(payloadType.toString(), MongoDbMessageStore.this.classLoader);
payload = this.read(payloadClass, payloadObject);
payload = read(payloadClass, payloadObject);
}
catch (Exception e) {
throw new IllegalStateException("failed to load class: " + payloadType, e);
@@ -602,86 +620,61 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
}
private static class UuidToDBObjectConverter implements Converter<UUID, DBObject> {
UuidToDBObjectConverter() {
@WritingConverter
private static class MessageHistoryToDocumentConverter implements Converter<MessageHistory, Document> {
MessageHistoryToDocumentConverter() {
super();
}
@Override
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 DBObjectToUUIDConverter implements Converter<DBObject, UUID> {
DBObjectToUUIDConverter() {
super();
}
@Override
public UUID convert(DBObject source) {
return UUID.fromString((String) source.get("_value"));
}
}
private static class MessageHistoryToDBObjectConverter implements Converter<MessageHistory, DBObject> {
MessageHistoryToDBObjectConverter() {
super();
}
@Override
public DBObject convert(MessageHistory source) {
BasicDBObject obj = new BasicDBObject();
obj.put("_class", MessageHistory.class.getName());
public Document convert(MessageHistory source) {
BasicDBList dbList = new BasicDBList();
for (Properties properties : source) {
BasicDBObject dbo = new BasicDBObject();
dbo.put(MessageHistory.NAME_PROPERTY, properties.getProperty(MessageHistory.NAME_PROPERTY));
dbo.put(MessageHistory.TYPE_PROPERTY, properties.getProperty(MessageHistory.TYPE_PROPERTY));
dbo.put(MessageHistory.TIMESTAMP_PROPERTY, properties.getProperty(MessageHistory.TIMESTAMP_PROPERTY));
dbList.add(dbo);
Document historyProperty = new Document()
.append(MessageHistory.NAME_PROPERTY, properties.getProperty(MessageHistory.NAME_PROPERTY))
.append(MessageHistory.TYPE_PROPERTY, properties.getProperty(MessageHistory.TYPE_PROPERTY))
.append(MessageHistory.TIMESTAMP_PROPERTY,
properties.getProperty(MessageHistory.TIMESTAMP_PROPERTY));
dbList.add(historyProperty);
}
obj.put("components", dbList);
return obj;
return new Document("components", dbList)
.append("_class", MessageHistory.class.getName());
}
}
private class DBObjectToGenericMessageConverter implements Converter<DBObject, GenericMessage<?>> {
@ReadingConverter
private class DocumentToGenericMessageConverter implements Converter<Document, GenericMessage<?>> {
DBObjectToGenericMessageConverter() {
DocumentToGenericMessageConverter() {
super();
}
@Override
public GenericMessage<?> convert(DBObject source) {
public GenericMessage<?> convert(Document source) {
@SuppressWarnings("unchecked")
Map<String, Object> headers =
MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
GenericMessage<?> message =
new GenericMessage<Object>(MongoDbMessageStore.this.converter.extractPayload(source), headers);
new GenericMessage<>(MongoDbMessageStore.this.converter.extractPayload(source), headers);
enhanceHeaders(message.getHeaders(), headers);
return message;
}
}
private final class DBObjectToMutableMessageConverter implements Converter<DBObject, MutableMessage<?>> {
@ReadingConverter
private final class DocumentToMutableMessageConverter implements Converter<Document, MutableMessage<?>> {
DBObjectToMutableMessageConverter() {
DocumentToMutableMessageConverter() {
super();
}
@Override
public MutableMessage<?> convert(DBObject source) {
public MutableMessage<?> convert(Document source) {
@SuppressWarnings("unchecked")
Map<String, Object> headers =
MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
@@ -694,14 +687,15 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
}
private class DBObjectToAdviceMessageConverter implements Converter<DBObject, AdviceMessage<?>> {
@ReadingConverter
private class DocumentToAdviceMessageConverter implements Converter<Document, AdviceMessage<?>> {
DBObjectToAdviceMessageConverter() {
DocumentToAdviceMessageConverter() {
super();
}
@Override
public AdviceMessage<?> convert(DBObject source) {
public AdviceMessage<?> convert(Document source) {
@SuppressWarnings("unchecked")
Map<String, Object> headers =
MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
@@ -709,8 +703,8 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
Message<?> inputMessage = null;
if (source.get("inputMessage") != null) {
DBObject inputMessageObject = (DBObject) source.get("inputMessage");
Object inputMessageType = inputMessageObject.get("_class");
Bson inputMessageObject = (Bson) source.get("inputMessage");
Object inputMessageType = asMap(inputMessageObject).get("_class");
try {
Class<?> messageClass = ClassUtils.forName(inputMessageType.toString(),
MongoDbMessageStore.this.classLoader);
@@ -722,7 +716,7 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
}
}
AdviceMessage<?> message = new AdviceMessage<Object>(
AdviceMessage<?> message = new AdviceMessage<>(
MongoDbMessageStore.this.converter.extractPayload(source), headers, inputMessage);
enhanceHeaders(message.getHeaders(), headers);
@@ -731,21 +725,22 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore
}
private class DBObjectToErrorMessageConverter implements Converter<DBObject, ErrorMessage> {
@ReadingConverter
private class DocumentToErrorMessageConverter implements Converter<Document, ErrorMessage> {
private final Converter<byte[], Object> deserializingConverter = new DeserializingConverter();
DBObjectToErrorMessageConverter() {
DocumentToErrorMessageConverter() {
super();
}
@Override
public ErrorMessage convert(DBObject source) {
public ErrorMessage convert(Document source) {
@SuppressWarnings("unchecked")
Map<String, Object> headers =
MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
Object payload = this.deserializingConverter.convert((byte[]) source.get("payload"));
Object payload = this.deserializingConverter.convert(((Binary) source.get("payload")).getData());
ErrorMessage message = new ErrorMessage((Throwable) payload, headers);
enhanceHeaders(message.getHeaders(), headers);

View File

@@ -0,0 +1,39 @@
/*
* Copyright 2016 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.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.mongodb.support;
import org.bson.types.Binary;
import org.springframework.core.convert.converter.Converter;
import org.springframework.core.serializer.support.DeserializingConverter;
import org.springframework.data.convert.ReadingConverter;
import org.springframework.messaging.Message;
/**
* @author Artem Bilan
* @since 5.0
*/
@ReadingConverter
public class BinaryToMessageConverter implements Converter<Binary, Message<?>> {
private final Converter<byte[], Object> deserializingConverter = new DeserializingConverter();
@Override
public Message<?> convert(Binary source) {
return (Message<?>) this.deserializingConverter.convert(source.getData());
}
}

View File

@@ -0,0 +1,39 @@
/*
* Copyright 2016 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.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.springframework.integration.mongodb.support;
import org.bson.types.Binary;
import org.springframework.core.convert.converter.Converter;
import org.springframework.core.serializer.support.SerializingConverter;
import org.springframework.data.convert.WritingConverter;
import org.springframework.messaging.Message;
/**
* @author Artem Bilan
* @since 5.0
*/
@WritingConverter
public class MessageToBinaryConverter implements Converter<Message<?>, Binary> {
private final Converter<Object, byte[]> serializingConverter = new SerializingConverter();
@Override
public Binary convert(Message<?> source) {
return new Binary(this.serializingConverter.convert(source));
}
}

View File

@@ -19,11 +19,14 @@ package org.springframework.integration.mongodb.support;
import java.util.HashSet;
import java.util.Set;
import org.bson.types.Binary;
import org.springframework.core.convert.TypeDescriptor;
import org.springframework.core.convert.converter.Converter;
import org.springframework.core.convert.converter.GenericConverter;
import org.springframework.core.serializer.support.DeserializingConverter;
import org.springframework.core.serializer.support.SerializingConverter;
import org.springframework.data.convert.ReadingConverter;
import org.springframework.data.convert.WritingConverter;
import org.springframework.messaging.Message;
/**
@@ -33,7 +36,11 @@ import org.springframework.messaging.Message;
* @author Artem Bilan
* @since 4.2.10
* @deprecated since 5.0 in favor of {@link MessageToBinaryConverter} and {@link BinaryToMessageConverter}
*/
@WritingConverter
@ReadingConverter
@Deprecated
public class MongoDbMessageBytesConverter implements GenericConverter {
private final Converter<Object, byte[]> serializingConverter = new SerializingConverter();
@@ -42,19 +49,19 @@ public class MongoDbMessageBytesConverter implements GenericConverter {
@Override
public Set<ConvertiblePair> getConvertibleTypes() {
Set<ConvertiblePair> convertiblePairs = new HashSet<ConvertiblePair>();
convertiblePairs.add(new ConvertiblePair(Message.class, byte[].class));
convertiblePairs.add(new ConvertiblePair(byte[].class, Message.class));
Set<ConvertiblePair> convertiblePairs = new HashSet<>();
convertiblePairs.add(new ConvertiblePair(Message.class, Binary.class));
convertiblePairs.add(new ConvertiblePair(Binary.class, Message.class));
return convertiblePairs;
}
@Override
public Object convert(Object source, TypeDescriptor sourceType, TypeDescriptor targetType) {
if (Message.class.isAssignableFrom(sourceType.getObjectType())) {
return this.serializingConverter.convert(source);
return new Binary(this.serializingConverter.convert(source));
}
else {
return this.deserializingConverter.convert((byte[]) source);
return this.deserializingConverter.convert(((Binary) source).getData());
}
}

View File

@@ -78,9 +78,14 @@
</int-mongodb:inbound-channel-adapter>
<int:transaction-synchronization-factory id="syncFactory">
<int:after-commit expression="@documentCleaner.remove(#mongoTemplate, payload, headers.mongo_collectionName)" />
<int:before-commit expression="@documentCleaner.remove(#mongoTemplate, payload, headers.mongo_collectionName)"/>
<int:after-commit channel="afterCommitChannel"/>
</int:transaction-synchronization-factory>
<int:channel id="afterCommitChannel">
<int:queue />
</int:channel>
<int-mongodb:inbound-channel-adapter id="mongoInboundAdapterWithConverter"
channel="replyChannel"
query="{'name' : 'Bob'}"

View File

@@ -22,9 +22,9 @@ import static org.junit.Assert.assertNull;
import java.util.List;
import org.bson.Document;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.parsing.BeanDefinitionParsingException;
@@ -45,7 +45,6 @@ import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
import com.mongodb.DBObject;
import com.mongodb.util.JSON;
/**
@@ -64,6 +63,9 @@ public class MongoDbInboundChannelAdapterIntegrationTests extends MongoDbAvailab
@Autowired
private QueueChannel replyChannel;
@Autowired
private QueueChannel afterCommitChannel;
@Autowired
@Qualifier("mongoInboundAdapter")
private SourcePollingChannelAdapter mongoInboundAdapter;
@@ -117,7 +119,7 @@ public class MongoDbInboundChannelAdapterIntegrationTests extends MongoDbAvailab
this.mongoInboundAdapterNamedFactory.start();
@SuppressWarnings("unchecked")
Message<List<DBObject>> message = (Message<List<DBObject>>) replyChannel.receive(10000);
Message<List<Document>> message = (Message<List<Document>>) replyChannel.receive(10000);
assertNotNull(message);
assertEquals("Bob", message.getPayload().get(0).get("name"));
@@ -185,6 +187,8 @@ public class MongoDbInboundChannelAdapterIntegrationTests extends MongoDbAvailab
this.inboundAdapterWithOnSuccessDisposition.stop();
assertNotNull(this.afterCommitChannel.receive(10000));
assertNull(this.mongoTemplate.findOne(new Query(Criteria.where("name").is("Bob")), Person.class, "data"));
this.replyChannel.purge(null);
}

View File

@@ -25,9 +25,10 @@ import static org.mockito.Mockito.verify;
import java.util.List;
import org.bson.Document;
import org.bson.conversions.Bson;
import org.junit.Test;
import org.mockito.Mockito;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.data.mongodb.MongoDbFactory;
import org.springframework.data.mongodb.core.MongoOperations;
@@ -39,7 +40,6 @@ import org.springframework.expression.common.LiteralExpression;
import org.springframework.integration.mongodb.rules.MongoDbAvailable;
import org.springframework.integration.mongodb.rules.MongoDbAvailableTests;
import com.mongodb.DBObject;
import com.mongodb.util.JSON;
/**
@@ -75,8 +75,7 @@ public class MongoDbMessageSourceTests extends MongoDbAvailableTests {
@Test
@MongoDbAvailable
public void validateSuccessfullQueryWithSinigleElementIfOneInListAsDbObject() throws Exception {
public void validateSuccessfulQueryWithSingleElementIfOneInListAsDbObject() throws Exception {
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
MongoTemplate template = new MongoTemplate(mongoDbFactory);
@@ -87,16 +86,16 @@ public class MongoDbMessageSourceTests extends MongoDbAvailableTests {
messageSource.setBeanFactory(mock(BeanFactory.class));
messageSource.afterPropertiesSet();
@SuppressWarnings("unchecked")
List<DBObject> results = ((List<DBObject>) messageSource.receive().getPayload());
List<Document> results = ((List<Document>) messageSource.receive().getPayload());
assertEquals(1, results.size());
DBObject resultObject = results.get(0);
Document resultObject = results.get(0);
assertEquals("Oleg", resultObject.get("name"));
}
@Test
@MongoDbAvailable
public void validateSuccessfullQueryWithSinigleElementIfOneInList() throws Exception {
public void validateSuccessfulQueryWithSingleElementIfOneInList() throws Exception {
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
@@ -118,7 +117,7 @@ public class MongoDbMessageSourceTests extends MongoDbAvailableTests {
@Test
@MongoDbAvailable
public void validateSuccessfullQueryWithSinigleElementIfOneInListAndSingleResult() throws Exception {
public void validateSuccessfulQueryWithSingleElementIfOneInListAndSingleResult() throws Exception {
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
@@ -140,7 +139,7 @@ public class MongoDbMessageSourceTests extends MongoDbAvailableTests {
@Test
@MongoDbAvailable
public void validateSuccessfullSubObjectQueryWithSinigleElementIfOneInList() throws Exception {
public void validateSuccessfulSubObjectQueryWithSingleElementIfOneInList() throws Exception {
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
@@ -161,7 +160,7 @@ public class MongoDbMessageSourceTests extends MongoDbAvailableTests {
@Test
@MongoDbAvailable
public void validateSuccessfullQueryWithMultipleElements() throws Exception {
public void validateSuccessfulQueryWithMultipleElements() throws Exception {
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
@@ -181,7 +180,7 @@ public class MongoDbMessageSourceTests extends MongoDbAvailableTests {
@Test
@MongoDbAvailable
public void validateSuccessfullQueryWithNullReturn() throws Exception {
public void validateSuccessfulQueryWithNullReturn() throws Exception {
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
@@ -200,7 +199,7 @@ public class MongoDbMessageSourceTests extends MongoDbAvailableTests {
@SuppressWarnings("unchecked")
@Test
@MongoDbAvailable
public void validateSuccessfullQueryWithCustomConverter() throws Exception {
public void validateSuccessfulQueryWithCustomConverter() throws Exception {
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
@@ -220,13 +219,13 @@ public class MongoDbMessageSourceTests extends MongoDbAvailableTests {
List<Person> persons = (List<Person>) messageSource.receive().getPayload();
assertEquals(3, persons.size());
verify(converter, times(3)).read((Class<Person>) Mockito.any(), Mockito.any(DBObject.class));
verify(converter, times(3)).read((Class<Person>) Mockito.any(), Mockito.any(Bson.class));
}
@SuppressWarnings("unchecked")
@Test
@MongoDbAvailable
public void validateSuccessfullQueryWithMongoTemplate() throws Exception {
public void validateSuccessfulQueryWithMongoTemplate() throws Exception {
MongoDbFactory mongoDbFactory = this.prepareMongoFactory();
@@ -247,7 +246,7 @@ public class MongoDbMessageSourceTests extends MongoDbAvailableTests {
List<Person> persons = (List<Person>) messageSource.receive().getPayload();
assertEquals(3, persons.size());
verify(converter, times(3)).read((Class<Person>) Mockito.any(), Mockito.any(DBObject.class));
verify(converter, times(3)).read((Class<Person>) Mockito.any(), Mockito.any(Bson.class));
}
@Test
@@ -265,11 +264,12 @@ public class MongoDbMessageSourceTests extends MongoDbAvailableTests {
messageSource.setExpectSingleResult(true);
messageSource.setBeanFactory(mock(BeanFactory.class));
messageSource.afterPropertiesSet();
DBObject result = (DBObject) messageSource.receive().getPayload();
Document result = (Document) messageSource.receive().getPayload();
Object id = result.get("_id");
result.put("company", "PepBoys");
template.save(result, "data");
result = (DBObject) messageSource.receive().getPayload();
result = (Document) messageSource.receive().getPayload();
assertEquals(id, result.get("_id"));
}
}

View File

@@ -22,9 +22,9 @@ import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import org.bson.conversions.Bson;
import org.junit.Test;
import org.mockito.Mockito;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.data.mongodb.MongoDbFactory;
import org.springframework.data.mongodb.core.MongoOperations;
@@ -39,8 +39,6 @@ import org.springframework.integration.mongodb.rules.MongoDbAvailableTests;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.messaging.Message;
import com.mongodb.DBObject;
/**
* @author Amol Nayak
* @author Oleg Zhurakousky
@@ -122,7 +120,7 @@ public class MongoDbStoringMessageHandlerTests extends MongoDbAvailableTests {
assertEquals("Bob", person.getName());
assertEquals("PA", person.getAddress().getState());
verify(converter, times(1)).write(Mockito.any(), Mockito.any(DBObject.class));
verify(converter, times(1)).write(Mockito.any(), Mockito.any(Bson.class));
}
@Test
@@ -148,6 +146,6 @@ public class MongoDbStoringMessageHandlerTests extends MongoDbAvailableTests {
assertEquals("Bob", person.getName());
assertEquals("PA", person.getAddress().getState());
verify(converter, times(1)).write(Mockito.any(), Mockito.any(DBObject.class));
verify(converter, times(1)).write(Mockito.any(), Mockito.any(Bson.class));
}
}

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2013 the original author or authors.
* Copyright 2002-2016 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.
@@ -22,7 +22,6 @@ import org.junit.rules.MethodRule;
import org.junit.runners.model.FrameworkMethod;
import org.junit.runners.model.Statement;
import com.mongodb.Mongo;
import com.mongodb.MongoClient;
import com.mongodb.MongoClientOptions;
import com.mongodb.ServerAddress;
@@ -32,6 +31,8 @@ import com.mongodb.ServerAddress;
*
* @author Oleg Zhurakousky
* @author Gary Russell
* @author Artem Bilan
*
* @since 2.1
*/
public final class MongoDbAvailableRule implements MethodRule {
@@ -46,9 +47,12 @@ public final class MongoDbAvailableRule implements MethodRule {
MongoDbAvailable mongoAvailable = method.getAnnotation(MongoDbAvailable.class);
if (mongoAvailable != null) {
try {
MongoClientOptions options = new MongoClientOptions.Builder().connectTimeout(100).build();
Mongo mongo = new MongoClient(ServerAddress.defaultHost(), options);
mongo.getDatabaseNames();
MongoClientOptions options = new MongoClientOptions.Builder()
.serverSelectionTimeout(100)
.build();
MongoClient mongo = new MongoClient(ServerAddress.defaultHost(), options);
mongo.listDatabaseNames()
.first();
}
catch (Exception e) {
logger.warn("MongoDb is not available. Skipping the test: " +

View File

@@ -16,10 +16,8 @@
package org.springframework.integration.mongodb.rules;
import com.mongodb.DBObject;
import com.mongodb.MongoClient;
import org.bson.conversions.Bson;
import org.junit.Rule;
import org.springframework.data.mapping.context.MappingContext;
import org.springframework.data.mongodb.MongoDbFactory;
import org.springframework.data.mongodb.core.MongoTemplate;
@@ -29,6 +27,8 @@ import org.springframework.data.mongodb.core.convert.MappingMongoConverter;
import org.springframework.data.mongodb.core.mapping.MongoPersistentEntity;
import org.springframework.data.mongodb.core.mapping.MongoPersistentProperty;
import com.mongodb.MongoClient;
/**
* Convenience base class that enables unit test methods to rely upon the {@link MongoDbAvailable} annotation.
*
@@ -148,12 +148,12 @@ public abstract class MongoDbAvailableTests {
}
@Override
public void write(Object source, DBObject target) {
public void write(Object source, Bson target) {
super.write(source, target);
}
@Override
public <S> S read(Class<S> clazz, DBObject source) {
public <S> S read(Class<S> clazz, Bson source) {
return super.read(clazz, source);
}

View File

@@ -378,9 +378,9 @@ public abstract class AbstractMongoDbMessageGroupStoreTests extends MongoDbAvail
MessageGroupStore store2 = this.getMessageGroupStore();
Message<?> message = new GenericMessage<String>("1");
store2.addMessagesToGroup(1, message);
store1.addMessagesToGroup(2, new GenericMessage<String>("2"));
store2.addMessagesToGroup(3, new GenericMessage<String>("3"));
store2.addMessagesToGroup("1", message);
store1.addMessagesToGroup("2", new GenericMessage<String>("2"));
store2.addMessagesToGroup("3", new GenericMessage<String>("3"));
MessageGroupStore store3 = this.getMessageGroupStore();
Iterator<MessageGroup> iterator = store3.iterator();
@@ -392,7 +392,7 @@ public abstract class AbstractMongoDbMessageGroupStoreTests extends MongoDbAvail
}
assertEquals(3, counter);
store2.removeMessagesFromGroup(1, message);
store2.removeMessagesFromGroup("1", message);
iterator = store3.iterator();
counter = 0;

View File

@@ -22,13 +22,14 @@ import static org.junit.Assert.assertThat;
import java.util.Map;
import org.bson.Document;
import org.hamcrest.Matchers;
import org.junit.Ignore;
import org.junit.Test;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.context.support.GenericApplicationContext;
import org.springframework.core.convert.converter.Converter;
import org.springframework.data.convert.ReadingConverter;
import org.springframework.data.mongodb.MongoDbFactory;
import org.springframework.data.mongodb.core.SimpleMongoDbFactory;
import org.springframework.integration.IntegrationMessageHeaderAccessor;
@@ -43,7 +44,6 @@ import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.util.StopWatch;
import com.mongodb.DBObject;
import com.mongodb.MongoClient;
/**
@@ -198,12 +198,15 @@ public class ConfigurableMongoDbMessageGroupStoreTests extends AbstractMongoDbMe
}
public static class MessageReadConverter implements Converter<DBObject, Message<?>> {
@ReadingConverter
public static class MessageReadConverter implements Converter<Document, Message<?>> {
@Override
@SuppressWarnings("unchecked")
public Message<?> convert(DBObject source) {
return MessageBuilder.withPayload(source.get("payload")).copyHeaders((Map<String, ?>) source.get("headers")).build();
public Message<?> convert(Document source) {
return MessageBuilder.withPayload(source.get("payload"))
.copyHeaders((Map<String, ?>) source.get("headers"))
.build();
}
}

View File

@@ -11,7 +11,7 @@
<beans:bean id="mongoConnectionFactory" class="org.springframework.data.mongodb.core.SimpleMongoDbFactory">
<beans:constructor-arg>
<beans:bean class="com.mongodb.Mongo"/>
<beans:bean class="com.mongodb.MongoClient"/>
</beans:constructor-arg>
<beans:constructor-arg value="test"/>
</beans:bean>

View File

@@ -11,7 +11,7 @@
<beans:bean id="mongoConnectionFactory" class="org.springframework.data.mongodb.core.SimpleMongoDbFactory">
<beans:constructor-arg>
<beans:bean class="com.mongodb.Mongo"/>
<beans:bean class="com.mongodb.MongoClient"/>
</beans:constructor-arg>
<beans:constructor-arg value="test"/>
</beans:bean>

View File

@@ -21,7 +21,6 @@ import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNotSame;
import static org.junit.Assert.assertThat;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
import java.util.Iterator;
import java.util.concurrent.TimeUnit;
@@ -29,7 +28,6 @@ import java.util.concurrent.TimeUnit;
import org.hamcrest.Matchers;
import org.junit.Rule;
import org.junit.Test;
import org.springframework.context.support.AbstractApplicationContext;
import org.springframework.context.support.ClassPathXmlApplicationContext;
import org.springframework.integration.context.IntegrationContextUtils;
@@ -88,16 +86,6 @@ public class DelayerHandlerRescheduleIntegrationTests extends MongoDbAvailableTe
(ThreadPoolTaskScheduler) IntegrationContextUtils.getTaskScheduler(context);
taskScheduler.shutdown();
taskScheduler.getScheduledExecutor().awaitTermination(10, TimeUnit.SECONDS);
context.destroy();
try {
context.getBean("input", MessageChannel.class);
fail("IllegalStateException expected");
}
catch (Exception e) {
assertTrue(e instanceof IllegalStateException);
assertTrue(e.getMessage().contains("BeanFactory not initialized or already closed - call 'refresh'"));
}
assertEquals(2, messageStore.messageGroupSize(delayerMessageGroupId));
@@ -115,6 +103,8 @@ public class DelayerHandlerRescheduleIntegrationTests extends MongoDbAvailableTe
.getOriginal();
assertThat(message1, Matchers.anyOf(Matchers.is(original1), Matchers.is(original2)));
context.destroy();
context.refresh();
PollableChannel output = context.getBean("output", PollableChannel.class);
@@ -129,6 +119,8 @@ public class DelayerHandlerRescheduleIntegrationTests extends MongoDbAvailableTe
Object payload2 = message.getPayload();
assertNotSame(payload1, payload2);
messageStore = context.getBean("messageStore", MessageGroupStore.class);
assertEquals(0, messageStore.messageGroupSize(delayerMessageGroupId));
context.destroy();
}

View File

@@ -22,7 +22,7 @@
<bean id="mongoConnectionFactory" class="org.springframework.data.mongodb.core.SimpleMongoDbFactory">
<constructor-arg>
<bean class="com.mongodb.Mongo"/>
<bean class="com.mongodb.MongoClient"/>
</constructor-arg>
<constructor-arg value="test"/>
</bean>

View File

@@ -22,7 +22,7 @@
<bean id="mongoConnectionFactory" class="org.springframework.data.mongodb.core.SimpleMongoDbFactory">
<constructor-arg>
<bean class="com.mongodb.Mongo"/>
<bean class="com.mongodb.MongoClient"/>
</constructor-arg>
<constructor-arg value="test"/>
</bean>