From 5829e1c1411a22b8c17e7c3f5d7a26a537a39dfe Mon Sep 17 00:00:00 2001 From: Rossen Stoyanchev Date: Mon, 12 Dec 2016 17:02:18 -0500 Subject: [PATCH] Polish method and field declaration order --- .../AbstractListenerWebSocketSession.java | 25 ++++++----- .../adapter/JettyWebSocketHandlerAdapter.java | 32 +++++++------- .../socket/adapter/JettyWebSocketSession.java | 30 ++++++------- .../TomcatWebSocketHandlerAdapter.java | 34 +++++++-------- .../adapter/TomcatWebSocketSession.java | 42 +++++++++---------- .../UndertowWebSocketHandlerAdapter.java | 4 +- .../adapter/UndertowWebSocketSession.java | 26 ++++++------ 7 files changed, 98 insertions(+), 95 deletions(-) diff --git a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/AbstractListenerWebSocketSession.java b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/AbstractListenerWebSocketSession.java index dbb4c53028..40aad9df42 100644 --- a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/AbstractListenerWebSocketSession.java +++ b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/AbstractListenerWebSocketSession.java @@ -48,8 +48,6 @@ public abstract class AbstractListenerWebSocketSession extends WebSocketSessi private static final int RECEIVE_BUFFER_SIZE = 8192; - private final AtomicBoolean sendCalled = new AtomicBoolean(); - private final String id; private final URI uri; @@ -58,6 +56,8 @@ public abstract class AbstractListenerWebSocketSession extends WebSocketSessi private volatile WebSocketSendProcessor sendProcessor; + private final AtomicBoolean sendCalled = new AtomicBoolean(); + public AbstractListenerWebSocketSession(T delegate, String id, URI uri) { super(delegate); @@ -104,13 +104,10 @@ public abstract class AbstractListenerWebSocketSession extends WebSocketSessi } /** - * Resume receiving new message(s) after demand is generated by the - * downstream Subscriber. - *

Note: if the underlying WebSocket API does not provide - * flow control for receiving messages, and this method should be a no-op - * and {@link #canSuspendReceiving()} should return {@code false}. + * Whether the underlying WebSocket API has flow control and can suspend and + * resume the receiving of messages. */ - protected abstract void resumeReceiving(); + protected abstract boolean canSuspendReceiving(); /** * Suspend receiving until received message(s) are processed and more demand @@ -122,16 +119,22 @@ public abstract class AbstractListenerWebSocketSession extends WebSocketSessi protected abstract void suspendReceiving(); /** - * Whether the underlying WebSocket API has flow control and can suspend and - * resume the receiving of messages. + * Resume receiving new message(s) after demand is generated by the + * downstream Subscriber. + *

Note: if the underlying WebSocket API does not provide + * flow control for receiving messages, and this method should be a no-op + * and {@link #canSuspendReceiving()} should return {@code false}. */ - protected abstract boolean canSuspendReceiving(); + protected abstract void resumeReceiving(); /** * Send the given WebSocket message. */ protected abstract boolean sendMessage(WebSocketMessage message) throws IOException; + + // WebSocketHandler adapter delegate methods + /** Handle a message callback from the WebSocketHandler adapter */ void handleMessage(Type type, WebSocketMessage message) { this.receivePublisher.handleMessage(message); diff --git a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/JettyWebSocketHandlerAdapter.java b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/JettyWebSocketHandlerAdapter.java index 9314dfe85d..4d5c08816f 100644 --- a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/JettyWebSocketHandlerAdapter.java +++ b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/JettyWebSocketHandlerAdapter.java @@ -53,12 +53,12 @@ public class JettyWebSocketHandlerAdapter { private static final ByteBuffer EMPTY_PAYLOAD = ByteBuffer.wrap(new byte[0]); - private final DataBufferFactory bufferFactory = new DefaultDataBufferFactory(false); - private final WebSocketHandler delegate; private JettyWebSocketSession session; + private final DataBufferFactory bufferFactory = new DefaultDataBufferFactory(false); + public JettyWebSocketHandlerAdapter(WebSocketHandler delegate) { Assert.notNull("WebSocketHandler is required"); @@ -102,20 +102,6 @@ public class JettyWebSocketHandlerAdapter { } } - @OnWebSocketClose - public void onWebSocketClose(int statusCode, String reason) { - if (this.session != null) { - this.session.handleClose(new CloseStatus(statusCode, reason)); - } - } - - @OnWebSocketError - public void onWebSocketError(Throwable cause) { - if (this.session != null) { - this.session.handleError(cause); - } - } - private WebSocketMessage toMessage(Type type, T message) { if (Type.TEXT.equals(type)) { byte[] bytes = ((String) message).getBytes(StandardCharsets.UTF_8); @@ -135,6 +121,20 @@ public class JettyWebSocketHandlerAdapter { } } + @OnWebSocketClose + public void onWebSocketClose(int statusCode, String reason) { + if (this.session != null) { + this.session.handleClose(new CloseStatus(statusCode, reason)); + } + } + + @OnWebSocketError + public void onWebSocketError(Throwable cause) { + if (this.session != null) { + this.session.handleError(cause); + } + } + private final class HandlerResultSubscriber implements Subscriber { diff --git a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/JettyWebSocketSession.java b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/JettyWebSocketSession.java index 7f23241f53..5f83418870 100644 --- a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/JettyWebSocketSession.java +++ b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/JettyWebSocketSession.java @@ -46,9 +46,18 @@ public class JettyWebSocketSession extends AbstractListenerWebSocketSession closeInternal(CloseStatus status) { - getDelegate().close(status.getCode(), status.getReason()); - return Mono.empty(); + protected boolean canSuspendReceiving() { + return false; + } + + @Override + protected void suspendReceiving() { + // No-op + } + + @Override + protected void resumeReceiving() { + // No-op } @Override @@ -76,18 +85,9 @@ public class JettyWebSocketSession extends AbstractListenerWebSocketSession closeInternal(CloseStatus status) { + getDelegate().close(status.getCode(), status.getReason()); + return Mono.empty(); } diff --git a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/TomcatWebSocketHandlerAdapter.java b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/TomcatWebSocketHandlerAdapter.java index 042c7da643..f1913a798d 100644 --- a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/TomcatWebSocketHandlerAdapter.java +++ b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/TomcatWebSocketHandlerAdapter.java @@ -45,12 +45,12 @@ import org.springframework.web.reactive.socket.WebSocketMessage.Type; */ public class TomcatWebSocketHandlerAdapter extends Endpoint { - private final DataBufferFactory bufferFactory = new DefaultDataBufferFactory(false); - private final WebSocketHandler delegate; private TomcatWebSocketSession session; + private final DataBufferFactory bufferFactory = new DefaultDataBufferFactory(false); + public TomcatWebSocketHandlerAdapter(WebSocketHandler delegate) { Assert.notNull("WebSocketHandler is required"); @@ -79,21 +79,6 @@ public class TomcatWebSocketHandlerAdapter extends Endpoint { this.delegate.handle(this.session).subscribe(resultSubscriber); } - @Override - public void onClose(Session session, CloseReason reason) { - if (this.session != null) { - int code = reason.getCloseCode().getCode(); - this.session.handleClose(new CloseStatus(code, reason.getReasonPhrase())); - } - } - - @Override - public void onError(Session session, Throwable exception) { - if (this.session != null) { - this.session.handleError(exception); - } - } - private WebSocketMessage toMessage(T message) { if (message instanceof String) { byte[] bytes = ((String) message).getBytes(StandardCharsets.UTF_8); @@ -112,6 +97,21 @@ public class TomcatWebSocketHandlerAdapter extends Endpoint { } } + @Override + public void onClose(Session session, CloseReason reason) { + if (this.session != null) { + int code = reason.getCloseCode().getCode(); + this.session.handleClose(new CloseStatus(code, reason.getReasonPhrase())); + } + } + + @Override + public void onError(Session session, Throwable exception) { + if (this.session != null) { + this.session.handleError(exception); + } + } + private final class HandlerResultSubscriber implements Subscriber { diff --git a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/TomcatWebSocketSession.java b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/TomcatWebSocketSession.java index a9698f6962..80aed582e3 100644 --- a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/TomcatWebSocketSession.java +++ b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/TomcatWebSocketSession.java @@ -47,15 +47,18 @@ public class TomcatWebSocketSession extends AbstractListenerWebSocketSession closeInternal(CloseStatus status) { - try { - getDelegate().close( - new CloseReason(CloseCodes.getCloseCode(status.getCode()), status.getReason())); - } - catch (IOException e) { - return Mono.error(e); - } - return Mono.empty(); + protected boolean canSuspendReceiving() { + return false; + } + + @Override + protected void suspendReceiving() { + // No-op + } + + @Override + protected void resumeReceiving() { + // No-op } @Override @@ -83,18 +86,15 @@ public class TomcatWebSocketSession extends AbstractListenerWebSocketSession closeInternal(CloseStatus status) { + try { + CloseReason.CloseCode code = CloseCodes.getCloseCode(status.getCode()); + getDelegate().close(new CloseReason(code, status.getReason())); + } + catch (IOException e) { + return Mono.error(e); + } + return Mono.empty(); } diff --git a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/UndertowWebSocketHandlerAdapter.java b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/UndertowWebSocketHandlerAdapter.java index 904c23e8bc..a0a17c3716 100644 --- a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/UndertowWebSocketHandlerAdapter.java +++ b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/UndertowWebSocketHandlerAdapter.java @@ -49,12 +49,12 @@ import io.undertow.websockets.spi.WebSocketHttpExchange; */ public class UndertowWebSocketHandlerAdapter implements WebSocketConnectionCallback { - private final DataBufferFactory bufferFactory = new DefaultDataBufferFactory(false); - private final WebSocketHandler delegate; private UndertowWebSocketSession session; + private final DataBufferFactory bufferFactory = new DefaultDataBufferFactory(false); + public UndertowWebSocketHandlerAdapter(WebSocketHandler delegate) { Assert.notNull("WebSocketHandler is required"); diff --git a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/UndertowWebSocketSession.java b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/UndertowWebSocketSession.java index c2a3a15b06..6d0b557753 100644 --- a/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/UndertowWebSocketSession.java +++ b/spring-web-reactive/src/main/java/org/springframework/web/reactive/socket/adapter/UndertowWebSocketSession.java @@ -49,25 +49,16 @@ public class UndertowWebSocketSession extends AbstractListenerWebSocketSession closeInternal(CloseStatus status) { - CloseMessage cm = new CloseMessage(status.getCode(), status.getReason()); - if (!getDelegate().isCloseFrameSent()) { - WebSockets.sendClose(cm, getDelegate(), null); - } - return Mono.empty(); - } - - protected void resumeReceiving() { - getDelegate().resumeReceives(); + protected boolean canSuspendReceiving() { + return true; } protected void suspendReceiving() { getDelegate().suspendReceives(); } - @Override - protected boolean canSuspendReceiving() { - return true; + protected void resumeReceiving() { + getDelegate().resumeReceives(); } @Override @@ -96,6 +87,15 @@ public class UndertowWebSocketSession extends AbstractListenerWebSocketSession closeInternal(CloseStatus status) { + CloseMessage cm = new CloseMessage(status.getCode(), status.getReason()); + if (!getDelegate().isCloseFrameSent()) { + WebSockets.sendClose(cm, getDelegate(), null); + } + return Mono.empty(); + } + private final class SendProcessorCallback implements WebSocketCallback {