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 0e59030fe2..167d1146bb 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 @@ -23,6 +23,7 @@ import org.springframework.security.core.Authentication; import org.springframework.stereotype.Service; import org.springframework.util.StringUtils; import org.springframework.web.socket.CloseStatus; +import org.springframework.web.socket.PongMessage; import org.springframework.web.socket.TextMessage; import org.springframework.web.socket.WebSocketSession; import org.springframework.web.socket.adapter.NativeWebSocketSession; @@ -111,6 +112,22 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr } } + @Override + protected void handlePongMessage(WebSocketSession session, PongMessage message) throws Exception { + try { + SessionMetaData sessionMd = internalSessionMap.get(session.getId()); + if (sessionMd != null) { + log.trace("[{}][{}] Processing pong response {}", sessionMd.sessionRef.getSecurityCtx().getTenantId(), session.getId(), message.getPayload()); + sessionMd.processPongMessage(System.currentTimeMillis()); + } else { + log.trace("[{}] Failed to find session", session.getId()); + session.close(CloseStatus.SERVER_ERROR.withReason("Session not found!")); + } + } catch (IOException e) { + log.warn("IO error", e); + } + } + @Override public void afterConnectionEstablished(WebSocketSession session) throws Exception { super.afterConnectionEstablished(session); @@ -213,19 +230,29 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr synchronized void sendPing(long currentTime) { try { if (currentTime - lastActivityTime >= pingTimeout) { + log.warn("[{}] Closing session due to ping timeout", session.getId()); + closeSession(CloseStatus.SESSION_NOT_RELIABLE); + } else { this.asyncRemote.sendPing(PING_MSG); - lastActivityTime = currentTime; } } 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); - } + closeSession(CloseStatus.SESSION_NOT_RELIABLE); } } + private void closeSession(CloseStatus reason) { + try { + close(this.sessionRef, reason); + } catch (IOException ioe) { + log.trace("[{}] Session transport error", session.getId(), ioe); + } + } + + synchronized void processPongMessage(long currentTime) { + lastActivityTime = currentTime; + } + synchronized void sendMsg(String msg) { if (isSending) { try { @@ -236,11 +263,7 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr } else { log.info("[{}][{}] Session closed due to queue error", sessionRef.getSecurityCtx().getTenantId(), session.getId()); } - try { - close(sessionRef, CloseStatus.POLICY_VIOLATION.withReason("Max pending updates limit reached!")); - } catch (IOException ioe) { - log.trace("[{}] Session transport error", session.getId(), ioe); - } + closeSession(CloseStatus.POLICY_VIOLATION.withReason("Max pending updates limit reached!")); } } else { isSending = true; @@ -253,11 +276,7 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr this.asyncRemote.sendText(msg, this); } catch (Exception e) { log.trace("[{}] Failed to send msg", session.getId(), e); - try { - close(this.sessionRef, CloseStatus.SESSION_NOT_RELIABLE); - } catch (IOException ioe) { - log.trace("[{}] Session transport error", session.getId(), ioe); - } + closeSession(CloseStatus.SESSION_NOT_RELIABLE); } } @@ -265,11 +284,7 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr public void onResult(SendResult result) { if (!result.isOK()) { log.trace("[{}] Failed to send msg", session.getId(), result.getException()); - try { - close(this.sessionRef, CloseStatus.SESSION_NOT_RELIABLE); - } catch (IOException ioe) { - log.trace("[{}] Session transport error", session.getId(), ioe); - } + closeSession(CloseStatus.SESSION_NOT_RELIABLE); } else { lastActivityTime = System.currentTimeMillis(); String msg = msgQueue.poll(); diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java index 73eed820ab..482c577091 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java @@ -146,6 +146,9 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi @Value("${server.ws.limits.max_subscriptions_per_public_user:0}") private int maxSubscriptionsPerPublicUser; + @Value("${server.ws.ping_timeout:30000}") + private long pingTimeout; + private ConcurrentMap> tenantSubscriptionsMap = new ConcurrentHashMap<>(); private ConcurrentMap> customerSubscriptionsMap = new ConcurrentHashMap<>(); private ConcurrentMap> regularUserSubscriptionsMap = new ConcurrentHashMap<>(); @@ -162,7 +165,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi executor = ThingsBoardExecutors.newWorkStealingPool(50, getClass()); pingExecutor = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("telemetry-web-socket-ping")); - pingExecutor.scheduleWithFixedDelay(this::sendPing, 10000, 10000, TimeUnit.MILLISECONDS); + pingExecutor.scheduleWithFixedDelay(this::sendPing, pingTimeout/3, pingTimeout/3, TimeUnit.MILLISECONDS); } @PreDestroy