From 46de69dbf2c11186ab250eee11f808b3e0581c15 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 1 Sep 2017 16:26:50 -0400 Subject: [PATCH] Honor the FluxSink State in the PollChPublisherAd The `PollableChannelPublisherAdapter` is based on the poll model of the `FluxSink` and iterate and poll downstream `PollableChannel` until there is an item or `n > 0`. Having `take(6)` we end up with the cancel from the downstream `Subscriber` after the batch is filled, but at the same time we continue to poll the upstream source because `n` is like `Long.MAX_VALUE`. * The proper way to interact is check for the `!sink.isCancelled()` as well. This way cancelled `sink` won't "steal" data from other subscribers Fix compatibility with the latest dependencies * Upgrade to Gradle 4.1 * Upgrade as much dependencies as possible * Fix MongoDB module to resolve deprecation in the latest driver * Increase receive timeout in the `ResequencerTests` * Restore generic argument for the method reference in the `ReactiveStreamsTests` * Fix `WebFluxInboundEndpoint` for the compatibility with the latest Reactor API --- build.gradle | 40 +++++++------- gradle/wrapper/gradle-wrapper.jar | Bin 54712 -> 54712 bytes gradle/wrapper/gradle-wrapper.properties | 4 +- .../channel/MessageChannelReactiveUtils.java | 2 +- .../aggregator/ResequencerTests.java | 7 +-- .../reactivestreams/ReactiveStreamsTests.java | 2 - ...InboundChannelAdapterIntegrationTests.java | 10 ++-- ...utboundChannelAdapterIntegrationTests.java | 51 +++++++++++++----- .../inbound/MongoDbMessageSourceTests.java | 5 +- .../inbound/WebFluxInboundEndpoint.java | 2 +- 10 files changed, 74 insertions(+), 49 deletions(-) diff --git a/build.gradle b/build.gradle index 91bab61a22..aeb95fb5b7 100644 --- a/build.gradle +++ b/build.gradle @@ -3,7 +3,7 @@ buildscript { maven { url 'https://repo.spring.io/plugins-release' } } dependencies { - classpath 'io.spring.gradle:dependency-management-plugin:1.0.2.RELEASE' + classpath 'io.spring.gradle:dependency-management-plugin:1.0.3.RELEASE' classpath 'io.spring.gradle:spring-io-plugin:0.0.8.RELEASE' classpath 'io.spring.gradle:docbook-reference-plugin:0.3.1' classpath 'org.asciidoctor:asciidoctor-gradle-plugin:1.5.0' @@ -87,9 +87,9 @@ subprojects { subproject -> } ext { - activeMqVersion = '5.14.5' + activeMqVersion = '5.15.0' aspectjVersion = '1.8.10' - apacheSshdVersion = '1.4.0' + apacheSshdVersion = '1.6.0' boonVersion = '0.34' commonsDbcp2Version = '2.1.1' commonsIoVersion = '2.4' @@ -97,51 +97,51 @@ subprojects { subproject -> curatorVersion = '2.11.1' derbyVersion = '10.13.1.1' eclipseLinkVersion = '2.6.4' - ftpServerVersion = '1.0.6' - groovyVersion = '2.4.10' + ftpServerVersion = '1.1.1' + groovyVersion = '2.4.12' guavaVersion = '20.0' hamcrestVersion = '1.3' hazelcastVersion = '3.8' hibernateVersion = '5.2.10.Final' hsqldbVersion = '2.4.0' - h2Version = '1.4.194' - jackson2Version = '2.9.0.pr4' + h2Version = '1.4.196' + jackson2Version = '2.9.0' javaxActivationVersion = '1.1.1' - javaxMailVersion = '1.6.0-rc1' + javaxMailVersion = '1.6.0' jedisVersion = '2.9.0' jmsApiVersion = '2.0.1' jpa21ApiVersion = '1.0.0.Final' jpaApiVersion = '2.1.1' jrubyVersion = '9.1.5.0' jschVersion = '0.1.54' - jsonpathVersion = '2.3.0' + jsonpathVersion = '2.4.0' junitVersion = '4.12' jythonVersion = '2.5.3' kryoShadedVersion = '3.0.3' log4jVersion = '1.2.17' - mockitoVersion = '2.7.22' - mysqlVersion = '5.1.41' + mockitoVersion = '2.9.0' + mysqlVersion = '6.0.6' pahoMqttClientVersion = '1.1.1' postgresVersion = '42.0.0' reactorNettyVersion = '0.7.0.BUILD-SNAPSHOT' reactorVersion = '3.1.0.BUILD-SNAPSHOT' - romeToolsVersion = '1.7.2' + romeToolsVersion = '1.7.4' servletApiVersion = '3.1.0' slf4jVersion = "1.7.25" smackVersion = '4.1.9' - springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.0.0.M5' - springDataJpaVersion = '2.0.0.RC2' - springDataMongoVersion = '2.0.0.RC2' - springDataRedisVersion = '2.0.0.RC2' - springGemfireVersion = '2.0.0.RC2' + springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.0.0.BUILD-SNAPSHOT' + springDataJpaVersion = '2.0.0.BUILD-SNAPSHOT' + springDataMongoVersion = '2.0.0.BUILD-SNAPSHOT' + springDataRedisVersion = '2.0.0.BUILD-SNAPSHOT' + springGemfireVersion = '2.0.0.BUILD-SNAPSHOT' springSecurityVersion = '5.0.0.BUILD-SNAPSHOT' - springSocialTwitterVersion = '2.0.0.M4' + springSocialTwitterVersion = '2.0.0.BUILD-SNAPSHOT' springRetryVersion = '1.2.0.RELEASE' springVersion = project.hasProperty('springVersion') ? project.springVersion : '5.0.0.BUILD-SNAPSHOT' springWsVersion = '2.4.0.RELEASE' - tomcatVersion = "8.5.16" + tomcatVersion = "8.5.20" xmlUnitVersion = '1.6' - xstreamVersion = '1.4.7' + xstreamVersion = '1.4.10' } eclipse { diff --git a/gradle/wrapper/gradle-wrapper.jar b/gradle/wrapper/gradle-wrapper.jar index bde50e136b69b14f820067e9a0b87aac662c7cc1..a4922b260836263ac7a10ff92b27d7fcb7bf905b 100644 GIT binary patch delta 28 hcmdn7nt8`+<_(_?vz$v+^WOaJaFHOGH96<14*<=C4Uqr< delta 28 hcmdn7nt8`+<_(_?vxKGp@Y?+CaFHOGH96<14*=7A4o3h0 diff --git a/gradle/wrapper/gradle-wrapper.properties b/gradle/wrapper/gradle-wrapper.properties index e70f80422e..54b71bd52c 100644 --- a/gradle/wrapper/gradle-wrapper.properties +++ b/gradle/wrapper/gradle-wrapper.properties @@ -1,6 +1,6 @@ -#Mon Jul 24 12:58:44 EDT 2017 +#Wed Sep 06 12:46:29 EDT 2017 distributionBase=GRADLE_USER_HOME distributionPath=wrapper/dists zipStoreBase=GRADLE_USER_HOME zipStorePath=wrapper/dists -distributionUrl=https\://services.gradle.org/distributions/gradle-4.0.1-bin.zip +distributionUrl=https\://services.gradle.org/distributions/gradle-4.1-bin.zip diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/MessageChannelReactiveUtils.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/MessageChannelReactiveUtils.java index 6d1a77db42..1e78752cf3 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/MessageChannelReactiveUtils.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/MessageChannelReactiveUtils.java @@ -106,7 +106,7 @@ public final class MessageChannelReactiveUtils { .>create(sink -> sink.onRequest(n -> { Message m; - while (n-- > 0 && (m = this.channel.receive()) != null) { + while (!sink.isCancelled() && n-- > 0 && (m = this.channel.receive()) != null) { sink.next((Message) m); } }), diff --git a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ResequencerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ResequencerTests.java index 473c207402..5991e894d3 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ResequencerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/aggregator/ResequencerTests.java @@ -262,8 +262,9 @@ public class ResequencerTests { assertNotNull(reply1); assertNotNull(reply2); assertNull(reply3); - ArrayList sequence = new ArrayList(Arrays.asList(new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber(), - new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber())); + ArrayList sequence = new ArrayList<>( + Arrays.asList(new IntegrationMessageHeaderAccessor(reply1).getSequenceNumber(), + new IntegrationMessageHeaderAccessor(reply2).getSequenceNumber())); Collections.sort(sequence); assertEquals("[1, 2]", sequence.toString()); // Once a group is expired, late messages are discarded immediately by default @@ -363,7 +364,7 @@ public class ResequencerTests { this.resequencer.handleMessage(message2); Message out1 = replyChannel.receive(10); assertNull(out1); - out1 = discardChannel.receive(1000); + out1 = discardChannel.receive(10000); assertNotNull(out1); Message out2 = discardChannel.receive(10); assertNotNull(out2); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java index 9a917c3c99..9b10aa5cff 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/reactivestreams/ReactiveStreamsTests.java @@ -33,7 +33,6 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.logging.Level; -import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; import org.reactivestreams.Publisher; @@ -105,7 +104,6 @@ public class ReactiveStreamsTests { } @Test - @Ignore public void testPollableReactiveFlow() throws Exception { this.inputChannel.send(new GenericMessage<>("1,2,3,4,5")); diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbInboundChannelAdapterIntegrationTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbInboundChannelAdapterIntegrationTests.java index d7684653da..3d824a72fa 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbInboundChannelAdapterIntegrationTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbInboundChannelAdapterIntegrationTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2017 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. @@ -32,7 +32,6 @@ import org.springframework.beans.factory.parsing.BeanDefinitionParsingException; import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.data.mongodb.core.MongoOperations; import org.springframework.data.mongodb.core.MongoTemplate; -import org.springframework.data.mongodb.core.query.BasicQuery; import org.springframework.data.mongodb.core.query.Criteria; import org.springframework.data.mongodb.core.query.Query; import org.springframework.integration.aop.AbstractMessageSourceAdvice; @@ -46,12 +45,11 @@ import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; -import com.mongodb.util.JSON; - /** * @author Oleg Zhurakousky * @author Artem Bilan * @author Yaron Yamin + * * @since 2.2 */ @ContextConfiguration @@ -168,6 +166,7 @@ public class MongoDbInboundChannelAdapterIntegrationTests extends MongoDbAvailab this.mongoInboundAdapterWithNamedCollection.stop(); this.replyChannel.purge(null); } + @Test @MongoDbAvailable public void testWithQueryExpression() throws Exception { @@ -181,6 +180,7 @@ public class MongoDbInboundChannelAdapterIntegrationTests extends MongoDbAvailab assertEquals("Bob", message.getPayload().get(0).getName()); this.mongoInboundAdapterWithQueryExpression.stop(); } + @Test @MongoDbAvailable public void testWithStringQueryExpression() throws Exception { @@ -273,7 +273,7 @@ public class MongoDbInboundChannelAdapterIntegrationTests extends MongoDbAvailab if (target instanceof List) { List documents = (List) target; for (Object document : documents) { - mongoOperations.remove(new BasicQuery(JSON.serialize(document)), collectionName); + mongoOperations.remove(document, collectionName); } } } diff --git a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbOutboundChannelAdapterIntegrationTests.java b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbOutboundChannelAdapterIntegrationTests.java index 3d6dcddd2c..82680333f1 100644 --- a/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbOutboundChannelAdapterIntegrationTests.java +++ b/spring-integration-mongodb/src/test/java/org/springframework/integration/mongodb/config/MongoDbOutboundChannelAdapterIntegrationTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2016 the original author or authors. + * Copyright 2002-2017 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. @@ -32,16 +32,18 @@ import org.springframework.messaging.MessageChannel; import org.springframework.messaging.support.GenericMessage; import com.mongodb.BasicDBObject; -import com.mongodb.util.JSON; + /** * @author Oleg Zhurakousky + * @author Artem Bilan */ public class MongoDbOutboundChannelAdapterIntegrationTests extends MongoDbAvailableTests { @Test @MongoDbAvailable public void testWithDefaultMongoFactory() throws Exception { - ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); + ClassPathXmlApplicationContext context = + new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); MessageChannel channel = context.getBean("simpleAdapter", MessageChannel.class); Message message = new GenericMessage(this.createPerson("Bob")); @@ -57,10 +59,14 @@ public class MongoDbOutboundChannelAdapterIntegrationTests extends MongoDbAvaila @MongoDbAvailable public void testWithNamedCollection() throws Exception { MongoDbFactory mongoDbFactory = this.prepareMongoFactory("foo"); - ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); + ClassPathXmlApplicationContext context = + new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); MessageChannel channel = context.getBean("simpleAdapterWithNamedCollection", MessageChannel.class); - Message message = MessageBuilder.withPayload(this.createPerson("Bob")).setHeader("collectionName", "foo").build(); + Message message = + MessageBuilder.withPayload(this.createPerson("Bob")) + .setHeader("collectionName", "foo") + .build(); channel.send(message); MongoTemplate template = new MongoTemplate(mongoDbFactory); @@ -72,10 +78,15 @@ public class MongoDbOutboundChannelAdapterIntegrationTests extends MongoDbAvaila @MongoDbAvailable public void testWithTemplate() throws Exception { MongoDbFactory mongoDbFactory = this.prepareMongoFactory("foo"); - ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); + ClassPathXmlApplicationContext context = + new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); MessageChannel channel = context.getBean("simpleAdapterWithTemplate", MessageChannel.class); - Message message = MessageBuilder.withPayload(this.createPerson("Bob")).setHeader("collectionName", "foo").build(); + Message message = + MessageBuilder.withPayload(this.createPerson("Bob")) + .setHeader("collectionName", "foo") + .build(); + channel.send(message); MongoTemplate template = new MongoTemplate(mongoDbFactory); @@ -87,13 +98,19 @@ public class MongoDbOutboundChannelAdapterIntegrationTests extends MongoDbAvaila @MongoDbAvailable public void testSavingDbObject() throws Exception { - BasicDBObject dbObject = (BasicDBObject) JSON.parse("{'foo' : 'bar'}"); + BasicDBObject dbObject = BasicDBObject.parse("{'foo' : 'bar'}"); MongoDbFactory mongoDbFactory = this.prepareMongoFactory("foo"); - ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); + ClassPathXmlApplicationContext context = + new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); MessageChannel channel = context.getBean("simpleAdapterWithTemplate", MessageChannel.class); - Message message = MessageBuilder.withPayload(dbObject).setHeader("collectionName", "foo").build(); + + Message message = + MessageBuilder.withPayload(dbObject) + .setHeader("collectionName", "foo") + .build(); + channel.send(message); MongoTemplate template = new MongoTemplate(mongoDbFactory); @@ -108,10 +125,16 @@ public class MongoDbOutboundChannelAdapterIntegrationTests extends MongoDbAvaila String object = "{'foo' : 'bar'}"; MongoDbFactory mongoDbFactory = this.prepareMongoFactory("foo"); - ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); + ClassPathXmlApplicationContext context = + new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); MessageChannel channel = context.getBean("simpleAdapterWithTemplate", MessageChannel.class); - Message message = MessageBuilder.withPayload(object).setHeader("collectionName", "foo").build(); + + Message message = + MessageBuilder.withPayload(object) + .setHeader("collectionName", "foo") + .build(); + channel.send(message); MongoTemplate template = new MongoTemplate(mongoDbFactory); @@ -122,7 +145,8 @@ public class MongoDbOutboundChannelAdapterIntegrationTests extends MongoDbAvaila @Test @MongoDbAvailable public void testWithMongoConverter() throws Exception { - ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); + ClassPathXmlApplicationContext context = + new ClassPathXmlApplicationContext("outbound-adapter-config.xml", this.getClass()); MessageChannel channel = context.getBean("simpleAdapterWithConverter", MessageChannel.class); Message message = new GenericMessage(this.createPerson("Bob")); @@ -133,4 +157,5 @@ public class MongoDbOutboundChannelAdapterIntegrationTests extends MongoDbAvaila assertNotNull(template.find(new BasicQuery("{'name' : 'Bob'}"), Person.class, "data")); context.close(); } + } 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 2454539ab8..ceac69a7e9 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 @@ -42,13 +42,14 @@ import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.mongodb.rules.MongoDbAvailable; import org.springframework.integration.mongodb.rules.MongoDbAvailableTests; -import com.mongodb.util.JSON; +import com.mongodb.BasicDBObject; /** * @author Amol Nayak * @author Oleg Zhurakousky * @author Gary Russell * @author Yaron Yamin + * @author Artem Bilan * * @since 2.2 * @@ -280,7 +281,7 @@ public class MongoDbMessageSourceTests extends MongoDbAvailableTests { MongoTemplate template = new MongoTemplate(mongoDbFactory); - template.save(JSON.parse("{'name' : 'Manny', 'id' : 1}"), "data"); + template.save(BasicDBObject.parse("{'name' : 'Manny', 'id' : 1}"), "data"); Expression queryExpression = new LiteralExpression("{'name' : 'Manny'}"); MongoDbMessageSource messageSource = new MongoDbMessageSource(mongoDbFactory, queryExpression); diff --git a/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/inbound/WebFluxInboundEndpoint.java b/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/inbound/WebFluxInboundEndpoint.java index 0bb4f2a526..f6d3221c60 100644 --- a/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/inbound/WebFluxInboundEndpoint.java +++ b/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/inbound/WebFluxInboundEndpoint.java @@ -156,7 +156,7 @@ public class WebFluxInboundEndpoint extends BaseHttpInboundEndpoint implements W return setStatusCode(exchange); } }) - .doOnTerminate((e, t) -> this.activeCount.decrementAndGet()); + .doOnTerminate(this.activeCount::decrementAndGet); }