From 2c42c9181d69b0ce0219599ee501932975c1c509 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 22 Jul 2020 09:24:05 -0400 Subject: [PATCH] Direct access to requester from ClientRSConnector Since `RSocketRequester` is now lazy load on client side there is no need to wrap it into a `Mono`. * Change `getRSocketRequester()` to a plain getter If there is a requirement to force connect to the server for receiving requests from there, the `ClientRSocketConnector.connect()` should be used --- .../rsocket/ClientRSocketConnector.java | 22 +++++-------------- .../outbound/RSocketOutboundGateway.java | 21 ++++++++---------- ...RSocketInboundGatewayIntegrationTests.java | 2 +- 3 files changed, 16 insertions(+), 29 deletions(-) diff --git a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ClientRSocketConnector.java b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ClientRSocketConnector.java index 30a650f035..5e1c94773f 100644 --- a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ClientRSocketConnector.java +++ b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/ClientRSocketConnector.java @@ -29,8 +29,6 @@ import org.springframework.util.MimeType; import io.rsocket.transport.ClientTransport; import io.rsocket.transport.netty.client.TcpClientTransport; import io.rsocket.transport.netty.client.WebsocketClientTransport; -import reactor.core.Disposable; -import reactor.core.publisher.Mono; /** * A client {@link AbstractRSocketConnector} extension to the RSocket connection. @@ -58,7 +56,7 @@ public class ClientRSocketConnector extends AbstractRSocketConnector { private boolean autoConnect; - private Mono rsocketRequesterMono; + private RSocketRequester rsocketRequester; /** * Instantiate a connector based on the {@link TcpClientTransport}. @@ -175,7 +173,7 @@ public class ClientRSocketConnector extends AbstractRSocketConnector { public void afterPropertiesSet() { super.afterPropertiesSet(); - RSocketRequester rsocketRequester = RSocketRequester.builder() + this.rsocketRequester = RSocketRequester.builder() .dataMimeType(getDataMimeType()) .metadataMimeType(getMetadataMimeType()) .rsocketStrategies(getRSocketStrategies()) @@ -186,11 +184,6 @@ public class ClientRSocketConnector extends AbstractRSocketConnector { connector.acceptor(this.rSocketMessageHandler.responder())) .apply((builder) -> this.setupMetadata.forEach(builder::setupMetadata)) .transport(this.clientTransport); - - this.rsocketRequesterMono = - Mono.just(rsocketRequester) - .doOnSubscribe((sub) -> rsocketRequester.rsocketClient().source().subscribe()) - .cache(); } @Override @@ -207,21 +200,18 @@ public class ClientRSocketConnector extends AbstractRSocketConnector { @Override public void destroy() { - this.rsocketRequesterMono - .flatMap((requester) -> requester.rsocketClient().source()) - .doOnNext(Disposable::dispose) - .subscribe(); + this.rsocketRequester.rsocketClient().dispose(); } /** * Perform subscription into the RSocket server for incoming requests. */ public void connect() { - this.rsocketRequesterMono.subscribe(); + this.rsocketRequester.rsocketClient().source().subscribe(); } - public Mono getRSocketRequester() { - return this.rsocketRequesterMono; + public RSocketRequester getRSocketRequester() { + return this.rsocketRequester; } } diff --git a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGateway.java b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGateway.java index 2651bac319..57aee27193 100644 --- a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGateway.java +++ b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGateway.java @@ -89,7 +89,7 @@ public class RSocketOutboundGateway extends AbstractReplyProducingMessageHandler private EvaluationContext evaluationContext; @Nullable - private Mono rsocketRequesterMono; + private RSocketRequester rsocketRequester; /** * Instantiate based on the provided RSocket endpoint {@code route} @@ -207,27 +207,24 @@ public class RSocketOutboundGateway extends AbstractReplyProducingMessageHandler super.doInit(); this.evaluationContext = ExpressionUtils.createStandardEvaluationContext(getBeanFactory()); if (this.clientRSocketConnector != null) { - this.rsocketRequesterMono = this.clientRSocketConnector.getRSocketRequester(); + this.rsocketRequester = this.clientRSocketConnector.getRSocketRequester(); } } @Override protected Object handleRequestMessage(Message requestMessage) { - RSocketRequester rsocketRequester = requestMessage.getHeaders() - .get(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER, RSocketRequester.class); - Mono requesterMono; - if (rsocketRequester != null) { - requesterMono = Mono.just(rsocketRequester); - } - else { - requesterMono = this.rsocketRequesterMono; + RSocketRequester requester = + requestMessage.getHeaders() + .get(RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER, RSocketRequester.class); + if (requester == null) { + requester = this.rsocketRequester; } - Assert.notNull(requesterMono, + Assert.notNull(requester, () -> "The 'RSocketRequester' must be configured via 'ClientRSocketConnector' or provided in the '" + RSocketRequesterMethodArgumentResolver.RSOCKET_REQUESTER_HEADER + "' request message headers."); - return requesterMono + return Mono.just(requester) .map((rSocketRequester) -> createRequestSpec(rSocketRequester, requestMessage)) .map((requestSpec) -> prepareRetrieveSpec(requestSpec, requestMessage)) .flatMap((retrieveSpec) -> performRetrieve(retrieveSpec, requestMessage)); diff --git a/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/inbound/RSocketInboundGatewayIntegrationTests.java b/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/inbound/RSocketInboundGatewayIntegrationTests.java index fb32b10c4f..9b2c40ff2f 100644 --- a/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/inbound/RSocketInboundGatewayIntegrationTests.java +++ b/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/inbound/RSocketInboundGatewayIntegrationTests.java @@ -97,7 +97,7 @@ public class RSocketInboundGatewayIntegrationTests { } else { this.clientRsocketRequester = - this.clientRSocketConnector.getRSocketRequester().block(Duration.ofSeconds(10)); + this.clientRSocketConnector.getRSocketRequester(); } }