From cc87c1b1b0849ee273ba52eeb1d75d8ed5ea3d76 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 27 Jun 2019 16:22:42 -0400 Subject: [PATCH] Optimize ServerRSocketConnector connection When RSocket client connects to the server there is no reason to wrap a `ConnectionSetupPayload` into a `Message` since we are not going to send it downstream * Refactor `IntegrationRSocket.handleConnectionSetupPayload()` just return a `Mono` for converted `ConnectionSetupPayload` * Ask for a `destination` and `RSocketRequester` from the `IntegrationRSocket` instead of message headers --- .../integration/rsocket/IntegrationRSocket.java | 10 ++++++---- .../rsocket/ServerRSocketConnector.java | 17 +++-------------- 2 files changed, 9 insertions(+), 18 deletions(-) diff --git a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/IntegrationRSocket.java b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/IntegrationRSocket.java index c1af1a0f46..308270a8a7 100644 --- a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/IntegrationRSocket.java +++ b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/IntegrationRSocket.java @@ -110,18 +110,20 @@ class IntegrationRSocket extends AbstractRSocket { this.bufferFactory = bufferFactory; } + RSocketRequester getRequester() { + return this.requester; + } + /** * Wrap the {@link ConnectionSetupPayload} with a {@link Message} and * delegate to {@link #handle(Payload)} for handling. * @param payload the connection payload * @return completion handle for success or error */ - Mono> handleConnectionSetupPayload(ConnectionSetupPayload payload) { - String destination = getDestination(payload); - MessageHeaders headers = createHeaders(destination, null); + Mono handleConnectionSetupPayload(ConnectionSetupPayload payload) { DataBuffer dataBuffer = retainDataAndReleasePayload(payload); int refCount = refCount(dataBuffer); - return Mono.just(MessageBuilder.createMessage(dataBuffer, headers)) + return Mono.just(dataBuffer) .doFinally(s -> { if (refCount(dataBuffer) == refCount) { DataBufferUtils.release(dataBuffer); diff --git a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ServerRSocketConnector.java b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ServerRSocketConnector.java index a5cf00f6c4..31e6963a48 100644 --- a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ServerRSocketConnector.java +++ b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ServerRSocketConnector.java @@ -29,12 +29,8 @@ import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisherAware; import org.springframework.core.io.buffer.DataBuffer; import org.springframework.lang.Nullable; -import org.springframework.messaging.MessageHeaders; -import org.springframework.messaging.handler.DestinationPatternsMessageCondition; import org.springframework.messaging.rsocket.RSocketRequester; -import org.springframework.messaging.rsocket.annotation.support.RSocketRequesterMethodArgumentResolver; import org.springframework.util.Assert; -import org.springframework.util.RouteMatcher; import io.rsocket.RSocketFactory; import io.rsocket.SocketAcceptor; @@ -181,17 +177,10 @@ public class ServerRSocketConnector extends AbstractRSocketConnector return (setupPayload, sendingRSocket) -> { IntegrationRSocket rsocket = createRSocket(setupPayload, sendingRSocket); return rsocket.handleConnectionSetupPayload(setupPayload) - .doOnNext((message) -> { - MessageHeaders messageHeaders = message.getHeaders(); - DataBuffer dataBuffer = message.getPayload(); - String destination = - messageHeaders.get(DestinationPatternsMessageCondition.LOOKUP_DESTINATION_HEADER, - RouteMatcher.Route.class) - .value(); + .doOnNext((dataBuffer) -> { + String destination = rsocket.getDestination(setupPayload); Object rsocketRequesterKey = this.clientRSocketKeyStrategy.apply(destination, dataBuffer); - RSocketRequester rsocketRequester = - messageHeaders.get(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER, - RSocketRequester.class); + RSocketRequester rsocketRequester = rsocket.getRequester(); this.clientRSocketRequesters.put(rsocketRequesterKey, rsocketRequester); RSocketConnectedEvent rSocketConnectedEvent = new RSocketConnectedEvent(rsocket, destination, dataBuffer, rsocketRequester);