From a62a572a6cf51a940c1d16867fc81148844dd0f0 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 22 Jan 2020 12:35:20 -0500 Subject: [PATCH] Add ReactiveMongoDbMessageSource polling test * Insert test data before running test * close an `ApplicationContext` after test --- .../ReactiveMongoDbMessageSourceTests.java | 83 ++++++++++++++++--- 1 file changed, 73 insertions(+), 10 deletions(-) diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/ReactiveMongoDbMessageSourceTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/ReactiveMongoDbMessageSourceTests.java index 49564fdc5c..65a26545e8 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/ReactiveMongoDbMessageSourceTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/inbound/ReactiveMongoDbMessageSourceTests.java @@ -32,14 +32,25 @@ import java.util.Optional; import org.bson.conversions.Bson; import org.junit.Test; import org.mockito.Mockito; +import org.reactivestreams.Publisher; import org.springframework.beans.factory.BeanFactory; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.annotation.AnnotationConfigApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; import org.springframework.data.mongodb.ReactiveMongoDatabaseFactory; import org.springframework.data.mongodb.core.ReactiveMongoTemplate; import org.springframework.data.mongodb.core.convert.MappingMongoConverter; import org.springframework.data.mongodb.core.mapping.MongoMappingContext; import org.springframework.expression.Expression; import org.springframework.expression.common.LiteralExpression; +import org.springframework.integration.channel.FluxMessageChannel; +import org.springframework.integration.config.EnableIntegration; +import org.springframework.integration.core.MessageSource; +import org.springframework.integration.dsl.IntegrationFlow; +import org.springframework.integration.dsl.IntegrationFlows; +import org.springframework.integration.dsl.Pollers; import org.springframework.integration.mongodb.rules.MongoDbAvailable; import org.springframework.integration.mongodb.rules.MongoDbAvailableTests; @@ -50,6 +61,7 @@ import reactor.test.StepVerifier; /** * @author David Turanski + * @author Artem Bilan * * @since 5.3 */ @@ -83,7 +95,7 @@ public class ReactiveMongoDbMessageSourceTests extends MongoDbAvailableTests { ReactiveMongoDatabaseFactory reactiveMongoDatabaseFactory = this.prepareReactiveMongoFactory(); ReactiveMongoTemplate template = new ReactiveMongoTemplate(reactiveMongoDatabaseFactory); - waitFor(template.save(this.createPerson(), "data")); + waitFor(template.save(createPerson(), "data")); Expression queryExpression = new LiteralExpression("{'name' : 'Oleg'}"); ReactiveMongoDbMessageSource messageSource = new ReactiveMongoDbMessageSource(reactiveMongoDatabaseFactory, @@ -103,7 +115,7 @@ public class ReactiveMongoDbMessageSourceTests extends MongoDbAvailableTests { ReactiveMongoDatabaseFactory reactiveMongoDatabaseFactory = this.prepareReactiveMongoFactory(); ReactiveMongoTemplate template = new ReactiveMongoTemplate(reactiveMongoDatabaseFactory); - waitFor(template.save(this.createPerson(), "data")); + waitFor(template.save(createPerson(), "data")); Expression queryExpression = new LiteralExpression("{'name' : 'Oleg'}"); ReactiveMongoDbMessageSource messageSource = new ReactiveMongoDbMessageSource(reactiveMongoDatabaseFactory, @@ -148,7 +160,7 @@ public class ReactiveMongoDbMessageSourceTests extends MongoDbAvailableTests { @MongoDbAvailable @SuppressWarnings("unchecked") public void validateSuccessfulQueryWithCustomConverter() { - MappingMongoConverter converter = new ReactiveTestMongoConverter(this.prepareReactiveMongoFactory(), + MappingMongoConverter converter = new ReactiveTestMongoConverter(prepareReactiveMongoFactory(), new MongoMappingContext()); converter.afterPropertiesSet(); converter = spy(converter); @@ -159,13 +171,30 @@ public class ReactiveMongoDbMessageSourceTests extends MongoDbAvailableTests { verify(converter, times(3)).read((Class) Mockito.any(), Mockito.any(Bson.class)); } + @Test + @MongoDbAvailable + public void validateWithConfiguredPollerFlow() { + ReactiveMongoDatabaseFactory reactiveMongoDatabaseFactory = prepareReactiveMongoFactory(); + ReactiveMongoTemplate template = new ReactiveMongoTemplate(reactiveMongoDatabaseFactory); + + waitFor(template.save(createPerson(), "data")); + + ConfigurableApplicationContext context = new AnnotationConfigApplicationContext(TestContext.class); + FluxMessageChannel output = context.getBean(FluxMessageChannel.class); + StepVerifier.create(output) + .assertNext( + message -> assertThat(((Person) message.getPayload()).getName()).isEqualTo("Oleg")) + .thenCancel() + .verify(); + + context.close(); + } + @Test @MongoDbAvailable @SuppressWarnings("unchecked") public void validatePipelineInModifyOut() { - - ReactiveMongoDatabaseFactory reactiveMongoDatabaseFactory = this.prepareReactiveMongoFactory(); - + ReactiveMongoDatabaseFactory reactiveMongoDatabaseFactory = prepareReactiveMongoFactory(); ReactiveMongoTemplate template = new ReactiveMongoTemplate(reactiveMongoDatabaseFactory); waitFor(template.save(BasicDBObject.parse("{'name' : 'Manny', 'id' : 1}"), "data")); @@ -185,7 +214,7 @@ public class ReactiveMongoDbMessageSourceTests extends MongoDbAvailableTests { } private Flux queryMultipleElements(Expression queryExpression) { - return this.queryMultipleElements(queryExpression, Optional.empty()); + return queryMultipleElements(queryExpression, Optional.empty()); } @SuppressWarnings("unchecked") @@ -193,9 +222,9 @@ public class ReactiveMongoDbMessageSourceTests extends MongoDbAvailableTests { ReactiveMongoDatabaseFactory reactiveMongoDatabaseFactory = this.prepareReactiveMongoFactory(); ReactiveMongoTemplate template = new ReactiveMongoTemplate(reactiveMongoDatabaseFactory); - waitFor(template.save(this.createPerson("Manny"), "data")); - waitFor(template.save(this.createPerson("Moe"), "data")); - waitFor(template.save(this.createPerson("Jack"), "data")); + waitFor(template.save(createPerson("Manny"), "data")); + waitFor(template.save(createPerson("Moe"), "data")); + waitFor(template.save(createPerson("Jack"), "data")); ReactiveMongoDbMessageSource messageSource = new ReactiveMongoDbMessageSource(reactiveMongoDatabaseFactory, queryExpression); @@ -211,4 +240,38 @@ public class ReactiveMongoDbMessageSourceTests extends MongoDbAvailableTests { return mono.block(Duration.ofSeconds(10)); } + @Configuration + @EnableIntegration + static class TestContext { + + @Bean + FluxMessageChannel output() { + return new FluxMessageChannel(); + } + + @Bean + public MessageSource> mongodbMessageSource(ReactiveMongoDatabaseFactory mongoDatabaseFactory) { + Expression queryExpression = new LiteralExpression("{'name' : 'Oleg'}"); + ReactiveMongoDbMessageSource reactiveMongoDbMessageSource = + new ReactiveMongoDbMessageSource(mongoDatabaseFactory, queryExpression); + reactiveMongoDbMessageSource.setEntityClass(Person.class); + return reactiveMongoDbMessageSource; + } + + @Bean + ReactiveMongoDatabaseFactory mongoDatabaseFactory() { + return MongoDbAvailableTests.REACTIVE_MONGO_DATABASE_FACTORY; + } + + @Bean + public IntegrationFlow pollingFlow(MessageSource> mongodbMessageSource) { + return IntegrationFlows + .from(mongodbMessageSource, c -> c.poller(Pollers.fixedDelay(100).maxMessagesPerPoll(1))) + .split() + .channel(output()) + .get(); + } + + } + }