Replace relevant code with lambda
See gh-1454
This commit is contained in:
@@ -170,13 +170,10 @@ public class JettyWebSocketClient extends AbstractWebSocketClient implements Lif
|
||||
final JettyWebSocketSession wsSession = new JettyWebSocketSession(attributes, user);
|
||||
final JettyWebSocketHandlerAdapter listener = new JettyWebSocketHandlerAdapter(wsHandler, wsSession);
|
||||
|
||||
Callable<WebSocketSession> connectTask = new Callable<WebSocketSession>() {
|
||||
@Override
|
||||
public WebSocketSession call() throws Exception {
|
||||
Future<Session> future = client.connect(listener, uri, request);
|
||||
future.get();
|
||||
return wsSession;
|
||||
}
|
||||
Callable<WebSocketSession> connectTask = () -> {
|
||||
Future<Session> future = client.connect(listener, uri, request);
|
||||
future.get();
|
||||
return wsSession;
|
||||
};
|
||||
|
||||
if (this.taskExecutor != null) {
|
||||
|
||||
@@ -99,20 +99,17 @@ public class AnnotatedEndpointConnectionManager extends ConnectionManagerSupport
|
||||
|
||||
@Override
|
||||
protected void openConnection() {
|
||||
this.taskExecutor.execute(new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
try {
|
||||
if (logger.isInfoEnabled()) {
|
||||
logger.info("Connecting to WebSocket at " + getUri());
|
||||
}
|
||||
Object endpointToUse = (endpoint != null) ? endpoint : endpointProvider.getHandler();
|
||||
session = webSocketContainer.connectToServer(endpointToUse, getUri());
|
||||
logger.info("Successfully connected to WebSocket");
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
logger.error("Failed to connect to WebSocket", ex);
|
||||
this.taskExecutor.execute(() -> {
|
||||
try {
|
||||
if (logger.isInfoEnabled()) {
|
||||
logger.info("Connecting to WebSocket at " + getUri());
|
||||
}
|
||||
Object endpointToUse = (endpoint != null) ? endpoint : endpointProvider.getHandler();
|
||||
session = webSocketContainer.connectToServer(endpointToUse, getUri());
|
||||
logger.info("Successfully connected to WebSocket");
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
logger.error("Failed to connect to WebSocket", ex);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@@ -130,21 +130,18 @@ public class EndpointConnectionManager extends ConnectionManagerSupport implemen
|
||||
|
||||
@Override
|
||||
protected void openConnection() {
|
||||
this.taskExecutor.execute(new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
try {
|
||||
if (logger.isInfoEnabled()) {
|
||||
logger.info("Connecting to WebSocket at " + getUri());
|
||||
}
|
||||
Endpoint endpointToUse = (endpoint != null) ? endpoint : endpointProvider.getHandler();
|
||||
ClientEndpointConfig endpointConfig = configBuilder.build();
|
||||
session = getWebSocketContainer().connectToServer(endpointToUse, endpointConfig, getUri());
|
||||
logger.info("Successfully connected to WebSocket");
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
logger.error("Failed to connect to WebSocket", ex);
|
||||
this.taskExecutor.execute(() -> {
|
||||
try {
|
||||
if (logger.isInfoEnabled()) {
|
||||
logger.info("Connecting to WebSocket at " + getUri());
|
||||
}
|
||||
Endpoint endpointToUse = (endpoint != null) ? endpoint : endpointProvider.getHandler();
|
||||
ClientEndpointConfig endpointConfig = configBuilder.build();
|
||||
session = getWebSocketContainer().connectToServer(endpointToUse, endpointConfig, getUri());
|
||||
logger.info("Successfully connected to WebSocket");
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
logger.error("Failed to connect to WebSocket", ex);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@@ -144,12 +144,9 @@ public class StandardWebSocketClient extends AbstractWebSocketClient {
|
||||
|
||||
final Endpoint endpoint = new StandardWebSocketHandlerAdapter(webSocketHandler, session);
|
||||
|
||||
Callable<WebSocketSession> connectTask = new Callable<WebSocketSession>() {
|
||||
@Override
|
||||
public WebSocketSession call() throws Exception {
|
||||
webSocketContainer.connectToServer(endpoint, endpointConfig, uri);
|
||||
return session;
|
||||
}
|
||||
Callable<WebSocketSession> connectTask = () -> {
|
||||
webSocketContainer.connectToServer(endpoint, endpointConfig, uri);
|
||||
return session;
|
||||
};
|
||||
|
||||
if (this.taskExecutor != null) {
|
||||
|
||||
@@ -110,12 +110,9 @@ public class WebSocketMessageBrokerStats {
|
||||
@Nullable
|
||||
private ScheduledFuture<?> initLoggingTask(long initialDelay) {
|
||||
if (logger.isInfoEnabled() && this.loggingPeriod > 0) {
|
||||
return this.sockJsTaskScheduler.scheduleAtFixedRate(new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
logger.info(WebSocketMessageBrokerStats.this.toString());
|
||||
}
|
||||
}, initialDelay, this.loggingPeriod, TimeUnit.MILLISECONDS);
|
||||
return this.sockJsTaskScheduler.scheduleAtFixedRate(()
|
||||
-> logger.info(WebSocketMessageBrokerStats.this.toString()),
|
||||
initialDelay, this.loggingPeriod, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -401,21 +401,18 @@ public class WebSocketStompClient extends StompClientSupport implements SmartLif
|
||||
public void onReadInactivity(final Runnable runnable, final long duration) {
|
||||
Assert.state(getTaskScheduler() != null, "No TaskScheduler configured");
|
||||
this.lastReadTime = System.currentTimeMillis();
|
||||
this.inactivityTasks.add(getTaskScheduler().scheduleWithFixedDelay(new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
if (System.currentTimeMillis() - lastReadTime > duration) {
|
||||
try {
|
||||
runnable.run();
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("ReadInactivityTask failure", ex);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}, duration / 2));
|
||||
this.inactivityTasks.add(getTaskScheduler().scheduleWithFixedDelay(() -> {
|
||||
if (System.currentTimeMillis() - lastReadTime > duration) {
|
||||
try {
|
||||
runnable.run();
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("ReadInactivityTask failure", ex);
|
||||
}
|
||||
}
|
||||
}
|
||||
}, duration / 2));
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -116,12 +116,7 @@ public abstract class AbstractClientSockJsSession implements WebSocketSession {
|
||||
* request.
|
||||
*/
|
||||
Runnable getTimeoutTask() {
|
||||
return new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
closeInternal(new CloseStatus(2007, "Transport timed out"));
|
||||
}
|
||||
};
|
||||
return () -> closeInternal(new CloseStatus(2007, "Transport timed out"));
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -99,35 +99,32 @@ public class RestTemplateXhrTransport extends AbstractXhrTransport {
|
||||
final URI receiveUrl, final HttpHeaders handshakeHeaders, final XhrClientSockJsSession session,
|
||||
final SettableListenableFuture<WebSocketSession> connectFuture) {
|
||||
|
||||
getTaskExecutor().execute(new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
HttpHeaders httpHeaders = transportRequest.getHttpRequestHeaders();
|
||||
XhrRequestCallback requestCallback = new XhrRequestCallback(handshakeHeaders);
|
||||
XhrRequestCallback requestCallbackAfterHandshake = new XhrRequestCallback(httpHeaders);
|
||||
XhrReceiveExtractor responseExtractor = new XhrReceiveExtractor(session);
|
||||
while (true) {
|
||||
if (session.isDisconnected()) {
|
||||
session.afterTransportClosed(null);
|
||||
break;
|
||||
getTaskExecutor().execute(() -> {
|
||||
HttpHeaders httpHeaders = transportRequest.getHttpRequestHeaders();
|
||||
XhrRequestCallback requestCallback = new XhrRequestCallback(handshakeHeaders);
|
||||
XhrRequestCallback requestCallbackAfterHandshake = new XhrRequestCallback(httpHeaders);
|
||||
XhrReceiveExtractor responseExtractor = new XhrReceiveExtractor(session);
|
||||
while (true) {
|
||||
if (session.isDisconnected()) {
|
||||
session.afterTransportClosed(null);
|
||||
break;
|
||||
}
|
||||
try {
|
||||
if (logger.isTraceEnabled()) {
|
||||
logger.trace("Starting XHR receive request, url=" + receiveUrl);
|
||||
}
|
||||
try {
|
||||
if (logger.isTraceEnabled()) {
|
||||
logger.trace("Starting XHR receive request, url=" + receiveUrl);
|
||||
}
|
||||
getRestTemplate().execute(receiveUrl, HttpMethod.POST, requestCallback, responseExtractor);
|
||||
requestCallback = requestCallbackAfterHandshake;
|
||||
getRestTemplate().execute(receiveUrl, HttpMethod.POST, requestCallback, responseExtractor);
|
||||
requestCallback = requestCallbackAfterHandshake;
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
if (!connectFuture.isDone()) {
|
||||
connectFuture.setException(ex);
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
if (!connectFuture.isDone()) {
|
||||
connectFuture.setException(ex);
|
||||
}
|
||||
else {
|
||||
session.handleTransportError(ex);
|
||||
session.afterTransportClosed(new CloseStatus(1006, ex.getMessage()));
|
||||
}
|
||||
break;
|
||||
else {
|
||||
session.handleTransportError(ex);
|
||||
session.afterTransportClosed(new CloseStatus(1006, ex.getMessage()));
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
@@ -368,27 +368,24 @@ public class TransportHandlingSockJsService extends AbstractSockJsService implem
|
||||
if (this.sessionCleanupTask != null) {
|
||||
return;
|
||||
}
|
||||
this.sessionCleanupTask = getTaskScheduler().scheduleAtFixedRate(new Runnable() {
|
||||
@Override
|
||||
public void run() {
|
||||
List<String> removedIds = new ArrayList<>();
|
||||
for (SockJsSession session : sessions.values()) {
|
||||
try {
|
||||
if (session.getTimeSinceLastActive() > getDisconnectDelay()) {
|
||||
sessions.remove(session.getId());
|
||||
removedIds.add(session.getId());
|
||||
session.close();
|
||||
}
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
// Could be part of normal workflow (e.g. browser tab closed)
|
||||
logger.debug("Failed to close " + session, ex);
|
||||
this.sessionCleanupTask = getTaskScheduler().scheduleAtFixedRate(() -> {
|
||||
List<String> removedIds = new ArrayList<>();
|
||||
for (SockJsSession session : sessions.values()) {
|
||||
try {
|
||||
if (session.getTimeSinceLastActive() > getDisconnectDelay()) {
|
||||
sessions.remove(session.getId());
|
||||
removedIds.add(session.getId());
|
||||
session.close();
|
||||
}
|
||||
}
|
||||
if (logger.isDebugEnabled() && !removedIds.isEmpty()) {
|
||||
logger.debug("Closed " + removedIds.size() + " sessions: " + removedIds);
|
||||
catch (Throwable ex) {
|
||||
// Could be part of normal workflow (e.g. browser tab closed)
|
||||
logger.debug("Failed to close " + session, ex);
|
||||
}
|
||||
}
|
||||
if (logger.isDebugEnabled() && !removedIds.isEmpty()) {
|
||||
logger.debug("Closed " + removedIds.size() + " sessions: " + removedIds);
|
||||
}
|
||||
}, getDisconnectDelay());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user