diff --git a/pom.xml b/pom.xml
index 43ca523..a7c92dd 100644
--- a/pom.xml
+++ b/pom.xml
@@ -11,12 +11,12 @@
org.springframework.cloud
spring-cloud-build
- 2.0.0.BUILD-SNAPSHOT
+ 2.0.1.BUILD-SNAPSHOT
1.8
- Bismuth-RELEASE
- 0.9-SNAPSHOT
+ Californium-BUILD-SNAPSHOT
+ 0.10.0.BUILD-SNAPSHOT
0.7.9
@@ -34,12 +34,12 @@
io.rsocket
rsocket-core
- 0.9-SNAPSHOT
+ ${rsocket.version}
io.rsocket
rsocket-transport-netty
- 0.9-SNAPSHOT
+ ${rsocket.version}
diff --git a/spring-cloud-sockets/pom.xml b/spring-cloud-sockets/pom.xml
index 2499759..c78048b 100644
--- a/spring-cloud-sockets/pom.xml
+++ b/spring-cloud-sockets/pom.xml
@@ -15,6 +15,10 @@
io.projectreactor
reactor-core
+
+ io.projectreactor.ipc
+ reactor-netty
+
org.springframework
spring-context
@@ -31,6 +35,16 @@
com.fasterxml.jackson.core
jackson-databind
+
+ org.springframework.boot
+ spring-boot-starter-logging
+ provided
+
+
+ io.projectreactor
+ reactor-test
+ test
+
org.springframework.boot
spring-boot-starter-test
diff --git a/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/DispatcherHandler.java b/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/DispatcherHandler.java
index 69662cb..7441fe4 100644
--- a/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/DispatcherHandler.java
+++ b/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/DispatcherHandler.java
@@ -27,8 +27,8 @@
import com.fasterxml.jackson.databind.ObjectMapper;
import io.rsocket.AbstractRSocket;
import io.rsocket.Payload;
-import io.rsocket.exceptions.ApplicationException;
-import io.rsocket.util.PayloadImpl;
+import io.rsocket.exceptions.ApplicationErrorException;
+import io.rsocket.util.DefaultPayload;
import org.reactivestreams.Publisher;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -82,6 +82,7 @@ public void afterPropertiesSet() throws Exception {
String[] beanNames = BeanFactoryUtils.beanNamesForTypeIncludingAncestors(this.applicationContext, Object.class);
for(String beanName : beanNames){
Class> beanType = this.applicationContext.getType(beanName);
+
if(beanType != null){
final Class> userType = ClassUtils.getUserClass(beanType);
ReflectionUtils.doWithMethods(userType, method -> {
@@ -131,10 +132,10 @@ public Mono requestResponse(Payload payload) {
Converter converter = converterFor(MimeType.valueOf(metadata.get("MIME_TYPE").textValue()));
Object converted = converter.read(ServiceUtils.toByteArray(payload.getData()), getActualType(handler.getInfo().getParameterType()));
Object result = handler.invoke(handler.getInfo().buildInvocationArguments(converted, null));
- Mono monoResult = monoOF(result);
+ Mono> monoResult = monoOF(result);
return monoResult.map(o -> {
byte[] data = converter.write(o);
- return new PayloadImpl(data);
+ return DefaultPayload.create(data);
});
}catch (Exception e){
@@ -149,20 +150,20 @@ public Flux requestStream(Payload payload) {
MethodHandler handler = handlerFor(metadata);
Converter converter = converterFor(MimeType.valueOf(metadata.get("MIME_TYPE").textValue()));
Object converted = converter.read(ServiceUtils.toByteArray(payload.getData()), getActualType(handler.getInfo().getParameterType()));
- Flux result = (Flux)handler.invoke(handler.getInfo().buildInvocationArguments(converted, null));
+ Flux> result = (Flux>)handler.invoke(handler.getInfo().buildInvocationArguments(converted, null));
return result.map(o ->
- new PayloadImpl(converter.write(o))
+ DefaultPayload.create(converter.write(o))
);
} catch (Exception e){
- return Flux.error(new ApplicationException("No path found for " + metadata.get("PATH").asText()));
+ return Flux.error(new ApplicationErrorException("No path found for " + metadata.get("PATH").asText()));
}
}
- private Mono monoOF(Object argument){
+ private Mono> monoOF(Object argument){
if(argument.getClass().isAssignableFrom(Mono.class)){
- return (Mono)argument;
+ return (Mono>)argument;
}else{
return Mono.just(argument);
}
@@ -180,9 +181,9 @@ public Flux requestChannel(Publisher payloads) {
Flux converted = flux.repeat().map(payload -> {
return converter.read(ServiceUtils.toByteArray(payload.getData()), getActualType( handler.getInfo().getParameterType()));
});
- Flux result = (Flux)handler.invoke(handler.getInfo().buildInvocationArguments(converted, null));
+ Flux> result = (Flux>)handler.invoke(handler.getInfo().buildInvocationArguments(converted, null));
return result.map(o ->
- new PayloadImpl(converter.write(o))
+ DefaultPayload.create(converter.write(o))
);
}catch (Exception e){
return Flux.error(e);
@@ -207,7 +208,7 @@ private MethodHandler handlerFor(JsonNode metadata){
.getMappingInfo()
.getPath().equals(metadata.get("PATH").asText()); })
.findFirst()
- .orElseThrow(() -> { return new ApplicationException("No handler found");} );
+ .orElseThrow(() -> { return new ApplicationErrorException("No handler found");} );
}
}
diff --git a/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/OneWayRemoteHandler.java b/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/OneWayRemoteHandler.java
index d49f5b5..0a2f488 100644
--- a/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/OneWayRemoteHandler.java
+++ b/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/OneWayRemoteHandler.java
@@ -20,7 +20,7 @@
import java.nio.ByteBuffer;
import io.rsocket.RSocket;
-import io.rsocket.util.PayloadImpl;
+import io.rsocket.util.DefaultPayload;
import org.springframework.cloud.reactive.socket.ServiceMethodInfo;
@@ -37,6 +37,9 @@ public OneWayRemoteHandler(RSocket socket, ServiceMethodInfo info) {
@Override
public Object doInvoke(Object argument) {
byte[] payload = payloadConverter.write(argument);
- return socket.fireAndForget(new PayloadImpl(ByteBuffer.wrap(payload), getMetadata()));
+ return socket.fireAndForget(DefaultPayload.create(
+ ByteBuffer.wrap(payload),
+ getMetadata()
+ ));
}
}
diff --git a/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/RequestManyRemoteHandler.java b/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/RequestManyRemoteHandler.java
index 13de88a..2fcf390 100644
--- a/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/RequestManyRemoteHandler.java
+++ b/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/RequestManyRemoteHandler.java
@@ -17,11 +17,10 @@
package org.springframework.cloud.reactive.socket.client;
-
import java.nio.ByteBuffer;
import io.rsocket.RSocket;
-import io.rsocket.util.PayloadImpl;
+import io.rsocket.util.DefaultPayload;
import org.springframework.cloud.reactive.socket.ServiceMethodInfo;
import org.springframework.cloud.reactive.socket.util.ServiceUtils;
@@ -37,7 +36,7 @@ public RequestManyRemoteHandler(RSocket socket, ServiceMethodInfo info) {
@Override
public Object doInvoke(Object argument) {
byte[] data = payloadConverter.write(argument);
- return socket.requestStream(new PayloadImpl(ByteBuffer.wrap(data), getMetadata()))
+ return socket.requestStream(DefaultPayload.create(ByteBuffer.wrap(data), getMetadata()))
.map(payload -> payloadConverter.read(ServiceUtils.toByteArray(payload.getData()), ServiceUtils.getActualType(info.getParameterType())));
}
}
diff --git a/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/RequestOneRemoteHandler.java b/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/RequestOneRemoteHandler.java
index 463d069..515c0f8 100644
--- a/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/RequestOneRemoteHandler.java
+++ b/spring-cloud-sockets/src/main/java/org/springframework/cloud/reactive/socket/client/RequestOneRemoteHandler.java
@@ -20,7 +20,7 @@
import java.nio.ByteBuffer;
import io.rsocket.RSocket;
-import io.rsocket.util.PayloadImpl;
+import io.rsocket.util.DefaultPayload;
import reactor.core.publisher.Mono;
import org.springframework.cloud.reactive.socket.ServiceMethodInfo;
@@ -38,7 +38,10 @@ public RequestOneRemoteHandler(RSocket socket, ServiceMethodInfo info) {
@Override
public Object doInvoke(Object argument) {
byte[] data = payloadConverter.write(argument);
- Mono monoResult = socket.requestResponse(new PayloadImpl(ByteBuffer.wrap(data), getMetadata()))
+ Mono monoResult = socket.requestResponse(DefaultPayload.create(
+ ByteBuffer.wrap(data),
+ getMetadata()
+ ))
.map(payload -> payloadConverter.read(ServiceUtils.toByteArray(payload.getData()), ServiceUtils.getActualType(info.getReturnType())));
if(Mono.class.isAssignableFrom(info.getReturnType().resolve())){
return monoResult;
diff --git a/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/DispatchHandlerTests.java b/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/DispatchHandlerTests.java
index 8cee144..e6e2866 100644
--- a/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/DispatchHandlerTests.java
+++ b/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/DispatchHandlerTests.java
@@ -23,7 +23,7 @@
import java.util.Map;
import java.util.concurrent.ArrayBlockingQueue;
-import io.rsocket.util.PayloadImpl;
+import io.rsocket.util.DefaultPayload;
import org.junit.Before;
import org.junit.Test;
import reactor.core.publisher.Flux;
@@ -65,7 +65,7 @@ public void setup() throws Exception{
@Test
public void oneWayHandler() throws Exception{
User user = new User("Mary", "blue");
- Mono result = this.handler.fireAndForget(new PayloadImpl(converter.write(user), getMetadataBytes(MimeType.valueOf("application/json") ,"/oneway")));
+ Mono result = this.handler.fireAndForget(DefaultPayload.create(converter.write(user), getMetadataBytes(MimeType.valueOf("application/json") ,"/oneway")));
User output = (User) resultsQueue.poll();
assertThat(output).isEqualTo(user);
}
@@ -74,7 +74,7 @@ public void oneWayHandler() throws Exception{
public void oneWaySerializable() throws Exception {
User user = new User("Mary", "blue");
SerializableConverter serializableConverter = new SerializableConverter();
- Mono result = this.handler.fireAndForget(new PayloadImpl(serializableConverter.write(user), getMetadataBytes(MimeType.valueOf("application/java-serialized-object") ,"/onewaybinary")));
+ Mono result = this.handler.fireAndForget(DefaultPayload.create(serializableConverter.write(user), getMetadataBytes(MimeType.valueOf("application/java-serialized-object") ,"/onewaybinary")));
User output = (User) resultsQueue.poll();
assertThat(output).isEqualTo(user);
}
@@ -82,7 +82,7 @@ public void oneWaySerializable() throws Exception {
@Test
public void oneWayWrongMimeType() throws Exception {
User user = new User("Mary", "blue");
- Mono result = this.handler.fireAndForget(new PayloadImpl(converter.write(user), getMetadataBytes(MimeType.valueOf("application/binary") ,"/oneway")));
+ Mono result = this.handler.fireAndForget(DefaultPayload.create(converter.write(user), getMetadataBytes(MimeType.valueOf("application/binary") ,"/oneway")));
result.doOnError(throwable -> resultsQueue.offer(throwable)).subscribe();
assertThat(resultsQueue.poll()).isInstanceOf(Throwable.class);
}
@@ -90,7 +90,7 @@ public void oneWayWrongMimeType() throws Exception {
@Test
public void requestOneHandler() throws Exception {
User user = new User("Mary", "red");
- Mono invocationResult = this.handler.requestResponse(new PayloadImpl(converter.write(user), getMetadataBytes(MimeType.valueOf("application/json") ,"/redblue")));
+ Mono invocationResult = this.handler.requestResponse(DefaultPayload.create(converter.write(user), getMetadataBytes(MimeType.valueOf("application/json") ,"/redblue")));
User result = invocationResult.map(payload -> {
return (User)converter.read(payload.getDataUtf8().getBytes(), User.class);
}).block();
@@ -102,7 +102,7 @@ public void requestOneHandler() throws Exception {
@Test
public void requestOneWrongPath() throws Exception {
User user = new User("Mary", "red");
- Mono invocationResult = this.handler.requestResponse(new PayloadImpl(converter.write(user), getMetadataBytes(MimeType.valueOf("application/json") ,"/notfound")));
+ Mono invocationResult = this.handler.requestResponse(DefaultPayload.create(converter.write(user), getMetadataBytes(MimeType.valueOf("application/json") ,"/notfound")));
invocationResult.doOnError(throwable -> { resultsQueue.offer(throwable);}).subscribe();
assertThat(resultsQueue.poll()).isInstanceOf(Throwable.class);
}
@@ -111,7 +111,7 @@ public void requestOneWrongPath() throws Exception {
@Test
public void requestMany() throws Exception {
Integer count = 10;
- Flux invocationResult = this.handler.requestStream(new PayloadImpl(converter.write(count), getMetadataBytes(MimeType.valueOf("application/json") ,"/requestMany")));
+ Flux invocationResult = this.handler.requestStream(DefaultPayload.create(converter.write(count), getMetadataBytes(MimeType.valueOf("application/json") ,"/requestMany")));
List results = invocationResult.map(payload -> (Integer)converter.read(payload.getDataUtf8().getBytes(), Integer.class)
).collectList().block();
@@ -122,7 +122,7 @@ public void requestMany() throws Exception {
public void requestStream() throws Exception {
Flux from = Flux.range(0,10);
Flux payloadFlux = from.map(integer -> {
- return new PayloadImpl(converter.write(integer), getMetadataBytes(MimeType.valueOf("application/json") ,"/requestStream"));
+ return DefaultPayload.create(converter.write(integer), getMetadataBytes(MimeType.valueOf("application/json") ,"/requestStream"));
});
Flux invocationResult = this.handler.requestChannel(payloadFlux);
diff --git a/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/ReactiveSocketClientTests.java b/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/ReactiveSocketClientTests.java
index 80ac52d..312e99c 100644
--- a/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/ReactiveSocketClientTests.java
+++ b/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/ReactiveSocketClientTests.java
@@ -21,7 +21,7 @@
import com.fasterxml.jackson.databind.JsonNode;
import io.rsocket.Payload;
import io.rsocket.RSocket;
-import io.rsocket.util.PayloadImpl;
+import io.rsocket.util.DefaultPayload;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.ArgumentCaptor;
@@ -81,7 +81,7 @@ public void requestOneClientTests() throws Exception {
ArgumentCaptor captor = ArgumentCaptor.forClass(Payload.class);
User user = new User("Alice","blue");
byte[] converted = converter.write(user);
- when(mockSocket.requestResponse(Mockito.any(Payload.class))).thenReturn(Mono.just(new PayloadImpl(converted)));
+ when(mockSocket.requestResponse(Mockito.any(Payload.class))).thenReturn(Mono.just(DefaultPayload.create(converted)));
client.create(user);
verify(mockSocket, times(1)).requestResponse(captor.capture());
Payload payload = captor.getValue();
diff --git a/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/ReactiveSocketsApplicationTests.java b/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/ReactiveSocketsApplicationTests.java
index 214387e..e110532 100644
--- a/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/ReactiveSocketsApplicationTests.java
+++ b/spring-cloud-sockets/src/test/java/org/springframework/cloud/reactive/socket/ReactiveSocketsApplicationTests.java
@@ -17,10 +17,74 @@
package org.springframework.cloud.reactive.socket;
+import java.time.Duration;
+
+import io.rsocket.RSocket;
+import io.rsocket.RSocketFactory;
+import io.rsocket.transport.netty.client.TcpClientTransport;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import reactor.core.publisher.Flux;
+import reactor.test.StepVerifier;
+
+import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.boot.test.context.SpringBootTest;
+import org.springframework.cloud.reactive.socket.annotation.EnableReactiveSockets;
+import org.springframework.cloud.reactive.socket.annotation.Payload;
+import org.springframework.cloud.reactive.socket.annotation.RequestManyMapping;
+import org.springframework.cloud.reactive.socket.client.ReactiveSocketClient;
+import org.springframework.test.context.junit4.SpringRunner;
+
/**
* @author Vinicius Carvalho
*/
+@RunWith(SpringRunner.class)
+@SpringBootTest
public class ReactiveSocketsApplicationTests {
+ @Test
+ public void integrityTest() {
+ RSocket rSocket = RSocketFactory.connect()
+ .transport(TcpClientTransport.create("localhost", 5000))
+ .start()
+ .block();
+
+ ReactiveSocketClient client = new ReactiveSocketClient(rSocket);
+ TestClient clientProxy = client.create(TestClient.class);
+
+ StepVerifier.create(Flux.merge(clientProxy.receiveStream1("a"),
+ clientProxy.receiveStream2("b")))
+ .expectSubscription()
+ .expectNext("a", "b")
+ .expectNextCount(5)
+ .thenCancel()
+ .verify();
+ }
+
+ @SpringBootApplication
+ @EnableReactiveSockets
+ public static class TestApplication {
+
+ @RequestManyMapping(value = "/stream1", mimeType = "application/json")
+ public Flux stream1(@Payload String a) {
+ return Flux.just(a)
+ .mergeWith(Flux.interval(Duration.ofMillis(100))
+ .map(i -> "1. Stream Message : [" + i + "]"));
+ }
+
+ @RequestManyMapping(value = "/stream2", mimeType = "application/json")
+ public Flux stream2(@Payload String b) {
+ return Flux.just(b)
+ .mergeWith(Flux.interval(Duration.ofMillis(500))
+ .map(i -> "2. Stream Message : [" + i + "]"));
+ }
+ }
+
+ public interface TestClient {
+ @RequestManyMapping(value = "/stream1", mimeType = "application/json")
+ Flux receiveStream1(String a);
+ @RequestManyMapping(value = "/stream1", mimeType = "application/json")
+ Flux receiveStream2(String b);
+ }
}