From 1c54a5abdf1dc846d8f0dc892040f260aab15664 Mon Sep 17 00:00:00 2001 From: Dmitry Kostyukov Date: Sun, 30 May 2021 01:49:38 +0300 Subject: [PATCH 1/3] Fix proxying websocket close status Fixes gh-1057 --- .../filter/WebsocketRoutingFilter.java | 6 +++- .../websocket/WebSocketIntegrationTests.java | 35 ++++++++++++++++--- 2 files changed, 35 insertions(+), 6 deletions(-) diff --git a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java index b3ff1eb3..d512d43d 100644 --- a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java +++ b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java @@ -203,6 +203,10 @@ public class WebsocketRoutingFilter implements GlobalFilter, Ordered { return client.execute(url, this.headers, new WebSocketHandler() { @Override public Mono handle(WebSocketSession proxySession) { + Mono serverClose = proxySession.closeStatus().filter(__ -> session.isOpen()) + .flatMap(session::close); + Mono proxyClose = session.closeStatus().filter(__ -> proxySession.isOpen()) + .flatMap(proxySession::close); // Use retain() for Reactor Netty Mono proxySessionSend = proxySession .send(session.receive().doOnNext(WebSocketMessage::retain)); @@ -210,7 +214,7 @@ public class WebsocketRoutingFilter implements GlobalFilter, Ordered { Mono serverSessionSend = session .send(proxySession.receive().doOnNext(WebSocketMessage::retain)); // .log("sessionSend", Level.FINE); - return Mono.zip(proxySessionSend, serverSessionSend).then(); + return Mono.zip(proxySessionSend, serverSessionSend, serverClose, proxyClose).then(); } /** diff --git a/spring-cloud-gateway-server/src/test/java/org/springframework/cloud/gateway/test/websocket/WebSocketIntegrationTests.java b/spring-cloud-gateway-server/src/test/java/org/springframework/cloud/gateway/test/websocket/WebSocketIntegrationTests.java index 43426a52..88b7bc17 100644 --- a/spring-cloud-gateway-server/src/test/java/org/springframework/cloud/gateway/test/websocket/WebSocketIntegrationTests.java +++ b/spring-cloud-gateway-server/src/test/java/org/springframework/cloud/gateway/test/websocket/WebSocketIntegrationTests.java @@ -36,6 +36,7 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.publisher.MonoProcessor; import reactor.core.publisher.ReplayProcessor; +import reactor.core.publisher.Sinks; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; @@ -101,6 +102,8 @@ public class WebSocketIntegrationTests { private int gatewayPort; + private static final Sinks.One serverCloseStatusSink = Sinks.one(); + private static Mono doSend(WebSocketSession session, Publisher output) { return session.send(output); // workaround for suspected RxNetty WebSocket client issue @@ -241,13 +244,25 @@ public class WebSocketIntegrationTests { } @Test - public void sessionClosing() throws Exception { - this.client.execute(getUrl("/close"), session -> { + public void serverClosing() throws Exception { + AtomicReference> closeStatus = new AtomicReference<>(); + this.client.execute(getUrl("/server-close"), session -> { logger.debug("Starting.."); + closeStatus.set(session.closeStatus()); return session.receive().doOnNext(s -> logger.debug("inbound " + s)).then().doFinally(signalType -> { logger.debug("Completed with: " + signalType); }); }).block(Duration.ofMillis(5000)); + assertThat(closeStatus.get().block(Duration.ofMillis(5000))) + .isEqualTo(CloseStatus.create(4999, "server-close")); + } + + @Test + public void clientClosing() throws Exception { + this.client.execute(getUrl("/client-close"), session -> session.close(CloseStatus.create(4999, "client-close"))) + .block(Duration.ofMillis(5000)); + assertThat(serverCloseStatusSink.asMono().block(Duration.ofMillis(5000))) + .isEqualTo(CloseStatus.create(4999, "client-close")); } @Configuration(proxyBeanMethods = false) @@ -279,7 +294,8 @@ public class WebSocketIntegrationTests { map.put("/echoForHttp", new EchoWebSocketHandler()); map.put("/sub-protocol", new SubProtocolWebSocketHandler()); map.put("/custom-header", new CustomHeaderHandler()); - map.put("/close", new SessionClosingHandler()); + map.put("/server-close", new ServerClosingHandler()); + map.put("/client-close", new ClientClosingHandler()); SimpleUrlHandlerMapping mapping = new SimpleUrlHandlerMapping(); mapping.setUrlMap(map); @@ -334,11 +350,20 @@ public class WebSocketIntegrationTests { } - private static class SessionClosingHandler implements WebSocketHandler { + private static class ServerClosingHandler implements WebSocketHandler { @Override public Mono handle(WebSocketSession session) { - return Flux.never().mergeWith(session.close(CloseStatus.GOING_AWAY)).then(); + return Flux.never().mergeWith(session.close(CloseStatus.create(4999, "server-close"))).then(); + } + + } + + private static class ClientClosingHandler implements WebSocketHandler { + + @Override + public Mono handle(WebSocketSession session) { + return session.closeStatus().doOnNext(serverCloseStatusSink::tryEmitValue).then(); } } From 48accb783f44533b6e9097fbfee21b1ea8a8cc7c Mon Sep 17 00:00:00 2001 From: spencergibb Date: Mon, 1 Nov 2021 15:39:34 -0400 Subject: [PATCH 2/3] Updates return value for Websocket return. From @rstoyanchev: I think serverClose and proxyClose don't need to be included in, and probably should be separated from the zip because they are sort of competing with the Mono from each WebSocketHandler, and zip will cancel the other publishers after one of them completes. --- .../cloud/gateway/filter/WebsocketRoutingFilter.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java index d512d43d..8286e7d3 100644 --- a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java +++ b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java @@ -214,8 +214,10 @@ public class WebsocketRoutingFilter implements GlobalFilter, Ordered { Mono serverSessionSend = session .send(proxySession.receive().doOnNext(WebSocketMessage::retain)); // .log("sessionSend", Level.FINE); - return Mono.zip(proxySessionSend, serverSessionSend, serverClose, proxyClose).then(); - } + // Ensure closeStatus from one propagates to the other + Mono.when(serverClose, proxyClose).subscribe(); + // Complete when both sessions are done + return Mono.zip(proxySessionSend, serverSessionSend).then(); } /** * Copy subProtocols so they are available downstream. From 0be612a5728794a2bd4eb2d3cd83d19669439dc6 Mon Sep 17 00:00:00 2001 From: spencergibb Date: Mon, 1 Nov 2021 15:40:08 -0400 Subject: [PATCH 3/3] formatting --- .../cloud/gateway/filter/WebsocketRoutingFilter.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java index 8286e7d3..c6e61b82 100644 --- a/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java +++ b/spring-cloud-gateway-server/src/main/java/org/springframework/cloud/gateway/filter/WebsocketRoutingFilter.java @@ -217,7 +217,8 @@ public class WebsocketRoutingFilter implements GlobalFilter, Ordered { // Ensure closeStatus from one propagates to the other Mono.when(serverClose, proxyClose).subscribe(); // Complete when both sessions are done - return Mono.zip(proxySessionSend, serverSessionSend).then(); } + return Mono.zip(proxySessionSend, serverSessionSend).then(); + } /** * Copy subProtocols so they are available downstream.