Browse Source

Merge pull request #10837 from thingsboard/master-hotfix-364

hotfix/3.6.4
pull/10841/head
Viacheslav Klimov 2 years ago
committed by GitHub
parent
commit
fc73177123
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 12
      application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java
  2. 12
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  3. 105
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java
  4. 39
      application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionModificationResult.java
  5. 13
      application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java
  6. 2
      application/src/main/java/org/thingsboard/server/service/ws/WebSocketMsgEndpoint.java
  7. 3
      application/src/main/java/org/thingsboard/server/service/ws/WebSocketService.java
  8. 2
      application/src/main/java/org/thingsboard/server/service/ws/WebSocketSessionRef.java

12
application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java

@ -539,6 +539,18 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements WebSocke
} }
} }
@Override
public boolean isOpen(String externalId) {
String internalId = externalSessionMap.get(externalId);
if (internalId != null) {
SessionMetaData sessionMd = getSessionMd(internalId);
if (sessionMd != null) {
return sessionMd.session.isOpen();
}
}
return false;
}
private boolean checkLimits(WebSocketSession session, WebSocketSessionRef sessionRef) throws IOException { private boolean checkLimits(WebSocketSession session, WebSocketSessionRef sessionRef) throws IOException {
var tenantProfileConfiguration = getTenantProfileConfiguration(sessionRef); var tenantProfileConfiguration = getTenantProfileConfiguration(sessionRef);
if (tenantProfileConfiguration == null) { if (tenantProfileConfiguration == null) {

12
application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java

@ -95,7 +95,8 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
private static final int DEFAULT_LIMIT = 100; private static final int DEFAULT_LIMIT = 100;
private final Map<String, Map<Integer, TbAbstractSubCtx>> subscriptionsBySessionId = new ConcurrentHashMap<>(); private final Map<String, Map<Integer, TbAbstractSubCtx>> subscriptionsBySessionId = new ConcurrentHashMap<>();
@Autowired @Lazy @Autowired
@Lazy
private WebSocketService wsService; private WebSocketService wsService;
@Autowired @Autowired
@ -728,7 +729,14 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
public void cancelAllSessionSubscriptions(String sessionId) { public void cancelAllSessionSubscriptions(String sessionId) {
Map<Integer, TbAbstractSubCtx> sessionSubs = subscriptionsBySessionId.remove(sessionId); Map<Integer, TbAbstractSubCtx> sessionSubs = subscriptionsBySessionId.remove(sessionId);
if (sessionSubs != null) { if (sessionSubs != null) {
sessionSubs.values().forEach(this::cleanupAndCancel); sessionSubs.values().forEach(sub -> {
try {
cleanupAndCancel(sub);
} catch (Exception e) {
log.warn("[{}] Failed to remove subscription {} due to ", sub.getTenantId(), sub, e);
}
}
);
} }
} }

105
application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java

@ -45,6 +45,7 @@ import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.discovery.event.ClusterTopologyChangeEvent; import org.thingsboard.server.queue.discovery.event.ClusterTopologyChangeEvent;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.ws.WebSocketService;
import org.thingsboard.server.service.ws.notification.sub.NotificationRequestUpdate; import org.thingsboard.server.service.ws.notification.sub.NotificationRequestUpdate;
import org.thingsboard.server.service.ws.notification.sub.NotificationsSubscriptionUpdate; import org.thingsboard.server.service.ws.notification.sub.NotificationsSubscriptionUpdate;
import org.thingsboard.server.service.ws.telemetry.sub.AlarmSubscriptionUpdate; import org.thingsboard.server.service.ws.telemetry.sub.AlarmSubscriptionUpdate;
@ -62,6 +63,8 @@ import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock; import java.util.concurrent.locks.ReentrantLock;
@ -84,18 +87,21 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
private final PartitionService partitionService; private final PartitionService partitionService;
private final TbClusterService clusterService; private final TbClusterService clusterService;
private final SubscriptionManagerService subscriptionManagerService; private final SubscriptionManagerService subscriptionManagerService;
private final WebSocketService webSocketService;
private ExecutorService tsCallBackExecutor; private ExecutorService tsCallBackExecutor;
private ScheduledExecutorService staleSessionCleanupExecutor;
public DefaultTbLocalSubscriptionService(AttributesService attrService, TimeseriesService tsService, TbServiceInfoProvider serviceInfoProvider, public DefaultTbLocalSubscriptionService(AttributesService attrService, TimeseriesService tsService, TbServiceInfoProvider serviceInfoProvider,
PartitionService partitionService, TbClusterService clusterService, PartitionService partitionService, TbClusterService clusterService,
@Lazy SubscriptionManagerService subscriptionManagerService) { @Lazy SubscriptionManagerService subscriptionManagerService, @Lazy WebSocketService webSocketService) {
this.attrService = attrService; this.attrService = attrService;
this.tsService = tsService; this.tsService = tsService;
this.serviceInfoProvider = serviceInfoProvider; this.serviceInfoProvider = serviceInfoProvider;
this.partitionService = partitionService; this.partitionService = partitionService;
this.clusterService = clusterService; this.clusterService = clusterService;
this.subscriptionManagerService = subscriptionManagerService; this.subscriptionManagerService = subscriptionManagerService;
this.webSocketService = webSocketService;
} }
private String serviceId; private String serviceId;
@ -108,6 +114,8 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
subscriptionUpdateExecutor = ThingsBoardExecutors.newWorkStealingPool(20, getClass()); subscriptionUpdateExecutor = ThingsBoardExecutors.newWorkStealingPool(20, getClass());
tsCallBackExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("ts-sub-callback")); tsCallBackExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("ts-sub-callback"));
serviceId = serviceInfoProvider.getServiceId(); serviceId = serviceInfoProvider.getServiceId();
staleSessionCleanupExecutor = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("stale-session-cleanup"));
staleSessionCleanupExecutor.scheduleWithFixedDelay(this::cleanupStaleSessions, 60, 60, TimeUnit.SECONDS);
} }
@PreDestroy @PreDestroy
@ -118,6 +126,9 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
if (tsCallBackExecutor != null) { if (tsCallBackExecutor != null) {
tsCallBackExecutor.shutdownNow(); tsCallBackExecutor.shutdownNow();
} }
if (staleSessionCleanupExecutor != null) {
staleSessionCleanupExecutor.shutdownNow();
}
} }
@Override @Override
@ -157,9 +168,18 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
TenantId tenantId = subscription.getTenantId(); TenantId tenantId = subscription.getTenantId();
EntityId entityId = subscription.getEntityId(); EntityId entityId = subscription.getEntityId();
log.debug("[{}][{}] Register subscription: {}", tenantId, entityId, subscription); log.debug("[{}][{}] Register subscription: {}", tenantId, entityId, subscription);
Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.computeIfAbsent(subscription.getSessionId(), k -> new ConcurrentHashMap<>()); SubscriptionModificationResult result;
sessionSubscriptions.put(subscription.getSubscriptionId(), subscription); subsLock.lock();
modifySubscription(tenantId, entityId, subscription, true); try {
Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.computeIfAbsent(subscription.getSessionId(), k -> new ConcurrentHashMap<>());
sessionSubscriptions.put(subscription.getSubscriptionId(), subscription);
result = modifySubscription(tenantId, entityId, subscription, true);
} finally {
subsLock.unlock();
}
if (result.hasEvent()) {
pushSubscriptionEvent(result);
}
} }
@Override @Override
@ -195,33 +215,49 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
@Override @Override
public void cancelSubscription(String sessionId, int subscriptionId) { public void cancelSubscription(String sessionId, int subscriptionId) {
log.debug("[{}][{}] Going to remove subscription.", sessionId, subscriptionId); log.debug("[{}][{}] Going to remove subscription.", sessionId, subscriptionId);
Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.get(sessionId); SubscriptionModificationResult result = null;
if (sessionSubscriptions != null) { subsLock.lock();
TbSubscription<?> subscription = sessionSubscriptions.remove(subscriptionId); try {
if (subscription != null) { Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.get(sessionId);
if (sessionSubscriptions.isEmpty()) { if (sessionSubscriptions != null) {
subscriptionsBySessionId.remove(sessionId); TbSubscription<?> subscription = sessionSubscriptions.remove(subscriptionId);
if (subscription != null) {
if (sessionSubscriptions.isEmpty()) {
subscriptionsBySessionId.remove(sessionId);
}
result = modifySubscription(subscription.getTenantId(), subscription.getEntityId(), subscription, false);
} else {
log.debug("[{}][{}] Subscription not found!", sessionId, subscriptionId);
} }
modifySubscription(subscription.getTenantId(), subscription.getEntityId(), subscription, false);
} else { } else {
log.debug("[{}][{}] Subscription not found!", sessionId, subscriptionId); log.debug("[{}] No session subscriptions found!", sessionId);
} }
} else { } finally {
log.debug("[{}] No session subscriptions found!", sessionId); subsLock.unlock();
}
if (result != null && result.hasEvent()) {
pushSubscriptionEvent(result);
} }
} }
@Override @Override
public void cancelAllSessionSubscriptions(String sessionId) { public void cancelAllSessionSubscriptions(String sessionId) {
log.debug("[{}] Going to remove session subscriptions.", sessionId); log.debug("[{}] Going to remove session subscriptions.", sessionId);
Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.remove(sessionId); List<SubscriptionModificationResult> results = new ArrayList<>();
if (sessionSubscriptions != null) { subsLock.lock();
for (TbSubscription<?> subscription : sessionSubscriptions.values()) { try {
modifySubscription(subscription.getTenantId(), subscription.getEntityId(), subscription, false); Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.remove(sessionId);
if (sessionSubscriptions != null) {
for (TbSubscription<?> subscription : sessionSubscriptions.values()) {
results.add(modifySubscription(subscription.getTenantId(), subscription.getEntityId(), subscription, false));
}
} else {
log.debug("[{}] No session subscriptions found!", sessionId);
} }
} else { } finally {
log.debug("[{}] No session subscriptions found!", sessionId); subsLock.unlock();
} }
results.stream().filter(SubscriptionModificationResult::hasEvent).forEach(this::pushSubscriptionEvent);
} }
@Override @Override
@ -384,10 +420,9 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
callback.onSuccess(); callback.onSuccess();
} }
private void modifySubscription(TenantId tenantId, EntityId entityId, TbSubscription<?> subscription, boolean add) { private SubscriptionModificationResult modifySubscription(TenantId tenantId, EntityId entityId, TbSubscription<?> subscription, boolean add) {
TbSubscription<?> missedUpdatesCandidate = null; TbSubscription<?> missedUpdatesCandidate = null;
TbEntitySubEvent event; TbEntitySubEvent event = null;
subsLock.lock();
try { try {
TbEntityLocalSubsInfo entitySubs = subscriptionsByEntityId.computeIfAbsent(entityId.getId(), id -> new TbEntityLocalSubsInfo(tenantId, entityId)); TbEntityLocalSubsInfo entitySubs = subscriptionsByEntityId.computeIfAbsent(entityId.getId(), id -> new TbEntityLocalSubsInfo(tenantId, entityId));
event = add ? entitySubs.add(subscription) : entitySubs.remove(subscription); event = add ? entitySubs.add(subscription) : entitySubs.remove(subscription);
@ -397,17 +432,23 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
} else if (add) { } else if (add) {
missedUpdatesCandidate = entitySubs.registerPendingSubscription(subscription, event); missedUpdatesCandidate = entitySubs.registerPendingSubscription(subscription, event);
} }
} finally { } catch (Exception e) {
subsLock.unlock(); log.warn("[{}][{}] Failed to {} subscription {} due to ", tenantId, entityId, add ? "add" : "remove", subscription, e);
} }
if (event != null) { return new SubscriptionModificationResult(tenantId, entityId, subscription, missedUpdatesCandidate, event);
log.trace("[{}][{}][{}] Event: {}", tenantId, entityId, subscription.getSubscriptionId(), event); }
pushSubEventToManagerService(tenantId, entityId, event);
private void pushSubscriptionEvent(SubscriptionModificationResult modificationResult) {
try {
TbEntitySubEvent event = modificationResult.getEvent();
log.trace("[{}][{}][{}] Event: {}", modificationResult.getTenantId(), modificationResult.getEntityId(), modificationResult.getSubscription().getSubscriptionId(), event);
pushSubEventToManagerService(modificationResult.getTenantId(), modificationResult.getEntityId(), event);
TbSubscription<?> missedUpdatesCandidate = modificationResult.getMissedUpdatesCandidate();
if (missedUpdatesCandidate != null) { if (missedUpdatesCandidate != null) {
checkMissedUpdates(missedUpdatesCandidate); checkMissedUpdates(missedUpdatesCandidate);
} }
} else { } catch (Exception e) {
log.trace("[{}][{}][{}] No changes detected.", tenantId, entityId, subscription.getSubscriptionId()); log.warn("[{}][{}] Failed to push subscription event {} due to ", modificationResult.getTenantId(), modificationResult.getEntityId(), modificationResult.getEvent(), e);
} }
} }
@ -518,4 +559,8 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
} }
} }
private void cleanupStaleSessions() {
subscriptionsBySessionId.keySet().forEach(webSocketService::cleanupIfStale);
}
} }

39
application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionModificationResult.java

@ -0,0 +1,39 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.subscription;
import lombok.Builder;
import lombok.Data;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
/**
* The modification result of entity subscription
*/
@Builder
@Data
public class SubscriptionModificationResult {
private TenantId tenantId;
private EntityId entityId;
private TbSubscription<?> subscription;
private TbSubscription<?> missedUpdatesCandidate;
private TbEntitySubEvent event;
public boolean hasEvent() {
return event != null;
}
}

13
application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java

@ -197,7 +197,8 @@ public class DefaultWebSocketService implements WebSocketService {
wsSessionsMap.put(sessionId, new WsSessionMetaData(sessionRef)); wsSessionsMap.put(sessionId, new WsSessionMetaData(sessionRef));
break; break;
case ERROR: case ERROR:
log.debug("[{}] Unknown websocket session error: {}. ", sessionId, event.getError().orElse(null)); log.debug("[{}] Unknown websocket session error: ", sessionId,
event.getError().orElse(new RuntimeException("No error specified")));
break; break;
case CLOSED: case CLOSED:
wsSessionsMap.remove(sessionId); wsSessionsMap.remove(sessionId);
@ -294,6 +295,16 @@ public class DefaultWebSocketService implements WebSocketService {
} }
} }
@Override
public void cleanupIfStale(String sessionId) {
if (!msgEndpoint.isOpen(sessionId)) {
log.info("[{}] Cleaning up stale session ", sessionId);
wsSessionsMap.remove(sessionId);
oldSubService.cancelAllSessionSubscriptions(sessionId);
entityDataSubService.cancelAllSessionSubscriptions(sessionId);
}
}
private void processSessionClose(WebSocketSessionRef sessionRef) { private void processSessionClose(WebSocketSessionRef sessionRef) {
var tenantProfileConfiguration = getTenantProfileConfiguration(sessionRef); var tenantProfileConfiguration = getTenantProfileConfiguration(sessionRef);
if (tenantProfileConfiguration != null) { if (tenantProfileConfiguration != null) {

2
application/src/main/java/org/thingsboard/server/service/ws/WebSocketMsgEndpoint.java

@ -29,4 +29,6 @@ public interface WebSocketMsgEndpoint {
void sendPing(WebSocketSessionRef sessionRef, long currentTime) throws IOException; void sendPing(WebSocketSessionRef sessionRef, long currentTime) throws IOException;
void close(WebSocketSessionRef sessionRef, CloseStatus withReason) throws IOException; void close(WebSocketSessionRef sessionRef, CloseStatus withReason) throws IOException;
boolean isOpen(String sessionId);
} }

3
application/src/main/java/org/thingsboard/server/service/ws/WebSocketService.java

@ -36,4 +36,7 @@ public interface WebSocketService {
void sendError(WebSocketSessionRef sessionRef, int subId, SubscriptionErrorCode errorCode, String errorMsg); void sendError(WebSocketSessionRef sessionRef, int subId, SubscriptionErrorCode errorCode, String errorMsg);
void close(String sessionId, CloseStatus status); void close(String sessionId, CloseStatus status);
void cleanupIfStale(String sessionId);
} }

2
application/src/main/java/org/thingsboard/server/service/ws/WebSocketSessionRef.java

@ -33,7 +33,7 @@ public class WebSocketSessionRef {
private static final long serialVersionUID = 1L; private static final long serialVersionUID = 1L;
private final String sessionId; private final String sessionId;
private SecurityUser securityCtx; private volatile SecurityUser securityCtx;
private final InetSocketAddress localAddress; private final InetSocketAddress localAddress;
private final InetSocketAddress remoteAddress; private final InetSocketAddress remoteAddress;
private final WebSocketSessionType sessionType; private final WebSocketSessionType sessionType;

Loading…
Cancel
Save