Browse Source

added upgrade script logic for all enrichment rule nodes && additional improvements to TbGetRelatedAttributeNode

pull/8661/head
ShvaykaD 3 years ago
parent
commit
5831876b8e
  1. 1
      application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java
  2. 139
      application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java
  3. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java
  4. 4
      dao/pom.xml
  5. 28
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java
  6. 19
      dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java
  7. 9
      dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
  8. 11
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java
  9. 2
      dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java
  10. 43
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/VersionedNode.java
  11. 6
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/VersionedNodeConfiguration.java
  12. 27
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNode.java
  13. 22
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/DataToFetch.java
  14. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java
  15. 101
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java
  16. 53
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java
  17. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java
  18. 38
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java
  19. 22
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java
  20. 22
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java
  21. 27
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java
  22. 22
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java
  23. 22
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java
  24. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java
  25. 22
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java
  26. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java
  27. 28
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java
  28. 25
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java
  29. 22
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java
  30. 41
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/ContactBasedEntityDetails.java
  31. 123
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/CalculateDeltaNodeTest.java
  32. 46
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java
  33. 26
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeTest.java
  34. 172
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java
  35. 46
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java
  36. 20
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeTest.java

1
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);

139
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<RuleChainId, TenantId>();
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<RuleChainId, TenantId> 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<String, DeviceProfileEntity> deviceProfileEntityDynamicConditionsUpdater =

2
common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java

@ -42,6 +42,8 @@ public interface TimeseriesService {
ListenableFuture<List<TsKvEntry>> findLatest(TenantId tenantId, EntityId entityId, Collection<String> keys);
List<TsKvEntry> findLatestSync(TenantId tenantId, EntityId entityId, Collection<String> keys);
ListenableFuture<List<TsKvEntry>> findAllLatest(TenantId tenantId, EntityId entityId);
ListenableFuture<Integer> save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry);

4
dao/pom.xml

@ -222,6 +222,10 @@
<groupId>com.jayway.jsonpath</groupId>
<artifactId>json-path</artifactId>
</dependency>
<dependency>
<groupId>org.thingsboard.rule-engine</groupId>
<artifactId>rule-engine-api</artifactId>
</dependency>
</dependencies>
<build>
<plugins>

28
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<Boolean, JsonNode> 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<Boolean, JsonNode> 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.");

19
dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java

@ -156,11 +156,12 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme
@Override
public ListenableFuture<TsKvEntry> 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;
}
}

9
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<TsKvEntry> findLatestSync(TenantId tenantId, EntityId entityId, Collection<String> keys) {
validate(entityId);
List<TsKvEntry> 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<List<TsKvEntry>> findAllLatest(TenantId tenantId, EntityId entityId) {
validate(entityId);

11
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 <T> ListenableFuture<T> findLatest(TenantId tenantId, EntityId entityId, String key, java.util.function.Function<TbResultSet, T> function) {
BoundStatementBuilder stmtBuilder = new BoundStatementBuilder(getFindLatestStmt().bind());
stmtBuilder.setString(0, entityId.getEntityType().name());

2
dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesLatestDao.java

@ -40,6 +40,8 @@ public interface TimeseriesLatestDao {
*/
ListenableFuture<TsKvEntry> findLatest(TenantId tenantId, EntityId entityId, String key);
TsKvEntry findLatestSync(TenantId tenantId, EntityId entityId, String key);
ListenableFuture<List<TsKvEntry>> findAllLatest(TenantId tenantId, EntityId entityId);
ListenableFuture<Void> saveLatest(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry);

43
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<Boolean, JsonNode> 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();
}
}

6
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntityDetails.java → 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();
}

27
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<ValueWithTs> fetchLatestValue(EntityId entityId) {
private ListenableFuture<ValueWithTs> 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<TsKvEntry> tsKvEntries = timeseriesService.findLatestSync(
ctx.getTenantId(),
entityId,
Collections.singletonList(config.getInputValueKey()));
return extractValue(tsKvEntries.get(0));
}
private ListenableFuture<ValueWithTs> 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);
}
}

22
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
}

4
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;
}

101
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<T extends EntityId> extends TbAbstractNodeWithFetchTo<TbGetEntityAttrNodeConfiguration> {
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<T extends EntityId> 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<String, String>();
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<Map<String, String>> collectMappedEntityFieldsAsync(TbContext ctx, EntityId entityId, HashMap<String, String> mappingsMap) {
return Futures.transform(EntitiesFieldsAsyncLoader.findAsync(ctx, entityId),
fieldsData -> {
var targetKeysToSourceValuesMap = new HashMap<String, String>();
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<List<KvEntry>> getAttributesAsync(TbContext ctx, EntityId entityId, List<String> attrKeys) {
@ -95,4 +149,27 @@ public abstract class TbAbstractGetEntityAttrNode<T extends EntityId> extends Tb
ctx.tellSuccess(transformMessage(msg, msgData, msgMetaData));
}
protected TbPair<Boolean, JsonNode> 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);
}
}

53
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<C extends TbAbstractGetEnti
protected abstract ListenableFuture<? extends ContactBased<I>> getContactBasedFuture(TbContext ctx, TbMsg msg);
protected void checkIfDetailsListIsNotEmptyOrElseThrow(List<EntityDetails> detailsList) throws TbNodeException {
protected void checkIfDetailsListIsNotEmptyOrElseThrow(List<ContactBasedEntityDetails> detailsList) throws TbNodeException {
if (detailsList == null || detailsList.isEmpty()) {
throw new TbNodeException("No entity details selected!");
}
@ -60,87 +60,64 @@ public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEnti
return Futures.immediateFuture(msg);
}
var msgMetaData = msg.getMetaData().copy();
setProperties(contactBased, messageData, msgMetaData);
fetchEntityDetailsToMsg(contactBased, messageData, msgMetaData);
return Futures.immediateFuture(transformMessage(msg, messageData, msgMetaData));
}, MoreExecutors.directExecutor());
}
private void setProperties(ContactBased<I> contactBased, ObjectNode messageData, TbMsgMetaData msgMetaData) {
String prefix = getPrefix();
String property;
String value;
for (var entityDetails : config.getDetailsList()) {
switch (entityDetails) {
private void fetchEntityDetailsToMsg(ContactBased<I> 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);
}
}

4
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<EntityDetails> detailsList;
private List<ContactBasedEntityDetails> detailsList;
}

38
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<C extends TbAbstractFetchToNodeConfiguration> implements TbNode {
public abstract class TbAbstractNodeWithFetchTo<C extends TbAbstractFetchToNodeConfiguration> 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<C extends TbAbstractFetchToNodeC
}
}
protected TbPair<Boolean, JsonNode> 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!");
}
}
}

22
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<TbFe
ctx.tellSuccess(transformedMsg);
}
@Override
public TbPair<Boolean, JsonNode> 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);
}
}

22
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<TbGetAttrib
return Futures.immediateFuture(msg.getOriginator());
}
@Override
public TbPair<Boolean, JsonNode> 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);
}
}

27
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<Cust
protected TbGetEntityAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
var config = TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class);
checkIfMappingIsNotEmptyOrElseThrow(config.getAttrMapping());
checkDataToFetchSupportedOrElseThrow(config.getDataToFetch());
return config;
}
@Override
protected void checkDataToFetchSupportedOrElseThrow(DataToFetch dataToFetch) throws TbNodeException {
if (dataToFetch == null || dataToFetch.equals(DataToFetch.FIELDS)) {
throw new TbNodeException("DataToFetch property has invalid value: " + dataToFetch +
". Only ATTRIBUTES and LATEST_TELEMETRY values supported!");
}
}
@Override
protected ListenableFuture<CustomerId> findEntityAsync(TbContext ctx, EntityId originator) {
return Futures.transformAsync(EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, originator),
@ -57,4 +71,17 @@ public class TbGetCustomerAttributeNode extends TbAbstractGetEntityAttrNode<Cust
);
}
@Override
public TbPair<Boolean, JsonNode> 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);
}
}

22
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<TbG
}
}
@Override
public TbPair<Boolean, JsonNode> 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);
}
}

22
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<TbGetDevice
ctx.getDbCallbackExecutor());
}
@Override
public TbPair<Boolean, JsonNode> 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);
}
}

4
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<TbGetEntityAttrNodeConfiguration> {
private Map<String, String> attrMapping;
private boolean isTelemetry;
private DataToFetch dataToFetch;
@Override
public TbGetEntityAttrNodeConfiguration defaultConfiguration() {
@ -35,7 +35,7 @@ public class TbGetEntityAttrNodeConfiguration extends TbAbstractFetchToNodeConfi
var attrMapping = new HashMap<String, String>();
attrMapping.putIfAbsent("alarmThreshold", "threshold");
configuration.setAttrMapping(attrMapping);
configuration.setTelemetry(false);
configuration.setDataToFetch(DataToFetch.ATTRIBUTES);
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}

22
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<TbGetOr
);
}
@Override
public TbPair<Boolean, JsonNode> 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);
}
}

2
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<String, String>();
attrMapping.putIfAbsent("serialNumber", "sn");
configuration.setAttrMapping(attrMapping);
configuration.setTelemetry(false);
configuration.setDataToFetch(DataToFetch.ATTRIBUTES);
configuration.setFetchTo(FetchTo.METADATA);
var relationsQuery = new RelationsQuery();

28
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<Entit
public TbGetRelatedAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
var config = TbNodeUtils.convert(configuration, TbGetRelatedAttrNodeConfiguration.class);
checkIfMappingIsNotEmptyOrElseThrow(config.getAttrMapping());
checkDataToFetchSupportedOrElseThrow(config.getDataToFetch());
return config;
}
@ -59,4 +67,24 @@ public class TbGetRelatedAttributeNode extends TbAbstractGetEntityAttrNode<Entit
ctx.getDbCallbackExecutor());
}
@Override
protected void checkDataToFetchSupportedOrElseThrow(DataToFetch dataToFetch) throws TbNodeException {
if (dataToFetch == null) {
throw new TbNodeException("DataToFetch property cannot be null! Supported values are: " + Arrays.toString(DataToFetch.values()));
}
}
@Override
public TbPair<Boolean, JsonNode> 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);
}
}

25
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<Tenant
public TbGetEntityAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
var config = TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class);
checkIfMappingIsNotEmptyOrElseThrow(config.getAttrMapping());
checkDataToFetchSupportedOrElseThrow(config.getDataToFetch());
return config;
}
@ -53,4 +57,25 @@ public class TbGetTenantAttributeNode extends TbAbstractGetEntityAttrNode<Tenant
return Futures.immediateFuture(ctx.getTenantId());
}
@Override
protected void checkDataToFetchSupportedOrElseThrow(DataToFetch dataToFetch) throws TbNodeException {
if (dataToFetch == null || dataToFetch.equals(DataToFetch.FIELDS)) {
throw new TbNodeException("DataToFetch property has invalid value: " + dataToFetch +
". Only ATTRIBUTES and LATEST_TELEMETRY values supported!");
}
}
@Override
public TbPair<Boolean, JsonNode> 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);
}
}

22
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<TbGet
return ctx.getTenantService().findTenantByIdAsync(ctx.getTenantId(), ctx.getTenantId());
}
@Override
public TbPair<Boolean, JsonNode> 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);
}
}

41
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;
}
}

123
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(

46
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<AttributeKvEntry> 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<AttributeKvEntry> 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<TsKvEntry> 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<TsKvEntry> 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;

26
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<EntityDetails> detailsList, EntityId originator) {
private void prepareMsgAndConfig(FetchTo fetchTo, List<ContactBasedEntityDetails> detailsList, EntityId originator) {
config.setDetailsList(detailsList);
config.setFetchTo(fetchTo);

172
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

46
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<AttributeKvEntry> 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<AttributeKvEntry> 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<TsKvEntry> 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<TsKvEntry> 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;

20
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<EntityDetails> detailsList) {
private void prepareMsgAndConfig(FetchTo fetchTo, List<ContactBasedEntityDetails> detailsList) {
config.setDetailsList(detailsList);
config.setFetchTo(fetchTo);

Loading…
Cancel
Save