From b8472717fbfd501d4edc210940e9438a4cece1b4 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 24 Feb 2023 18:38:36 +0100 Subject: [PATCH 1/7] WebSocketConfiguration TbWebSocketHandler: fixed double instance of handler created, removed static from session maps --- .../server/config/WebSocketConfiguration.java | 17 +++++++++++------ .../controller/plugin/TbWebSocketHandler.java | 16 ++++++++-------- 2 files changed, 19 insertions(+), 14 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java b/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java index 47355d6c44..9a6ffb0e32 100644 --- a/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java +++ b/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java @@ -15,6 +15,8 @@ */ package org.thingsboard.server.config; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.http.HttpStatus; @@ -40,11 +42,15 @@ import java.util.Map; @Configuration @TbCoreComponent @EnableWebSocket +@RequiredArgsConstructor +@Slf4j public class WebSocketConfiguration implements WebSocketConfigurer { public static final String WS_PLUGIN_PREFIX = "/api/ws/plugins/"; private static final String WS_PLUGIN_MAPPING = WS_PLUGIN_PREFIX + "**"; + private final WebSocketHandler wsHandler; + @Bean public ServletServerContainerFactoryBean createWebSocketContainer() { ServletServerContainerFactoryBean container = new ServletServerContainerFactoryBean(); @@ -55,7 +61,11 @@ public class WebSocketConfiguration implements WebSocketConfigurer { @Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { - registry.addHandler(wsHandler(), WS_PLUGIN_MAPPING).setAllowedOriginPatterns("*") + if (!(wsHandler instanceof TbWebSocketHandler)) { + log.error("TbWebSocketHandler expected but [{}] provided", wsHandler); + throw new RuntimeException("TbWebSocketHandler expected but " + wsHandler + " provided"); + } + registry.addHandler(wsHandler, WS_PLUGIN_MAPPING).setAllowedOriginPatterns("*") .addInterceptors(new HttpSessionHandshakeInterceptor(), new HandshakeInterceptor() { @Override @@ -82,11 +92,6 @@ public class WebSocketConfiguration implements WebSocketConfigurer { }); } - @Bean - public WebSocketHandler wsHandler() { - return new TbWebSocketHandler(); - } - protected SecurityUser getCurrentUser() throws ThingsboardException { Authentication authentication = SecurityContextHolder.getContext().getAuthentication(); if (authentication != null && authentication.getPrincipal() instanceof SecurityUser) { 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 c689ae3593..f43607b24c 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 @@ -67,8 +67,8 @@ import static org.thingsboard.server.service.telemetry.DefaultTelemetryWebSocket @Slf4j public class TbWebSocketHandler extends TextWebSocketHandler implements TelemetryWebSocketMsgEndpoint { - private static final ConcurrentMap internalSessionMap = new ConcurrentHashMap<>(); - private static final ConcurrentMap externalSessionMap = new ConcurrentHashMap<>(); + private final ConcurrentMap internalSessionMap = new ConcurrentHashMap<>(); + private final ConcurrentMap externalSessionMap = new ConcurrentHashMap<>(); @Autowired @@ -82,13 +82,13 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr @Value("${server.ws.ping_timeout:30000}") private long pingTimeout; - private ConcurrentMap blacklistedSessions = new ConcurrentHashMap<>(); - private ConcurrentMap perSessionUpdateLimits = new ConcurrentHashMap<>(); + private final ConcurrentMap blacklistedSessions = new ConcurrentHashMap<>(); + private final ConcurrentMap perSessionUpdateLimits = new ConcurrentHashMap<>(); - private ConcurrentMap> tenantSessionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> customerSessionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> regularUserSessionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> publicUserSessionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> tenantSessionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> customerSessionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> regularUserSessionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> publicUserSessionsMap = new ConcurrentHashMap<>(); @Override public void handleTextMessage(WebSocketSession session, TextMessage message) { From 435e851d44c2dce7180db1aa1c8d2381128e1cb4 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 24 Feb 2023 16:21:44 +0100 Subject: [PATCH 2/7] test wsClient - volatile --- .../thingsboard/server/controller/AbstractControllerTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java index 06cee7b71d..32d1f65765 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java @@ -52,7 +52,7 @@ public abstract class AbstractControllerTest extends AbstractNotifyEntityTest { @LocalServerPort protected int wsPort; - private TbTestWebSocketClient wsClient; // lazy + private volatile TbTestWebSocketClient wsClient; // lazy public TbTestWebSocketClient getWsClient() { if (wsClient == null) { From cd881d8a68e5ef9c0337d18b277712333005451a Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 24 Feb 2023 16:24:01 +0100 Subject: [PATCH 3/7] TbTestWebSocketClient improvements --- .../controller/TbTestWebSocketClient.java | 21 +++++++++++++++---- 1 file changed, 17 insertions(+), 4 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java b/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java index db26e9f6df..1f9ac4ea69 100644 --- a/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java +++ b/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java @@ -47,6 +47,7 @@ import java.util.concurrent.TimeUnit; @Slf4j public class TbTestWebSocketClient extends WebSocketClient { + private static final long TIMEOUT = TimeUnit.SECONDS.toMillis(30); private volatile String lastMsg; private volatile CountDownLatch reply; private volatile CountDownLatch update; @@ -87,12 +88,14 @@ public class TbTestWebSocketClient extends WebSocketClient { } public void registerWaitForUpdate(int count) { + log.debug("registerWaitForUpdate [{}]", count); lastMsg = null; update = new CountDownLatch(count); } @Override public void send(String text) throws NotYetConnectedException { + log.debug("send [{}]", text); reply = new CountDownLatch(1); super.send(text); } @@ -110,21 +113,31 @@ public class TbTestWebSocketClient extends WebSocketClient { } public String waitForUpdate() { - return waitForUpdate(TimeUnit.SECONDS.toMillis(3)); + return waitForUpdate(TIMEOUT); } public String waitForUpdate(long ms) { + log.debug("waitForUpdate [{}]", ms); try { - update.await(ms, TimeUnit.MILLISECONDS); + if (!update.await(ms, TimeUnit.MILLISECONDS)) { + log.warn("Failed to await update (waiting time [{}]ms elapsed)", ms, new RuntimeException("stacktrace")); + } } catch (InterruptedException e) { - log.warn("Failed to await reply", e); + log.warn("Failed to await update", e); } return lastMsg; } public String waitForReply() { + return waitForReply(TIMEOUT); + } + + public String waitForReply(long ms) { + log.debug("waitForReply [{}]", ms); try { - reply.await(3, TimeUnit.SECONDS); + if (!reply.await(ms, TimeUnit.MILLISECONDS)) { + log.warn("Failed to await reply (waiting time [{}]ms elapsed)", ms, new RuntimeException("stacktrace")); + } } catch (InterruptedException e) { log.warn("Failed to await reply", e); } From 433a5dcbcf2e4802b8fe5d1a182c4d592b1e1e2c Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 24 Feb 2023 16:38:26 +0100 Subject: [PATCH 4/7] BaseWebsocketApiTest - improved assertion output and logging --- .../server/controller/BaseWebsocketApiTest.java | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java index 09fb5c505d..ce354749f3 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java @@ -548,7 +548,7 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { SingleEntityFilter entityFilter = new SingleEntityFilter(); entityFilter.setSingleEntity(tenantId); - assertThatNoException().isThrownBy(() -> { + assertThatNoException().as("subscribeForAttributes").isThrownBy(() -> { JsonNode update = getWsClient().subscribeForAttributes(tenantId, TbAttributeSubscriptionScope.SERVER_SCOPE.name(), List.of("attr")); assertThat(update.get("errorMsg").isNull()).isTrue(); assertThat(update.get("errorCode").asInt()).isEqualTo(SubscriptionErrorCode.NO_ERROR.getCode()); @@ -560,7 +560,7 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { new BaseAttributeKvEntry(System.currentTimeMillis(), new StringDataEntry("attr", expectedAttrValue)) )); JsonNode update = JacksonUtil.toJsonNode(getWsClient().waitForUpdate()); - assertThat(update).isNotNull(); + assertThat(update).as("waitForUpdate").isNotNull(); assertThat(update.get("data").get("attr").get(0).get(1).asText()).isEqualTo(expectedAttrValue); } @@ -569,15 +569,17 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { tsService.saveAndNotify(device.getTenantId(), null, device.getId(), tsData, 0, new FutureCallback() { @Override public void onSuccess(@Nullable Void result) { + log.debug("sendTelemetry callback onSuccess"); latch.countDown(); } @Override public void onFailure(Throwable t) { + log.error("Failed to send telemetry", t); latch.countDown(); } }); - latch.await(3, TimeUnit.SECONDS); + assertThat(latch.await(TIMEOUT, TimeUnit.SECONDS)).as("await sendTelemetry callback"); } private void sendAttributes(Device device, TbAttributeSubscriptionScope scope, List attrData) throws InterruptedException { @@ -589,14 +591,16 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { tsService.saveAndNotify(tenantId, entityId, scope.name(), attrData, new FutureCallback() { @Override public void onSuccess(@Nullable Void result) { + log.debug("sendAttributes callback onSuccess"); latch.countDown(); } @Override public void onFailure(Throwable t) { + log.error("Failed to sendAttributes", t); latch.countDown(); } }); - latch.await(3, TimeUnit.SECONDS); + assertThat(latch.await(TIMEOUT, TimeUnit.SECONDS)).as("await sendAttributes callback").isTrue(); } } From f706fbe78432f7e02db15218a35dae82fb5e18e3 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 24 Feb 2023 16:18:32 +0100 Subject: [PATCH 5/7] DefaultTelemetryWebSocketService - on subscribe response order fixed (will send response after subscription service called) --- .../telemetry/DefaultTelemetryWebSocketService.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) 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 a2a0d0fd8e..92be9982aa 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 @@ -450,7 +450,6 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi @Override public void onSuccess(List data) { List attributesData = data.stream().map(d -> new BasicTsKvEntry(d.getLastUpdateTs(), d)).collect(Collectors.toList()); - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); Map subState = new HashMap<>(keys.size()); keys.forEach(key -> subState.put(key, 0L)); @@ -470,6 +469,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) .build(); oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); } @Override @@ -550,7 +550,6 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi @Override public void onSuccess(List data) { List attributesData = data.stream().map(d -> new BasicTsKvEntry(d.getLastUpdateTs(), d)).collect(Collectors.toList()); - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); Map subState = new HashMap<>(attributesData.size()); attributesData.forEach(v -> subState.put(v.getKey(), v.getTs())); @@ -568,6 +567,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) .scope(scope).build(); oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); } @Override @@ -636,7 +636,6 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi FutureCallback> callback = new FutureCallback>() { @Override public void onSuccess(List data) { - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); Map subState = new HashMap<>(data.size()); data.forEach(v -> subState.put(v.getKey(), v.getTs())); @@ -650,6 +649,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .allKeys(true) .keyStates(subState).build(); oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); } @Override @@ -673,7 +673,6 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi return new FutureCallback<>() { @Override public void onSuccess(List data) { - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); Map subState = new HashMap<>(keys.size()); keys.forEach(key -> subState.put(key, startTs)); data.forEach(v -> subState.put(v.getKey(), v.getTs())); @@ -688,6 +687,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .allKeys(false) .keyStates(subState).build(); oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); } @Override From 56788caa1a78e3a2bf95d57a4044daa2314cc015 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 28 Feb 2023 18:58:39 +0100 Subject: [PATCH 6/7] DefaultTelemetryWebSocketService - removed static ObjectMapper, maps is final --- .../DefaultTelemetryWebSocketService.java | 15 +++++++-------- 1 file changed, 7 insertions(+), 8 deletions(-) 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 92be9982aa..65ab630685 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 @@ -16,7 +16,6 @@ package org.thingsboard.server.service.telemetry; import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.base.Function; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; @@ -27,6 +26,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.springframework.web.socket.CloseStatus; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.DataConstants; @@ -112,7 +112,6 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi private static final Aggregation DEFAULT_AGGREGATION = Aggregation.NONE; private static final int UNKNOWN_SUBSCRIPTION_ID = 0; private static final String PROCESSING_MSG = "[{}] Processing: {}"; - private static final ObjectMapper jsonMapper = new ObjectMapper(); private static final String FAILED_TO_FETCH_DATA = "Failed to fetch data!"; private static final String FAILED_TO_FETCH_ATTRIBUTES = "Failed to fetch attributes!"; private static final String SESSION_META_DATA_NOT_FOUND = "Session meta-data not found!"; @@ -147,10 +146,10 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi @Value("${server.ws.ping_timeout:30000}") private long pingTimeout; - private ConcurrentMap> tenantSubscriptionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> customerSubscriptionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> regularUserSubscriptionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> publicUserSubscriptionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> tenantSubscriptionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> customerSubscriptionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> regularUserSubscriptionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> publicUserSubscriptionsMap = new ConcurrentHashMap<>(); private ExecutorService executor; private String serviceId; @@ -204,7 +203,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi } try { - TelemetryPluginCmdsWrapper cmdsWrapper = jsonMapper.readValue(msg, TelemetryPluginCmdsWrapper.class); + TelemetryPluginCmdsWrapper cmdsWrapper = JacksonUtil.OBJECT_MAPPER.readValue(msg, TelemetryPluginCmdsWrapper.class); if (cmdsWrapper != null) { if (cmdsWrapper.getAttrSubCmds() != null) { cmdsWrapper.getAttrSubCmds().forEach(cmd -> { @@ -793,7 +792,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi private void sendWsMsg(TelemetryWebSocketSessionRef sessionRef, int cmdId, Object update) { try { - String msg = jsonMapper.writeValueAsString(update); + String msg = JacksonUtil.OBJECT_MAPPER.writeValueAsString(update); executor.submit(() -> { try { msgEndpoint.send(sessionRef, cmdId, msg); From 9c6c06cb0b38801ca9cb3d6e79e5eefc07c657f0 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 28 Feb 2023 20:01:14 +0100 Subject: [PATCH 7/7] maintain message order for web socket: subscribed report (always first) further updates (never first) --- .../DefaultTelemetryWebSocketService.java | 86 ++++++++++++++++--- 1 file changed, 73 insertions(+), 13 deletions(-) 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 65ab630685..ffde6192ec 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 @@ -95,6 +95,8 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import java.util.function.Consumer; import java.util.stream.Collectors; @@ -456,6 +458,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi TbAttributeSubscriptionScope scope = StringUtils.isEmpty(cmd.getScope()) ? TbAttributeSubscriptionScope.ANY_SCOPE : TbAttributeSubscriptionScope.valueOf(cmd.getScope()); + Lock subLock = new ReentrantLock(); TbAttributeSubscription sub = TbAttributeSubscription.builder() .serviceId(serviceId) .sessionId(sessionId) @@ -465,10 +468,24 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .allKeys(false) .keyStates(subState) .scope(scope) - .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) + .updateConsumer((sessionId, update) -> { + subLock.lock(); + try { + sendWsMsg(sessionId, update); + } finally { + subLock.unlock(); + } + }) .build(); - oldSubService.addSubscription(sub); - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); + + subLock.lock(); + try{ + oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); + } finally { + subLock.unlock(); + } + } @Override @@ -555,6 +572,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi TbAttributeSubscriptionScope scope = StringUtils.isEmpty(cmd.getScope()) ? TbAttributeSubscriptionScope.ANY_SCOPE : TbAttributeSubscriptionScope.valueOf(cmd.getScope()); + Lock subLock = new ReentrantLock(); TbAttributeSubscription sub = TbAttributeSubscription.builder() .serviceId(serviceId) .sessionId(sessionId) @@ -563,10 +581,24 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .entityId(entityId) .allKeys(true) .keyStates(subState) - .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) - .scope(scope).build(); - oldSubService.addSubscription(sub); - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); + .updateConsumer((sessionId, update) -> { + subLock.lock(); + try { + sendWsMsg(sessionId, update); + } finally { + subLock.unlock(); + } + }) + .scope(scope) + .build(); + + subLock.lock(); + try { + oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); + } finally { + subLock.unlock(); + } } @Override @@ -638,17 +670,31 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi Map subState = new HashMap<>(data.size()); data.forEach(v -> subState.put(v.getKey(), v.getTs())); + Lock subLock = new ReentrantLock(); TbTimeseriesSubscription sub = TbTimeseriesSubscription.builder() .serviceId(serviceId) .sessionId(sessionId) .subscriptionId(cmd.getCmdId()) .tenantId(sessionRef.getSecurityCtx().getTenantId()) .entityId(entityId) - .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) + .updateConsumer((sessionId, update) -> { + subLock.lock(); + try { + sendWsMsg(sessionId, update); + } finally { + subLock.unlock(); + } + }) .allKeys(true) .keyStates(subState).build(); - oldSubService.addSubscription(sub); - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); + + subLock.lock(); + try { + oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); + } finally { + subLock.unlock(); + } } @Override @@ -676,17 +722,31 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi keys.forEach(key -> subState.put(key, startTs)); data.forEach(v -> subState.put(v.getKey(), v.getTs())); + Lock subLock = new ReentrantLock(); TbTimeseriesSubscription sub = TbTimeseriesSubscription.builder() .serviceId(serviceId) .sessionId(sessionId) .subscriptionId(cmd.getCmdId()) .tenantId(sessionRef.getSecurityCtx().getTenantId()) .entityId(entityId) - .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) + .updateConsumer((sessionId, update) -> { + subLock.lock(); + try { + sendWsMsg(sessionId, update); + } finally { + subLock.unlock(); + } + }) .allKeys(false) .keyStates(subState).build(); - oldSubService.addSubscription(sub); - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); + + subLock.lock(); + try{ + oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); + } finally { + subLock.unlock(); + } } @Override