Add ReactiveMongoDbMessageSource polling test
* Insert test data before running test * close an `ApplicationContext` after test
This commit is contained in:
@@ -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<Person>) 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<Person> 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<Publisher<?>> 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<Publisher<?>> mongodbMessageSource) {
|
||||
return IntegrationFlows
|
||||
.from(mongodbMessageSource, c -> c.poller(Pollers.fixedDelay(100).maxMessagesPerPoll(1)))
|
||||
.split()
|
||||
.channel(output())
|
||||
.get();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user