From dd9648d99905b10e514797c23a081e95ab4ce093 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Fri, 10 Jun 2022 18:47:52 +0300 Subject: [PATCH 1/3] Fixing incorrect check for missing related rule chains #1 --- .../server/controller/EdgeController.java | 3 ++- .../constructor/RuleChainMsgConstructor.java | 25 +++++++++---------- .../server/dao/edge/EdgeService.java | 2 +- .../server/dao/edge/EdgeServiceImpl.java | 18 +++++++------ .../dao/service/BaseEdgeServiceTest.java | 25 +++++++++++++++++++ 5 files changed, 50 insertions(+), 23 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/EdgeController.java b/application/src/main/java/org/thingsboard/server/controller/EdgeController.java index 5d9737db23..1dfbc83996 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EdgeController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EdgeController.java @@ -32,6 +32,7 @@ import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.ResponseBody; import org.springframework.web.bind.annotation.ResponseStatus; import org.springframework.web.bind.annotation.RestController; +import org.thingsboard.rule.engine.flow.TbRuleChainInputNode; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.edge.Edge; @@ -557,7 +558,7 @@ public class EdgeController extends BaseController { edgeId = checkNotNull(edgeId); SecurityUser user = getCurrentUser(); TenantId tenantId = user.getTenantId(); - return edgeService.findMissingToRelatedRuleChains(tenantId, edgeId); + return edgeService.findMissingToRelatedRuleChains(tenantId, edgeId, TbRuleChainInputNode.class.getName()); } catch (Exception e) { throw handleException(e); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainMsgConstructor.java index 0fc62cfe41..db67ecf822 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainMsgConstructor.java @@ -21,7 +21,9 @@ import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.rule.engine.flow.TbRuleChainInputNode; import org.thingsboard.rule.engine.flow.TbRuleChainInputNodeConfiguration; +import org.thingsboard.rule.engine.flow.TbRuleChainOutputNode; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.rule.NodeConnectionInfo; import org.thingsboard.server.common.data.rule.RuleChain; @@ -49,9 +51,8 @@ import java.util.stream.Collectors; @TbCoreComponent public class RuleChainMsgConstructor { - private static final ObjectMapper objectMapper = new ObjectMapper(); - private static final String RULE_CHAIN_INPUT_NODE = "org.thingsboard.rule.engine.flow.TbRuleChainInputNode"; - private static final String TB_RULE_CHAIN_OUTPUT_NODE = "org.thingsboard.rule.engine.flow.TbRuleChainOutputNode"; + private static final String RULE_CHAIN_INPUT_NODE = TbRuleChainInputNode.class.getName(); + private static final String TB_RULE_CHAIN_OUTPUT_NODE = TbRuleChainOutputNode.class.getName(); public RuleChainUpdateMsg constructRuleChainUpdatedMsg(RuleChainId edgeRootRuleChainId, UpdateMsgType msgType, RuleChain ruleChain) { RuleChainUpdateMsg.Builder builder = RuleChainUpdateMsg.newBuilder() @@ -210,13 +211,11 @@ public class RuleChainMsgConstructor { private List filterNodes_V_3_3_0(List nodes) { List result = new ArrayList<>(); for (RuleNode node : nodes) { - switch (node.getType()) { - case RULE_CHAIN_INPUT_NODE: - case TB_RULE_CHAIN_OUTPUT_NODE: - log.trace("Skipping not supported rule node {}", node); - break; - default: - result.add(node); + if (RULE_CHAIN_INPUT_NODE.equals(node.getType()) + || TB_RULE_CHAIN_OUTPUT_NODE.equals(node.getType())) { + log.trace("Skipping not supported rule node {}", node); + } else { + result.add(node); } } return result; @@ -280,7 +279,7 @@ public class RuleChainMsgConstructor { .setTargetRuleChainIdMSB(ruleChainConnectionInfo.getTargetRuleChainId().getId().getMostSignificantBits()) .setTargetRuleChainIdLSB(ruleChainConnectionInfo.getTargetRuleChainId().getId().getLeastSignificantBits()) .setType(ruleChainConnectionInfo.getType()) - .setAdditionalInfo(objectMapper.writeValueAsString(additionalInfo)) + .setAdditionalInfo(JacksonUtil.OBJECT_MAPPER.writeValueAsString(additionalInfo)) .build(); } @@ -291,8 +290,8 @@ public class RuleChainMsgConstructor { .setType(node.getType()) .setName(node.getName()) .setDebugMode(node.isDebugMode()) - .setConfiguration(objectMapper.writeValueAsString(node.getConfiguration())) - .setAdditionalInfo(objectMapper.writeValueAsString(node.getAdditionalInfo())) + .setConfiguration(JacksonUtil.OBJECT_MAPPER.writeValueAsString(node.getConfiguration())) + .setAdditionalInfo(JacksonUtil.OBJECT_MAPPER.writeValueAsString(node.getAdditionalInfo())) .build(); } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java index ea95923275..55f573a710 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java @@ -84,5 +84,5 @@ public interface EdgeService { PageData findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink); - String findMissingToRelatedRuleChains(TenantId tenantId, EdgeId edgeId); + String findMissingToRelatedRuleChains(TenantId tenantId, EdgeId edgeId, String tbRuleChainInputNodeName); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java index 271c74a3e1..ab091632f0 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.dao.edge; -import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.base.Function; @@ -26,9 +25,8 @@ import lombok.extern.slf4j.Slf4j; import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; -import org.springframework.transaction.annotation.Propagation; -import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.event.TransactionalEventListener; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.StringUtils; @@ -80,8 +78,6 @@ public class EdgeServiceImpl extends AbstractCachedEntityService edgeRuleChains = findEdgeRuleChains(tenantId, edgeId); List edgeRuleChainIds = edgeRuleChains.stream().map(IdBased::getId).collect(Collectors.toList()); - ObjectNode result = mapper.createObjectNode(); + ObjectNode result = JacksonUtil.OBJECT_MAPPER.createObjectNode(); for (RuleChain edgeRuleChain : edgeRuleChains) { + // ruleChainService. + // loadRuleChainMetaData(edgeRuleChain.getTenantId(), edgeRuleChain.getId()) + // .getNodes() + // .get(11) + // .getConfiguration() + // .get("ruleChainId") List connectionInfos = ruleChainService.loadRuleChainMetaData(edgeRuleChain.getTenantId(), edgeRuleChain.getId()).getRuleChainConnections(); if (connectionInfos != null && !connectionInfos.isEmpty()) { @@ -471,7 +473,7 @@ public class EdgeServiceImpl extends AbstractCachedEntityService Date: Mon, 13 Jun 2022 13:10:50 +0300 Subject: [PATCH 2/3] Fixing incorrect check for missing related rule chains #2 --- .../server/dao/edge/EdgeServiceImpl.java | 22 +++++----- .../dao/service/BaseEdgeServiceTest.java | 41 +++++++++++++++++-- 2 files changed, 47 insertions(+), 16 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java index ab091632f0..9eee9e4774 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java @@ -46,7 +46,7 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntitySearchDirection; import org.thingsboard.server.common.data.rule.RuleChain; -import org.thingsboard.server.common.data.rule.RuleChainConnectionInfo; +import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.dao.entity.AbstractCachedEntityService; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.relation.RelationService; @@ -62,6 +62,7 @@ import java.util.Collections; import java.util.Comparator; import java.util.List; import java.util.Optional; +import java.util.UUID; import java.util.stream.Collectors; import static org.thingsboard.server.dao.DaoUtil.toUUIDs; @@ -449,22 +450,19 @@ public class EdgeServiceImpl extends AbstractCachedEntityService edgeRuleChains = findEdgeRuleChains(tenantId, edgeId); List edgeRuleChainIds = edgeRuleChains.stream().map(IdBased::getId).collect(Collectors.toList()); ObjectNode result = JacksonUtil.OBJECT_MAPPER.createObjectNode(); for (RuleChain edgeRuleChain : edgeRuleChains) { - // ruleChainService. - // loadRuleChainMetaData(edgeRuleChain.getTenantId(), edgeRuleChain.getId()) - // .getNodes() - // .get(11) - // .getConfiguration() - // .get("ruleChainId") - List connectionInfos = - ruleChainService.loadRuleChainMetaData(edgeRuleChain.getTenantId(), edgeRuleChain.getId()).getRuleChainConnections(); - if (connectionInfos != null && !connectionInfos.isEmpty()) { + List ruleNodes = + ruleChainService.loadRuleChainMetaData(edgeRuleChain.getTenantId(), edgeRuleChain.getId()).getNodes(); + if (ruleNodes != null && !ruleNodes.isEmpty()) { List connectedRuleChains = - connectionInfos.stream().map(RuleChainConnectionInfo::getTargetRuleChainId).collect(Collectors.toList()); + ruleNodes.stream() + .filter(rn -> rn.getType().equals(tbRuleChainInputNodeClassName)) + .map(rn -> new RuleChainId(UUID.fromString(rn.getConfiguration().get("ruleChainId").asText()))) + .collect(Collectors.toList()); List missingRuleChains = new ArrayList<>(); for (RuleChainId connectedRuleChain : connectedRuleChains) { if (!edgeRuleChainIds.contains(connectedRuleChain)) { diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeServiceTest.java index 7b88adf50f..9c7c53d587 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeServiceTest.java @@ -16,11 +16,13 @@ package org.thingsboard.server.dao.service; import com.datastax.oss.driver.api.core.uuid.Uuids; +import com.fasterxml.jackson.databind.node.ObjectNode; import org.apache.commons.lang3.RandomStringUtils; import org.junit.After; import org.junit.Assert; import org.junit.Before; import org.junit.Test; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.Tenant; @@ -30,10 +32,13 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.rule.RuleChain; +import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.rule.RuleChainType; +import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.dao.exception.DataValidationException; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collections; import java.util.List; @@ -614,17 +619,45 @@ public abstract class BaseEdgeServiceTest extends AbstractServiceTest { ruleChain.setName("Rule Chain #1"); ruleChain.setType(RuleChainType.EDGE); RuleChain ruleChain1 = ruleChainService.saveRuleChain(ruleChain); - ruleChainService.assignRuleChainToEdge(tenantId, ruleChain1.getId(), savedEdge.getId()); ruleChain = new RuleChain(); ruleChain.setTenantId(tenantId); ruleChain.setName("Rule Chain #2"); ruleChain.setType(RuleChainType.EDGE); RuleChain ruleChain2 = ruleChainService.saveRuleChain(ruleChain); - ruleChainService.assignRuleChainToEdge(tenantId, ruleChain2.getId(), savedEdge.getId()); - String missingToRelatedRuleChains = edgeService.findMissingToRelatedRuleChains(tenantId, savedEdge.getId()); - Assert.assertEquals("[]", missingToRelatedRuleChains); + ruleChain = new RuleChain(); + ruleChain.setTenantId(tenantId); + ruleChain.setName("Rule Chain #3"); + ruleChain.setType(RuleChainType.EDGE); + RuleChain ruleChain3 = ruleChainService.saveRuleChain(ruleChain); + + RuleNode ruleNode1 = new RuleNode(); + ruleNode1.setName("Input rule node 1"); + ruleNode1.setType("org.thingsboard.rule.engine.flow.TbRuleChainInputNode"); + ObjectNode configuration = JacksonUtil.OBJECT_MAPPER.createObjectNode(); + configuration.put("ruleChainId", ruleChain1.getUuidId().toString()); + ruleNode1.setConfiguration(configuration); + + RuleNode ruleNode2 = new RuleNode(); + ruleNode2.setName("Input rule node 2"); + ruleNode2.setType("org.thingsboard.rule.engine.flow.TbRuleChainInputNode"); + configuration = JacksonUtil.OBJECT_MAPPER.createObjectNode(); + configuration.put("ruleChainId", ruleChain2.getUuidId().toString()); + ruleNode2.setConfiguration(configuration); + + RuleChainMetaData ruleChainMetaData3 = new RuleChainMetaData(); + ruleChainMetaData3.setNodes(Arrays.asList(ruleNode1, ruleNode2)); + ruleChainMetaData3.setFirstNodeIndex(0); + ruleChainMetaData3.setRuleChainId(ruleChain3.getId()); + ruleChainService.saveRuleChainMetaData(tenantId, ruleChainMetaData3); + + ruleChainService.assignRuleChainToEdge(tenantId, ruleChain3.getId(), savedEdge.getId()); + + String missingToRelatedRuleChains = edgeService.findMissingToRelatedRuleChains(tenantId, + savedEdge.getId(), + "org.thingsboard.rule.engine.flow.TbRuleChainInputNode"); + Assert.assertEquals("{\"Rule Chain #3\":[\"Rule Chain #1\",\"Rule Chain #2\"]}", missingToRelatedRuleChains); } } \ No newline at end of file From 2a3126822e1e86358d2b2e8874c6ccf0aad562f7 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 13 Jun 2022 13:52:40 +0300 Subject: [PATCH 3/3] Minor code review changes --- .../main/java/org/thingsboard/server/dao/edge/EdgeService.java | 2 +- .../org/thingsboard/server/dao/service/BaseEdgeServiceTest.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java index 55f573a710..c0c72552f6 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java @@ -84,5 +84,5 @@ public interface EdgeService { PageData findRelatedEdgeIdsByEntityId(TenantId tenantId, EntityId entityId, PageLink pageLink); - String findMissingToRelatedRuleChains(TenantId tenantId, EdgeId edgeId, String tbRuleChainInputNodeName); + String findMissingToRelatedRuleChains(TenantId tenantId, EdgeId edgeId, String tbRuleChainInputNodeClassName); } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeServiceTest.java index 9c7c53d587..093bd386a7 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseEdgeServiceTest.java @@ -660,4 +660,4 @@ public abstract class BaseEdgeServiceTest extends AbstractServiceTest { Assert.assertEquals("{\"Rule Chain #3\":[\"Rule Chain #1\",\"Rule Chain #2\"]}", missingToRelatedRuleChains); } -} \ No newline at end of file +}