From 70a60abc0229ea59ca5249589eff8c88852f4351 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Fri, 6 Sep 2024 16:45:18 +0300 Subject: [PATCH 1/4] Fixed issue with Edge Root Rule Chain cache for related edges --- .../service/edge/EdgeEventSourcingListener.java | 11 +++++++++++ .../service/edge/RelatedEdgesSourcingListener.java | 2 ++ .../server/dao/edge/BaseRelatedEdgesService.java | 4 ++++ .../server/dao/rule/BaseRuleChainService.java | 6 ++---- 4 files changed, 19 insertions(+), 4 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java index 840750a781..3134a9d0dc 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java @@ -34,8 +34,10 @@ import org.thingsboard.server.common.data.alarm.AlarmComment; import org.thingsboard.server.common.data.alarm.EntityAlarm; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.domain.Domain; +import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; +import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; @@ -134,6 +136,15 @@ public class EdgeEventSourcingListener { return; } try { + if (event.getEntityId().getEntityType().equals(EntityType.RULE_CHAIN) && event.getEdgeId() != null && event.getActionType().equals(ActionType.ASSIGNED_TO_EDGE)) { + try { + Edge edge = JacksonUtil.fromString(event.getBody(), Edge.class); + if (edge != null && new RuleChainId(event.getEntityId().getId()).equals(edge.getRootRuleChainId())) { + log.trace("skipping ASSIGNED_TO_EDGE event of RULE_CHAIN entity in case Edge Root Rule Chain: {}", event); + return; + } + } catch (Exception ignored) {} + } log.trace("[{}] ActionEntityEvent called: {}", event.getTenantId(), event); tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), event.getEdgeId(), event.getEntityId(), event.getBody(), null, EdgeUtils.getEdgeEventActionTypeByActionType(event.getActionType()), diff --git a/application/src/main/java/org/thingsboard/server/service/edge/RelatedEdgesSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/edge/RelatedEdgesSourcingListener.java index 851f59c3ab..2132b9e4a6 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/RelatedEdgesSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/RelatedEdgesSourcingListener.java @@ -55,6 +55,7 @@ public class RelatedEdgesSourcingListener { @TransactionalEventListener(fallbackExecution = true) public void handleEvent(ActionEntityEvent event) { executorService.submit(() -> { + log.trace("[{}] ActionEntityEvent called: {}", event.getTenantId(), event); try { switch (event.getActionType()) { case ASSIGNED_TO_EDGE, UNASSIGNED_FROM_EDGE -> @@ -69,6 +70,7 @@ public class RelatedEdgesSourcingListener { @TransactionalEventListener(fallbackExecution = true) public void handleEvent(DeleteEntityEvent event) { executorService.submit(() -> { + log.trace("[{}] DeleteEntityEvent called: {}", event.getTenantId(), event); try { relatedEdgesService.publishRelatedEdgeIdsEvictEvent(event.getTenantId(), event.getEntityId()); } catch (Exception e) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/BaseRelatedEdgesService.java b/dao/src/main/java/org/thingsboard/server/dao/edge/BaseRelatedEdgesService.java index 265831640c..bedee455f4 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/BaseRelatedEdgesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/BaseRelatedEdgesService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.edge; +import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; @@ -30,6 +31,7 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.entity.AbstractCachedEntityService; @Service +@Slf4j public class BaseRelatedEdgesService extends AbstractCachedEntityService implements RelatedEdgesService { public static final int RELATED_EDGES_CACHE_ITEMS = 1000; @@ -47,6 +49,7 @@ public class BaseRelatedEdgesService extends AbstractCachedEntityService findEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink) { + log.trace("Executing findEdgeIdsByEntityId, tenantId [{}], entityId [{}], pageLink [{}]", tenantId, entityId, pageLink); if (!pageLink.equals(FIRST_PAGE)) { return edgeService.findEdgeIdsByTenantIdAndEntityId(tenantId, entityId, pageLink); } @@ -56,6 +59,7 @@ public class BaseRelatedEdgesService extends AbstractCachedEntityService Date: Mon, 9 Sep 2024 10:55:03 +0300 Subject: [PATCH 2/4] EdgeEventSourcingListener - return in case errors during processing ASSIGNED_TO_EDGE rule chain event --- .../server/service/edge/EdgeEventSourcingListener.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java index 3134a9d0dc..17ceb17ed1 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java @@ -140,10 +140,12 @@ public class EdgeEventSourcingListener { try { Edge edge = JacksonUtil.fromString(event.getBody(), Edge.class); if (edge != null && new RuleChainId(event.getEntityId().getId()).equals(edge.getRootRuleChainId())) { - log.trace("skipping ASSIGNED_TO_EDGE event of RULE_CHAIN entity in case Edge Root Rule Chain: {}", event); + log.trace("[{}] skipping ASSIGNED_TO_EDGE event of RULE_CHAIN entity in case Edge Root Rule Chain: {}", event.getTenantId(), event); return; } - } catch (Exception ignored) {} + } catch (Exception ignored) { + return; + } } log.trace("[{}] ActionEntityEvent called: {}", event.getTenantId(), event); tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), event.getEdgeId(), event.getEntityId(), From ab2fb38e050f479ac4b8b00cca4825c31ba62f8c Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 9 Sep 2024 13:41:40 +0300 Subject: [PATCH 3/4] RuleChainEdgeTest - added testUpdateRootRuleChain --- .../server/edge/AbstractEdgeTest.java | 18 +++++++++--------- .../server/edge/RuleChainEdgeTest.java | 15 +++++++++++++++ 2 files changed, 24 insertions(+), 9 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java index b3a1dd790e..34f7da02e6 100644 --- a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java @@ -163,15 +163,10 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { } private RuleChainId getEdgeRootRuleChainId() throws Exception { - List edgeRuleChains = doGetTypedWithPageLink("/api/edge/" + edge.getUuidId() + "/ruleChains?", - new TypeReference>() { - }, new PageLink(100)).getData(); - for (RuleChain edgeRuleChain : edgeRuleChains) { - if (edgeRuleChain.isRoot()) { - return edgeRuleChain.getId(); - } - } - throw new RuntimeException("Root rule chain not found"); + return doGetTypedWithPageLink("/api/ruleChains?type={type}&", new TypeReference>() {}, + new PageLink(100, 0, "Edge Root Rule Chain"), + "EDGE") + .getData().get(0).getId(); } @After @@ -196,6 +191,11 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { Asset savedAsset = saveAsset("Edge Asset 1"); + RuleChainId rootRuleChainId = getEdgeRootRuleChainId(); + RuleChainMetaData rootRuleChainMetadata = doGet("/api/ruleChain/" + rootRuleChainId.getId().toString() + "/metadata", RuleChainMetaData.class); + rootRuleChainMetadata.getNodes().forEach(n -> n.setDebugMode(true)); + doPost("/api/ruleChain/metadata", rootRuleChainMetadata, RuleChainMetaData.class); + edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class); doPost("/api/edge/" + edge.getUuidId() diff --git a/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java index c694df84f8..6dafeefa80 100644 --- a/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java @@ -185,6 +185,21 @@ public class RuleChainEdgeTest extends AbstractEdgeTest { return doPost("/api/ruleChain/metadata", ruleChainMetaData, RuleChainMetaData.class); } + @Test + public void testUpdateRootRuleChain() throws Exception { + RuleChainMetaData rootRuleChainMetadata = doGet("/api/ruleChain/" + edge.getRootRuleChainId().getId().toString() + "/metadata", RuleChainMetaData.class); + + rootRuleChainMetadata.getNodes().forEach(n -> n.setDebugMode(true)); + edgeImitator.expectMessageAmount(2); + doPost("/api/ruleChain/metadata", rootRuleChainMetadata, RuleChainMetaData.class); + Assert.assertTrue(edgeImitator.waitForMessages()); + + Optional ruleChainUpdateMsgOpt = edgeImitator.findMessageByType(RuleChainUpdateMsg.class); + Assert.assertTrue(ruleChainUpdateMsgOpt.isPresent()); + Optional ruleChainMetadataUpdateMsgOpt = edgeImitator.findMessageByType(RuleChainMetadataUpdateMsg.class); + Assert.assertTrue(ruleChainMetadataUpdateMsgOpt.isPresent()); + } + @Test public void testSetRootRuleChain() throws Exception { // create rule chain From 6098f19dc23163185ffb7a10e72d942cb87e6e53 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 9 Sep 2024 13:51:49 +0300 Subject: [PATCH 4/4] Refactoring to remove duplicates --- .../thingsboard/server/edge/AbstractEdgeTest.java | 15 +++++++++++---- .../server/edge/RuleChainEdgeTest.java | 5 +---- 2 files changed, 12 insertions(+), 8 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java index 34f7da02e6..5484c29305 100644 --- a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java @@ -103,6 +103,7 @@ import org.thingsboard.server.gen.edge.v1.UserUpdateMsg; import java.util.ArrayList; import java.util.List; import java.util.Optional; +import java.util.Random; import java.util.TreeMap; import java.util.UUID; import java.util.concurrent.TimeUnit; @@ -124,6 +125,8 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { protected EdgeImitator edgeImitator; protected Edge edge; + private Random random = new Random(); + @Autowired protected EdgeEventService edgeEventService; @@ -191,10 +194,7 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { Asset savedAsset = saveAsset("Edge Asset 1"); - RuleChainId rootRuleChainId = getEdgeRootRuleChainId(); - RuleChainMetaData rootRuleChainMetadata = doGet("/api/ruleChain/" + rootRuleChainId.getId().toString() + "/metadata", RuleChainMetaData.class); - rootRuleChainMetadata.getNodes().forEach(n -> n.setDebugMode(true)); - doPost("/api/ruleChain/metadata", rootRuleChainMetadata, RuleChainMetaData.class); + updateRootRuleChainMetadata(); edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class); @@ -207,6 +207,13 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { TimeUnit.MILLISECONDS.sleep(1000); } + protected void updateRootRuleChainMetadata() throws Exception { + RuleChainId rootRuleChainId = getEdgeRootRuleChainId(); + RuleChainMetaData rootRuleChainMetadata = doGet("/api/ruleChain/" + rootRuleChainId.getId().toString() + "/metadata", RuleChainMetaData.class); + rootRuleChainMetadata.getNodes().forEach(n -> n.setDebugMode(random.nextBoolean())); + doPost("/api/ruleChain/metadata", rootRuleChainMetadata, RuleChainMetaData.class); + } + protected void extendDeviceProfileData(DeviceProfile deviceProfile) { DeviceProfileData profileData = deviceProfile.getProfileData(); List alarms = new ArrayList<>(); diff --git a/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java index 6dafeefa80..eaf4468156 100644 --- a/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java @@ -187,11 +187,8 @@ public class RuleChainEdgeTest extends AbstractEdgeTest { @Test public void testUpdateRootRuleChain() throws Exception { - RuleChainMetaData rootRuleChainMetadata = doGet("/api/ruleChain/" + edge.getRootRuleChainId().getId().toString() + "/metadata", RuleChainMetaData.class); - - rootRuleChainMetadata.getNodes().forEach(n -> n.setDebugMode(true)); edgeImitator.expectMessageAmount(2); - doPost("/api/ruleChain/metadata", rootRuleChainMetadata, RuleChainMetaData.class); + updateRootRuleChainMetadata(); Assert.assertTrue(edgeImitator.waitForMessages()); Optional ruleChainUpdateMsgOpt = edgeImitator.findMessageByType(RuleChainUpdateMsg.class);