diff --git a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java index a07bddc676..acade3afba 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java @@ -217,9 +217,9 @@ public class DefaultDataUpdateService implements DataUpdateService { 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 ..."); + log.info("Starting enrichment rule nodes update ..."); updateEnrichmentRuleNodes(); - log.info("Finished enrichment rule nodes update ..."); + log.info("Finished enrichment rule nodes update!"); break; default: throw new RuntimeException("Unable to update data, unsupported fromVersion: " + fromVersion); @@ -247,12 +247,12 @@ public class DefaultDataUpdateService implements DataUpdateService { 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!", + log.error("Failed to update rule node: [{}] with id: [{}] Node configuration is null! Skipping this node!", ruleNodeType, ruleNode.getId()); continue; } if (!configuration.isObject()) { - log.error("Unable to update [{}] rule node with ID [{}]! Node configuration is not an object! Skipping this node!", + log.error("Failed to update rule node: [{}] with id: [{}] Node configuration is not an object! Skipping this node!", ruleNodeType, ruleNode.getId()); continue; } @@ -265,8 +265,10 @@ public class DefaultDataUpdateService implements DataUpdateService { } 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()); + log.error("Failed to updated rule node: [{}] with id: [{}] " + + "Reason: fetchToMetadata property has unexpected value: {} Allowed values: true or false!", + ruleNodeType, ruleNode.getId(), fetchToMetadata); + continue; } configObjectNode.remove("fetchToMetadata"); } @@ -277,8 +279,10 @@ public class DefaultDataUpdateService implements DataUpdateService { } 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()); + log.error("Failed to updated rule node: [{}] with id: [{}] " + + "Reason: fetchToData property has unexpected value: {} Allowed values: true or false!", + ruleNodeType, ruleNode.getId(), fetchToData); + continue; } configObjectNode.remove("fetchToData"); } @@ -289,16 +293,28 @@ public class DefaultDataUpdateService implements DataUpdateService { } 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()); + log.error("Failed to updated rule node: [{}] with id: [{}] " + + "Reason: addToMetadata property has unexpected value: {} Allowed values: true or false!", + ruleNodeType, ruleNode.getId(), addToMetadata); + continue; } configObjectNode.remove("addToMetadata"); } - configObjectNode.put("fetchTo", fetchTo.toString()); + configObjectNode.put("fetchTo", fetchTo.name()); ruleNode.setConfiguration(configObjectNode); - ruleChainIdToTenantId.computeIfAbsent(ruleNode.getRuleChainId(), - ruleChainId -> ruleChainService.findRuleChainById(TenantId.SYS_TENANT_ID, ruleNode.getRuleChainId()).getTenantId()); - ruleChainService.saveRuleNode(ruleChainIdToTenantId.get(ruleNode.getRuleChainId()), ruleNode); + RuleChainId ruleChainId = ruleNode.getRuleChainId(); + TenantId tenantId = ruleChainIdToTenantId.computeIfAbsent(ruleChainId, + id -> { + RuleChain ruleChain = ruleChainService.findRuleChainById(TenantId.SYS_TENANT_ID, id); + if (ruleChain == null) { + log.error("Failed to find rule chain by id: [{}], ruleNodeId: [{}]", ruleChainId, ruleNode.getId()); + return null; + } + return ruleChain.getTenantId(); + }); + if (tenantId != null) { + ruleChainService.saveRuleNode(tenantId, ruleNode); + } } }); } catch (Exception e) { diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeConfiguration.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeConfiguration.java index 8c11ddc110..12a25b517c 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeConfiguration.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeConfiguration.java @@ -16,5 +16,7 @@ package org.thingsboard.rule.engine.api; public interface NodeConfiguration { + T defaultConfiguration(); + } diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbNode.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbNode.java index a0fbdaf8e1..b857307247 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbNode.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbNode.java @@ -24,6 +24,7 @@ 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; @@ -33,4 +34,5 @@ public interface TbNode { default void onPartitionChangeMsg(TbContext ctx, PartitionChangeMsg msg) { } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java index 63ddf9be5f..14c27a7ac7 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java @@ -19,5 +19,7 @@ import lombok.Data; @Data public abstract class TbAbstractFetchToNodeConfiguration { + private FetchTo fetchTo; + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java index 0c40ace712..5cc2da4d7a 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java @@ -19,7 +19,9 @@ 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.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.TbNodeConfiguration; @@ -31,14 +33,13 @@ import org.thingsboard.server.common.data.kv.BasicTsKvEntry; 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.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; 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.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.stream.Collectors; @@ -48,6 +49,7 @@ 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; +@Slf4j public abstract class TbAbstractGetAttributesNode extends TbAbstractNodeWithFetchTo { private static final String VALUE = "value"; private static final String TS = "ts"; @@ -58,98 +60,81 @@ public abstract class TbAbstractGetAttributesNode safePutAttributes(ctx, msg, entityId), - t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); - } catch (Throwable th) { - ctx.tellFailure(msg, th); - } + ctx.checkTenantEntity(msg.getOriginator()); + var msgDataAsObjectNode = FetchTo.DATA.equals(fetchTo) ? getMsgDataAsObjectNode(msg) : null; + withCallback( + findEntityIdAsync(ctx, msg), + entityId -> safePutAttributes(ctx, msg, msgDataAsObjectNode, entityId), + t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } protected abstract ListenableFuture findEntityIdAsync(TbContext ctx, TbMsg msg); - private void safePutAttributes(TbContext ctx, TbMsg msg, T entityId) { - if (entityId == null || entityId.isNullUid()) { - ctx.tellFailure(msg, new NoSuchElementException("Did not find entity! Msg ID: " + msg.getId())); - return; - } - ObjectNode msgDataNode; - if (FetchTo.DATA.equals(fetchTo)) { - msgDataNode = getMsgDataAsObjectNode(msg); - } else { - msgDataNode = null; - } - var failuresMap = new ConcurrentHashMap>(); - ListenableFuture>>> allFutures = Futures.allAsList( - getLatestTelemetry(ctx, entityId, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresMap), - getAttrAsync(ctx, entityId, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg), failuresMap), - getAttrAsync(ctx, entityId, SHARED_SCOPE, TbNodeUtils.processPatterns(config.getSharedAttributeNames(), msg), failuresMap), - getAttrAsync(ctx, entityId, SERVER_SCOPE, TbNodeUtils.processPatterns(config.getServerAttributeNames(), msg), failuresMap) + private void safePutAttributes(TbContext ctx, TbMsg msg, ObjectNode msgDataNode, T entityId) { + Set>> failuresPairSet = ConcurrentHashMap.newKeySet(); + var getKvEntryPairFutures = Futures.allAsList( + getLatestTelemetry(ctx, entityId, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresPairSet), + getAttrAsync(ctx, entityId, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg), failuresPairSet), + getAttrAsync(ctx, entityId, SHARED_SCOPE, TbNodeUtils.processPatterns(config.getSharedAttributeNames(), msg), failuresPairSet), + getAttrAsync(ctx, entityId, SERVER_SCOPE, TbNodeUtils.processPatterns(config.getServerAttributeNames(), msg), failuresPairSet) ); - withCallback(allFutures, futuresList -> { + withCallback(getKvEntryPairFutures, futuresList -> { var msgMetaData = msg.getMetaData().copy(); - futuresList.stream().filter(Objects::nonNull).forEach(kvEntriesMap -> { - kvEntriesMap.forEach((keyScope, kvEntryList) -> { - var prefix = getPrefix(keyScope); - kvEntryList.forEach(kvEntry -> { - 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()); - } - }); + futuresList.stream().filter(Objects::nonNull).forEach(kvEntriesPair -> { + var keyScope = kvEntriesPair.getFirst(); + var kvEntryList = kvEntriesPair.getSecond(); + var prefix = getPrefix(keyScope); + kvEntryList.forEach(kvEntry -> { + String targetKey = prefix + kvEntry.getKey(); + enrichMessage(msgDataNode, msgMetaData, kvEntry, targetKey); }); }); - - 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()) { + TbMsg outMsg = transformMessage(msg, msgDataNode, msgMetaData); + if (failuresPairSet.isEmpty()) { ctx.tellSuccess(outMsg); } else { - ctx.tellFailure(outMsg, reportFailures(failuresMap)); + ctx.tellFailure(outMsg, reportFailures(failuresPairSet)); } }, t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } - private ListenableFuture>> getAttrAsync(TbContext ctx, EntityId entityId, String scope, List keys, ConcurrentHashMap> failuresMap) { + private ListenableFuture>> getAttrAsync( + TbContext ctx, + EntityId entityId, + String scope, + List keys, + Set>> failuresPairSet + ) { if (CollectionUtils.isEmpty(keys)) { return Futures.immediateFuture(null); } 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)); + List nonExistentKeys = getNonExistentKeys(attributeKvEntryList, keys); + failuresPairSet.add(new TbPair<>(scope, nonExistentKeys)); } - var mapAttributeKvEntry = new HashMap>(); - mapAttributeKvEntry.put(scope, attributeKvEntryList); - return mapAttributeKvEntry; + return new TbPair<>(scope, attributeKvEntryList); }, MoreExecutors.directExecutor()); } - private ListenableFuture>> getLatestTelemetry(TbContext ctx, EntityId entityId, List keys, ConcurrentHashMap> failuresMap) { + private ListenableFuture>> getLatestTelemetry(TbContext ctx, EntityId entityId, List keys, Set>> failuresPairSet) { if (CollectionUtils.isEmpty(keys)) { return Futures.immediateFuture(null); } ListenableFuture> latestTelemetryFutures = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, keys); return Futures.transform(latestTelemetryFutures, tsKvEntries -> { var listTsKvEntry = new ArrayList(); + var nonExistentKeys = new ArrayList(); tsKvEntries.forEach(tsKvEntry -> { if (tsKvEntry.getValue() == null) { if (isTellFailureIfAbsent) { - computeFailuresMap(LATEST_TS, failuresMap, tsKvEntry.getKey()); + nonExistentKeys.add(tsKvEntry.getKey()); } } else if (getLatestValueWithTs) { listTsKvEntry.add(getValueWithTs(tsKvEntry)); @@ -157,9 +142,10 @@ public abstract class TbAbstractGetAttributesNode>(); - mapTsKvEntry.put(LATEST_TS, listTsKvEntry); - return mapTsKvEntry; + if (isTellFailureIfAbsent && !nonExistentKeys.isEmpty()) { + failuresPairSet.add(new TbPair<>(LATEST_TS, nonExistentKeys)); + } + return new TbPair<>(LATEST_TS, listTsKvEntry); }, MoreExecutors.directExecutor()); } @@ -187,31 +173,19 @@ public abstract class TbAbstractGetAttributesNode getNotExistingKeys(List existingAttributesKvEntry, List allKeys) { + private List getNonExistentKeys(List existingAttributesKvEntry, List allKeys) { List existingKeys = existingAttributesKvEntry.stream().map(KvEntry::getKey).collect(Collectors.toList()); return allKeys.stream().filter(key -> !existingKeys.contains(key)).collect(Collectors.toList()); } - private void computeFailuresMap(String scope, ConcurrentHashMap> failuresMap, String key) { - List failures = failuresMap.computeIfAbsent(scope, k -> new ArrayList<>()); - failures.add(key); - } - - private RuntimeException reportFailures(ConcurrentHashMap> failuresMap) { + private RuntimeException reportFailures(Set>> failuresPairSet) { 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"); - } - if (failuresMap.containsKey(SERVER_SCOPE)) { - errorMessage.append("\t").append("[" + SERVER_SCOPE + "]:").append(failuresMap.get(SERVER_SCOPE).toString()).append("\n"); - } - if (failuresMap.containsKey(SHARED_SCOPE)) { - errorMessage.append("\t").append("[" + SHARED_SCOPE + "]:").append(failuresMap.get(SHARED_SCOPE).toString()).append("\n"); - } - if (failuresMap.containsKey(LATEST_TS)) { - errorMessage.append("\t").append("[" + LATEST_TS + "]:").append(failuresMap.get(LATEST_TS).toString()).append("\n"); - } - failuresMap.clear(); + failuresPairSet.forEach(failurePair -> { + String scope = failurePair.getFirst(); + List nonExistentKeys = failurePair.getSecond(); + errorMessage.append("\t").append("[").append(scope).append("]:").append(nonExistentKeys.toString()).append("\n"); + }); + failuresPairSet.clear(); return new RuntimeException(errorMessage.toString()); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java index 874cee2e75..0a650bf3d9 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java @@ -20,8 +20,8 @@ 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.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.kv.KvEntry; @@ -30,7 +30,6 @@ 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; @@ -38,35 +37,31 @@ import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; @Slf4j public abstract class TbAbstractGetEntityAttrNode extends TbAbstractNodeWithFetchTo { + @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()); + var msgDataAsObjectNode = FetchTo.DATA.equals(fetchTo) ? getMsgDataAsObjectNode(msg) : null; withCallback(findEntityAsync(ctx, msg.getOriginator()), - entityId -> safeGetAttributes(ctx, msg, entityId, msgDataAsJsonNode), + entityId -> safeGetAttributes(ctx, msg, entityId, msgDataAsObjectNode), t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } protected abstract ListenableFuture 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; + protected void checkIfMappingIsNotEmptyOrThrow(TbGetEntityAttrNodeConfiguration config) throws TbNodeException { + if (config.getAttrMapping().isEmpty()) { + throw new TbNodeException("At least one attribute mapping should be specified!"); } + } - Map mappingsMap = new HashMap<>(); + private void safeGetAttributes(TbContext ctx, TbMsg msg, T entityId, ObjectNode msgDataAsJsonNode) { + var mappingsMap = new HashMap(); config.getAttrMapping().forEach((key, value) -> { String patternProcessedSourceKey = TbNodeUtils.processPattern(key, msg); String patternProcessedTargetKey = TbNodeUtils.processPattern(value, msg); mappingsMap.put(patternProcessedSourceKey, patternProcessedTargetKey); }); - var sourceKeys = List.copyOf(mappingsMap.keySet()); withCallback(config.isTelemetry() ? getLatestTelemetryAsync(ctx, entityId, sourceKeys) : getAttributesAsync(ctx, entityId, sourceKeys), data -> putDataAndTell(ctx, msg, data, mappingsMap, msgDataAsJsonNode), @@ -91,20 +86,13 @@ public abstract class TbAbstractGetEntityAttrNode extends Tb MoreExecutors.directExecutor()); } - private void putDataAndTell(TbContext ctx, TbMsg msg, List data, Map map, ObjectNode msgDataAsJsonNode) { + private void putDataAndTell(TbContext ctx, TbMsg msg, List data, Map map, ObjectNode msgData) { + var msgMetaData = msg.getMetaData().copy(); 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); + enrichMessage(msgData, msgMetaData, entry, targetKey); } + ctx.tellSuccess(transformMessage(msg, msgData, msgMetaData)); } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java index 4644db2ff7..0016234a91 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNode.java @@ -15,159 +15,131 @@ */ 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 com.google.gson.Gson; -import com.google.gson.JsonElement; -import com.google.gson.JsonObject; -import com.google.gson.JsonParser; -import com.google.gson.reflect.TypeToken; -import lombok.AllArgsConstructor; -import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.api.TbContext; -import org.thingsboard.rule.engine.util.EntityDetails; +import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.server.common.data.ContactBased; +import org.thingsboard.server.common.data.id.UUIDBased; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; -import java.lang.reflect.Type; -import java.util.Map; - import static org.thingsboard.common.util.DonAsynchron.withCallback; @Slf4j -public abstract class TbAbstractGetEntityDetailsNode extends TbAbstractNodeWithFetchTo { - private static final Gson gson = new Gson(); - private static final Type TYPE = new TypeToken>() { - }.getType(); +public abstract class TbAbstractGetEntityDetailsNode extends TbAbstractNodeWithFetchTo { @Override public void onMsg(TbContext ctx, TbMsg msg) { - withCallback(getDetails(ctx, msg), + ctx.checkTenantEntity(msg.getOriginator()); + var msgDataAsObjectNode = FetchTo.DATA.equals(fetchTo) ? getMsgDataAsObjectNode(msg) : null; + withCallback(getDetails(ctx, msg, msgDataAsObjectNode), ctx::tellSuccess, t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } - protected abstract ListenableFuture getDetails(TbContext ctx, TbMsg msg); + protected abstract String getPrefix(); - protected abstract ListenableFuture getContactBasedListenableFuture(TbContext ctx, TbMsg msg); + protected abstract ListenableFuture> getContactBasedFuture(TbContext ctx, TbMsg msg); - protected MessageData getDataAsJson(TbMsg msg) { - 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 { - throw new IllegalArgumentException("Unsupported fetchTo value!"); + protected void checkIfDetailsListIsNotEmptyOrThrow(C configuration) throws TbNodeException { + if (configuration.getDetailsList().isEmpty()) { + throw new TbNodeException("No entity details selected!"); } } - protected ListenableFuture getTbMsgListenableFuture(TbContext ctx, TbMsg msg, MessageData messageData, String prefix) { - if (config.getDetailsList().isEmpty()) { - return Futures.immediateFuture(msg); - } else { - ListenableFuture contactBasedListenableFuture = getContactBasedListenableFuture(ctx, msg); - ListenableFuture resultObject = addContactProperties(messageData.getData(), contactBasedListenableFuture, prefix); - return transformMsg(ctx, msg, resultObject, messageData); - } - } - - private ListenableFuture transformMsg(TbContext ctx, TbMsg msg, ListenableFuture propertiesFuture, MessageData messageData) { - return Futures.transformAsync(propertiesFuture, jsonElement -> { - if (jsonElement == null) { - return Futures.immediateFuture(null); - } else if (messageData.getDataSource().equals(DataSource.METADATA)) { - Map metadataMap = gson.fromJson(jsonElement.toString(), TYPE); - return Futures.immediateFuture(ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), new TbMsgMetaData(metadataMap), msg.getData())); - } else { - return Futures.immediateFuture(ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), gson.toJson(jsonElement))); - } - }, MoreExecutors.directExecutor()); - } - - private ListenableFuture addContactProperties(JsonElement data, ListenableFuture entityFuture, String prefix) { - return Futures.transformAsync(entityFuture, contactBased -> { + private ListenableFuture getDetails(TbContext ctx, TbMsg msg, ObjectNode messageData) { + ListenableFuture> contactBasedFuture = getContactBasedFuture(ctx, msg); + return Futures.transformAsync(contactBasedFuture, contactBased -> { if (contactBased == null) { - return Futures.immediateFuture(null); - } else { - JsonElement jsonElement = null; - for (EntityDetails entityDetails : this.config.getDetailsList()) { - jsonElement = setProperties(contactBased, data, entityDetails, prefix); - } - return Futures.immediateFuture(jsonElement); + return Futures.immediateFuture(msg); } + var msgMetaData = msg.getMetaData().copy(); + setProperties(contactBased, messageData, msgMetaData); + return Futures.immediateFuture(transformMessage(msg, messageData, msgMetaData)); }, MoreExecutors.directExecutor()); } - private JsonElement setProperties(ContactBased entity, JsonElement data, EntityDetails entityDetails, String prefix) { - JsonObject dataAsObject = data.getAsJsonObject(); - switch (entityDetails) { - case ID: - dataAsObject.addProperty(prefix + "id", entity.getId().toString()); - break; - case TITLE: - dataAsObject.addProperty(prefix + "title", entity.getName()); - break; - case ADDRESS: - if (entity.getAddress() != null) { - dataAsObject.addProperty(prefix + "address", entity.getAddress()); - } - break; - case ADDRESS2: - if (entity.getAddress2() != null) { - dataAsObject.addProperty(prefix + "address2", entity.getAddress2()); - } - break; - case CITY: - if (entity.getCity() != null) dataAsObject.addProperty(prefix + "city", entity.getCity()); - break; - case COUNTRY: - if (entity.getCountry() != null) - dataAsObject.addProperty(prefix + "country", entity.getCountry()); - break; - case STATE: - if (entity.getState() != null) { - dataAsObject.addProperty(prefix + "state", entity.getState()); - } - break; - case EMAIL: - if (entity.getEmail() != null) { - dataAsObject.addProperty(prefix + "email", entity.getEmail()); - } - break; - case PHONE: - if (entity.getPhone() != null) { - dataAsObject.addProperty(prefix + "phone", entity.getPhone()); - } - break; - case ZIP: - if (entity.getZip() != null) { - dataAsObject.addProperty(prefix + "zip", entity.getZip()); - } - break; - case ADDITIONAL_INFO: - if (entity.getAdditionalInfo().hasNonNull("description")) { - dataAsObject.addProperty(prefix + "additionalInfo", entity.getAdditionalInfo().get("description").asText()); - } - break; + private void setProperties(ContactBased contactBased, ObjectNode messageData, TbMsgMetaData msgMetaData) { + String prefix = getPrefix(); + String property; + String value; + for (var entityDetails : config.getDetailsList()) { + switch (entityDetails) { + case ID: + property = prefix + "id"; + value = contactBased.getId().getId().toString(); + setDetail(property, value, messageData, msgMetaData); + break; + case TITLE: + property = prefix + "title"; + value = contactBased.getName(); + setDetail(property, value, messageData, msgMetaData); + break; + case ADDRESS: + property = prefix + "address"; + value = contactBased.getAddress(); + setDetail(property, value, messageData, msgMetaData); + break; + case ADDRESS2: + property = prefix + "address2"; + value = contactBased.getAddress2(); + setDetail(property, value, messageData, msgMetaData); + break; + case CITY: + property = prefix + "city"; + value = contactBased.getCity(); + setDetail(property, value, messageData, msgMetaData); + break; + case COUNTRY: + property = prefix + "country"; + value = contactBased.getCountry(); + setDetail(property, value, messageData, msgMetaData); + break; + case STATE: + property = prefix + "state"; + value = contactBased.getState(); + setDetail(property, value, messageData, msgMetaData); + break; + case EMAIL: + property = prefix + "email"; + value = contactBased.getEmail(); + setDetail(property, value, messageData, msgMetaData); + break; + case PHONE: + property = prefix + "phone"; + value = contactBased.getPhone(); + setDetail(property, value, messageData, msgMetaData); + break; + case ZIP: + property = prefix + "zip"; + value = contactBased.getZip(); + setDetail(property, value, messageData, msgMetaData); + break; + case ADDITIONAL_INFO: + if (contactBased.getAdditionalInfo().hasNonNull("description")) { + property = prefix + "additionalInfo"; + value = contactBased.getAdditionalInfo().get("description").asText(); + setDetail(property, value, messageData, msgMetaData); + } + break; + } } - return dataAsObject; } - @Data - @AllArgsConstructor - private static class MessageData { - private JsonElement data; - private DataSource dataSource; + private void setDetail(String property, String value, ObjectNode messageData, TbMsgMetaData msgMetaData) { + if (value == null) { + return; + } + if (FetchTo.METADATA.equals(fetchTo)) { + msgMetaData.putValue(property, value); + } + if (FetchTo.DATA.equals(fetchTo)) { + messageData.put(property, value); + } } - private enum DataSource { - DATA, METADATA - } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java index 1da798ae82..4bf150ebff 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDetailsNodeConfiguration.java @@ -24,5 +24,7 @@ import java.util.List; @Data @EqualsAndHashCode(callSuper = true) public abstract class TbAbstractGetEntityDetailsNodeConfiguration extends TbAbstractFetchToNodeConfiguration { + private List detailsList; + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java index 3cf6300b77..a9c204215c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java @@ -17,14 +17,24 @@ package org.thingsboard.rule.engine.metadata; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.util.concurrent.AsyncFunction; +import com.google.common.util.concurrent.Futures; +import lombok.extern.slf4j.Slf4j; 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.data.id.EntityId; +import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgMetaData; +import java.util.NoSuchElementException; + +@Slf4j public abstract class TbAbstractNodeWithFetchTo implements TbNode { + protected C config; protected FetchTo fetchTo; @@ -32,7 +42,7 @@ public abstract class TbAbstractNodeWithFetchTo AsyncFunction checkIfEntityIsPresentOrThrow(String message) { + return id -> { + if (id == null || id.isNullUid()) { + return Futures.immediateFailedFuture(new NoSuchElementException(message)); + } + return Futures.immediateFuture(id); + }; + } + protected ObjectNode getMsgDataAsObjectNode(TbMsg msg) { JsonNode msgDataNode = JacksonUtil.toJsonNode(msg.getData()); - if (!msgDataNode.isObject()) { + if (msgDataNode == null || !msgDataNode.isObject()) { throw new IllegalArgumentException("Message body is not an object!"); } return (ObjectNode) msgDataNode; } + + protected void enrichMessage(ObjectNode msgData, TbMsgMetaData metaData, KvEntry kvEntry, String targetKey) { + if (FetchTo.DATA.equals(fetchTo)) { + JacksonUtil.addKvEntry(msgData, kvEntry, targetKey); + } else if (FetchTo.METADATA.equals(fetchTo)) { + metaData.putValue(targetKey, kvEntry.getValueAsString()); + } + } + + protected TbMsg transformMessage(TbMsg msg, ObjectNode msgDataNode, TbMsgMetaData msgMetaData) { + switch (fetchTo) { + case DATA: + return TbMsg.transformMsgData(msg, JacksonUtil.toString(msgDataNode)); + case METADATA: + return TbMsg.transformMsg(msg, msgMetaData); + default: + log.debug("Unexpected FetchTo value: {}. Allowed values: {}", fetchTo, FetchTo.values()); + return msg; + } + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java index 0a8ca6b255..f1dba47a90 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.RuleNode; @@ -43,6 +44,7 @@ import java.util.concurrent.ExecutionException; uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeFetchDeviceCredentialsConfig") public class TbFetchDeviceCredentialsNode extends TbAbstractNodeWithFetchTo { + private static final String CREDENTIALS = "credentials"; private static final String CREDENTIALS_TYPE = "credentialsType"; @@ -55,6 +57,7 @@ public class TbFetchDeviceCredentialsNode extends TbAbstractNodeWithFetchTo findEntityIdAsync(TbContext ctx, TbMsg msg) { - ctx.checkTenantEntity(msg.getOriginator()); return Futures.immediateFuture(msg.getOriginator()); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java index 89c4eb8086..2450796b59 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNodeConfiguration.java @@ -28,6 +28,7 @@ import java.util.List; @Data @EqualsAndHashCode(callSuper = true) public class TbGetAttributesNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration { + private List clientAttributeNames; private List sharedAttributeNames; private List serverAttributeNames; @@ -49,4 +50,5 @@ public class TbGetAttributesNodeConfiguration extends TbAbstractFetchToNodeConfi configuration.setFetchTo(FetchTo.METADATA); return configuration; } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java index d00de4dd0e..5f5f73c01c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; @@ -30,22 +31,30 @@ import org.thingsboard.server.common.data.plugin.ComponentType; type = ComponentType.ENRICHMENT, name = "customer attributes", configClazz = TbGetEntityAttrNodeConfiguration.class, - 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. " + + nodeDescription = "Add Originators Customer Attributes or Latest Telemetry into Message or Metadata", + nodeDetails = "Enrich the Message or Metadata with the corresponding customer's latest attributes or telemetry value. " + "The customer is selected based on the originator of the message: device, asset, etc. " + "
" + "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 TbAbstractGetEntityAttrNode { + + private static final String CUSTOMER_NOT_FOUND_MESSAGE = "Failed to find customer for entity with id %s and type %s"; + @Override - protected ListenableFuture findEntityAsync(TbContext ctx, EntityId originator) { - ctx.checkTenantEntity(originator); - return EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, originator); + protected TbGetEntityAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { + var config = TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class); + checkIfMappingIsNotEmptyOrThrow(config); + return config; } @Override - protected TbGetEntityAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { - return TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class); + protected ListenableFuture findEntityAsync(TbContext ctx, EntityId originator) { + return Futures.transformAsync(EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, originator), + checkIfEntityIsPresentOrThrow(String.format(CUSTOMER_NOT_FOUND_MESSAGE, originator.getId(), originator.getEntityType().getNormalName())), + ctx.getDbCallbackExecutor() + ); } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java index 1d7c53e002..810a6c2cd3 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java @@ -29,6 +29,7 @@ import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.HasCustomerId; import org.thingsboard.server.common.data.HasName; 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.EdgeId; import org.thingsboard.server.common.data.id.EntityId; @@ -37,6 +38,8 @@ import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; +import java.util.NoSuchElementException; + @Slf4j @RuleNode(type = ComponentType.ENRICHMENT, name = "customer details", @@ -47,28 +50,24 @@ import org.thingsboard.server.common.msg.TbMsg; "If the originator of the message is not assigned to Customer, or originator type is not supported - Message will be forwarded to Failure chain, otherwise, Success chain will be used.", uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeEntityDetailsConfig") -public class TbGetCustomerDetailsNode extends TbAbstractGetEntityDetailsNode { +public class TbGetCustomerDetailsNode extends TbAbstractGetEntityDetailsNode { + private static final String CUSTOMER_PREFIX = "customer_"; @Override protected TbGetCustomerDetailsNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { - return TbNodeUtils.convert(configuration, TbGetCustomerDetailsNodeConfiguration.class); + var config = TbNodeUtils.convert(configuration, TbGetCustomerDetailsNodeConfiguration.class); + checkIfDetailsListIsNotEmptyOrThrow(config); + return config; } @Override - protected ListenableFuture getDetails(TbContext ctx, TbMsg msg) { - ctx.checkTenantEntity(msg.getOriginator()); - return getTbMsgListenableFuture(ctx, msg, getDataAsJson(msg), CUSTOMER_PREFIX); + protected String getPrefix() { + return CUSTOMER_PREFIX; } @Override - protected ListenableFuture getContactBasedListenableFuture(TbContext ctx, TbMsg msg) { - return Futures.transformAsync(getCustomer(ctx, msg), customer -> - customer == null ? Futures.immediateFuture(null) : Futures.immediateFuture(customer), - MoreExecutors.directExecutor()); - } - - private ListenableFuture getCustomer(TbContext ctx, TbMsg msg) { + protected ListenableFuture> getContactBasedFuture(TbContext ctx, TbMsg msg) { switch (msg.getOriginator().getEntityType()) { case DEVICE: return Futures.transformAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), new DeviceId(msg.getOriginator().getId())), @@ -86,7 +85,7 @@ public class TbGetCustomerDetailsNode extends TbAbstractGetEntityDetailsNode getCustomerFuture(ctx, edge, msg.getOriginator()), MoreExecutors.directExecutor()); default: - throw new RuntimeException("Entity with entityType '" + msg.getOriginator().getEntityType() + "' is not supported."); + return Futures.immediateFailedFuture(new NoSuchElementException("Entity with entityType '" + msg.getOriginator().getEntityType() + "' is not supported.")); } } @@ -96,7 +95,7 @@ public class TbGetCustomerDetailsNode extends TbAbstractGetEntityDetailsNode { + @Override public TbGetCustomerDetailsNodeConfiguration defaultConfiguration() { var configuration = new TbGetCustomerDetailsNodeConfiguration(); configuration.setDetailsList(Collections.emptyList()); - configuration.setFetchTo(FetchTo.METADATA); + configuration.setFetchTo(FetchTo.DATA); return configuration; } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java index 96230ae28f..fdbb412e6b 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.api.RuleNode; @@ -39,6 +40,9 @@ import org.thingsboard.server.common.msg.TbMsg; uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeDeviceAttributesConfig") public class TbGetDeviceAttrNode extends TbAbstractGetAttributesNode { + + private static final String RELATED_DEVICE_NOT_FOUND_MESSAGE = "Failed to find related device to message originator using relation query specified in the configuration!"; + @Override protected TbGetDeviceAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { return TbNodeUtils.convert(configuration, TbGetDeviceAttrNodeConfiguration.class); @@ -46,7 +50,10 @@ public class TbGetDeviceAttrNode extends TbAbstractGetAttributesNode findEntityIdAsync(TbContext ctx, TbMsg msg) { - ctx.checkTenantEntity(msg.getOriginator()); - return EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctx, msg.getOriginator(), config.getDeviceRelationsQuery()); + return Futures.transformAsync( + EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctx, msg.getOriginator(), config.getDeviceRelationsQuery()), + checkIfEntityIsPresentOrThrow(RELATED_DEVICE_NOT_FOUND_MESSAGE), + ctx.getDbCallbackExecutor()); } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java index 108a1f5017..60a2c9cdf7 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java @@ -26,6 +26,7 @@ import java.util.Collections; @Data @EqualsAndHashCode(callSuper = true) public class TbGetDeviceAttrNodeConfiguration extends TbGetAttributesNodeConfiguration { + private DeviceRelationsQuery deviceRelationsQuery; @Override @@ -49,4 +50,5 @@ public class TbGetDeviceAttrNodeConfiguration extends TbGetAttributesNodeConfigu return configuration; } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java index fccbdfb39f..18134a5040 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java @@ -25,17 +25,19 @@ import java.util.Map; @Data @EqualsAndHashCode(callSuper = true) public class TbGetEntityAttrNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration { + private Map attrMapping; - private boolean isTelemetry = false; + private boolean isTelemetry; @Override public TbGetEntityAttrNodeConfiguration defaultConfiguration() { var configuration = new TbGetEntityAttrNodeConfiguration(); var attrMapping = new HashMap(); - attrMapping.putIfAbsent("serialNumber", "sn"); + attrMapping.putIfAbsent("alarmThreshold", "threshold"); configuration.setAttrMapping(attrMapping); configuration.setTelemetry(false); configuration.setFetchTo(FetchTo.METADATA); return configuration; } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsConfiguration.java index 3c0498f6ba..e5d428100b 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsConfiguration.java @@ -25,6 +25,7 @@ import java.util.Map; @Data @EqualsAndHashCode(callSuper = true) public class TbGetOriginatorFieldsConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration { + private Map fieldsMapping; private boolean ignoreNullStrings; @@ -39,4 +40,5 @@ public class TbGetOriginatorFieldsConfiguration extends TbAbstractFetchToNodeCon configuration.setFetchTo(FetchTo.METADATA); return configuration; } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java index 7dee45aae0..f520425cbc 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java @@ -15,11 +15,9 @@ */ 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 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.TbNodeConfiguration; @@ -29,8 +27,8 @@ import org.thingsboard.rule.engine.util.EntitiesFieldsAsyncLoader; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgMetaData; -import java.util.Collections; import java.util.HashMap; import java.util.Map; @@ -48,60 +46,53 @@ import static org.thingsboard.common.util.DonAsynchron.withCallback; uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeOriginatorFieldsConfig") public class TbGetOriginatorFieldsNode extends TbAbstractNodeWithFetchTo { + @Override protected TbGetOriginatorFieldsConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { - return TbNodeUtils.convert(configuration, TbGetOriginatorFieldsConfiguration.class); + var getOriginatorFieldsConfiguration = TbNodeUtils.convert(configuration, TbGetOriginatorFieldsConfiguration.class); + if (config.getFieldsMapping().isEmpty()) { + throw new TbNodeException("At least one field mapping should be specified!"); + } + return getOriginatorFieldsConfiguration; } @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()); + var msgDataAsObjectNode = FetchTo.DATA.equals(fetchTo) ? getMsgDataAsObjectNode(msg) : null; withCallback(collectMappedEntityFieldsAsync(ctx, msg.getOriginator()), targetKeysToSourceValuesMap -> { + TbMsgMetaData msgMetaData = msg.getMetaData().copy(); for (var entry : targetKeysToSourceValuesMap.entrySet()) { var targetKeyName = entry.getKey(); var sourceFieldValue = entry.getValue(); if (FetchTo.DATA.equals(fetchTo)) { - msgDataAsJsonNode.put(targetKeyName, sourceFieldValue); + msgDataAsObjectNode.put(targetKeyName, sourceFieldValue); } else if (FetchTo.METADATA.equals(fetchTo)) { - msg.getMetaData().putValue(targetKeyName, sourceFieldValue); + msgMetaData.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); - } + TbMsg outMsg = transformMessage(msg, msgDataAsObjectNode, msgMetaData); + ctx.tellSuccess(outMsg); }, t -> ctx.tellFailure(msg, t), MoreExecutors.directExecutor()); } private ListenableFuture> collectMappedEntityFieldsAsync(TbContext ctx, EntityId entityId) { - if (config.getFieldsMapping().isEmpty()) { - return Futures.immediateFuture(Collections.emptyMap()); - } else { - return Futures.transform(EntitiesFieldsAsyncLoader.findAsync(ctx, entityId), - fieldsData -> { - var targetKeysToSourceValuesMap = new HashMap(); - 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 Futures.transform(EntitiesFieldsAsyncLoader.findAsync(ctx, entityId), + fieldsData -> { + var targetKeysToSourceValuesMap = new HashMap(); + 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 targetKeysToSourceValuesMap; - }, ctx.getDbCallbackExecutor() - ); - } + } + return targetKeysToSourceValuesMap; + }, ctx.getDbCallbackExecutor() + ); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java index 9387858811..26491509b8 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java @@ -28,6 +28,7 @@ import java.util.HashMap; @Data @EqualsAndHashCode(callSuper = true) public class TbGetRelatedAttrNodeConfiguration extends TbGetEntityAttrNodeConfiguration { + private RelationsQuery relationsQuery; @Override @@ -37,14 +38,16 @@ public class TbGetRelatedAttrNodeConfiguration extends TbGetEntityAttrNodeConfig attrMapping.putIfAbsent("serialNumber", "sn"); configuration.setAttrMapping(attrMapping); configuration.setTelemetry(false); + configuration.setFetchTo(FetchTo.METADATA); var relationsQuery = new RelationsQuery(); + var relationEntityTypeFilter = new RelationEntityTypeFilter(EntityRelation.CONTAINS_TYPE, Collections.emptyList()); relationsQuery.setDirection(EntitySearchDirection.FROM); relationsQuery.setMaxLevel(1); - var relationEntityTypeFilter = new RelationEntityTypeFilter(EntityRelation.CONTAINS_TYPE, Collections.emptyList()); relationsQuery.setFilters(Collections.singletonList(relationEntityTypeFilter)); configuration.setRelationsQuery(relationsQuery); return configuration; } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java index d722d1c62b..b535d050f4 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; @@ -39,15 +40,23 @@ import org.thingsboard.server.common.data.plugin.ComponentType; uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeRelatedAttributesConfig") public class TbGetRelatedAttributeNode extends TbAbstractGetEntityAttrNode { + + private static final String RELATED_ENTITY_NOT_FOUND_MESSAGE = "Failed to find related entity to message originator using relation query specified in the configuration!"; + @Override public TbGetRelatedAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { - return TbNodeUtils.convert(configuration, TbGetRelatedAttrNodeConfiguration.class); + var config = TbNodeUtils.convert(configuration, TbGetRelatedAttrNodeConfiguration.class); + checkIfMappingIsNotEmptyOrThrow(config); + return config; } @Override public ListenableFuture findEntityAsync(TbContext ctx, EntityId originator) { - ctx.checkTenantEntity(originator); var relatedAttrConfig = (TbGetRelatedAttrNodeConfiguration) config; - return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, originator, relatedAttrConfig.getRelationsQuery()); + return Futures.transformAsync( + EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, originator, relatedAttrConfig.getRelationsQuery()), + checkIfEntityIsPresentOrThrow(RELATED_ENTITY_NOT_FOUND_MESSAGE), + ctx.getDbCallbackExecutor()); } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java index dcc96e82c1..d729c6159f 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNode.java @@ -40,14 +40,17 @@ import org.thingsboard.server.common.data.plugin.ComponentType; uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeTenantAttributesConfig") public class TbGetTenantAttributeNode extends TbAbstractGetEntityAttrNode { + @Override - public ListenableFuture findEntityAsync(TbContext ctx, EntityId originator) { - ctx.checkTenantEntity(originator); - return Futures.immediateFuture(ctx.getTenantId()); + public TbGetEntityAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { + var config = TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class); + checkIfMappingIsNotEmptyOrThrow(config); + return config; } @Override - public TbGetEntityAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { - return TbNodeUtils.convert(configuration, TbGetEntityAttrNodeConfiguration.class); + public ListenableFuture findEntityAsync(TbContext ctx, EntityId originator) { + return Futures.immediateFuture(ctx.getTenantId()); } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java index c89259d9ad..645106ab10 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNode.java @@ -15,9 +15,7 @@ */ 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.RuleNode; import org.thingsboard.rule.engine.api.TbContext; @@ -25,6 +23,7 @@ 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.ContactBased; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; @@ -38,25 +37,25 @@ import org.thingsboard.server.common.msg.TbMsg; "If the originator of the message is not assigned to Tenant, or originator type is not supported - Message will be forwarded to Failure chain, otherwise, Success chain will be used.", uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeEntityDetailsConfig") -public class TbGetTenantDetailsNode extends TbAbstractGetEntityDetailsNode { +public class TbGetTenantDetailsNode extends TbAbstractGetEntityDetailsNode { + private static final String TENANT_PREFIX = "tenant_"; @Override protected TbGetTenantDetailsNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { - return TbNodeUtils.convert(configuration, TbGetTenantDetailsNodeConfiguration.class); + var config = TbNodeUtils.convert(configuration, TbGetTenantDetailsNodeConfiguration.class); + checkIfDetailsListIsNotEmptyOrThrow(config); + return config; } @Override - protected ListenableFuture getDetails(TbContext ctx, TbMsg msg) { - ctx.checkTenantEntity(msg.getOriginator()); - return getTbMsgListenableFuture(ctx, msg, getDataAsJson(msg), TENANT_PREFIX); + protected String getPrefix() { + return TENANT_PREFIX; } @Override - protected ListenableFuture getContactBasedListenableFuture(TbContext ctx, TbMsg msg) { - ctx.checkTenantEntity(msg.getOriginator()); - return Futures.transformAsync(ctx.getTenantService().findTenantByIdAsync(ctx.getTenantId(), ctx.getTenantId()), tenant -> - tenant == null ? Futures.immediateFuture(null) : Futures.immediateFuture(tenant), - MoreExecutors.directExecutor()); + protected ListenableFuture> getContactBasedFuture(TbContext ctx, TbMsg msg) { + return ctx.getTenantService().findTenantByIdAsync(ctx.getTenantId(), ctx.getTenantId()); } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeConfiguration.java index e3608abbbd..c8d74c6170 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTenantDetailsNodeConfiguration.java @@ -24,11 +24,13 @@ import java.util.Collections; @Data @EqualsAndHashCode(callSuper = true) public class TbGetTenantDetailsNodeConfiguration extends TbAbstractGetEntityDetailsNodeConfiguration implements NodeConfiguration { + @Override public TbGetTenantDetailsNodeConfiguration defaultConfiguration() { TbGetTenantDetailsNodeConfiguration configuration = new TbGetTenantDetailsNodeConfiguration(); configuration.setDetailsList(Collections.emptyList()); - configuration.setFetchTo(FetchTo.METADATA); + configuration.setFetchTo(FetchTo.DATA); return configuration; } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java index bc693d9b97..b2a1599e8c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java @@ -36,12 +36,13 @@ import static org.thingsboard.common.util.DonAsynchron.withCallback; * Created by ashvayka on 19.01.18. */ @Slf4j -public abstract class TbAbstractTransformNode implements TbNode { - private TbTransformNodeConfiguration config; +public abstract class TbAbstractTransformNode implements TbNode { + + protected C config; @Override - public void init(TbContext context, TbNodeConfiguration configuration) throws TbNodeException { - this.config = TbNodeUtils.convert(configuration, TbTransformNodeConfiguration.class); + public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { + config = loadNodeConfiguration(ctx, configuration); } @Override @@ -52,18 +53,12 @@ public abstract class TbAbstractTransformNode implements TbNode { MoreExecutors.directExecutor()); } + protected abstract C loadNodeConfiguration(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException; + protected void transformFailure(TbContext ctx, TbMsg msg, Throwable t) { ctx.tellFailure(msg, t); } - protected void transformSuccess(TbContext ctx, TbMsg msg, TbMsg m) { - if (m != null) { - ctx.tellSuccess(m); - } else { - ctx.tellFailure(msg, new RuntimeException("Message is null!")); - } - } - protected void transformSuccess(TbContext ctx, TbMsg msg, List msgs) { if (msgs != null && !msgs.isEmpty()) { if (msgs.size() == 1) { @@ -89,7 +84,4 @@ public abstract class TbAbstractTransformNode implements TbNode { protected abstract ListenableFuture> transform(TbContext ctx, TbMsg msg); - public void setConfig(TbTransformNodeConfiguration config) { - this.config = config; - } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java index 911b811032..5cc5dfe23f 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java @@ -37,6 +37,7 @@ import org.thingsboard.server.common.msg.TbMsg; import java.util.Collections; import java.util.HashSet; import java.util.List; +import java.util.NoSuchElementException; @Slf4j @RuleNode( @@ -51,31 +52,36 @@ import java.util.List; configDirective = "tbTransformationNodeChangeOriginatorConfig", icon = "find_replace" ) -public class TbChangeOriginatorNode extends TbAbstractTransformNode { +public class TbChangeOriginatorNode extends TbAbstractTransformNode { - protected static final String CUSTOMER_SOURCE = "CUSTOMER"; - protected static final String TENANT_SOURCE = "TENANT"; - protected static final String RELATED_SOURCE = "RELATED"; - protected static final String ALARM_ORIGINATOR_SOURCE = "ALARM_ORIGINATOR"; - protected static final String ENTITY_SOURCE = "ENTITY"; - - private TbChangeOriginatorNodeConfiguration config; + private static final String CUSTOMER_SOURCE = "CUSTOMER"; + private static final String TENANT_SOURCE = "TENANT"; + private static final String RELATED_SOURCE = "RELATED"; + private static final String ALARM_ORIGINATOR_SOURCE = "ALARM_ORIGINATOR"; + private static final String ENTITY_SOURCE = "ENTITY"; @Override - public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { - this.config = TbNodeUtils.convert(configuration, TbChangeOriginatorNodeConfiguration.class); + protected TbChangeOriginatorNodeConfiguration loadNodeConfiguration(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { + var config = TbNodeUtils.convert(configuration, TbChangeOriginatorNodeConfiguration.class); validateConfig(config); - setConfig(config); + return config; } @Override protected ListenableFuture> transform(TbContext ctx, TbMsg msg) { - ListenableFuture newOriginator = getNewOriginator(ctx, msg); - return Futures.transform(newOriginator, n -> { - if (n == null || n.isNullUid()) { - return null; + ListenableFuture newOriginatorFuture = getNewOriginator(ctx, msg); + return Futures.transformAsync(newOriginatorFuture, newOriginator -> { + if (newOriginator == null || newOriginator.isNullUid()) { + return Futures.immediateFailedFuture(new NoSuchElementException("Failed to find new originator!")); } - return Collections.singletonList((ctx.transformMsg(msg, msg.getType(), n, msg.getMetaData(), msg.getData()))); + return Futures.immediateFuture( + Collections.singletonList( + ctx.transformMsg( + msg, + msg.getType(), + newOriginator, + msg.getMetaData(), + msg.getData()))); }, ctx.getDbCallbackExecutor()); } @@ -129,7 +135,6 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode { } EntitiesByNameAndTypeLoader.checkEntityType(EntityType.valueOf(conf.getEntityType())); } - } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeConfiguration.java index cc85933a48..473e83b91c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeConfiguration.java @@ -25,7 +25,9 @@ import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter; import java.util.Collections; @Data -public class TbChangeOriginatorNodeConfiguration extends TbTransformNodeConfiguration implements NodeConfiguration { +public class TbChangeOriginatorNodeConfiguration implements NodeConfiguration { + + private static final String CUSTOMER_SOURCE = "CUSTOMER"; private String originatorSource; @@ -36,7 +38,7 @@ public class TbChangeOriginatorNodeConfiguration extends TbTransformNodeConfigur @Override public TbChangeOriginatorNodeConfiguration defaultConfiguration() { TbChangeOriginatorNodeConfiguration configuration = new TbChangeOriginatorNodeConfiguration(); - configuration.setOriginatorSource(TbChangeOriginatorNode.CUSTOMER_SOURCE); + configuration.setOriginatorSource(CUSTOMER_SOURCE); RelationsQuery relationsQuery = new RelationsQuery(); relationsQuery.setDirection(EntitySearchDirection.FROM); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNode.java index 5c55cfb46d..ba78970c54 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNode.java @@ -43,17 +43,16 @@ import java.util.List; uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbTransformationNodeScriptConfig" ) -public class TbTransformMsgNode extends TbAbstractTransformNode { +public class TbTransformMsgNode extends TbAbstractTransformNode { - private TbTransformMsgNodeConfiguration config; private ScriptEngine scriptEngine; @Override - public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { - this.config = TbNodeUtils.convert(configuration, TbTransformMsgNodeConfiguration.class); + protected TbTransformMsgNodeConfiguration loadNodeConfiguration(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { + var config = TbNodeUtils.convert(configuration, TbTransformMsgNodeConfiguration.class); scriptEngine = ctx.createScriptEngine(config.getScriptLang(), ScriptLanguage.TBEL.equals(config.getScriptLang()) ? config.getTbelScript() : config.getJsScript()); - setConfig(config); + return config; } @Override @@ -62,12 +61,6 @@ public class TbTransformMsgNode extends TbAbstractTransformNode { return scriptEngine.executeUpdateAsync(msg); } - @Override - protected void transformSuccess(TbContext ctx, TbMsg msg, TbMsg m) { - ctx.logJsEvalResponse(); - super.transformSuccess(ctx, msg, m); - } - @Override protected void transformFailure(TbContext ctx, TbMsg msg, Throwable t) { ctx.logJsEvalFailure(); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNodeConfiguration.java index 2d4aeac161..34d465f49c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNodeConfiguration.java @@ -20,7 +20,7 @@ import org.thingsboard.rule.engine.api.NodeConfiguration; import org.thingsboard.server.common.data.script.ScriptLanguage; @Data -public class TbTransformMsgNodeConfiguration extends TbTransformNodeConfiguration implements NodeConfiguration { +public class TbTransformMsgNodeConfiguration implements NodeConfiguration { private ScriptLanguage scriptLang; private String jsScript; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformNodeConfiguration.java deleted file mode 100644 index 160bd6d8d1..0000000000 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformNodeConfiguration.java +++ /dev/null @@ -1,23 +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.transform; - -import lombok.Data; - -@Data -public class TbTransformNodeConfiguration { - -} diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java index d36162196e..52fb011614 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java @@ -29,9 +29,7 @@ import org.thingsboard.server.common.data.id.UserId; public class EntitiesCustomerIdAsyncLoader { - public static ListenableFuture findEntityIdAsync(TbContext ctx, EntityId original) { - switch (original.getEntityType()) { case CUSTOMER: return Futures.immediateFuture((CustomerId) original); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoader.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoader.java index 0d54e5ddc1..15dd7f85b1 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoader.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoader.java @@ -17,8 +17,6 @@ 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; @@ -34,36 +32,37 @@ 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.NoSuchElementException; import java.util.function.Function; -@Slf4j public class EntitiesFieldsAsyncLoader { + public static ListenableFuture findAsync(TbContext ctx, EntityId originatorId) { switch (originatorId.getEntityType()) { case TENANT: return toEntityFieldsDataAsync(ctx.getTenantService().findTenantByIdAsync(ctx.getTenantId(), (TenantId) originatorId), - EntityFieldsData::new); + EntityFieldsData::new, ctx); case CUSTOMER: return toEntityFieldsDataAsync(ctx.getCustomerService().findCustomerByIdAsync(ctx.getTenantId(), (CustomerId) originatorId), - EntityFieldsData::new); + EntityFieldsData::new, ctx); case USER: return toEntityFieldsDataAsync(ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), (UserId) originatorId), - EntityFieldsData::new); + EntityFieldsData::new, ctx); case ASSET: return toEntityFieldsDataAsync(ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), (AssetId) originatorId), - EntityFieldsData::new); + EntityFieldsData::new, ctx); case DEVICE: return toEntityFieldsDataAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), (DeviceId) originatorId), - EntityFieldsData::new); + EntityFieldsData::new, ctx); case ALARM: return toEntityFieldsDataAsync(ctx.getAlarmService().findAlarmByIdAsync(ctx.getTenantId(), (AlarmId) originatorId), - EntityFieldsData::new); + EntityFieldsData::new, ctx); case RULE_CHAIN: return toEntityFieldsDataAsync(ctx.getRuleChainService().findRuleChainByIdAsync(ctx.getTenantId(), (RuleChainId) originatorId), - EntityFieldsData::new); + EntityFieldsData::new, ctx); case ENTITY_VIEW: return toEntityFieldsDataAsync(ctx.getEntityViewService().findEntityViewByIdAsync(ctx.getTenantId(), (EntityViewId) originatorId), - EntityFieldsData::new); + EntityFieldsData::new, ctx); default: return Futures.immediateFailedFuture(new TbNodeException("Unexpected originator EntityType: " + originatorId.getEntityType())); } @@ -71,10 +70,12 @@ public class EntitiesFieldsAsyncLoader { private static > ListenableFuture toEntityFieldsDataAsync( ListenableFuture future, - Function converter + Function converter, + TbContext ctx ) { return Futures.transformAsync(future, in -> in != null ? Futures.immediateFuture(converter.apply(in)) - : Futures.immediateFailedFuture(new TbNodeException("Entity not found!")), MoreExecutors.directExecutor()); + : Futures.immediateFailedFuture(new NoSuchElementException("Entity not found!")), ctx.getDbCallbackExecutor()); } + } diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java index 1f7057050e..6af846d36d 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java @@ -53,6 +53,8 @@ import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; @RunWith(MockitoJUnitRunner.class) public class TbChangeOriginatorNodeTest { + private static final String CUSTOMER_SOURCE = "CUSTOMER"; + private TbChangeOriginatorNode node; @Mock @@ -158,7 +160,7 @@ public class TbChangeOriginatorNodeTest { public void init() throws TbNodeException { TbChangeOriginatorNodeConfiguration config = new TbChangeOriginatorNodeConfiguration(); - config.setOriginatorSource(TbChangeOriginatorNode.CUSTOMER_SOURCE); + config.setOriginatorSource(CUSTOMER_SOURCE); ObjectMapper mapper = new ObjectMapper(); TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config));