Polish SockJS client
This commit is contained in:
@@ -17,6 +17,7 @@
|
||||
package org.springframework.web.socket.sockjs.client;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.net.URI;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.HashMap;
|
||||
@@ -27,7 +28,6 @@ import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.LinkedBlockingQueue;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import javax.servlet.Filter;
|
||||
import javax.servlet.FilterChain;
|
||||
import javax.servlet.FilterConfig;
|
||||
@@ -39,7 +39,6 @@ import javax.servlet.http.HttpServletResponse;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.junit.BeforeClass;
|
||||
@@ -56,6 +55,7 @@ import org.springframework.tests.TestGroup;
|
||||
import org.springframework.util.concurrent.ListenableFutureCallback;
|
||||
import org.springframework.web.context.support.AnnotationConfigWebApplicationContext;
|
||||
import org.springframework.web.socket.TextMessage;
|
||||
import org.springframework.web.socket.WebSocketHttpHeaders;
|
||||
import org.springframework.web.socket.WebSocketSession;
|
||||
import org.springframework.web.socket.WebSocketTestServer;
|
||||
import org.springframework.web.socket.config.annotation.EnableWebSocket;
|
||||
@@ -66,7 +66,10 @@ import org.springframework.web.socket.server.HandshakeHandler;
|
||||
import org.springframework.web.socket.server.RequestUpgradeStrategy;
|
||||
import org.springframework.web.socket.server.support.DefaultHandshakeHandler;
|
||||
|
||||
import static org.junit.Assert.*;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
import static org.junit.Assert.fail;
|
||||
|
||||
/**
|
||||
* Abstract base class for integration tests using the
|
||||
@@ -90,7 +93,7 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
|
||||
private AnnotationConfigWebApplicationContext wac;
|
||||
|
||||
private ErrorFilter errorFilter;
|
||||
private TestFilter testFilter;
|
||||
|
||||
private String baseUrl;
|
||||
|
||||
@@ -104,12 +107,12 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
@Before
|
||||
public void setup() throws Exception {
|
||||
logger.debug("Setting up '" + this.testName.getMethodName() + "'");
|
||||
this.errorFilter = new ErrorFilter();
|
||||
this.testFilter = new TestFilter();
|
||||
this.wac = new AnnotationConfigWebApplicationContext();
|
||||
this.wac.register(TestConfig.class, upgradeStrategyConfigClass());
|
||||
this.server = createWebSocketTestServer();
|
||||
this.server.setup();
|
||||
this.server.deployConfig(this.wac, this.errorFilter);
|
||||
this.server.deployConfig(this.wac, this.testFilter);
|
||||
// Set ServletContext in WebApplicationContext after deployment but before
|
||||
// starting the server.
|
||||
this.wac.setServletContext(this.server.getServletContext());
|
||||
@@ -178,25 +181,25 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void receiveOneMessageWebSocket() throws Exception {
|
||||
testReceiveOneMessage(createWebSocketTransport());
|
||||
testReceiveOneMessage(createWebSocketTransport(), null);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void receiveOneMessageXhrStreaming() throws Exception {
|
||||
testReceiveOneMessage(createXhrTransport());
|
||||
testReceiveOneMessage(createXhrTransport(), null);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void receiveOneMessageXhr() throws Exception {
|
||||
AbstractXhrTransport xhrTransport = createXhrTransport();
|
||||
xhrTransport.setXhrStreamingDisabled(true);
|
||||
testReceiveOneMessage(xhrTransport);
|
||||
testReceiveOneMessage(xhrTransport, null);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void infoRequestFailure() throws Exception {
|
||||
TestClientHandler handler = new TestClientHandler();
|
||||
this.errorFilter.responseStatusMap.put("/info", 500);
|
||||
this.testFilter.sendErrorMap.put("/info", 500);
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
initSockJsClient(createWebSocketTransport());
|
||||
this.sockJsClient.doHandshake(handler, this.baseUrl + "/echo").addCallback(
|
||||
@@ -204,6 +207,7 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
@Override
|
||||
public void onSuccess(WebSocketSession result) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onFailure(Throwable ex) {
|
||||
latch.countDown();
|
||||
@@ -215,8 +219,8 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
|
||||
@Test
|
||||
public void fallbackAfterTransportFailure() throws Exception {
|
||||
this.errorFilter.responseStatusMap.put("/websocket", 200);
|
||||
this.errorFilter.responseStatusMap.put("/xhr_streaming", 500);
|
||||
this.testFilter.sendErrorMap.put("/websocket", 200);
|
||||
this.testFilter.sendErrorMap.put("/xhr_streaming", 500);
|
||||
TestClientHandler handler = new TestClientHandler();
|
||||
initSockJsClient(createWebSocketTransport(), createXhrTransport());
|
||||
WebSocketSession session = this.sockJsClient.doHandshake(handler, this.baseUrl + "/echo").get();
|
||||
@@ -229,8 +233,8 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
@Test(timeout = 5000)
|
||||
public void fallbackAfterConnectTimeout() throws Exception {
|
||||
TestClientHandler clientHandler = new TestClientHandler();
|
||||
this.errorFilter.sleepDelayMap.put("/xhr_streaming", 10000L);
|
||||
this.errorFilter.responseStatusMap.put("/xhr_streaming", 503);
|
||||
this.testFilter.sleepDelayMap.put("/xhr_streaming", 10000L);
|
||||
this.testFilter.sendErrorMap.put("/xhr_streaming", 503);
|
||||
initSockJsClient(createXhrTransport());
|
||||
this.sockJsClient.setConnectTimeoutScheduler(this.wac.getBean(ThreadPoolTaskScheduler.class));
|
||||
WebSocketSession clientSession = sockJsClient.doHandshake(clientHandler, this.baseUrl + "/echo").get();
|
||||
@@ -261,10 +265,12 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
session.close();
|
||||
}
|
||||
|
||||
private void testReceiveOneMessage(Transport transport) throws Exception {
|
||||
private void testReceiveOneMessage(Transport transport, WebSocketHttpHeaders headers)
|
||||
throws Exception {
|
||||
|
||||
TestClientHandler clientHandler = new TestClientHandler();
|
||||
initSockJsClient(transport);
|
||||
this.sockJsClient.doHandshake(clientHandler, this.baseUrl + "/test").get();
|
||||
this.sockJsClient.doHandshake(clientHandler, headers, new URI(this.baseUrl + "/test")).get();
|
||||
TestServerHandler serverHandler = this.wac.getBean(TestServerHandler.class);
|
||||
|
||||
assertNotNull("afterConnectionEstablished should have been called", clientHandler.session);
|
||||
@@ -275,6 +281,22 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
clientHandler.awaitMessage(message, 5000);
|
||||
}
|
||||
|
||||
private static void awaitEvent(Supplier<Boolean> condition, long timeToWait, String description) {
|
||||
long timeToSleep = 200;
|
||||
for (int i = 0 ; i < Math.floor(timeToWait / timeToSleep); i++) {
|
||||
if (condition.get()) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
Thread.sleep(timeToSleep);
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
throw new IllegalStateException("Interrupted while waiting for " + description, e);
|
||||
}
|
||||
}
|
||||
throw new IllegalStateException("Timed out waiting for " + description);
|
||||
}
|
||||
|
||||
|
||||
@Configuration
|
||||
@EnableWebSocket
|
||||
@@ -296,23 +318,6 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
private static void awaitEvent(Supplier<Boolean> condition, long timeToWait, String description) {
|
||||
long timeToSleep = 200;
|
||||
for (int i = 0 ; i < Math.floor(timeToWait / timeToSleep); i++) {
|
||||
if (condition.get()) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
Thread.sleep(timeToSleep);
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
throw new IllegalStateException("Interrupted while waiting for " + description, e);
|
||||
}
|
||||
}
|
||||
throw new IllegalStateException("Timed out waiting for " + description);
|
||||
}
|
||||
|
||||
private static class TestClientHandler extends TextWebSocketHandler {
|
||||
|
||||
private final BlockingQueue<TextMessage> receivedMessages = new LinkedBlockingQueue<>();
|
||||
@@ -356,6 +361,14 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
}
|
||||
}
|
||||
|
||||
private static class EchoHandler extends TextWebSocketHandler {
|
||||
|
||||
@Override
|
||||
protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception {
|
||||
session.sendMessage(message);
|
||||
}
|
||||
}
|
||||
|
||||
private static class TestServerHandler extends TextWebSocketHandler {
|
||||
|
||||
private WebSocketSession session;
|
||||
@@ -371,24 +384,23 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
}
|
||||
}
|
||||
|
||||
private static class EchoHandler extends TextWebSocketHandler {
|
||||
private static class TestFilter implements Filter {
|
||||
|
||||
@Override
|
||||
protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception {
|
||||
session.sendMessage(message);
|
||||
}
|
||||
}
|
||||
|
||||
private static class ErrorFilter implements Filter {
|
||||
|
||||
private final Map<String, Integer> responseStatusMap = new HashMap<>();
|
||||
private final List<ServletRequest> requests = new ArrayList<>();
|
||||
|
||||
private final Map<String, Long> sleepDelayMap = new HashMap<>();
|
||||
|
||||
private final Map<String, Integer> sendErrorMap = new HashMap<>();
|
||||
|
||||
|
||||
@Override
|
||||
public void doFilter(ServletRequest req, ServletResponse resp, FilterChain chain) throws IOException, ServletException {
|
||||
public void doFilter(ServletRequest request, ServletResponse response, FilterChain chain)
|
||||
throws IOException, ServletException {
|
||||
|
||||
this.requests.add(request);
|
||||
|
||||
for (String suffix : this.sleepDelayMap.keySet()) {
|
||||
if (((HttpServletRequest) req).getRequestURI().endsWith(suffix)) {
|
||||
if (((HttpServletRequest) request).getRequestURI().endsWith(suffix)) {
|
||||
try {
|
||||
Thread.sleep(this.sleepDelayMap.get(suffix));
|
||||
break;
|
||||
@@ -398,17 +410,17 @@ public abstract class AbstractSockJsIntegrationTests {
|
||||
}
|
||||
}
|
||||
}
|
||||
for (String suffix : this.responseStatusMap.keySet()) {
|
||||
if (((HttpServletRequest) req).getRequestURI().endsWith(suffix)) {
|
||||
((HttpServletResponse) resp).sendError(this.responseStatusMap.get(suffix));
|
||||
for (String suffix : this.sendErrorMap.keySet()) {
|
||||
if (((HttpServletRequest) request).getRequestURI().endsWith(suffix)) {
|
||||
((HttpServletResponse) response).sendError(this.sendErrorMap.get(suffix));
|
||||
return;
|
||||
}
|
||||
}
|
||||
chain.doFilter(req, resp);
|
||||
chain.doFilter(request, response);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void init(FilterConfig filterConfig) throws ServletException {
|
||||
public void init(FilterConfig filterConfig) {
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
Reference in New Issue
Block a user