Browse Source

Minor refactoring for subscription service

pull/10837/head
ViacheslavKlimov 2 years ago
parent
commit
cdc9bda92c
  1. 38
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionModificationResult.java

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

@ -168,17 +168,17 @@ 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);
ModifySubscriptionResult modifySubscriptionResult; SubscriptionModificationResult result;
subsLock.lock(); subsLock.lock();
try { try {
Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.computeIfAbsent(subscription.getSessionId(), k -> new ConcurrentHashMap<>()); Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.computeIfAbsent(subscription.getSessionId(), k -> new ConcurrentHashMap<>());
sessionSubscriptions.put(subscription.getSubscriptionId(), subscription); sessionSubscriptions.put(subscription.getSubscriptionId(), subscription);
modifySubscriptionResult = modifySubscription(tenantId, entityId, subscription, true); result = modifySubscription(tenantId, entityId, subscription, true);
} finally { } finally {
subsLock.unlock(); subsLock.unlock();
} }
if (modifySubscriptionResult.hasEvent()) { if (result.hasEvent()) {
pushSubscriptionEvent(modifySubscriptionResult); pushSubscriptionEvent(result);
} }
} }
@ -215,7 +215,7 @@ 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);
ModifySubscriptionResult modifySubscriptionResult = null; SubscriptionModificationResult result = null;
subsLock.lock(); subsLock.lock();
try { try {
Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.get(sessionId); Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.get(sessionId);
@ -225,7 +225,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
if (sessionSubscriptions.isEmpty()) { if (sessionSubscriptions.isEmpty()) {
subscriptionsBySessionId.remove(sessionId); subscriptionsBySessionId.remove(sessionId);
} }
modifySubscriptionResult = modifySubscription(subscription.getTenantId(), subscription.getEntityId(), subscription, false); result = modifySubscription(subscription.getTenantId(), subscription.getEntityId(), subscription, false);
} else { } else {
log.debug("[{}][{}] Subscription not found!", sessionId, subscriptionId); log.debug("[{}][{}] Subscription not found!", sessionId, subscriptionId);
} }
@ -235,21 +235,21 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
} finally { } finally {
subsLock.unlock(); subsLock.unlock();
} }
if (modifySubscriptionResult != null && modifySubscriptionResult.hasEvent()) { if (result != null && result.hasEvent()) {
pushSubscriptionEvent(modifySubscriptionResult); 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);
List<ModifySubscriptionResult> result = new ArrayList<>(); List<SubscriptionModificationResult> results = new ArrayList<>();
subsLock.lock(); subsLock.lock();
try { try {
Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.remove(sessionId); Map<Integer, TbSubscription<?>> sessionSubscriptions = subscriptionsBySessionId.remove(sessionId);
if (sessionSubscriptions != null) { if (sessionSubscriptions != null) {
for (TbSubscription<?> subscription : sessionSubscriptions.values()) { for (TbSubscription<?> subscription : sessionSubscriptions.values()) {
result.add(modifySubscription(subscription.getTenantId(), subscription.getEntityId(), subscription, false)); results.add(modifySubscription(subscription.getTenantId(), subscription.getEntityId(), subscription, false));
} }
} else { } else {
log.debug("[{}] No session subscriptions found!", sessionId); log.debug("[{}] No session subscriptions found!", sessionId);
@ -257,7 +257,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
} finally { } finally {
subsLock.unlock(); subsLock.unlock();
} }
result.stream().filter(ModifySubscriptionResult::hasEvent).forEach(this::pushSubscriptionEvent); results.stream().filter(SubscriptionModificationResult::hasEvent).forEach(this::pushSubscriptionEvent);
} }
@Override @Override
@ -420,7 +420,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
callback.onSuccess(); callback.onSuccess();
} }
private ModifySubscriptionResult 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 = null; TbEntitySubEvent event = null;
try { try {
@ -435,20 +435,20 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}][{}] Failed to {} subscription {} due to ", tenantId, entityId, add ? "add" : "remove", subscription, e); log.warn("[{}][{}] Failed to {} subscription {} due to ", tenantId, entityId, add ? "add" : "remove", subscription, e);
} }
return new ModifySubscriptionResult(tenantId, entityId, subscription, missedUpdatesCandidate, event); return new SubscriptionModificationResult(tenantId, entityId, subscription, missedUpdatesCandidate, event);
} }
private void pushSubscriptionEvent(ModifySubscriptionResult modifySubResult) { private void pushSubscriptionEvent(SubscriptionModificationResult modificationResult) {
try { try {
TbEntitySubEvent event = modifySubResult.getEvent(); TbEntitySubEvent event = modificationResult.getEvent();
log.trace("[{}][{}][{}] Event: {}", modifySubResult.getTenantId(), modifySubResult.getEntityId(), modifySubResult.getSubscription().getSubscriptionId(), event); log.trace("[{}][{}][{}] Event: {}", modificationResult.getTenantId(), modificationResult.getEntityId(), modificationResult.getSubscription().getSubscriptionId(), event);
pushSubEventToManagerService(modifySubResult.getTenantId(), modifySubResult.getEntityId(), event); pushSubEventToManagerService(modificationResult.getTenantId(), modificationResult.getEntityId(), event);
TbSubscription<?> missedUpdatesCandidate = modifySubResult.getMissedUpdatesCandidate(); TbSubscription<?> missedUpdatesCandidate = modificationResult.getMissedUpdatesCandidate();
if (missedUpdatesCandidate != null) { if (missedUpdatesCandidate != null) {
checkMissedUpdates(missedUpdatesCandidate); checkMissedUpdates(missedUpdatesCandidate);
} }
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}][{}] Failed to push subscription event {} due to ", modifySubResult.getTenantId(), modifySubResult.getEntityId(), modifySubResult.getEvent(), e); log.warn("[{}][{}] Failed to push subscription event {} due to ", modificationResult.getTenantId(), modificationResult.getEntityId(), modificationResult.getEvent(), e);
} }
} }

2
application/src/main/java/org/thingsboard/server/service/subscription/ModifySubscriptionResult.java → application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionModificationResult.java

@ -25,7 +25,7 @@ import org.thingsboard.server.common.data.id.TenantId;
*/ */
@Builder @Builder
@Data @Data
public class ModifySubscriptionResult { public class SubscriptionModificationResult {
private TenantId tenantId; private TenantId tenantId;
private EntityId entityId; private EntityId entityId;
Loading…
Cancel
Save