From 444c1f9913d2650d4ffc4362cc4e2615f997eb4f Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 3 Sep 2019 14:22:36 -0400 Subject: [PATCH] RSocket: Add support for RoutingMetadata Related to https://github.com/spring-projects/spring-framework/issues/23137 The metadata in Spring Messaging for RSockets now supports any arbitrary objects for setup payload, including composition. * Switch the `ClientRSocketConnector` to fully delegate to the `RSocketRequester.Builder` inheriting possible metadata encoding/decoding in the target `RSocketRequester` implementation * Turn off a default `dataMimeType` from the `MimeTypeUtils.TEXT_PLAIN` to the `null` by default relying on the encoder/decoder logic in the target RSocket wrappers * Expose more delegating options in the `ClientRSocketConnector`, like `setupRouteVars`, `setupMetadata` --- build.gradle | 6 +- .../rsocket/AbstractRSocketConnector.java | 2 +- .../rsocket/ClientRSocketConnector.java | 123 ++++++++++++------ ...RSocketInboundGatewayIntegrationTests.java | 2 +- ...SocketOutboundGatewayIntegrationTests.java | 28 +--- 5 files changed, 92 insertions(+), 69 deletions(-) diff --git a/build.gradle b/build.gradle index 79be9f34de..493afa6e80 100644 --- a/build.gradle +++ b/build.gradle @@ -84,11 +84,11 @@ ext { mysqlVersion = '8.0.16' pahoMqttClientVersion = '1.2.0' postgresVersion = '42.2.6' - reactorNettyVersion = '0.9.0.M3' - reactorVersion = '3.3.0.M3' + reactorNettyVersion = '0.9.0.BUILD-SNAPSHOT' + reactorVersion = '3.3.0.BUILD-SNAPSHOT' resilience4jVersion = '0.16.0' romeToolsVersion = '1.12.1' - rsocketVersion = '1.0.0-RC3-SNAPSHOT' + rsocketVersion = '1.0.0-RC3' servletApiVersion = '4.0.1' smackVersion = '4.3.4' springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '2.2.0.BUILD-SNAPSHOT' diff --git a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/AbstractRSocketConnector.java b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/AbstractRSocketConnector.java index 4cd297fe3f..a1bad52276 100644 --- a/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/AbstractRSocketConnector.java +++ b/spring-integration-rsocket/src/main/java/org/springframework/integration/rsocket/AbstractRSocketConnector.java @@ -53,7 +53,7 @@ public abstract class AbstractRSocketConnector protected final IntegrationRSocketMessageHandler rSocketMessageHandler; // NOSONAR - final - private MimeType dataMimeType = MimeTypeUtils.TEXT_PLAIN; + private MimeType dataMimeType; private MimeType metadataMimeType = MimeTypeUtils.parseMimeType(WellKnownMimeType.MESSAGE_RSOCKET_COMPOSITE_METADATA.toString()); 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 d7104bcf73..6aa338231b 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 @@ -17,19 +17,18 @@ package org.springframework.integration.rsocket; import java.net.URI; -import java.util.function.Consumer; +import java.util.Arrays; +import java.util.LinkedHashMap; +import java.util.Map; +import org.springframework.messaging.rsocket.ClientRSocketFactoryConfigurer; import org.springframework.messaging.rsocket.RSocketRequester; import org.springframework.util.Assert; +import org.springframework.util.MimeType; -import io.rsocket.Payload; -import io.rsocket.RSocket; -import io.rsocket.RSocketFactory; import io.rsocket.transport.ClientTransport; import io.rsocket.transport.netty.client.TcpClientTransport; import io.rsocket.transport.netty.client.WebsocketClientTransport; -import io.rsocket.util.DefaultPayload; -import io.rsocket.util.EmptyPayload; import reactor.core.Disposable; import reactor.core.publisher.Mono; @@ -40,22 +39,26 @@ import reactor.core.publisher.Mono; * * @since 5.2 * - * @see RSocketFactory.ClientRSocketFactory + * @see io.rsocket.RSocketFactory.ClientRSocketFactory * @see RSocketRequester */ public class ClientRSocketConnector extends AbstractRSocketConnector { private final ClientTransport clientTransport; - private Consumer factoryConfigurer = (clientRSocketFactory) -> { }; + private final Map setupMetadata = new LinkedHashMap<>(4); - private String connectRoute; + private ClientRSocketFactoryConfigurer factoryConfigurer = (clientRSocketFactory) -> { }; - private String connectData = ""; + private Object setupData; + + private String setupRoute; + + private Object[] setupRouteVars = new Object[0]; private boolean autoConnect; - private Mono rsocketMono; + private Mono rsocketRequesterMono; /** * Instantiate a connector based on the {@link TcpClientTransport}. @@ -79,6 +82,7 @@ public class ClientRSocketConnector extends AbstractRSocketConnector { /** * Instantiate a connector based on the provided {@link ClientTransport}. * @param clientTransport the {@link ClientTransport} to use. + * @see RSocketRequester.Builder#connect(ClientTransport) */ public ClientRSocketConnector(ClientTransport clientTransport) { super(new IntegrationRSocketMessageHandler()); @@ -87,47 +91,83 @@ public class ClientRSocketConnector extends AbstractRSocketConnector { } /** - * Specify a {@link Consumer} for configuring a {@link RSocketFactory.ClientRSocketFactory}. - * @param factoryConfigurer the {@link Consumer} to configure the {@link RSocketFactory.ClientRSocketFactory}. + * Callback to configure the {@code ClientRSocketFactory} directly. + * Note: this class adds extra {@link ClientRSocketFactoryConfigurer} to the + * target {@link RSocketRequester} to populate a reference to an internal + * {@link IntegrationRSocketMessageHandler#responder()}. + * This overrides possible external + * {@link io.rsocket.RSocketFactory.ClientRSocketFactory#acceptor(io.rsocket.SocketAcceptor)} + * @param factoryConfigurer the {@link ClientRSocketFactoryConfigurer} to + * configure the {@link io.rsocket.RSocketFactory.ClientRSocketFactory}. + * @see RSocketRequester.Builder#rsocketFactory(ClientRSocketFactoryConfigurer) */ - public void setFactoryConfigurer(Consumer factoryConfigurer) { + public void setFactoryConfigurer(ClientRSocketFactoryConfigurer factoryConfigurer) { Assert.notNull(factoryConfigurer, "'factoryConfigurer' must not be null"); this.factoryConfigurer = factoryConfigurer; } /** - * Configure a route for server RSocket endpoint. - * @param connectRoute the route to connect to. + * Set the route for the setup payload. + * @param setupRoute the route to connect to. + * @see RSocketRequester.Builder#setupRoute(String, Object...) */ - public void setConnectRoute(String connectRoute) { - this.connectRoute = connectRoute; + public void setSetupRoute(String setupRoute) { + Assert.notNull(setupRoute, "'setupRoute' must not be null"); + this.setupRoute = setupRoute; } /** - * Configure a data for connect. - * Defaults to empty string. - * @param connectData the data for connect frame. + * Set the variables for route template to expand with. + * @param setupRouteVars the route to connect to. + * @see RSocketRequester.Builder#setupRoute(String, Object...) */ - public void setConnectData(String connectData) { - Assert.notNull(connectData, "'connectData' must not be null"); - this.connectData = connectData; + public void setSetupRouteVariables(Object... setupRouteVars) { + Assert.notNull(setupRouteVars, "'setupRouteVars' must not be null"); + this.setupRouteVars = Arrays.copyOf(setupRouteVars, setupRouteVars.length); + } + + /** + * Add metadata to the setup payload. Composite metadata must be + * in use if this is called more than once or in addition to + * {@link #setSetupRoute(String)}. + * @param setupMetadata the map of metadata to use. + * @see RSocketRequester.Builder#setupMetadata(Object, MimeType) + */ + public void setSetupMetadata(Map setupMetadata) { + Assert.notNull(setupMetadata, "'setupMetadata' must not be null"); + this.setupMetadata.clear(); + this.setupMetadata.putAll(setupMetadata); + } + + /** + * Set the data for the setup payload. + * @param setupData the data for connect frame. + * @see RSocketRequester.Builder#setupData(Object) + */ + public void setSetupData(Object setupData) { + Assert.notNull(setupData, "'setupData' must not be null"); + this.setupData = setupData; } @Override public void afterPropertiesSet() { super.afterPropertiesSet(); - RSocketFactory.ClientRSocketFactory clientFactory = - RSocketFactory.connect() - .dataMimeType(getDataMimeType().toString()) - .metadataMimeType(getMetadataMimeType().toString()); - this.factoryConfigurer.accept(clientFactory); - clientFactory.acceptor(this.rSocketMessageHandler.responder()); - Payload connectPayload = EmptyPayload.INSTANCE; - if (this.connectRoute != null) { - connectPayload = DefaultPayload.create(this.connectData, this.connectRoute); - } - clientFactory.setupPayload(connectPayload); - this.rsocketMono = clientFactory.transport(this.clientTransport).start().cache(); + + RSocketRequester.Builder rsocketRequesterBuilder = + RSocketRequester.builder() + .dataMimeType(getDataMimeType()) + .metadataMimeType(getMetadataMimeType()) + .rsocketStrategies(getRSocketStrategies()) + .setupData(this.setupData) + .setupRoute(this.setupRoute, this.setupRouteVars) + .rsocketFactory(this.factoryConfigurer) + .rsocketFactory((rsocketFactory) -> + rsocketFactory.acceptor(this.rSocketMessageHandler.responder())); + this.setupMetadata.forEach(rsocketRequesterBuilder::setupMetadata); + this.rsocketRequesterMono = + rsocketRequesterBuilder + .connect(this.clientTransport) + .cache(); } @Override @@ -144,7 +184,8 @@ public class ClientRSocketConnector extends AbstractRSocketConnector { @Override public void destroy() { - this.rsocketMono + this.rsocketRequesterMono + .map(RSocketRequester::rsocket) .doOnNext(Disposable::dispose) .subscribe(); } @@ -153,15 +194,11 @@ public class ClientRSocketConnector extends AbstractRSocketConnector { * Perform subscription into the RSocket server for incoming requests. */ public void connect() { - this.rsocketMono.subscribe(); + this.rsocketRequesterMono.subscribe(); } public Mono getRSocketRequester() { - return this.rsocketMono - .map((rsocket) -> - RSocketRequester - .wrap(rsocket, getDataMimeType(), getMetadataMimeType(), getRSocketStrategies())) - .cache(); + return this.rsocketRequesterMono; } } 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 d1c60a7040..42d04eea00 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 @@ -230,7 +230,7 @@ public class RSocketInboundGatewayIntegrationTests { clientRSocketConnector.setMetadataMimeType(new MimeType("message", "x.rsocket.routing.v0")); clientRSocketConnector.setFactoryConfigurer((factory) -> factory.frameDecoder(PayloadDecoder.ZERO_COPY)); clientRSocketConnector.setRSocketStrategies(rsocketStrategies()); - clientRSocketConnector.setConnectRoute("clientConnect"); + clientRSocketConnector.setSetupRoute("clientConnect"); return clientRSocketConnector; } diff --git a/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGatewayIntegrationTests.java b/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGatewayIntegrationTests.java index f39abf6a1a..cbe2b380f4 100644 --- a/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGatewayIntegrationTests.java +++ b/spring-integration-rsocket/src/test/java/org/springframework/integration/rsocket/outbound/RSocketOutboundGatewayIntegrationTests.java @@ -19,7 +19,6 @@ package org.springframework.integration.rsocket.outbound; import static org.assertj.core.api.Assertions.assertThat; import java.time.Duration; -import java.util.Collections; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; @@ -64,10 +63,8 @@ import io.netty.buffer.PooledByteBufAllocator; import io.rsocket.RSocket; import io.rsocket.RSocketFactory; import io.rsocket.frame.decoder.PayloadDecoder; -import io.rsocket.transport.netty.client.TcpClientTransport; import io.rsocket.transport.netty.server.CloseableChannel; import io.rsocket.transport.netty.server.TcpServerTransport; -import io.rsocket.util.DefaultPayload; import reactor.core.Disposable; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; @@ -515,33 +512,22 @@ public class RSocketOutboundGatewayIntegrationTests { @EnableIntegration public static class ClientConfig extends CommonConfig { - @Bean - public RSocketMessageHandler messageHandler() { - RSocketMessageHandler handler = new RSocketMessageHandler(); - handler.setRSocketStrategies(rsocketStrategies()); - handler.setHandlers(Collections.singletonList(controller())); - return handler; - } - @Bean(destroyMethod = "dispose") @Nullable public RSocket rsocketForServerRequests() { - return RSocketFactory.connect() - .setupPayload(DefaultPayload.create("", "clientConnect")) - .dataMimeType("text/plain") - .metadataMimeType("message/x.rsocket.routing.v0") - .frameDecoder(PayloadDecoder.ZERO_COPY) - .acceptor(messageHandler().responder()) - .transport(TcpClientTransport.create("localhost", server.address().getPort())) - .start() - .block(); + + return RSocketRequester.builder() + .setupRoute("clientConnect") + .rsocketFactory(RSocketMessageHandler.clientResponder(rsocketStrategies(), controller())) + .connectTcp("localhost", server.address().getPort()) + .block() + .rsocket(); } @Bean public ClientRSocketConnector clientRSocketConnector() { ClientRSocketConnector clientRSocketConnector = new ClientRSocketConnector("localhost", server.address().getPort()); - clientRSocketConnector.setFactoryConfigurer((factory) -> factory.frameDecoder(PayloadDecoder.ZERO_COPY)); clientRSocketConnector.setRSocketStrategies(rsocketStrategies()); return clientRSocketConnector; }