From 287c1fd48a882b9c5ec04ac1145d42b6f82db86a Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Wed, 5 Oct 2022 10:09:01 +0300 Subject: [PATCH 01/16] add to msgData --- .../metadata/TbAbstractGetAttributesNode.java | 64 +++++++++++++++---- .../engine/metadata/TbGetAttributesNode.java | 6 +- .../TbGetAttributesNodeConfiguration.java | 2 + .../TbGetDeviceAttrNodeConfiguration.java | 1 + 4 files changed, 56 insertions(+), 17 deletions(-) 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 b2c029589c..6f90078cdb 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 @@ -17,6 +17,7 @@ package org.thingsboard.rule.engine.metadata; import com.fasterxml.jackson.core.JsonParser; import com.fasterxml.jackson.core.json.JsonWriteFeature; +import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; @@ -25,6 +26,7 @@ import com.google.common.util.concurrent.MoreExecutors; import com.google.gson.JsonParseException; 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; @@ -57,10 +59,12 @@ public abstract class TbAbstractGetAttributesNode> failuresMap = new ConcurrentHashMap<>(); ListenableFuture> allFutures = Futures.allAsList( - putLatestTelemetry(ctx, entityId, msg, LATEST_TS, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresMap), - putAttrAsync(ctx, entityId, msg, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg), failuresMap, "cs_"), - putAttrAsync(ctx, entityId, msg, SHARED_SCOPE, TbNodeUtils.processPatterns(config.getSharedAttributeNames(), msg), failuresMap, "shared_"), - putAttrAsync(ctx, entityId, msg, SERVER_SCOPE, TbNodeUtils.processPatterns(config.getServerAttributeNames(), msg), failuresMap, "ss_") + putLatestTelemetry(ctx, entityId, msg, LATEST_TS, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresMap, msgNewData), + putAttrAsync(ctx, entityId, msg, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg), failuresMap, "cs_", msgNewData), + putAttrAsync(ctx, entityId, msg, SHARED_SCOPE, TbNodeUtils.processPatterns(config.getSharedAttributeNames(), msg), failuresMap, "shared_", msgNewData), + putAttrAsync(ctx, entityId, msg, SERVER_SCOPE, TbNodeUtils.processPatterns(config.getServerAttributeNames(), msg), failuresMap, "ss_", msgNewData) ); withCallback(allFutures, i -> { if (!failuresMap.isEmpty()) { throw reportFailures(failuresMap); } - ctx.tellSuccess(msg); + if (fetchToData) { + ObjectNode msgDataNode = (ObjectNode) JacksonUtil.toJsonNode(msg.getData()); + msgDataNode.setAll(msgNewData); + ctx.tellSuccess(TbMsg.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), JacksonUtil.toString(msgDataNode))); + } else { + ctx.tellSuccess(msg); + } }, t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } - private ListenableFuture putAttrAsync(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap, String prefix) { + private ListenableFuture putAttrAsync(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap, String prefix, ObjectNode msgData) { if (CollectionUtils.isEmpty(keys)) { return Futures.immediateFuture(null); } @@ -109,7 +127,13 @@ public abstract class TbAbstractGetAttributesNode { if (!CollectionUtils.isEmpty(attributeKvEntryList)) { List existingAttributesKvEntry = attributeKvEntryList.stream().filter(attributeKvEntry -> keys.contains(attributeKvEntry.getKey())).collect(Collectors.toList()); - existingAttributesKvEntry.forEach(kvEntry -> msg.getMetaData().putValue(prefix + kvEntry.getKey(), kvEntry.getValueAsString())); + existingAttributesKvEntry.forEach(kvEntry -> { + if (fetchToData) { + msgData.put(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); + } else { + msg.getMetaData().putValue(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); + } + }); if (existingAttributesKvEntry.size() != keys.size() && BooleanUtils.toBooleanDefaultIfNull(this.config.isTellFailureIfAbsent(), true)) { getNotExistingKeys(existingAttributesKvEntry, keys).forEach(key -> computeFailuresMap(scope, failuresMap, key)); } @@ -122,7 +146,7 @@ public abstract class TbAbstractGetAttributesNode putLatestTelemetry(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap) { + private ListenableFuture putLatestTelemetry(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap, ObjectNode msgData) { if (CollectionUtils.isEmpty(keys)) { return Futures.immediateFuture(null); } @@ -134,16 +158,24 @@ public abstract class TbAbstractGetAttributesNode getNotExistingKeys(List existingAttributesKvEntry, List allKeys) { 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 431fe64c09..6388f9c44e 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 @@ -34,9 +34,9 @@ import org.thingsboard.server.common.msg.TbMsg; @RuleNode(type = ComponentType.ENRICHMENT, name = "originator attributes", configClazz = TbGetAttributesNodeConfiguration.class, - nodeDescription = "Add Message Originator Attributes or Latest Telemetry into Message Metadata", - nodeDetails = "If Attributes enrichment configured, CLIENT/SHARED/SERVER attributes are added into Message metadata " + - "with specific prefix: cs/shared/ss. Latest telemetry value added into metadata without prefix. " + + nodeDescription = "Add Message Originator Attributes or Latest Telemetry into Message Metadata/data", + nodeDetails = "If Attributes enrichment configured, CLIENT/SHARED/SERVER attributes are added into Message metadata or data " + + "with specific prefix: cs/shared/ss. Latest telemetry value added into metadata/data without prefix. " + "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"}, 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 4fc892296f..67766e5b51 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 @@ -35,6 +35,7 @@ public class TbGetAttributesNodeConfiguration implements NodeConfiguration Date: Fri, 7 Oct 2022 10:27:58 +0300 Subject: [PATCH 02/16] refactor node --- .../thingsboard/common/util/JacksonUtil.java | 22 +++ .../metadata/TbAbstractGetAttributesNode.java | 146 +++++++++--------- 2 files changed, 98 insertions(+), 70 deletions(-) diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index 9c4b68d3c0..46d1881ab7 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java @@ -24,6 +24,8 @@ import com.fasterxml.jackson.databind.SerializationFeature; import com.fasterxml.jackson.databind.json.JsonMapper; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; +import org.thingsboard.server.common.data.kv.DataType; +import org.thingsboard.server.common.data.kv.KvEntry; import java.io.IOException; import java.util.ArrayList; @@ -212,4 +214,24 @@ public class JacksonUtil { } } } + + public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry, String prefix) { + if (kvEntry.getDataType() == DataType.BOOLEAN) { + kvEntry.getBooleanValue().ifPresent(value -> entityNode.put(prefix + kvEntry.getKey(), value)); + } else if (kvEntry.getDataType() == DataType.DOUBLE) { + kvEntry.getDoubleValue().ifPresent(value -> entityNode.put(prefix + kvEntry.getKey(), value)); + } else if (kvEntry.getDataType() == DataType.LONG) { + kvEntry.getLongValue().ifPresent(value -> entityNode.put(prefix + kvEntry.getKey(), value)); + } else if (kvEntry.getDataType() == DataType.JSON) { + if (kvEntry.getJsonValue().isPresent()) { + entityNode.set(prefix + kvEntry.getKey(), JacksonUtil.valueToTree(kvEntry.getJsonValue().get())); + } + } else { + entityNode.put(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); + } + } + + public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry) { + addKvEntry(entityNode, kvEntry, ""); + } } 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 6f90078cdb..89c92305ca 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 @@ -34,13 +34,19 @@ 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.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.JsonDataEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgMetaData; import java.io.IOException; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +import java.util.Map; +import java.util.Objects; import java.util.concurrent.ConcurrentHashMap; import java.util.stream.Collectors; @@ -60,11 +66,15 @@ public abstract class TbAbstractGetAttributesNode> failuresMap = new ConcurrentHashMap<>(); - ListenableFuture> allFutures = Futures.allAsList( - putLatestTelemetry(ctx, entityId, msg, LATEST_TS, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresMap, msgNewData), - putAttrAsync(ctx, entityId, msg, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg), failuresMap, "cs_", msgNewData), - putAttrAsync(ctx, entityId, msg, SHARED_SCOPE, TbNodeUtils.processPatterns(config.getSharedAttributeNames(), msg), failuresMap, "shared_", msgNewData), - putAttrAsync(ctx, entityId, msg, SERVER_SCOPE, TbNodeUtils.processPatterns(config.getServerAttributeNames(), msg), failuresMap, "ss_", msgNewData) + ListenableFuture>>> allFutures = Futures.allAsList( + getLatestTelemetry(ctx, entityId, msg, LATEST_TS, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresMap), + getAttrAsync(ctx, entityId, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg), failuresMap), + getAttrAsync(ctx, entityId, SHARED_SCOPE, TbNodeUtils.processPatterns(config.getSharedAttributeNames(), msg), failuresMap), + getAttrAsync(ctx, entityId, SERVER_SCOPE, TbNodeUtils.processPatterns(config.getServerAttributeNames(), msg), failuresMap) ); - withCallback(allFutures, i -> { + withCallback(allFutures, futuresList -> { if (!failuresMap.isEmpty()) { throw reportFailures(failuresMap); } + TbMsgMetaData msgMetaData = msg.getMetaData(); + futuresList.stream().filter(Objects::nonNull).forEach(kvEntriesMap -> { + kvEntriesMap.forEach((keyScope, kvEntryList) -> { + String prefix = getPrefix(keyScope); + kvEntryList.forEach(kvEntry -> { + if (fetchToData) { + JacksonUtil.addKvEntry((ObjectNode) msgDataNode, kvEntry, prefix); + } else { + msgMetaData.putValue(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); + } + }); + }); + }); if (fetchToData) { - ObjectNode msgDataNode = (ObjectNode) JacksonUtil.toJsonNode(msg.getData()); - msgDataNode.setAll(msgNewData); ctx.tellSuccess(TbMsg.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), JacksonUtil.toString(msgDataNode))); } else { ctx.tellSuccess(msg); @@ -119,100 +139,86 @@ public abstract class TbAbstractGetAttributesNode ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } - private ListenableFuture putAttrAsync(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap, String prefix, ObjectNode msgData) { + private ListenableFuture>> getAttrAsync(TbContext ctx, EntityId entityId, String scope, List keys, ConcurrentHashMap> failuresMap) { if (CollectionUtils.isEmpty(keys)) { return Futures.immediateFuture(null); } ListenableFuture> attributeKvEntryListFuture = ctx.getAttributesService().find(ctx.getTenantId(), entityId, scope, keys); return Futures.transform(attributeKvEntryListFuture, attributeKvEntryList -> { - if (!CollectionUtils.isEmpty(attributeKvEntryList)) { - List existingAttributesKvEntry = attributeKvEntryList.stream().filter(attributeKvEntry -> keys.contains(attributeKvEntry.getKey())).collect(Collectors.toList()); - existingAttributesKvEntry.forEach(kvEntry -> { - if (fetchToData) { - msgData.put(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); - } else { - msg.getMetaData().putValue(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); - } - }); - if (existingAttributesKvEntry.size() != keys.size() && BooleanUtils.toBooleanDefaultIfNull(this.config.isTellFailureIfAbsent(), true)) { - getNotExistingKeys(existingAttributesKvEntry, keys).forEach(key -> computeFailuresMap(scope, failuresMap, key)); - } - } else { - if (BooleanUtils.toBooleanDefaultIfNull(this.config.isTellFailureIfAbsent(), true)) { - keys.forEach(key -> computeFailuresMap(scope, failuresMap, key)); - } + if (isTellFailureIfAbsent && attributeKvEntryList.size() != keys.size()) { + getNotExistingKeys(attributeKvEntryList, keys).forEach(key -> computeFailuresMap(scope, failuresMap, key)); } - return null; + Map> mapAttributeKvEntry = new HashMap<>(); + mapAttributeKvEntry.put(scope, attributeKvEntryList); + return mapAttributeKvEntry; }, MoreExecutors.directExecutor()); } - private ListenableFuture putLatestTelemetry(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap, ObjectNode msgData) { + private ListenableFuture>> getLatestTelemetry(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap) { if (CollectionUtils.isEmpty(keys)) { return Futures.immediateFuture(null); } - ListenableFuture> latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, keys); - return Futures.transform(latest, l -> { - l.forEach(r -> { - boolean getLatestValueWithTs = BooleanUtils.toBooleanDefaultIfNull(this.config.isGetLatestValueWithTs(), false); - if (BooleanUtils.toBooleanDefaultIfNull(this.config.isTellFailureIfAbsent(), true)) { - if (r.getValue() == null) { - computeFailuresMap(scope, failuresMap, r.getKey()); - } else if (getLatestValueWithTs) { - putValueWithTs(msg, r, msgData); - } else { - if (fetchToData) { - msgData.put(r.getKey(), r.getValueAsString()); - } else { - msg.getMetaData().putValue(r.getKey(), r.getValueAsString()); - } + ListenableFuture> latestTelemetryFutures = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, keys); + return Futures.transform(latestTelemetryFutures, tsKvEntries -> { + List listTsKvEntry = new ArrayList<>(); + tsKvEntries.forEach(tsKvEntry -> { + if (tsKvEntry.getValue() == null) { + if (isTellFailureIfAbsent) { + computeFailuresMap(scope, failuresMap, tsKvEntry.getKey()); } + } else if (getLatestValueWithTs) { + listTsKvEntry.add(getValueWithTs(tsKvEntry)); } else { - if (r.getValue() != null) { - if (getLatestValueWithTs) { - putValueWithTs(msg, r, msgData); - } else { - if (fetchToData) { - msgData.put(r.getKey(), r.getValueAsString()); - } else { - msg.getMetaData().putValue(r.getKey(), r.getValueAsString()); - } - } - } + listTsKvEntry.add(new BasicTsKvEntry(tsKvEntry.getTs(), tsKvEntry)); } }); - return null; + Map> mapTsKvEntry = new HashMap<>(); + mapTsKvEntry.put(scope, listTsKvEntry); + return mapTsKvEntry; }, MoreExecutors.directExecutor()); } - private void putValueWithTs(TbMsg msg, TsKvEntry r, ObjectNode msgData) { - ObjectNode value = mapper.createObjectNode(); - value.put(TS, r.getTs()); - switch (r.getDataType()) { + private TsKvEntry getValueWithTs(TsKvEntry tsKvEntry) { + ObjectNode value = JacksonUtil.newObjectNode(); + value.put(TS, tsKvEntry.getTs()); + switch (tsKvEntry.getDataType()) { case STRING: - value.put(VALUE, r.getValueAsString()); + value.put(VALUE, tsKvEntry.getValueAsString()); break; case LONG: - value.put(VALUE, r.getLongValue().get()); + value.put(VALUE, tsKvEntry.getLongValue().get()); break; case BOOLEAN: - value.put(VALUE, r.getBooleanValue().get()); + value.put(VALUE, tsKvEntry.getBooleanValue().get()); break; case DOUBLE: - value.put(VALUE, r.getDoubleValue().get()); + value.put(VALUE, tsKvEntry.getDoubleValue().get()); break; case JSON: try { - value.set(VALUE, mapper.readTree(r.getJsonValue().get())); + value.set(VALUE, mapper.readTree(tsKvEntry.getJsonValue().get())); } catch (IOException e) { - throw new JsonParseException("Can't parse jsonValue: " + r.getJsonValue().get(), e); + throw new JsonParseException("Can't parse jsonValue: " + tsKvEntry.getJsonValue().get(), e); } break; } - if (fetchToData) { - msgData.putIfAbsent(r.getKey(), value); - } else { - msg.getMetaData().putValue(r.getKey(), value.toString()); + return new BasicTsKvEntry(tsKvEntry.getTs(), new JsonDataEntry(tsKvEntry.getKey(), value.toString())); + } + + private String getPrefix(String scope) { + String prefix = ""; + switch (scope) { + case CLIENT_SCOPE: + prefix = "cs_"; + break; + case SHARED_SCOPE: + prefix = "shared_"; + break; + case SERVER_SCOPE: + prefix = "ss_"; + break; } + return prefix; } private List getNotExistingKeys(List existingAttributesKvEntry, List allKeys) { From 0ec08b8448d75fd85ce10ff9d1b0c00156888c0b Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Fri, 7 Oct 2022 12:13:48 +0300 Subject: [PATCH 03/16] code review --- .../thingsboard/common/util/JacksonUtil.java | 14 +++--- .../metadata/TbAbstractGetAttributesNode.java | 49 ++++--------------- 2 files changed, 16 insertions(+), 47 deletions(-) diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index 46d1881ab7..1bfd85351c 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java @@ -215,23 +215,23 @@ public class JacksonUtil { } } - public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry, String prefix) { + public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry, String key) { if (kvEntry.getDataType() == DataType.BOOLEAN) { - kvEntry.getBooleanValue().ifPresent(value -> entityNode.put(prefix + kvEntry.getKey(), value)); + kvEntry.getBooleanValue().ifPresent(value -> entityNode.put(key, value)); } else if (kvEntry.getDataType() == DataType.DOUBLE) { - kvEntry.getDoubleValue().ifPresent(value -> entityNode.put(prefix + kvEntry.getKey(), value)); + kvEntry.getDoubleValue().ifPresent(value -> entityNode.put(key, value)); } else if (kvEntry.getDataType() == DataType.LONG) { - kvEntry.getLongValue().ifPresent(value -> entityNode.put(prefix + kvEntry.getKey(), value)); + kvEntry.getLongValue().ifPresent(value -> entityNode.put(key, value)); } else if (kvEntry.getDataType() == DataType.JSON) { if (kvEntry.getJsonValue().isPresent()) { - entityNode.set(prefix + kvEntry.getKey(), JacksonUtil.valueToTree(kvEntry.getJsonValue().get())); + entityNode.set(key, JacksonUtil.valueToTree(kvEntry.getJsonValue().get())); } } else { - entityNode.put(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); + entityNode.put(key, kvEntry.getValueAsString()); } } public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry) { - addKvEntry(entityNode, kvEntry, ""); + addKvEntry(entityNode, kvEntry, kvEntry.getKey()); } } 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 89c92305ca..c516811e9c 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 @@ -15,15 +15,11 @@ */ package org.thingsboard.rule.engine.metadata; -import com.fasterxml.jackson.core.JsonParser; -import com.fasterxml.jackson.core.json.JsonWriteFeature; import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; -import com.google.gson.JsonParseException; import org.apache.commons.collections.CollectionUtils; import org.apache.commons.lang3.BooleanUtils; import org.thingsboard.common.util.JacksonUtil; @@ -39,9 +35,7 @@ import org.thingsboard.server.common.data.kv.JsonDataEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.common.msg.TbMsgMetaData; -import java.io.IOException; import java.util.ArrayList; import java.util.HashMap; import java.util.List; @@ -59,8 +53,6 @@ import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE; public abstract class TbAbstractGetAttributesNode implements TbNode { - private static ObjectMapper mapper = new ObjectMapper(); - private static final String VALUE = "value"; private static final String TS = "ts"; @@ -73,10 +65,8 @@ public abstract class TbAbstractGetAttributesNode> failuresMap = new ConcurrentHashMap<>(); - ListenableFuture>>> allFutures = Futures.allAsList( + ListenableFuture>>> allFutures = Futures.allAsList( getLatestTelemetry(ctx, entityId, msg, LATEST_TS, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresMap), getAttrAsync(ctx, entityId, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg), failuresMap), getAttrAsync(ctx, entityId, SHARED_SCOPE, TbNodeUtils.processPatterns(config.getSharedAttributeNames(), msg), failuresMap), @@ -118,15 +108,14 @@ public abstract class TbAbstractGetAttributesNode { kvEntriesMap.forEach((keyScope, kvEntryList) -> { String prefix = getPrefix(keyScope); kvEntryList.forEach(kvEntry -> { if (fetchToData) { - JacksonUtil.addKvEntry((ObjectNode) msgDataNode, kvEntry, prefix); + JacksonUtil.addKvEntry((ObjectNode) msgDataNode, kvEntry, prefix + kvEntry.getKey()); } else { - msgMetaData.putValue(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); + msg.getMetaData().putValue(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); } }); }); @@ -139,7 +128,7 @@ public abstract class TbAbstractGetAttributesNode ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } - private ListenableFuture>> getAttrAsync(TbContext ctx, EntityId entityId, String scope, List keys, ConcurrentHashMap> failuresMap) { + private ListenableFuture>> getAttrAsync(TbContext ctx, EntityId entityId, String scope, List keys, ConcurrentHashMap> failuresMap) { if (CollectionUtils.isEmpty(keys)) { return Futures.immediateFuture(null); } @@ -148,13 +137,13 @@ public abstract class TbAbstractGetAttributesNode computeFailuresMap(scope, failuresMap, key)); } - Map> mapAttributeKvEntry = new HashMap<>(); + Map> mapAttributeKvEntry = new HashMap<>(); mapAttributeKvEntry.put(scope, attributeKvEntryList); return mapAttributeKvEntry; }, MoreExecutors.directExecutor()); } - private ListenableFuture>> getLatestTelemetry(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap) { + private ListenableFuture>> getLatestTelemetry(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap) { if (CollectionUtils.isEmpty(keys)) { return Futures.immediateFuture(null); } @@ -172,7 +161,7 @@ public abstract class TbAbstractGetAttributesNode> mapTsKvEntry = new HashMap<>(); + Map> mapTsKvEntry = new HashMap<>(); mapTsKvEntry.put(scope, listTsKvEntry); return mapTsKvEntry; }, MoreExecutors.directExecutor()); @@ -181,27 +170,7 @@ public abstract class TbAbstractGetAttributesNode Date: Fri, 7 Oct 2022 13:07:12 +0300 Subject: [PATCH 04/16] add JacksonUtil.toJsonNode --- .../src/main/java/org/thingsboard/common/util/JacksonUtil.java | 2 +- .../rule/engine/metadata/TbAbstractGetAttributesNode.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index 1bfd85351c..55d3c254d7 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java @@ -224,7 +224,7 @@ public class JacksonUtil { kvEntry.getLongValue().ifPresent(value -> entityNode.put(key, value)); } else if (kvEntry.getDataType() == DataType.JSON) { if (kvEntry.getJsonValue().isPresent()) { - entityNode.set(key, JacksonUtil.valueToTree(kvEntry.getJsonValue().get())); + entityNode.set(key, JacksonUtil.toJsonNode(kvEntry.getJsonValue().get())); } } else { entityNode.put(key, kvEntry.getValueAsString()); 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 c516811e9c..2e525b56ed 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 @@ -171,7 +171,7 @@ public abstract class TbAbstractGetAttributesNode Date: Mon, 17 Oct 2022 12:36:13 +0300 Subject: [PATCH 05/16] add description --- .../rule/engine/metadata/TbGetAttributesNode.java | 6 +++--- .../rule/engine/metadata/TbGetDeviceAttrNode.java | 6 +++--- 2 files changed, 6 insertions(+), 6 deletions(-) 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 6388f9c44e..472b84c803 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 @@ -34,9 +34,9 @@ import org.thingsboard.server.common.msg.TbMsg; @RuleNode(type = ComponentType.ENRICHMENT, name = "originator attributes", configClazz = TbGetAttributesNodeConfiguration.class, - nodeDescription = "Add Message Originator Attributes or Latest Telemetry into Message Metadata/data", - nodeDetails = "If Attributes enrichment configured, CLIENT/SHARED/SERVER attributes are added into Message metadata or data " + - "with specific prefix: cs/shared/ss. Latest telemetry value added into metadata/data without prefix. " + + nodeDescription = "Add Message Originator Attributes or Latest Telemetry into Message Data or Metadata", + 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 " + "metadata.cs_temperature or metadata.shared_limit ", uiResources = {"static/rulenode/rulenode-core-config.js"}, 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 5a48c8f95e..f38b5eccc7 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 @@ -31,9 +31,9 @@ import org.thingsboard.server.common.msg.TbMsg; @RuleNode(type = ComponentType.ENRICHMENT, name = "related device attributes", configClazz = TbGetDeviceAttrNodeConfiguration.class, - nodeDescription = "Add Originators Related Device Attributes and Latest Telemetry value into Message Metadata", - nodeDetails = "If Attributes enrichment configured, CLIENT/SHARED/SERVER attributes are added into Message metadata " + - "with specific prefix: cs/shared/ss. Latest telemetry value added into metadata without prefix. " + + nodeDescription = "Add Originators Related Device Attributes and Latest Telemetry value into Message Data or Metadata", + 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 " + "metadata.cs_temperature or metadata.shared_limit ", uiResources = {"static/rulenode/rulenode-core-config.js"}, From 4c277daadd8a3941d903d65ae095c0a73d3e7518 Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Fri, 28 Oct 2022 11:05:28 +0300 Subject: [PATCH 06/16] add test for TbAbstractGetAttributesNode --- .../metadata/TbAbstractGetAttributesNode.java | 2 +- .../TbAbstractGetAttributesNodeTest.java | 296 ++++++++++++++++++ 2 files changed, 297 insertions(+), 1 deletion(-) create mode 100644 rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNodeTest.java 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 2e525b56ed..e364ea2c18 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 @@ -93,7 +93,7 @@ public abstract class TbAbstractGetAttributesNode clientAttributes; + private List serverAttributes; + private List sharedAttributes; + private List tsKeys; + + @Before + public void before() throws TbNodeException { + dbExecutor = new AbstractListeningExecutor() { + @Override + protected int getThreadPollSize() { + return 3; + } + }; + dbExecutor.init(); + + Mockito.reset(ctx); + Mockito.reset(attributesService); + Mockito.reset(tsService); + + Mockito.reset(ctx); + Mockito.reset(attributesService); + Mockito.reset(tsService); + + lenient().when(ctx.getAttributesService()).thenReturn(attributesService); + lenient().when(ctx.getTimeseriesService()).thenReturn(tsService); + lenient().when(ctx.getTenantId()).thenReturn(tenantId); + lenient().when(ctx.getDbCallbackExecutor()).thenReturn(dbExecutor); + + clientAttributes = getAttributeNames("client"); + serverAttributes = getAttributeNames("server"); + sharedAttributes = getAttributeNames("shared"); + tsKeys = List.of("temperature", "humidity", "unknown"); + + Mockito.when(attributesService.find(tenantId, originator, DataConstants.CLIENT_SCOPE, clientAttributes)) + .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(clientAttributes))); + + + Mockito.when(attributesService.find(tenantId, originator, DataConstants.SERVER_SCOPE, serverAttributes)) + .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(serverAttributes))); + + + Mockito.when(attributesService.find(tenantId, originator, DataConstants.SHARED_SCOPE, sharedAttributes)) + .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(sharedAttributes))); + + Mockito.when(tsService.findLatest(tenantId, originator, tsKeys)) + .thenReturn(Futures.immediateFuture(getListTelemetryKvEntry(tsKeys))); + } + + @After + public void after() { + dbExecutor.destroy(); + } + + @Test + public void fetchToMetadata_whenOnMsg_then_success() throws Exception { + TbGetAttributesNode node = initNode(false, false, false); + TbMsg msg = getTbMsg(originator); + node.onMsg(ctx, msg); + + TbMsg resultMsg = checkMsg(); + TbMsgMetaData msgMetaData = resultMsg.getMetaData(); + + //check attributes + checkAttributes(clientAttributes, "cs_", false, msgMetaData, null); + checkAttributes(serverAttributes, "ss_", false, msgMetaData, null); + checkAttributes(sharedAttributes, "shared_", false, msgMetaData, null); + + //check timeseries + checkTs(tsKeys, false, false, msgMetaData, null); + } + + @Test + public void fetchToData_whenOnMsg_then_success() throws Exception { + TbGetAttributesNode node = initNode(true, true, false); + TbMsg msg = getTbMsg(originator); + node.onMsg(ctx, msg); + + TbMsg resultMsg = checkMsg(); + JsonNode msgData = JacksonUtil.toJsonNode(resultMsg.getData()); + + //check attributes + checkAttributes(clientAttributes, "cs_", true, null, msgData); + checkAttributes(serverAttributes, "ss_", true, null, msgData); + checkAttributes(sharedAttributes, "shared_", true, null, msgData); + + //check timeseries with ts + checkTs(tsKeys, true, true, null, msgData); + } + + @Test + public void fetchToData_whenOnMsg_then_failure() throws Exception { + TbGetAttributesNode node = initNode(true, true, true); + TbMsg msg = getTbMsg(originator); + node.onMsg(ctx, msg); + + ArgumentCaptor newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); + ArgumentCaptor exceptionCaptor = ArgumentCaptor.forClass(Exception.class); + Mockito.verify(ctx, never()).tellSuccess(any()); + Mockito.verify(ctx, Mockito.timeout(5000)).tellFailure(newMsgCaptor.capture(), exceptionCaptor.capture()); + + Assert.assertSame(newMsgCaptor.getValue(), msg); + Assert.assertNotNull(exceptionCaptor.getValue()); + } + + @Test + public void fetchToData_whenOnMsg_then_data_not_object_failure() throws Exception { + TbGetAttributesNode node = initNode(true, true, true); + TbMsg msg = TbMsg.newMsg("TEST", originator, new TbMsgMetaData(), "[]"); + node.onMsg(ctx, msg); + + ArgumentCaptor newMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); + ArgumentCaptor exceptionCaptor = ArgumentCaptor.forClass(Exception.class); + Mockito.verify(ctx, never()).tellSuccess(any()); + Mockito.verify(ctx, Mockito.timeout(5000)).tellFailure(newMsgCaptor.capture(), exceptionCaptor.capture()); + + Assert.assertSame(newMsgCaptor.getValue(), msg); + Assert.assertNotNull(exceptionCaptor.getValue()); + } + + private TbMsg checkMsg() { + ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); + Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); + + TbMsg resultMsg = msgCaptor.getValue(); + Assert.assertNotNull(resultMsg); + Assert.assertNotNull(resultMsg.getMetaData()); + Assert.assertNotNull(resultMsg.getData()); + return resultMsg; + } + + private void checkAttributes(List attributes, String prefix, boolean fetchToData, TbMsgMetaData msgMetaData, JsonNode msgData) { + attributes.stream() + .filter(attribute -> !attribute.equals("unknown")) + .forEach(attribute -> { + String result; + if (fetchToData) { + result = msgData.get(prefix + attribute).asText(); + } else { + result = msgMetaData.getValue(prefix + attribute); + } + Assert.assertNotNull(result); + Assert.assertEquals(attribute + "_value", result); + }); + } + + private void checkTs(List tsKeys, boolean fetchToData, boolean getLatestValueWithTs, TbMsgMetaData msgMetaData, JsonNode msgData) { + long tsValue = 1L; + for (String key : tsKeys) { + if (key.equals("unknown")) { + continue; + } + String result; + if (fetchToData) { + if (getLatestValueWithTs) { + JsonNode resultTs = msgData.get(key); + Assert.assertNotNull(resultTs); + Assert.assertNotNull(resultTs.get("value")); + Assert.assertNotNull(resultTs.get("ts")); + result = resultTs.get("value").asText(); + } else { + result = msgData.get(key).asText(); + } + } else { + result = msgMetaData.getValue(key); + } + Assert.assertNotNull(result); + Assert.assertEquals(String.valueOf(tsValue), result); + tsValue++; + } + } + + private TbGetAttributesNode initNode(boolean fetchToData, 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.setGetLatestValueWithTs(getLatestValueWithTs); + config.setTellFailureIfAbsent(isTellFailureIfAbsent); + TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + TbGetAttributesNode node = new TbGetAttributesNode(); + node.init(ctx, nodeConfiguration); + return node; + } + + private TbMsg getTbMsg(EntityId entityId) { + ObjectNode msgData = JacksonUtil.newObjectNode(); + msgData.put("shared_attr_data", "shared_attr_3"); + + TbMsgMetaData msgMetaData = new TbMsgMetaData(); + msgMetaData.putValue("client_attr_metadata", "client_attr_3"); + msgMetaData.putValue("server_attr_metadata", "server_attr_3"); + + return TbMsg.newMsg("TEST", entityId, msgMetaData, msgData.toString()); + } + + private List getAttributeNames(String prefix) { + return List.of(prefix + "_attr_1", prefix + "_attr_2", prefix + "_attr_3", "unknown"); + } + + private List getListAttributeKvEntry(List attributes) { + List attributeKvEntries = new ArrayList<>(); + attributes.stream().filter(attribute -> !attribute.equals("unknown")).forEach(attribute -> attributeKvEntries.add(new BaseAttributeKvEntry(System.currentTimeMillis(), new StringDataEntry(attribute, attribute + "_value")))); + return attributeKvEntries; + } + + private List getListTelemetryKvEntry(List keys) { + long value = 1L; + List kvEntries = new ArrayList<>(); + for (String key : keys) { + if (key.equals("unknown")) { + continue; + } + kvEntries.add(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry(key, value))); + value++; + } + return kvEntries; + } + +} \ No newline at end of file From 36651d231da2db33ede7bd4e1a624c8c7e9aa70b Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Fri, 28 Oct 2022 11:06:00 +0300 Subject: [PATCH 07/16] mvn license:format --- .../rule/engine/metadata/TbAbstractGetAttributesNodeTest.java | 1 - 1 file changed, 1 deletion(-) 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 fc8f4ea5c9..e2914c8e60 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 @@ -55,7 +55,6 @@ import java.util.List; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.lenient; import static org.mockito.Mockito.never; -import static org.mockito.Mockito.times; @RunWith(MockitoJUnitRunner.class) public class TbAbstractGetAttributesNodeTest { From 513d0f307587be117e8bb60da7eddc16c80f3c9f Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Mon, 31 Oct 2022 14:10:23 +0200 Subject: [PATCH 08/16] refactor code review --- .../rule/engine/metadata/TbAbstractGetAttributesNode.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) 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 e364ea2c18..75f8c7732c 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 @@ -35,6 +35,7 @@ import org.thingsboard.server.common.data.kv.JsonDataEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgMetaData; import java.util.ArrayList; import java.util.HashMap; @@ -108,6 +109,7 @@ public abstract class TbAbstractGetAttributesNode { kvEntriesMap.forEach((keyScope, kvEntryList) -> { String prefix = getPrefix(keyScope); @@ -115,15 +117,15 @@ public abstract class TbAbstractGetAttributesNode ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } From 691339c5c5c907bbe9dbba8d0cce898dec991245 Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Mon, 31 Oct 2022 16:03:27 +0200 Subject: [PATCH 09/16] update JacksonUtil method --- .../src/main/java/org/thingsboard/common/util/JacksonUtil.java | 2 +- .../rule/engine/metadata/TbAbstractGetAttributesNodeTest.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index 55d3c254d7..9fc83374a2 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java @@ -224,7 +224,7 @@ public class JacksonUtil { kvEntry.getLongValue().ifPresent(value -> entityNode.put(key, value)); } else if (kvEntry.getDataType() == DataType.JSON) { if (kvEntry.getJsonValue().isPresent()) { - entityNode.set(key, JacksonUtil.toJsonNode(kvEntry.getJsonValue().get())); + entityNode.set(key, toJsonNode(kvEntry.getJsonValue().get())); } } else { entityNode.put(key, kvEntry.getValueAsString()); 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 e2914c8e60..a066c79343 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 @@ -292,4 +292,4 @@ public class TbAbstractGetAttributesNodeTest { return kvEntries; } -} \ No newline at end of file +} From 1d8733e55ad8b46345e4c5a4c6cb2cc244d7edb8 Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Tue, 1 Nov 2022 16:54:25 +0200 Subject: [PATCH 10/16] return mapper --- .../metadata/TbAbstractGetAttributesNode.java | 43 ++++++++++++++++--- 1 file changed, 38 insertions(+), 5 deletions(-) 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 75f8c7732c..f50fe58dd5 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 @@ -15,7 +15,10 @@ */ package org.thingsboard.rule.engine.metadata; +import com.fasterxml.jackson.core.JsonParser; +import com.fasterxml.jackson.core.json.JsonWriteFeature; import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; @@ -31,12 +34,14 @@ 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.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.DataType; import org.thingsboard.server.common.data.kv.JsonDataEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; +import java.io.IOException; import java.util.ArrayList; import java.util.HashMap; import java.util.List; @@ -54,6 +59,8 @@ import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE; public abstract class TbAbstractGetAttributesNode implements TbNode { + private static ObjectMapper mapper = new ObjectMapper(); + private static final String VALUE = "value"; private static final String TS = "ts"; @@ -65,6 +72,8 @@ public abstract class TbAbstractGetAttributesNode { if (fetchToData) { - JacksonUtil.addKvEntry((ObjectNode) msgDataNode, kvEntry, prefix + kvEntry.getKey()); + addKvEntryToJson((ObjectNode) msgDataNode, kvEntry, prefix + kvEntry.getKey()); } else { msgMetaData.putValue(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); } @@ -170,10 +179,10 @@ public abstract class TbAbstractGetAttributesNode entityNode.put(key, value)); + } else if (kvEntry.getDataType() == DataType.DOUBLE) { + kvEntry.getDoubleValue().ifPresent(value -> entityNode.put(key, value)); + } else if (kvEntry.getDataType() == DataType.LONG) { + kvEntry.getLongValue().ifPresent(value -> entityNode.put(key, value)); + } else if (kvEntry.getDataType() == DataType.JSON) { + if (kvEntry.getJsonValue().isPresent()) { + entityNode.set(key, toJsonNode(kvEntry.getJsonValue().get())); + } + } else { + entityNode.put(key, kvEntry.getValueAsString()); + } + } + + public static JsonNode toJsonNode(String value) { + try { + return mapper.readTree(value); + } catch (IOException e) { + throw new IllegalArgumentException(e); + } + } + private List getNotExistingKeys(List existingAttributesKvEntry, List allKeys) { List existingKeys = existingAttributesKvEntry.stream().map(KvEntry::getKey).collect(Collectors.toList()); return allKeys.stream().filter(key -> !existingKeys.contains(key)).collect(Collectors.toList()); From b74e5dfc4572da6de50579157cc8b4cf88de1fec Mon Sep 17 00:00:00 2001 From: Yuriy Lytvynchuk Date: Tue, 1 Nov 2022 16:56:08 +0200 Subject: [PATCH 11/16] change public to prv --- .../rule/engine/metadata/TbAbstractGetAttributesNode.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 f50fe58dd5..c0564a98fe 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 @@ -201,7 +201,7 @@ public abstract class TbAbstractGetAttributesNode entityNode.put(key, value)); } else if (kvEntry.getDataType() == DataType.DOUBLE) { @@ -217,7 +217,7 @@ public abstract class TbAbstractGetAttributesNode Date: Tue, 1 Nov 2022 17:03:01 +0200 Subject: [PATCH 12/16] add JsonParseException --- .../rule/engine/metadata/TbAbstractGetAttributesNode.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) 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 c0564a98fe..42a0666614 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 @@ -23,6 +23,7 @@ import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; +import com.google.gson.JsonParseException; import org.apache.commons.collections.CollectionUtils; import org.apache.commons.lang3.BooleanUtils; import org.thingsboard.common.util.JacksonUtil; @@ -221,7 +222,7 @@ public abstract class TbAbstractGetAttributesNode Date: Wed, 2 Nov 2022 11:44:46 +0200 Subject: [PATCH 13/16] added ALLOW_UNQUOTED_FIELD_NAMES_MAPPER to JacksonUtil --- .../thingsboard/common/util/JacksonUtil.java | 32 +++++++++++--- .../metadata/TbAbstractGetAttributesNode.java | 42 ++---------------- .../engine/metadata/TbGetTelemetryNode.java | 43 +++---------------- 3 files changed, 36 insertions(+), 81 deletions(-) diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index 9fc83374a2..dca796ee1a 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java @@ -1,12 +1,12 @@ /** * Copyright © 2016-2022 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,7 +15,9 @@ */ package org.thingsboard.common.util; +import com.fasterxml.jackson.core.JsonParser; import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.core.json.JsonWriteFeature; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.MapperFeature; @@ -46,6 +48,10 @@ public class JacksonUtil { .configure(SerializationFeature.ORDER_MAP_ENTRIES_BY_KEYS, true) .configure(MapperFeature.SORT_PROPERTIES_ALPHABETICALLY, true) .build(); + public static ObjectMapper ALLOW_UNQUOTED_FIELD_NAMES_MAPPER = JsonMapper.builder() + .configure(JsonWriteFeature.QUOTE_FIELD_NAMES.mappedFeature(), false) + .configure(JsonParser.Feature.ALLOW_UNQUOTED_FIELD_NAMES, true) + .build(); public static T convertValue(Object fromValue, Class toValueType) { try { @@ -119,11 +125,15 @@ public class JacksonUtil { } public static JsonNode toJsonNode(String value) { + return toJsonNode(value, OBJECT_MAPPER); + } + + public static JsonNode toJsonNode(String value, ObjectMapper mapper) { if (value == null || value.isEmpty()) { return null; } try { - return OBJECT_MAPPER.readTree(value); + return mapper.readTree(value); } catch (IOException e) { throw new IllegalArgumentException(e); } @@ -138,7 +148,11 @@ public class JacksonUtil { } public static ObjectNode newObjectNode() { - return OBJECT_MAPPER.createObjectNode(); + return newObjectNode(OBJECT_MAPPER); + } + + public static ObjectNode newObjectNode(ObjectMapper mapper) { + return mapper.createObjectNode(); } public static T clone(T value) { @@ -216,6 +230,10 @@ public class JacksonUtil { } public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry, String key) { + addKvEntry(entityNode, kvEntry, key, OBJECT_MAPPER); + } + + public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry, String key, ObjectMapper mapper) { if (kvEntry.getDataType() == DataType.BOOLEAN) { kvEntry.getBooleanValue().ifPresent(value -> entityNode.put(key, value)); } else if (kvEntry.getDataType() == DataType.DOUBLE) { @@ -224,7 +242,7 @@ public class JacksonUtil { kvEntry.getLongValue().ifPresent(value -> entityNode.put(key, value)); } else if (kvEntry.getDataType() == DataType.JSON) { if (kvEntry.getJsonValue().isPresent()) { - entityNode.set(key, toJsonNode(kvEntry.getJsonValue().get())); + entityNode.set(key, toJsonNode(kvEntry.getJsonValue().get(), mapper)); } } else { entityNode.put(key, kvEntry.getValueAsString()); 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 42a0666614..1b94bbfd2b 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 @@ -15,15 +15,11 @@ */ package org.thingsboard.rule.engine.metadata; -import com.fasterxml.jackson.core.JsonParser; -import com.fasterxml.jackson.core.json.JsonWriteFeature; import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; -import com.google.gson.JsonParseException; import org.apache.commons.collections.CollectionUtils; import org.apache.commons.lang3.BooleanUtils; import org.thingsboard.common.util.JacksonUtil; @@ -35,14 +31,12 @@ 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.BasicTsKvEntry; -import org.thingsboard.server.common.data.kv.DataType; import org.thingsboard.server.common.data.kv.JsonDataEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; -import java.io.IOException; import java.util.ArrayList; import java.util.HashMap; import java.util.List; @@ -60,8 +54,6 @@ import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE; public abstract class TbAbstractGetAttributesNode implements TbNode { - private static ObjectMapper mapper = new ObjectMapper(); - private static final String VALUE = "value"; private static final String TS = "ts"; @@ -73,8 +65,6 @@ public abstract class TbAbstractGetAttributesNode { if (fetchToData) { - addKvEntryToJson((ObjectNode) msgDataNode, kvEntry, prefix + kvEntry.getKey()); + JacksonUtil.addKvEntry((ObjectNode) msgDataNode, kvEntry, prefix + kvEntry.getKey(), JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER); } else { msgMetaData.putValue(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); } @@ -180,9 +170,9 @@ public abstract class TbAbstractGetAttributesNode entityNode.put(key, value)); - } else if (kvEntry.getDataType() == DataType.DOUBLE) { - kvEntry.getDoubleValue().ifPresent(value -> entityNode.put(key, value)); - } else if (kvEntry.getDataType() == DataType.LONG) { - kvEntry.getLongValue().ifPresent(value -> entityNode.put(key, value)); - } else if (kvEntry.getDataType() == DataType.JSON) { - if (kvEntry.getJsonValue().isPresent()) { - entityNode.set(key, toJsonNode(kvEntry.getJsonValue().get())); - } - } else { - entityNode.put(key, kvEntry.getValueAsString()); - } - } - - private static JsonNode toJsonNode(String value) { - try { - return mapper.readTree(value); - } catch (IOException e) { - throw new JsonParseException("Can't parse jsonValue: " + value, e); - } - } - private List getNotExistingKeys(List existingAttributesKvEntry, List allKeys) { List existingKeys = existingAttributesKvEntry.stream().map(KvEntry::getKey).collect(Collectors.toList()); return allKeys.stream().filter(key -> !existingKeys.contains(key)).collect(Collectors.toList()); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java index 77298cba77..045555eae8 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java @@ -15,25 +15,22 @@ */ package org.thingsboard.rule.engine.metadata; -import com.fasterxml.jackson.core.JsonParser; -import com.fasterxml.jackson.core.json.JsonWriteFeature; -import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.ListenableFuture; -import com.google.gson.JsonParseException; import lombok.Data; import lombok.NoArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.server.common.data.StringUtils; import org.apache.commons.lang3.math.NumberUtils; import org.thingsboard.common.util.DonAsynchron; +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.StringUtils; import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; @@ -41,7 +38,6 @@ import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; -import java.io.IOException; import java.util.List; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; @@ -75,7 +71,6 @@ public class TbGetTelemetryNode implements TbNode { private TbGetTelemetryNodeConfiguration config; private List tsKeyNames; private int limit; - private ObjectMapper mapper; private String fetchMode; private String orderByFetchAll; private Aggregation aggregation; @@ -91,10 +86,6 @@ public class TbGetTelemetryNode implements TbNode { orderByFetchAll = ASC_ORDER; } aggregation = parseAggregationConfig(config.getAggregation()); - - mapper = new ObjectMapper(); - mapper.configure(JsonWriteFeature.QUOTE_FIELD_NAMES.mappedFeature(), false); - mapper.configure(JsonParser.Feature.ALLOW_UNQUOTED_FIELD_NAMES, true); } Aggregation parseAggregationConfig(String aggName) { @@ -146,7 +137,7 @@ public class TbGetTelemetryNode implements TbNode { } private void process(List entries, TbMsg msg, List keys) { - ObjectNode resultNode = mapper.createObjectNode(); + ObjectNode resultNode = JacksonUtil.newObjectNode(JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER); if (FETCH_MODE_ALL.equals(fetchMode)) { entries.forEach(entry -> processArray(resultNode, entry)); } else { @@ -169,36 +160,16 @@ public class TbGetTelemetryNode implements TbNode { ArrayNode arrayNode = (ArrayNode) node.get(entry.getKey()); arrayNode.add(buildNode(entry)); } else { - ArrayNode arrayNode = mapper.createArrayNode(); + ArrayNode arrayNode = JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER.createArrayNode(); arrayNode.add(buildNode(entry)); node.set(entry.getKey(), arrayNode); } } private ObjectNode buildNode(TsKvEntry entry) { - ObjectNode obj = mapper.createObjectNode() - .put("ts", entry.getTs()); - switch (entry.getDataType()) { - case STRING: - obj.put("value", entry.getValueAsString()); - break; - case LONG: - obj.put("value", entry.getLongValue().get()); - break; - case BOOLEAN: - obj.put("value", entry.getBooleanValue().get()); - break; - case DOUBLE: - obj.put("value", entry.getDoubleValue().get()); - break; - case JSON: - try { - obj.set("value", mapper.readTree(entry.getJsonValue().get())); - } catch (IOException e) { - throw new JsonParseException("Can't parse jsonValue: " + entry.getJsonValue().get(), e); - } - break; - } + ObjectNode obj = JacksonUtil.newObjectNode(JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER); + obj.put("ts", entry.getTs()); + JacksonUtil.addKvEntry(obj, entry, "value", JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER); return obj; } From 41a5f8d9e4be9a9c47f0508ee44f2be90c4530cd Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Wed, 2 Nov 2022 11:49:05 +0200 Subject: [PATCH 14/16] fix license for JacksonUtil --- .../java/org/thingsboard/common/util/JacksonUtil.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index dca796ee1a..1f21435bd3 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java @@ -1,12 +1,12 @@ /** * Copyright © 2016-2022 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. From e94d305f61c03737065012a31080bd0885cf9277 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Thu, 3 Nov 2022 12:31:31 +0200 Subject: [PATCH 15/16] added JacksonUtilTest for ALLOW_UNQUOTED_FIELD_NAMES_MAPPER and updated getAttributes tests --- .../thingsboard/common/util/JacksonUtil.java | 7 +- .../common/util/JacksonUtilTest.java | 35 +++++++ .../metadata/TbAbstractGetAttributesNode.java | 22 +++-- .../TbAbstractGetAttributesNodeTest.java | 95 ++++++++++++++----- 4 files changed, 122 insertions(+), 37 deletions(-) create mode 100644 common/util/src/test/java/org/thingsboard/common/util/JacksonUtilTest.java diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index 1f21435bd3..f80af0789e 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java @@ -229,6 +229,10 @@ public class JacksonUtil { } } + public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry) { + addKvEntry(entityNode, kvEntry, kvEntry.getKey()); + } + public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry, String key) { addKvEntry(entityNode, kvEntry, key, OBJECT_MAPPER); } @@ -249,7 +253,4 @@ public class JacksonUtil { } } - public static void addKvEntry(ObjectNode entityNode, KvEntry kvEntry) { - addKvEntry(entityNode, kvEntry, kvEntry.getKey()); - } } diff --git a/common/util/src/test/java/org/thingsboard/common/util/JacksonUtilTest.java b/common/util/src/test/java/org/thingsboard/common/util/JacksonUtilTest.java new file mode 100644 index 0000000000..22cd1b3ff6 --- /dev/null +++ b/common/util/src/test/java/org/thingsboard/common/util/JacksonUtilTest.java @@ -0,0 +1,35 @@ +/** + * Copyright © 2016-2022 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.common.util; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; +import org.junit.Assert; +import org.junit.Test; + +public class JacksonUtilTest { + + @Test + public void allow_unquoted_field_mapper_test() { + String data = "{data: 123}"; + JsonNode actualResult = JacksonUtil.toJsonNode(data, JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER); // should be: {"data": 123} + ObjectNode expectedResult = JacksonUtil.newObjectNode(); + expectedResult.put("data", 123); // {"data": 123} + Assert.assertEquals(expectedResult, actualResult); + Assert.assertThrows(IllegalArgumentException.class, () -> JacksonUtil.toJsonNode(data)); // syntax exception due to missing quotes in the field name! + } + +} \ No newline at end of file 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 1b94bbfd2b..32f2d16484 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 @@ -16,6 +16,7 @@ package org.thingsboard.rule.engine.metadata; import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; @@ -91,16 +92,19 @@ public abstract class TbAbstractGetAttributesNode> failuresMap = new ConcurrentHashMap<>(); ListenableFuture>>> allFutures = Futures.allAsList( - getLatestTelemetry(ctx, entityId, msg, LATEST_TS, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresMap), + getLatestTelemetry(ctx, entityId, TbNodeUtils.processPatterns(config.getLatestTsKeyNames(), msg), failuresMap), getAttrAsync(ctx, entityId, CLIENT_SCOPE, TbNodeUtils.processPatterns(config.getClientAttributeNames(), msg), failuresMap), getAttrAsync(ctx, entityId, SHARED_SCOPE, TbNodeUtils.processPatterns(config.getSharedAttributeNames(), msg), failuresMap), getAttrAsync(ctx, entityId, SERVER_SCOPE, TbNodeUtils.processPatterns(config.getServerAttributeNames(), msg), failuresMap) @@ -114,10 +118,11 @@ public abstract class TbAbstractGetAttributesNode { String prefix = getPrefix(keyScope); kvEntryList.forEach(kvEntry -> { + String key = prefix + kvEntry.getKey(); if (fetchToData) { - JacksonUtil.addKvEntry((ObjectNode) msgDataNode, kvEntry, prefix + kvEntry.getKey(), JacksonUtil.ALLOW_UNQUOTED_FIELD_NAMES_MAPPER); + JacksonUtil.addKvEntry((ObjectNode) msgDataNode, kvEntry, key); } else { - msgMetaData.putValue(prefix + kvEntry.getKey(), kvEntry.getValueAsString()); + msgMetaData.putValue(key, kvEntry.getValueAsString()); } }); }); @@ -145,7 +150,7 @@ public abstract class TbAbstractGetAttributesNode>> getLatestTelemetry(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, ConcurrentHashMap> failuresMap) { + private ListenableFuture>> getLatestTelemetry(TbContext ctx, EntityId entityId, List keys, ConcurrentHashMap> failuresMap) { if (CollectionUtils.isEmpty(keys)) { return Futures.immediateFuture(null); } @@ -155,7 +160,7 @@ public abstract class TbAbstractGetAttributesNode { if (tsKvEntry.getValue() == null) { if (isTellFailureIfAbsent) { - computeFailuresMap(scope, failuresMap, tsKvEntry.getKey()); + computeFailuresMap(LATEST_TS, failuresMap, tsKvEntry.getKey()); } } else if (getLatestValueWithTs) { listTsKvEntry.add(getValueWithTs(tsKvEntry)); @@ -164,7 +169,7 @@ public abstract class TbAbstractGetAttributesNode> mapTsKvEntry = new HashMap<>(); - mapTsKvEntry.put(scope, listTsKvEntry); + mapTsKvEntry.put(LATEST_TS, listTsKvEntry); return mapTsKvEntry; }, MoreExecutors.directExecutor()); } @@ -172,7 +177,8 @@ public abstract class TbAbstractGetAttributesNode serverAttributes; private List sharedAttributes; private List tsKeys; + private long ts; @Before public void before() throws TbNodeException { @@ -104,20 +106,21 @@ public class TbAbstractGetAttributesNodeTest { serverAttributes = getAttributeNames("server"); sharedAttributes = getAttributeNames("shared"); tsKeys = List.of("temperature", "humidity", "unknown"); + ts = System.currentTimeMillis(); Mockito.when(attributesService.find(tenantId, originator, DataConstants.CLIENT_SCOPE, clientAttributes)) - .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(clientAttributes))); + .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(clientAttributes, ts))); Mockito.when(attributesService.find(tenantId, originator, DataConstants.SERVER_SCOPE, serverAttributes)) - .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(serverAttributes))); + .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(serverAttributes, ts))); Mockito.when(attributesService.find(tenantId, originator, DataConstants.SHARED_SCOPE, sharedAttributes)) - .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(sharedAttributes))); + .thenReturn(Futures.immediateFuture(getListAttributeKvEntry(sharedAttributes, ts))); Mockito.when(tsService.findLatest(tenantId, originator, tsKeys)) - .thenReturn(Futures.immediateFuture(getListTelemetryKvEntry(tsKeys))); + .thenReturn(Futures.immediateFuture(getListTsKvEntry(tsKeys, ts))); } @After @@ -143,8 +146,44 @@ public class TbAbstractGetAttributesNodeTest { checkTs(tsKeys, false, false, msgMetaData, null); } + @Test + public void fetchToMetadata_latestWithTs_whenOnMsg_then_success() throws Exception { + TbGetAttributesNode node = initNode(false, true, false); + TbMsg msg = getTbMsg(originator); + node.onMsg(ctx, msg); + + TbMsg resultMsg = checkMsg(); + TbMsgMetaData msgMetaData = resultMsg.getMetaData(); + + //check attributes + checkAttributes(clientAttributes, "cs_", false, msgMetaData, null); + checkAttributes(serverAttributes, "ss_", false, msgMetaData, null); + checkAttributes(sharedAttributes, "shared_", false, msgMetaData, null); + + //check timeseries with ts + checkTs(tsKeys, false, true, msgMetaData, null); + } + @Test public void fetchToData_whenOnMsg_then_success() throws Exception { + TbGetAttributesNode node = initNode(true, false, false); + TbMsg msg = getTbMsg(originator); + node.onMsg(ctx, msg); + + TbMsg resultMsg = checkMsg(); + JsonNode msgData = JacksonUtil.toJsonNode(resultMsg.getData()); + + //check attributes + checkAttributes(clientAttributes, "cs_", true, null, msgData); + checkAttributes(serverAttributes, "ss_", true, null, msgData); + checkAttributes(sharedAttributes, "shared_", true, null, msgData); + + //check timeseries + checkTs(tsKeys, true, false, null, msgData); + } + + @Test + public void fetchToData_latestWithTs_whenOnMsg_then_success() throws Exception { TbGetAttributesNode node = initNode(true, true, false); TbMsg msg = getTbMsg(originator); node.onMsg(ctx, msg); @@ -218,28 +257,26 @@ public class TbAbstractGetAttributesNodeTest { } private void checkTs(List tsKeys, boolean fetchToData, boolean getLatestValueWithTs, TbMsgMetaData msgMetaData, JsonNode msgData) { - long tsValue = 1L; + long value = 1L; for (String key : tsKeys) { if (key.equals("unknown")) { continue; } - String result; + String actualValue; + String expectedValue; + if (getLatestValueWithTs) { + expectedValue = "{\"ts\":" + ts + ",\"value\":{\"data\":" + value + "}}"; + } else { + expectedValue = "{\"data\":" + value + "}"; + } if (fetchToData) { - if (getLatestValueWithTs) { - JsonNode resultTs = msgData.get(key); - Assert.assertNotNull(resultTs); - Assert.assertNotNull(resultTs.get("value")); - Assert.assertNotNull(resultTs.get("ts")); - result = resultTs.get("value").asText(); - } else { - result = msgData.get(key).asText(); - } + actualValue = JacksonUtil.toString(msgData.get(key)); } else { - result = msgMetaData.getValue(key); + actualValue = msgMetaData.getValue(key); } - Assert.assertNotNull(result); - Assert.assertEquals(String.valueOf(tsValue), result); - tsValue++; + Assert.assertNotNull(actualValue); + Assert.assertEquals(expectedValue, actualValue); + value++; } } @@ -273,20 +310,26 @@ public class TbAbstractGetAttributesNodeTest { return List.of(prefix + "_attr_1", prefix + "_attr_2", prefix + "_attr_3", "unknown"); } - private List getListAttributeKvEntry(List attributes) { - List attributeKvEntries = new ArrayList<>(); - attributes.stream().filter(attribute -> !attribute.equals("unknown")).forEach(attribute -> attributeKvEntries.add(new BaseAttributeKvEntry(System.currentTimeMillis(), new StringDataEntry(attribute, attribute + "_value")))); - return attributeKvEntries; + private List getListAttributeKvEntry(List attributes, long ts) { + return attributes.stream() + .filter(attribute -> !attribute.equals("unknown")) + .map(attribute -> toAttributeKvEntry(ts, attribute)) + .collect(Collectors.toList()); + } + + private BaseAttributeKvEntry toAttributeKvEntry(long ts, String attribute) { + return new BaseAttributeKvEntry(ts, new StringDataEntry(attribute, attribute + "_value")); } - private List getListTelemetryKvEntry(List keys) { + private List getListTsKvEntry(List keys, long ts) { long value = 1L; List kvEntries = new ArrayList<>(); for (String key : keys) { if (key.equals("unknown")) { continue; } - kvEntries.add(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry(key, value))); + String dataValue = "{\"data\":" + value + "}"; + kvEntries.add(new BasicTsKvEntry(ts, new JsonDataEntry(key, dataValue))); value++; } return kvEntries; From a5691f9347b8b54cb8d7df461de3b96d10506fcd Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Thu, 3 Nov 2022 15:58:22 +0200 Subject: [PATCH 16/16] Fix for the attributes mapper --- .../rule/engine/metadata/TbAbstractGetAttributesNode.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 32f2d16484..d9352f6e4e 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 @@ -175,9 +175,9 @@ public abstract class TbAbstractGetAttributesNode