diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java index 11e3172168..61875f1277 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java @@ -25,6 +25,10 @@ import com.google.gson.JsonPrimitive; import com.google.gson.reflect.TypeToken; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.rule.engine.action.TbSaveToCustomCassandraTableNode; +import org.thingsboard.rule.engine.aws.lambda.TbAwsLambdaNode; +import org.thingsboard.rule.engine.rest.TbSendRestApiCallReplyNode; +import org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode; import org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode; import org.thingsboard.server.common.adaptor.JsonConverter; import org.thingsboard.server.common.data.Customer; @@ -117,11 +121,25 @@ import org.thingsboard.server.gen.edge.v1.WidgetsBundleUpdateMsg; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.service.edge.rpc.utils.EdgeVersionUtils; +import java.util.Iterator; import java.util.List; +import java.util.Map; +import java.util.Set; import java.util.UUID; @Slf4j public class EdgeMsgConstructorUtils { + public static final Map NODE_TO_IGNORE_PARAM_FOR_OLD_EDGE_VERSION = Map.of( + TbMsgTimeseriesNode.class.getName(), "processingSettings", + TbMsgAttributesNode.class.getName(), "processingSettings", + TbSaveToCustomCassandraTableNode.class.getName(), "defaultTtl" + ); + + //added in edge version 3.8.0 + public static final Set MISSING_NODES_IN_VERSION_37 = Set.of( + TbSendRestApiCallReplyNode.class.getName(), + TbAwsLambdaNode.class.getName() + ); public static AlarmUpdateMsg constructAlarmUpdatedMsg(UpdateMsgType msgType, Alarm alarm) { return AlarmUpdateMsg.newBuilder().setMsgType(msgType) @@ -431,25 +449,36 @@ public class EdgeMsgConstructorUtils { } private static String filterMetadataForOldEdgeVersions(RuleChainMetaData ruleChainMetaData, EdgeVersion edgeVersion) { - if (EdgeVersionUtils.isEdgeVersionOlderThan(edgeVersion, EdgeVersion.V_3_9_0)) { - JsonNode jsonNode = JacksonUtil.valueToTree(ruleChainMetaData); - JsonNode nodes = jsonNode.get("nodes"); + JsonNode jsonNode = JacksonUtil.valueToTree(ruleChainMetaData); + JsonNode nodes = jsonNode.get("nodes"); + + if (EdgeVersionUtils.isEdgeVersionOlderThan(edgeVersion, EdgeVersion.V_3_8_0)) { + Iterator iterator = nodes.iterator(); + while (iterator.hasNext()) { + JsonNode node = iterator.next(); - for (JsonNode node : nodes) { - if (node.isObject()) { - removeIncompatibleFields((ObjectNode) node); + String type = node.get("type").asText(); + if (MISSING_NODES_IN_VERSION_37.contains(type)) { + iterator.remove(); } } + } + + if (EdgeVersionUtils.isEdgeVersionOlderThan(edgeVersion, EdgeVersion.V_3_9_0)) { + nodes.forEach(EdgeMsgConstructorUtils::changeRuleNodeConfigForOldEdgeVersion); + return JacksonUtil.toString(jsonNode); } else { return JacksonUtil.toString(ruleChainMetaData); } } - private static void removeIncompatibleFields(ObjectNode node) { - if (TbMsgTimeseriesNode.class.getName().equals(node.get("type").asText())) { - if (node.has("configuration") && node.get("configuration").isObject()) { - ((ObjectNode) node.get("configuration")).remove("processingSettings"); + private static void changeRuleNodeConfigForOldEdgeVersion(JsonNode node) { + if (node.isObject()) { + JsonNode configurationNode = node.get("configuration"); + if (configurationNode != null && configurationNode.isObject() && + NODE_TO_IGNORE_PARAM_FOR_OLD_EDGE_VERSION.containsKey(node.get("type").asText())) { + ((ObjectNode) configurationNode).remove(NODE_TO_IGNORE_PARAM_FOR_OLD_EDGE_VERSION.get(node.get("type").asText())); } } } diff --git a/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java b/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java index 5e7cfcfc28..84ccfdb3ee 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java +++ b/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java @@ -15,9 +15,20 @@ */ package org.thingsboard.server.service.edge; +import lombok.extern.slf4j.Slf4j; import org.junit.Assert; import org.junit.Test; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.rule.engine.action.TbSaveToCustomCassandraTableNode; +import org.thingsboard.rule.engine.action.TbSaveToCustomCassandraTableNodeConfiguration; +import org.thingsboard.rule.engine.api.NodeConfiguration; +import org.thingsboard.rule.engine.aws.lambda.TbAwsLambdaNode; +import org.thingsboard.rule.engine.aws.lambda.TbAwsLambdaNodeConfiguration; +import org.thingsboard.rule.engine.rest.TbSendRestApiCallReplyNode; +import org.thingsboard.rule.engine.rest.TbSendRestApiCallReplyNodeConfiguration; +import org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode; +import org.thingsboard.rule.engine.telemetry.TbMsgAttributesNodeConfiguration; +import org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode; import org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.rule.RuleNode; @@ -26,66 +37,129 @@ import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.service.edge.rpc.utils.EdgeVersionUtils; -import java.util.Collections; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.Map; +import static org.thingsboard.server.service.edge.EdgeMsgConstructorUtils.MISSING_NODES_IN_VERSION_37; +import static org.thingsboard.server.service.edge.EdgeMsgConstructorUtils.NODE_TO_IGNORE_PARAM_FOR_OLD_EDGE_VERSION; + +@Slf4j public class EdgeMsgConstructorUtilsTest { private static final int CONFIGURATION_VERSION = 5; + public static final List TEST_SUPPORTED_EDGE_VERSIONS = Arrays.asList( + EdgeVersion.V_4_0_0, EdgeVersion.V_3_9_0, EdgeVersion.V_3_8_0, EdgeVersion.V_3_7_0 + ); + + private static final Map CONFIG_TO_NODE_NAME = Map.of( + new TbMsgTimeseriesNodeConfiguration(), TbMsgTimeseriesNode.class.getName(), + new TbMsgAttributesNodeConfiguration(), TbMsgAttributesNode.class.getName(), + new TbSaveToCustomCassandraTableNodeConfiguration(), TbSaveToCustomCassandraTableNode.class.getName() + ); + + private static final Map NODE_TO_CONFIG_PARAMS_COUNT = Map.of( + TbMsgTimeseriesNode.class.getName(), 3, + TbMsgAttributesNode.class.getName(), 5, + TbSaveToCustomCassandraTableNode.class.getName(), 3 + ); + + private static final Map CONFIG_TO_MISS_NODE_FOR_OLD_EDGE = Map.of( + new TbSendRestApiCallReplyNodeConfiguration(), TbSendRestApiCallReplyNode.class.getName(), + new TbAwsLambdaNodeConfiguration(), TbAwsLambdaNode.class.getName() + ); + @Test - public void testRuleChainMetadataUpdateMsgForAllEdgeVersions() { + public void testRuleChainMetadataUpdateMsgForOldEdgeVersions() { // GIVEN - RuleChainMetaData metaData = createIncompatibleRuleNodesForOldEdge(); - - // WHEN - RuleNode ruleNode_V_4_0_0 = getRuleNodeFromMetadataUpdateMessage(metaData, EdgeVersion.V_4_0_0); - RuleNode ruleNode_V_3_9_0 = getRuleNodeFromMetadataUpdateMessage(metaData, EdgeVersion.V_3_9_0); - RuleNode ruleNode_V_3_8_0 = getRuleNodeFromMetadataUpdateMessage(metaData, EdgeVersion.V_3_8_0); - RuleNode ruleNode_V_3_7_0 = getRuleNodeFromMetadataUpdateMessage(metaData, EdgeVersion.V_3_7_0); - - // THEN - assertRuleNodeConfiguration(ruleNode_V_4_0_0, EdgeVersion.V_4_0_0); - assertRuleNodeConfiguration(ruleNode_V_3_9_0, EdgeVersion.V_3_9_0); - assertRuleNodeConfiguration(ruleNode_V_3_8_0, EdgeVersion.V_3_8_0); - assertRuleNodeConfiguration(ruleNode_V_3_7_0, EdgeVersion.V_3_7_0); + RuleChainMetaData metaData = createMetadataWithProblemNodes(CONFIG_TO_NODE_NAME); + + TEST_SUPPORTED_EDGE_VERSIONS.forEach(edgeVersion -> { + // WHEN + List ruleNodes = getRuleNodesFromUpdateMsg(metaData, edgeVersion); + + // THEN + validateRuleNodeConfig(ruleNodes, edgeVersion); + }); } - private RuleChainMetaData createIncompatibleRuleNodesForOldEdge() { + @Test + public void testRuleChainMetadataWithMissingNodeForOldEdgeVersions() { + // GIVEN + RuleChainMetaData metaData = createMetadataWithProblemNodes(CONFIG_TO_MISS_NODE_FOR_OLD_EDGE); + + TEST_SUPPORTED_EDGE_VERSIONS.forEach(edgeVersion -> { + // WHEN + List ruleNodes = getRuleNodesFromUpdateMsg(metaData, edgeVersion); + + // THEN + boolean isOldEdge = EdgeVersionUtils.isEdgeVersionOlderThan(edgeVersion, EdgeVersion.V_3_8_0); + + if (isOldEdge) { + Assert.assertTrue("Rule Node must be empty", ruleNodes.isEmpty()); + } else { + Assert.assertEquals(MISSING_NODES_IN_VERSION_37.size(), ruleNodes.size()); + } + }); + } + + private RuleChainMetaData createMetadataWithProblemNodes(Map nodeMap) { RuleChainMetaData ruleChainMetaData = new RuleChainMetaData(); + List ruleNodes = new ArrayList<>(); + + nodeMap.entrySet().forEach(configToNodeName -> { + RuleNode ruleNode = new RuleNode(); - RuleNode ruleNode1 = new RuleNode(); - ruleNode1.setName("TbMsgTimeseriesNode"); - ruleNode1.setType(org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode.class.getName()); - ruleNode1.setConfigurationVersion(CONFIGURATION_VERSION); - ruleNode1.setConfiguration(JacksonUtil.valueToTree(new TbMsgTimeseriesNodeConfiguration().defaultConfiguration())); + ruleNode.setName(configToNodeName.getValue()); + ruleNode.setType(configToNodeName.getValue()); + ruleNode.setConfigurationVersion(CONFIGURATION_VERSION); + ruleNode.setConfiguration(JacksonUtil.valueToTree(configToNodeName.getKey().defaultConfiguration())); + + ruleNodes.add(ruleNode); + }); ruleChainMetaData.setFirstNodeIndex(0); - ruleChainMetaData.setNodes(Collections.singletonList(ruleNode1)); + ruleChainMetaData.setNodes(ruleNodes); return ruleChainMetaData; } - - private RuleNode getRuleNodeFromMetadataUpdateMessage(RuleChainMetaData metaData, EdgeVersion edgeVersion) { + private List getRuleNodesFromUpdateMsg(RuleChainMetaData metaData, EdgeVersion edgeVersion) { RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = EdgeMsgConstructorUtils.constructRuleChainMetadataUpdatedMsg(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, metaData, edgeVersion); RuleChainMetaData ruleChainMetaData = JacksonUtil.fromString(ruleChainMetadataUpdateMsg.getEntity(), RuleChainMetaData.class, true); Assert.assertNotNull("RuleChainMetaData is null", ruleChainMetaData); - RuleNode ruleNode = ruleChainMetaData.getNodes().stream().findFirst().orElse(null); - Assert.assertNotNull("RuleNode is null for Edge version " + edgeVersion, ruleNode); - Assert.assertNotNull("Configuration is null for Edge version " + edgeVersion, ruleNode.getConfiguration()); + return ruleChainMetaData.getNodes(); + } - return ruleNode; + private void validateRuleNodeConfig(List ruleNodes, EdgeVersion edgeVersion) { + ruleNodes.forEach(ruleNode -> { + int ruleNodeConfigAmount = NODE_TO_CONFIG_PARAMS_COUNT.get(ruleNode.getName()); + + boolean isOldEdge = EdgeVersionUtils.isEdgeVersionOlderThan(edgeVersion, EdgeVersion.V_3_9_0); + int expectedConfigAmount = isOldEdge ? ruleNodeConfigAmount - 1 : ruleNodeConfigAmount; + boolean includeConfigParam = !isOldEdge; + + validateParams(ruleNode, expectedConfigAmount, includeConfigParam); + }); } - private void assertRuleNodeConfiguration(RuleNode ruleNode, EdgeVersion edgeVersion) { - if (EdgeVersionUtils.isEdgeVersionOlderThan(edgeVersion, EdgeVersion.V_3_9_0)) { - Assert.assertEquals("Unexpected config size", 2, ruleNode.getConfiguration().size()); - Assert.assertFalse("Unexpected field 'processingSettings'", ruleNode.getConfiguration().has("processingSettings")); - }else{ - Assert.assertEquals("Unexpected config size", 3, ruleNode.getConfiguration().size()); - Assert.assertTrue("Missing field 'processingSettings'", ruleNode.getConfiguration().has("processingSettings")); - } + private void validateParams(RuleNode ruleNode, int expectedConfigAmount, boolean includeConfigParam) { + String ignoreConfigParam = NODE_TO_IGNORE_PARAM_FOR_OLD_EDGE_VERSION.get(ruleNode.getName()); + + Assert.assertEquals( + String.format("Expected %d config params for ruleNode '%s', but found %d", expectedConfigAmount, ruleNode.getName(), ruleNode.getConfiguration().size()), + expectedConfigAmount, ruleNode.getConfiguration().size() + ); + + boolean hasIgnoredField = ruleNode.getConfiguration().has(ignoreConfigParam); + Assert.assertEquals( + String.format("Field '%s' for ruleNode '%s' should %s be present", ignoreConfigParam, ruleNode.getName(), includeConfigParam ? "not" : ""), + includeConfigParam, hasIgnoredField + ); } + }