From e6af887972304224b05c7fc8b9e9169d5acda98e Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 21 Jan 2016 17:01:43 -0500 Subject: [PATCH] Upgrade to Rector-2.5 --- build.gradle | 8 ++------ .../AmqpInboundGatewayParserTests-context.xml | 4 +++- ...essageChannel.java => ReactiveChannel.java} | 12 ++++++------ .../config/ConsumerEndpointFactoryBean.java | 14 +++++--------- .../AbstractMethodAnnotationPostProcessor.java | 10 +++++----- .../integration/endpoint/ReactiveEndpoint.java | 18 +++++++++++++----- .../gateway/GatewayProxyFactoryBean.java | 5 +++-- ...nelTests.java => ReactiveChannelTests.java} | 8 ++++---- .../config/xml/GatewayParserTests.java | 4 ++-- .../configuration/EnableIntegrationTests.java | 6 +++--- .../integration/gateway/AsyncGatewayTests.java | 10 +++++----- .../routingslip/RoutingSlipTests.java | 8 ++++---- 12 files changed, 55 insertions(+), 52 deletions(-) rename spring-integration-core/src/main/java/org/springframework/integration/channel/{ReactiveMessageChannel.java => ReactiveChannel.java} (82%) rename spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/{ReactiveMessageChannelTests.java => ReactiveChannelTests.java} (92%) diff --git a/build.gradle b/build.gradle index 805079343d..f1c2f5e353 100644 --- a/build.gradle +++ b/build.gradle @@ -293,10 +293,7 @@ project('spring-integration-core') { exclude group: 'org.springframework', module: 'spring-core' } // compile ("org.springframework.cloud:spring-cloud-cluster-core:$springCloudClusterVersion", optional) - compile ("io.projectreactor:reactor-stream:$reactorVersion") { - optional it - exclude group: 'org.slf4j', module: 'slf4j-api' - } + compile "io.projectreactor:reactor-core:$reactorVersion" compile("com.fasterxml.jackson.core:jackson-databind:$jackson2Version", optional) compile("com.jayway.jsonpath:json-path:$jsonpathVersion") { optional it @@ -309,8 +306,7 @@ project('spring-integration-core') { // testCompile ("net.openhft:chronicle:$chronicleVersion") // testCompile ("io.projectreactor:reactor-chronicle:$reactorVersion") - testCompile "org.testng:testng:6.8.21" - testCompile "org.reactivestreams:reactive-streams-tck:1.0.0" + testCompile "io.projectreactor:reactor-stream:$reactorVersion" } } diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpInboundGatewayParserTests-context.xml b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpInboundGatewayParserTests-context.xml index b7da686bf8..f8c7e18ee1 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpInboundGatewayParserTests-context.xml +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpInboundGatewayParserTests-context.xml @@ -12,7 +12,9 @@ + connection-factory="rabbitConnectionFactory" message-converter="testConverter"> + + diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveMessageChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveChannel.java similarity index 82% rename from spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveMessageChannel.java rename to spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveChannel.java index ef521cd897..7b44f5c9db 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveMessageChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveChannel.java @@ -23,24 +23,24 @@ import org.reactivestreams.Subscriber; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; -import reactor.core.subscription.ReactiveSession; -import reactor.rx.broadcast.Broadcaster; +import reactor.core.publisher.ProcessorGroup; +import reactor.core.subscriber.ReactiveSession; /** * @author Artem Bilan * @since 5.0 */ -public class ReactiveMessageChannel implements MessageChannel, Publisher> { +public class ReactiveChannel implements MessageChannel, Publisher> { private final Processor, Message> processor; private final ReactiveSession> reactiveSession; - public ReactiveMessageChannel() { - this(Broadcaster.passthrough()); + public ReactiveChannel() { + this(ProcessorGroup.>sync().get()); } - public ReactiveMessageChannel(Processor, Message> processor) { + public ReactiveChannel(Processor, Message> processor) { this.processor = processor; this.reactiveSession = ReactiveSession.create(processor); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java b/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java index fa919b454e..06bbbb20b0 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/ConsumerEndpointFactoryBean.java @@ -21,7 +21,6 @@ import java.util.List; import org.aopalliance.aop.Advice; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; -import org.reactivestreams.Publisher; import org.reactivestreams.Subscriber; import org.springframework.aop.framework.Advised; @@ -56,7 +55,8 @@ import org.springframework.util.Assert; import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; -import reactor.core.subscriber.SubscriberFactory; +import reactor.core.subscriber.Subscribers; + /** * @author Mark Fisher @@ -289,20 +289,16 @@ public class ConsumerEndpointFactoryBean pollingConsumer.setBeanFactory(this.beanFactory); this.endpoint = pollingConsumer; } - else if (channel instanceof Publisher) { - Publisher> publisher = (Publisher>) channel; + else { Subscriber> subscriber; if (this.handler instanceof Subscriber) { subscriber = (Subscriber>) this.handler; } else { //TODO errorConsumer, completeConsumer - subscriber = SubscriberFactory.consumer(this.handler::handleMessage); + subscriber = Subscribers.consumer(this.handler::handleMessage); } - this.endpoint = new ReactiveEndpoint(publisher, subscriber); - } - else { - throw new IllegalArgumentException("unsupported channel type: [" + channel.getClass() + "]"); + this.endpoint = new ReactiveEndpoint(channel, subscriber); } this.endpoint.setBeanName(this.beanName); this.endpoint.setBeanFactory(this.beanFactory); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java index 32edf75988..4e5c7d9900 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/annotation/AbstractMethodAnnotationPostProcessor.java @@ -80,7 +80,7 @@ import org.springframework.util.CollectionUtils; import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; -import reactor.core.subscriber.SubscriberFactory; +import reactor.core.subscriber.Subscribers; /** * Base class for Method-level annotation post-processors. @@ -306,7 +306,8 @@ public abstract class AbstractMethodAnnotationPostProcessor annotations) { + @SuppressWarnings("unchecked") + protected AbstractEndpoint doCreateEndpoint(MessageHandler handler, MessageChannel inputChannel,List annotations) { AbstractEndpoint endpoint; if (inputChannel instanceof PollableChannel) { PollingConsumer pollingConsumer = new PollingConsumer((PollableChannel) inputChannel, handler); @@ -318,16 +319,15 @@ public abstract class AbstractMethodAnnotationPostProcessor> publisher = (Publisher>) inputChannel; Subscriber> subscriber; if (handler instanceof Subscriber) { subscriber = (Subscriber>) handler; } else { //TODO errorConsumer, completeConsumer - subscriber = SubscriberFactory.consumer(handler::handleMessage); + subscriber = Subscribers.consumer(handler::handleMessage); } - endpoint = new ReactiveEndpoint(publisher, subscriber); + endpoint = new ReactiveEndpoint(inputChannel, subscriber); } else { endpoint = new EventDrivenConsumer((SubscribableChannel) inputChannel, handler); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveEndpoint.java b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveEndpoint.java index e681c17919..b6e5bd994e 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveEndpoint.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveEndpoint.java @@ -20,9 +20,11 @@ import org.reactivestreams.Publisher; import org.reactivestreams.Subscriber; import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; import org.springframework.util.Assert; -import reactor.core.subscription.ReactiveSession; +import reactor.core.subscriber.ReactiveSession; + /** * @author Artem Bilan @@ -36,11 +38,17 @@ public class ReactiveEndpoint extends AbstractEndpoint { private ReactiveSession> reactiveSession; - public ReactiveEndpoint(Publisher> inputChannel, Subscriber> subscriber) { - Assert.isInstanceOf(Publisher.class, inputChannel, - "The 'inputChannel', must implement org.reactivestreams.Publisher."); + @SuppressWarnings("unchecked") + public ReactiveEndpoint(MessageChannel inputChannel, Subscriber> subscriber) { + Assert.notNull(inputChannel); Assert.notNull(subscriber); - this.inputChannel = inputChannel; + if (inputChannel instanceof Publisher) { + this.inputChannel = (Publisher>) inputChannel; + } + else { + //TODO: Wrap all other channels to the Publisher + this.inputChannel = null; + } this.subscriber = subscriber; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java b/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java index d0d33ce2d4..b14913a17a 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/gateway/GatewayProxyFactoryBean.java @@ -66,7 +66,8 @@ import org.springframework.util.ObjectUtils; import org.springframework.util.ReflectionUtils; import org.springframework.util.StringUtils; -import reactor.rx.Promises; +import reactor.core.publisher.Mono; + /** * Generates a proxy for the provided service interface to enable interaction @@ -400,7 +401,7 @@ public class GatewayProxyFactoryBean extends AbstractEndpoint } } if (reactorPresent && Publisher.class.isAssignableFrom(returnType)) { - return Promises.task(() -> new AsyncInvocationTask(invocation)); + return Mono.fromCallable(new AsyncInvocationTask(invocation)); } return this.doInvoke(invocation, true); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveMessageChannelTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveChannelTests.java similarity index 92% rename from spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveMessageChannelTests.java rename to spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveChannelTests.java index 34cc9cd3a8..912f5c59c5 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveMessageChannelTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/reactive/ReactiveChannelTests.java @@ -26,7 +26,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.integration.annotation.ServiceActivator; import org.springframework.integration.channel.QueueChannel; -import org.springframework.integration.channel.ReactiveMessageChannel; +import org.springframework.integration.channel.ReactiveChannel; import org.springframework.integration.config.EnableIntegration; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; @@ -35,7 +35,7 @@ import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; -import reactor.Processors; +import reactor.core.publisher.Processors; /** * @author Artem Bilan @@ -44,7 +44,7 @@ import reactor.Processors; @ContextConfiguration @RunWith(SpringJUnit4ClassRunner.class) @DirtiesContext -public class ReactiveMessageChannelTests { +public class ReactiveChannelTests { @Autowired private MessageChannel reactiveChannel; @@ -71,7 +71,7 @@ public class ReactiveMessageChannelTests { @Bean public MessageChannel reactiveChannel() { - return new ReactiveMessageChannel(Processors.queue()); + return new ReactiveChannel(Processors.queue()); } @ServiceActivator(inputChannel = "reactiveChannel") diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/GatewayParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/GatewayParserTests.java index 5ac4f2d57a..78bcfa612b 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/GatewayParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/GatewayParserTests.java @@ -65,7 +65,7 @@ import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; -import reactor.rx.Promises; +import reactor.rx.Promise; /** * @author Mark Fisher @@ -175,7 +175,7 @@ public class GatewayParserTests { this.startResponder(requestChannel, replyChannel); TestService service = context.getBean("promise", TestService.class); Publisher> result = service.promise("foo"); - Message reply = Promises.from(result).await(1, TimeUnit.SECONDS); + Message reply = Promise.from(result).await(1, TimeUnit.SECONDS); assertEquals("foo", reply.getPayload()); assertNotNull(TestUtils.getPropertyValue(context.getBean("&promise"), "asyncExecutor")); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java index 6ff8522d96..731c80b494 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/configuration/EnableIntegrationTests.java @@ -127,7 +127,7 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import org.springframework.test.context.support.AnnotationConfigContextLoader; import org.springframework.util.MultiValueMap; -import reactor.rx.Streams; +import reactor.rx.Stream; /** * @author Artem Bilan @@ -617,11 +617,11 @@ public class EnableIntegrationTests { final AtomicReference> ref = new AtomicReference>(); final CountDownLatch consumeLatch = new CountDownLatch(1); - Streams.just("1", "2", "3", "4", "5") + Stream.just("1", "2", "3", "4", "5") .map(Integer::parseInt) .flatMap(this.testGateway::multiply) .toList() - .onSuccess(integers -> { + .doOnSuccess(integers -> { ref.set(integers); consumeLatch.countDown(); }); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/gateway/AsyncGatewayTests.java b/spring-integration-core/src/test/java/org/springframework/integration/gateway/AsyncGatewayTests.java index 810a5a8111..8dfae0d6ef 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/gateway/AsyncGatewayTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/gateway/AsyncGatewayTests.java @@ -47,7 +47,7 @@ import org.springframework.util.concurrent.ListenableFuture; import org.springframework.util.concurrent.ListenableFutureCallback; import reactor.fn.Consumer; -import reactor.rx.Promises; +import reactor.rx.Promise; /** * @author Mark Fisher @@ -252,7 +252,7 @@ public class AsyncGatewayTests { proxyFactory.afterPropertiesSet(); TestEchoService service = (TestEchoService) proxyFactory.getObject(); Publisher> promise = service.returnMessagePromise("foo"); - Object result = Promises.from(promise).await(10, TimeUnit.SECONDS); + Object result = Promise.from(promise).await(10, TimeUnit.SECONDS); assertEquals("foobar", ((Message) result).getPayload()); } @@ -268,7 +268,7 @@ public class AsyncGatewayTests { proxyFactory.afterPropertiesSet(); TestEchoService service = (TestEchoService) proxyFactory.getObject(); Publisher promise = service.returnStringPromise("foo"); - Object result = Promises.from(promise).await(10, TimeUnit.SECONDS); + Object result = Promise.from(promise).await(10, TimeUnit.SECONDS); assertEquals("foobar", result); } @@ -284,7 +284,7 @@ public class AsyncGatewayTests { proxyFactory.afterPropertiesSet(); TestEchoService service = (TestEchoService) proxyFactory.getObject(); Publisher promise = service.returnSomethingPromise("foo"); - Object result = Promises.from(promise).await(10, TimeUnit.SECONDS); + Object result = Promise.from(promise).await(10, TimeUnit.SECONDS); assertNotNull(result); assertEquals("foobar", result); } @@ -305,7 +305,7 @@ public class AsyncGatewayTests { final AtomicReference result = new AtomicReference(); final CountDownLatch latch = new CountDownLatch(1); - Promises.from(promise).onSuccess(new Consumer() { + Promise.from(promise).doOnSuccess(new Consumer() { @Override public void accept(String s) { result.set(s); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests.java b/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests.java index 699c752b09..79b9c01b96 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/routingslip/RoutingSlipTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2015 the original author or authors. + * Copyright 2014-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. @@ -60,7 +60,7 @@ import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import org.springframework.test.context.support.GenericXmlContextLoader; -import reactor.rx.Streams; +import reactor.rx.Stream; /** * @author Artem Bilan @@ -189,11 +189,11 @@ public class RoutingSlipTests { public RoutingSlipRouteStrategy routeStrategy() { return (requestMessage, reply) -> requestMessage.getPayload() instanceof String ? new FixedSubscriberChannel(m -> - Streams.just((String) m.getPayload()) + Stream.just((String) m.getPayload()) .map(String::toUpperCase) .consume(v -> messagingTemplate().convertAndSend(resultsChannel(), v))) : new FixedSubscriberChannel(m -> - Streams.just((Integer) m.getPayload()) + Stream.just((Integer) m.getPayload()) .map(v -> v * 2) .consume(v -> messagingTemplate().convertAndSend(resultsChannel(), v))); }