From f72c058030271ad42d6ede7f7bccb5630dee2c76 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 17 Jan 2022 12:44:06 +0200 Subject: [PATCH] Suppress and handle WS exceptions on cluster rebalancing --- ...efaultTbEntityDataSubscriptionService.java | 33 ++++++++++++++----- .../DefaultTelemetryWebSocketService.java | 12 +++++++ .../telemetry/TelemetryWebSocketService.java | 2 ++ .../core/ws/telemetry-websocket.service.ts | 2 +- 4 files changed, 39 insertions(+), 10 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java index ac831d43ba..b35a83cdf1 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java @@ -27,6 +27,7 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Lazy; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; +import org.springframework.web.socket.CloseStatus; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.TenantId; @@ -205,21 +206,30 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc ListenableFuture historyFuture; if (cmd.getHistoryCmd() != null) { log.trace("[{}][{}] Going to process history command: {}", session.getSessionId(), cmd.getCmdId(), cmd.getHistoryCmd()); - historyFuture = handleHistoryCmd(ctx, cmd.getHistoryCmd()); + try { + historyFuture = handleHistoryCmd(ctx, cmd.getHistoryCmd()); + } catch (RuntimeException e) { + handleWsCmdRuntimeException(ctx.getSessionId(), e, cmd); + return; + } } else { historyFuture = Futures.immediateFuture(ctx); } Futures.addCallback(historyFuture, new FutureCallback<>() { @Override public void onSuccess(@Nullable TbEntityDataSubCtx theCtx) { - if (cmd.getLatestCmd() != null) { - handleLatestCmd(theCtx, cmd.getLatestCmd()); - } else if (cmd.getTsCmd() != null) { - handleTimeSeriesCmd(theCtx, cmd.getTsCmd()); - } else if (!theCtx.isInitialDataSent()) { - EntityDataUpdate update = new EntityDataUpdate(theCtx.getCmdId(), theCtx.getData(), null, theCtx.getMaxEntitiesPerDataSubscription()); - wsService.sendWsMsg(theCtx.getSessionId(), update); - theCtx.setInitialDataSent(true); + try { + if (cmd.getLatestCmd() != null) { + handleLatestCmd(theCtx, cmd.getLatestCmd()); + } else if (cmd.getTsCmd() != null) { + handleTimeSeriesCmd(theCtx, cmd.getTsCmd()); + } else if (!theCtx.isInitialDataSent()) { + EntityDataUpdate update = new EntityDataUpdate(theCtx.getCmdId(), theCtx.getData(), null, theCtx.getMaxEntitiesPerDataSubscription()); + wsService.sendWsMsg(theCtx.getSessionId(), update); + theCtx.setInitialDataSent(true); + } + } catch (RuntimeException e) { + handleWsCmdRuntimeException(theCtx.getSessionId(), e, cmd); } } @@ -230,6 +240,11 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc }, wsCallBackExecutor); } + private void handleWsCmdRuntimeException(String sessionId, RuntimeException e, EntityDataCmd cmd) { + log.debug("[{}] Failed to process ws cmd: {}", sessionId, cmd, e); + wsService.close(sessionId, CloseStatus.SERVICE_RESTARTED); + } + @Override public void handleCmd(TelemetryWebSocketSessionRef session, EntityCountCmd cmd) { TbEntityCountSubCtx ctx = getSubCtx(session.getSessionId(), cmd.getCmdId()); 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 dfc2185807..dcea91f6eb 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 @@ -307,6 +307,18 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi } } + @Override + public void close(String sessionId, CloseStatus status) { + WsSessionMetaData md = wsSessionsMap.get(sessionId); + if (md != null) { + try { + msgEndpoint.close(md.getSessionRef(), status); + } catch (IOException e) { + log.warn("[{}] Failed to send session close: {}", sessionId, e); + } + } + } + private void processSessionClose(TelemetryWebSocketSessionRef sessionRef) { String sessionId = "[" + sessionRef.getSessionId() + "]"; if (maxSubscriptionsPerTenant > 0) { diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java index c8cdbd06af..b665c3bd2c 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.telemetry; +import org.springframework.web.socket.CloseStatus; import org.thingsboard.server.service.telemetry.cmd.v2.CmdUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.DataUpdate; import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; @@ -32,4 +33,5 @@ public interface TelemetryWebSocketService { void sendWsMsg(String sessionId, CmdUpdate update); + void close(String sessionId, CloseStatus status); } diff --git a/ui-ngx/src/app/core/ws/telemetry-websocket.service.ts b/ui-ngx/src/app/core/ws/telemetry-websocket.service.ts index a3fc6eb767..2463e35d8a 100644 --- a/ui-ngx/src/app/core/ws/telemetry-websocket.service.ts +++ b/ui-ngx/src/app/core/ws/telemetry-websocket.service.ts @@ -332,7 +332,7 @@ export class TelemetryWebsocketService implements TelemetryService { } private onClose(closeEvent: CloseEvent) { - if (closeEvent && closeEvent.code > 1001 && closeEvent.code !== 1006) { + if (closeEvent && closeEvent.code > 1001 && closeEvent.code !== 1006 && closeEvent.code !== 1012) { this.showWsError(closeEvent.code, closeEvent.reason); } this.isOpening = false;