From 33ec463f23d09b4588245338cfa2b3f1144b343a Mon Sep 17 00:00:00 2001 From: Spring Builds Date: Tue, 20 Sep 2022 13:18:32 +0000 Subject: [PATCH 1/3] Next development version (v1.0.3-SNAPSHOT) --- gradle.properties | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/gradle.properties b/gradle.properties index 110b552a..92e89c43 100644 --- a/gradle.properties +++ b/gradle.properties @@ -1,4 +1,4 @@ -version=1.0.2-SNAPSHOT +version=1.0.3-SNAPSHOT org.gradle.caching=true org.gradle.daemon=true From 0faa63beeaef0a722df28c0aeab57a49d30dd031 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Iv=C3=A1n=20Rodr=C3=ADguez=20Murillo?= Date: Mon, 10 Oct 2022 17:07:04 +0100 Subject: [PATCH 2/3] Allow use of loadbalanced RSocketRequester See gh-498 --- .../DefaultRSocketGraphQlClientBuilder.java | 38 +++++++++++++++++-- .../graphql/client/RSocketGraphQlClient.java | 18 +++++++++ 2 files changed, 52 insertions(+), 4 deletions(-) diff --git a/spring-graphql/src/main/java/org/springframework/graphql/client/DefaultRSocketGraphQlClientBuilder.java b/spring-graphql/src/main/java/org/springframework/graphql/client/DefaultRSocketGraphQlClientBuilder.java index ee652bbe..1d4933e4 100644 --- a/spring-graphql/src/main/java/org/springframework/graphql/client/DefaultRSocketGraphQlClientBuilder.java +++ b/spring-graphql/src/main/java/org/springframework/graphql/client/DefaultRSocketGraphQlClientBuilder.java @@ -17,11 +17,15 @@ package org.springframework.graphql.client; import java.net.URI; +import java.util.List; import java.util.function.Consumer; +import io.rsocket.loadbalance.LoadbalanceStrategy; +import io.rsocket.loadbalance.LoadbalanceTarget; import io.rsocket.transport.ClientTransport; import io.rsocket.transport.netty.client.TcpClientTransport; import io.rsocket.transport.netty.client.WebsocketClientTransport; +import org.reactivestreams.Publisher; import reactor.core.publisher.Mono; import org.springframework.lang.Nullable; @@ -45,6 +49,10 @@ final class DefaultRSocketGraphQlClientBuilder private final RSocketRequester.Builder requesterBuilder; + private Publisher> targetPublisher; + + private LoadbalanceStrategy loadbalanceStrategy; + @Nullable private ClientTransport clientTransport; @@ -111,6 +119,13 @@ final class DefaultRSocketGraphQlClientBuilder return this; } + @Override + public DefaultRSocketGraphQlClientBuilder transports(Publisher> targetPublisher, LoadbalanceStrategy loadbalanceStrategy) { + this.targetPublisher = targetPublisher; + this.loadbalanceStrategy = loadbalanceStrategy; + return this; + } + @Override public RSocketGraphQlClient build() { @@ -120,13 +135,20 @@ final class DefaultRSocketGraphQlClientBuilder builder.encoders(encoders -> setJsonEncoder(CodecDelegate.findJsonEncoder(encoders))); }); - Assert.state(this.clientTransport != null, "Neither WebSocket nor TCP networking configured"); - RSocketRequester requester = this.requesterBuilder.transport(this.clientTransport); + RSocketRequester requester; + + if (this.targetPublisher != null && this.loadbalanceStrategy != null) { + requester = this.requesterBuilder.transports(this.targetPublisher, this.loadbalanceStrategy); + } else { + Assert.state(this.clientTransport != null, "Neither WebSocket nor TCP networking configured"); + requester = this.requesterBuilder.transport(this.clientTransport); + } RSocketGraphQlTransport graphQlTransport = new RSocketGraphQlTransport(this.route, requester, getJsonDecoder()); return new DefaultRSocketGraphQlClient( super.buildGraphQlClient(graphQlTransport), requester, - this.requesterBuilder, this.clientTransport, this.route, getBuilderInitializer()); + this.requesterBuilder, this.clientTransport, this.targetPublisher, this.loadbalanceStrategy, + this.route, getBuilderInitializer()); } @@ -141,19 +163,26 @@ final class DefaultRSocketGraphQlClientBuilder private final ClientTransport clientTransport; + private final Publisher> targetPublisher; + + private final LoadbalanceStrategy loadbalanceStrategy; + private final String route; private final Consumer> builderInitializer; DefaultRSocketGraphQlClient( GraphQlClient graphQlClient, RSocketRequester requester, RSocketRequester.Builder requesterBuilder, - ClientTransport clientTransport, String route, Consumer> builderInitializer) { + ClientTransport clientTransport, Publisher> targetPublisher, LoadbalanceStrategy loadbalanceStrategy, + String route, Consumer> builderInitializer) { super(graphQlClient); this.requester = requester; this.requesterBuilder = requesterBuilder; this.clientTransport = clientTransport; + this.targetPublisher = targetPublisher; + this.loadbalanceStrategy = loadbalanceStrategy; this.route = route; this.builderInitializer = builderInitializer; } @@ -174,6 +203,7 @@ final class DefaultRSocketGraphQlClientBuilder public RSocketGraphQlClient.Builder mutate() { DefaultRSocketGraphQlClientBuilder builder = new DefaultRSocketGraphQlClientBuilder(this.requesterBuilder); builder.clientTransport(this.clientTransport); + builder.transports(this.targetPublisher, this.loadbalanceStrategy); builder.route(this.route); this.builderInitializer.accept(builder); return builder; diff --git a/spring-graphql/src/main/java/org/springframework/graphql/client/RSocketGraphQlClient.java b/spring-graphql/src/main/java/org/springframework/graphql/client/RSocketGraphQlClient.java index 448a57af..6c15c0b7 100644 --- a/spring-graphql/src/main/java/org/springframework/graphql/client/RSocketGraphQlClient.java +++ b/spring-graphql/src/main/java/org/springframework/graphql/client/RSocketGraphQlClient.java @@ -17,10 +17,14 @@ package org.springframework.graphql.client; import java.net.URI; +import java.util.List; import java.util.function.Consumer; import io.rsocket.core.RSocketClient; +import io.rsocket.loadbalance.LoadbalanceStrategy; +import io.rsocket.loadbalance.LoadbalanceTarget; import io.rsocket.transport.ClientTransport; +import org.reactivestreams.Publisher; import reactor.core.publisher.Mono; import org.springframework.messaging.rsocket.RSocketRequester; @@ -129,6 +133,20 @@ public interface RSocketGraphQlClient extends GraphQlClient { */ B rsocketRequester(Consumer requester); + /** + * Build an {@link RSocketRequester} with an + * {@link io.rsocket.loadbalance.LoadbalanceRSocketClient} that will + * connect to one of the given targets selected through the given + * {@link io.rsocket.loadbalance.LoadbalanceRSocketClient}. + * @param targetPublisher a {@code Publisher} that supplies a list of + * target transports to loadbalance against; the given list may be + * periodically updated by the {@code Publisher}. + * @param loadbalanceStrategy the strategy to use for selecting from + * the list of loadbalance targets. + * @return the same builder instance + */ + B transports(Publisher> targetPublisher, LoadbalanceStrategy loadbalanceStrategy); + /** * Build the {@code RSocketGraphQlClient} instance. */ From f3d07c1bd86fbf8fbc20a1a9144162edded1a514 Mon Sep 17 00:00:00 2001 From: rstoyanchev Date: Thu, 13 Oct 2022 17:02:08 +0100 Subject: [PATCH 3/3] Polishing contribution Closes gh-498 --- .../DefaultRSocketGraphQlClientBuilder.java | 60 ++++++++++++------- .../graphql/client/RSocketGraphQlClient.java | 41 +++++++------ 2 files changed, 63 insertions(+), 38 deletions(-) diff --git a/spring-graphql/src/main/java/org/springframework/graphql/client/DefaultRSocketGraphQlClientBuilder.java b/spring-graphql/src/main/java/org/springframework/graphql/client/DefaultRSocketGraphQlClientBuilder.java index 1d4933e4..0dc645d4 100644 --- a/spring-graphql/src/main/java/org/springframework/graphql/client/DefaultRSocketGraphQlClientBuilder.java +++ b/spring-graphql/src/main/java/org/springframework/graphql/client/DefaultRSocketGraphQlClientBuilder.java @@ -49,8 +49,10 @@ final class DefaultRSocketGraphQlClientBuilder private final RSocketRequester.Builder requesterBuilder; + @Nullable private Publisher> targetPublisher; + @Nullable private LoadbalanceStrategy loadbalanceStrategy; @Nullable @@ -95,8 +97,17 @@ final class DefaultRSocketGraphQlClientBuilder } @Override - public DefaultRSocketGraphQlClientBuilder clientTransport(ClientTransport clientTransport) { - this.clientTransport = clientTransport; + public DefaultRSocketGraphQlClientBuilder clientTransport(ClientTransport transport) { + this.clientTransport = transport; + return this; + } + + @Override + public DefaultRSocketGraphQlClientBuilder clientTransports( + Publisher> publisher, LoadbalanceStrategy strategy) { + + this.targetPublisher = publisher; + this.loadbalanceStrategy = strategy; return this; } @@ -114,15 +125,8 @@ final class DefaultRSocketGraphQlClientBuilder } @Override - public DefaultRSocketGraphQlClientBuilder rsocketRequester(Consumer requesterConsumer) { - requesterConsumer.accept(this.requesterBuilder); - return this; - } - - @Override - public DefaultRSocketGraphQlClientBuilder transports(Publisher> targetPublisher, LoadbalanceStrategy loadbalanceStrategy) { - this.targetPublisher = targetPublisher; - this.loadbalanceStrategy = loadbalanceStrategy; + public DefaultRSocketGraphQlClientBuilder rsocketRequester(Consumer consumer) { + consumer.accept(this.requesterBuilder); return this; } @@ -137,13 +141,18 @@ final class DefaultRSocketGraphQlClientBuilder RSocketRequester requester; - if (this.targetPublisher != null && this.loadbalanceStrategy != null) { - requester = this.requesterBuilder.transports(this.targetPublisher, this.loadbalanceStrategy); - } else { - Assert.state(this.clientTransport != null, "Neither WebSocket nor TCP networking configured"); + if (this.clientTransport != null) { requester = this.requesterBuilder.transport(this.clientTransport); } - RSocketGraphQlTransport graphQlTransport = new RSocketGraphQlTransport(this.route, requester, getJsonDecoder()); + else if (this.targetPublisher != null && this.loadbalanceStrategy != null) { + requester = this.requesterBuilder.transports(this.targetPublisher, this.loadbalanceStrategy); + } + else { + throw new IllegalStateException("Neither ClientTransport, nor Loadbalance targets and strategy"); + } + + RSocketGraphQlTransport graphQlTransport = + new RSocketGraphQlTransport(this.route, requester, getJsonDecoder()); return new DefaultRSocketGraphQlClient( super.buildGraphQlClient(graphQlTransport), requester, @@ -161,10 +170,13 @@ final class DefaultRSocketGraphQlClientBuilder private final RSocketRequester.Builder requesterBuilder; + @Nullable private final ClientTransport clientTransport; + @Nullable private final Publisher> targetPublisher; + @Nullable private final LoadbalanceStrategy loadbalanceStrategy; private final String route; @@ -172,8 +184,10 @@ final class DefaultRSocketGraphQlClientBuilder private final Consumer> builderInitializer; DefaultRSocketGraphQlClient( - GraphQlClient graphQlClient, RSocketRequester requester, RSocketRequester.Builder requesterBuilder, - ClientTransport clientTransport, Publisher> targetPublisher, LoadbalanceStrategy loadbalanceStrategy, + GraphQlClient graphQlClient, + RSocketRequester requester, RSocketRequester.Builder requesterBuilder, + @Nullable ClientTransport clientTransport, + @Nullable Publisher> targetPublisher, @Nullable LoadbalanceStrategy strategy, String route, Consumer> builderInitializer) { super(graphQlClient); @@ -182,7 +196,7 @@ final class DefaultRSocketGraphQlClientBuilder this.requesterBuilder = requesterBuilder; this.clientTransport = clientTransport; this.targetPublisher = targetPublisher; - this.loadbalanceStrategy = loadbalanceStrategy; + this.loadbalanceStrategy = strategy; this.route = route; this.builderInitializer = builderInitializer; } @@ -202,8 +216,12 @@ final class DefaultRSocketGraphQlClientBuilder @Override public RSocketGraphQlClient.Builder mutate() { DefaultRSocketGraphQlClientBuilder builder = new DefaultRSocketGraphQlClientBuilder(this.requesterBuilder); - builder.clientTransport(this.clientTransport); - builder.transports(this.targetPublisher, this.loadbalanceStrategy); + if (this.clientTransport != null) { + builder.clientTransport(this.clientTransport); + } + if (this.targetPublisher != null && this.loadbalanceStrategy != null) { + builder.clientTransports(this.targetPublisher, this.loadbalanceStrategy); + } builder.route(this.route); this.builderInitializer.accept(builder); return builder; diff --git a/spring-graphql/src/main/java/org/springframework/graphql/client/RSocketGraphQlClient.java b/spring-graphql/src/main/java/org/springframework/graphql/client/RSocketGraphQlClient.java index 6c15c0b7..133975d8 100644 --- a/spring-graphql/src/main/java/org/springframework/graphql/client/RSocketGraphQlClient.java +++ b/spring-graphql/src/main/java/org/springframework/graphql/client/RSocketGraphQlClient.java @@ -82,7 +82,9 @@ public interface RSocketGraphQlClient extends GraphQlClient { interface Builder> extends GraphQlClient.Builder { /** - * Select TCP as the underlying network protocol. + * Select TCP as the underlying network protocol. This delegates to + * {@link RSocketRequester.Builder#tcp(String, int)} to create the + * {@code RSocketRequester} instance. * @param host the remote host to connect to * @param port the remote port to connect to * @return the same builder instance @@ -90,19 +92,38 @@ public interface RSocketGraphQlClient extends GraphQlClient { B tcp(String host, int port); /** - * Select WebSocket as the underlying network protocol. + * Select WebSocket as the underlying network protocol. This delegates to + * {@link RSocketRequester.Builder#websocket(URI)} to create the + * {@code RSocketRequester} instance. * @param uri the URL for the WebSocket handshake * @return the same builder instance */ B webSocket(URI uri); /** - * Use a given {@link ClientTransport} to communicate with the remote server. + * Use a given {@link ClientTransport} to communicate with the remote + * server. This delegates to + * {@link RSocketRequester.Builder#transport(ClientTransport)} to create + * the {@code RSocketRequester} instance. * @param clientTransport the transport to use * @return the same builder instance */ B clientTransport(ClientTransport clientTransport); + /** + * Use a {@link Publisher} of {@link LoadbalanceTarget}s, each of which + * contains a {@link ClientTransport}. This delegates to + * {@link RSocketRequester.Builder#transports(Publisher, LoadbalanceStrategy)} + * to create the {@code RSocketRequester} instance. + * @param targetPublisher supplies list of targets to loadbalance against; + * the targets are replaced when the given {@code Publisher} emits again. + * @param loadbalanceStrategy the strategy to use for selecting from + * the list of targets. + * @return the same builder instance + * @since 1.0.3 + */ + B clientTransports(Publisher> targetPublisher, LoadbalanceStrategy loadbalanceStrategy); + /** * Customize the format of data payloads for the connection. *

By default, this is set to {@code "application/graphql+json"} but @@ -133,20 +154,6 @@ public interface RSocketGraphQlClient extends GraphQlClient { */ B rsocketRequester(Consumer requester); - /** - * Build an {@link RSocketRequester} with an - * {@link io.rsocket.loadbalance.LoadbalanceRSocketClient} that will - * connect to one of the given targets selected through the given - * {@link io.rsocket.loadbalance.LoadbalanceRSocketClient}. - * @param targetPublisher a {@code Publisher} that supplies a list of - * target transports to loadbalance against; the given list may be - * periodically updated by the {@code Publisher}. - * @param loadbalanceStrategy the strategy to use for selecting from - * the list of loadbalance targets. - * @return the same builder instance - */ - B transports(Publisher> targetPublisher, LoadbalanceStrategy loadbalanceStrategy); - /** * Build the {@code RSocketGraphQlClient} instance. */