From 15322d5a62f001ae56860a3edc55e92e403e7f42 Mon Sep 17 00:00:00 2001 From: Spencer Gibb Date: Tue, 16 Apr 2019 13:19:29 -0400 Subject: [PATCH] Renames server package to core since boot now manages the server. --- .../GatewayRSocketAutoConfiguration.java | 4 +- .../{server => core}/GatewayExchange.java | 2 +- .../{server => core}/GatewayFilter.java | 2 +- .../{server => core}/GatewayFilterChain.java | 2 +- .../{server => core}/GatewayPredicate.java | 2 +- .../{server => core}/GatewayRSocket.java | 14 ++--- ...GatewayServerRSocketFactoryCustomizer.java | 2 +- .../PendingRequestRSocket.java | 8 +-- .../cloud/gateway/rsocket/route/Route.java | 4 +- .../cloud/gateway/rsocket/route/Routes.java | 2 +- .../socketacceptor/GatewaySocketAcceptor.java | 2 +- .../GatewayRSocketAutoConfigurationTests.java | 2 +- .../GatewayRSocketIntegrationTests.java | 6 +- .../{server => core}/GatewayRSocketTests.java | 2 +- .../GatewaySocketAcceptorTests.java | 2 +- .../gateway/rsocket/test/PingPongApp.java | 58 +++++++++---------- 16 files changed, 54 insertions(+), 60 deletions(-) rename spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/{server => core}/GatewayExchange.java (97%) rename spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/{server => core}/GatewayFilter.java (93%) rename spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/{server => core}/GatewayFilterChain.java (96%) rename spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/{server => core}/GatewayPredicate.java (93%) rename spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/{server => core}/GatewayRSocket.java (94%) rename spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/{server => core}/GatewayServerRSocketFactoryCustomizer.java (98%) rename spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/{server => core}/PendingRequestRSocket.java (93%) rename spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/{server => core}/GatewayRSocketIntegrationTests.java (94%) rename spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/{server => core}/GatewayRSocketTests.java (99%) diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfiguration.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfiguration.java index 21f856da..c7528de4 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfiguration.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfiguration.java @@ -28,8 +28,8 @@ import org.springframework.cloud.gateway.rsocket.registry.Registry; import org.springframework.cloud.gateway.rsocket.registry.RegistryRoutes; import org.springframework.cloud.gateway.rsocket.registry.RegistrySocketAcceptorFilter; import org.springframework.cloud.gateway.rsocket.route.Routes; -import org.springframework.cloud.gateway.rsocket.server.GatewayRSocket; -import org.springframework.cloud.gateway.rsocket.server.GatewayServerRSocketFactoryCustomizer; +import org.springframework.cloud.gateway.rsocket.core.GatewayRSocket; +import org.springframework.cloud.gateway.rsocket.core.GatewayServerRSocketFactoryCustomizer; import org.springframework.cloud.gateway.rsocket.socketacceptor.GatewaySocketAcceptor; import org.springframework.cloud.gateway.rsocket.socketacceptor.SocketAcceptorFilter; import org.springframework.cloud.gateway.rsocket.socketacceptor.SocketAcceptorPredicate; diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayExchange.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayExchange.java similarity index 97% rename from spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayExchange.java rename to spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayExchange.java index 15670bcd..c5a9a30f 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayExchange.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayExchange.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.gateway.rsocket.server; +package org.springframework.cloud.gateway.rsocket.core; import io.micrometer.core.instrument.Tags; import io.rsocket.Payload; diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayFilter.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayFilter.java similarity index 93% rename from spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayFilter.java rename to spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayFilter.java index 692a6cb3..7fb2f848 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayFilter.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayFilter.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.gateway.rsocket.server; +package org.springframework.cloud.gateway.rsocket.core; import org.springframework.cloud.gateway.rsocket.filter.RSocketFilter; diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayFilterChain.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayFilterChain.java similarity index 96% rename from spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayFilterChain.java rename to spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayFilterChain.java index 689ba2a9..7fdad76b 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayFilterChain.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayFilterChain.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.gateway.rsocket.server; +package org.springframework.cloud.gateway.rsocket.core; import java.util.List; diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayPredicate.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayPredicate.java similarity index 93% rename from spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayPredicate.java rename to spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayPredicate.java index bcf5fc6c..a01a4944 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayPredicate.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayPredicate.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.gateway.rsocket.server; +package org.springframework.cloud.gateway.rsocket.core; import org.springframework.cloud.gateway.rsocket.support.AsyncPredicate; diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayRSocket.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocket.java similarity index 94% rename from spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayRSocket.java rename to spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocket.java index 643d74df..6aa98137 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayRSocket.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocket.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.gateway.rsocket.server; +package org.springframework.cloud.gateway.rsocket.core; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Function; @@ -43,12 +43,12 @@ import org.springframework.cloud.gateway.rsocket.support.Metadata; import org.springframework.util.Assert; import org.springframework.util.StringUtils; -import static org.springframework.cloud.gateway.rsocket.server.GatewayExchange.ROUTE_ATTR; -import static org.springframework.cloud.gateway.rsocket.server.GatewayExchange.Type.FIRE_AND_FORGET; -import static org.springframework.cloud.gateway.rsocket.server.GatewayExchange.Type.REQUEST_CHANNEL; -import static org.springframework.cloud.gateway.rsocket.server.GatewayExchange.Type.REQUEST_RESPONSE; -import static org.springframework.cloud.gateway.rsocket.server.GatewayExchange.Type.REQUEST_STREAM; -import static org.springframework.cloud.gateway.rsocket.server.GatewayFilterChain.executeFilterChain; +import static org.springframework.cloud.gateway.rsocket.core.GatewayExchange.ROUTE_ATTR; +import static org.springframework.cloud.gateway.rsocket.core.GatewayExchange.Type.FIRE_AND_FORGET; +import static org.springframework.cloud.gateway.rsocket.core.GatewayExchange.Type.REQUEST_CHANNEL; +import static org.springframework.cloud.gateway.rsocket.core.GatewayExchange.Type.REQUEST_RESPONSE; +import static org.springframework.cloud.gateway.rsocket.core.GatewayExchange.Type.REQUEST_STREAM; +import static org.springframework.cloud.gateway.rsocket.core.GatewayFilterChain.executeFilterChain; /** * Acts as a proxy to other registered sockets. Creates a GatewayExchange and attempts to diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayServerRSocketFactoryCustomizer.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayServerRSocketFactoryCustomizer.java similarity index 98% rename from spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayServerRSocketFactoryCustomizer.java rename to spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayServerRSocketFactoryCustomizer.java index c627fb59..e79c8060 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/GatewayServerRSocketFactoryCustomizer.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/GatewayServerRSocketFactoryCustomizer.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.gateway.rsocket.server; +package org.springframework.cloud.gateway.rsocket.core; import java.util.Arrays; import java.util.List; diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/PendingRequestRSocket.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/PendingRequestRSocket.java similarity index 93% rename from spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/PendingRequestRSocket.java rename to spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/PendingRequestRSocket.java index e23976de..a8dda50b 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/server/PendingRequestRSocket.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/core/PendingRequestRSocket.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.gateway.rsocket.server; +package org.springframework.cloud.gateway.rsocket.core; import java.util.function.Consumer; import java.util.function.Function; @@ -38,9 +38,9 @@ import org.springframework.cloud.gateway.rsocket.registry.Registry.RegisteredEve import org.springframework.cloud.gateway.rsocket.route.Route; import org.springframework.cloud.gateway.rsocket.support.Metadata; -import static org.springframework.cloud.gateway.rsocket.server.GatewayExchange.ROUTE_ATTR; -import static org.springframework.cloud.gateway.rsocket.server.GatewayExchange.Type.REQUEST_STREAM; -import static org.springframework.cloud.gateway.rsocket.server.GatewayFilterChain.executeFilterChain; +import static org.springframework.cloud.gateway.rsocket.core.GatewayExchange.ROUTE_ATTR; +import static org.springframework.cloud.gateway.rsocket.core.GatewayExchange.Type.REQUEST_STREAM; +import static org.springframework.cloud.gateway.rsocket.core.GatewayFilterChain.executeFilterChain; public class PendingRequestRSocket extends AbstractRSocket implements ResponderRSocket, Consumer { diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/route/Route.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/route/Route.java index 3ace7355..c9a3d29e 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/route/Route.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/route/Route.java @@ -23,8 +23,8 @@ import java.util.Collections; import java.util.List; import java.util.Objects; -import org.springframework.cloud.gateway.rsocket.server.GatewayExchange; -import org.springframework.cloud.gateway.rsocket.server.GatewayFilter; +import org.springframework.cloud.gateway.rsocket.core.GatewayExchange; +import org.springframework.cloud.gateway.rsocket.core.GatewayFilter; import org.springframework.cloud.gateway.rsocket.support.AsyncPredicate; import org.springframework.cloud.gateway.rsocket.support.Metadata; import org.springframework.core.Ordered; diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/route/Routes.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/route/Routes.java index a54d8a17..47bdeb9e 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/route/Routes.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/route/Routes.java @@ -21,7 +21,7 @@ import org.apache.commons.logging.LogFactory; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; -import org.springframework.cloud.gateway.rsocket.server.GatewayExchange; +import org.springframework.cloud.gateway.rsocket.core.GatewayExchange; /** * @author Spencer Gibb diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/socketacceptor/GatewaySocketAcceptor.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/socketacceptor/GatewaySocketAcceptor.java index 1559bb71..8b9eeed7 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/socketacceptor/GatewaySocketAcceptor.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/socketacceptor/GatewaySocketAcceptor.java @@ -30,7 +30,7 @@ import reactor.core.publisher.Mono; import org.springframework.cloud.gateway.rsocket.autoconfigure.GatewayRSocketProperties; import org.springframework.cloud.gateway.rsocket.metrics.MicrometerResponderRSocket; -import org.springframework.cloud.gateway.rsocket.server.GatewayRSocket; +import org.springframework.cloud.gateway.rsocket.core.GatewayRSocket; import org.springframework.cloud.gateway.rsocket.support.Metadata; public class GatewaySocketAcceptor implements SocketAcceptor { diff --git a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfigurationTests.java b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfigurationTests.java index 82dfeab2..ea30fa9d 100644 --- a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfigurationTests.java +++ b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfigurationTests.java @@ -25,7 +25,7 @@ import org.springframework.boot.test.context.runner.ReactiveWebApplicationContex import org.springframework.cloud.gateway.rsocket.registry.Registry; import org.springframework.cloud.gateway.rsocket.registry.RegistryRoutes; import org.springframework.cloud.gateway.rsocket.registry.RegistrySocketAcceptorFilter; -import org.springframework.cloud.gateway.rsocket.server.GatewayServerRSocketFactoryCustomizer; +import org.springframework.cloud.gateway.rsocket.core.GatewayServerRSocketFactoryCustomizer; import org.springframework.cloud.gateway.rsocket.socketacceptor.GatewaySocketAcceptor; import org.springframework.cloud.gateway.rsocket.socketacceptor.SocketAcceptorPredicate; import org.springframework.cloud.gateway.rsocket.socketacceptor.SocketAcceptorPredicateFilter; diff --git a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/server/GatewayRSocketIntegrationTests.java b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocketIntegrationTests.java similarity index 94% rename from spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/server/GatewayRSocketIntegrationTests.java rename to spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocketIntegrationTests.java index 4187071e..526a2039 100644 --- a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/server/GatewayRSocketIntegrationTests.java +++ b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocketIntegrationTests.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.gateway.rsocket.server; +package org.springframework.cloud.gateway.rsocket.core; import java.time.Duration; @@ -26,7 +26,7 @@ import reactor.test.StepVerifier; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.rsocket.RSocketProperties; -import org.springframework.boot.rsocket.netty.NettyRSocketBootstrap; +import org.springframework.boot.rsocket.server.RSocketServerBootstrap; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.context.SpringBootTest.WebEnvironment; import org.springframework.cloud.gateway.rsocket.test.PingPongApp; @@ -56,7 +56,7 @@ public class GatewayRSocketIntegrationTests { private PingPongApp.MySocketAcceptorFilter mySocketAcceptorFilter; @Autowired - private NettyRSocketBootstrap server; + private RSocketServerBootstrap server; @BeforeClass public static void init() { diff --git a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/server/GatewayRSocketTests.java b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocketTests.java similarity index 99% rename from spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/server/GatewayRSocketTests.java rename to spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocketTests.java index 50d7ce7c..c7c0fee5 100644 --- a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/server/GatewayRSocketTests.java +++ b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocketTests.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.gateway.rsocket.server; +package org.springframework.cloud.gateway.rsocket.core; import java.time.Duration; import java.util.Arrays; diff --git a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/socketacceptor/GatewaySocketAcceptorTests.java b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/socketacceptor/GatewaySocketAcceptorTests.java index f8b769ae..2bbd9ff5 100644 --- a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/socketacceptor/GatewaySocketAcceptorTests.java +++ b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/socketacceptor/GatewaySocketAcceptorTests.java @@ -31,7 +31,7 @@ import org.junit.Test; import reactor.core.publisher.Mono; import org.springframework.cloud.gateway.rsocket.autoconfigure.GatewayRSocketProperties; -import org.springframework.cloud.gateway.rsocket.server.GatewayRSocket; +import org.springframework.cloud.gateway.rsocket.core.GatewayRSocket; import org.springframework.cloud.gateway.rsocket.support.Metadata; import static java.util.Collections.singletonList; diff --git a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/test/PingPongApp.java b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/test/PingPongApp.java index e669b0fa..5c7537f7 100644 --- a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/test/PingPongApp.java +++ b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/test/PingPongApp.java @@ -42,9 +42,9 @@ import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.event.ApplicationReadyEvent; -import org.springframework.cloud.gateway.rsocket.server.GatewayExchange; -import org.springframework.cloud.gateway.rsocket.server.GatewayFilter; -import org.springframework.cloud.gateway.rsocket.server.GatewayFilterChain; +import org.springframework.cloud.gateway.rsocket.core.GatewayExchange; +import org.springframework.cloud.gateway.rsocket.core.GatewayFilter; +import org.springframework.cloud.gateway.rsocket.core.GatewayFilterChain; import org.springframework.cloud.gateway.rsocket.socketacceptor.SocketAcceptorExchange; import org.springframework.cloud.gateway.rsocket.socketacceptor.SocketAcceptorFilter; import org.springframework.cloud.gateway.rsocket.socketacceptor.SocketAcceptorFilterChain; @@ -140,39 +140,33 @@ public class PingPongApp { DefaultPayload.create(EMPTY_BUFFER, announcementMetadata)) .addClientPlugin(interceptor) .transport(TcpClientTransport.create(gatewayPort)) // proxy - .start().flatMapMany(socket -> { - Flux pong = socket.requestChannel( - Flux.interval(Duration.ofSeconds(1)).map(i -> { - ByteBuf data = ByteBufUtil.writeUtf8( - ByteBufAllocator.DEFAULT, "ping" + id); - ByteBuf routingMetadata = Metadata.from("pong") - .encode(); - return DefaultPayload.create(data, routingMetadata); - }).onBackpressureDrop(payload -> log.debug( - "Dropped payload " + payload.getDataUtf8())) // this - // is - // needed - // in - // case - // pong - // is - // not - // available - // yet - ).map(Payload::getDataUtf8).doOnNext(str -> { - int received = pongsReceived.incrementAndGet(); - log.info("received " + str + "(" + received + ") in Ping" - + id); - }).doFinally(signal -> socket.dispose()); - if (take != null) { - return pong.take(take); - } - return pong; - }); + .start().flatMapMany(socket -> doPing(take, socket)); pongFlux.subscribe(); } + Publisher doPing(Integer take, RSocket socket) { + Flux pong = socket.requestChannel( + Flux.interval(Duration.ofSeconds(1)).map(i -> { + ByteBuf data = ByteBufUtil.writeUtf8( + ByteBufAllocator.DEFAULT, "ping" + id); + ByteBuf routingMetadata = Metadata.from("pong") + .encode(); + return DefaultPayload.create(data, routingMetadata); + // onBackpressue is needed in case pong is not available yet + }).onBackpressureDrop(payload -> log.debug( + "Dropped payload " + payload.getDataUtf8())) + ).map(Payload::getDataUtf8).doOnNext(str -> { + int received = pongsReceived.incrementAndGet(); + log.info("received " + str + "(" + received + ") in Ping" + + id); + }).doFinally(signal -> socket.dispose()); + if (take != null) { + return pong.take(take); + } + return pong; + } + public Flux getPongFlux() { return pongFlux; }