Browse Source

Merge pull request #15 from dskarzh/feature/enrichment-rule-nodes-improvements

Feature/enrichment rule nodes improvements
pull/8661/head
Shvaika Dmytro 4 years ago
committed by GitHub
parent
commit
dbf454f0ff
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 1
      application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java
  2. 98
      application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java
  3. 2
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeConfiguration.java
  4. 8
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbNode.java
  5. 21
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchTo.java
  6. 23
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java
  7. 75
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java
  8. 110
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java
  9. 30
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java
  10. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java
  11. 50
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java
  12. 109
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java
  13. 40
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java
  14. 11
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeConfiguration.java
  15. 15
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java
  16. 10
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java
  17. 17
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java
  18. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java
  19. 7
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeConfiguration.java
  20. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java
  21. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java
  22. 13
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java
  23. 10
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsConfiguration.java
  24. 80
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java
  25. 14
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java
  26. 25
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java
  27. 20
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java
  28. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java
  29. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeConfiguration.java
  30. 6
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java
  31. 34
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoader.java
  32. 20
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractAttributeNodeTest.java
  33. 58
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNodeTest.java
  34. 6
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeTest.java
  35. 73
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java
  36. 376
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNodeTest.java
  37. 75
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java
  38. 77
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java
  39. 246
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoaderTest.java
  40. 7
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java

1
application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java

@ -247,6 +247,7 @@ public class ThingsboardInstallService {
case "3.4.4":
log.info("Upgrading ThingsBoard from version 3.4.4 to 3.5.0 ...");
databaseEntitiesUpgradeService.upgradeDatabase("3.4.4");
dataUpdateService.updateData("3.4.4");
log.info("Updating system data...");
systemDataLoaderService.updateSystemWidgets();
if (!getEnv("SKIP_DEFAULT_NOTIFICATION_CONFIGS_CREATION", false)) {

98
application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java

@ -28,6 +28,16 @@ import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.flow.TbRuleChainInputNode;
import org.thingsboard.rule.engine.flow.TbRuleChainInputNodeConfiguration;
import org.thingsboard.rule.engine.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;
@ -46,6 +56,7 @@ import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageDataIterable;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.common.data.query.DynamicValue;
@ -83,6 +94,7 @@ import org.thingsboard.server.service.install.TbRuleEngineQueueConfigService;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.atomic.AtomicLong;
@ -203,11 +215,97 @@ public class DefaultDataUpdateService implements DataUpdateService {
log.info("Skipping edge events migration");
}
break;
case "3.4.4":
log.info("Updating data from version 3.4.4 to 3.5.0 ...");
log.info("Started enrichment rule nodes update ...");
updateEnrichmentRuleNodes();
log.info("Finished enrichment rule nodes update ...");
break;
default:
throw new RuntimeException("Unable to update data, unsupported fromVersion: " + fromVersion);
}
}
private void updateEnrichmentRuleNodes() {
try {
var ruleNodeTypesToUpdate = List.of(
TbGetOriginatorFieldsNode.class.getName(),
TbFetchDeviceCredentialsNode.class.getName(),
TbGetAttributesNode.class.getName(),
TbGetDeviceAttrNode.class.getName(),
TbGetRelatedAttributeNode.class.getName(),
TbGetTenantAttributeNode.class.getName(),
TbGetCustomerAttributeNode.class.getName(),
TbGetCustomerDetailsNode.class.getName(),
TbGetTenantDetailsNode.class.getName()
);
var ruleChainIdToTenantId = new HashMap<RuleChainId, TenantId>();
ruleNodeTypesToUpdate.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("Unable to update [{}] rule node with ID [{}]! Node configuration is null! Skipping this node!",
ruleNodeType, ruleNode.getId());
continue;
}
if (!configuration.isObject()) {
log.error("Unable to update [{}] rule node with ID [{}]! Node configuration is not an object! Skipping this node!",
ruleNodeType, ruleNode.getId());
continue;
}
var configObjectNode = (ObjectNode) configuration;
var fetchTo = FetchTo.METADATA;
if (configObjectNode.has("fetchToMetadata")) {
var fetchToMetadata = configObjectNode.get("fetchToMetadata").asText();
if ("true".equals(fetchToMetadata)) {
fetchTo = FetchTo.METADATA;
} else if ("false".equals(fetchToMetadata)) {
fetchTo = FetchTo.DATA;
} else {
log.error("[fetchToMetadata] property has unexpected value: {}! Expected true or false! Skipping this node ID[{}]!",
fetchToMetadata, ruleNode.getId());
}
configObjectNode.remove("fetchToMetadata");
}
if (configObjectNode.has("fetchToData")) {
var fetchToData = configObjectNode.get("fetchToData").asText();
if ("true".equals(fetchToData)) {
fetchTo = FetchTo.DATA;
} else if ("false".equals(fetchToData)) {
fetchTo = FetchTo.METADATA;
} else {
log.error("[fetchToData] property has unexpected value: {}! Expected true or false! Skipping this node ID[{}]!",
fetchToData, ruleNode.getId());
}
configObjectNode.remove("fetchToData");
}
if (configObjectNode.has("addToMetadata")) {
var addToMetadata = configObjectNode.get("addToMetadata").asText();
if ("true".equals(addToMetadata)) {
fetchTo = FetchTo.METADATA;
} else if ("false".equals(addToMetadata)) {
fetchTo = FetchTo.DATA;
} else {
log.error("[addToMetadata] property has unexpected value: {}! Skipping Expected true or false! Skipping this node ID[{}]!",
addToMetadata, ruleNode.getId());
}
configObjectNode.remove("addToMetadata");
}
configObjectNode.put("fetchTo", fetchTo.toString());
ruleNode.setConfiguration(configObjectNode);
ruleChainIdToTenantId.computeIfAbsent(ruleNode.getRuleChainId(),
ruleChainId -> ruleChainService.findRuleChainById(TenantId.SYS_TENANT_ID, ruleNode.getRuleChainId()).getTenantId());
ruleChainService.saveRuleNode(ruleChainIdToTenantId.get(ruleNode.getRuleChainId()), ruleNode);
}
});
} catch (Exception e) {
log.error("Unexpected error during enrichment rule nodes updating!", e);
}
}
private final PaginatedUpdater<String, DeviceProfileEntity> deviceProfileEntityDynamicConditionsUpdater =
new PaginatedUpdater<>() {

2
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeConfiguration.java

@ -16,7 +16,5 @@
package org.thingsboard.rule.engine.api;
public interface NodeConfiguration<T extends NodeConfiguration> {
T defaultConfiguration();
}

8
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbNode.java

@ -24,13 +24,13 @@ import java.util.concurrent.ExecutionException;
* Created by ashvayka on 19.01.18.
*/
public interface TbNode {
void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException;
void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException, TbNodeException;
default void destroy() {}
default void onPartitionChangeMsg(TbContext ctx, PartitionChangeMsg msg) {}
default void destroy() {
}
default void onPartitionChangeMsg(TbContext ctx, PartitionChangeMsg msg) {
}
}

21
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchTo.java

@ -0,0 +1,21 @@
/**
* 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 FetchTo {
DATA,
METADATA
}

23
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java

@ -0,0 +1,23 @@
/**
* 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;
import lombok.Data;
@Data
public abstract class TbAbstractFetchToNodeConfiguration {
private FetchTo fetchTo;
}

75
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java

@ -15,17 +15,13 @@
*/
package org.thingsboard.rule.engine.metadata;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
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 org.apache.commons.collections.CollectionUtils;
import org.apache.commons.lang3.BooleanUtils;
import org.thingsboard.common.util.JacksonUtil;
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.util.TbNodeUtils;
@ -36,43 +32,35 @@ import org.thingsboard.server.common.data.kv.JsonDataEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.NoSuchElementException;
import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE;
import static org.thingsboard.server.common.data.DataConstants.CLIENT_SCOPE;
import static org.thingsboard.server.common.data.DataConstants.LATEST_TS;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE;
public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeConfiguration, T extends EntityId> implements TbNode {
public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeConfiguration, T extends EntityId> extends TbAbstractNodeWithFetchTo<C> {
private static final String VALUE = "value";
private static final String TS = "ts";
protected C config;
private boolean fetchToData;
private boolean isTellFailureIfAbsent;
private boolean getLatestValueWithTs;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = loadGetAttributesNodeConfig(configuration);
this.fetchToData = config.isFetchToData();
this.getLatestValueWithTs = config.isGetLatestValueWithTs();
this.isTellFailureIfAbsent = BooleanUtils.toBooleanDefaultIfNull(this.config.isTellFailureIfAbsent(), true);
super.init(ctx, configuration);
getLatestValueWithTs = config.isGetLatestValueWithTs();
isTellFailureIfAbsent = config.isTellFailureIfAbsent();
}
protected abstract C loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException;
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws TbNodeException {
try {
@ -89,20 +77,16 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
private void safePutAttributes(TbContext ctx, TbMsg msg, T entityId) {
if (entityId == null || entityId.isNullUid()) {
ctx.tellNext(msg, FAILURE);
ctx.tellFailure(msg, new NoSuchElementException("Did not find entity! Msg ID: " + msg.getId()));
return;
}
JsonNode msgDataNode;
if (fetchToData) {
msgDataNode = JacksonUtil.toJsonNode(msg.getData());
if (!msgDataNode.isObject()) {
ctx.tellFailure(msg, new IllegalArgumentException("Msg body is not an object!"));
return;
}
ObjectNode msgDataNode;
if (FetchTo.DATA.equals(fetchTo)) {
msgDataNode = getMsgDataAsObjectNode(msg);
} else {
msgDataNode = null;
}
ConcurrentHashMap<String, List<String>> failuresMap = new ConcurrentHashMap<>();
var failuresMap = new ConcurrentHashMap<String, List<String>>();
ListenableFuture<List<Map<String, ? extends List<? extends KvEntry>>>> allFutures = Futures.allAsList(
getLatestTelemetry(ctx, entityId, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresMap),
getAttrAsync(ctx, entityId, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg), failuresMap),
@ -110,23 +94,28 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
getAttrAsync(ctx, entityId, SERVER_SCOPE, TbNodeUtils.processPatterns(config.getServerAttributeNames(), msg), failuresMap)
);
withCallback(allFutures, futuresList -> {
TbMsgMetaData msgMetaData = msg.getMetaData().copy();
var msgMetaData = msg.getMetaData().copy();
futuresList.stream().filter(Objects::nonNull).forEach(kvEntriesMap -> {
kvEntriesMap.forEach((keyScope, kvEntryList) -> {
String prefix = getPrefix(keyScope);
var prefix = getPrefix(keyScope);
kvEntryList.forEach(kvEntry -> {
String key = prefix + kvEntry.getKey();
if (fetchToData) {
JacksonUtil.addKvEntry((ObjectNode) msgDataNode, kvEntry, key);
} else {
var key = prefix + kvEntry.getKey();
if (FetchTo.DATA.equals(fetchTo)) {
JacksonUtil.addKvEntry(msgDataNode, kvEntry, key);
} else if (FetchTo.METADATA.equals(fetchTo)) {
msgMetaData.putValue(key, kvEntry.getValueAsString());
}
});
});
});
TbMsg outMsg = fetchToData ?
TbMsg.transformMsgData(msg, JacksonUtil.toString(msgDataNode)) :
TbMsg.transformMsg(msg, msgMetaData);
TbMsg outMsg = null;
if (FetchTo.DATA.equals(fetchTo)) {
outMsg = TbMsg.transformMsgData(msg, JacksonUtil.toString(msgDataNode));
} else if (FetchTo.METADATA.equals(fetchTo)) {
outMsg = TbMsg.transformMsg(msg, msgMetaData);
}
if (failuresMap.isEmpty()) {
ctx.tellSuccess(outMsg);
} else {
@ -139,12 +128,12 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
if (CollectionUtils.isEmpty(keys)) {
return Futures.immediateFuture(null);
}
ListenableFuture<List<AttributeKvEntry>> attributeKvEntryListFuture = ctx.getAttributesService().find(ctx.getTenantId(), entityId, scope, keys);
var attributeKvEntryListFuture = ctx.getAttributesService().find(ctx.getTenantId(), entityId, scope, keys);
return Futures.transform(attributeKvEntryListFuture, attributeKvEntryList -> {
if (isTellFailureIfAbsent && attributeKvEntryList.size() != keys.size()) {
getNotExistingKeys(attributeKvEntryList, keys).forEach(key -> computeFailuresMap(scope, failuresMap, key));
}
Map<String, List<AttributeKvEntry>> mapAttributeKvEntry = new HashMap<>();
var mapAttributeKvEntry = new HashMap<String, List<AttributeKvEntry>>();
mapAttributeKvEntry.put(scope, attributeKvEntryList);
return mapAttributeKvEntry;
}, MoreExecutors.directExecutor());
@ -156,7 +145,7 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
}
ListenableFuture<List<TsKvEntry>> latestTelemetryFutures = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, keys);
return Futures.transform(latestTelemetryFutures, tsKvEntries -> {
List<TsKvEntry> listTsKvEntry = new ArrayList<>();
var listTsKvEntry = new ArrayList<TsKvEntry>();
tsKvEntries.forEach(tsKvEntry -> {
if (tsKvEntry.getValue() == null) {
if (isTellFailureIfAbsent) {
@ -168,22 +157,22 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
listTsKvEntry.add(new BasicTsKvEntry(tsKvEntry.getTs(), tsKvEntry));
}
});
Map<String, List<TsKvEntry>> mapTsKvEntry = new HashMap<>();
var mapTsKvEntry = new HashMap<String, List<TsKvEntry>>();
mapTsKvEntry.put(LATEST_TS, listTsKvEntry);
return mapTsKvEntry;
}, MoreExecutors.directExecutor());
}
private TsKvEntry getValueWithTs(TsKvEntry tsKvEntry) {
ObjectMapper mapper = fetchToData ? JacksonUtil.OBJECT_MAPPER : JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER;
ObjectNode value = JacksonUtil.newObjectNode(mapper);
var mapper = FetchTo.DATA.equals(fetchTo) ? JacksonUtil.OBJECT_MAPPER : JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER;
var value = JacksonUtil.newObjectNode(mapper);
value.put(TS, tsKvEntry.getTs());
JacksonUtil.addKvEntry(value, tsKvEntry, VALUE, mapper);
return new BasicTsKvEntry(tsKvEntry.getTs(), new JsonDataEntry(tsKvEntry.getKey(), value.toString()));
}
private String getPrefix(String scope) {
String prefix = "";
var prefix = "";
switch (scope) {
case CLIENT_SCOPE:
prefix = "cs_";
@ -209,7 +198,7 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
}
private RuntimeException reportFailures(ConcurrentHashMap<String, List<String>> failuresMap) {
StringBuilder errorMessage = new StringBuilder("The following attribute/telemetry keys is not present in the DB: ").append("\n");
var errorMessage = new StringBuilder("The following attribute/telemetry keys is not present in the DB: ").append("\n");
if (failuresMap.containsKey(CLIENT_SCOPE)) {
errorMessage.append("\t").append("[" + CLIENT_SCOPE + "]:").append(failuresMap.get(CLIENT_SCOPE).toString()).append("\n");
}

110
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java

@ -0,0 +1,110 @@
/**
* 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;
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.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.msg.TbMsg;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.NoSuchElementException;
import java.util.stream.Collectors;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@Slf4j
public abstract class TbAbstractGetEntityAttrNode<T extends EntityId> extends TbAbstractNodeWithFetchTo<TbGetEntityAttrNodeConfiguration> {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
ObjectNode msgDataAsJsonNode;
if (FetchTo.DATA.equals(fetchTo)) {
msgDataAsJsonNode = getMsgDataAsObjectNode(msg);
} else {
msgDataAsJsonNode = null;
}
ctx.checkTenantEntity(msg.getOriginator());
withCallback(findEntityAsync(ctx, msg.getOriginator()),
entityId -> safeGetAttributes(ctx, msg, entityId, msgDataAsJsonNode),
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
}
protected abstract ListenableFuture<T> findEntityAsync(TbContext ctx, EntityId originator);
private void safeGetAttributes(TbContext ctx, TbMsg msg, T entityId, ObjectNode msgDataAsJsonNode) {
if (entityId == null || entityId.isNullUid()) {
ctx.tellFailure(msg, new NoSuchElementException("Did not find entity! Msg ID: " + msg.getId()));
return;
}
Map<String, String> mappingsMap = new HashMap<>();
config.getAttrMapping().forEach((key, value) -> {
String patternProcessedSourceKey = TbNodeUtils.processPattern(key, msg);
String patternProcessedTargetKey = TbNodeUtils.processPattern(value, msg);
mappingsMap.put(patternProcessedSourceKey, patternProcessedTargetKey);
});
var sourceKeys = List.copyOf(mappingsMap.keySet());
withCallback(config.isTelemetry() ? getLatestTelemetryAsync(ctx, entityId, sourceKeys) : getAttributesAsync(ctx, entityId, sourceKeys),
data -> putDataAndTell(ctx, msg, data, mappingsMap, msgDataAsJsonNode),
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
}
private ListenableFuture<List<KvEntry>> getAttributesAsync(TbContext ctx, EntityId entityId, List<String> attrKeys) {
var latest = ctx.getAttributesService().find(ctx.getTenantId(), entityId, SERVER_SCOPE, attrKeys);
return Futures.transform(latest, l ->
l.stream()
.map(i -> (KvEntry) i)
.collect(Collectors.toList()),
MoreExecutors.directExecutor());
}
private ListenableFuture<List<KvEntry>> getLatestTelemetryAsync(TbContext ctx, EntityId entityId, List<String> timeseriesKeys) {
var latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, timeseriesKeys);
return Futures.transform(latest, l ->
l.stream()
.map(i -> (KvEntry) i)
.collect(Collectors.toList()),
MoreExecutors.directExecutor());
}
private void putDataAndTell(TbContext ctx, TbMsg msg, List<? extends KvEntry> data, Map<String, String> map, ObjectNode msgDataAsJsonNode) {
for (KvEntry entry : data) {
String targetKey = map.get(entry.getKey());
String value = entry.getValueAsString();
if (FetchTo.DATA.equals(fetchTo)) {
msgDataAsJsonNode.put(targetKey, value);
} else if (FetchTo.METADATA.equals(fetchTo)) {
msg.getMetaData().putValue(targetKey, value);
}
}
if (FetchTo.DATA.equals(fetchTo)) {
ctx.tellSuccess(TbMsg.transformMsgData(msg, JacksonUtil.toString(msgDataAsJsonNode)));
} else if (FetchTo.METADATA.equals(fetchTo)) {
ctx.tellSuccess(msg);
}
}
}

30
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java

@ -27,9 +27,6 @@ import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
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.util.EntityDetails;
import org.thingsboard.server.common.data.ContactBased;
import org.thingsboard.server.common.msg.TbMsg;
@ -41,20 +38,11 @@ import java.util.Map;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
@Slf4j
public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEntityDetailsNodeConfiguration> implements TbNode {
public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEntityDetailsNodeConfiguration> extends TbAbstractNodeWithFetchTo<C> {
private static final Gson gson = new Gson();
private static final JsonParser jsonParser = new JsonParser();
private static final Type TYPE = new TypeToken<Map<String, String>>() {
}.getType();
protected C config;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = loadGetEntityDetailsNodeConfiguration(configuration);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
withCallback(getDetails(ctx, msg),
@ -62,22 +50,26 @@ public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEnti
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
}
protected abstract C loadGetEntityDetailsNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException;
protected abstract ListenableFuture<TbMsg> getDetails(TbContext ctx, TbMsg msg);
protected abstract ListenableFuture<ContactBased> getContactBasedListenableFuture(TbContext ctx, TbMsg msg);
protected MessageData getDataAsJson(TbMsg msg) {
if (this.config.isAddToMetadata()) {
if (fetchTo == FetchTo.METADATA) {
return new MessageData(gson.toJsonTree(msg.getMetaData().getData(), TYPE), DataSource.METADATA);
} else if (fetchTo == FetchTo.DATA) {
var msgDataJsonElement = JsonParser.parseString(msg.getData());
if (!msgDataJsonElement.isJsonObject()) {
throw new IllegalArgumentException("Message body is not an object!");
}
return new MessageData(msgDataJsonElement, DataSource.DATA);
} else {
return new MessageData(jsonParser.parse(msg.getData()), DataSource.DATA);
throw new IllegalArgumentException("Unsupported fetchTo value!");
}
}
protected ListenableFuture<TbMsg> getTbMsgListenableFuture(TbContext ctx, TbMsg msg, MessageData messageData, String prefix) {
if (this.config.getDetailsList().isEmpty()) {
if (config.getDetailsList().isEmpty()) {
return Futures.immediateFuture(msg);
} else {
ListenableFuture<ContactBased> contactBasedListenableFuture = getContactBasedListenableFuture(ctx, msg);
@ -178,6 +170,4 @@ public abstract class TbAbstractGetEntityDetailsNode<C extends TbAbstractGetEnti
private enum DataSource {
DATA, METADATA
}
}

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java

@ -16,16 +16,13 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.util.EntityDetails;
import java.util.List;
@Data
public abstract class TbAbstractGetEntityDetailsNodeConfiguration {
@EqualsAndHashCode(callSuper = true)
public abstract class TbAbstractGetEntityDetailsNodeConfiguration extends TbAbstractFetchToNodeConfiguration {
private List<EntityDetails> detailsList;
private boolean addToMetadata;
}

50
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java

@ -0,0 +1,50 @@
/**
* 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;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import org.thingsboard.common.util.JacksonUtil;
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.server.common.msg.TbMsg;
public abstract class TbAbstractNodeWithFetchTo<C extends TbAbstractFetchToNodeConfiguration> implements TbNode {
protected C config;
protected FetchTo fetchTo;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
config = loadNodeConfiguration(configuration);
if (config.getFetchTo() == null) {
throw new TbNodeException("FetchTo cannot be NULL!");
} else {
fetchTo = config.getFetchTo();
}
}
protected abstract C loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException;
protected ObjectNode getMsgDataAsObjectNode(TbMsg msg) {
JsonNode msgDataNode = JacksonUtil.toJsonNode(msg.getData());
if (!msgDataNode.isObject()) {
throw new IllegalArgumentException("Message body is not an object!");
}
return (ObjectNode) msgDataNode;
}
}

109
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java

@ -1,109 +0,0 @@
/**
* 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;
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.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.util.TbNodeUtils;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.msg.TbMsg;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@Slf4j
public abstract class TbEntityGetAttrNode<T extends EntityId> implements TbNode {
private TbGetEntityAttrNodeConfiguration config;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
try {
withCallback(findEntityAsync(ctx, msg.getOriginator()),
entityId -> safeGetAttributes(ctx, msg, entityId),
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
} catch (Throwable th) {
ctx.tellFailure(msg, th);
}
}
private void safeGetAttributes(TbContext ctx, TbMsg msg, T entityId) {
if (entityId == null || entityId.isNullUid()) {
ctx.tellNext(msg, FAILURE);
return;
}
Map<String, String> mappingsMap = new HashMap<>();
config.getAttrMapping().forEach((key, value) -> {
String processPatternKey = TbNodeUtils.processPattern(key, msg);
String processPatternValue = TbNodeUtils.processPattern(value, msg);
mappingsMap.put(processPatternKey, processPatternValue);
});
List<String> keys = List.copyOf(mappingsMap.keySet());
withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId, keys) : getAttributesAsync(ctx, entityId, keys),
attributes -> putAttributesAndTell(ctx, msg, attributes, mappingsMap),
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
}
private ListenableFuture<List<KvEntry>> getAttributesAsync(TbContext ctx, EntityId entityId, List<String> attrKeys) {
ListenableFuture<List<AttributeKvEntry>> latest = ctx.getAttributesService().find(ctx.getTenantId(), entityId, SERVER_SCOPE, attrKeys);
return Futures.transform(latest, l ->
l.stream().map(i -> (KvEntry) i).collect(Collectors.toList()), MoreExecutors.directExecutor());
}
private ListenableFuture<List<KvEntry>> getLatestTelemetry(TbContext ctx, EntityId entityId, List<String> timeseriesKeys) {
ListenableFuture<List<TsKvEntry>> latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, timeseriesKeys);
return Futures.transform(latest, l ->
l.stream().map(i -> (KvEntry) i).collect(Collectors.toList()), MoreExecutors.directExecutor());
}
private void putAttributesAndTell(TbContext ctx, TbMsg msg, List<? extends KvEntry> attributes, Map<String, String> map) {
attributes.forEach(r -> {
String attrName = map.get(r.getKey());
msg.getMetaData().putValue(attrName, r.getValueAsString());
});
ctx.tellSuccess(msg);
}
protected abstract ListenableFuture<T> findEntityAsync(TbContext ctx, EntityId originator);
public void setConfig(TbGetEntityAttrNodeConfiguration config) {
this.config = config;
}
}

40
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java

@ -15,24 +15,18 @@
*/
package org.thingsboard.rule.engine.metadata;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.RuleNode;
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.util.TbNodeUtils;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.concurrent.ExecutionException;
@ -48,39 +42,37 @@ import java.util.concurrent.ExecutionException;
"- send Message via <code>Failure</code> chain, otherwise <code>Success</code> chain is used.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeFetchDeviceCredentialsConfig")
public class TbFetchDeviceCredentialsNode implements TbNode {
public class TbFetchDeviceCredentialsNode extends TbAbstractNodeWithFetchTo<TbFetchDeviceCredentialsNodeConfiguration> {
private static final String CREDENTIALS = "credentials";
private static final String CREDENTIALS_TYPE = "credentialsType";
TbFetchDeviceCredentialsNodeConfiguration config;
boolean fetchToMetadata;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbFetchDeviceCredentialsNodeConfiguration.class);
this.fetchToMetadata = config.isFetchToMetadata();
protected TbFetchDeviceCredentialsNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbFetchDeviceCredentialsNodeConfiguration.class);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException, TbNodeException {
EntityId originator = msg.getOriginator();
var originator = msg.getOriginator();
ctx.checkTenantEntity(originator);
if (!EntityType.DEVICE.equals(originator.getEntityType())) {
ctx.tellFailure(msg, new RuntimeException("Unsupported originator type: " + originator.getEntityType() + "!"));
return;
}
DeviceId deviceId = new DeviceId(msg.getOriginator().getId());
DeviceCredentials deviceCredentials = ctx.getDeviceCredentialsService().findDeviceCredentialsByDeviceId(ctx.getTenantId(), deviceId);
var deviceId = new DeviceId(msg.getOriginator().getId());
var deviceCredentials = ctx.getDeviceCredentialsService().findDeviceCredentialsByDeviceId(ctx.getTenantId(), deviceId);
if (deviceCredentials == null) {
ctx.tellFailure(msg, new RuntimeException("Failed to get Device Credentials for device: " + deviceId + "!"));
return;
}
TbMsg transformedMsg;
DeviceCredentialsType credentialsType = deviceCredentials.getCredentialsType();
JsonNode credentialsInfo = ctx.getDeviceCredentialsService().toCredentialsInfo(deviceCredentials);
if (fetchToMetadata) {
TbMsgMetaData metaData = msg.getMetaData();
TbMsg transformedMsg = null;
var credentialsType = deviceCredentials.getCredentialsType();
var credentialsInfo = ctx.getDeviceCredentialsService().toCredentialsInfo(deviceCredentials);
if (FetchTo.METADATA.equals(fetchTo)) {
var metaData = msg.getMetaData();
metaData.putValue(CREDENTIALS_TYPE, credentialsType.name());
if (credentialsType.equals(DeviceCredentialsType.ACCESS_TOKEN) || credentialsType.equals(DeviceCredentialsType.X509_CERTIFICATE)) {
metaData.putValue(CREDENTIALS, credentialsInfo.asText());
@ -88,8 +80,8 @@ public class TbFetchDeviceCredentialsNode implements TbNode {
metaData.putValue(CREDENTIALS, JacksonUtil.toString(credentialsInfo));
}
transformedMsg = TbMsg.transformMsg(msg, msg.getType(), originator, metaData, msg.getData());
} else {
ObjectNode data = (ObjectNode) JacksonUtil.toJsonNode(msg.getData());
} else if (FetchTo.DATA.equals(fetchTo)) {
var data = getMsgDataAsObjectNode(msg);
data.put(CREDENTIALS_TYPE, credentialsType.name());
data.set(CREDENTIALS, credentialsInfo);
transformedMsg = TbMsg.transformMsg(msg, msg.getType(), originator, msg.getMetaData(), JacksonUtil.toString(data));

11
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeConfiguration.java

@ -17,18 +17,17 @@ package org.thingsboard.rule.engine.metadata;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
@Data
@EqualsAndHashCode(callSuper = true)
@JsonIgnoreProperties(ignoreUnknown = true)
public class TbFetchDeviceCredentialsNodeConfiguration implements NodeConfiguration<TbFetchDeviceCredentialsNodeConfiguration> {
private boolean fetchToMetadata;
public class TbFetchDeviceCredentialsNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration<TbFetchDeviceCredentialsNodeConfiguration> {
@Override
public TbFetchDeviceCredentialsNodeConfiguration defaultConfiguration() {
TbFetchDeviceCredentialsNodeConfiguration configuration = new TbFetchDeviceCredentialsNodeConfiguration();
configuration.setFetchToMetadata(true);
var configuration = new TbFetchDeviceCredentialsNodeConfiguration();
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

15
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java

@ -32,25 +32,24 @@ import org.thingsboard.server.common.msg.TbMsg;
*/
@Slf4j
@RuleNode(type = ComponentType.ENRICHMENT,
name = "originator attributes",
configClazz = TbGetAttributesNodeConfiguration.class,
nodeDescription = "Enrich the message body or metadata with the originator attributes and/or timeseries data",
nodeDetails = "If Attributes enrichment configured, <b>CLIENT/SHARED/SERVER</b> attributes are added into Message data/metadata " +
name = "originator attributes",
configClazz = TbGetAttributesNodeConfiguration.class,
nodeDescription = "Enrich the message body or metadata with the originator attributes and/or timeseries data",
nodeDetails = "If Attributes enrichment configured, <b>CLIENT/SHARED/SERVER</b> attributes are added into Message data/metadata " +
"with specific prefix: <i>cs/shared/ss</i>. Latest telemetry value added into Message data/metadata without prefix. " +
"To access those attributes in other nodes this template can be used " +
"To access those attributes in other nodes this template can be used " +
"<code>metadata.cs_temperature</code> or <code>metadata.shared_limit</code> ",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeOriginatorAttributesConfig")
public class TbGetAttributesNode extends TbAbstractGetAttributesNode<TbGetAttributesNodeConfiguration, EntityId> {
@Override
protected TbGetAttributesNodeConfiguration loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException {
protected TbGetAttributesNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetAttributesNodeConfiguration.class);
}
@Override
protected ListenableFuture<EntityId> findEntityIdAsync(TbContext ctx, TbMsg msg) {
ctx.checkTenantEntity(msg.getOriginator());
return Futures.immediateFuture(msg.getOriginator());
}
}

10
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java

@ -16,6 +16,7 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import java.util.Collections;
@ -25,8 +26,8 @@ import java.util.List;
* Created by ashvayka on 19.01.18.
*/
@Data
public class TbGetAttributesNodeConfiguration implements NodeConfiguration<TbGetAttributesNodeConfiguration> {
@EqualsAndHashCode(callSuper = true)
public class TbGetAttributesNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration<TbGetAttributesNodeConfiguration> {
private List<String> clientAttributeNames;
private List<String> sharedAttributeNames;
private List<String> serverAttributeNames;
@ -35,18 +36,17 @@ public class TbGetAttributesNodeConfiguration implements NodeConfiguration<TbGet
private boolean tellFailureIfAbsent;
private boolean getLatestValueWithTs;
private boolean fetchToData;
@Override
public TbGetAttributesNodeConfiguration defaultConfiguration() {
TbGetAttributesNodeConfiguration configuration = new TbGetAttributesNodeConfiguration();
var configuration = new TbGetAttributesNodeConfiguration();
configuration.setClientAttributeNames(Collections.emptyList());
configuration.setSharedAttributeNames(Collections.emptyList());
configuration.setServerAttributeNames(Collections.emptyList());
configuration.setLatestTsKeyNames(Collections.emptyList());
configuration.setTellFailureIfAbsent(true);
configuration.setGetLatestValueWithTs(false);
configuration.setFetchToData(false);
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

17
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java

@ -18,6 +18,9 @@ package org.thingsboard.rule.engine.metadata;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.rule.engine.api.RuleNode;
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.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;
@ -25,20 +28,24 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
@RuleNode(
type = ComponentType.ENRICHMENT,
name="customer attributes",
name = "customer attributes",
configClazz = TbGetEntityAttrNodeConfiguration.class,
nodeDescription = "Add Originators Customer Attributes or Latest Telemetry into Message Metadata",
nodeDetails = "Enrich the message metadata with the corresponding customer's latest attributes or telemetry value. " +
nodeDescription = "Add Originators Customer Attributes or Latest Telemetry into Message Metadata/Data",
nodeDetails = "Enrich the Message Metadata/Data with the corresponding customer's latest attributes or telemetry value. " +
"The customer is selected based on the originator of the message: device, asset, etc. " +
"</br>" +
"Useful when you store some parameters on the customer level and would like to use them for message processing.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeCustomerAttributesConfig")
public class TbGetCustomerAttributeNode extends TbEntityGetAttrNode<CustomerId> {
public class TbGetCustomerAttributeNode extends TbAbstractGetEntityAttrNode<CustomerId> {
@Override
protected ListenableFuture<CustomerId> findEntityAsync(TbContext ctx, EntityId originator) {
ctx.checkTenantEntity(originator);
return EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, originator);
}
@Override
protected TbGetEntityAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class);
}
}

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java

@ -48,16 +48,16 @@ import org.thingsboard.server.common.msg.TbMsg;
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeEntityDetailsConfig")
public class TbGetCustomerDetailsNode extends TbAbstractGetEntityDetailsNode<TbGetCustomerDetailsNodeConfiguration> {
private static final String CUSTOMER_PREFIX = "customer_";
@Override
protected TbGetCustomerDetailsNodeConfiguration loadGetEntityDetailsNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
protected TbGetCustomerDetailsNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetCustomerDetailsNodeConfiguration.class);
}
@Override
protected ListenableFuture<TbMsg> getDetails(TbContext ctx, TbMsg msg) {
ctx.checkTenantEntity(msg.getOriginator());
return getTbMsgListenableFuture(ctx, msg, getDataAsJson(msg), CUSTOMER_PREFIX);
}

7
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeConfiguration.java

@ -16,18 +16,19 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import java.util.Collections;
@Data
@EqualsAndHashCode(callSuper = true)
public class TbGetCustomerDetailsNodeConfiguration extends TbAbstractGetEntityDetailsNodeConfiguration implements NodeConfiguration<TbGetCustomerDetailsNodeConfiguration> {
@Override
public TbGetCustomerDetailsNodeConfiguration defaultConfiguration() {
TbGetCustomerDetailsNodeConfiguration configuration = new TbGetCustomerDetailsNodeConfiguration();
var configuration = new TbGetCustomerDetailsNodeConfiguration();
configuration.setDetailsList(Collections.emptyList());
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java

@ -39,15 +39,14 @@ import org.thingsboard.server.common.msg.TbMsg;
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeDeviceAttributesConfig")
public class TbGetDeviceAttrNode extends TbAbstractGetAttributesNode<TbGetDeviceAttrNodeConfiguration, DeviceId> {
@Override
protected TbGetDeviceAttrNodeConfiguration loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException {
protected TbGetDeviceAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetDeviceAttrNodeConfiguration.class);
}
@Override
protected ListenableFuture<DeviceId> findEntityIdAsync(TbContext ctx, TbMsg msg) {
ctx.checkTenantEntity(msg.getOriginator());
return EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctx, msg.getOriginator(), config.getDeviceRelationsQuery());
}
}

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java

@ -16,6 +16,7 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.data.DeviceRelationsQuery;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
@ -23,22 +24,22 @@ import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import java.util.Collections;
@Data
@EqualsAndHashCode(callSuper = true)
public class TbGetDeviceAttrNodeConfiguration extends TbGetAttributesNodeConfiguration {
private DeviceRelationsQuery deviceRelationsQuery;
@Override
public TbGetDeviceAttrNodeConfiguration defaultConfiguration() {
TbGetDeviceAttrNodeConfiguration configuration = new TbGetDeviceAttrNodeConfiguration();
var configuration = new TbGetDeviceAttrNodeConfiguration();
configuration.setClientAttributeNames(Collections.emptyList());
configuration.setSharedAttributeNames(Collections.emptyList());
configuration.setServerAttributeNames(Collections.emptyList());
configuration.setLatestTsKeyNames(Collections.emptyList());
configuration.setTellFailureIfAbsent(true);
configuration.setGetLatestValueWithTs(false);
configuration.setFetchToData(false);
configuration.setFetchTo(FetchTo.METADATA);
DeviceRelationsQuery deviceRelationsQuery = new DeviceRelationsQuery();
var deviceRelationsQuery = new DeviceRelationsQuery();
deviceRelationsQuery.setDirection(EntitySearchDirection.FROM);
deviceRelationsQuery.setMaxLevel(1);
deviceRelationsQuery.setRelationType(EntityRelation.CONTAINS_TYPE);

13
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java

@ -16,25 +16,26 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
@Data
public class TbGetEntityAttrNodeConfiguration implements NodeConfiguration<TbGetEntityAttrNodeConfiguration> {
@EqualsAndHashCode(callSuper = true)
public class TbGetEntityAttrNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration<TbGetEntityAttrNodeConfiguration> {
private Map<String, String> attrMapping;
private boolean isTelemetry = false;
@Override
public TbGetEntityAttrNodeConfiguration defaultConfiguration() {
TbGetEntityAttrNodeConfiguration configuration = new TbGetEntityAttrNodeConfiguration();
Map<String, String> attrMapping = new HashMap<>();
attrMapping.putIfAbsent("temperature", "tempo");
var configuration = new TbGetEntityAttrNodeConfiguration();
var attrMapping = new HashMap<String, String>();
attrMapping.putIfAbsent("serialNumber", "sn");
configuration.setAttrMapping(attrMapping);
configuration.setTelemetry(false);
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

10
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsConfiguration.java

@ -16,25 +16,27 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import java.util.HashMap;
import java.util.Map;
@Data
public class TbGetOriginatorFieldsConfiguration implements NodeConfiguration<TbGetOriginatorFieldsConfiguration> {
@EqualsAndHashCode(callSuper = true)
public class TbGetOriginatorFieldsConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration<TbGetOriginatorFieldsConfiguration> {
private Map<String, String> fieldsMapping;
private boolean ignoreNullStrings;
@Override
public TbGetOriginatorFieldsConfiguration defaultConfiguration() {
TbGetOriginatorFieldsConfiguration configuration = new TbGetOriginatorFieldsConfiguration();
Map<String, String> fieldsMapping = new HashMap<>();
var configuration = new TbGetOriginatorFieldsConfiguration();
var fieldsMapping = new HashMap<String, String>();
fieldsMapping.put("name", "originatorName");
fieldsMapping.put("type", "originatorType");
configuration.setFieldsMapping(fieldsMapping);
configuration.setIgnoreNullStrings(false);
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

80
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java

@ -15,13 +15,13 @@
*/
package org.thingsboard.rule.engine.metadata;
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.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.RuleNode;
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.util.TbNodeUtils;
@ -30,56 +30,78 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
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,
nodeDescription = "Add Message Originator fields values into Message Metadata",
nodeDetails = "Will fetch fields values specified in mapping. If specified field is not part of originator fields it will be ignored.",
nodeDescription = "Add Message Originator fields values into Message Metadata or Message Data",
nodeDetails = "Will fetch fields values specified in mapping. If specified field is not part of originator fields it will be ignored. " +
"This node supports only following originator types: TENANT, CUSTOMER, USER, ASSET, DEVICE, ALARM, RULE_CHAIN, ENTITY_VIEW.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeOriginatorFieldsConfig")
public class TbGetOriginatorFieldsNode implements TbNode {
private TbGetOriginatorFieldsConfiguration config;
private boolean ignoreNullStrings;
public class TbGetOriginatorFieldsNode extends TbAbstractNodeWithFetchTo<TbGetOriginatorFieldsConfiguration> {
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
config = TbNodeUtils.convert(configuration, TbGetOriginatorFieldsConfiguration.class);
ignoreNullStrings = config.isIgnoreNullStrings();
protected TbGetOriginatorFieldsConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetOriginatorFieldsConfiguration.class);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
try {
withCallback(putEntityFields(ctx, msg.getOriginator(), msg),
i -> ctx.tellSuccess(msg), t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
} catch (Throwable th) {
ctx.tellFailure(msg, th);
ObjectNode msgDataAsJsonNode;
if (FetchTo.DATA.equals(fetchTo)) {
msgDataAsJsonNode = getMsgDataAsObjectNode(msg);
} else {
msgDataAsJsonNode = null;
}
ctx.checkTenantEntity(msg.getOriginator());
withCallback(collectMappedEntityFieldsAsync(ctx, msg.getOriginator()),
targetKeysToSourceValuesMap -> {
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)) {
msg.getMetaData().putValue(targetKeyName, sourceFieldValue);
}
}
if (FetchTo.DATA.equals(fetchTo)) {
ctx.tellSuccess(TbMsg.transformMsgData(msg, JacksonUtil.toString(msgDataAsJsonNode)));
} else if (FetchTo.METADATA.equals(fetchTo)) {
ctx.tellSuccess(msg);
}
},
t -> ctx.tellFailure(msg, t),
MoreExecutors.directExecutor());
}
private ListenableFuture<Void> putEntityFields(TbContext ctx, EntityId entityId, TbMsg msg) {
private ListenableFuture<Map<String, String>> collectMappedEntityFieldsAsync(TbContext ctx, EntityId entityId) {
if (config.getFieldsMapping().isEmpty()) {
return Futures.immediateFuture(null);
return Futures.immediateFuture(Collections.emptyMap());
} else {
return Futures.transform(EntitiesFieldsAsyncLoader.findAsync(ctx, entityId),
data -> {
config.getFieldsMapping().forEach((field, metaKey) -> {
String val = data.getFieldValue(field, ignoreNullStrings);
if (val != null) {
msg.getMetaData().putValue(metaKey, val);
fieldsData -> {
var targetKeysToSourceValuesMap = new HashMap<String, String>();
for (var mappingEntry : config.getFieldsMapping().entrySet()) {
var sourceFieldName = mappingEntry.getKey();
var targetKeyName = mappingEntry.getValue();
var sourceFieldValue = fieldsData.getFieldValue(sourceFieldName, config.isIgnoreNullStrings());
if (sourceFieldValue != null) {
targetKeysToSourceValuesMap.put(targetKeyName, sourceFieldValue);
}
});
return null;
}, MoreExecutors.directExecutor()
}
return targetKeysToSourceValuesMap;
}, ctx.getDbCallbackExecutor()
);
}
}
}

14
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java

@ -16,6 +16,7 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.data.RelationsQuery;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
@ -23,25 +24,24 @@ import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
@Data
@EqualsAndHashCode(callSuper = true)
public class TbGetRelatedAttrNodeConfiguration extends TbGetEntityAttrNodeConfiguration {
private RelationsQuery relationsQuery;
@Override
public TbGetRelatedAttrNodeConfiguration defaultConfiguration() {
TbGetRelatedAttrNodeConfiguration configuration = new TbGetRelatedAttrNodeConfiguration();
Map<String, String> attrMapping = new HashMap<>();
attrMapping.putIfAbsent("temperature", "tempo");
var configuration = new TbGetRelatedAttrNodeConfiguration();
var attrMapping = new HashMap<String, String>();
attrMapping.putIfAbsent("serialNumber", "sn");
configuration.setAttrMapping(attrMapping);
configuration.setTelemetry(false);
RelationsQuery relationsQuery = new RelationsQuery();
var relationsQuery = new RelationsQuery();
relationsQuery.setDirection(EntitySearchDirection.FROM);
relationsQuery.setMaxLevel(1);
RelationEntityTypeFilter relationEntityTypeFilter = new RelationEntityTypeFilter(EntityRelation.CONTAINS_TYPE, Collections.emptyList());
var relationEntityTypeFilter = new RelationEntityTypeFilter(EntityRelation.CONTAINS_TYPE, Collections.emptyList());
relationsQuery.setFilters(Collections.singletonList(relationEntityTypeFilter));
configuration.setRelationsQuery(relationsQuery);

25
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java

@ -27,30 +27,27 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
@RuleNode(
type = ComponentType.ENRICHMENT,
name="related attributes",
name = "related attributes",
configClazz = TbGetRelatedAttrNodeConfiguration.class,
nodeDescription = "Add Originators Related Entity Attributes or Latest Telemetry into Message Metadata",
nodeDescription = "Add Originators Related Entity Attributes or Latest Telemetry into Message Metadata/Data",
nodeDetails = "Related Entity found using configured relation direction and Relation Type. " +
"If multiple Related Entities are found, only first Entity is used for attributes enrichment, other entities are discarded. " +
"If Attributes enrichment configured, server scope attributes are added into Message metadata. " +
"If Latest Telemetry enrichment configured, latest telemetry added into metadata. " +
"If Attributes enrichment configured, server scope attributes are added into Message Metadata/Data. " +
"If Latest Telemetry enrichment configured, latest telemetry added into Metadata/Data. " +
"To access those attributes in other nodes this template can be used " +
"<code>metadata.temperature</code>.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeRelatedAttributesConfig")
public class TbGetRelatedAttributeNode extends TbEntityGetAttrNode<EntityId> {
private TbGetRelatedAttrNodeConfiguration config;
public class TbGetRelatedAttributeNode extends TbAbstractGetEntityAttrNode<EntityId> {
@Override
public void init(TbContext context, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbGetRelatedAttrNodeConfiguration.class);
setConfig(config);
public TbGetRelatedAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetRelatedAttrNodeConfiguration.class);
}
@Override
protected ListenableFuture<EntityId> findEntityAsync(TbContext ctx, EntityId originator) {
return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, originator, config.getRelationsQuery());
public ListenableFuture<EntityId> findEntityAsync(TbContext ctx, EntityId originator) {
ctx.checkTenantEntity(originator);
var relatedAttrConfig = (TbGetRelatedAttrNodeConfiguration) config;
return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, originator, relatedAttrConfig.getRelationsQuery());
}
}

20
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java

@ -20,6 +20,9 @@ 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;
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.TenantId;
import org.thingsboard.server.common.data.plugin.ComponentType;
@ -27,21 +30,24 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
@Slf4j
@RuleNode(
type = ComponentType.ENRICHMENT,
name="tenant attributes",
name = "tenant attributes",
configClazz = TbGetEntityAttrNodeConfiguration.class,
nodeDescription = "Add Originators Tenant Attributes or Latest Telemetry into Message Metadata",
nodeDetails = "If Attributes enrichment configured, server scope attributes are added into Message metadata. " +
"If Latest Telemetry enrichment configured, latest telemetry added into metadata. " +
nodeDescription = "Add Originators Tenant Attributes or Latest Telemetry into Message Metadata/Data",
nodeDetails = "If Attributes enrichment configured, server scope attributes are added into Message Metadata/Data. " +
"If Latest Telemetry enrichment configured, latest telemetry added into Metadata/Data. " +
"To access those attributes in other nodes this template can be used " +
"<code>metadata.temperature</code>.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeTenantAttributesConfig")
public class TbGetTenantAttributeNode extends TbEntityGetAttrNode<TenantId> {
public class TbGetTenantAttributeNode extends TbAbstractGetEntityAttrNode<TenantId> {
@Override
protected ListenableFuture<TenantId> findEntityAsync(TbContext ctx, EntityId originator) {
public ListenableFuture<TenantId> findEntityAsync(TbContext ctx, EntityId originator) {
ctx.checkTenantEntity(originator);
return Futures.immediateFuture(ctx.getTenantId());
}
@Override
public TbGetEntityAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class);
}
}

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java

@ -39,16 +39,16 @@ import org.thingsboard.server.common.msg.TbMsg;
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeEntityDetailsConfig")
public class TbGetTenantDetailsNode extends TbAbstractGetEntityDetailsNode<TbGetTenantDetailsNodeConfiguration> {
private static final String TENANT_PREFIX = "tenant_";
@Override
protected TbGetTenantDetailsNodeConfiguration loadGetEntityDetailsNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
protected TbGetTenantDetailsNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetTenantDetailsNodeConfiguration.class);
}
@Override
protected ListenableFuture<TbMsg> getDetails(TbContext ctx, TbMsg msg) {
ctx.checkTenantEntity(msg.getOriginator());
return getTbMsgListenableFuture(ctx, msg, getDataAsJson(msg), TENANT_PREFIX);
}

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeConfiguration.java

@ -16,18 +16,19 @@
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import java.util.Collections;
@Data
@EqualsAndHashCode(callSuper = true)
public class TbGetTenantDetailsNodeConfiguration extends TbAbstractGetEntityDetailsNodeConfiguration implements NodeConfiguration<TbGetTenantDetailsNodeConfiguration> {
@Override
public TbGetTenantDetailsNodeConfiguration defaultConfiguration() {
TbGetTenantDetailsNodeConfiguration configuration = new TbGetTenantDetailsNodeConfiguration();
configuration.setDetailsList(Collections.emptyList());
configuration.setFetchTo(FetchTo.METADATA);
return configuration;
}
}

6
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java

@ -31,14 +31,12 @@ import org.thingsboard.server.common.msg.queue.TbMsgCallback;
import java.util.List;
import static org.thingsboard.common.util.DonAsynchron.withCallback;
import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE;
/**
* Created by ashvayka on 19.01.18.
*/
@Slf4j
public abstract class TbAbstractTransformNode implements TbNode {
private TbTransformNodeConfiguration config;
@Override
@ -62,7 +60,7 @@ public abstract class TbAbstractTransformNode implements TbNode {
if (m != null) {
ctx.tellSuccess(m);
} else {
ctx.tellNext(msg, FAILURE);
ctx.tellFailure(msg, new RuntimeException("Message is null!"));
}
}
@ -85,7 +83,7 @@ public abstract class TbAbstractTransformNode implements TbNode {
msgs.forEach(newMsg -> ctx.enqueueForTellNext(newMsg, TbRelationTypes.SUCCESS, wrapper::onSuccess, wrapper::onFailure));
}
} else {
ctx.tellNext(msg, FAILURE);
ctx.tellFailure(msg, new RuntimeException("Message or messages list are empty!"));
}
}

34
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoader.java

@ -18,6 +18,7 @@ package org.thingsboard.rule.engine.util;
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.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.BaseData;
@ -30,47 +31,50 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityViewId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.id.UserId;
import java.util.function.Function;
@Slf4j
public class EntitiesFieldsAsyncLoader {
public static ListenableFuture<EntityFieldsData> findAsync(TbContext ctx, EntityId original) {
switch (original.getEntityType()) {
public static ListenableFuture<EntityFieldsData> findAsync(TbContext ctx, EntityId originatorId) {
switch (originatorId.getEntityType()) {
case TENANT:
return getAsync(ctx.getTenantService().findTenantByIdAsync(ctx.getTenantId(), (TenantId) original),
return toEntityFieldsDataAsync(ctx.getTenantService().findTenantByIdAsync(ctx.getTenantId(), (TenantId) originatorId),
EntityFieldsData::new);
case CUSTOMER:
return getAsync(ctx.getCustomerService().findCustomerByIdAsync(ctx.getTenantId(), (CustomerId) original),
return toEntityFieldsDataAsync(ctx.getCustomerService().findCustomerByIdAsync(ctx.getTenantId(), (CustomerId) originatorId),
EntityFieldsData::new);
case USER:
return getAsync(ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), (UserId) original),
return toEntityFieldsDataAsync(ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), (UserId) originatorId),
EntityFieldsData::new);
case ASSET:
return getAsync(ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), (AssetId) original),
return toEntityFieldsDataAsync(ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), (AssetId) originatorId),
EntityFieldsData::new);
case DEVICE:
return getAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), (DeviceId) original),
return toEntityFieldsDataAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), (DeviceId) originatorId),
EntityFieldsData::new);
case ALARM:
return getAsync(ctx.getAlarmService().findAlarmByIdAsync(ctx.getTenantId(), (AlarmId) original),
return toEntityFieldsDataAsync(ctx.getAlarmService().findAlarmByIdAsync(ctx.getTenantId(), (AlarmId) originatorId),
EntityFieldsData::new);
case RULE_CHAIN:
return getAsync(ctx.getRuleChainService().findRuleChainByIdAsync(ctx.getTenantId(), (RuleChainId) original),
return toEntityFieldsDataAsync(ctx.getRuleChainService().findRuleChainByIdAsync(ctx.getTenantId(), (RuleChainId) originatorId),
EntityFieldsData::new);
case ENTITY_VIEW:
return getAsync(ctx.getEntityViewService().findEntityViewByIdAsync(ctx.getTenantId(), (EntityViewId) original),
return toEntityFieldsDataAsync(ctx.getEntityViewService().findEntityViewByIdAsync(ctx.getTenantId(), (EntityViewId) originatorId),
EntityFieldsData::new);
default:
return Futures.immediateFailedFuture(new TbNodeException("Unexpected original EntityType " + original.getEntityType()));
return Futures.immediateFailedFuture(new TbNodeException("Unexpected originator EntityType: " + originatorId.getEntityType()));
}
}
private static <T extends BaseData> ListenableFuture<EntityFieldsData> getAsync(
ListenableFuture<T> future, Function<T, EntityFieldsData> converter) {
private static <T extends BaseData<? extends UUIDBased>> ListenableFuture<EntityFieldsData> toEntityFieldsDataAsync(
ListenableFuture<T> future,
Function<T, EntityFieldsData> converter
) {
return Futures.transformAsync(future, in -> in != null ?
Futures.immediateFuture(converter.apply(in))
: Futures.immediateFailedFuture(new RuntimeException("Entity not found!")), MoreExecutors.directExecutor());
: Futures.immediateFailedFuture(new TbNodeException("Entity not found!")), MoreExecutors.directExecutor());
}
}

20
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/AbstractAttributeNodeTest.java → rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractAttributeNodeTest.java

@ -52,7 +52,9 @@ import org.thingsboard.server.dao.user.UserService;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.NoSuchElementException;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue;
import static org.mockito.ArgumentMatchers.any;
@ -61,11 +63,10 @@ import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.ArgumentMatchers.same;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@RunWith(MockitoJUnitRunner.class)
public abstract class AbstractAttributeNodeTest {
public abstract class TbAbstractAttributeNodeTest {
final CustomerId customerId = new CustomerId(Uuids.timeBased());
final TenantId tenantId = TenantId.fromUUID(Uuids.timeBased());
final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased());
@ -86,9 +87,9 @@ public abstract class AbstractAttributeNodeTest {
DeviceService deviceService;
TbMsg msg;
Map<String, String> metaData;
TbEntityGetAttrNode node;
TbAbstractGetEntityAttrNode node;
void init(TbEntityGetAttrNode node) throws TbNodeException {
void init(TbAbstractGetEntityAttrNode node) throws TbNodeException {
ObjectMapper mapper = JacksonUtil.OBJECT_MAPPER;
TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(getTbNodeConfig()));
@ -117,7 +118,6 @@ public abstract class AbstractAttributeNodeTest {
}
void errorThrownIfCannotLoadAttributesAsync(User user) {
msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
when(ctx.getAttributesService()).thenReturn(attributesService);
@ -137,7 +137,10 @@ public abstract class AbstractAttributeNodeTest {
msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
node.onMsg(ctx, msg);
verify(ctx).tellNext(msg, FAILURE);
var exceptionCaptor = ArgumentCaptor.forClass(NoSuchElementException.class);
verify(ctx).tellFailure(eq(msg), exceptionCaptor.capture());
assertThat(exceptionCaptor.getValue().getMessage()).contains("Did not find entity! Msg ID: ");
assertTrue(msg.getMetaData().getData().isEmpty());
}
@ -168,7 +171,7 @@ public abstract class AbstractAttributeNodeTest {
ObjectMapper mapper = JacksonUtil.OBJECT_MAPPER;
TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(getTbNodeConfigForTelemetry()));
TbEntityGetAttrNode node = getEmptyNode();
TbAbstractGetEntityAttrNode node = getEmptyNode();
node.init(null, nodeConfiguration);
msg = TbMsg.newMsg("DEVICE", device.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
@ -210,10 +213,11 @@ public abstract class AbstractAttributeNodeTest {
conf.put(keyAttrConf, valueAttrConf);
config.setAttrMapping(conf);
config.setTelemetry(isTelemetry);
config.setFetchTo(FetchTo.METADATA);
return config;
}
protected abstract TbEntityGetAttrNode getEmptyNode();
protected abstract TbAbstractGetEntityAttrNode getEmptyNode();
abstract EntityId getEntityId();

58
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNodeTest.java

@ -130,7 +130,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToMetadata_whenOnMsg_then_success() throws Exception {
TbGetAttributesNode node = initNode(false, false, false);
TbGetAttributesNode node = initNode(FetchTo.METADATA, false, false);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -138,9 +138,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg resultMsg = checkMsg(true);
//check attributes
checkAttributes(resultMsg, false, "cs_", clientAttributes);
checkAttributes(resultMsg, false, "ss_", serverAttributes);
checkAttributes(resultMsg, false, "shared_", sharedAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "cs_", clientAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "ss_", serverAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "shared_", sharedAttributes);
//check timeseries
checkTs(resultMsg, false, false, tsKeys);
@ -148,7 +148,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToMetadata_latestWithTs_whenOnMsg_then_success() throws Exception {
TbGetAttributesNode node = initNode(false, true, false);
TbGetAttributesNode node = initNode(FetchTo.METADATA, true, false);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -156,9 +156,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg resultMsg = checkMsg(true);
//check attributes
checkAttributes(resultMsg, false, "cs_", clientAttributes);
checkAttributes(resultMsg, false, "ss_", serverAttributes);
checkAttributes(resultMsg, false, "shared_", sharedAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "cs_", clientAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "ss_", serverAttributes);
checkAttributes(resultMsg, FetchTo.METADATA, "shared_", sharedAttributes);
//check timeseries with ts
checkTs(resultMsg, false, true, tsKeys);
@ -166,7 +166,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToData_whenOnMsg_then_success() throws Exception {
TbGetAttributesNode node = initNode(true, false, false);
TbGetAttributesNode node = initNode(FetchTo.DATA, false, false);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -174,9 +174,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg resultMsg = checkMsg(true);
//check attributes
checkAttributes(resultMsg, true, "cs_", clientAttributes);
checkAttributes(resultMsg, true, "ss_", serverAttributes);
checkAttributes(resultMsg, true, "shared_", sharedAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "cs_", clientAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "ss_", serverAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "shared_", sharedAttributes);
//check timeseries
checkTs(resultMsg, true, false, tsKeys);
@ -184,7 +184,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToData_latestWithTs_whenOnMsg_then_success() throws Exception {
TbGetAttributesNode node = initNode(true, true, false);
TbGetAttributesNode node = initNode(FetchTo.DATA, true, false);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -192,9 +192,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg resultMsg = checkMsg(true);
//check attributes
checkAttributes(resultMsg, true, "cs_", clientAttributes);
checkAttributes(resultMsg, true, "ss_", serverAttributes);
checkAttributes(resultMsg, true, "shared_", sharedAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "cs_", clientAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "ss_", serverAttributes);
checkAttributes(resultMsg, FetchTo.DATA, "shared_", sharedAttributes);
//check timeseries with ts
checkTs(resultMsg, true, true, tsKeys);
@ -202,7 +202,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToMetadata_whenOnMsg_then_failure() throws Exception {
TbGetAttributesNode node = initNode(false, false, true);
TbGetAttributesNode node = initNode(FetchTo.METADATA, false, true);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -210,9 +210,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg actualMsg = checkMsg(false);
//check attributes
checkAttributes(actualMsg, false, "cs_", clientAttributes);
checkAttributes(actualMsg, false, "ss_", serverAttributes);
checkAttributes(actualMsg, false, "shared_", sharedAttributes);
checkAttributes(actualMsg, FetchTo.METADATA, "cs_", clientAttributes);
checkAttributes(actualMsg, FetchTo.METADATA, "ss_", serverAttributes);
checkAttributes(actualMsg, FetchTo.METADATA, "shared_", sharedAttributes);
//check timeseries with ts
checkTs(actualMsg, false, false, tsKeys);
@ -220,7 +220,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToData_whenOnMsg_then_failure() throws Exception {
TbGetAttributesNode node = initNode(true, true, true);
TbGetAttributesNode node = initNode(FetchTo.DATA, true, true);
TbMsg msg = getTbMsg(originator);
node.onMsg(ctx, msg);
@ -228,9 +228,9 @@ public class TbAbstractGetAttributesNodeTest {
TbMsg actualMsg = checkMsg(false);
//check attributes
checkAttributes(actualMsg, true, "cs_", clientAttributes);
checkAttributes(actualMsg, true, "ss_", serverAttributes);
checkAttributes(actualMsg, true, "shared_", sharedAttributes);
checkAttributes(actualMsg, FetchTo.DATA, "cs_", clientAttributes);
checkAttributes(actualMsg, FetchTo.DATA, "ss_", serverAttributes);
checkAttributes(actualMsg, FetchTo.DATA, "shared_", sharedAttributes);
//check timeseries with ts
checkTs(actualMsg, true, true, tsKeys);
@ -238,7 +238,7 @@ public class TbAbstractGetAttributesNodeTest {
@Test
public void fetchToData_whenOnMsg_and_data_is_not_object_then_failure() throws Exception {
TbGetAttributesNode node = initNode(true, true, true);
TbGetAttributesNode node = initNode(FetchTo.DATA, true, true);
TbMsg msg = TbMsg.newMsg("TEST", originator, new TbMsgMetaData(), "[]");
node.onMsg(ctx, msg);
@ -272,13 +272,13 @@ public class TbAbstractGetAttributesNodeTest {
return resultMsg;
}
private void checkAttributes(TbMsg actualMsg, boolean fetchToData, String prefix, List<String> attributes) {
private void checkAttributes(TbMsg actualMsg, FetchTo fetchTo, String prefix, List<String> attributes) {
JsonNode msgData = JacksonUtil.toJsonNode(actualMsg.getData());
attributes.stream()
.filter(attribute -> !attribute.equals("unknown"))
.forEach(attribute -> {
String result;
if (fetchToData) {
if (FetchTo.DATA.equals(fetchTo)) {
result = msgData.get(prefix + attribute).asText();
} else {
result = actualMsg.getMetaData().getValue(prefix + attribute);
@ -313,13 +313,13 @@ public class TbAbstractGetAttributesNodeTest {
}
}
private TbGetAttributesNode initNode(boolean fetchToData, boolean getLatestValueWithTs, boolean isTellFailureIfAbsent) throws TbNodeException {
private TbGetAttributesNode initNode(FetchTo fetchTo, boolean getLatestValueWithTs, boolean isTellFailureIfAbsent) throws TbNodeException {
TbGetAttributesNodeConfiguration config = new TbGetAttributesNodeConfiguration();
config.setClientAttributeNames(List.of("client_attr_1", "client_attr_2", "${client_attr_metadata}", "unknown"));
config.setServerAttributeNames(List.of("server_attr_1", "server_attr_2", "${server_attr_metadata}", "unknown"));
config.setSharedAttributeNames(List.of("shared_attr_1", "shared_attr_2", "$[shared_attr_data]", "unknown"));
config.setLatestTsKeyNames(List.of("temperature", "humidity", "unknown"));
config.setFetchToData(fetchToData);
config.setFetchTo(fetchTo);
config.setGetLatestValueWithTs(getLatestValueWithTs);
config.setTellFailureIfAbsent(isTellFailureIfAbsent);
TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config));

6
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeTest.java

@ -64,7 +64,7 @@ public class TbFetchDeviceCredentialsNodeTest {
callback = mock(TbMsgCallback.class);
ctx = mock(TbContext.class);
config = new TbFetchDeviceCredentialsNodeConfiguration().defaultConfiguration();
config.setFetchToMetadata(true);
config.setFetchTo(FetchTo.METADATA);
nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config));
node = spy(new TbFetchDeviceCredentialsNode());
node.init(ctx, nodeConfiguration);
@ -89,13 +89,13 @@ public class TbFetchDeviceCredentialsNodeTest {
@Test
void givenDefaultConfig_whenInit_thenOK() {
assertThat(node.config).isEqualTo(config);
assertThat(node.fetchToMetadata).isEqualTo(true);
assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA);
}
@Test
void givenDefaultConfig_whenVerify_thenOK() {
TbFetchDeviceCredentialsNodeConfiguration defaultConfig = new TbFetchDeviceCredentialsNodeConfiguration().defaultConfiguration();
assertThat(defaultConfig.isFetchToMetadata()).isEqualTo(true);
assertThat(defaultConfig.getFetchTo()).isEqualTo(FetchTo.METADATA);
}
@Test

73
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java

@ -16,28 +16,47 @@
package org.thingsboard.rule.engine.metadata;
import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.Futures;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.ArgumentCaptor;
import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.List;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyCollection;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@RunWith(MockitoJUnitRunner.class)
public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest {
public class TbGetCustomerAttributeNodeTest extends TbAbstractAttributeNodeTest {
User user = new User();
Asset asset = new Asset();
Device device = new Device();
@ -56,7 +75,7 @@ public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest {
}
@Override
protected TbEntityGetAttrNode getEmptyNode() {
protected TbAbstractGetEntityAttrNode getEmptyNode() {
return new TbGetCustomerAttributeNode();
}
@ -65,6 +84,31 @@ public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest {
return customerId;
}
@Test
public void errorThrownIfFetchToIsNull() {
var node = new TbGetCustomerAttributeNode();
var config = new TbGetEntityAttrNodeConfiguration().defaultConfiguration();
config.setFetchTo(null);
var nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
var exception = assertThrows(TbNodeException.class, () -> node.init(ctx, nodeConfiguration));
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be NULL!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfMsgDataIsNotAnObjectAndFetchToData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", new CustomerId(UUID.randomUUID()), new TbMsgMetaData(), "[]");
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctx, msg));
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfCannotLoadAttributes() {
mockFindUser(user);
@ -89,6 +133,29 @@ public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest {
entityAttributeAddedInMetadata(customerId, "CUSTOMER");
}
@Test
public void customerAttributeAddedInData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
msg = TbMsg.newMsg("CUSTOMER", customerId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
List<AttributeKvEntry> attributes = Lists.newArrayList(new BaseAttributeKvEntry(new StringDataEntry("temperature", "high"), 1L));
when(ctx.getAttributesService()).thenReturn(attributesService);
when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection()))
.thenReturn(Futures.immediateFuture(attributes));
node.onMsg(ctx, msg);
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctx, times(1)).tellSuccess(actualMessageCaptor.capture());
var expectedMsgData = "{\"answer\":\"high\"}";
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData);
}
@Test
public void usersCustomerAttributesFetched() {
mockFindUser(user);

376
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNodeTest.java

@ -0,0 +1,376 @@
/**
* 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;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import org.jetbrains.annotations.NotNull;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.JacksonUtil;
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.server.common.data.Device;
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.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.device.DeviceService;
import java.util.Collections;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.Callable;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
public class TbGetOriginatorFieldsNodeTest {
private static final EntityId DUMMY_ENTITY_ID = new DeviceId(UUID.randomUUID());
public static final ListeningExecutor DB_EXECUTOR = new ListeningExecutor() {
@Override
public <T> ListenableFuture<T> executeAsync(Callable<T> task) {
try {
return Futures.immediateFuture(task.call());
} catch (Exception e) {
throw new RuntimeException(e);
}
}
@Override
public void execute(@NotNull Runnable command) {
command.run();
}
};
@Mock
private TbContext ctxMock;
@Mock
private DeviceService deviceService;
private TbGetOriginatorFieldsNode node;
private TbGetOriginatorFieldsConfiguration config;
private TbNodeConfiguration nodeConfiguration;
private TbMsg msg;
@BeforeEach
public void setUp() {
config = new TbGetOriginatorFieldsConfiguration();
node = new TbGetOriginatorFieldsNode();
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
}
@Test
public void givenConfigWithNullFetchTo_whenOnInit_thenException() {
// GIVEN
config = config.defaultConfiguration();
config.setFetchTo(null);
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
// WHEN
var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration));
// THEN
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be NULL!");
verify(ctxMock, never()).tellSuccess(any());
}
@Test
public void givenDefaultConfig_whenInit_thenOK() throws TbNodeException {
// GIVEN
config = config.defaultConfiguration();
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
// WHEN
node.init(ctxMock, nodeConfiguration);
// THEN
assertThat(node.config).isEqualTo(config);
assertThat(config.getFieldsMapping()).isEqualTo(Map.of(
"name", "originatorName",
"type", "originatorType"));
assertThat(config.isIgnoreNullStrings()).isEqualTo(false);
assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA);
}
@Test
public void givenCustomConfig_whenInit_thenOK() throws TbNodeException {
// GIVEN
config.setFieldsMapping(Map.of(
"sourceField1", "targetKey1",
"sourceField2", "targetKey2",
"sourceField3", "targetKey3"));
config.setIgnoreNullStrings(true);
config.setFetchTo(FetchTo.DATA);
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
// WHEN
node.init(ctxMock, nodeConfiguration);
// THEN
assertThat(node.config).isEqualTo(config);
assertThat(config.getFieldsMapping()).isEqualTo(Map.of(
"sourceField1", "targetKey1",
"sourceField2", "targetKey2",
"sourceField3", "targetKey3"));
assertThat(config.isIgnoreNullStrings()).isEqualTo(true);
assertThat(node.fetchTo).isEqualTo(FetchTo.DATA);
}
@Test
public void givenMsgDataIsNotAnJsonObjectAndFetchToData_whenOnMsg_thenException() {
// GIVEN
node.fetchTo = FetchTo.DATA;
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", DUMMY_ENTITY_ID, new TbMsgMetaData(), "[]");
// WHEN
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctxMock, msg));
// THEN
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!");
verify(ctxMock, never()).tellSuccess(any());
}
@Test
public void givenEntityThatDoesNotBelongToTheCurrentTenant_whenOnMsg_thenException() {
// SETUP
var expectedExceptionMessage = "Entity with id: '" + DUMMY_ENTITY_ID +
"' specified in the configuration doesn't belong to the current tenant.";
// GIVEN
doThrow(new RuntimeException(expectedExceptionMessage)).when(ctxMock).checkTenantEntity(DUMMY_ENTITY_ID);
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", DUMMY_ENTITY_ID, new TbMsgMetaData(), "{}");
// WHEN
var exception = assertThrows(RuntimeException.class, () -> node.onMsg(ctxMock, msg));
// THEN
assertThat(exception.getMessage()).isEqualTo(expectedExceptionMessage);
verify(ctxMock, never()).tellSuccess(any());
}
@Test
public void givenValidMsgAndFetchToData_whenOnMsg_thenShouldTellSuccessAndFetchToData() {
// GIVEN
var device = new Device();
device.setId((DeviceId) DUMMY_ENTITY_ID);
device.setName("Test device");
device.setType("Test device type");
config.setFieldsMapping(Map.of(
"name", "originatorName",
"type", "originatorType",
"label", "originatorLabel"));
config.setIgnoreNullStrings(true);
config.setFetchTo(FetchTo.DATA);
node.config = config;
node.fetchTo = FetchTo.DATA;
var msgMetaData = new TbMsgMetaData();
var msgData = "{\"temp\":42,\"humidity\":77}";
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_ENTITY_ID, msgMetaData, msgData);
when(ctxMock.getDeviceService()).thenReturn(deviceService);
when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(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,\"originatorName\":\"Test device\",\"originatorType\":\"Test device type\"}";
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData);
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msgMetaData);
}
@Test
public void givenValidMsgAndFetchToMetaData_whenOnMsg_thenShouldTellSuccessAndFetchToMetaData() {
// GIVEN
var device = new Device();
device.setId((DeviceId) DUMMY_ENTITY_ID);
device.setName("Test device");
device.setType("Test device type");
config.setFieldsMapping(Map.of(
"name", "originatorName",
"type", "originatorType",
"label", "originatorLabel"));
config.setIgnoreNullStrings(true);
config.setFetchTo(FetchTo.METADATA);
node.config = config;
node.fetchTo = FetchTo.METADATA;
var msgMetaData = new TbMsgMetaData(Map.of(
"testKey1", "testValue1",
"testKey2", "123"));
var msgData = "[\"value1\",\"value2\"]";
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_ENTITY_ID, msgMetaData, msgData);
when(ctxMock.getDeviceService()).thenReturn(deviceService);
when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(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(
"testKey1", "testValue1",
"testKey2", "123",
"originatorName", "Test device",
"originatorType", "Test device type"
));
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msgData);
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData);
}
@Test
public void givenNullEntityFieldsAndIgnoreNullStringsFalse_whenOnMsg_thenShouldTellSuccessAndFetchNullField() {
// GIVEN
var device = new Device();
device.setId((DeviceId) DUMMY_ENTITY_ID);
device.setName("Test device");
device.setType("Test device type");
config.setFieldsMapping(Map.of(
"name", "originatorName",
"type", "originatorType",
"label", "originatorLabel"));
config.setIgnoreNullStrings(false);
config.setFetchTo(FetchTo.METADATA);
node.config = config;
node.fetchTo = FetchTo.METADATA;
var msgMetaData = new TbMsgMetaData(Map.of(
"testKey1", "testValue1",
"testKey2", "123"));
var msgData = "[\"value1\",\"value2\"]";
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_ENTITY_ID, msgMetaData, msgData);
when(ctxMock.getDeviceService()).thenReturn(deviceService);
when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(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(
"testKey1", "testValue1",
"testKey2", "123",
"originatorName", "Test device",
"originatorType", "Test device type",
"originatorLabel", "null"
));
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msgData);
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData);
}
@Test
public void givenEmptyFieldsMapping_whenOnMsg_thenShouldTellSuccessWithSameMsg() {
// GIVEN
config.setFieldsMapping(Collections.emptyMap());
config.setIgnoreNullStrings(false);
config.setFetchTo(FetchTo.METADATA);
node.config = config;
node.fetchTo = FetchTo.METADATA;
var msgMetaData = new TbMsgMetaData(Map.of(
"testKey1", "testValue1",
"testKey2", "123"));
var msgData = "[\"value1\",\"value2\"]";
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_ENTITY_ID, msgMetaData, msgData);
// 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(
"testKey1", "testValue1",
"testKey2", "123"
));
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msgData);
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData);
}
@Test
public void givenUnsupportedEntityType_whenOnMsg_thenShouldTellFailureWithSameMsg() {
// GIVEN
config.setFieldsMapping(Map.of(
"name", "originatorName",
"type", "originatorType",
"label", "originatorLabel"));
config.setIgnoreNullStrings(false);
config.setFetchTo(FetchTo.METADATA);
node.config = config;
node.fetchTo = FetchTo.METADATA;
var msgMetaData = new TbMsgMetaData(Map.of(
"testKey1", "testValue1",
"testKey2", "123"));
var msgData = "[\"value1\",\"value2\"]";
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", new DashboardId(UUID.randomUUID()), msgMetaData, msgData);
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN
node.onMsg(ctxMock, msg);
// THEN
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctxMock, times(1)).tellFailure(actualMessageCaptor.capture(), any());
verify(ctxMock, never()).tellSuccess(any());
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msgData);
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msgMetaData);
}
}

75
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java

@ -15,12 +15,16 @@
*/
package org.thingsboard.rule.engine.metadata;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.Futures;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.User;
@ -29,7 +33,13 @@ import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.relation.RelationService;
import java.util.HashMap;
@ -37,11 +47,19 @@ import java.util.List;
import java.util.Map;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyCollection;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@RunWith(MockitoJUnitRunner.class)
public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest {
public class TbGetRelatedAttributeNodeTest extends TbAbstractAttributeNodeTest {
User user = new User();
Asset asset = new Asset();
Device device = new Device();
@ -69,7 +87,7 @@ public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest {
}
@Override
protected TbEntityGetAttrNode getEmptyNode() {
protected TbAbstractGetEntityAttrNode getEmptyNode() {
return new TbGetRelatedAttributeNode();
}
@ -90,6 +108,7 @@ public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest {
conf.put(keyAttrConf, valueAttrConf);
config.setAttrMapping(conf);
config.setTelemetry(isTelemetry);
config.setFetchTo(FetchTo.METADATA);
return config;
}
@ -98,6 +117,31 @@ public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest {
return customerId;
}
@Test
public void errorThrownIfFetchToIsNull() {
var node = new TbGetRelatedAttributeNode();
var config = new TbGetRelatedAttrNodeConfiguration().defaultConfiguration();
config.setFetchTo(null);
var nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
var exception = assertThrows(TbNodeException.class, () -> node.init(ctx, nodeConfiguration));
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be NULL!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfMsgDataIsNotAnObjectAndFetchToData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", new DeviceId(UUID.randomUUID()), new TbMsgMetaData(), "[]");
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctx, msg));
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfCannotLoadAttributes() {
entityRelation.setFrom(user.getId());
@ -130,6 +174,33 @@ public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest {
entityAttributeAddedInMetadata(customerId, "CUSTOMER");
}
@Test
public void customerAttributeAddedInData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
entityRelation.setFrom(customerId);
entityRelation.setTo(customerId);
when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation)));
msg = TbMsg.newMsg("CUSTOMER", customerId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
List<AttributeKvEntry> attributes = Lists.newArrayList(new BaseAttributeKvEntry(new StringDataEntry("temperature", "high"), 1L));
when(ctx.getAttributesService()).thenReturn(attributesService);
when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection()))
.thenReturn(Futures.immediateFuture(attributes));
node.onMsg(ctx, msg);
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctx, times(1)).tellSuccess(actualMessageCaptor.capture());
var expectedMsgData = "{\"answer\":\"high\"}";
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData);
}
@Test
public void usersCustomerAttributesFetched() {
entityRelation.setFrom(user.getId());

77
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java

@ -15,10 +15,15 @@
*/
package org.thingsboard.rule.engine.metadata;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.Futures;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.ArgumentCaptor;
import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.User;
@ -26,15 +31,31 @@ import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.List;
import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyCollection;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@RunWith(MockitoJUnitRunner.class)
public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
public class TbGetTenantAttributeNodeTest extends TbAbstractAttributeNodeTest {
User user = new User();
Asset asset = new Asset();
Device device = new Device();
@ -56,7 +77,7 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
}
@Override
protected TbEntityGetAttrNode getEmptyNode() {
protected TbAbstractGetEntityAttrNode getEmptyNode() {
return new TbGetTenantAttributeNode();
}
@ -65,6 +86,31 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
return tenantId;
}
@Test
public void errorThrownIfFetchToIsNull() {
var node = new TbGetTenantAttributeNode();
var config = new TbGetEntityAttrNodeConfiguration().defaultConfiguration();
config.setFetchTo(null);
var nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
var exception = assertThrows(TbNodeException.class, () -> node.init(ctx, nodeConfiguration));
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be NULL!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfMsgDataIsNotAnObjectAndFetchToData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", new TenantId(UUID.randomUUID()), new TbMsgMetaData(), "[]");
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctx, msg));
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!");
verify(ctx, never()).tellSuccess(any());
}
@Test
public void errorThrownIfCannotLoadAttributes() {
errorThrownIfCannotLoadAttributes(user);
@ -86,6 +132,29 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest {
entityAttributeAddedInMetadata(tenantId, "TENANT");
}
@Test
public void customerAttributeAddedInData() {
node.fetchTo = FetchTo.DATA;
node.config.setFetchTo(FetchTo.DATA);
msg = TbMsg.newMsg("TENANT", tenantId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId);
List<AttributeKvEntry> attributes = Lists.newArrayList(new BaseAttributeKvEntry(new StringDataEntry("temperature", "high"), 1L));
when(ctx.getAttributesService()).thenReturn(attributesService);
when(attributesService.find(any(), eq(tenantId), eq(SERVER_SCOPE), anyCollection()))
.thenReturn(Futures.immediateFuture(attributes));
node.onMsg(ctx, msg);
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctx, times(1)).tellSuccess(actualMessageCaptor.capture());
var expectedMsgData = "{\"answer\":\"high\"}";
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData);
}
@Test
public void usersCustomerAttributesFetched() {
usersCustomerAttributesFetched(user);

246
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoaderTest.java

@ -0,0 +1,246 @@
/**
* 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 com.google.common.util.concurrent.Futures;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.rule.engine.api.RuleEngineAlarmService;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.BaseData;
import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityFieldsData;
import org.thingsboard.server.common.data.EntityType;
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.alarm.Alarm;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.AlarmId;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.id.EntityViewId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.rule.RuleChain;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.customer.CustomerService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.entityview.EntityViewService;
import org.thingsboard.server.dao.rule.RuleChainService;
import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.dao.user.UserService;
import java.util.EnumSet;
import java.util.UUID;
import java.util.concurrent.ExecutionException;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
public class EntitiesFieldsAsyncLoaderTest {
private static EnumSet<EntityType> SUPPORTED_ENTITY_TYPES;
private static UUID RANDOM_UUID;
private static TenantId TENANT_ID;
@Mock
private TbContext ctxMock;
@Mock
private TenantService tenantServiceMock;
@Mock
private CustomerService customerServiceMock;
@Mock
private UserService userServiceMock;
@Mock
private AssetService assetServiceMock;
@Mock
private DeviceService deviceServiceMock;
@Mock
private RuleEngineAlarmService ruleEngineAlarmServiceMock;
@Mock
private RuleChainService ruleChainServiceMock;
@Mock
private EntityViewService entityViewServiceMock;
@BeforeAll
public static void setup() {
RANDOM_UUID = UUID.randomUUID();
TENANT_ID = new TenantId(UUID.randomUUID());
SUPPORTED_ENTITY_TYPES = EnumSet.of(
EntityType.TENANT,
EntityType.CUSTOMER,
EntityType.USER,
EntityType.ASSET,
EntityType.DEVICE,
EntityType.ALARM,
EntityType.RULE_CHAIN,
EntityType.ENTITY_VIEW
);
}
@Test
public void givenSupportedEntityTypes_whenFindAsync_thenOK() throws ExecutionException, InterruptedException {
for (var entityType : SUPPORTED_ENTITY_TYPES) {
var entityId = EntityIdFactory.getByTypeAndUuid(entityType, RANDOM_UUID);
initMocks(entityType, false);
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
var actualEntityFieldsData = EntitiesFieldsAsyncLoader.findAsync(ctxMock, entityId).get();
var expectedEntityFieldsData = new EntityFieldsData(getEntityFromEntityId(entityId));
Assertions.assertEquals(expectedEntityFieldsData, actualEntityFieldsData);
}
}
@Test
public void givenUnsupportedEntityTypes_whenFindAsync_thenException() {
for (var entityType : EntityType.values()) {
if (!SUPPORTED_ENTITY_TYPES.contains(entityType)) {
var entityId = EntityIdFactory.getByTypeAndUuid(entityType, RANDOM_UUID);
var expectedExceptionMsg = "org.thingsboard.rule.engine.api.TbNodeException: Unexpected originator EntityType: " + entityType;
var exception = assertThrows(ExecutionException.class,
() -> EntitiesFieldsAsyncLoader.findAsync(ctxMock, entityId).get());
assertInstanceOf(TbNodeException.class, exception.getCause());
assertThat(exception.getMessage()).isEqualTo(expectedExceptionMsg);
}
}
}
@Test
public void givenSupportedTypeButEntityDoesNotExist_whenFindAsync_thenException() {
for (var entityType : SUPPORTED_ENTITY_TYPES) {
var entityId = EntityIdFactory.getByTypeAndUuid(entityType, RANDOM_UUID);
initMocks(entityType, true);
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
var expectedExceptionMsg = "org.thingsboard.rule.engine.api.TbNodeException: Entity not found!";
var exception = assertThrows(ExecutionException.class,
() -> EntitiesFieldsAsyncLoader.findAsync(ctxMock, entityId).get());
assertInstanceOf(TbNodeException.class, exception.getCause());
assertThat(exception.getMessage()).isEqualTo(expectedExceptionMsg);
}
}
private void initMocks(EntityType entityType, boolean entityDoesNotExist) {
switch (entityType) {
case TENANT:
var tenant = Futures.immediateFuture(entityDoesNotExist ? null : new Tenant(new TenantId(RANDOM_UUID)));
when(ctxMock.getTenantService()).thenReturn(tenantServiceMock);
doReturn(tenant).when(tenantServiceMock).findTenantByIdAsync(eq(TENANT_ID), any());
break;
case CUSTOMER:
var customer = Futures.immediateFuture(entityDoesNotExist ? null : new Customer(new CustomerId(RANDOM_UUID)));
when(ctxMock.getCustomerService()).thenReturn(customerServiceMock);
doReturn(customer).when(customerServiceMock).findCustomerByIdAsync(eq(TENANT_ID), any());
break;
case USER:
var user = Futures.immediateFuture(entityDoesNotExist ? null : new User(new UserId(RANDOM_UUID)));
when(ctxMock.getUserService()).thenReturn(userServiceMock);
doReturn(user).when(userServiceMock).findUserByIdAsync(eq(TENANT_ID), any());
break;
case ASSET:
var asset = Futures.immediateFuture(entityDoesNotExist ? null : new Asset(new AssetId(RANDOM_UUID)));
when(ctxMock.getAssetService()).thenReturn(assetServiceMock);
doReturn(asset).when(assetServiceMock).findAssetByIdAsync(eq(TENANT_ID), any());
break;
case DEVICE:
var device = Futures.immediateFuture(entityDoesNotExist ? null : new Device(new DeviceId(RANDOM_UUID)));
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock);
doReturn(device).when(deviceServiceMock).findDeviceByIdAsync(eq(TENANT_ID), any());
break;
case ALARM:
var alarm = Futures.immediateFuture(entityDoesNotExist ? null : new Alarm(new AlarmId(RANDOM_UUID)));
when(ctxMock.getAlarmService()).thenReturn(ruleEngineAlarmServiceMock);
doReturn(alarm).when(ruleEngineAlarmServiceMock).findAlarmByIdAsync(eq(TENANT_ID), any());
break;
case RULE_CHAIN:
var ruleChain = Futures.immediateFuture(entityDoesNotExist ? null : new RuleChain(new RuleChainId(RANDOM_UUID)));
when(ctxMock.getRuleChainService()).thenReturn(ruleChainServiceMock);
doReturn(ruleChain).when(ruleChainServiceMock).findRuleChainByIdAsync(eq(TENANT_ID), any());
break;
case ENTITY_VIEW:
var entityView = Futures.immediateFuture(entityDoesNotExist ? null : new EntityView(new EntityViewId(RANDOM_UUID)));
when(ctxMock.getEntityViewService()).thenReturn(entityViewServiceMock);
doReturn(entityView).when(entityViewServiceMock).findEntityViewByIdAsync(eq(TENANT_ID), any());
break;
default:
throw new RuntimeException("Unexpected EntityType: " + entityType);
}
}
private BaseData<? extends UUIDBased> getEntityFromEntityId(EntityId entityId) {
switch (entityId.getEntityType()) {
case TENANT:
return new Tenant((TenantId) entityId);
case CUSTOMER:
return new Customer((CustomerId) entityId);
case USER:
return new User((UserId) entityId);
case ASSET:
return new Asset((AssetId) entityId);
case DEVICE:
return new Device((DeviceId) entityId);
case ALARM:
return new Alarm((AlarmId) entityId);
case RULE_CHAIN:
return new RuleChain((RuleChainId) entityId);
case ENTITY_VIEW:
return new EntityView((EntityViewId) entityId);
default:
throw new RuntimeException("Unexpected EntityType: " + entityId.getEntityType());
}
}
}

7
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java

@ -61,7 +61,6 @@ import org.thingsboard.server.common.data.rule.RuleChain;
import org.thingsboard.server.common.data.rule.RuleNode;
import org.thingsboard.server.common.data.widget.WidgetType;
import org.thingsboard.server.common.data.widget.WidgetsBundle;
import org.thingsboard.server.dao.alarm.AlarmCommentService;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.customer.CustomerService;
import org.thingsboard.server.dao.dashboard.DashboardService;
@ -104,8 +103,6 @@ public class TenantIdLoaderTest {
@Mock
private RuleEngineAlarmService alarmService;
@Mock
private AlarmCommentService alarmCommentService;
@Mock
private RuleChainService ruleChainService;
@Mock
private EntityViewService entityViewService;
@ -359,9 +356,8 @@ public class TenantIdLoaderTest {
doReturn(notificationRule).when(notificationRuleService).findNotificationRuleById(eq(tenantId), any());
break;
default:
throw new RuntimeException("Unexpected original EntityType " + entityType);
throw new RuntimeException("Unexpected originator EntityType " + entityType);
}
}
private EntityId getEntityId(EntityType entityType) {
@ -397,5 +393,4 @@ public class TenantIdLoaderTest {
public void test_findEntityIdAsync_other_tenant() {
checkTenant(new TenantId(UUID.randomUUID()), false);
}
}

Loading…
Cancel
Save