diff --git a/build.gradle b/build.gradle
index 0ddeca18b6..805079343d 100644
--- a/build.gradle
+++ b/build.gradle
@@ -73,16 +73,10 @@ subprojects { subproject ->
}
compileJava {
- sourceCompatibility = 1.6
- targetCompatibility = 1.6
+ sourceCompatibility = 1.8
+ targetCompatibility = 1.8
}
- compileTestJava {
- sourceCompatibility = 1.6
- targetCompatibility = 1.6
- }
-
-
ext {
activeMqVersion = '5.13.2'
aspectjVersion = '1.8.9'
@@ -123,7 +117,7 @@ subprojects { subproject ->
openJpaVersion = '2.4.0'
pahoMqttClientVersion = '1.0.2'
postgresVersion = '9.1-901-1.jdbc4'
- reactorVersion = '2.0.8.RELEASE'
+ reactorVersion = '2.5.0.BUILD-SNAPSHOT'
romeToolsVersion = '1.6.0'
servletApiVersion = '3.1.0'
slf4jVersion = "1.7.21"
@@ -289,11 +283,6 @@ project('spring-integration-amqp') {
project('spring-integration-core') {
description = 'Spring Integration Core'
- compileTestJava {
- sourceCompatibility = 1.8
- targetCompatibility = 1.8
- }
-
dependencies {
compile "org.springframework:spring-core:$springVersion"
compile "org.springframework:spring-aop:$springVersion"
@@ -317,7 +306,11 @@ project('spring-integration-core') {
compile("com.esotericsoftware:kryo-shaded:$kryoShadedVersion", optional)
testCompile ("org.aspectj:aspectjweaver:$aspectjVersion")
- testCompile ("net.openhft:chronicle:$chronicleVersion")
+// 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"
}
}
diff --git a/gradle.properties b/gradle.properties
index 3f49498608..9136e989dd 100644
--- a/gradle.properties
+++ b/gradle.properties
@@ -1,2 +1,2 @@
-version=4.3.2.BUILD-SNAPSHOT
+version=5.0.0.BUILD-SNAPSHOT
org.gradle.daemon=true
diff --git a/spring-integration-core/src/main/java/org/springframework/integration/annotation/MessagingGateway.java b/spring-integration-core/src/main/java/org/springframework/integration/annotation/MessagingGateway.java
index 489c66cfb8..6f9280b611 100644
--- a/spring-integration-core/src/main/java/org/springframework/integration/annotation/MessagingGateway.java
+++ b/spring-integration-core/src/main/java/org/springframework/integration/annotation/MessagingGateway.java
@@ -127,17 +127,4 @@ public @interface MessagingGateway {
*/
String mapper() default "";
- /**
- * Provide a reference to an {@link reactor.Environment}
- * to use for any of the interface methods that have a {@link reactor.rx.Promise} return type.
- * This {@code reactor.core.Environment} will only be used for those async methods; the sync methods
- * will be invoked in the caller's thread.
- *
This attribute is required in case of {@link reactor.rx.Promise} usage.
- * @return the suggested reactor Environment bean name.
- * @since 4.1
- * @deprecated with no-op in favor of global JVM-wide Reactor configuration.
- */
- @Deprecated
- String reactorEnvironment() default "";
-
}
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/ReactiveMessageChannel.java
new file mode 100644
index 0000000000..ef521cd897
--- /dev/null
+++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/ReactiveMessageChannel.java
@@ -0,0 +1,65 @@
+/*
+ * Copyright 2015 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.channel;
+
+import org.reactivestreams.Processor;
+import org.reactivestreams.Publisher;
+import org.reactivestreams.Subscriber;
+
+import org.springframework.messaging.Message;
+import org.springframework.messaging.MessageChannel;
+
+import reactor.core.subscription.ReactiveSession;
+import reactor.rx.broadcast.Broadcaster;
+
+/**
+ * @author Artem Bilan
+ * @since 5.0
+ */
+public class ReactiveMessageChannel implements MessageChannel, Publisher> {
+
+ private final Processor, Message>> processor;
+
+ private final ReactiveSession> reactiveSession;
+
+ public ReactiveMessageChannel() {
+ this(Broadcaster.passthrough());
+ }
+
+ public ReactiveMessageChannel(Processor, Message>> processor) {
+ this.processor = processor;
+ this.reactiveSession = ReactiveSession.create(processor);
+ }
+
+ @Override
+ public boolean send(Message> message) {
+ return send(message, -1);
+ }
+
+ @Override
+ @SuppressWarnings("unchecked")
+ public boolean send(Message> message, long timeout) {
+ this.reactiveSession.submit(message, timeout);
+ return true;
+ }
+
+ @Override
+ public void subscribe(Subscriber super Message>> subscriber) {
+ this.processor.subscribe(subscriber);
+ }
+
+}
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 96bd24ca1a..fa919b454e 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,6 +21,8 @@ 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;
import org.springframework.aop.framework.ProxyFactory;
@@ -39,9 +41,11 @@ import org.springframework.integration.context.IntegrationObjectSupport;
import org.springframework.integration.endpoint.AbstractEndpoint;
import org.springframework.integration.endpoint.EventDrivenConsumer;
import org.springframework.integration.endpoint.PollingConsumer;
+import org.springframework.integration.endpoint.ReactiveEndpoint;
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
import org.springframework.integration.handler.advice.HandleMessageAdvice;
import org.springframework.integration.scheduling.PollerMetadata;
+import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.PollableChannel;
@@ -52,6 +56,8 @@ import org.springframework.util.Assert;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;
+import reactor.core.subscriber.SubscriberFactory;
+
/**
* @author Mark Fisher
* @author Oleg Zhurakousky
@@ -240,6 +246,7 @@ public class ConsumerEndpointFactoryBean
return this.endpoint.getClass();
}
+ @SuppressWarnings("unchecked")
private void initializeEndpoint() throws Exception {
synchronized (this.initializationMonitor) {
if (this.initialized) {
@@ -282,6 +289,18 @@ public class ConsumerEndpointFactoryBean
pollingConsumer.setBeanFactory(this.beanFactory);
this.endpoint = pollingConsumer;
}
+ else if (channel instanceof Publisher) {
+ Publisher> publisher = (Publisher>) channel;
+ Subscriber> subscriber;
+ if (this.handler instanceof Subscriber) {
+ subscriber = (Subscriber>) this.handler;
+ }
+ else {
+ //TODO errorConsumer, completeConsumer
+ subscriber = SubscriberFactory.consumer(this.handler::handleMessage);
+ }
+ this.endpoint = new ReactiveEndpoint(publisher, subscriber);
+ }
else {
throw new IllegalArgumentException("unsupported channel type: [" + channel.getClass() + "]");
}
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 dfc7496cae..32edf75988 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
@@ -26,6 +26,8 @@ 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.TargetSource;
import org.springframework.aop.framework.Advised;
@@ -53,6 +55,7 @@ import org.springframework.integration.endpoint.AbstractEndpoint;
import org.springframework.integration.endpoint.AbstractPollingEndpoint;
import org.springframework.integration.endpoint.EventDrivenConsumer;
import org.springframework.integration.endpoint.PollingConsumer;
+import org.springframework.integration.endpoint.ReactiveEndpoint;
import org.springframework.integration.endpoint.SourcePollingChannelAdapter;
import org.springframework.integration.handler.AbstractMessageProducingHandler;
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
@@ -61,6 +64,7 @@ import org.springframework.integration.router.AbstractMessageRouter;
import org.springframework.integration.scheduling.PollerMetadata;
import org.springframework.integration.support.channel.BeanFactoryChannelResolver;
import org.springframework.integration.util.MessagingAnnotationUtils;
+import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.PollableChannel;
@@ -76,6 +80,8 @@ import org.springframework.util.CollectionUtils;
import org.springframework.util.ObjectUtils;
import org.springframework.util.StringUtils;
+import reactor.core.subscriber.SubscriberFactory;
+
/**
* Base class for Method-level annotation post-processors.
*
@@ -311,7 +317,21 @@ 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);
+ }
+ endpoint = new ReactiveEndpoint(publisher, subscriber);
+ }
+ else {
+ endpoint = new EventDrivenConsumer((SubscribableChannel) inputChannel, handler);
+ }
}
return endpoint;
}
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
new file mode 100644
index 0000000000..e681c17919
--- /dev/null
+++ b/spring-integration-core/src/main/java/org/springframework/integration/endpoint/ReactiveEndpoint.java
@@ -0,0 +1,58 @@
+/*
+ * 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.endpoint;
+
+import org.reactivestreams.Publisher;
+import org.reactivestreams.Subscriber;
+
+import org.springframework.messaging.Message;
+import org.springframework.util.Assert;
+
+import reactor.core.subscription.ReactiveSession;
+
+/**
+ * @author Artem Bilan
+ * @since 5.0
+ */
+public class ReactiveEndpoint extends AbstractEndpoint {
+
+ private final Publisher> inputChannel;
+
+ private final Subscriber> subscriber;
+
+ private ReactiveSession> reactiveSession;
+
+ public ReactiveEndpoint(Publisher> inputChannel, Subscriber> subscriber) {
+ Assert.isInstanceOf(Publisher.class, inputChannel,
+ "The 'inputChannel', must implement org.reactivestreams.Publisher.");
+ Assert.notNull(subscriber);
+ this.inputChannel = inputChannel;
+ this.subscriber = subscriber;
+ }
+
+ @Override
+ protected void doStart() {
+ this.reactiveSession = ReactiveSession.create(this.subscriber);
+ this.inputChannel.subscribe(this.reactiveSession);
+ }
+
+ @Override
+ protected void doStop() {
+ this.reactiveSession.finish();
+ }
+
+}
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 b7e63f47a5..d0d33ce2d4 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
@@ -29,6 +29,7 @@ import java.util.concurrent.Future;
import org.aopalliance.intercept.MethodInterceptor;
import org.aopalliance.intercept.MethodInvocation;
+import org.reactivestreams.Publisher;
import org.springframework.aop.framework.ProxyFactory;
import org.springframework.aop.support.AopUtils;
@@ -65,8 +66,6 @@ import org.springframework.util.ObjectUtils;
import org.springframework.util.ReflectionUtils;
import org.springframework.util.StringUtils;
-import reactor.Environment;
-import reactor.rx.Promise;
import reactor.rx.Promises;
/**
@@ -288,19 +287,6 @@ public class GatewayProxyFactoryBean extends AbstractEndpoint
this.globalMethodMetadata = globalMethodMetadata;
}
- /**
- * Set the Reactor {@link Environment} to be used for processing methods with a
- * {@link Promise} return type. (Required when any such methods are declared on the
- * service interface).
- * @param reactorEnvironment the Reactor Environment.
- * @since 4.1
- * @deprecated with no-op in favor of global JVM-wide Reactor configuration.
- */
- @Deprecated
- public void setReactorEnvironment(Object reactorEnvironment) {
-
- }
-
@Override
public void setBeanClassLoader(ClassLoader beanClassLoader) {
this.beanClassLoader = beanClassLoader;
@@ -413,9 +399,8 @@ public class GatewayProxyFactoryBean extends AbstractEndpoint
}
}
}
- if (reactorPresent && Promise.class.isAssignableFrom(returnType)) {
- return Promises.