From 506151185fc34c4abc3f2c006bb95ca3288f4d9d Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Tue, 2 Mar 2021 16:42:59 +0200 Subject: [PATCH] added ping for WS --- .../controller/plugin/TbWebSocketHandler.java | 46 ++++++++++++++++++- .../src/main/resources/thingsboard.yml | 1 + 2 files changed, 46 insertions(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java b/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java index ade9797be0..7c6a2f8c9d 100644 --- a/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java +++ b/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java @@ -41,16 +41,25 @@ import org.thingsboard.server.service.telemetry.TelemetryWebSocketMsgEndpoint; import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; -import javax.websocket.*; +import javax.annotation.PreDestroy; +import javax.websocket.RemoteEndpoint; +import javax.websocket.SendHandler; +import javax.websocket.SendResult; +import javax.websocket.Session; import java.io.IOException; import java.net.URI; +import java.nio.ByteBuffer; import java.security.InvalidParameterException; import java.util.Queue; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.Executors; import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; @Service @TbCoreComponent @@ -79,6 +88,9 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr @Value("${server.ws.limits.max_updates_per_session:}") private String perSessionUpdatesConfiguration; + @Value("${server.ws.ping_timeout:30000}") + private long pingTimeout; + private ConcurrentMap blacklistedSessions = new ConcurrentHashMap<>(); private ConcurrentMap perSessionUpdateLimits = new ConcurrentHashMap<>(); @@ -87,6 +99,8 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr private ConcurrentMap> regularUserSessionsMap = new ConcurrentHashMap<>(); private ConcurrentMap> publicUserSessionsMap = new ConcurrentHashMap<>(); + private ScheduledExecutorService pingExecutor = Executors.newSingleThreadScheduledExecutor(); + @Override public void handleTextMessage(WebSocketSession session, TextMessage message) { try { @@ -120,6 +134,9 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr return; } internalSessionMap.put(internalSessionId, new SessionMetaData(session, sessionRef, maxMsgQueuePerSession)); + + pingExecutor.scheduleWithFixedDelay(() -> internalSessionMap.get(internalSessionId).sendPing(), 10000, 10000, TimeUnit.MILLISECONDS); + externalSessionMap.put(externalSessionId, internalSessionId); processInWebSocketService(sessionRef, SessionEvent.onEstablished()); log.info("[{}][{}][{}] Session is opened", sessionRef.getSecurityCtx().getTenantId(), externalSessionId, session.getId()); @@ -189,6 +206,8 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr private volatile boolean isSending = false; private final Queue msgQueue; + private final AtomicLong lastActivityTime; + SessionMetaData(WebSocketSession session, TelemetryWebSocketSessionRef sessionRef, int maxMsgQueuePerSession) { super(); this.session = session; @@ -196,12 +215,30 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr this.asyncRemote = nativeSession.getAsyncRemote(); this.sessionRef = sessionRef; this.msgQueue = new LinkedBlockingQueue<>(maxMsgQueuePerSession); + this.lastActivityTime = new AtomicLong(System.currentTimeMillis()); + } + + synchronized void sendPing() { + try { + if (System.currentTimeMillis() - lastActivityTime.get() >= pingTimeout) { + this.asyncRemote.sendPing(ByteBuffer.wrap(new byte[]{})); + lastActivityTime.set(System.currentTimeMillis()); + } + } catch (Exception e) { + log.trace("[{}] Failed to send ping msg", session.getId(), e); + try { + close(this.sessionRef, CloseStatus.SESSION_NOT_RELIABLE); + } catch (IOException ioe) { + log.trace("[{}] Session transport error", session.getId(), ioe); + } + } } synchronized void sendMsg(String msg) { if (isSending) { try { msgQueue.add(msg); + lastActivityTime.set(System.currentTimeMillis()); } catch (RuntimeException e) { if (log.isTraceEnabled()) { log.trace("[{}][{}] Session closed due to queue error", sessionRef.getSecurityCtx().getTenantId(), session.getId(), e); @@ -393,4 +430,11 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr } } + @PreDestroy + private void destroy() { + if (pingExecutor != null) { + pingExecutor.shutdownNow(); + } + } + } \ No newline at end of file diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 7097acd4dd..f750000750 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -34,6 +34,7 @@ server: log_controller_error_stack_trace: "${HTTP_LOG_CONTROLLER_ERROR_STACK_TRACE:false}" ws: send_timeout: "${TB_SERVER_WS_SEND_TIMEOUT:5000}" + ping_timeout: "${TB_SERVER_WS_PING_TIMEOUT:30000}" limits: # Limit the amount of sessions and subscriptions available on each server. Put values to zero to disable particular limitation max_sessions_per_tenant: "${TB_SERVER_WS_TENANT_RATE_LIMITS_MAX_SESSIONS_PER_TENANT:0}"