From 5831876b8ec89e35d233e6e121c421866c7ce4aa Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Tue, 16 May 2023 15:57:30 +0300 Subject: [PATCH] added upgrade script logic for all enrichment rule nodes && additional improvements to TbGetRelatedAttributeNode --- .../install/ThingsboardInstallService.java | 1 + .../update/DefaultDataUpdateService.java | 139 +++++--------- .../dao/timeseries/TimeseriesService.java | 2 + dao/pom.xml | 4 + .../server/dao/rule/BaseRuleChainService.java | 28 +++ .../dao/sqlts/SqlTimeseriesLatestDao.java | 19 +- .../dao/timeseries/BaseTimeseriesService.java | 9 + .../CassandraBaseTimeseriesLatestDao.java | 11 ++ .../dao/timeseries/TimeseriesLatestDao.java | 2 + .../rule/engine/api/VersionedNode.java | 43 +++++ .../api/VersionedNodeConfiguration.java} | 6 +- .../engine/metadata/CalculateDeltaNode.java | 27 ++- .../rule/engine/metadata/DataToFetch.java | 22 +++ .../TbAbstractFetchToNodeConfiguration.java | 4 +- .../metadata/TbAbstractGetEntityAttrNode.java | 101 ++++++++-- .../TbAbstractGetEntityDetailsNode.java | 53 ++---- ...ractGetEntityDetailsNodeConfiguration.java | 4 +- .../metadata/TbAbstractNodeWithFetchTo.java | 38 +++- .../TbFetchDeviceCredentialsNode.java | 22 +++ .../engine/metadata/TbGetAttributesNode.java | 22 +++ .../metadata/TbGetCustomerAttributeNode.java | 27 +++ .../metadata/TbGetCustomerDetailsNode.java | 22 +++ .../engine/metadata/TbGetDeviceAttrNode.java | 22 +++ .../TbGetEntityAttrNodeConfiguration.java | 4 +- .../metadata/TbGetOriginatorFieldsNode.java | 22 +++ .../TbGetRelatedAttrNodeConfiguration.java | 2 +- .../metadata/TbGetRelatedAttributeNode.java | 28 +++ .../metadata/TbGetTenantAttributeNode.java | 25 +++ .../metadata/TbGetTenantDetailsNode.java | 22 +++ .../util/ContactBasedEntityDetails.java | 41 +++++ .../metadata/CalculateDeltaNodeTest.java | 123 ++++++------- .../TbGetCustomerAttributeNodeTest.java | 46 ++++- .../TbGetCustomerDetailsNodeTest.java | 26 +-- .../TbGetRelatedAttributeNodeTest.java | 172 ++++++++++++++---- .../TbGetTenantAttributeNodeTest.java | 46 ++++- .../metadata/TbGetTenantDetailsNodeTest.java | 20 +- 36 files changed, 891 insertions(+), 314 deletions(-) create mode 100644 rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/VersionedNode.java rename rule-engine/{rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntityDetails.java => rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/VersionedNodeConfiguration.java} (79%) create mode 100644 rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/DataToFetch.java create mode 100644 rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/ContactBasedEntityDetails.java diff --git a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java index 1f163aa876..66867cd90d 100644 --- a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java +++ b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java @@ -262,6 +262,7 @@ public class ThingsboardInstallService { log.info("Updating system data..."); systemDataLoaderService.updateSystemWidgets(); //TODO update CacheCleanupService on the next version upgrade + break; default: throw new RuntimeException("Unable to upgrade ThingsBoard, unsupported fromVersion: " + upgradeFromVersion); diff --git a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java index 315c2e71e1..e5bb70cbbd 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java @@ -26,9 +26,18 @@ import org.springframework.context.annotation.Lazy; import org.springframework.context.annotation.Profile; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.rule.engine.api.VersionedNode; import org.thingsboard.rule.engine.flow.TbRuleChainInputNode; import org.thingsboard.rule.engine.flow.TbRuleChainInputNodeConfiguration; -import org.thingsboard.rule.engine.metadata.FetchTo; +import org.thingsboard.rule.engine.metadata.TbFetchDeviceCredentialsNode; +import org.thingsboard.rule.engine.metadata.TbGetAttributesNode; +import org.thingsboard.rule.engine.metadata.TbGetCustomerAttributeNode; +import org.thingsboard.rule.engine.metadata.TbGetCustomerDetailsNode; +import org.thingsboard.rule.engine.metadata.TbGetDeviceAttrNode; +import org.thingsboard.rule.engine.metadata.TbGetOriginatorFieldsNode; +import org.thingsboard.rule.engine.metadata.TbGetRelatedAttributeNode; +import org.thingsboard.rule.engine.metadata.TbGetTenantAttributeNode; +import org.thingsboard.rule.engine.metadata.TbGetTenantDetailsNode; import org.thingsboard.rule.engine.profile.TbDeviceProfileNode; import org.thingsboard.rule.engine.profile.TbDeviceProfileNodeConfiguration; import org.thingsboard.server.common.data.DataConstants; @@ -209,7 +218,7 @@ public class DefaultDataUpdateService implements DataUpdateService { case "3.5.0": log.info("Updating data from version 3.5.0 to 3.5.1 ..."); log.info("Starting enrichment rule nodes update ..."); - updateEnrichmentRuleNodes(); + upgradeEnrichmentRuleNodesWithFetchTo(); log.info("Finished enrichment rule nodes update!"); break; default: @@ -217,108 +226,46 @@ public class DefaultDataUpdateService implements DataUpdateService { } } - private void updateEnrichmentRuleNodes() { + private void upgradeEnrichmentRuleNodesWithFetchTo() { try { var ruleChainIdToTenantId = new HashMap(); - var allNodesToUpdate = List.of( - "org.thingsboard.rule.engine.metadata.TbGetOriginatorFieldsNode", - "org.thingsboard.rule.engine.metadata.TbGetRelatedAttributeNode", - "org.thingsboard.rule.engine.metadata.TbGetTenantAttributeNode", - "org.thingsboard.rule.engine.metadata.TbGetCustomerAttributeNode", - "org.thingsboard.rule.engine.metadata.TbGetAttributesNode", - "org.thingsboard.rule.engine.metadata.TbGetDeviceAttrNode", - "org.thingsboard.rule.engine.metadata.TbGetCustomerDetailsNode", - "org.thingsboard.rule.engine.metadata.TbGetTenantDetailsNode", - "org.thingsboard.rule.engine.metadata.TbFetchDeviceCredentialsNode" - ); - allNodesToUpdate.forEach(ruleNodeType -> { - var ruleNodes = new PageDataIterable<>( - pageLink -> ruleChainService.findAllRuleNodesByType(ruleNodeType, pageLink), 1024 - ); - for (var ruleNode : ruleNodes) { - var configuration = ruleNode.getConfiguration(); - if (configuration == null) { - log.error("Failed to update rule node: [{}] with id: [{}] Node configuration is null! Skipping this node!", - ruleNodeType, ruleNode.getId()); - continue; - } - if (!configuration.isObject()) { - log.error("Failed to update rule node: [{}] with id: [{}] Node configuration is not an object! Skipping this node!", - ruleNodeType, ruleNode.getId()); - continue; - } - var configObjectNode = (ObjectNode) configuration; - - FetchTo fetchTo; - - switch (ruleNodeType) { - case "org.thingsboard.rule.engine.metadata.TbGetAttributesNode": - case "org.thingsboard.rule.engine.metadata.TbGetDeviceAttrNode": - fetchTo = checkEnrichmentNodeFetchProperty(configObjectNode, "fetchToData", FetchTo.DATA, FetchTo.METADATA); - break; - case "org.thingsboard.rule.engine.metadata.TbGetCustomerDetailsNode": - case "org.thingsboard.rule.engine.metadata.TbGetTenantDetailsNode": - fetchTo = checkEnrichmentNodeFetchProperty(configObjectNode, "addToMetadata", FetchTo.METADATA, FetchTo.DATA); - break; - case "org.thingsboard.rule.engine.metadata.TbFetchDeviceCredentialsNode": - fetchTo = checkEnrichmentNodeFetchProperty(configObjectNode, "fetchToMetadata", FetchTo.METADATA, FetchTo.DATA); - break; - case "org.thingsboard.rule.engine.metadata.TbGetOriginatorFieldsNode": - case "org.thingsboard.rule.engine.metadata.TbGetRelatedAttributeNode": - case "org.thingsboard.rule.engine.metadata.TbGetTenantAttributeNode": - case "org.thingsboard.rule.engine.metadata.TbGetCustomerAttributeNode": - fetchTo = FetchTo.METADATA; - break; - default: - log.error("Failed to update rule node: [{}] with id: [{}] " + - "Reason: Unexpected rule node type!", ruleNodeType, ruleNode.getId()); - continue; - } - - if (fetchTo == null) { - log.error("Failed to update rule node: [{}] with id: [{}]", ruleNodeType, ruleNode.getId()); - continue; - } - - configObjectNode.put("fetchTo", fetchTo.name()); - ruleNode.setConfiguration(configObjectNode); - var ruleChainId = ruleNode.getRuleChainId(); - var tenantId = ruleChainIdToTenantId.computeIfAbsent(ruleChainId, - id -> { - RuleChain ruleChain = ruleChainService.findRuleChainById(TenantId.SYS_TENANT_ID, id); - if (ruleChain == null) { - log.error("Failed to find rule chain by id: [{}], ruleNodeId: [{}]", ruleChainId, ruleNode.getId()); - return null; - } - return ruleChain.getTenantId(); - }); - if (tenantId != null) { - ruleChainService.saveRuleNode(tenantId, ruleNode); - } - } - }); + upgradeRuleNode(ruleChainIdToTenantId, new TbGetOriginatorFieldsNode()); + upgradeRuleNode(ruleChainIdToTenantId, new TbGetRelatedAttributeNode()); + upgradeRuleNode(ruleChainIdToTenantId, new TbGetTenantAttributeNode()); + upgradeRuleNode(ruleChainIdToTenantId, new TbGetCustomerAttributeNode()); + upgradeRuleNode(ruleChainIdToTenantId, new TbGetAttributesNode()); + upgradeRuleNode(ruleChainIdToTenantId, new TbGetDeviceAttrNode()); + upgradeRuleNode(ruleChainIdToTenantId, new TbGetCustomerDetailsNode()); + upgradeRuleNode(ruleChainIdToTenantId, new TbGetTenantDetailsNode()); + upgradeRuleNode(ruleChainIdToTenantId, new TbFetchDeviceCredentialsNode()); } catch (Exception e) { log.error("Unexpected error during enrichment rule nodes updating!", e); } } - private FetchTo checkEnrichmentNodeFetchProperty(ObjectNode config, String property, FetchTo ifTrue, FetchTo ifFalse) { - if (config.has(property)) { - var value = config.get(property).asText(); - if ("true".equals(value)) { - config.remove(property); - return ifTrue; - } else if ("false".equals(value)) { - config.remove(property); - return ifFalse; - } else { - log.error(property + " property has unexpected value: {} Allowed values: true or false!", value); - return null; + private void upgradeRuleNode(HashMap ruleChainIdToTenantId, VersionedNode versionedNode) { + var ruleNodes = new PageDataIterable<>( + pageLink -> ruleChainService.findAllRuleNodesByType(versionedNode.getClass().getName(), pageLink), 1024 + ); + ruleNodes.forEach(ruleNode -> { + var upgradeRuleNodeConfigurationResult = versionedNode.upgrade(ruleNode.getId(), ruleNode.getConfiguration()); + if (upgradeRuleNodeConfigurationResult.getFirst()) { + ruleNode.setConfiguration(upgradeRuleNodeConfigurationResult.getSecond()); + var ruleChainId = ruleNode.getRuleChainId(); + var tenantId = ruleChainIdToTenantId.computeIfAbsent(ruleChainId, + id -> { + RuleChain ruleChain = ruleChainService.findRuleChainById(TenantId.SYS_TENANT_ID, id); + if (ruleChain == null) { + log.error("Failed to find rule chain by id: [{}], ruleNodeId: [{}]", ruleChainId, ruleNode.getId()); + return null; + } + return ruleChain.getTenantId(); + }); + if (tenantId != null) { + ruleChainService.saveRuleNode(tenantId, ruleNode); + } } - } else { - log.error(property + " property is not present!"); - return null; - } + }); } private final PaginatedUpdater deviceProfileEntityDynamicConditionsUpdater = diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java index 186813f476..c2bc997235 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java @@ -42,6 +42,8 @@ public interface TimeseriesService { ListenableFuture> findLatest(TenantId tenantId, EntityId entityId, Collection keys); + List findLatestSync(TenantId tenantId, EntityId entityId, Collection keys); + ListenableFuture> findAllLatest(TenantId tenantId, EntityId entityId); ListenableFuture save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry); diff --git a/dao/pom.xml b/dao/pom.xml index 4fc8a3c08d..29a4a439c2 100644 --- a/dao/pom.xml +++ b/dao/pom.xml @@ -222,6 +222,10 @@ com.jayway.jsonpath json-path + + org.thingsboard.rule-engine + rule-engine-api + diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java index 76947c5043..a3fe6b39ea 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java @@ -27,6 +27,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.rule.engine.api.VersionedNode; import org.thingsboard.server.common.data.BaseData; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.edge.Edge; @@ -52,6 +53,7 @@ import org.thingsboard.server.common.data.rule.RuleChainUpdateResult; import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.common.data.rule.RuleNodeUpdateResult; import org.thingsboard.server.common.data.util.ReflectionUtils; +import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.entity.EntityCountService; import org.thingsboard.server.dao.exception.DataValidationException; @@ -60,6 +62,7 @@ import org.thingsboard.server.dao.service.PaginatedRemover; import org.thingsboard.server.dao.service.Validator; import org.thingsboard.server.dao.service.validator.RuleChainDataValidator; +import java.lang.reflect.InvocationTargetException; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; @@ -186,6 +189,11 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC RuleNode savedNode = ruleNodeDao.save(tenantId, node); relations.add(new EntityRelation(ruleChainMetaData.getRuleChainId(), savedNode.getId(), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.RULE_CHAIN)); + TbPair upgradeResult = upgradeRuleNode(savedNode); + if (upgradeResult.getFirst()) { + savedNode.setConfiguration(upgradeResult.getSecond()); + savedNode = ruleNodeDao.save(tenantId, savedNode); + } int index = nodes.indexOf(node); nodes.set(index, savedNode); ruleNodeIndexMap.put(savedNode.getId(), index); @@ -255,6 +263,26 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC return RuleChainUpdateResult.successful(updatedRuleNodes); } + private TbPair upgradeRuleNode(RuleNode node) { + var configuration = node.getConfiguration(); + String ruleNodeClassName = node.getType(); + try { + var ruleNodeClass = Class.forName(ruleNodeClassName).getConstructor().newInstance(); + if (ruleNodeClass instanceof VersionedNode) { + VersionedNode versionedNode = (VersionedNode) ruleNodeClass; + return versionedNode.upgrade(node.getId(), configuration); + } + } catch (InstantiationException | + IllegalAccessException | + InvocationTargetException | + NoSuchMethodException | + ClassNotFoundException e + ) { + log.warn("Failed to upgrade rule node due to: ", e); + } + return new TbPair<>(false, node.getConfiguration()); + } + @Override public RuleChainMetaData loadRuleChainMetaData(TenantId tenantId, RuleChainId ruleChainId) { Validator.validateId(ruleChainId, "Incorrect rule chain id."); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java index b58f0dc7c7..48b5b7b19a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java @@ -156,11 +156,12 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme @Override public ListenableFuture findLatest(TenantId tenantId, EntityId entityId, String key) { - TsKvEntry latest = doFindLatest(entityId, key); - if (latest == null) { - latest = new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry(key, null)); - } - return Futures.immediateFuture(latest); + return Futures.immediateFuture(getLatestTsKvEntry(entityId, key)); + } + + @Override + public TsKvEntry findLatestSync(TenantId tenantId, EntityId entityId, String key) { + return getLatestTsKvEntry(entityId, key); } @Override @@ -268,4 +269,12 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme return tsLatestQueue.add(latestEntity); } + private TsKvEntry getLatestTsKvEntry(EntityId entityId, String key) { + TsKvEntry latest = doFindLatest(entityId, key); + if (latest == null) { + latest = new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry(key, null)); + } + return latest; + } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java index 89aad061a7..34bcd15c0a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java @@ -132,6 +132,15 @@ public class BaseTimeseriesService implements TimeseriesService { return Futures.allAsList(futures); } + @Override + public List findLatestSync(TenantId tenantId, EntityId entityId, Collection keys) { + validate(entityId); + List latestEntries = Lists.newArrayListWithExpectedSize(keys.size()); + keys.forEach(key -> Validator.validateString(key, "Incorrect key " + key)); + keys.forEach(key -> latestEntries.add(timeseriesLatestDao.findLatestSync(tenantId, entityId, key))); + return latestEntries; + } + @Override public ListenableFuture> findAllLatest(TenantId tenantId, EntityId entityId) { validate(entityId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java index bab6ea8736..6c093509d0 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java @@ -44,6 +44,7 @@ import org.thingsboard.server.dao.util.NoSqlTsLatestDao; import java.util.Collections; import java.util.List; import java.util.Optional; +import java.util.concurrent.ExecutionException; import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.literal; @@ -69,6 +70,16 @@ public class CassandraBaseTimeseriesLatestDao extends AbstractCassandraBaseTimes return findLatest(tenantId, entityId, key, rs -> convertResultToTsKvEntry(key, rs.one())); } + @Override + public TsKvEntry findLatestSync(TenantId tenantId, EntityId entityId, String key) { + try { + return findLatest(tenantId, entityId, key, rs -> convertResultToTsKvEntry(key, rs.one())).get(); + } catch (InterruptedException | ExecutionException e) { + log.error("[{}][{}] Failed to get latest entry for key: {}", tenantId, entityId, key, e); + throw new RuntimeException(e); + } + } + private ListenableFuture findLatest(TenantId tenantId, EntityId entityId, String key, java.util.function.Function function) { BoundStatementBuilder stmtBuilder = new BoundStatementBuilder(getFindLatestStmt().bind()); stmtBuilder.setString(0, entityId.getEntityType().name()); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java index c55abeed7e..4c46d67a46 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java @@ -40,6 +40,8 @@ public interface TimeseriesLatestDao { */ ListenableFuture findLatest(TenantId tenantId, EntityId entityId, String key); + TsKvEntry findLatestSync(TenantId tenantId, EntityId entityId, String key); + ListenableFuture> findAllLatest(TenantId tenantId, EntityId entityId); ListenableFuture saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/VersionedNode.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/VersionedNode.java new file mode 100644 index 0000000000..51286c00c4 --- /dev/null +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/VersionedNode.java @@ -0,0 +1,43 @@ +/** + * Copyright © 2016-2023 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.rule.engine.api; + +import com.fasterxml.jackson.databind.JsonNode; +import org.thingsboard.server.common.data.id.RuleNodeId; +import org.thingsboard.server.common.data.util.TbPair; + +public interface VersionedNode { + + String VERSION_PROPERTY_NAME = "version"; + + TbPair upgrade(RuleNodeId ruleNodeId, JsonNode oldConfiguration); + + default int getVersionOrElseThrowTbNodeException(RuleNodeId ruleNodeId, JsonNode oldConfiguration) throws TbNodeException { + if (oldConfiguration == null) { + throw new TbNodeException("Rule node: [" + this.getClass().getName() + "] " + + "with id: [" + ruleNodeId + "] has null configuration!"); + } else if (!oldConfiguration.isObject()) { + throw new TbNodeException("Rule node: [" + this.getClass().getName() + "] " + + "with id: [" + ruleNodeId + "] has non json object configuration!"); + + } + if (!oldConfiguration.has(VERSION_PROPERTY_NAME)) { + return 0; + } + return oldConfiguration.get(VERSION_PROPERTY_NAME).asInt(); + } + +} diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntityDetails.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/VersionedNodeConfiguration.java similarity index 79% rename from rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntityDetails.java rename to rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/VersionedNodeConfiguration.java index 1ca00f2e5b..09e38282e2 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntityDetails.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/VersionedNodeConfiguration.java @@ -13,10 +13,10 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.rule.engine.util; +package org.thingsboard.rule.engine.api; -public enum EntityDetails { +public interface VersionedNodeConfiguration { - ID, TITLE, COUNTRY, CITY, STATE, ZIP, ADDRESS, ADDRESS2, PHONE, EMAIL, ADDITIONAL_INFO + int getVersion(); } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java index 1b88f79ef5..291c2cf316 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java @@ -26,7 +26,6 @@ import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; -import org.thingsboard.rule.engine.api.TbRelationTypes; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.kv.TsKvEntry; @@ -38,6 +37,7 @@ import org.thingsboard.server.dao.timeseries.TimeseriesService; import java.math.BigDecimal; import java.math.RoundingMode; import java.util.Collections; +import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -68,7 +68,6 @@ public class CalculateDeltaNode implements TbNode { this.ctx = ctx; this.timeseriesService = ctx.getTimeseriesService(); this.useCache = config.isUseCache(); - if (useCache) { cache = new ConcurrentHashMap<>(); } @@ -81,9 +80,6 @@ public class CalculateDeltaNode implements TbNode { return; } JsonNode json = JacksonUtil.toJsonNode(msg.getData()); - if (!json.isObject()) { - throw new IllegalArgumentException("Message body is not an object!"); - } String inputKey = config.getInputValueKey(); if (!json.has(inputKey)) { ctx.tellNext(msg, "Other"); @@ -101,7 +97,7 @@ public class CalculateDeltaNode implements TbNode { BigDecimal delta = BigDecimal.valueOf(previousData != null ? currentValue - previousData.value : 0.0); if (config.isTellFailureIfDeltaIsNegative() && delta.doubleValue() < 0) { - ctx.tellNext(msg, TbRelationTypes.FAILURE); + ctx.tellFailure(msg, new IllegalArgumentException("Delta value is negative!")); return; } @@ -132,18 +128,29 @@ public class CalculateDeltaNode implements TbNode { } } - private ListenableFuture fetchLatestValue(EntityId entityId) { + private ListenableFuture fetchLatestValueAsync(EntityId entityId) { return Futures.transform(timeseriesService.findLatest(ctx.getTenantId(), entityId, Collections.singletonList(config.getInputValueKey())), list -> extractValue(list.get(0)) , ctx.getDbCallbackExecutor()); } + private ValueWithTs fetchLatestValue(EntityId entityId) { + List tsKvEntries = timeseriesService.findLatestSync( + ctx.getTenantId(), + entityId, + Collections.singletonList(config.getInputValueKey())); + return extractValue(tsKvEntries.get(0)); + } + private ListenableFuture getLastValue(EntityId entityId) { - ValueWithTs latestValue; - if (useCache && (latestValue = cache.get(entityId)) != null) { + if (useCache) { + ValueWithTs latestValue; + if ((latestValue = cache.get(entityId)) == null) { + latestValue = fetchLatestValue(entityId); + } return Futures.immediateFuture(latestValue); } else { - return fetchLatestValue(entityId); + return fetchLatestValueAsync(entityId); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/DataToFetch.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/DataToFetch.java new file mode 100644 index 0000000000..4685ae542c --- /dev/null +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/DataToFetch.java @@ -0,0 +1,22 @@ +/** + * Copyright © 2016-2023 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.rule.engine.metadata; + +public enum DataToFetch { + + ATTRIBUTES, LATEST_TELEMETRY, FIELDS + +} diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java index 14c27a7ac7..98970cc4eb 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java @@ -16,10 +16,12 @@ package org.thingsboard.rule.engine.metadata; import lombok.Data; +import org.thingsboard.rule.engine.api.VersionedNodeConfiguration; @Data -public abstract class TbAbstractFetchToNodeConfiguration { +public abstract class TbAbstractFetchToNodeConfiguration implements VersionedNodeConfiguration { private FetchTo fetchTo; + private int version = 1; } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java index 220ed79abf..168319b7d6 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; @@ -23,9 +24,13 @@ import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; +import org.thingsboard.rule.engine.util.EntitiesFieldsAsyncLoader; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgMetaData; import java.util.HashMap; import java.util.List; @@ -38,11 +43,14 @@ import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; @Slf4j public abstract class TbAbstractGetEntityAttrNode extends TbAbstractNodeWithFetchTo { + protected final static String DATA_TO_FETCH_PROPERTY_NAME = "dataToFetch"; + protected static final String OLD_PROPERTY_NAME = "telemetry"; + @Override public void onMsg(TbContext ctx, TbMsg msg) { var msgDataAsObjectNode = FetchTo.DATA.equals(fetchTo) ? getMsgDataAsObjectNode(msg) : null; withCallback(findEntityAsync(ctx, msg.getOriginator()), - entityId -> safeGetAttributes(ctx, msg, entityId, msgDataAsObjectNode), + entityId -> getData(ctx, msg, entityId, msgDataAsObjectNode), t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } @@ -54,18 +62,64 @@ public abstract class TbAbstractGetEntityAttrNode extends Tb } } - private void safeGetAttributes(TbContext ctx, TbMsg msg, T entityId, ObjectNode msgDataAsJsonNode) { + protected abstract void checkDataToFetchSupportedOrElseThrow(DataToFetch dataToFetch) throws TbNodeException; + + private void getData(TbContext ctx, TbMsg msg, T entityId, ObjectNode msgDataAsJsonNode) { var mappingsMap = new HashMap(); - config.getAttrMapping().forEach((key, value) -> { - String patternProcessedSourceKey = TbNodeUtils.processPattern(key, msg); - String patternProcessedTargetKey = TbNodeUtils.processPattern(value, msg); - mappingsMap.put(patternProcessedSourceKey, patternProcessedTargetKey); - }); - var sourceKeys = List.copyOf(mappingsMap.keySet()); - withCallback(config.isTelemetry() ? getLatestTelemetryAsync(ctx, entityId, sourceKeys) : getAttributesAsync(ctx, entityId, sourceKeys), - data -> putDataAndTell(ctx, msg, data, mappingsMap, msgDataAsJsonNode), - t -> ctx.tellFailure(msg, t), - MoreExecutors.directExecutor()); + if (DataToFetch.FIELDS.equals(config.getDataToFetch())) { + config.getAttrMapping().forEach((sourceField, targetKey) -> { + String patternProcessedTargetKey = TbNodeUtils.processPattern(targetKey, msg); + mappingsMap.put(sourceField, patternProcessedTargetKey); + }); + withCallback(collectMappedEntityFieldsAsync(ctx, entityId, mappingsMap), + targetKeysToSourceValuesMap -> { + TbMsgMetaData msgMetaData = msg.getMetaData().copy(); + for (var entry : targetKeysToSourceValuesMap.entrySet()) { + var targetKeyName = entry.getKey(); + var sourceFieldValue = entry.getValue(); + if (FetchTo.DATA.equals(fetchTo)) { + msgDataAsJsonNode.put(targetKeyName, sourceFieldValue); + } else if (FetchTo.METADATA.equals(fetchTo)) { + msgMetaData.putValue(targetKeyName, sourceFieldValue); + } + } + TbMsg outMsg = transformMessage(msg, msgDataAsJsonNode, msgMetaData); + ctx.tellSuccess(outMsg); + }, + t -> ctx.tellFailure(msg, t), + MoreExecutors.directExecutor()); + } else { + config.getAttrMapping().forEach((sourceKey, targetKey) -> { + String patternProcessedSourceKey = TbNodeUtils.processPattern(sourceKey, msg); + String patternProcessedTargetKey = TbNodeUtils.processPattern(targetKey, msg); + mappingsMap.put(patternProcessedSourceKey, patternProcessedTargetKey); + }); + var sourceKeys = List.copyOf(mappingsMap.keySet()); + withCallback(DataToFetch.LATEST_TELEMETRY.equals(config.getDataToFetch()) ? + getLatestTelemetryAsync(ctx, entityId, sourceKeys) : + getAttributesAsync(ctx, entityId, sourceKeys), + data -> putDataAndTell(ctx, msg, data, mappingsMap, msgDataAsJsonNode), + t -> ctx.tellFailure(msg, t), + MoreExecutors.directExecutor()); + } + + } + + private ListenableFuture> collectMappedEntityFieldsAsync(TbContext ctx, EntityId entityId, HashMap mappingsMap) { + return Futures.transform(EntitiesFieldsAsyncLoader.findAsync(ctx, entityId), + fieldsData -> { + var targetKeysToSourceValuesMap = new HashMap(); + for (var mappingEntry : mappingsMap.entrySet()) { + var sourceFieldName = mappingEntry.getKey(); + var targetKeyName = mappingEntry.getValue(); + var sourceFieldValue = fieldsData.getFieldValue(sourceFieldName, true); + if (sourceFieldValue != null) { + targetKeysToSourceValuesMap.put(targetKeyName, sourceFieldValue); + } + } + return targetKeysToSourceValuesMap; + }, ctx.getDbCallbackExecutor() + ); } private ListenableFuture> getAttributesAsync(TbContext ctx, EntityId entityId, List attrKeys) { @@ -95,4 +149,27 @@ public abstract class TbAbstractGetEntityAttrNode extends Tb ctx.tellSuccess(transformMessage(msg, msgData, msgMetaData)); } + protected TbPair upgradeToUseFetchToAndDataToFetch(RuleNodeId ruleNodeId, JsonNode oldConfiguration) throws TbNodeException { + var newConfigObjectNode = (ObjectNode) oldConfiguration; + if (!newConfigObjectNode.has(OLD_PROPERTY_NAME)) { + throw new TbNodeException("Rule node: [" + this.getClass().getName() + "] " + + "with id: [" + ruleNodeId + "] doesn't have property: [" + OLD_PROPERTY_NAME + "]"); + } + var value = newConfigObjectNode.get(OLD_PROPERTY_NAME).asText(); + if ("true".equals(value)) { + newConfigObjectNode.remove(OLD_PROPERTY_NAME); + newConfigObjectNode.put(DATA_TO_FETCH_PROPERTY_NAME, DataToFetch.LATEST_TELEMETRY.name()); + } else if ("false".equals(value)) { + newConfigObjectNode.remove(OLD_PROPERTY_NAME); + newConfigObjectNode.put(DATA_TO_FETCH_PROPERTY_NAME, DataToFetch.ATTRIBUTES.name()); + } else { + throw new TbNodeException("Rule node: [" + this.getClass().getName() + "] " + + "with id: [" + ruleNodeId + "] has property: [" + OLD_PROPERTY_NAME + "] " + + "with unexpected value: [" + value + "] Allowed values: true or false!"); + } + newConfigObjectNode.put(FETCH_TO_PROPERTY_NAME, FetchTo.METADATA.name()); + newConfigObjectNode.put(VERSION_PROPERTY_NAME, 1); + return new TbPair<>(true, newConfigObjectNode); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java index 48d4415693..14227cadc4 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java @@ -22,7 +22,7 @@ import com.google.common.util.concurrent.MoreExecutors; import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeException; -import org.thingsboard.rule.engine.util.EntityDetails; +import org.thingsboard.rule.engine.util.ContactBasedEntityDetails; import org.thingsboard.server.common.data.ContactBased; import org.thingsboard.server.common.data.id.UUIDBased; import org.thingsboard.server.common.msg.TbMsg; @@ -47,7 +47,7 @@ public abstract class TbAbstractGetEntityDetailsNode> getContactBasedFuture(TbContext ctx, TbMsg msg); - protected void checkIfDetailsListIsNotEmptyOrElseThrow(List detailsList) throws TbNodeException { + protected void checkIfDetailsListIsNotEmptyOrElseThrow(List detailsList) throws TbNodeException { if (detailsList == null || detailsList.isEmpty()) { throw new TbNodeException("No entity details selected!"); } @@ -60,87 +60,64 @@ public abstract class TbAbstractGetEntityDetailsNode contactBased, ObjectNode messageData, TbMsgMetaData msgMetaData) { - String prefix = getPrefix(); - String property; - String value; - for (var entityDetails : config.getDetailsList()) { - switch (entityDetails) { + private void fetchEntityDetailsToMsg(ContactBased contactBased, ObjectNode messageData, TbMsgMetaData msgMetaData) { + String value = null; + for (var entityDetail : config.getDetailsList()) { + switch (entityDetail) { case ID: - property = prefix + "id"; value = contactBased.getId().getId().toString(); - setDetail(property, value, messageData, msgMetaData); break; case TITLE: - property = prefix + "title"; value = contactBased.getName(); - setDetail(property, value, messageData, msgMetaData); break; case ADDRESS: - property = prefix + "address"; value = contactBased.getAddress(); - setDetail(property, value, messageData, msgMetaData); break; case ADDRESS2: - property = prefix + "address2"; value = contactBased.getAddress2(); - setDetail(property, value, messageData, msgMetaData); break; case CITY: - property = prefix + "city"; value = contactBased.getCity(); - setDetail(property, value, messageData, msgMetaData); break; case COUNTRY: - property = prefix + "country"; value = contactBased.getCountry(); - setDetail(property, value, messageData, msgMetaData); break; case STATE: - property = prefix + "state"; value = contactBased.getState(); - setDetail(property, value, messageData, msgMetaData); break; case EMAIL: - property = prefix + "email"; value = contactBased.getEmail(); - setDetail(property, value, messageData, msgMetaData); break; case PHONE: - property = prefix + "phone"; value = contactBased.getPhone(); - setDetail(property, value, messageData, msgMetaData); break; case ZIP: - property = prefix + "zip"; value = contactBased.getZip(); - setDetail(property, value, messageData, msgMetaData); break; case ADDITIONAL_INFO: if (contactBased.getAdditionalInfo().hasNonNull("description")) { - property = prefix + "additionalInfo"; value = contactBased.getAdditionalInfo().get("description").asText(); - setDetail(property, value, messageData, msgMetaData); } break; } + if (value == null) { + continue; + } + setDetail(entityDetail.getRuleEngineName(), value, messageData, msgMetaData); } } private void setDetail(String property, String value, ObjectNode messageData, TbMsgMetaData msgMetaData) { - if (value == null) { - return; - } + String fieldName = getPrefix() + property; if (FetchTo.METADATA.equals(fetchTo)) { - msgMetaData.putValue(property, value); - } - if (FetchTo.DATA.equals(fetchTo)) { - messageData.put(property, value); + msgMetaData.putValue(fieldName, value); + } else if (FetchTo.DATA.equals(fetchTo)) { + messageData.put(fieldName, value); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java index 4bf150ebff..d7f24d39b1 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java @@ -17,7 +17,7 @@ package org.thingsboard.rule.engine.metadata; import lombok.Data; import lombok.EqualsAndHashCode; -import org.thingsboard.rule.engine.util.EntityDetails; +import org.thingsboard.rule.engine.util.ContactBasedEntityDetails; import java.util.List; @@ -25,6 +25,6 @@ import java.util.List; @EqualsAndHashCode(callSuper = true) public abstract class TbAbstractGetEntityDetailsNodeConfiguration extends TbAbstractFetchToNodeConfiguration { - private List detailsList; + private List detailsList; } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java index 7cb7f747e6..229b541863 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.AsyncFunction; import com.google.common.util.concurrent.Futures; @@ -24,15 +25,20 @@ import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; +import org.thingsboard.rule.engine.api.VersionedNode; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import java.util.NoSuchElementException; @Slf4j -public abstract class TbAbstractNodeWithFetchTo implements TbNode { +public abstract class TbAbstractNodeWithFetchTo implements TbNode, VersionedNode { + + protected final static String FETCH_TO_PROPERTY_NAME = "fetchTo"; protected C config; protected FetchTo fetchTo; @@ -86,4 +92,34 @@ public abstract class TbAbstractNodeWithFetchTo upgradeRuleNodesWithOldPropertyToUseFetchTo( + RuleNodeId ruleNodeId, + JsonNode oldConfiguration, + String oldProperty, + String ifTrue, + String ifFalse + ) throws TbNodeException { + var newConfigObjectNode = (ObjectNode) oldConfiguration; + if (!newConfigObjectNode.has(oldProperty)) { + throw new TbNodeException("Rule node: [" + this.getClass().getName() + "] " + + "with id: [" + ruleNodeId + "] doesn't have property: [" + oldProperty + "]"); + } + var value = newConfigObjectNode.get(oldProperty).asText(); + if ("true".equals(value)) { + newConfigObjectNode.remove(oldProperty); + newConfigObjectNode.put(FETCH_TO_PROPERTY_NAME, ifTrue); + newConfigObjectNode.put(VERSION_PROPERTY_NAME, 1); + return new TbPair<>(true, newConfigObjectNode); + } else if ("false".equals(value)) { + newConfigObjectNode.remove(oldProperty); + newConfigObjectNode.put(FETCH_TO_PROPERTY_NAME, ifFalse); + newConfigObjectNode.put(VERSION_PROPERTY_NAME, 1); + return new TbPair<>(true, newConfigObjectNode); + } else { + throw new TbNodeException("Rule node: [" + this.getClass().getName() + "] " + + "with id: [" + ruleNodeId + "] has property: [" + oldProperty + "] " + + "with unexpected value: [" + value + "] Allowed values: true or false!"); + } + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java index 0d2b110cb9..514ad20c5d 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.JsonNode; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.RuleNode; @@ -24,8 +25,10 @@ import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.security.DeviceCredentialsType; +import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; import java.util.concurrent.ExecutionException; @@ -85,4 +88,23 @@ public class TbFetchDeviceCredentialsNode extends TbAbstractNodeWithFetchTo upgrade(RuleNodeId ruleNodeId, JsonNode oldConfiguration) { + try { + int oldVersion = getVersionOrElseThrowTbNodeException(ruleNodeId, oldConfiguration); + if (oldVersion == 0) { + return upgradeRuleNodesWithOldPropertyToUseFetchTo( + ruleNodeId, + oldConfiguration, + "fetchToMetadata", + FetchTo.METADATA.name(), + FetchTo.DATA.name() + ); + } + } catch (TbNodeException e) { + log.warn(e.getMessage()); + } + return new TbPair<>(false, oldConfiguration); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java index 4d306e93ed..a4c2a9f3a1 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; @@ -24,7 +25,9 @@ import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.plugin.ComponentType; +import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; /** @@ -53,4 +56,23 @@ public class TbGetAttributesNode extends TbAbstractGetAttributesNode upgrade(RuleNodeId ruleNodeId, JsonNode oldConfiguration) { + try { + int oldVersion = getVersionOrElseThrowTbNodeException(ruleNodeId, oldConfiguration); + if (oldVersion == 0) { + return upgradeRuleNodesWithOldPropertyToUseFetchTo( + ruleNodeId, + oldConfiguration, + "fetchToData", + FetchTo.DATA.name(), + FetchTo.METADATA.name() + ); + } + } catch (TbNodeException e) { + log.warn(e.getMessage()); + } + return new TbPair<>(false, oldConfiguration); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java index dc878ab8ad..25d4cd8994 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java @@ -15,8 +15,10 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; +import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; @@ -25,8 +27,11 @@ import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.util.EntitiesCustomerIdAsyncLoader; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.plugin.ComponentType; +import org.thingsboard.server.common.data.util.TbPair; +@Slf4j @RuleNode( type = ComponentType.ENRICHMENT, name = "customer attributes", @@ -46,9 +51,18 @@ public class TbGetCustomerAttributeNode extends TbAbstractGetEntityAttrNode findEntityAsync(TbContext ctx, EntityId originator) { return Futures.transformAsync(EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, originator), @@ -57,4 +71,17 @@ public class TbGetCustomerAttributeNode extends TbAbstractGetEntityAttrNode upgrade(RuleNodeId ruleNodeId, JsonNode oldConfiguration) { + try { + int oldVersion = getVersionOrElseThrowTbNodeException(ruleNodeId, oldConfiguration); + if (oldVersion == 0) { + return upgradeToUseFetchToAndDataToFetch(ruleNodeId, oldConfiguration); + } + } catch (TbNodeException e) { + log.warn(e.getMessage()); + } + return new TbPair<>(false, oldConfiguration); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java index 65f11facf5..1026d8ca2b 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; @@ -32,8 +33,10 @@ import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityViewId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.plugin.ComponentType; +import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; import java.util.NoSuchElementException; @@ -103,4 +106,23 @@ public class TbGetCustomerDetailsNode extends TbAbstractGetEntityDetailsNode upgrade(RuleNodeId ruleNodeId, JsonNode oldConfiguration) { + try { + int oldVersion = getVersionOrElseThrowTbNodeException(ruleNodeId, oldConfiguration); + if (oldVersion == 0) { + return upgradeRuleNodesWithOldPropertyToUseFetchTo( + ruleNodeId, + oldConfiguration, + "addToMetadata", + FetchTo.METADATA.name(), + FetchTo.DATA.name() + ); + } + } catch (TbNodeException e) { + log.warn(e.getMessage()); + } + return new TbPair<>(false, oldConfiguration); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java index fdbb412e6b..dd689ed38d 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; @@ -25,7 +26,9 @@ import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.util.EntitiesRelatedDeviceIdAsyncLoader; import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.plugin.ComponentType; +import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; @Slf4j @@ -56,4 +59,23 @@ public class TbGetDeviceAttrNode extends TbAbstractGetAttributesNode upgrade(RuleNodeId ruleNodeId, JsonNode oldConfiguration) { + try { + int oldVersion = getVersionOrElseThrowTbNodeException(ruleNodeId, oldConfiguration); + if (oldVersion == 0) { + return upgradeRuleNodesWithOldPropertyToUseFetchTo( + ruleNodeId, + oldConfiguration, + "fetchToData", + FetchTo.DATA.name(), + FetchTo.METADATA.name() + ); + } + } catch (TbNodeException e) { + log.warn(e.getMessage()); + } + return new TbPair<>(false, oldConfiguration); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java index 18134a5040..eec7b497e1 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java @@ -27,7 +27,7 @@ import java.util.Map; public class TbGetEntityAttrNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration { private Map attrMapping; - private boolean isTelemetry; + private DataToFetch dataToFetch; @Override public TbGetEntityAttrNodeConfiguration defaultConfiguration() { @@ -35,7 +35,7 @@ public class TbGetEntityAttrNodeConfiguration extends TbAbstractFetchToNodeConfi var attrMapping = new HashMap(); attrMapping.putIfAbsent("alarmThreshold", "threshold"); configuration.setAttrMapping(attrMapping); - configuration.setTelemetry(false); + configuration.setDataToFetch(DataToFetch.ATTRIBUTES); configuration.setFetchTo(FetchTo.METADATA); return configuration; } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java index a75e0e50a0..93244ef477 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java @@ -15,9 +15,12 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; +import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; @@ -25,7 +28,9 @@ import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.util.EntitiesFieldsAsyncLoader; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.plugin.ComponentType; +import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; @@ -37,6 +42,7 @@ import static org.thingsboard.common.util.DonAsynchron.withCallback; /** * Created by ashvayka on 19.01.18. */ +@Slf4j @RuleNode(type = ComponentType.ENRICHMENT, name = "originator fields", configClazz = TbGetOriginatorFieldsConfiguration.class, @@ -95,4 +101,20 @@ public class TbGetOriginatorFieldsNode extends TbAbstractNodeWithFetchTo upgrade(RuleNodeId ruleNodeId, JsonNode oldConfiguration) { + try { + int oldVersion = getVersionOrElseThrowTbNodeException(ruleNodeId, oldConfiguration); + if (oldVersion == 0) { + var newConfigObjectNode = (ObjectNode) oldConfiguration; + newConfigObjectNode.put(FETCH_TO_PROPERTY_NAME, FetchTo.METADATA.name()); + newConfigObjectNode.put(VERSION_PROPERTY_NAME, 1); + return new TbPair<>(true, newConfigObjectNode); + } + } catch (TbNodeException e) { + log.warn(e.getMessage()); + } + return new TbPair<>(false, oldConfiguration); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java index 26491509b8..3f4d95a08f 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java @@ -37,7 +37,7 @@ public class TbGetRelatedAttrNodeConfiguration extends TbGetEntityAttrNodeConfig var attrMapping = new HashMap(); attrMapping.putIfAbsent("serialNumber", "sn"); configuration.setAttrMapping(attrMapping); - configuration.setTelemetry(false); + configuration.setDataToFetch(DataToFetch.ATTRIBUTES); configuration.setFetchTo(FetchTo.METADATA); var relationsQuery = new RelationsQuery(); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java index 50cd2a53d3..d7fd22d306 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java @@ -15,8 +15,10 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; +import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; @@ -24,8 +26,13 @@ import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.util.EntitiesRelatedEntityIdAsyncLoader; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.plugin.ComponentType; +import org.thingsboard.server.common.data.util.TbPair; +import java.util.Arrays; + +@Slf4j @RuleNode( type = ComponentType.ENRICHMENT, name = "related attributes", @@ -47,6 +54,7 @@ public class TbGetRelatedAttributeNode extends TbAbstractGetEntityAttrNode upgrade(RuleNodeId ruleNodeId, JsonNode oldConfiguration) { + try { + int oldVersion = getVersionOrElseThrowTbNodeException(ruleNodeId, oldConfiguration); + if (oldVersion == 0) { + return upgradeToUseFetchToAndDataToFetch(ruleNodeId, oldConfiguration); + } + } catch (TbNodeException e) { + log.warn(e.getMessage()); + } + return new TbPair<>(false, oldConfiguration); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java index e32a093e9b..e2ebc6a016 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; @@ -24,8 +25,10 @@ import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.plugin.ComponentType; +import org.thingsboard.server.common.data.util.TbPair; @Slf4j @RuleNode( @@ -45,6 +48,7 @@ public class TbGetTenantAttributeNode extends TbAbstractGetEntityAttrNode upgrade(RuleNodeId ruleNodeId, JsonNode oldConfiguration) { + try { + int oldVersion = getVersionOrElseThrowTbNodeException(ruleNodeId, oldConfiguration); + if (oldVersion == 0) { + return upgradeToUseFetchToAndDataToFetch(ruleNodeId, oldConfiguration); + } + } catch (TbNodeException e) { + log.warn(e.getMessage()); + } + return new TbPair<>(false, oldConfiguration); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java index 5e3db845dc..31371b8b1c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.api.RuleNode; @@ -23,8 +24,10 @@ import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.Tenant; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.plugin.ComponentType; +import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; @Slf4j @@ -58,4 +61,23 @@ public class TbGetTenantDetailsNode extends TbAbstractGetEntityDetailsNode upgrade(RuleNodeId ruleNodeId, JsonNode oldConfiguration) { + try { + int oldVersion = getVersionOrElseThrowTbNodeException(ruleNodeId, oldConfiguration); + if (oldVersion == 0) { + return upgradeRuleNodesWithOldPropertyToUseFetchTo( + ruleNodeId, + oldConfiguration, + "addToMetadata", + FetchTo.METADATA.name(), + FetchTo.DATA.name() + ); + } + } catch (TbNodeException e) { + log.warn(e.getMessage()); + } + return new TbPair<>(false, oldConfiguration); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/ContactBasedEntityDetails.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/ContactBasedEntityDetails.java new file mode 100644 index 0000000000..83aaff6a6d --- /dev/null +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/ContactBasedEntityDetails.java @@ -0,0 +1,41 @@ +/** + * Copyright © 2016-2023 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.rule.engine.util; + +import lombok.Getter; + +public enum ContactBasedEntityDetails { + + ID("id"), + TITLE("title"), + COUNTRY("country"), + CITY("city"), + STATE("state"), + ZIP("zip"), + ADDRESS("address"), + ADDRESS2("address2"), + PHONE("phone"), + EMAIL("email"), + ADDITIONAL_INFO("additionalInfo"); + + @Getter + private final String ruleEngineName; + + ContactBasedEntityDetails(String ruleEngineName) { + this.ruleEngineName = ruleEngineName; + } + +} diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeTest.java index 27fb5bf785..86765446e0 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeTest.java @@ -18,6 +18,7 @@ package org.thingsboard.rule.engine.metadata; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; +import org.assertj.core.api.Assertions; import org.jetbrains.annotations.NotNull; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -127,6 +128,22 @@ public class CalculateDeltaNodeTest { verify(ctxMock, never()).tellFailure(any(), any()); } + @Test + public void givenInvalidMsgDataType_whenOnMsg_thenShouldTellNextOther() { + // GIVEN + var msgData = "[]"; + var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); + + // WHEN + node.onMsg(ctxMock, msg); + + // THEN + verify(ctxMock, times(1)).tellNext(eq(msg), eq("Other")); + verify(ctxMock, never()).tellSuccess(any()); + verify(ctxMock, never()).tellFailure(any(), any()); + } + + @Test public void givenInputKeyIsNotPresent_whenOnMsg_thenShouldTellNextOther() { // GIVEN @@ -142,15 +159,16 @@ public class CalculateDeltaNodeTest { } @Test - public void givenDoubleValue_whenOnMsg_thenShouldTellSuccess() throws TbNodeException { + public void givenDoubleValue_whenOnMsgAndCachingOff_thenShouldTellSuccess() throws TbNodeException { // GIVEN config.setRound(1); config.setInputValueKey("temperature"); config.setOutputValueKey("temp_delta"); + config.setUseCache(false); nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); node.init(ctxMock, nodeConfiguration); - mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new DoubleDataEntry("temperature", 40.5))); + mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new DoubleDataEntry("temperature", 40.5))); var msgData = "{\"temperature\": 42,\"airPressure\":123}"; var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); @@ -172,14 +190,15 @@ public class CalculateDeltaNodeTest { } @Test - public void givenLongStringValue_whenOnMsg_thenShouldTellSuccess() throws TbNodeException { + public void givenLongStringValue_whenOnMsgAndCachingOff_thenShouldTellSuccess() throws TbNodeException { // GIVEN config.setInputValueKey("temperature"); config.setOutputValueKey("temp_delta"); + config.setUseCache(false); nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); node.init(ctxMock, nodeConfiguration); - mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("temperature", 40L))); + mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("temperature", 40L))); var msgData = "{\"temperature\": 42,\"airPressure\":123}"; var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); @@ -201,14 +220,15 @@ public class CalculateDeltaNodeTest { } @Test - public void givenValidStringValue_whenOnMsg_thenShouldTellSuccess() throws TbNodeException { + public void givenValidStringValue_whenOnMsgAndCachingOff_thenShouldTellSuccess() throws TbNodeException { // GIVEN config.setInputValueKey("temperature"); config.setOutputValueKey("temp_delta"); + config.setUseCache(false); nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); node.init(ctxMock, nodeConfiguration); - mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry("temperature", "40.0"))); + mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry("temperature", "40.0"))); var msgData = "{\"temperature\": 42,\"airPressure\":123}"; var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); @@ -289,14 +309,15 @@ public class CalculateDeltaNodeTest { } @Test - public void givenLastValueIsNull_whenOnMsh_thenDeltaShouldBeZero() throws TbNodeException { + public void givenLastValueIsNull_whenOnMsgAndCachingOff_thenDeltaShouldBeZero() throws TbNodeException { // GIVEN config.setInputValueKey("temperature"); config.setOutputValueKey("temp_delta"); + config.setUseCache(false); nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); node.init(ctxMock, nodeConfiguration); - mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new DoubleDataEntry("temperature", null))); + mockFindLatestAsync(new BasicTsKvEntry(System.currentTimeMillis(), new DoubleDataEntry("temperature", null))); var msgData = "{\"temperature\": 42,\"airPressure\":123}"; var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); @@ -332,24 +353,6 @@ public class CalculateDeltaNodeTest { // WHEN node.onMsg(ctxMock, msg); - // THEN - verify(ctxMock, times(1)).tellNext(msg, "Failure"); - verify(ctxMock, never()).tellSuccess(any()); - verify(ctxMock, never()).tellFailure(any(), any()); - verify(ctxMock, never()).tellNext(any(), anySet()); - } - - @Test - public void givenInvalidStringValue_whenOnMsg_thenException() { - // GIVEN - mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry("pulseCounter", "high"))); - - var msgData = "{\"pulseCounter\":\"123\"}"; - var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); - - // WHEN - node.onMsg(ctxMock, msg); - // THEN var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); var actualExceptionCaptor = ArgumentCaptor.forClass(Exception.class); @@ -359,40 +362,41 @@ public class CalculateDeltaNodeTest { verify(ctxMock, never()).tellNext(any(), anyString()); verify(ctxMock, never()).tellNext(any(), anySet()); - var expectedExceptionMsg = "Calculation failed. Unable to parse value [high] of telemetry [pulseCounter] to Double"; + var expectedExceptionMsg = "Delta value is negative!"; var actualException = actualExceptionCaptor.getValue(); assertEquals(msg, actualMsgCaptor.getValue()); assertInstanceOf(IllegalArgumentException.class, actualException); assertEquals(expectedExceptionMsg, actualException.getMessage()); + } @Test - public void givenBooleanValue_whenOnMsg_thenException() { + public void givenInvalidStringValue_whenOnMsg_thenException() { // GIVEN - mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry("pulseCounter", false))); + mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry("pulseCounter", "high"))); - var msgData = "{\"pulseCounter\":true}"; + var msgData = "{\"pulseCounter\":\"123\"}"; var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); - // WHEN - node.onMsg(ctxMock, msg); - - // THEN - var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); - var actualExceptionCaptor = ArgumentCaptor.forClass(Exception.class); + // WHEN-THEN + Assertions.assertThatThrownBy(() -> node.onMsg(ctxMock, msg)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Calculation failed. Unable to parse value [high] of telemetry [pulseCounter] to Double"); + } - verify(ctxMock, times(1)).tellFailure(actualMsgCaptor.capture(), actualExceptionCaptor.capture()); - verify(ctxMock, never()).tellSuccess(any()); - verify(ctxMock, never()).tellNext(any(), anyString()); - verify(ctxMock, never()).tellNext(any(), anySet()); + @Test + public void givenBooleanValue_whenOnMsg_thenException() { + // GIVEN + mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry("pulseCounter", false))); - var expectedExceptionMsg = "Calculation failed. Boolean values are not supported!"; - var actualException = actualExceptionCaptor.getValue(); + var msgData = "{\"pulseCounter\":true}"; + var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); - assertEquals(msg, actualMsgCaptor.getValue()); - assertInstanceOf(IllegalArgumentException.class, actualException); - assertEquals(expectedExceptionMsg, actualException.getMessage()); + // WHEN-THEN + Assertions.assertThatThrownBy(() -> node.onMsg(ctxMock, msg)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Calculation failed. Boolean values are not supported!"); } @Test @@ -403,27 +407,20 @@ public class CalculateDeltaNodeTest { var msgData = "{\"pulseCounter\":{\"isActive\":true}}"; var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); - // WHEN - node.onMsg(ctxMock, msg); - - // THEN - var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); - var actualExceptionCaptor = ArgumentCaptor.forClass(Exception.class); - - verify(ctxMock, times(1)).tellFailure(actualMsgCaptor.capture(), actualExceptionCaptor.capture()); - verify(ctxMock, never()).tellSuccess(any()); - verify(ctxMock, never()).tellNext(any(), anyString()); - verify(ctxMock, never()).tellNext(any(), anySet()); - - var expectedExceptionMsg = "Calculation failed. JSON values are not supported!"; - var actualException = actualExceptionCaptor.getValue(); - - assertEquals(msg, actualMsgCaptor.getValue()); - assertInstanceOf(IllegalArgumentException.class, actualException); - assertEquals(expectedExceptionMsg, actualException.getMessage()); + // WHEN-THEN + Assertions.assertThatThrownBy(() -> node.onMsg(ctxMock, msg)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Calculation failed. JSON values are not supported!"); } private void mockFindLatest(TsKvEntry tsKvEntry) { + when(ctxMock.getTenantId()).thenReturn(TENANT_ID); + when(timeseriesServiceMock.findLatestSync( + eq(TENANT_ID), eq(DUMMY_DEVICE_ORIGINATOR), argThat(new ListMatcher<>(List.of(tsKvEntry.getKey()))) + )).thenReturn(List.of(tsKvEntry)); + } + + private void mockFindLatestAsync(TsKvEntry tsKvEntry) { when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); when(ctxMock.getTenantId()).thenReturn(TENANT_ID); when(timeseriesServiceMock.findLatest( diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java index b1d72c4db8..bc1ba87322 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java @@ -134,6 +134,34 @@ public class TbGetCustomerAttributeNodeTest { verify(ctxMock, never()).tellSuccess(any()); } + @Test + public void givenConfigWithNullDataToFetch_whenInit_thenException() { + // GIVEN + config.setDataToFetch(null); + nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + + // WHEN + var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); + + // THEN + assertThat(exception.getMessage()).isEqualTo("DataToFetch property has invalid value: null. Only ATTRIBUTES and LATEST_TELEMETRY values supported!"); + verify(ctxMock, never()).tellSuccess(any()); + } + + @Test + public void givenConfigWithUnsupportedDataToFetch_whenInit_thenException() { + // GIVEN + config.setDataToFetch(DataToFetch.FIELDS); + nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + + // WHEN + var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); + + // THEN + assertThat(exception.getMessage()).isEqualTo("DataToFetch property has invalid value: FIELDS. Only ATTRIBUTES and LATEST_TELEMETRY values supported!"); + verify(ctxMock, never()).tellSuccess(any()); + } + @Test public void givenDefaultConfig_whenInit_thenOK() throws TbNodeException { // GIVEN @@ -144,7 +172,7 @@ public class TbGetCustomerAttributeNodeTest { // THEN assertThat(node.config).isEqualTo(config); assertThat(config.getAttrMapping()).isEqualTo(Map.of("alarmThreshold", "threshold")); - assertThat(config.isTelemetry()).isEqualTo(false); + assertThat(config.getDataToFetch()).isEqualTo(DataToFetch.ATTRIBUTES); assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA); } @@ -155,7 +183,7 @@ public class TbGetCustomerAttributeNodeTest { "sourceAttr1", "targetKey1", "sourceAttr2", "targetKey2", "sourceAttr3", "targetKey3")); - config.setTelemetry(true); + config.setDataToFetch(DataToFetch.LATEST_TELEMETRY); config.setFetchTo(FetchTo.DATA); nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); @@ -168,7 +196,7 @@ public class TbGetCustomerAttributeNodeTest { "sourceAttr1", "targetKey1", "sourceAttr2", "targetKey2", "sourceAttr3", "targetKey3")); - assertThat(config.isTelemetry()).isEqualTo(true); + assertThat(config.getDataToFetch()).isEqualTo(DataToFetch.LATEST_TELEMETRY); assertThat(node.fetchTo).isEqualTo(FetchTo.DATA); } @@ -245,7 +273,7 @@ public class TbGetCustomerAttributeNodeTest { var device = new Device(new DeviceId(UUID.randomUUID())); device.setCustomerId(CUSTOMER_ID); - prepareMsgAndConfig(FetchTo.DATA, false, device.getId()); + prepareMsgAndConfig(FetchTo.DATA, DataToFetch.ATTRIBUTES, device.getId()); List attributesList = List.of( new BaseAttributeKvEntry(new StringDataEntry("sourceKey1", "sourceValue1"), 1L), @@ -292,7 +320,7 @@ public class TbGetCustomerAttributeNodeTest { var user = new User(new UserId(UUID.randomUUID())); user.setCustomerId(CUSTOMER_ID); - prepareMsgAndConfig(FetchTo.METADATA, false, user.getId()); + prepareMsgAndConfig(FetchTo.METADATA, DataToFetch.ATTRIBUTES, user.getId()); List attributesList = List.of( new BaseAttributeKvEntry(new StringDataEntry("sourceKey1", "sourceValue1"), 1L), @@ -338,7 +366,7 @@ public class TbGetCustomerAttributeNodeTest { // GIVEN var customer = new Customer(new CustomerId(UUID.randomUUID())); - prepareMsgAndConfig(FetchTo.DATA, true, customer.getId()); + prepareMsgAndConfig(FetchTo.DATA, DataToFetch.LATEST_TELEMETRY, customer.getId()); List timeseriesList = List.of( new BasicTsKvEntry(1L, new StringDataEntry("sourceKey1", "sourceValue1")), @@ -382,7 +410,7 @@ public class TbGetCustomerAttributeNodeTest { var asset = new Asset(new AssetId(UUID.randomUUID())); asset.setCustomerId(new CustomerId(UUID.randomUUID())); - prepareMsgAndConfig(FetchTo.METADATA, true, asset.getId()); + prepareMsgAndConfig(FetchTo.METADATA, DataToFetch.LATEST_TELEMETRY, asset.getId()); List timeseriesList = List.of( new BasicTsKvEntry(1L, new StringDataEntry("sourceKey1", "sourceValue1")), @@ -423,12 +451,12 @@ public class TbGetCustomerAttributeNodeTest { assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData); } - private void prepareMsgAndConfig(FetchTo fetchTo, boolean isTelemetry, EntityId originator) { + private void prepareMsgAndConfig(FetchTo fetchTo, DataToFetch dataToFetch, EntityId originator) { config.setAttrMapping(Map.of( "sourceKey1", "targetKey1", "${metaDataPattern1}", "$[messageBodyPattern1]", "$[messageBodyPattern2]", "${metaDataPattern2}")); - config.setTelemetry(isTelemetry); + config.setDataToFetch(dataToFetch); config.setFetchTo(fetchTo); node.config = config; diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeTest.java index a9b32f710a..fc61b0d2b7 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeTest.java @@ -29,7 +29,7 @@ import org.thingsboard.common.util.ListeningExecutor; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; -import org.thingsboard.rule.engine.util.EntityDetails; +import org.thingsboard.rule.engine.util.ContactBasedEntityDetails; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.Device; @@ -130,7 +130,7 @@ public class TbGetCustomerDetailsNodeTest { @Test public void givenConfigWithNullFetchTo_whenInit_thenException() { // GIVEN - config.setDetailsList(List.of(EntityDetails.ID)); + config.setDetailsList(List.of(ContactBasedEntityDetails.ID)); config.setFetchTo(null); nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); @@ -151,7 +151,7 @@ public class TbGetCustomerDetailsNodeTest { @Test public void givenCustomConfig_whenInit_thenOK() throws TbNodeException { // GIVEN - config.setDetailsList(List.of(EntityDetails.ID, EntityDetails.PHONE)); + config.setDetailsList(List.of(ContactBasedEntityDetails.ID, ContactBasedEntityDetails.PHONE)); config.setFetchTo(FetchTo.METADATA); nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); @@ -160,7 +160,7 @@ public class TbGetCustomerDetailsNodeTest { // THEN assertThat(node.config).isEqualTo(config); - assertThat(config.getDetailsList()).isEqualTo(List.of(EntityDetails.ID, EntityDetails.PHONE)); + assertThat(config.getDetailsList()).isEqualTo(List.of(ContactBasedEntityDetails.ID, ContactBasedEntityDetails.PHONE)); assertThat(config.getFetchTo()).isEqualTo(FetchTo.METADATA); assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA); } @@ -186,7 +186,7 @@ public class TbGetCustomerDetailsNodeTest { device.setId(new DeviceId(UUID.randomUUID())); device.setCustomerId(customer.getId()); - prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.values()), device.getId()); + prepareMsgAndConfig(FetchTo.DATA, List.of(ContactBasedEntityDetails.values()), device.getId()); when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); when(deviceServiceMock.findDeviceByIdAsync(eq(TENANT_ID), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); @@ -228,7 +228,7 @@ public class TbGetCustomerDetailsNodeTest { asset.setId(new AssetId(UUID.randomUUID())); asset.setCustomerId(customer.getId()); - prepareMsgAndConfig(FetchTo.METADATA, List.of(EntityDetails.ID, EntityDetails.TITLE, EntityDetails.PHONE), asset.getId()); + prepareMsgAndConfig(FetchTo.METADATA, List.of(ContactBasedEntityDetails.ID, ContactBasedEntityDetails.TITLE, ContactBasedEntityDetails.PHONE), asset.getId()); when(ctxMock.getAssetService()).thenReturn(assetServiceMock); when(assetServiceMock.findAssetByIdAsync(eq(TENANT_ID), eq(asset.getId()))).thenReturn(Futures.immediateFuture(asset)); @@ -266,7 +266,7 @@ public class TbGetCustomerDetailsNodeTest { user.setId(new UserId(UUID.randomUUID())); user.setCustomerId(customer.getId()); - prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2), user.getId()); + prepareMsgAndConfig(FetchTo.DATA, List.of(ContactBasedEntityDetails.ZIP, ContactBasedEntityDetails.ADDRESS, ContactBasedEntityDetails.ADDRESS2), user.getId()); when(ctxMock.getUserService()).thenReturn(userServiceMock); when(userServiceMock.findUserByIdAsync(eq(TENANT_ID), eq(user.getId()))).thenReturn(Futures.immediateFuture(user)); @@ -295,7 +295,7 @@ public class TbGetCustomerDetailsNodeTest { edge.setId(new EdgeId(UUID.randomUUID())); edge.setCustomerId(customer.getId()); - prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2), edge.getId()); + prepareMsgAndConfig(FetchTo.DATA, List.of(ContactBasedEntityDetails.ZIP, ContactBasedEntityDetails.ADDRESS, ContactBasedEntityDetails.ADDRESS2), edge.getId()); when(ctxMock.getTenantId()).thenReturn(TENANT_ID); @@ -327,7 +327,7 @@ public class TbGetCustomerDetailsNodeTest { edge.setId(new EdgeId(UUID.randomUUID())); edge.setCustomerId(customer.getId()); - prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2), edge.getId()); + prepareMsgAndConfig(FetchTo.DATA, List.of(ContactBasedEntityDetails.ZIP, ContactBasedEntityDetails.ADDRESS, ContactBasedEntityDetails.ADDRESS2), edge.getId()); when(ctxMock.getTenantId()).thenReturn(TENANT_ID); @@ -356,7 +356,7 @@ public class TbGetCustomerDetailsNodeTest { device.setId(new DeviceId(UUID.randomUUID())); device.setName("Thermostat"); - prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2), device.getId()); + prepareMsgAndConfig(FetchTo.DATA, List.of(ContactBasedEntityDetails.ZIP, ContactBasedEntityDetails.ADDRESS, ContactBasedEntityDetails.ADDRESS2), device.getId()); when(ctxMock.getTenantId()).thenReturn(TENANT_ID); @@ -394,7 +394,7 @@ public class TbGetCustomerDetailsNodeTest { device.setId(new DeviceId(UUID.randomUUID())); device.setCustomerId(customer.getId()); - prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ADDITIONAL_INFO), device.getId()); + prepareMsgAndConfig(FetchTo.DATA, List.of(ContactBasedEntityDetails.ADDITIONAL_INFO), device.getId()); when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); when(deviceServiceMock.findDeviceByIdAsync(eq(TENANT_ID), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); @@ -422,7 +422,7 @@ public class TbGetCustomerDetailsNodeTest { var dashboard = new Dashboard(); dashboard.setId(new DashboardId(UUID.randomUUID())); - prepareMsgAndConfig(FetchTo.METADATA, List.of(EntityDetails.STATE), dashboard.getId()); + prepareMsgAndConfig(FetchTo.METADATA, List.of(ContactBasedEntityDetails.STATE), dashboard.getId()); // WHEN node.onMsg(ctxMock, msg); @@ -444,7 +444,7 @@ public class TbGetCustomerDetailsNodeTest { assertThat(actualException.getMessage()).isEqualTo("Entity with entityType 'DASHBOARD' is not supported."); } - private void prepareMsgAndConfig(FetchTo fetchTo, List detailsList, EntityId originator) { + private void prepareMsgAndConfig(FetchTo fetchTo, List detailsList, EntityId originator) { config.setDetailsList(detailsList); config.setFetchTo(fetchTo); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java index a033f18cf8..f5684d1786 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java @@ -32,19 +32,9 @@ import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.data.RelationsQuery; -import org.thingsboard.server.common.data.Customer; -import org.thingsboard.server.common.data.Dashboard; -import org.thingsboard.server.common.data.Device; -import org.thingsboard.server.common.data.EntityView; -import org.thingsboard.server.common.data.Tenant; -import org.thingsboard.server.common.data.User; -import org.thingsboard.server.common.data.id.CustomerId; -import org.thingsboard.server.common.data.id.DashboardId; -import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.EntityViewId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.id.UserId; +import org.thingsboard.server.common.data.*; +import org.thingsboard.server.common.data.asset.Asset; +import org.thingsboard.server.common.data.id.*; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; @@ -55,15 +45,13 @@ import org.thingsboard.server.common.data.relation.EntitySearchDirection; import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.common.msg.session.SessionMsgType; import org.thingsboard.server.dao.attributes.AttributesService; +import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.timeseries.TimeseriesService; -import java.util.Collections; -import java.util.List; -import java.util.Map; -import java.util.NoSuchElementException; -import java.util.UUID; +import java.util.*; import java.util.concurrent.Callable; import static org.assertj.core.api.Assertions.assertThat; @@ -108,6 +96,8 @@ public class TbGetRelatedAttributeNodeTest { private TimeseriesService timeseriesServiceMock; @Mock private RelationService relationServiceMock; + @Mock + private DeviceService deviceServiceMock; private TbGetRelatedAttributeNode node; private TbGetRelatedAttrNodeConfiguration config; private TbNodeConfiguration nodeConfiguration; @@ -136,6 +126,20 @@ public class TbGetRelatedAttributeNodeTest { verify(ctxMock, never()).tellSuccess(any()); } + @Test + public void givenConfigWithNullDataToFetch_whenInit_thenException() { + // GIVEN + config.setDataToFetch(null); + nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + + // WHEN + var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); + + // THEN + assertThat(exception.getMessage()).isEqualTo("DataToFetch property cannot be null! Supported values are: " + Arrays.toString(DataToFetch.values())); + verify(ctxMock, never()).tellSuccess(any()); + } + @Test public void givenDefaultConfig_whenInit_thenOK() throws TbNodeException { // GIVEN @@ -147,7 +151,7 @@ public class TbGetRelatedAttributeNodeTest { var nodeConfig = (TbGetRelatedAttrNodeConfiguration) node.config; assertThat(nodeConfig).isEqualTo(config); assertThat(nodeConfig.getAttrMapping()).isEqualTo(Map.of("serialNumber", "sn")); - assertThat(nodeConfig.isTelemetry()).isEqualTo(false); + assertThat(nodeConfig.getDataToFetch()).isEqualTo(DataToFetch.ATTRIBUTES); assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA); var relationsQuery = new RelationsQuery(); @@ -166,7 +170,7 @@ public class TbGetRelatedAttributeNodeTest { "sourceAttr1", "targetKey1", "sourceAttr2", "targetKey2", "sourceAttr3", "targetKey3")); - config.setTelemetry(true); + config.setDataToFetch(DataToFetch.LATEST_TELEMETRY); config.setFetchTo(FetchTo.DATA); var relationsQuery = new RelationsQuery(); @@ -189,7 +193,7 @@ public class TbGetRelatedAttributeNodeTest { "sourceAttr2", "targetKey2", "sourceAttr3", "targetKey3" )); - assertThat(nodeConfig.isTelemetry()).isEqualTo(true); + assertThat(nodeConfig.getDataToFetch()).isEqualTo(DataToFetch.LATEST_TELEMETRY); assertThat(node.fetchTo).isEqualTo(FetchTo.DATA); assertThat(nodeConfig.getRelationsQuery()).isEqualTo(relationsQuery); } @@ -227,7 +231,7 @@ public class TbGetRelatedAttributeNodeTest { @Test public void givenDidNotFindEntity_whenOnMsg_thenShouldTellFailure() { // GIVEN - prepareMsgAndConfig(FetchTo.METADATA, false, DUMMY_DEVICE_ORIGINATOR); + prepareMsgAndConfig(FetchTo.METADATA, DataToFetch.ATTRIBUTES, DUMMY_DEVICE_ORIGINATOR); when(ctxMock.getTenantId()).thenReturn(TENANT_ID); @@ -263,7 +267,7 @@ public class TbGetRelatedAttributeNodeTest { var customer = new Customer(new CustomerId(UUID.randomUUID())); var user = new User(new UserId(UUID.randomUUID())); - prepareMsgAndConfig(FetchTo.DATA, false, customer.getId()); + prepareMsgAndConfig(FetchTo.DATA, DataToFetch.ATTRIBUTES, customer.getId()); entityRelation.setFrom(customer.getId()); entityRelation.setTo(user.getId()); @@ -314,7 +318,7 @@ public class TbGetRelatedAttributeNodeTest { var firstCustomer = new Customer(new CustomerId(UUID.randomUUID())); var secondCustomer = new Customer(new CustomerId(UUID.randomUUID())); - prepareMsgAndConfig(FetchTo.METADATA, false, firstCustomer.getId()); + prepareMsgAndConfig(FetchTo.METADATA, DataToFetch.ATTRIBUTES, firstCustomer.getId()); entityRelation.setFrom(firstCustomer.getId()); entityRelation.setTo(secondCustomer.getId()); @@ -365,7 +369,7 @@ public class TbGetRelatedAttributeNodeTest { var dashboard = new Dashboard(new DashboardId(UUID.randomUUID())); var entityView = new EntityView(new EntityViewId(UUID.randomUUID())); - prepareMsgAndConfig(FetchTo.DATA, true, dashboard.getId()); + prepareMsgAndConfig(FetchTo.DATA, DataToFetch.LATEST_TELEMETRY, dashboard.getId()); entityRelation.setFrom(dashboard.getId()); entityRelation.setTo(entityView.getId()); @@ -416,7 +420,7 @@ public class TbGetRelatedAttributeNodeTest { var tenant = new Tenant(new TenantId(UUID.randomUUID())); var device = new Device(new DeviceId(UUID.randomUUID())); - prepareMsgAndConfig(FetchTo.METADATA, true, tenant.getId()); + prepareMsgAndConfig(FetchTo.METADATA, DataToFetch.LATEST_TELEMETRY, tenant.getId()); entityRelation.setFrom(tenant.getId()); entityRelation.setTo(device.getId()); @@ -461,24 +465,114 @@ public class TbGetRelatedAttributeNodeTest { assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData); } - private void prepareMsgAndConfig(FetchTo fetchTo, boolean isTelemetry, EntityId originator) { - config.setAttrMapping(Map.of( - "sourceKey1", "targetKey1", - "${metaDataPattern1}", "$[messageBodyPattern1]", - "$[messageBodyPattern2]", "${metaDataPattern2}")); - config.setTelemetry(isTelemetry); - config.setFetchTo(fetchTo); + @Test + public void givenFetchFieldsToData_whenOnMsg_thenShouldFetchFieldsToData() { + // GIVEN + var device = new Device(); + device.setId(new DeviceId(UUID.randomUUID())); + device.setName("Device Name"); + var asset = new Asset(new AssetId(UUID.randomUUID())); + + prepareMsgAndConfig(FetchTo.DATA, DataToFetch.FIELDS, asset.getId()); + + entityRelation.setFrom(asset.getId()); + entityRelation.setTo(device.getId()); + entityRelation.setType(EntityRelation.CONTAINS_TYPE); + + when(ctxMock.getTenantId()).thenReturn(TENANT_ID); + + when(ctxMock.getRelationService()).thenReturn(relationServiceMock); + doReturn(Futures.immediateFuture(List.of(entityRelation))).when(relationServiceMock).findByQuery(eq(TENANT_ID), any()); + + when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); + when(ctxMock.getTenantId()).thenReturn(TENANT_ID); + when(deviceServiceMock.findDeviceById(eq(TENANT_ID), eq(device.getId()))).thenReturn(device); + + when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); + + // WHEN + node.onMsg(ctxMock, msg); + + // THEN + var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); + verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture()); + verify(ctxMock, never()).tellFailure(any(), any()); + + var expectedMsgData = "{\"temp\":42,\"humidity\":77,\"messageBodyPattern\":\"relatedEntityId\"," + + "\"relatedEntityId\":\"" + device.getId().getId() + "\",\"relatedEntityName\":\"" + device.getName() + "\"}"; + assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData); + assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); + } + + @Test + public void givenFetchFieldsToMetadata_whenOnMsg_thenShouldFetchFieldsToMetadata() { + // GIVEN + var device = new Device(); + device.setId(new DeviceId(UUID.randomUUID())); + device.setName("Device Name"); + var asset = new Asset(new AssetId(UUID.randomUUID())); + + prepareMsgAndConfig(FetchTo.METADATA, DataToFetch.FIELDS, asset.getId()); + + entityRelation.setFrom(asset.getId()); + entityRelation.setTo(device.getId()); + entityRelation.setType(EntityRelation.CONTAINS_TYPE); + + when(ctxMock.getTenantId()).thenReturn(TENANT_ID); + + when(ctxMock.getRelationService()).thenReturn(relationServiceMock); + doReturn(Futures.immediateFuture(List.of(entityRelation))).when(relationServiceMock).findByQuery(eq(TENANT_ID), any()); + + when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); + when(ctxMock.getTenantId()).thenReturn(TENANT_ID); + when(deviceServiceMock.findDeviceById(eq(TENANT_ID), eq(device.getId()))).thenReturn(device); + + when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); + + // WHEN + node.onMsg(ctxMock, msg); + + // THEN + var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); + verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture()); + verify(ctxMock, never()).tellFailure(any(), any()); + + var expectedMsgMetadata = new TbMsgMetaData(Map.of( + "metaDataPattern", "relatedEntityName", + "relatedEntityId", device.getId().getId().toString(), + "relatedEntityName", device.getName() + )); + + assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); + assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetadata); + } + + private void prepareMsgAndConfig(FetchTo fetchTo, DataToFetch dataToFetch, EntityId originator) { + + config.setDataToFetch(dataToFetch); + config.setFetchTo(fetchTo); node.config = config; node.fetchTo = fetchTo; - var msgMetaData = new TbMsgMetaData(); - msgMetaData.putValue("metaDataPattern1", "sourceKey2"); - msgMetaData.putValue("metaDataPattern2", "targetKey3"); - - var msgData = "{\"temp\":42,\"humidity\":77,\"messageBodyPattern1\":\"targetKey2\",\"messageBodyPattern2\":\"sourceKey3\"}"; + String msgData; + if (dataToFetch.equals(DataToFetch.FIELDS)) { + config.setAttrMapping(Map.of( + "id", "$[messageBodyPattern]", + "name", "${metaDataPattern}")); + msgMetaData.putValue("metaDataPattern", "relatedEntityName"); + msgData = "{\"temp\":42,\"humidity\":77,\"messageBodyPattern\":\"relatedEntityId\"}"; + } else { + config.setAttrMapping(Map.of( + "sourceKey1", "targetKey1", + "${metaDataPattern1}", "$[messageBodyPattern1]", + "$[messageBodyPattern2]", "${metaDataPattern2}")); + msgMetaData.putValue("metaDataPattern1", "sourceKey2"); + msgMetaData.putValue("metaDataPattern2", "targetKey3"); + msgData = "{\"temp\":42,\"humidity\":77,\"messageBodyPattern1\":\"targetKey2\",\"messageBodyPattern2\":\"sourceKey3\"}"; + } - msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", originator, msgMetaData, msgData); + msg = TbMsg.newMsg(SessionMsgType.POST_TELEMETRY_REQUEST.name(), originator, msgMetaData, msgData); } @RequiredArgsConstructor diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java index 0f1863255c..19641b7000 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java @@ -115,6 +115,34 @@ public class TbGetTenantAttributeNodeTest { verify(ctxMock, never()).tellSuccess(any()); } + @Test + public void givenConfigWithNullDataToFetch_whenInit_thenException() { + // GIVEN + config.setDataToFetch(null); + nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + + // WHEN + var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); + + // THEN + assertThat(exception.getMessage()).isEqualTo("DataToFetch property has invalid value: null. Only ATTRIBUTES and LATEST_TELEMETRY values supported!"); + verify(ctxMock, never()).tellSuccess(any()); + } + + @Test + public void givenConfigWithUnsupportedDataToFetch_whenInit_thenException() { + // GIVEN + config.setDataToFetch(DataToFetch.FIELDS); + nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + + // WHEN + var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); + + // THEN + assertThat(exception.getMessage()).isEqualTo("DataToFetch property has invalid value: FIELDS. Only ATTRIBUTES and LATEST_TELEMETRY values supported!"); + verify(ctxMock, never()).tellSuccess(any()); + } + @Test public void givenDefaultConfig_whenInit_thenOK() throws TbNodeException { // GIVEN @@ -125,7 +153,7 @@ public class TbGetTenantAttributeNodeTest { // THEN assertThat(node.config).isEqualTo(config); assertThat(config.getAttrMapping()).isEqualTo(Map.of("alarmThreshold", "threshold")); - assertThat(config.isTelemetry()).isEqualTo(false); + assertThat(config.getDataToFetch()).isEqualTo(DataToFetch.ATTRIBUTES); assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA); } @@ -136,7 +164,7 @@ public class TbGetTenantAttributeNodeTest { "sourceAttr1", "targetKey1", "sourceAttr2", "targetKey2", "sourceAttr3", "targetKey3")); - config.setTelemetry(true); + config.setDataToFetch(DataToFetch.LATEST_TELEMETRY); config.setFetchTo(FetchTo.DATA); nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); @@ -149,7 +177,7 @@ public class TbGetTenantAttributeNodeTest { "sourceAttr1", "targetKey1", "sourceAttr2", "targetKey2", "sourceAttr3", "targetKey3")); - assertThat(config.isTelemetry()).isEqualTo(true); + assertThat(config.getDataToFetch()).isEqualTo(DataToFetch.LATEST_TELEMETRY); assertThat(node.fetchTo).isEqualTo(FetchTo.DATA); } @@ -188,7 +216,7 @@ public class TbGetTenantAttributeNodeTest { // GIVEN var deviceId = new DeviceId(UUID.randomUUID()); - prepareMsgAndConfig(FetchTo.DATA, false, deviceId); + prepareMsgAndConfig(FetchTo.DATA, DataToFetch.ATTRIBUTES, deviceId); List attributesList = List.of( new BaseAttributeKvEntry(new StringDataEntry("sourceKey1", "sourceValue1"), 1L), @@ -229,7 +257,7 @@ public class TbGetTenantAttributeNodeTest { @Test public void givenFetchAttributesToMetaData_whenOnMsg_thenShouldFetchAttributesToMetaData() { // GIVEN - prepareMsgAndConfig(FetchTo.METADATA, false, TENANT_ID); + prepareMsgAndConfig(FetchTo.METADATA, DataToFetch.ATTRIBUTES, TENANT_ID); List attributesList = List.of( new BaseAttributeKvEntry(new StringDataEntry("sourceKey1", "sourceValue1"), 1L), @@ -272,7 +300,7 @@ public class TbGetTenantAttributeNodeTest { // GIVEN var customerId = new CustomerId(UUID.randomUUID()); - prepareMsgAndConfig(FetchTo.DATA, true, customerId); + prepareMsgAndConfig(FetchTo.DATA, DataToFetch.LATEST_TELEMETRY, customerId); List timeseries = List.of( new BasicTsKvEntry(1L, new StringDataEntry("sourceKey1", "sourceValue1")), @@ -315,7 +343,7 @@ public class TbGetTenantAttributeNodeTest { // GIVEN var ruleChainId = new RuleChainId(UUID.randomUUID()); - prepareMsgAndConfig(FetchTo.METADATA, true, ruleChainId); + prepareMsgAndConfig(FetchTo.METADATA, DataToFetch.LATEST_TELEMETRY, ruleChainId); List timeseries = List.of( new BasicTsKvEntry(1L, new StringDataEntry("sourceKey1", "sourceValue1")), @@ -353,12 +381,12 @@ public class TbGetTenantAttributeNodeTest { assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData); } - private void prepareMsgAndConfig(FetchTo fetchTo, boolean isTelemetry, EntityId originator) { + private void prepareMsgAndConfig(FetchTo fetchTo, DataToFetch dataToFetch, EntityId originator) { config.setAttrMapping(Map.of( "sourceKey1", "targetKey1", "${metaDataPattern1}", "$[messageBodyPattern1]", "$[messageBodyPattern2]", "${metaDataPattern2}")); - config.setTelemetry(isTelemetry); + config.setDataToFetch(dataToFetch); config.setFetchTo(fetchTo); node.config = config; diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeTest.java index 423c770e27..e5333cb459 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeTest.java @@ -26,7 +26,7 @@ import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; -import org.thingsboard.rule.engine.util.EntityDetails; +import org.thingsboard.rule.engine.util.ContactBasedEntityDetails; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; @@ -83,7 +83,7 @@ public class TbGetTenantDetailsNodeTest { @Test public void givenConfigWithNullFetchTo_whenInit_thenException() { // GIVEN - config.setDetailsList(List.of(EntityDetails.ID)); + config.setDetailsList(List.of(ContactBasedEntityDetails.ID)); config.setFetchTo(null); nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); @@ -105,7 +105,7 @@ public class TbGetTenantDetailsNodeTest { @Test public void givenCustomConfig_whenInit_thenOK() throws TbNodeException { // GIVEN - config.setDetailsList(List.of(EntityDetails.ID, EntityDetails.PHONE)); + config.setDetailsList(List.of(ContactBasedEntityDetails.ID, ContactBasedEntityDetails.PHONE)); config.setFetchTo(FetchTo.METADATA); nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); @@ -114,7 +114,7 @@ public class TbGetTenantDetailsNodeTest { // THEN assertThat(node.config).isEqualTo(config); - assertThat(config.getDetailsList()).isEqualTo(List.of(EntityDetails.ID, EntityDetails.PHONE)); + assertThat(config.getDetailsList()).isEqualTo(List.of(ContactBasedEntityDetails.ID, ContactBasedEntityDetails.PHONE)); assertThat(config.getFetchTo()).isEqualTo(FetchTo.METADATA); assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA); } @@ -136,7 +136,7 @@ public class TbGetTenantDetailsNodeTest { @Test public void givenAllEntityDetailsAndFetchToData_whenOnMsg_thenShouldTellSuccessAndFetchAllToData() { // GIVEN - prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.values())); + prepareMsgAndConfig(FetchTo.DATA, List.of(ContactBasedEntityDetails.values())); mockFindTenant(); @@ -169,7 +169,7 @@ public class TbGetTenantDetailsNodeTest { @Test public void givenSomeEntityDetailsAndFetchToMetadata_whenOnMsg_thenShouldTellSuccessAndFetchSomeToMetaData() { // GIVEN - prepareMsgAndConfig(FetchTo.METADATA, List.of(EntityDetails.ID, EntityDetails.TITLE, EntityDetails.PHONE)); + prepareMsgAndConfig(FetchTo.METADATA, List.of(ContactBasedEntityDetails.ID, ContactBasedEntityDetails.TITLE, ContactBasedEntityDetails.PHONE)); mockFindTenant(); @@ -198,7 +198,7 @@ public class TbGetTenantDetailsNodeTest { tenant.setAddress(null); tenant.setAddress2(null); - prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2)); + prepareMsgAndConfig(FetchTo.DATA, List.of(ContactBasedEntityDetails.ZIP, ContactBasedEntityDetails.ADDRESS, ContactBasedEntityDetails.ADDRESS2)); mockFindTenant(); @@ -218,7 +218,7 @@ public class TbGetTenantDetailsNodeTest { @Test public void givenDidNotFindTenant_whenOnMsg_thenShouldTellSuccessAndFetchNothingToData() { // GIVEN - prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2)); + prepareMsgAndConfig(FetchTo.DATA, List.of(ContactBasedEntityDetails.ZIP, ContactBasedEntityDetails.ADDRESS, ContactBasedEntityDetails.ADDRESS2)); when(ctxMock.getTenantId()).thenReturn(tenant.getId()); when(ctxMock.getTenantService()).thenReturn(tenantServiceMock); @@ -242,7 +242,7 @@ public class TbGetTenantDetailsNodeTest { // GIVEN tenant.setAdditionalInfo(JacksonUtil.toJsonNode("{\"someProperty\":\"someValue\",\"description\":null}")); - prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ADDITIONAL_INFO)); + prepareMsgAndConfig(FetchTo.DATA, List.of(ContactBasedEntityDetails.ADDITIONAL_INFO)); mockFindTenant(); @@ -259,7 +259,7 @@ public class TbGetTenantDetailsNodeTest { assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); } - private void prepareMsgAndConfig(FetchTo fetchTo, List detailsList) { + private void prepareMsgAndConfig(FetchTo fetchTo, List detailsList) { config.setDetailsList(detailsList); config.setFetchTo(fetchTo);