From c0a507c36cbe072e72b8c460fa0c9e50e7737761 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 30 Nov 2016 07:32:51 -0500 Subject: [PATCH] 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 --- .travis.yml | 6 + build.gradle | 14 +- .../handler/LambdaMessageProcessor.java | 2 +- .../expression/ExpressionUtilsTests.java | 16 +- ...tOnceFileListFilterExternalStoreTests.java | 2 +- .../CacheListeningMessageProducer.java | 10 +- .../ContinuousQueryMessageProducer.java | 4 +- .../metadata/GemfireMetadataStore.java | 6 +- .../outbound/CacheWritingMessageHandler.java | 6 +- .../gemfire/store/GemfireMessageStore.java | 4 +- .../gemfire/util/GemfireLockRegistry.java | 6 +- .../gemfire/TestCacheListenerLogger.java | 4 +- ...boundChannelAdapterParserTests-context.xml | 2 +- ...boundChannelAdapterParserTests-context.xml | 2 +- .../gemfire/fork/CacheServerProcess.java | 12 +- .../CacheListeningMessageProducerTests.java | 2 +- .../ContinuousQueryMessageProducerTests.java | 12 +- .../inbound/CqInboundChannelAdapterTests.java | 4 +- .../GemfireInboundChannelAdapterTests.java | 4 +- .../metadata/GemfireMetadataStoreTests.java | 6 +- .../CacheWritingMessageHandlerTests.java | 6 +- .../GemfireOutboundChannelAdapterTests.java | 2 +- ...ayerHandlerRescheduleIntegrationTests.java | 6 +- .../gemfire/store/GemfireGroupStoreTests.java | 6 +- .../store/GemfireMessageStoreTests.java | 6 +- .../metadata/MongoDbMetadataStore.java | 8 +- ...stractConfigurableMongoDbMessageStore.java | 7 +- .../ConfigurableMongoDbMessageStore.java | 8 +- .../mongodb/store/MongoDbMessageStore.java | 213 +++++++++--------- .../support/BinaryToMessageConverter.java | 39 ++++ .../support/MessageToBinaryConverter.java | 39 ++++ .../support/MongoDbMessageBytesConverter.java | 17 +- ...ChannelAdapterIntegrationTests-context.xml | 7 +- ...InboundChannelAdapterIntegrationTests.java | 10 +- .../inbound/MongoDbMessageSourceTests.java | 34 +-- .../MongoDbStoringMessageHandlerTests.java | 8 +- .../mongodb/rules/MongoDbAvailableRule.java | 14 +- .../mongodb/rules/MongoDbAvailableTests.java | 10 +- ...AbstractMongoDbMessageGroupStoreTests.java | 8 +- ...igurableMongoDbMessageGroupStoreTests.java | 13 +- ...leIntegrationConfigurableTests-context.xml | 2 +- ...dlerRescheduleIntegrationTests-context.xml | 2 +- ...ayerHandlerRescheduleIntegrationTests.java | 16 +- .../mongodb/store/mongo-aggregator-config.xml | 2 +- .../mongo-aggregator-configurable-config.xml | 2 +- 45 files changed, 346 insertions(+), 263 deletions(-) create mode 100644 spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/BinaryToMessageConverter.java create mode 100644 spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MessageToBinaryConverter.java diff --git a/.travis.yml b/.travis.yml index 815c5397e5..93a98e3b40 100644 --- a/.travis.yml +++ b/.travis.yml @@ -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: diff --git a/build.gradle b/build.gradle index 38f992bf95..24e5d1c5e5 100644 --- a/build.gradle +++ b/build.gradle @@ -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' } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/LambdaMessageProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/LambdaMessageProcessor.java index dd94340559..66333220c5 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/LambdaMessageProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/LambdaMessageProcessor.java @@ -86,7 +86,7 @@ public class LambdaMessageProcessor implements MessageProcessor, 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; } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/expression/ExpressionUtilsTests.java b/spring-integration-core/src/test/java/org/springframework/integration/expression/ExpressionUtilsTests.java index 1bca90c83b..405cb47a9a 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/expression/ExpressionUtilsTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/expression/ExpressionUtilsTests.java @@ -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")); } } diff --git a/spring-integration-file/src/test/java/org/springframework/integration/file/filters/PersistentAcceptOnceFileListFilterExternalStoreTests.java b/spring-integration-file/src/test/java/org/springframework/integration/file/filters/PersistentAcceptOnceFileListFilterExternalStoreTests.java index daab0a8104..0bf7b65135 100644 --- a/spring-integration-file/src/test/java/org/springframework/integration/file/filters/PersistentAcceptOnceFileListFilterExternalStoreTests.java +++ b/spring-integration-file/src/test/java/org/springframework/integration/file/filters/PersistentAcceptOnceFileListFilterExternalStoreTests.java @@ -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 diff --git a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducer.java b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducer.java index 259fc25375..39885dbe9f 100644 --- a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducer.java +++ b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducer.java @@ -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 diff --git a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducer.java b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducer.java index bb31062d24..821df4405e 100644 --- a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducer.java +++ b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducer.java @@ -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 diff --git a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/metadata/GemfireMetadataStore.java b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/metadata/GemfireMetadataStore.java index 2894d9ee1b..913a22cc18 100644 --- a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/metadata/GemfireMetadataStore.java +++ b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/metadata/GemfireMetadataStore.java @@ -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}. diff --git a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/outbound/CacheWritingMessageHandler.java b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/outbound/CacheWritingMessageHandler.java index 9e164f40ac..64988b0780 100644 --- a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/outbound/CacheWritingMessageHandler.java +++ b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/outbound/CacheWritingMessageHandler.java @@ -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 diff --git a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/store/GemfireMessageStore.java b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/store/GemfireMessageStore.java index 70f4fa2323..bd5cbc8081 100644 --- a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/store/GemfireMessageStore.java +++ b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/store/GemfireMessageStore.java @@ -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 diff --git a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/util/GemfireLockRegistry.java b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/util/GemfireLockRegistry.java index a3605e9e2d..66e770450e 100644 --- a/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/util/GemfireLockRegistry.java +++ b/spring-integration-gemfire/src/main/java/org/springframework/integration/gemfire/util/GemfireLockRegistry.java @@ -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. diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/TestCacheListenerLogger.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/TestCacheListenerLogger.java index f3f7afd700..9d2690d797 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/TestCacheListenerLogger.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/TestCacheListenerLogger.java @@ -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; diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/config/xml/GemfireInboundChannelAdapterParserTests-context.xml b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/config/xml/GemfireInboundChannelAdapterParserTests-context.xml index 439344d27c..46288c8ac5 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/config/xml/GemfireInboundChannelAdapterParserTests-context.xml +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/config/xml/GemfireInboundChannelAdapterParserTests-context.xml @@ -8,7 +8,7 @@ http://www.springframework.org/schema/beans/spring-beans.xsd"> - + - + diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/fork/CacheServerProcess.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/fork/CacheServerProcess.java index a869146638..cb2415cd3e 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/fork/CacheServerProcess.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/fork/CacheServerProcess.java @@ -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 diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducerTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducerTests.java index 5062165e85..7cbe00af22 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducerTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/CacheListeningMessageProducerTests.java @@ -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 diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducerTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducerTests.java index afd9e64bc0..2bad709124 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducerTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/ContinuousQueryMessageProducerTests.java @@ -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]; diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/CqInboundChannelAdapterTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/CqInboundChannelAdapterTests.java index 295827a0b3..bf4aaf0429 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/CqInboundChannelAdapterTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/CqInboundChannelAdapterTests.java @@ -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 diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/GemfireInboundChannelAdapterTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/GemfireInboundChannelAdapterTests.java index 0993e3023e..92fd8ff650 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/GemfireInboundChannelAdapterTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/inbound/GemfireInboundChannelAdapterTests.java @@ -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 diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/metadata/GemfireMetadataStoreTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/metadata/GemfireMetadataStoreTests.java index 5f325ba9d8..b9568c1f67 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/metadata/GemfireMetadataStoreTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/metadata/GemfireMetadataStoreTests.java @@ -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 diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/outbound/CacheWritingMessageHandlerTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/outbound/CacheWritingMessageHandlerTests.java index a372e5a56c..e1a8ebf2de 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/outbound/CacheWritingMessageHandlerTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/outbound/CacheWritingMessageHandlerTests.java @@ -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 diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/outbound/GemfireOutboundChannelAdapterTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/outbound/GemfireOutboundChannelAdapterTests.java index 55eaa667d2..7cd92bfff8 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/outbound/GemfireOutboundChannelAdapterTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/outbound/GemfireOutboundChannelAdapterTests.java @@ -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 diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/DelayerHandlerRescheduleIntegrationTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/DelayerHandlerRescheduleIntegrationTests.java index 56e045d8d6..19fa27b91f 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/DelayerHandlerRescheduleIntegrationTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/DelayerHandlerRescheduleIntegrationTests.java @@ -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; /** diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireGroupStoreTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireGroupStoreTests.java index 4caba29de3..83716d258b 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireGroupStoreTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireGroupStoreTests.java @@ -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; diff --git a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireMessageStoreTests.java b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireMessageStoreTests.java index f3460d8f0a..1a3e3fcc3f 100644 --- a/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireMessageStoreTests.java +++ b/spring-integration-gemfire/src/test/java/org/springframework/integration/gemfire/store/GemfireMessageStoreTests.java @@ -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 diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/metadata/MongoDbMetadataStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/metadata/MongoDbMetadataStore.java index 1ba0107a8b..1165a24793 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/metadata/MongoDbMetadataStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/metadata/MongoDbMetadataStore.java @@ -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 entry = new HashMap<>(); + final Map 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; } } diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java index ce45cce81f..ffd6afc621 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java @@ -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 customConverters = new ArrayList(); - customConverters.add(new MongoDbMessageBytesConverter()); + customConverters.add(new MessageToBinaryConverter()); + customConverters.add(new BinaryToMessageConverter()); this.mappingMongoConverter.setCustomConversions(new CustomConversions(customConverters)); this.mappingMongoConverter.afterPropertiesSet(); } diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java index 58ce4adf15..7aadd3224f 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java @@ -259,9 +259,8 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb List messageGroups = new ArrayList(); 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 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(); } diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java index 03f21d382d..3db2eba8e5 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java @@ -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 ids = new ArrayList(); + Collection 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 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 iterator() { - List messageGroups = new ArrayList(); + List 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 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 messageWrappers = this.template.find(query, MessageWrapper.class, this.collectionName); - List> messages = new ArrayList>(); - 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 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 customConverters = new ArrayList(); - 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 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 read(Class clazz, DBObject source) { + @SuppressWarnings({"unchecked"}) + public S read(Class clazz, Bson source) { if (!MessageWrapper.class.equals(clazz)) { return super.read(clazz, source); } if (source != null) { + Map 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 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 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 { - UuidToDBObjectConverter() { + @WritingConverter + private static class MessageHistoryToDocumentConverter implements Converter { + + 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 { - - DBObjectToUUIDConverter() { - super(); - } - - @Override - public UUID convert(DBObject source) { - return UUID.fromString((String) source.get("_value")); - } - } - - - private static class MessageHistoryToDBObjectConverter implements Converter { - - 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> { + @ReadingConverter + private class DocumentToGenericMessageConverter implements Converter> { - DBObjectToGenericMessageConverter() { + DocumentToGenericMessageConverter() { super(); } @Override - public GenericMessage convert(DBObject source) { + public GenericMessage convert(Document source) { @SuppressWarnings("unchecked") Map headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map) source.get("headers")); GenericMessage message = - new GenericMessage(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> { + @ReadingConverter + private final class DocumentToMutableMessageConverter implements Converter> { - - DBObjectToMutableMessageConverter() { + DocumentToMutableMessageConverter() { super(); } @Override - public MutableMessage convert(DBObject source) { + public MutableMessage convert(Document source) { @SuppressWarnings("unchecked") Map headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map) source.get("headers")); @@ -694,14 +687,15 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore } - private class DBObjectToAdviceMessageConverter implements Converter> { + @ReadingConverter + private class DocumentToAdviceMessageConverter implements Converter> { - DBObjectToAdviceMessageConverter() { + DocumentToAdviceMessageConverter() { super(); } @Override - public AdviceMessage convert(DBObject source) { + public AdviceMessage convert(Document source) { @SuppressWarnings("unchecked") Map headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map) 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( + 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 { + @ReadingConverter + private class DocumentToErrorMessageConverter implements Converter { private final Converter deserializingConverter = new DeserializingConverter(); - DBObjectToErrorMessageConverter() { + DocumentToErrorMessageConverter() { super(); } @Override - public ErrorMessage convert(DBObject source) { + public ErrorMessage convert(Document source) { @SuppressWarnings("unchecked") Map headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map) 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); diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/BinaryToMessageConverter.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/BinaryToMessageConverter.java new file mode 100644 index 0000000000..a2a169ffb3 --- /dev/null +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/BinaryToMessageConverter.java @@ -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> { + + private final Converter deserializingConverter = new DeserializingConverter(); + + @Override + public Message convert(Binary source) { + return (Message) this.deserializingConverter.convert(source.getData()); + } + +} diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MessageToBinaryConverter.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MessageToBinaryConverter.java new file mode 100644 index 0000000000..2916bcbbda --- /dev/null +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MessageToBinaryConverter.java @@ -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, Binary> { + + private final Converter serializingConverter = new SerializingConverter(); + + @Override + public Binary convert(Message source) { + return new Binary(this.serializingConverter.convert(source)); + } + +} diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MongoDbMessageBytesConverter.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MongoDbMessageBytesConverter.java index 6b8395b3f2..4afcd8f541 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MongoDbMessageBytesConverter.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/support/MongoDbMessageBytesConverter.java @@ -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 serializingConverter = new SerializingConverter(); @@ -42,19 +49,19 @@ public class MongoDbMessageBytesConverter implements GenericConverter { @Override public Set getConvertibleTypes() { - Set convertiblePairs = new HashSet(); - convertiblePairs.add(new ConvertiblePair(Message.class, byte[].class)); - convertiblePairs.add(new ConvertiblePair(byte[].class, Message.class)); + Set 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()); } } diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbInboundChannelAdapterIntegrationTests-context.xml b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbInboundChannelAdapterIntegrationTests-context.xml index 434178bccc..71d225204e 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbInboundChannelAdapterIntegrationTests-context.xml +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbInboundChannelAdapterIntegrationTests-context.xml @@ -78,9 +78,14 @@ - + + + + + + > message = (Message>) replyChannel.receive(10000); + Message> message = (Message>) 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); } diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/MongoDbMessageSourceTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/MongoDbMessageSourceTests.java index 615b9a3fe4..ea2fae34ed 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/MongoDbMessageSourceTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/MongoDbMessageSourceTests.java @@ -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 results = ((List) messageSource.receive().getPayload()); + List results = ((List) 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 persons = (List) messageSource.receive().getPayload(); assertEquals(3, persons.size()); - verify(converter, times(3)).read((Class) Mockito.any(), Mockito.any(DBObject.class)); + verify(converter, times(3)).read((Class) 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 persons = (List) messageSource.receive().getPayload(); assertEquals(3, persons.size()); - verify(converter, times(3)).read((Class) Mockito.any(), Mockito.any(DBObject.class)); + verify(converter, times(3)).read((Class) 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")); } + } diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/outbound/MongoDbStoringMessageHandlerTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/outbound/MongoDbStoringMessageHandlerTests.java index 58b7a74c24..4c57581aa3 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/outbound/MongoDbStoringMessageHandlerTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/outbound/MongoDbStoringMessageHandlerTests.java @@ -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)); } } diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/rules/MongoDbAvailableRule.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/rules/MongoDbAvailableRule.java index fed15f3e79..bd352c8c76 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/rules/MongoDbAvailableRule.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/rules/MongoDbAvailableRule.java @@ -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: " + diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/rules/MongoDbAvailableTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/rules/MongoDbAvailableTests.java index 496d706344..079cb18ff3 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/rules/MongoDbAvailableTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/rules/MongoDbAvailableTests.java @@ -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 read(Class clazz, DBObject source) { + public S read(Class clazz, Bson source) { return super.read(clazz, source); } diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/AbstractMongoDbMessageGroupStoreTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/AbstractMongoDbMessageGroupStoreTests.java index b1412f3522..99fbf9bb5f 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/AbstractMongoDbMessageGroupStoreTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/AbstractMongoDbMessageGroupStoreTests.java @@ -378,9 +378,9 @@ public abstract class AbstractMongoDbMessageGroupStoreTests extends MongoDbAvail MessageGroupStore store2 = this.getMessageGroupStore(); Message message = new GenericMessage("1"); - store2.addMessagesToGroup(1, message); - store1.addMessagesToGroup(2, new GenericMessage("2")); - store2.addMessagesToGroup(3, new GenericMessage("3")); + store2.addMessagesToGroup("1", message); + store1.addMessagesToGroup("2", new GenericMessage("2")); + store2.addMessagesToGroup("3", new GenericMessage("3")); MessageGroupStore store3 = this.getMessageGroupStore(); Iterator 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; diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageGroupStoreTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageGroupStoreTests.java index 24a81c6ec6..b27877076a 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageGroupStoreTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageGroupStoreTests.java @@ -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> { + @ReadingConverter + public static class MessageReadConverter implements Converter> { @Override @SuppressWarnings("unchecked") - public Message convert(DBObject source) { - return MessageBuilder.withPayload(source.get("payload")).copyHeaders((Map) source.get("headers")).build(); + public Message convert(Document source) { + return MessageBuilder.withPayload(source.get("payload")) + .copyHeaders((Map) source.get("headers")) + .build(); } } diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationConfigurableTests-context.xml b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationConfigurableTests-context.xml index 87cfc91b9d..5fb2c546f5 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationConfigurableTests-context.xml +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationConfigurableTests-context.xml @@ -11,7 +11,7 @@ - + diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationTests-context.xml b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationTests-context.xml index 074ca447b1..9746fe59f7 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationTests-context.xml +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationTests-context.xml @@ -11,7 +11,7 @@ - + diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationTests.java index 579033c72e..14099ac7b1 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/DelayerHandlerRescheduleIntegrationTests.java @@ -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(); } diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-config.xml b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-config.xml index cb3231ae9c..821b099b0f 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-config.xml +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-config.xml @@ -22,7 +22,7 @@ - + diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-configurable-config.xml b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-configurable-config.xml index 7aa0279ccc..dfdb01316c 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-configurable-config.xml +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/store/mongo-aggregator-configurable-config.xml @@ -22,7 +22,7 @@ - +