Browse Source

Suppress and handle WS exceptions on cluster rebalancing

pull/5921/head
Andrii Shvaika 5 years ago
parent
commit
f72c058030
  1. 33
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  2. 12
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java
  3. 2
      application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java
  4. 2
      ui-ngx/src/app/core/ws/telemetry-websocket.service.ts

33
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.context.annotation.Lazy;
import org.springframework.scheduling.annotation.Scheduled; import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.web.socket.CloseStatus;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -205,21 +206,30 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
ListenableFuture<TbEntityDataSubCtx> historyFuture; ListenableFuture<TbEntityDataSubCtx> historyFuture;
if (cmd.getHistoryCmd() != null) { if (cmd.getHistoryCmd() != null) {
log.trace("[{}][{}] Going to process history command: {}", session.getSessionId(), cmd.getCmdId(), cmd.getHistoryCmd()); 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 { } else {
historyFuture = Futures.immediateFuture(ctx); historyFuture = Futures.immediateFuture(ctx);
} }
Futures.addCallback(historyFuture, new FutureCallback<>() { Futures.addCallback(historyFuture, new FutureCallback<>() {
@Override @Override
public void onSuccess(@Nullable TbEntityDataSubCtx theCtx) { public void onSuccess(@Nullable TbEntityDataSubCtx theCtx) {
if (cmd.getLatestCmd() != null) { try {
handleLatestCmd(theCtx, cmd.getLatestCmd()); if (cmd.getLatestCmd() != null) {
} else if (cmd.getTsCmd() != null) { handleLatestCmd(theCtx, cmd.getLatestCmd());
handleTimeSeriesCmd(theCtx, cmd.getTsCmd()); } else if (cmd.getTsCmd() != null) {
} else if (!theCtx.isInitialDataSent()) { handleTimeSeriesCmd(theCtx, cmd.getTsCmd());
EntityDataUpdate update = new EntityDataUpdate(theCtx.getCmdId(), theCtx.getData(), null, theCtx.getMaxEntitiesPerDataSubscription()); } else if (!theCtx.isInitialDataSent()) {
wsService.sendWsMsg(theCtx.getSessionId(), update); EntityDataUpdate update = new EntityDataUpdate(theCtx.getCmdId(), theCtx.getData(), null, theCtx.getMaxEntitiesPerDataSubscription());
theCtx.setInitialDataSent(true); 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); }, 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 @Override
public void handleCmd(TelemetryWebSocketSessionRef session, EntityCountCmd cmd) { public void handleCmd(TelemetryWebSocketSessionRef session, EntityCountCmd cmd) {
TbEntityCountSubCtx ctx = getSubCtx(session.getSessionId(), cmd.getCmdId()); TbEntityCountSubCtx ctx = getSubCtx(session.getSessionId(), cmd.getCmdId());

12
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) { private void processSessionClose(TelemetryWebSocketSessionRef sessionRef) {
String sessionId = "[" + sessionRef.getSessionId() + "]"; String sessionId = "[" + sessionRef.getSessionId() + "]";
if (maxSubscriptionsPerTenant > 0) { if (maxSubscriptionsPerTenant > 0) {

2
application/src/main/java/org/thingsboard/server/service/telemetry/TelemetryWebSocketService.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.service.telemetry; 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.CmdUpdate;
import org.thingsboard.server.service.telemetry.cmd.v2.DataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.DataUpdate;
import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate;
@ -32,4 +33,5 @@ public interface TelemetryWebSocketService {
void sendWsMsg(String sessionId, CmdUpdate update); void sendWsMsg(String sessionId, CmdUpdate update);
void close(String sessionId, CloseStatus status);
} }

2
ui-ngx/src/app/core/ws/telemetry-websocket.service.ts

@ -332,7 +332,7 @@ export class TelemetryWebsocketService implements TelemetryService {
} }
private onClose(closeEvent: CloseEvent) { 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.showWsError(closeEvent.code, closeEvent.reason);
} }
this.isOpening = false; this.isOpening = false;

Loading…
Cancel
Save