From f29e1f5feff55921339b6f0d34483919039830f3 Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Thu, 23 Mar 2023 12:44:23 +0200 Subject: [PATCH] Add fetch to logic, fix and add tests for some rule nodes --- .../rule/engine/api/NodeConfiguration.java | 2 - .../thingsboard/rule/engine/api/TbNode.java | 8 +- .../rule/engine/metadata/FetchTo.java | 21 + .../TbAbstractFetchToNodeConfiguration.java | 23 ++ .../metadata/TbAbstractGetAttributesNode.java | 46 +-- .../metadata/TbAbstractGetEntityAttrNode.java | 110 +++++ .../TbAbstractGetEntityDetailsNode.java | 25 +- ...ractGetEntityDetailsNodeConfiguration.java | 9 +- .../metadata/TbAbstractNodeWithFetchTo.java | 50 +++ .../engine/metadata/TbEntityGetAttrNode.java | 109 ----- .../TbFetchDeviceCredentialsNode.java | 36 +- ...tchDeviceCredentialsNodeConfiguration.java | 9 +- .../engine/metadata/TbGetAttributesNode.java | 15 +- .../TbGetAttributesNodeConfiguration.java | 9 +- .../metadata/TbGetCustomerAttributeNode.java | 17 +- .../metadata/TbGetCustomerDetailsNode.java | 5 +- ...TbGetCustomerDetailsNodeConfiguration.java | 5 +- .../engine/metadata/TbGetDeviceAttrNode.java | 5 +- .../TbGetDeviceAttrNodeConfiguration.java | 5 +- .../TbGetEntityAttrNodeConfiguration.java | 7 +- .../TbGetOriginatorFieldsConfiguration.java | 6 +- .../metadata/TbGetOriginatorFieldsNode.java | 80 ++-- .../TbGetRelatedAttrNodeConfiguration.java | 3 +- .../metadata/TbGetRelatedAttributeNode.java | 25 +- .../metadata/TbGetTenantAttributeNode.java | 18 +- .../metadata/TbGetTenantDetailsNode.java | 4 +- .../TbGetTenantDetailsNodeConfiguration.java | 5 +- .../util/EntitiesFieldsAsyncLoader.java | 42 +- ....java => TbAbstractAttributeNodeTest.java} | 20 +- .../TbAbstractGetAttributesNodeTest.java | 58 +-- .../TbFetchDeviceCredentialsNodeTest.java | 6 +- .../TbGetCustomerAttributeNodeTest.java | 81 +++- .../TbGetOriginatorFieldsNodeTest.java | 376 ++++++++++++++++++ .../TbGetRelatedAttributeNodeTest.java | 83 +++- .../TbGetTenantAttributeNodeTest.java | 85 +++- .../util/EntitiesFieldsAsyncLoaderTest.java | 231 +++++++++++ .../rule/engine/util/TenantIdLoaderTest.java | 7 +- 37 files changed, 1277 insertions(+), 369 deletions(-) create mode 100644 rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchTo.java create mode 100644 rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractFetchToNodeConfiguration.java create mode 100644 rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java create mode 100644 rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java delete mode 100644 rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java rename rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/{AbstractAttributeNodeTest.java => TbAbstractAttributeNodeTest.java} (96%) create mode 100644 rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNodeTest.java create mode 100644 rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoaderTest.java 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 12a25b517c..8c11ddc110 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,7 +16,5 @@ 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 18044eb05c..a0fbdaf8e1 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,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) { + } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchTo.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/FetchTo.java new file mode 100644 index 0000000000..dcd8c81477 --- /dev/null +++ b/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 +} 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 new file mode 100644 index 0000000000..63ddf9be5f --- /dev/null +++ b/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; +} 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 e5ef68e421..378e7d1358 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 @@ -22,10 +22,8 @@ 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; @@ -53,26 +51,19 @@ 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 implements TbNode { - +public abstract class TbAbstractGetAttributesNode extends TbAbstractNodeWithFetchTo { 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 { @@ -92,13 +83,9 @@ public abstract class TbAbstractGetAttributesNode { String key = prefix + kvEntry.getKey(); - if (fetchToData) { - JacksonUtil.addKvEntry((ObjectNode) msgDataNode, kvEntry, key); - } else { + 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 { @@ -175,7 +167,7 @@ public abstract class TbAbstractGetAttributesNode 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()); + withCallback(findEntityAsync(ctx, msg.getOriginator()), + entityId -> safeGetAttributes(ctx, msg, entityId, msgDataAsJsonNode), + 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.tellNext(msg, FAILURE); + return; + } + + Map 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> getAttributesAsync(TbContext ctx, EntityId entityId, List 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> getLatestTelemetryAsync(TbContext ctx, EntityId entityId, List 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 data, Map 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); + } + } +} 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 315da5f16a..21918029ef 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 @@ -27,11 +27,9 @@ 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.data.id.UUIDBased; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; @@ -41,20 +39,11 @@ import java.util.Map; import static org.thingsboard.common.util.DonAsynchron.withCallback; @Slf4j -public abstract class TbAbstractGetEntityDetailsNode implements TbNode { - +public abstract class TbAbstractGetEntityDetailsNode extends TbAbstractNodeWithFetchTo { private static final Gson gson = new Gson(); - private static final JsonParser jsonParser = new JsonParser(); private static final Type TYPE = new TypeToken>() { }.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 +51,20 @@ public abstract class TbAbstractGetEntityDetailsNode ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } - protected abstract C loadGetEntityDetailsNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException; - protected abstract ListenableFuture getDetails(TbContext ctx, TbMsg msg); protected abstract ListenableFuture getContactBasedListenableFuture(TbContext ctx, TbMsg msg); protected MessageData getDataAsJson(TbMsg msg) { - if (this.config.isAddToMetadata()) { + if (config.getFetchTo() == FetchTo.METADATA) { return new MessageData(gson.toJsonTree(msg.getMetaData().getData(), TYPE), DataSource.METADATA); } else { - return new MessageData(jsonParser.parse(msg.getData()), DataSource.DATA); + return new MessageData(JsonParser.parseString(msg.getData()), DataSource.DATA); } } protected ListenableFuture getTbMsgListenableFuture(TbContext ctx, TbMsg msg, MessageData messageData, String prefix) { - if (this.config.getDetailsList().isEmpty()) { + if (config.getDetailsList().isEmpty()) { return Futures.immediateFuture(msg); } else { ListenableFuture contactBasedListenableFuture = getContactBasedListenableFuture(ctx, msg); @@ -178,6 +165,4 @@ public abstract class TbAbstractGetEntityDetailsNode detailsList; - - private boolean addToMetadata; - } 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 new file mode 100644 index 0000000000..3cf6300b77 --- /dev/null +++ b/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 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; + } +} diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java deleted file mode 100644 index 8031939c7d..0000000000 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java +++ /dev/null @@ -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 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 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 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> getAttributesAsync(TbContext ctx, EntityId entityId, List attrKeys) { - ListenableFuture> 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> getLatestTelemetry(TbContext ctx, EntityId entityId, List timeseriesKeys) { - ListenableFuture> 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 attributes, Map map) { - attributes.forEach(r -> { - String attrName = map.get(r.getKey()); - msg.getMetaData().putValue(attrName, r.getValueAsString()); - }); - ctx.tellSuccess(msg); - } - - protected abstract ListenableFuture findEntityAsync(TbContext ctx, EntityId originator); - - public void setConfig(TbGetEntityAttrNodeConfiguration config) { - this.config = config; - } - -} 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 2aa3b7a837..abded26923 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,24 +15,19 @@ */ 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 +43,36 @@ import java.util.concurrent.ExecutionException; "- send Message via Failure chain, otherwise Success chain is used.", uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeFetchDeviceCredentialsConfig") -public class TbFetchDeviceCredentialsNode implements TbNode { - +public class TbFetchDeviceCredentialsNode extends TbAbstractNodeWithFetchTo { 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(); 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,7 +80,7 @@ public class TbFetchDeviceCredentialsNode implements TbNode { metaData.putValue(CREDENTIALS, JacksonUtil.toString(credentialsInfo)); } transformedMsg = TbMsg.transformMsg(msg, msg.getType(), originator, metaData, msg.getData()); - } else { + } else if (FetchTo.DATA.equals(fetchTo)) { ObjectNode data = (ObjectNode) JacksonUtil.toJsonNode(msg.getData()); data.put(CREDENTIALS_TYPE, credentialsType.name()); data.set(CREDENTIALS, credentialsInfo); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeConfiguration.java index 66fddb0770..5f8b5921f8 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeConfiguration.java +++ b/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 { - - private boolean fetchToMetadata; - +public class TbFetchDeviceCredentialsNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration { @Override public TbFetchDeviceCredentialsNodeConfiguration defaultConfiguration() { TbFetchDeviceCredentialsNodeConfiguration configuration = new TbFetchDeviceCredentialsNodeConfiguration(); - configuration.setFetchToMetadata(true); + configuration.setFetchTo(FetchTo.METADATA); return configuration; } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java index 7e42105e5d..0cd8b9cade 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java @@ -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, CLIENT/SHARED/SERVER 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, CLIENT/SHARED/SERVER attributes are added into Message data/metadata " + "with specific prefix: cs/shared/ss. 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 " + "metadata.cs_temperature or metadata.shared_limit ", uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeOriginatorAttributesConfig") public class TbGetAttributesNode extends TbAbstractGetAttributesNode { - @Override - protected TbGetAttributesNodeConfiguration loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException { + protected TbGetAttributesNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { return TbNodeUtils.convert(configuration, TbGetAttributesNodeConfiguration.class); } @Override protected ListenableFuture 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 e587fb9c7a..7c2b515668 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 @@ -16,8 +16,8 @@ package org.thingsboard.rule.engine.metadata; import lombok.Data; +import lombok.EqualsAndHashCode; import org.thingsboard.rule.engine.api.NodeConfiguration; - import java.util.Collections; import java.util.List; @@ -25,8 +25,8 @@ import java.util.List; * Created by ashvayka on 19.01.18. */ @Data -public class TbGetAttributesNodeConfiguration implements NodeConfiguration { - +@EqualsAndHashCode(callSuper = true) +public class TbGetAttributesNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration { private List clientAttributeNames; private List sharedAttributeNames; private List serverAttributeNames; @@ -35,7 +35,6 @@ public class TbGetAttributesNodeConfiguration implements NodeConfiguration" + "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 { - +public class TbGetCustomerAttributeNode extends TbAbstractGetEntityAttrNode { @Override protected ListenableFuture 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); + } } 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 3afc66c4dc..cdcf522876 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 @@ -33,6 +33,7 @@ import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityViewId; +import org.thingsboard.server.common.data.id.UUIDBased; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; @@ -48,16 +49,16 @@ import org.thingsboard.server.common.msg.TbMsg; uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeEntityDetailsConfig") public class TbGetCustomerDetailsNode extends TbAbstractGetEntityDetailsNode { - 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 getDetails(TbContext ctx, TbMsg msg) { + ctx.checkTenantEntity(msg.getOriginator()); return getTbMsgListenableFuture(ctx, msg, getDataAsJson(msg), CUSTOMER_PREFIX); } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeConfiguration.java index 0d74d49942..91587e18f2 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeConfiguration.java +++ b/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 { - - @Override public TbGetCustomerDetailsNodeConfiguration defaultConfiguration() { TbGetCustomerDetailsNodeConfiguration configuration = new TbGetCustomerDetailsNodeConfiguration(); configuration.setDetailsList(Collections.emptyList()); + configuration.setFetchTo(FetchTo.METADATA); 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 b0a221e46e..96230ae28f 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 @@ -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 { - @Override - protected TbGetDeviceAttrNodeConfiguration loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException { + protected TbGetDeviceAttrNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { return TbNodeUtils.convert(configuration, TbGetDeviceAttrNodeConfiguration.class); } @Override protected ListenableFuture findEntityIdAsync(TbContext ctx, TbMsg msg) { + ctx.checkTenantEntity(msg.getOriginator()); return EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctx, msg.getOriginator(), config.getDeviceRelationsQuery()); } - } 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 ac5b98134a..f3ba3651cc 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 @@ -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,8 +24,8 @@ 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 @@ -36,7 +37,7 @@ public class TbGetDeviceAttrNodeConfiguration extends TbGetAttributesNodeConfigu configuration.setLatestTsKeyNames(Collections.emptyList()); configuration.setTellFailureIfAbsent(true); configuration.setGetLatestValueWithTs(false); - configuration.setFetchToData(false); + configuration.setFetchTo(FetchTo.METADATA); DeviceRelationsQuery deviceRelationsQuery = new DeviceRelationsQuery(); deviceRelationsQuery.setDirection(EntitySearchDirection.FROM); 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 be2eab5bb0..12a94e227f 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 @@ -16,15 +16,15 @@ 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 { - +@EqualsAndHashCode(callSuper = true) +public class TbGetEntityAttrNodeConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration { private Map attrMapping; private boolean isTelemetry = false; @@ -35,6 +35,7 @@ public class TbGetEntityAttrNodeConfiguration implements NodeConfiguration { - +@EqualsAndHashCode(callSuper = true) +public class TbGetOriginatorFieldsConfiguration extends TbAbstractFetchToNodeConfiguration implements NodeConfiguration { private Map fieldsMapping; private boolean ignoreNullStrings; @@ -35,6 +36,7 @@ public class TbGetOriginatorFieldsConfiguration implements NodeConfiguration { @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 putEntityFields(TbContext ctx, EntityId entityId, TbMsg msg) { + private ListenableFuture> 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(); + 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() ); } } - } 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 489e04f60e..8df69c67af 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 @@ -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; @@ -26,8 +27,8 @@ import java.util.HashMap; import java.util.Map; @Data +@EqualsAndHashCode(callSuper = true) public class TbGetRelatedAttrNodeConfiguration extends TbGetEntityAttrNodeConfiguration { - private RelationsQuery relationsQuery; @Override 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 b7ef24f36b..d722d1c62b 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 @@ -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 " + "metadata.temperature.", uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeRelatedAttributesConfig") - -public class TbGetRelatedAttributeNode extends TbEntityGetAttrNode { - - private TbGetRelatedAttrNodeConfiguration config; - +public class TbGetRelatedAttributeNode extends TbAbstractGetEntityAttrNode { @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 findEntityAsync(TbContext ctx, EntityId originator) { - return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, originator, config.getRelationsQuery()); + public ListenableFuture findEntityAsync(TbContext ctx, EntityId originator) { + ctx.checkTenantEntity(originator); + var relatedAttrConfig = (TbGetRelatedAttrNodeConfiguration) config; + return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, originator, relatedAttrConfig.getRelationsQuery()); } } 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 ab422eece8..87b60dfdaf 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 @@ -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; @@ -29,19 +32,22 @@ import org.thingsboard.server.common.data.plugin.ComponentType; type = ComponentType.ENRICHMENT, 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 " + "metadata.temperature.", uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeTenantAttributesConfig") -public class TbGetTenantAttributeNode extends TbEntityGetAttrNode { - +public class TbGetTenantAttributeNode extends TbAbstractGetEntityAttrNode { @Override - protected ListenableFuture findEntityAsync(TbContext ctx, EntityId originator) { + public ListenableFuture 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); + } } 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 52c32a5a1e..983b17e61e 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 @@ -25,6 +25,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.UUIDBased; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; @@ -39,11 +40,10 @@ import org.thingsboard.server.common.msg.TbMsg; uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbEnrichmentNodeEntityDetailsConfig") public class TbGetTenantDetailsNode extends TbAbstractGetEntityDetailsNode { - 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); } 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 7770d6acb1..e3608abbbd 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 @@ -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 { - - @Override public TbGetTenantDetailsNodeConfiguration defaultConfiguration() { TbGetTenantDetailsNodeConfiguration configuration = new TbGetTenantDetailsNodeConfiguration(); configuration.setDetailsList(Collections.emptyList()); + configuration.setFetchTo(FetchTo.METADATA); return configuration; } } 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 cb9ee7c8b6..9d78005e23 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 @@ -1,12 +1,12 @@ /** * 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 - * + *

+ * 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. @@ -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 findAsync(TbContext ctx, EntityId original) { - switch (original.getEntityType()) { + public static ListenableFuture 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 ListenableFuture getAsync( - ListenableFuture future, Function converter) { + private static > ListenableFuture toEntityFieldsDataAsync( + ListenableFuture future, + Function 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()); } } diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/AbstractAttributeNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractAttributeNodeTest.java similarity index 96% rename from rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/AbstractAttributeNodeTest.java rename to rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractAttributeNodeTest.java index 45e4f1c27f..a3949587b9 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/AbstractAttributeNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractAttributeNodeTest.java @@ -1,12 +1,12 @@ /** * 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 - * + *

+ * 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. @@ -65,7 +65,7 @@ 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 +86,9 @@ public abstract class AbstractAttributeNodeTest { DeviceService deviceService; TbMsg msg; Map 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 +117,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); @@ -168,7 +167,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 +209,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(); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNodeTest.java index cfcf1c64a1..853818f226 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNodeTest.java +++ b/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 attributes) { + private void checkAttributes(TbMsg actualMsg, FetchTo fetchTo, String prefix, List 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)); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeTest.java index 1d291545ba..344ce20c20 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbFetchDeviceCredentialsNodeTest.java +++ b/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 diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java index 10b19add34..f18ac3299b 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java @@ -1,12 +1,12 @@ /** * 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 - * + *

+ * 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. @@ -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 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); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNodeTest.java new file mode 100644 index 0000000000..77c2cd3130 --- /dev/null +++ b/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 ListenableFuture executeAsync(Callable 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); + } +} \ No newline at end of file diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java index 8a41805d26..ef0e94d545 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java @@ -1,12 +1,12 @@ /** * 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 - * + *

+ * 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. @@ -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 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()); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java index 510e97dd9c..ce4dde791a 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java @@ -1,12 +1,12 @@ /** * 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 - * + *

+ * 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. @@ -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 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); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoaderTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoaderTest.java new file mode 100644 index 0000000000..3374e8e011 --- /dev/null +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoaderTest.java @@ -0,0 +1,231 @@ +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 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 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()); + } + } +} diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java index 6f937d7e18..4675784bc8 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java @@ -56,7 +56,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; @@ -94,8 +93,6 @@ public class TenantIdLoaderTest { @Mock private RuleEngineAlarmService alarmService; @Mock - private AlarmCommentService alarmCommentService; - @Mock private RuleChainService ruleChainService; @Mock private EntityViewService entityViewService; @@ -312,9 +309,8 @@ public class TenantIdLoaderTest { break; default: - throw new RuntimeException("Unexpected original EntityType " + entityType); + throw new RuntimeException("Unexpected originator EntityType " + entityType); } - } private EntityId getEntityId(EntityType entityType) { @@ -350,5 +346,4 @@ public class TenantIdLoaderTest { public void test_findEntityIdAsync_other_tenant() { checkTenant(new TenantId(UUID.randomUUID()), false); } - }