Polishing
This commit is contained in:
@@ -129,8 +129,9 @@ public class WebSocketConnectionManager extends ConnectionManagerSupport {
|
||||
|
||||
@Override
|
||||
protected void openConnection() {
|
||||
|
||||
logger.info("Connecting to WebSocket at " + getUri());
|
||||
if (logger.isInfoEnabled()) {
|
||||
logger.info("Connecting to WebSocket at " + getUri());
|
||||
}
|
||||
|
||||
ListenableFuture<WebSocketSession> future =
|
||||
this.client.doHandshake(this.webSocketHandler, this.headers, getUri());
|
||||
@@ -142,8 +143,8 @@ public class WebSocketConnectionManager extends ConnectionManagerSupport {
|
||||
logger.info("Successfully connected");
|
||||
}
|
||||
@Override
|
||||
public void onFailure(Throwable t) {
|
||||
logger.error("Failed to connect", t);
|
||||
public void onFailure(Throwable ex) {
|
||||
logger.error("Failed to connect", ex);
|
||||
}
|
||||
});
|
||||
}
|
||||
@@ -155,7 +156,7 @@ public class WebSocketConnectionManager extends ConnectionManagerSupport {
|
||||
|
||||
@Override
|
||||
protected boolean isConnected() {
|
||||
return ((this.webSocketSession != null) && (this.webSocketSession.isOpen()));
|
||||
return (this.webSocketSession != null && this.webSocketSession.isOpen());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -16,8 +16,16 @@
|
||||
|
||||
package org.springframework.web.socket.sockjs.client;
|
||||
|
||||
import java.net.URI;
|
||||
import java.security.Principal;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Date;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.springframework.http.HttpHeaders;
|
||||
import org.springframework.scheduling.TaskScheduler;
|
||||
import org.springframework.util.Assert;
|
||||
@@ -29,17 +37,8 @@ import org.springframework.web.socket.sockjs.SockJsTransportFailureException;
|
||||
import org.springframework.web.socket.sockjs.frame.SockJsMessageCodec;
|
||||
import org.springframework.web.socket.sockjs.transport.TransportType;
|
||||
|
||||
import java.net.URI;
|
||||
import java.security.Principal;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Date;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
/**
|
||||
* A default implementation of
|
||||
* {@link org.springframework.web.socket.sockjs.client.TransportRequest
|
||||
* TransportRequest}.
|
||||
* A default implementation of {@link TransportRequest}.
|
||||
*
|
||||
* @author Rossen Stoyanchev
|
||||
* @since 4.1
|
||||
@@ -176,26 +175,24 @@ class DefaultTransportRequest implements TransportRequest {
|
||||
|
||||
private final AtomicBoolean handled = new AtomicBoolean(false);
|
||||
|
||||
|
||||
public ConnectCallback(WebSocketHandler handler, SettableListenableFuture<WebSocketSession> future) {
|
||||
this.handler = handler;
|
||||
this.future = future;
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public void onSuccess(WebSocketSession session) {
|
||||
if (this.handled.compareAndSet(false, true)) {
|
||||
this.future.set(session);
|
||||
}
|
||||
else {
|
||||
else if (logger.isErrorEnabled()) {
|
||||
logger.error("Connect success/failure already handled for " + DefaultTransportRequest.this);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onFailure(Throwable failure) {
|
||||
handleFailure(failure, false);
|
||||
public void onFailure(Throwable ex) {
|
||||
handleFailure(ex, false);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -203,20 +200,20 @@ class DefaultTransportRequest implements TransportRequest {
|
||||
handleFailure(null, true);
|
||||
}
|
||||
|
||||
private void handleFailure(Throwable failure, boolean isTimeoutFailure) {
|
||||
private void handleFailure(Throwable ex, boolean isTimeoutFailure) {
|
||||
if (this.handled.compareAndSet(false, true)) {
|
||||
if (isTimeoutFailure) {
|
||||
String message = "Connect timed out for " + DefaultTransportRequest.this;
|
||||
logger.error(message);
|
||||
failure = new SockJsTransportFailureException(message, getSockJsUrlInfo().getSessionId(), null);
|
||||
ex = new SockJsTransportFailureException(message, getSockJsUrlInfo().getSessionId(), null);
|
||||
}
|
||||
if (fallbackRequest != null) {
|
||||
logger.error(DefaultTransportRequest.this + " failed. Falling back on next transport.", failure);
|
||||
logger.error(DefaultTransportRequest.this + " failed. Falling back on next transport.", ex);
|
||||
fallbackRequest.connect(this.handler, this.future);
|
||||
}
|
||||
else {
|
||||
logger.error("No more fallback transports after " + DefaultTransportRequest.this, failure);
|
||||
this.future.setException(failure);
|
||||
logger.error("No more fallback transports after " + DefaultTransportRequest.this, ex);
|
||||
this.future.setException(ex);
|
||||
}
|
||||
if (isTimeoutFailure) {
|
||||
try {
|
||||
@@ -224,15 +221,16 @@ class DefaultTransportRequest implements TransportRequest {
|
||||
runnable.run();
|
||||
}
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
logger.error("Transport failed to run timeout tasks for " + DefaultTransportRequest.this, ex);
|
||||
catch (Throwable ex2) {
|
||||
logger.error("Transport failed to run timeout tasks for " + DefaultTransportRequest.this, ex2);
|
||||
}
|
||||
}
|
||||
}
|
||||
else {
|
||||
logger.error("Connect success/failure events already took place for " +
|
||||
DefaultTransportRequest.this + ". Ignoring this additional failure event.", failure);
|
||||
DefaultTransportRequest.this + ". Ignoring this additional failure event.", ex);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -16,8 +16,14 @@
|
||||
|
||||
package org.springframework.web.socket.sockjs.client;
|
||||
|
||||
import java.net.URI;
|
||||
import java.util.Arrays;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.springframework.context.Lifecycle;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.concurrent.ListenableFuture;
|
||||
@@ -32,11 +38,6 @@ import org.springframework.web.socket.client.WebSocketClient;
|
||||
import org.springframework.web.socket.handler.TextWebSocketHandler;
|
||||
import org.springframework.web.socket.sockjs.transport.TransportType;
|
||||
|
||||
import java.net.URI;
|
||||
import java.util.Arrays;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
/**
|
||||
* A SockJS {@link Transport} that uses a
|
||||
* {@link org.springframework.web.socket.client.WebSocketClient WebSocketClient}.
|
||||
@@ -59,11 +60,6 @@ public class WebSocketTransport implements Transport, Lifecycle {
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public List<TransportType> getTransportTypes() {
|
||||
return Arrays.asList(TransportType.WEBSOCKET);
|
||||
}
|
||||
|
||||
/**
|
||||
* Return the configured {@code WebSocketClient}.
|
||||
*/
|
||||
@@ -71,6 +67,38 @@ public class WebSocketTransport implements Transport, Lifecycle {
|
||||
return this.webSocketClient;
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<TransportType> getTransportTypes() {
|
||||
return Arrays.asList(TransportType.WEBSOCKET);
|
||||
}
|
||||
|
||||
@Override
|
||||
public ListenableFuture<WebSocketSession> connect(TransportRequest request, WebSocketHandler handler) {
|
||||
final SettableListenableFuture<WebSocketSession> future = new SettableListenableFuture<WebSocketSession>();
|
||||
WebSocketClientSockJsSession session = new WebSocketClientSockJsSession(request, handler, future);
|
||||
handler = new ClientSockJsWebSocketHandler(session);
|
||||
request.addTimeoutTask(session.getTimeoutTask());
|
||||
|
||||
URI url = request.getTransportUrl();
|
||||
WebSocketHttpHeaders headers = new WebSocketHttpHeaders(request.getHandshakeHeaders());
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Starting WebSocket session url=" + url);
|
||||
}
|
||||
this.webSocketClient.doHandshake(handler, headers, url).addCallback(
|
||||
new ListenableFutureCallback<WebSocketSession>() {
|
||||
@Override
|
||||
public void onSuccess(WebSocketSession webSocketSession) {
|
||||
// WebSocket session ready, SockJS Session not yet
|
||||
}
|
||||
@Override
|
||||
public void onFailure(Throwable ex) {
|
||||
future.setException(ex);
|
||||
}
|
||||
});
|
||||
return future;
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public void start() {
|
||||
if (!isRunning()) {
|
||||
@@ -106,32 +134,6 @@ public class WebSocketTransport implements Transport, Lifecycle {
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public ListenableFuture<WebSocketSession> connect(TransportRequest request, WebSocketHandler handler) {
|
||||
final SettableListenableFuture<WebSocketSession> future = new SettableListenableFuture<WebSocketSession>();
|
||||
WebSocketClientSockJsSession session = new WebSocketClientSockJsSession(request, handler, future);
|
||||
handler = new ClientSockJsWebSocketHandler(session);
|
||||
request.addTimeoutTask(session.getTimeoutTask());
|
||||
|
||||
URI url = request.getTransportUrl();
|
||||
WebSocketHttpHeaders headers = new WebSocketHttpHeaders(request.getHandshakeHeaders());
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Starting WebSocket session url=" + url);
|
||||
}
|
||||
this.webSocketClient.doHandshake(handler, headers, url).addCallback(
|
||||
new ListenableFutureCallback<WebSocketSession>() {
|
||||
@Override
|
||||
public void onSuccess(WebSocketSession webSocketSession) {
|
||||
// WebSocket session ready, SockJS Session not yet
|
||||
}
|
||||
@Override
|
||||
public void onFailure(Throwable t) {
|
||||
future.setException(t);
|
||||
}
|
||||
});
|
||||
return future;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "WebSocketTransport[client=" + this.webSocketClient + "]";
|
||||
@@ -144,7 +146,6 @@ public class WebSocketTransport implements Transport, Lifecycle {
|
||||
|
||||
private final AtomicInteger connectCount = new AtomicInteger(0);
|
||||
|
||||
|
||||
private ClientSockJsWebSocketHandler(WebSocketClientSockJsSession session) {
|
||||
Assert.notNull(session);
|
||||
this.sockJsSession = session;
|
||||
|
||||
@@ -16,6 +16,25 @@
|
||||
|
||||
package org.springframework.web.socket.sockjs.client;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.BlockingQueue;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.LinkedBlockingQueue;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import javax.servlet.Filter;
|
||||
import javax.servlet.FilterChain;
|
||||
import javax.servlet.FilterConfig;
|
||||
import javax.servlet.ServletException;
|
||||
import javax.servlet.ServletRequest;
|
||||
import javax.servlet.ServletResponse;
|
||||
import javax.servlet.http.HttpServletRequest;
|
||||
import javax.servlet.http.HttpServletResponse;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.junit.After;
|
||||
@@ -24,6 +43,7 @@ import org.junit.Ignore;
|
||||
import org.junit.Rule;
|
||||
import org.junit.Test;
|
||||
import org.junit.rules.TestName;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
@@ -41,29 +61,7 @@ import org.springframework.web.socket.server.HandshakeHandler;
|
||||
import org.springframework.web.socket.server.RequestUpgradeStrategy;
|
||||
import org.springframework.web.socket.server.support.DefaultHandshakeHandler;
|
||||
|
||||
import javax.servlet.Filter;
|
||||
import javax.servlet.FilterChain;
|
||||
import javax.servlet.FilterConfig;
|
||||
import javax.servlet.ServletException;
|
||||
import javax.servlet.ServletRequest;
|
||||
import javax.servlet.ServletResponse;
|
||||
import javax.servlet.http.HttpServletRequest;
|
||||
import javax.servlet.http.HttpServletResponse;
|
||||
import java.io.IOException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.BlockingQueue;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.LinkedBlockingQueue;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.fail;
|
||||
import static org.junit.Assert.*;
|
||||
|
||||
/**
|
||||
* Integration tests using the
|
||||
@@ -188,10 +186,9 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
new ListenableFutureCallback<WebSocketSession>() {
|
||||
@Override
|
||||
public void onSuccess(WebSocketSession result) {
|
||||
|
||||
}
|
||||
@Override
|
||||
public void onFailure(Throwable t) {
|
||||
public void onFailure(Throwable ex) {
|
||||
latch.countDown();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,6 +28,7 @@ import java.util.concurrent.LinkedBlockingDeque;
|
||||
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.core.task.SyncTaskExecutor;
|
||||
import org.springframework.http.HttpHeaders;
|
||||
import org.springframework.http.HttpMethod;
|
||||
@@ -136,8 +137,8 @@ public class RestTemplateXhrTransportTests {
|
||||
public void onSuccess(WebSocketSession result) {
|
||||
}
|
||||
@Override
|
||||
public void onFailure(Throwable actual) {
|
||||
if (actual == expected) {
|
||||
public void onFailure(Throwable ex) {
|
||||
if (ex == expected) {
|
||||
latch.countDown();
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user