From 6d657fed466d3e46cc1db3eb9f170aeedc31ffb1 Mon Sep 17 00:00:00 2001 From: van-vanich Date: Fri, 12 Nov 2021 14:02:51 +0200 Subject: [PATCH 1/8] add process pattern for rule node: - customer attribute - tenant attribute - related attribute also, refactoring old tests for customer attribute and add tests related and tenant attribute. --- .../engine/metadata/TbEntityGetAttrNode.java | 25 +- .../TbGetCustomerAttributeNodeTest.java | 49 +-- .../TbGetRelatedAttributeNodeTest.java | 299 ++++++++++++++++++ .../TbGetTenantAttributeNodeTest.java | 278 ++++++++++++++++ 4 files changed, 627 insertions(+), 24 deletions(-) create mode 100644 rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java create mode 100644 rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java 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 index 79b2edf856..2143639a2c 100644 --- 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 @@ -30,12 +30,14 @@ 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.ArrayList; +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.rule.engine.api.TbRelationTypes.SUCCESS; import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; @Slf4j @@ -84,10 +86,29 @@ public abstract class TbEntityGetAttrNode implements TbNode private void putAttributesAndTell(TbContext ctx, TbMsg msg, List attributes) { + log.info("attr " + attributes.toString()); + log.info("conf attr " + config.getAttrMapping().toString()); + List attrProcessPattern = new ArrayList<>(); + log.info("msg {}", msg); + log.info("result process {}", attrProcessPattern); + Map updConf = new HashMap<>(); + config.getAttrMapping().forEach((key, value) -> { + String processPattern = TbNodeUtils.processPattern(key, msg); + updConf.put(processPattern, value); + attrProcessPattern.add(processPattern); + }); + attributes.forEach(r -> { - String attrName = config.getAttrMapping().get(r.getKey()); + log.info("r {}", r); + log.info("rkey {}", r.getKey()); + log.info("index {}", attrProcessPattern.indexOf(r.getKey())); + log.info("getByKey {}", updConf.get(r.getKey())); + String attrName = updConf.get(r.getKey()); + log.info("attrName {}", attrName); msg.getMetaData().putValue(attrName, r.getValueAsString()); + log.info(msg.getMetaData().toString()); }); + log.info(msg.toString()); ctx.tellSuccess(msg); } 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 ca88290808..6720be90de 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 @@ -86,20 +86,24 @@ public class TbGetCustomerAttributeNodeTest { private DeviceService deviceService; private TbMsg msg; + private Map metaData; - private RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); - private RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); + private final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); + private final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); @Before public void init() throws TbNodeException { TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); - Map attrMapping = new HashMap<>(); - attrMapping.putIfAbsent("temperature", "tempo"); - config.setAttrMapping(attrMapping); + Map conf = new HashMap<>(); + conf.put("${word}", "result"); + config.setAttrMapping(conf); config.setTelemetry(false); ObjectMapper mapper = new ObjectMapper(); TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + metaData = new HashMap<>(); + metaData.putIfAbsent("word", "temperature"); + node = new TbGetCustomerAttributeNode(); node.init(null, nodeConfiguration); } @@ -111,13 +115,13 @@ public class TbGetCustomerAttributeNodeTest { User user = new User(); user.setCustomerId(customerId); - msg = TbMsg.newMsg( "USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); when(ctx.getUserService()).thenReturn(userService); when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), eq(Collections.singleton("temperature")))) + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) .thenThrow(new IllegalStateException("something wrong")); node.onMsg(ctx, msg); @@ -136,13 +140,13 @@ public class TbGetCustomerAttributeNodeTest { User user = new User(); user.setCustomerId(customerId); - msg = TbMsg.newMsg( "USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); when(ctx.getUserService()).thenReturn(userService); when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), eq(Collections.singleton("temperature")))) + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) .thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); node.onMsg(ctx, msg); @@ -161,7 +165,7 @@ public class TbGetCustomerAttributeNodeTest { User user = new User(); user.setCustomerId(customerId); - msg = TbMsg.newMsg( "USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON,"{}", ruleChainId, ruleNodeId); + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); when(ctx.getUserService()).thenReturn(userService); when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(null)); @@ -175,7 +179,7 @@ public class TbGetCustomerAttributeNodeTest { @Test public void customerAttributeAddedInMetadata() { CustomerId customerId = new CustomerId(Uuids.timeBased()); - msg = TbMsg.newMsg( "CUSTOMER", customerId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + msg = TbMsg.newMsg("CUSTOMER", customerId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); entityAttributeFetched(customerId); } @@ -186,7 +190,7 @@ public class TbGetCustomerAttributeNodeTest { User user = new User(); user.setCustomerId(customerId); - msg = TbMsg.newMsg( "USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); when(ctx.getUserService()).thenReturn(userService); when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); @@ -201,7 +205,7 @@ public class TbGetCustomerAttributeNodeTest { Asset asset = new Asset(); asset.setCustomerId(customerId); - msg = TbMsg.newMsg( "USER", assetId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + msg = TbMsg.newMsg("USER", assetId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); when(ctx.getAssetService()).thenReturn(assetService); when(assetService.findAssetByIdAsync(any(), eq(assetId))).thenReturn(Futures.immediateFuture(asset)); @@ -216,7 +220,7 @@ public class TbGetCustomerAttributeNodeTest { Device device = new Device(); device.setCustomerId(customerId); - msg = TbMsg.newMsg( "USER", deviceId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); when(ctx.getDeviceService()).thenReturn(deviceService); when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); @@ -227,9 +231,10 @@ public class TbGetCustomerAttributeNodeTest { @Test public void deviceCustomerTelemetryFetched() throws TbNodeException { TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); - Map attrMapping = new HashMap<>(); - attrMapping.putIfAbsent("temperature", "tempo"); - config.setAttrMapping(attrMapping); + + Map conf = new HashMap<>(); + conf.put("${word}", "result"); + config.setAttrMapping(conf); config.setTelemetry(true); ObjectMapper mapper = new ObjectMapper(); TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); @@ -243,7 +248,7 @@ public class TbGetCustomerAttributeNodeTest { Device device = new Device(); device.setCustomerId(customerId); - msg = TbMsg.newMsg( "USER", deviceId, new TbMsgMetaData(), TbMsgDataType.JSON,"{}", ruleChainId, ruleNodeId); + msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); when(ctx.getDeviceService()).thenReturn(deviceService); when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); @@ -251,23 +256,23 @@ public class TbGetCustomerAttributeNodeTest { List timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); when(ctx.getTimeseriesService()).thenReturn(timeseriesService); - when(timeseriesService.findLatest(any(), eq(customerId), eq(Collections.singleton("temperature")))) + when(timeseriesService.findLatest(any(), eq(customerId), eq(Collections.singleton("${word}")))) .thenReturn(Futures.immediateFuture(timeseries)); node.onMsg(ctx, msg); verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("tempo"), "highest"); + assertEquals(msg.getMetaData().getValue("result"), "highest"); } private void entityAttributeFetched(CustomerId customerId) { 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), eq(Collections.singleton("temperature")))) + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) .thenReturn(Futures.immediateFuture(attributes)); node.onMsg(ctx, msg); verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("tempo"), "high"); + assertEquals(msg.getMetaData().getValue("result"), "high"); } } 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 new file mode 100644 index 0000000000..820d3118d4 --- /dev/null +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java @@ -0,0 +1,299 @@ +/** + * Copyright © 2016-2021 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.datastax.oss.driver.api.core.uuid.Uuids; +import com.fasterxml.jackson.databind.ObjectMapper; +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.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.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.RuleChainId; +import org.thingsboard.server.common.data.id.RuleNodeId; +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.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +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.asset.AssetService; +import org.thingsboard.server.dao.attributes.AttributesService; +import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.dao.relation.RelationService; +import org.thingsboard.server.dao.timeseries.TimeseriesService; +import org.thingsboard.server.dao.user.UserService; + +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.same; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; +import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; + +@RunWith(MockitoJUnitRunner.Silent.class) +public class TbGetRelatedAttributeNodeTest { + private final CustomerId customerId = new CustomerId(Uuids.timeBased()); + private final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); + private final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); + private TbGetRelatedAttributeNode node; + @Mock + private TbContext ctx; + @Mock + private AttributesService attributesService; + @Mock + private TimeseriesService timeseriesService; + @Mock + private UserService userService; + @Mock + private AssetService assetService; + @Mock + private DeviceService deviceService; + @Mock + private RelationService relationService; + private TbMsg msg; + private Map metaData; + private EntityRelation entityRelation; + + @Before + public void init() throws TbNodeException { + TbGetRelatedAttrNodeConfiguration config = new TbGetRelatedAttrNodeConfiguration(); + config = config.defaultConfiguration(); + Map conf = new HashMap<>(); + conf.put("${word}", "result"); + config.setAttrMapping(conf); + config.setTelemetry(false); + ObjectMapper mapper = new ObjectMapper(); + TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + + metaData = new HashMap<>(); + metaData.putIfAbsent("word", "temperature"); + + entityRelation = new EntityRelation(); + entityRelation.setTo(customerId); + entityRelation.setType(EntityRelation.CONTAINS_TYPE); + when(ctx.getRelationService()).thenReturn(relationService); + + node = new TbGetRelatedAttributeNode(); + node.init(null, nodeConfiguration); + } + + @Test + public void errorThrownIfCannotLoadAttributes() { + UserId userId = new UserId(Uuids.timeBased()); + User user = new User(); + user.setCustomerId(customerId); + entityRelation.setFrom(userId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); + + when(ctx.getAttributesService()).thenReturn(attributesService); + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) + .thenThrow(new IllegalStateException("something wrong")); + + node.onMsg(ctx, msg); + final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); + verify(ctx).tellFailure(same(msg), captor.capture()); + + Throwable value = captor.getValue(); + assertEquals("something wrong", value.getMessage()); + assertTrue(msg.getMetaData().getData().isEmpty()); + } + + @Test + public void errorThrownIfCannotLoadAttributesAsync() { + UserId userId = new UserId(Uuids.timeBased()); + User user = new User(); + user.setCustomerId(customerId); + entityRelation.setFrom(userId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); + + when(ctx.getAttributesService()).thenReturn(attributesService); + + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) + .thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); + + node.onMsg(ctx, msg); + final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); + verify(ctx).tellFailure(same(msg), captor.capture()); + + Throwable value = captor.getValue(); + assertEquals("something wrong", value.getMessage()); + assertTrue(msg.getMetaData().getData().isEmpty()); + } + + @Test + public void failedChainUsedIfCustomerCannotBeFound() { + UserId userId = new UserId(Uuids.timeBased()); + User user = new User(); + user.setCustomerId(customerId); + entityRelation.setFrom(customerId); + entityRelation.setTo(null); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(null)); + + + node.onMsg(ctx, msg); + verify(ctx).tellNext(msg, FAILURE); + assertTrue(msg.getMetaData().getData().isEmpty()); + + entityRelation.setTo(customerId); + } + + @Test + public void customerAttributeAddedInMetadata() { + entityRelation.setFrom(customerId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + msg = TbMsg.newMsg("CUSTOMER", customerId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + entityAttributeFetched(customerId); + } + + @Test + public void usersCustomerAttributesFetched() { + UserId userId = new UserId(Uuids.timeBased()); + User user = new User(); + user.setCustomerId(customerId); + entityRelation.setFrom(userId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); + + entityAttributeFetched(customerId); + } + + @Test + public void assetsCustomerAttributesFetched() { + AssetId assetId = new AssetId(Uuids.timeBased()); + Asset asset = new Asset(); + asset.setCustomerId(customerId); + entityRelation.setFrom(assetId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + + msg = TbMsg.newMsg("USER", assetId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getAssetService()).thenReturn(assetService); + when(assetService.findAssetByIdAsync(any(), eq(assetId))).thenReturn(Futures.immediateFuture(asset)); + + entityAttributeFetched(customerId); + } + + @Test + public void deviceCustomerAttributesFetched() { + DeviceId deviceId = new DeviceId(Uuids.timeBased()); + Device device = new Device(); + device.setCustomerId(customerId); + entityRelation.setFrom(deviceId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + + msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getDeviceService()).thenReturn(deviceService); + when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); + + entityAttributeFetched(customerId); + } + + @Test + public void deviceCustomerTelemetryFetched() throws TbNodeException { + TbGetRelatedAttrNodeConfiguration config = new TbGetRelatedAttrNodeConfiguration(); + config = config.defaultConfiguration(); + + Map conf = new HashMap<>(); + conf.put("${word}", "result"); + config.setAttrMapping(conf); + config.setTelemetry(true); + ObjectMapper mapper = new ObjectMapper(); + TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + + node = new TbGetRelatedAttributeNode(); + node.init(null, nodeConfiguration); + + + DeviceId deviceId = new DeviceId(Uuids.timeBased()); + Device device = new Device(); + device.setCustomerId(customerId); + + entityRelation.setFrom(deviceId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + + msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getDeviceService()).thenReturn(deviceService); + when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); + + List timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); + + when(ctx.getTimeseriesService()).thenReturn(timeseriesService); + when(timeseriesService.findLatest(any(), eq(customerId), eq(Collections.singleton("${word}")))) + .thenReturn(Futures.immediateFuture(timeseries)); + + node.onMsg(ctx, msg); + verify(ctx).tellSuccess(msg); + assertEquals(msg.getMetaData().getValue("result"), "highest"); + } + + private void entityAttributeFetched(CustomerId customerId) { + 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), eq(Collections.singleton("${word}")))) + .thenReturn(Futures.immediateFuture(attributes)); + + node.onMsg(ctx, msg); + verify(ctx).tellSuccess(msg); + assertEquals(msg.getMetaData().getValue("result"), "high"); + } +} \ No newline at end of file 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 new file mode 100644 index 0000000000..1e111d150c --- /dev/null +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java @@ -0,0 +1,278 @@ +/** + * Copyright © 2016-2021 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.datastax.oss.driver.api.core.uuid.Uuids; +import com.fasterxml.jackson.databind.ObjectMapper; +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.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.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.RuleChainId; +import org.thingsboard.server.common.data.id.RuleNodeId; +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.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +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.asset.AssetService; +import org.thingsboard.server.dao.attributes.AttributesService; +import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.dao.timeseries.TimeseriesService; +import org.thingsboard.server.dao.user.UserService; + +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.same; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; +import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; + +@RunWith(MockitoJUnitRunner.Silent.class) +public class TbGetTenantAttributeNodeTest { + private TbGetTenantAttributeNode node; + + @Mock + private TbContext ctx; + + @Mock + private AttributesService attributesService; + @Mock + private TimeseriesService timeseriesService; + @Mock + private UserService userService; + @Mock + private AssetService assetService; + @Mock + private DeviceService deviceService; + + private TbMsg msg; + private Map metaData; + + private final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); + private final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); + + @Before + public void init() throws TbNodeException { + TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); + Map conf = new HashMap<>(); + conf.put("${word}", "result"); + config.setAttrMapping(conf); + config.setTelemetry(false); + ObjectMapper mapper = new ObjectMapper(); + TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + + metaData = new HashMap<>(); + metaData.putIfAbsent("word", "temperature"); + + node = new TbGetTenantAttributeNode(); + node.init(null, nodeConfiguration); + } + + @Test + public void errorThrownIfCannotLoadAttributes() { + UserId userId = new UserId(Uuids.timeBased()); + TenantId tenantId = new TenantId(Uuids.timeBased()); + User user = new User(); + user.setTenantId(tenantId); + + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); + + when(ctx.getAttributesService()).thenReturn(attributesService); + when(attributesService.find(any(), eq(tenantId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) + .thenThrow(new IllegalStateException("something wrong")); + + node.onMsg(ctx, msg); + final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); + verify(ctx).tellFailure(same(msg), captor.capture()); + + Throwable value = captor.getValue(); + assertEquals("something wrong", value.getMessage()); + assertTrue(msg.getMetaData().getData().isEmpty()); + } + + @Test + public void errorThrownIfCannotLoadAttributesAsync() { + UserId userId = new UserId(Uuids.timeBased()); + TenantId tenantId = new TenantId(Uuids.timeBased()); + User user = new User(); + user.setTenantId(tenantId); + + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); + + when(ctx.getAttributesService()).thenReturn(attributesService); + when(attributesService.find(any(), eq(tenantId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) + .thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); + + node.onMsg(ctx, msg); + final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); + verify(ctx).tellFailure(same(msg), captor.capture()); + + Throwable value = captor.getValue(); + assertEquals("something wrong", value.getMessage()); + assertTrue(msg.getMetaData().getData().isEmpty()); + } + + @Test + public void failedChainUsedIfCustomerCannotBeFound() { + UserId userId = new UserId(Uuids.timeBased()); + CustomerId customerId = new CustomerId(Uuids.timeBased()); + User user = new User(); + user.setCustomerId(customerId); + + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(null)); + + + node.onMsg(ctx, msg); + verify(ctx).tellNext(msg, FAILURE); + assertTrue(msg.getMetaData().getData().isEmpty()); + } + + @Test + public void customerAttributeAddedInMetadata() { + TenantId tenantId = new TenantId(Uuids.timeBased()); + msg = TbMsg.newMsg("TENANT", tenantId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + entityAttributeFetched(tenantId); + } + + @Test + public void usersCustomerAttributesFetched() { + UserId userId = new UserId(Uuids.timeBased()); + TenantId tenantId = new TenantId(Uuids.timeBased()); + User user = new User(); + user.setTenantId(tenantId); + + msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); + + entityAttributeFetched(tenantId); + } + + @Test + public void assetsCustomerAttributesFetched() { + AssetId assetId = new AssetId(Uuids.timeBased()); + TenantId tenantId = new TenantId(Uuids.timeBased()); + Asset asset = new Asset(); + asset.setTenantId(tenantId); + + msg = TbMsg.newMsg("USER", assetId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getAssetService()).thenReturn(assetService); + when(assetService.findAssetByIdAsync(any(), eq(assetId))).thenReturn(Futures.immediateFuture(asset)); + + entityAttributeFetched(tenantId); + } + + @Test + public void deviceCustomerAttributesFetched() { + DeviceId deviceId = new DeviceId(Uuids.timeBased()); + TenantId tenantId = new TenantId(Uuids.timeBased()); + Device device = new Device(); + device.setTenantId(tenantId); + + msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getDeviceService()).thenReturn(deviceService); + when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); + + entityAttributeFetched(tenantId); + } + + @Test + public void deviceCustomerTelemetryFetched() throws TbNodeException { + TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); + + Map conf = new HashMap<>(); + conf.put("${word}", "result"); + config.setAttrMapping(conf); + config.setTelemetry(true); + ObjectMapper mapper = new ObjectMapper(); + TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + + node = new TbGetTenantAttributeNode(); + node.init(null, nodeConfiguration); + + + DeviceId deviceId = new DeviceId(Uuids.timeBased()); + TenantId tenantId = new TenantId(Uuids.timeBased()); + Device device = new Device(); + device.setTenantId(tenantId); + + msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getDeviceService()).thenReturn(deviceService); + when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); + + List timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); + + when(ctx.getTimeseriesService()).thenReturn(timeseriesService); + when(timeseriesService.findLatest(any(), eq(tenantId), eq(Collections.singleton("${word}")))) + .thenReturn(Futures.immediateFuture(timeseries)); + + node.onMsg(ctx, msg); + verify(ctx).tellSuccess(msg); + assertEquals(msg.getMetaData().getValue("result"), "highest"); + } + + private void entityAttributeFetched(TenantId customerId) { + 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), eq(Collections.singleton("${word}")))) + .thenReturn(Futures.immediateFuture(attributes)); + + node.onMsg(ctx, msg); + verify(ctx).tellSuccess(msg); + assertEquals(msg.getMetaData().getValue("result"), "high"); + } +} \ No newline at end of file From 3fec3aa21c1e85e52f036734a2a51afe8b1ffb53 Mon Sep 17 00:00:00 2001 From: van-vanich Date: Fri, 12 Nov 2021 15:59:24 +0200 Subject: [PATCH 2/8] refactoring and fix bugs --- .../engine/metadata/TbEntityGetAttrNode.java | 23 ++++++------------- .../TbGetCustomerAttributeNodeTest.java | 18 ++++++--------- .../TbGetRelatedAttributeNodeTest.java | 10 ++++---- .../TbGetTenantAttributeNodeTest.java | 18 ++++++--------- 4 files changed, 26 insertions(+), 43 deletions(-) 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 index 2143639a2c..f2c52b1ba5 100644 --- 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 @@ -67,30 +67,28 @@ public abstract class TbEntityGetAttrNode implements TbNode return; } - withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId) : getAttributesAsync(ctx, entityId), + withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId, msg) : getAttributesAsync(ctx, entityId, msg), attributes -> putAttributesAndTell(ctx, msg, attributes), t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } - private ListenableFuture> getAttributesAsync(TbContext ctx, EntityId entityId) { - ListenableFuture> latest = ctx.getAttributesService().find(ctx.getTenantId(), entityId, SERVER_SCOPE, config.getAttrMapping().keySet()); + private ListenableFuture> getAttributesAsync(TbContext ctx, EntityId entityId, TbMsg msg) { + List attrProcess = TbNodeUtils.processPatterns(List.copyOf(config.getAttrMapping().keySet()), msg); + ListenableFuture> latest = ctx.getAttributesService().find(ctx.getTenantId(), entityId, SERVER_SCOPE, attrProcess); return Futures.transform(latest, l -> l.stream().map(i -> (KvEntry) i).collect(Collectors.toList()), MoreExecutors.directExecutor()); } - private ListenableFuture> getLatestTelemetry(TbContext ctx, EntityId entityId) { - ListenableFuture> latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, config.getAttrMapping().keySet()); + private ListenableFuture> getLatestTelemetry(TbContext ctx, EntityId entityId, TbMsg msg) { + List latestProcess = TbNodeUtils.processPatterns(List.copyOf(config.getAttrMapping().keySet()), msg); + ListenableFuture> latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, latestProcess); 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) { - log.info("attr " + attributes.toString()); - log.info("conf attr " + config.getAttrMapping().toString()); List attrProcessPattern = new ArrayList<>(); - log.info("msg {}", msg); - log.info("result process {}", attrProcessPattern); Map updConf = new HashMap<>(); config.getAttrMapping().forEach((key, value) -> { String processPattern = TbNodeUtils.processPattern(key, msg); @@ -99,16 +97,9 @@ public abstract class TbEntityGetAttrNode implements TbNode }); attributes.forEach(r -> { - log.info("r {}", r); - log.info("rkey {}", r.getKey()); - log.info("index {}", attrProcessPattern.indexOf(r.getKey())); - log.info("getByKey {}", updConf.get(r.getKey())); String attrName = updConf.get(r.getKey()); - log.info("attrName {}", attrName); msg.getMetaData().putValue(attrName, r.getValueAsString()); - log.info(msg.getMetaData().toString()); }); - log.info(msg.toString()); ctx.tellSuccess(msg); } 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 6720be90de..b7db9b6995 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 @@ -51,7 +51,6 @@ import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.user.UserService; -import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -59,6 +58,7 @@ import java.util.Map; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyCollection; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.same; import static org.mockito.Mockito.verify; @@ -69,11 +69,11 @@ import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; @RunWith(MockitoJUnitRunner.class) public class TbGetCustomerAttributeNodeTest { + private final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); + private final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); private TbGetCustomerAttributeNode node; - @Mock private TbContext ctx; - @Mock private AttributesService attributesService; @Mock @@ -84,13 +84,9 @@ public class TbGetCustomerAttributeNodeTest { private AssetService assetService; @Mock private DeviceService deviceService; - private TbMsg msg; private Map metaData; - private final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); - private final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); - @Before public void init() throws TbNodeException { TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); @@ -121,7 +117,7 @@ public class TbGetCustomerAttributeNodeTest { when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) .thenThrow(new IllegalStateException("something wrong")); node.onMsg(ctx, msg); @@ -146,7 +142,7 @@ public class TbGetCustomerAttributeNodeTest { when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) .thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); node.onMsg(ctx, msg); @@ -256,7 +252,7 @@ public class TbGetCustomerAttributeNodeTest { List timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); when(ctx.getTimeseriesService()).thenReturn(timeseriesService); - when(timeseriesService.findLatest(any(), eq(customerId), eq(Collections.singleton("${word}")))) + when(timeseriesService.findLatest(any(), eq(customerId), anyCollection())) .thenReturn(Futures.immediateFuture(timeseries)); node.onMsg(ctx, msg); @@ -268,7 +264,7 @@ public class TbGetCustomerAttributeNodeTest { 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), eq(Collections.singleton("${word}")))) + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) .thenReturn(Futures.immediateFuture(attributes)); node.onMsg(ctx, msg); 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 820d3118d4..35340bed0a 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 @@ -53,7 +53,6 @@ import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.user.UserService; -import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -61,6 +60,7 @@ import java.util.Map; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyCollection; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.same; import static org.mockito.Mockito.verify; @@ -129,7 +129,7 @@ public class TbGetRelatedAttributeNodeTest { when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) .thenThrow(new IllegalStateException("something wrong")); node.onMsg(ctx, msg); @@ -156,7 +156,7 @@ public class TbGetRelatedAttributeNodeTest { when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) .thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); node.onMsg(ctx, msg); @@ -277,7 +277,7 @@ public class TbGetRelatedAttributeNodeTest { List timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); when(ctx.getTimeseriesService()).thenReturn(timeseriesService); - when(timeseriesService.findLatest(any(), eq(customerId), eq(Collections.singleton("${word}")))) + when(timeseriesService.findLatest(any(), eq(customerId), anyCollection())) .thenReturn(Futures.immediateFuture(timeseries)); node.onMsg(ctx, msg); @@ -289,7 +289,7 @@ public class TbGetRelatedAttributeNodeTest { 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), eq(Collections.singleton("${word}")))) + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) .thenReturn(Futures.immediateFuture(attributes)); node.onMsg(ctx, msg); 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 1e111d150c..783d786117 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 @@ -52,7 +52,6 @@ import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.user.UserService; -import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -60,6 +59,7 @@ import java.util.Map; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyCollection; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.same; import static org.mockito.Mockito.verify; @@ -69,11 +69,11 @@ import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; @RunWith(MockitoJUnitRunner.Silent.class) public class TbGetTenantAttributeNodeTest { + private final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); + private final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); private TbGetTenantAttributeNode node; - @Mock private TbContext ctx; - @Mock private AttributesService attributesService; @Mock @@ -84,13 +84,9 @@ public class TbGetTenantAttributeNodeTest { private AssetService assetService; @Mock private DeviceService deviceService; - private TbMsg msg; private Map metaData; - private final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); - private final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); - @Before public void init() throws TbNodeException { TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); @@ -121,7 +117,7 @@ public class TbGetTenantAttributeNodeTest { when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(tenantId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) + when(attributesService.find(any(), eq(tenantId), eq(SERVER_SCOPE), anyCollection())) .thenThrow(new IllegalStateException("something wrong")); node.onMsg(ctx, msg); @@ -146,7 +142,7 @@ public class TbGetTenantAttributeNodeTest { when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(tenantId), eq(SERVER_SCOPE), eq(Collections.singleton("${word}")))) + when(attributesService.find(any(), eq(tenantId), eq(SERVER_SCOPE), anyCollection())) .thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); node.onMsg(ctx, msg); @@ -256,7 +252,7 @@ public class TbGetTenantAttributeNodeTest { List timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); when(ctx.getTimeseriesService()).thenReturn(timeseriesService); - when(timeseriesService.findLatest(any(), eq(tenantId), eq(Collections.singleton("${word}")))) + when(timeseriesService.findLatest(any(), eq(tenantId), anyCollection())) .thenReturn(Futures.immediateFuture(timeseries)); node.onMsg(ctx, msg); @@ -268,7 +264,7 @@ public class TbGetTenantAttributeNodeTest { 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), eq(Collections.singleton("${word}")))) + when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) .thenReturn(Futures.immediateFuture(attributes)); node.onMsg(ctx, msg); From 4f845a82803b42ea6546abb65d9b5dcc8c1dbd06 Mon Sep 17 00:00:00 2001 From: van-vanich Date: Fri, 12 Nov 2021 16:00:55 +0200 Subject: [PATCH 3/8] remove code, which use in hard debugging --- .../thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java | 2 -- 1 file changed, 2 deletions(-) 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 index f2c52b1ba5..c8df7b4a92 100644 --- 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 @@ -88,12 +88,10 @@ public abstract class TbEntityGetAttrNode implements TbNode private void putAttributesAndTell(TbContext ctx, TbMsg msg, List attributes) { - List attrProcessPattern = new ArrayList<>(); Map updConf = new HashMap<>(); config.getAttrMapping().forEach((key, value) -> { String processPattern = TbNodeUtils.processPattern(key, msg); updConf.put(processPattern, value); - attrProcessPattern.add(processPattern); }); attributes.forEach(r -> { From 866238acda28d91c3e5f8b3add9377da8b9868e0 Mon Sep 17 00:00:00 2001 From: van-vanich Date: Mon, 22 Nov 2021 11:45:58 +0200 Subject: [PATCH 4/8] first version refactoring test code --- .../metadata/AbstractAttributeNodeTest.java | 227 +++++++++++++++ .../TbGetCustomerAttributeNodeTest.java | 246 ++++------------ .../TbGetRelatedAttributeNodeTest.java | 269 +++++------------- .../TbGetTenantAttributeNodeTest.java | 248 ++++------------ 4 files changed, 392 insertions(+), 598 deletions(-) create mode 100644 rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/AbstractAttributeNodeTest.java 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/AbstractAttributeNodeTest.java new file mode 100644 index 0000000000..97e9ddf30b --- /dev/null +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/AbstractAttributeNodeTest.java @@ -0,0 +1,227 @@ +/** + * Copyright © 2016-2021 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.datastax.oss.driver.api.core.uuid.Uuids; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.common.collect.Lists; +import com.google.common.util.concurrent.Futures; +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.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.User; +import org.thingsboard.server.common.data.asset.Asset; +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.RuleChainId; +import org.thingsboard.server.common.data.id.RuleNodeId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +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.asset.AssetService; +import org.thingsboard.server.dao.attributes.AttributesService; +import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.dao.timeseries.TimeseriesService; +import org.thingsboard.server.dao.user.UserService; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyCollection; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.same; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; +import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; + +@RunWith(MockitoJUnitRunner.class) +public abstract class AbstractAttributeNodeTest { + final CustomerId customerId = new CustomerId(Uuids.timeBased()); + final TenantId tenantId = new TenantId(Uuids.timeBased()); + final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); + final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); + final String keyAttrConf = "${word}"; + final String valueAttrConf = "result"; + @Mock + TbContext ctx; + @Mock + AttributesService attributesService; + @Mock + TimeseriesService timeseriesService; + @Mock + UserService userService; + @Mock + AssetService assetService; + @Mock + DeviceService deviceService; + TbMsg msg; + Map metaData; + TbEntityGetAttrNode node; + + void init(TbEntityGetAttrNode node) throws TbNodeException { + ObjectMapper mapper = JacksonUtil.OBJECT_MAPPER; + TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(getTbNodeConfig())); + + metaData = new HashMap<>(); + metaData.putIfAbsent("word", "temperature"); + + this.node = node; + this.node.init(null, nodeConfiguration); + } + + void errorThrownIfCannotLoadAttributes(User user) { + msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(user)); + + when(ctx.getAttributesService()).thenReturn(attributesService); + when(attributesService.find(any(), eq(getEntityId()), eq(SERVER_SCOPE), anyCollection())) + .thenThrow(new IllegalStateException("something wrong")); + + node.onMsg(ctx, msg); + final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); + verify(ctx).tellFailure(same(msg), captor.capture()); + + Throwable value = captor.getValue(); + assertEquals("something wrong", value.getMessage()); + assertTrue(msg.getMetaData().getData().isEmpty()); + } + + void errorThrownIfCannotLoadAttributesAsync(User user) { + + msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(user)); + + when(ctx.getAttributesService()).thenReturn(attributesService); + when(attributesService.find(any(), eq(getEntityId()), eq(SERVER_SCOPE), anyCollection())) + .thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); + + node.onMsg(ctx, msg); + final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); + verify(ctx).tellFailure(same(msg), captor.capture()); + + Throwable value = captor.getValue(); + assertEquals("something wrong", value.getMessage()); + assertTrue(msg.getMetaData().getData().isEmpty()); + } + + void failedChainUsedIfCustomerCannotBeFound(User user) { + msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(null)); + + node.onMsg(ctx, msg); + verify(ctx).tellNext(msg, FAILURE); + assertTrue(msg.getMetaData().getData().isEmpty()); + } + + void entityAttributeAddedInMetadata(EntityId entityId, String type) { + msg = TbMsg.newMsg(type, entityId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + entityAttributeFetched(getEntityId()); + } + + void usersCustomerAttributesFetched(User user) { + msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(user)); + + entityAttributeFetched(getEntityId()); + } + + void assetsCustomerAttributesFetched(Asset asset) { + msg = TbMsg.newMsg("ASSET", asset.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getAssetService()).thenReturn(assetService); + when(assetService.findAssetByIdAsync(any(), eq(asset.getId()))).thenReturn(Futures.immediateFuture(asset)); + + entityAttributeFetched(getEntityId()); + } + + void deviceCustomerAttributesFetched(Device device) { + msg = TbMsg.newMsg("USER", device.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getDeviceService()).thenReturn(deviceService); + when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); + + entityAttributeFetched(getEntityId()); + } + + void deviceCustomerTelemetryFetched(Device device) throws TbNodeException { + ObjectMapper mapper = JacksonUtil.OBJECT_MAPPER; + TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(getTbNodeConfigFotTelemetry())); + + TbEntityGetAttrNode node = getEmptyNode(); + node.init(null, nodeConfiguration); + + msg = TbMsg.newMsg("USER", device.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + when(ctx.getDeviceService()).thenReturn(deviceService); + when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); + + List timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); + + when(ctx.getTimeseriesService()).thenReturn(timeseriesService); + when(timeseriesService.findLatest(any(), eq(getEntityId()), anyCollection())) + .thenReturn(Futures.immediateFuture(timeseries)); + + node.onMsg(ctx, msg); + verify(ctx).tellSuccess(msg); + assertEquals(msg.getMetaData().getValue("result"), "highest"); + } + + void entityAttributeFetched(EntityId entityId) { + List attributes = Lists.newArrayList(new BaseAttributeKvEntry(new StringDataEntry("temperature", "high"), 1L)); + + when(ctx.getAttributesService()).thenReturn(attributesService); + when(attributesService.find(any(), eq(entityId), eq(SERVER_SCOPE), anyCollection())) + .thenReturn(Futures.immediateFuture(attributes)); + + node.onMsg(ctx, msg); + verify(ctx).tellSuccess(msg); + assertEquals(msg.getMetaData().getValue("result"), "high"); + } + + protected abstract TbEntityGetAttrNode getEmptyNode(); + + abstract T getTbNodeConfig(); + + abstract T getTbNodeConfigFotTelemetry(); + + abstract EntityId getEntityId(); +} 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 b7db9b6995..3ed2e35702 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 @@ -16,259 +16,109 @@ package org.thingsboard.rule.engine.metadata; import com.datastax.oss.driver.api.core.uuid.Uuids; -import com.fasterxml.jackson.databind.ObjectMapper; -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.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.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.RuleChainId; -import org.thingsboard.server.common.data.id.RuleNodeId; +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.BasicTsKvEntry; -import org.thingsboard.server.common.data.kv.StringDataEntry; -import org.thingsboard.server.common.data.kv.TsKvEntry; -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.asset.AssetService; -import org.thingsboard.server.dao.attributes.AttributesService; -import org.thingsboard.server.dao.device.DeviceService; -import org.thingsboard.server.dao.timeseries.TimeseriesService; -import org.thingsboard.server.dao.user.UserService; import java.util.HashMap; -import java.util.List; import java.util.Map; - -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.anyCollection; -import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.ArgumentMatchers.same; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.when; -import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; -import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; +import java.util.UUID; @RunWith(MockitoJUnitRunner.class) -public class TbGetCustomerAttributeNodeTest { - - private final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); - private final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); - private TbGetCustomerAttributeNode node; - @Mock - private TbContext ctx; - @Mock - private AttributesService attributesService; - @Mock - private TimeseriesService timeseriesService; - @Mock - private UserService userService; - @Mock - private AssetService assetService; - @Mock - private DeviceService deviceService; - private TbMsg msg; - private Map metaData; +public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest { + User user = new User(); + Asset asset = new Asset(); + Device device = new Device(); @Before - public void init() throws TbNodeException { + public void initDataForTests() throws TbNodeException { + init(new TbGetCustomerAttributeNode()); + user.setCustomerId(customerId); + user.setId(new UserId(UUID.randomUUID())); + + asset.setCustomerId(customerId); + asset.setId(new AssetId(UUID.randomUUID())); + + device.setCustomerId(customerId); + device.setId(new DeviceId(Uuids.timeBased())); + } + + @Override + protected TbEntityGetAttrNode getEmptyNode() { + return new TbGetCustomerAttributeNode(); + } + + @Override + T getTbNodeConfig() { TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); Map conf = new HashMap<>(); - conf.put("${word}", "result"); + conf.put(keyAttrConf, valueAttrConf); config.setAttrMapping(conf); config.setTelemetry(false); - ObjectMapper mapper = new ObjectMapper(); - TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + return (T) config; + } - metaData = new HashMap<>(); - metaData.putIfAbsent("word", "temperature"); + @Override + T getTbNodeConfigFotTelemetry() { + TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); + Map conf = new HashMap<>(); + conf.put(keyAttrConf, valueAttrConf); + config.setAttrMapping(conf); + config.setTelemetry(true); + return (T) config; + } - node = new TbGetCustomerAttributeNode(); - node.init(null, nodeConfiguration); + @Override + EntityId getEntityId() { + return customerId; } @Test public void errorThrownIfCannotLoadAttributes() { - UserId userId = new UserId(Uuids.timeBased()); - CustomerId customerId = new CustomerId(Uuids.timeBased()); - User user = new User(); - user.setCustomerId(customerId); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); - - when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) - .thenThrow(new IllegalStateException("something wrong")); - - node.onMsg(ctx, msg); - final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); - verify(ctx).tellFailure(same(msg), captor.capture()); - - Throwable value = captor.getValue(); - assertEquals("something wrong", value.getMessage()); - assertTrue(msg.getMetaData().getData().isEmpty()); + errorThrownIfCannotLoadAttributes(user); } @Test public void errorThrownIfCannotLoadAttributesAsync() { - UserId userId = new UserId(Uuids.timeBased()); - CustomerId customerId = new CustomerId(Uuids.timeBased()); - User user = new User(); - user.setCustomerId(customerId); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); - - when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) - .thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); - - node.onMsg(ctx, msg); - final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); - verify(ctx).tellFailure(same(msg), captor.capture()); - - Throwable value = captor.getValue(); - assertEquals("something wrong", value.getMessage()); - assertTrue(msg.getMetaData().getData().isEmpty()); + errorThrownIfCannotLoadAttributesAsync(user); } @Test public void failedChainUsedIfCustomerCannotBeFound() { - UserId userId = new UserId(Uuids.timeBased()); - CustomerId customerId = new CustomerId(Uuids.timeBased()); - User user = new User(); - user.setCustomerId(customerId); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(null)); - - - node.onMsg(ctx, msg); - verify(ctx).tellNext(msg, FAILURE); - assertTrue(msg.getMetaData().getData().isEmpty()); + failedChainUsedIfCustomerCannotBeFound(user); } @Test public void customerAttributeAddedInMetadata() { - CustomerId customerId = new CustomerId(Uuids.timeBased()); - msg = TbMsg.newMsg("CUSTOMER", customerId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - entityAttributeFetched(customerId); + entityAttributeAddedInMetadata(customerId, "CUSTOMER"); } @Test public void usersCustomerAttributesFetched() { - UserId userId = new UserId(Uuids.timeBased()); - CustomerId customerId = new CustomerId(Uuids.timeBased()); - User user = new User(); - user.setCustomerId(customerId); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); - - entityAttributeFetched(customerId); + usersCustomerAttributesFetched(user); } @Test public void assetsCustomerAttributesFetched() { - AssetId assetId = new AssetId(Uuids.timeBased()); - CustomerId customerId = new CustomerId(Uuids.timeBased()); - Asset asset = new Asset(); - asset.setCustomerId(customerId); - - msg = TbMsg.newMsg("USER", assetId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getAssetService()).thenReturn(assetService); - when(assetService.findAssetByIdAsync(any(), eq(assetId))).thenReturn(Futures.immediateFuture(asset)); - - entityAttributeFetched(customerId); + assetsCustomerAttributesFetched(asset); } @Test public void deviceCustomerAttributesFetched() { - DeviceId deviceId = new DeviceId(Uuids.timeBased()); - CustomerId customerId = new CustomerId(Uuids.timeBased()); - Device device = new Device(); - device.setCustomerId(customerId); - - msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getDeviceService()).thenReturn(deviceService); - when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); - - entityAttributeFetched(customerId); + deviceCustomerAttributesFetched(device); } @Test public void deviceCustomerTelemetryFetched() throws TbNodeException { - TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); - - Map conf = new HashMap<>(); - conf.put("${word}", "result"); - config.setAttrMapping(conf); - config.setTelemetry(true); - ObjectMapper mapper = new ObjectMapper(); - TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); - - node = new TbGetCustomerAttributeNode(); - node.init(null, nodeConfiguration); - - - DeviceId deviceId = new DeviceId(Uuids.timeBased()); - CustomerId customerId = new CustomerId(Uuids.timeBased()); - Device device = new Device(); - device.setCustomerId(customerId); - - msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getDeviceService()).thenReturn(deviceService); - when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); - - List timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); - - when(ctx.getTimeseriesService()).thenReturn(timeseriesService); - when(timeseriesService.findLatest(any(), eq(customerId), anyCollection())) - .thenReturn(Futures.immediateFuture(timeseries)); - - node.onMsg(ctx, msg); - verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("result"), "highest"); - } - - private void entityAttributeFetched(CustomerId customerId) { - 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); - verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("result"), "high"); + deviceCustomerTelemetryFetched(device); } } 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 35340bed0a..abb647c78c 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 @@ -15,285 +15,152 @@ */ package org.thingsboard.rule.engine.metadata; -import com.datastax.oss.driver.api.core.uuid.Uuids; -import com.fasterxml.jackson.databind.ObjectMapper; -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.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.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.RuleChainId; -import org.thingsboard.server.common.data.id.RuleNodeId; +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.BasicTsKvEntry; -import org.thingsboard.server.common.data.kv.StringDataEntry; -import org.thingsboard.server.common.data.kv.TsKvEntry; 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.asset.AssetService; -import org.thingsboard.server.dao.attributes.AttributesService; -import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.relation.RelationService; -import org.thingsboard.server.dao.timeseries.TimeseriesService; -import org.thingsboard.server.dao.user.UserService; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.UUID; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.anyCollection; -import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.ArgumentMatchers.same; -import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; -import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; -import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; @RunWith(MockitoJUnitRunner.Silent.class) -public class TbGetRelatedAttributeNodeTest { - private final CustomerId customerId = new CustomerId(Uuids.timeBased()); - private final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); - private final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); - private TbGetRelatedAttributeNode node; - @Mock - private TbContext ctx; - @Mock - private AttributesService attributesService; - @Mock - private TimeseriesService timeseriesService; - @Mock - private UserService userService; - @Mock - private AssetService assetService; - @Mock - private DeviceService deviceService; +public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest { + User user = new User(); + Asset asset = new Asset(); + Device device = new Device(); @Mock private RelationService relationService; - private TbMsg msg; - private Map metaData; private EntityRelation entityRelation; @Before - public void init() throws TbNodeException { + public void initDataForTests() throws TbNodeException { + init(new TbGetRelatedAttributeNode()); + entityRelation = new EntityRelation(); + entityRelation.setTo(customerId); + entityRelation.setType(EntityRelation.CONTAINS_TYPE); + when(ctx.getRelationService()).thenReturn(relationService); + + user.setCustomerId(customerId); + user.setId(new UserId(UUID.randomUUID())); + entityRelation.setFrom(user.getId()); + + asset.setCustomerId(customerId); + asset.setId(new AssetId(UUID.randomUUID())); + + device.setCustomerId(customerId); + device.setId(new DeviceId(UUID.randomUUID())); + } + + @Override + protected TbEntityGetAttrNode getEmptyNode() { + return new TbGetRelatedAttributeNode(); + } + + @Override + T getTbNodeConfig() { TbGetRelatedAttrNodeConfiguration config = new TbGetRelatedAttrNodeConfiguration(); config = config.defaultConfiguration(); Map conf = new HashMap<>(); - conf.put("${word}", "result"); + conf.put(keyAttrConf, valueAttrConf); config.setAttrMapping(conf); config.setTelemetry(false); - ObjectMapper mapper = new ObjectMapper(); - TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); - - metaData = new HashMap<>(); - metaData.putIfAbsent("word", "temperature"); + return (T) config; + } - entityRelation = new EntityRelation(); - entityRelation.setTo(customerId); - entityRelation.setType(EntityRelation.CONTAINS_TYPE); - when(ctx.getRelationService()).thenReturn(relationService); + @Override + T getTbNodeConfigFotTelemetry() { + TbGetRelatedAttrNodeConfiguration config = new TbGetRelatedAttrNodeConfiguration(); + config = config.defaultConfiguration(); + Map conf = new HashMap<>(); + conf.put(keyAttrConf, valueAttrConf); + config.setAttrMapping(conf); + config.setTelemetry(true); + return (T) config; + } - node = new TbGetRelatedAttributeNode(); - node.init(null, nodeConfiguration); + @Override + EntityId getEntityId() { + return customerId; } @Test public void errorThrownIfCannotLoadAttributes() { - UserId userId = new UserId(Uuids.timeBased()); - User user = new User(); - user.setCustomerId(customerId); - entityRelation.setFrom(userId); + entityRelation.setFrom(user.getId()); + entityRelation.setTo(customerId); when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); - - when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) - .thenThrow(new IllegalStateException("something wrong")); - - node.onMsg(ctx, msg); - final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); - verify(ctx).tellFailure(same(msg), captor.capture()); - - Throwable value = captor.getValue(); - assertEquals("something wrong", value.getMessage()); - assertTrue(msg.getMetaData().getData().isEmpty()); + errorThrownIfCannotLoadAttributes(user); } @Test public void errorThrownIfCannotLoadAttributesAsync() { - UserId userId = new UserId(Uuids.timeBased()); - User user = new User(); - user.setCustomerId(customerId); - entityRelation.setFrom(userId); + entityRelation.setFrom(user.getId()); + entityRelation.setTo(customerId); when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); - - when(ctx.getAttributesService()).thenReturn(attributesService); - - when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) - .thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); - - node.onMsg(ctx, msg); - final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); - verify(ctx).tellFailure(same(msg), captor.capture()); - - Throwable value = captor.getValue(); - assertEquals("something wrong", value.getMessage()); - assertTrue(msg.getMetaData().getData().isEmpty()); + errorThrownIfCannotLoadAttributesAsync(user); } @Test public void failedChainUsedIfCustomerCannotBeFound() { - UserId userId = new UserId(Uuids.timeBased()); - User user = new User(); - user.setCustomerId(customerId); entityRelation.setFrom(customerId); entityRelation.setTo(null); when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(null)); - - - node.onMsg(ctx, msg); - verify(ctx).tellNext(msg, FAILURE); - assertTrue(msg.getMetaData().getData().isEmpty()); - - entityRelation.setTo(customerId); + failedChainUsedIfCustomerCannotBeFound(user); } @Test public void customerAttributeAddedInMetadata() { 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); - entityAttributeFetched(customerId); + entityAttributeAddedInMetadata(customerId, "CUSTOMER"); } @Test public void usersCustomerAttributesFetched() { - UserId userId = new UserId(Uuids.timeBased()); - User user = new User(); - user.setCustomerId(customerId); - entityRelation.setFrom(userId); + entityRelation.setFrom(user.getId()); + entityRelation.setTo(customerId); when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); - - entityAttributeFetched(customerId); + usersCustomerAttributesFetched(user); } @Test public void assetsCustomerAttributesFetched() { - AssetId assetId = new AssetId(Uuids.timeBased()); - Asset asset = new Asset(); - asset.setCustomerId(customerId); - entityRelation.setFrom(assetId); + entityRelation.setFrom(asset.getId()); + entityRelation.setTo(customerId); when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); - - msg = TbMsg.newMsg("USER", assetId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getAssetService()).thenReturn(assetService); - when(assetService.findAssetByIdAsync(any(), eq(assetId))).thenReturn(Futures.immediateFuture(asset)); - - entityAttributeFetched(customerId); + assetsCustomerAttributesFetched(asset); } @Test public void deviceCustomerAttributesFetched() { - DeviceId deviceId = new DeviceId(Uuids.timeBased()); - Device device = new Device(); - device.setCustomerId(customerId); - entityRelation.setFrom(deviceId); + entityRelation.setFrom(device.getId()); + entityRelation.setTo(customerId); when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); - - msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getDeviceService()).thenReturn(deviceService); - when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); - - entityAttributeFetched(customerId); + deviceCustomerAttributesFetched(device); } @Test public void deviceCustomerTelemetryFetched() throws TbNodeException { - TbGetRelatedAttrNodeConfiguration config = new TbGetRelatedAttrNodeConfiguration(); - config = config.defaultConfiguration(); - - Map conf = new HashMap<>(); - conf.put("${word}", "result"); - config.setAttrMapping(conf); - config.setTelemetry(true); - ObjectMapper mapper = new ObjectMapper(); - TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); - - node = new TbGetRelatedAttributeNode(); - node.init(null, nodeConfiguration); - - - DeviceId deviceId = new DeviceId(Uuids.timeBased()); - Device device = new Device(); - device.setCustomerId(customerId); - - entityRelation.setFrom(deviceId); + entityRelation.setFrom(device.getId()); + entityRelation.setTo(customerId); when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); - - msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getDeviceService()).thenReturn(deviceService); - when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); - - List timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); - - when(ctx.getTimeseriesService()).thenReturn(timeseriesService); - when(timeseriesService.findLatest(any(), eq(customerId), anyCollection())) - .thenReturn(Futures.immediateFuture(timeseries)); - - node.onMsg(ctx, msg); - verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("result"), "highest"); - } - - private void entityAttributeFetched(CustomerId customerId) { - 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); - verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("result"), "high"); + deviceCustomerTelemetryFetched(device); } } \ No newline at end of file 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 783d786117..63cee9fb46 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 @@ -15,260 +15,110 @@ */ package org.thingsboard.rule.engine.metadata; -import com.datastax.oss.driver.api.core.uuid.Uuids; -import com.fasterxml.jackson.databind.ObjectMapper; -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.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.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.RuleChainId; -import org.thingsboard.server.common.data.id.RuleNodeId; -import org.thingsboard.server.common.data.id.TenantId; +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.BasicTsKvEntry; -import org.thingsboard.server.common.data.kv.StringDataEntry; -import org.thingsboard.server.common.data.kv.TsKvEntry; -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.asset.AssetService; -import org.thingsboard.server.dao.attributes.AttributesService; -import org.thingsboard.server.dao.device.DeviceService; -import org.thingsboard.server.dao.timeseries.TimeseriesService; -import org.thingsboard.server.dao.user.UserService; import java.util.HashMap; -import java.util.List; import java.util.Map; +import java.util.UUID; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.anyCollection; -import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.ArgumentMatchers.same; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.when; -import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; -import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; +@RunWith(MockitoJUnitRunner.class) +public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest { -@RunWith(MockitoJUnitRunner.Silent.class) -public class TbGetTenantAttributeNodeTest { - private final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); - private final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); - private TbGetTenantAttributeNode node; - @Mock - private TbContext ctx; - @Mock - private AttributesService attributesService; - @Mock - private TimeseriesService timeseriesService; - @Mock - private UserService userService; - @Mock - private AssetService assetService; - @Mock - private DeviceService deviceService; - private TbMsg msg; - private Map metaData; + User user = new User(); + Asset asset = new Asset(); + Device device = new Device(); @Before - public void init() throws TbNodeException { + public void initDataForTests() throws TbNodeException { + init(new TbGetTenantAttributeNode()); + user.setTenantId(tenantId); + user.setId(new UserId(UUID.randomUUID())); + + asset.setTenantId(tenantId); + asset.setId(new AssetId(UUID.randomUUID())); + + device.setTenantId(tenantId); + device.setId(new DeviceId(UUID.randomUUID())); + } + + @Override + protected TbEntityGetAttrNode getEmptyNode() { + return new TbGetTenantAttributeNode(); + } + + @Override + T getTbNodeConfig() { TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); Map conf = new HashMap<>(); - conf.put("${word}", "result"); + conf.put(keyAttrConf, valueAttrConf); config.setAttrMapping(conf); config.setTelemetry(false); - ObjectMapper mapper = new ObjectMapper(); - TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + return (T) config; + } - metaData = new HashMap<>(); - metaData.putIfAbsent("word", "temperature"); + @Override + T getTbNodeConfigFotTelemetry() { + TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); + Map conf = new HashMap<>(); + conf.put(keyAttrConf, valueAttrConf); + config.setAttrMapping(conf); + config.setTelemetry(true); + return (T) config; + } - node = new TbGetTenantAttributeNode(); - node.init(null, nodeConfiguration); + @Override + EntityId getEntityId() { + return tenantId; } @Test public void errorThrownIfCannotLoadAttributes() { - UserId userId = new UserId(Uuids.timeBased()); - TenantId tenantId = new TenantId(Uuids.timeBased()); - User user = new User(); - user.setTenantId(tenantId); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); - - when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(tenantId), eq(SERVER_SCOPE), anyCollection())) - .thenThrow(new IllegalStateException("something wrong")); - - node.onMsg(ctx, msg); - final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); - verify(ctx).tellFailure(same(msg), captor.capture()); - - Throwable value = captor.getValue(); - assertEquals("something wrong", value.getMessage()); - assertTrue(msg.getMetaData().getData().isEmpty()); + errorThrownIfCannotLoadAttributes(user); } @Test public void errorThrownIfCannotLoadAttributesAsync() { - UserId userId = new UserId(Uuids.timeBased()); - TenantId tenantId = new TenantId(Uuids.timeBased()); - User user = new User(); - user.setTenantId(tenantId); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); - - when(ctx.getAttributesService()).thenReturn(attributesService); - when(attributesService.find(any(), eq(tenantId), eq(SERVER_SCOPE), anyCollection())) - .thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); - - node.onMsg(ctx, msg); - final ArgumentCaptor captor = ArgumentCaptor.forClass(Throwable.class); - verify(ctx).tellFailure(same(msg), captor.capture()); - - Throwable value = captor.getValue(); - assertEquals("something wrong", value.getMessage()); - assertTrue(msg.getMetaData().getData().isEmpty()); + errorThrownIfCannotLoadAttributesAsync(user); } @Test public void failedChainUsedIfCustomerCannotBeFound() { - UserId userId = new UserId(Uuids.timeBased()); - CustomerId customerId = new CustomerId(Uuids.timeBased()); - User user = new User(); - user.setCustomerId(customerId); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(null)); - - - node.onMsg(ctx, msg); - verify(ctx).tellNext(msg, FAILURE); - assertTrue(msg.getMetaData().getData().isEmpty()); + failedChainUsedIfCustomerCannotBeFound(user); } @Test public void customerAttributeAddedInMetadata() { - TenantId tenantId = new TenantId(Uuids.timeBased()); - msg = TbMsg.newMsg("TENANT", tenantId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - entityAttributeFetched(tenantId); + entityAttributeAddedInMetadata(tenantId, "TENANT"); } @Test public void usersCustomerAttributesFetched() { - UserId userId = new UserId(Uuids.timeBased()); - TenantId tenantId = new TenantId(Uuids.timeBased()); - User user = new User(); - user.setTenantId(tenantId); - - msg = TbMsg.newMsg("USER", userId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); - - entityAttributeFetched(tenantId); + usersCustomerAttributesFetched(user); } @Test public void assetsCustomerAttributesFetched() { - AssetId assetId = new AssetId(Uuids.timeBased()); - TenantId tenantId = new TenantId(Uuids.timeBased()); - Asset asset = new Asset(); - asset.setTenantId(tenantId); - - msg = TbMsg.newMsg("USER", assetId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getAssetService()).thenReturn(assetService); - when(assetService.findAssetByIdAsync(any(), eq(assetId))).thenReturn(Futures.immediateFuture(asset)); - - entityAttributeFetched(tenantId); + assetsCustomerAttributesFetched(asset); } @Test public void deviceCustomerAttributesFetched() { - DeviceId deviceId = new DeviceId(Uuids.timeBased()); - TenantId tenantId = new TenantId(Uuids.timeBased()); - Device device = new Device(); - device.setTenantId(tenantId); - - msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getDeviceService()).thenReturn(deviceService); - when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); - - entityAttributeFetched(tenantId); + deviceCustomerAttributesFetched(device); } @Test public void deviceCustomerTelemetryFetched() throws TbNodeException { - TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); - - Map conf = new HashMap<>(); - conf.put("${word}", "result"); - config.setAttrMapping(conf); - config.setTelemetry(true); - ObjectMapper mapper = new ObjectMapper(); - TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); - - node = new TbGetTenantAttributeNode(); - node.init(null, nodeConfiguration); - - - DeviceId deviceId = new DeviceId(Uuids.timeBased()); - TenantId tenantId = new TenantId(Uuids.timeBased()); - Device device = new Device(); - device.setTenantId(tenantId); - - msg = TbMsg.newMsg("USER", deviceId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getDeviceService()).thenReturn(deviceService); - when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); - - List timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); - - when(ctx.getTimeseriesService()).thenReturn(timeseriesService); - when(timeseriesService.findLatest(any(), eq(tenantId), anyCollection())) - .thenReturn(Futures.immediateFuture(timeseries)); - - node.onMsg(ctx, msg); - verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("result"), "highest"); - } - - private void entityAttributeFetched(TenantId customerId) { - 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); - verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("result"), "high"); + deviceCustomerTelemetryFetched(device); } } \ No newline at end of file From a0e06abe1f5b6d1c6168a92580c5cfc3f9e042b2 Mon Sep 17 00:00:00 2001 From: van-vanich Date: Mon, 22 Nov 2021 13:55:24 +0200 Subject: [PATCH 5/8] refactoring tests and some program --- .../rule/engine/metadata/TbEntityGetAttrNode.java | 9 ++++----- .../engine/metadata/TbGetRelatedAttributeNodeTest.java | 2 +- .../engine/metadata/TbGetTenantAttributeNodeTest.java | 2 +- 3 files changed, 6 insertions(+), 7 deletions(-) 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 index c8df7b4a92..757d0d7ca3 100644 --- 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 @@ -67,20 +67,19 @@ public abstract class TbEntityGetAttrNode implements TbNode return; } - withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId, msg) : getAttributesAsync(ctx, entityId, msg), + List latestProcess = TbNodeUtils.processPatterns(List.copyOf(config.getAttrMapping().keySet()), msg); + withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId, latestProcess) : getAttributesAsync(ctx, entityId, latestProcess), attributes -> putAttributesAndTell(ctx, msg, attributes), t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } - private ListenableFuture> getAttributesAsync(TbContext ctx, EntityId entityId, TbMsg msg) { - List attrProcess = TbNodeUtils.processPatterns(List.copyOf(config.getAttrMapping().keySet()), msg); + private ListenableFuture> getAttributesAsync(TbContext ctx, EntityId entityId, List attrProcess) { ListenableFuture> latest = ctx.getAttributesService().find(ctx.getTenantId(), entityId, SERVER_SCOPE, attrProcess); return Futures.transform(latest, l -> l.stream().map(i -> (KvEntry) i).collect(Collectors.toList()), MoreExecutors.directExecutor()); } - private ListenableFuture> getLatestTelemetry(TbContext ctx, EntityId entityId, TbMsg msg) { - List latestProcess = TbNodeUtils.processPatterns(List.copyOf(config.getAttrMapping().keySet()), msg); + private ListenableFuture> getLatestTelemetry(TbContext ctx, EntityId entityId, List latestProcess) { ListenableFuture> latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, latestProcess); return Futures.transform(latest, l -> l.stream().map(i -> (KvEntry) i).collect(Collectors.toList()), MoreExecutors.directExecutor()); 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 abb647c78c..d2067ca893 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 @@ -163,4 +163,4 @@ public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest { when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); deviceCustomerTelemetryFetched(device); } -} \ No newline at end of file +} 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 63cee9fb46..e297722928 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 @@ -121,4 +121,4 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest { public void deviceCustomerTelemetryFetched() throws TbNodeException { deviceCustomerTelemetryFetched(device); } -} \ No newline at end of file +} From 436e19d6524532765edb063f83115c9861fcdba8 Mon Sep 17 00:00:00 2001 From: van-vanich Date: Tue, 23 Nov 2021 16:14:22 +0200 Subject: [PATCH 6/8] refactoring and improve code --- .../engine/metadata/TbEntityGetAttrNode.java | 31 +++++---- .../metadata/AbstractAttributeNodeTest.java | 64 +++++++++++-------- .../TbGetCustomerAttributeNodeTest.java | 35 ++++------ .../TbGetRelatedAttributeNodeTest.java | 24 +++---- .../TbGetTenantAttributeNodeTest.java | 35 ++++------ 5 files changed, 89 insertions(+), 100 deletions(-) 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 index 757d0d7ca3..747549cd77 100644 --- 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 @@ -30,7 +30,6 @@ 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.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -67,34 +66,34 @@ public abstract class TbEntityGetAttrNode implements TbNode return; } - List latestProcess = TbNodeUtils.processPatterns(List.copyOf(config.getAttrMapping().keySet()), msg); - withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId, latestProcess) : getAttributesAsync(ctx, entityId, latestProcess), - attributes -> putAttributesAndTell(ctx, msg, attributes), + Map mappingsMap = new HashMap<>(); + config.getAttrMapping().forEach((key, value) -> { + String processPattern = TbNodeUtils.processPattern(key, msg); + mappingsMap.put(processPattern, value); + }); + + 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 attrProcess) { - ListenableFuture> latest = ctx.getAttributesService().find(ctx.getTenantId(), entityId, SERVER_SCOPE, attrProcess); + 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 latestProcess) { - ListenableFuture> latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, latestProcess); + 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 updConf = new HashMap<>(); - config.getAttrMapping().forEach((key, value) -> { - String processPattern = TbNodeUtils.processPattern(key, msg); - updConf.put(processPattern, value); - }); - + private void putAttributesAndTell(TbContext ctx, TbMsg msg, List attributes, Map map) { attributes.forEach(r -> { - String attrName = updConf.get(r.getKey()); + String attrName = map.get(r.getKey()); msg.getMetaData().putValue(attrName, r.getValueAsString()); }); ctx.tellSuccess(msg); 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/AbstractAttributeNodeTest.java index 97e9ddf30b..005b40b4e7 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/AbstractAttributeNodeTest.java @@ -19,6 +19,7 @@ import com.datastax.oss.driver.api.core.uuid.Uuids; import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.collect.Lists; import com.google.common.util.concurrent.Futures; +import org.jetbrains.annotations.NotNull; import org.junit.runner.RunWith; import org.mockito.ArgumentCaptor; import org.mockito.Mock; @@ -59,7 +60,6 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyCollection; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.same; -import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; @@ -103,9 +103,6 @@ public abstract class AbstractAttributeNodeTest { void errorThrownIfCannotLoadAttributes(User user) { msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(user)); - when(ctx.getAttributesService()).thenReturn(attributesService); when(attributesService.find(any(), eq(getEntityId()), eq(SERVER_SCOPE), anyCollection())) .thenThrow(new IllegalStateException("something wrong")); @@ -123,9 +120,6 @@ public abstract class AbstractAttributeNodeTest { msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(user)); - when(ctx.getAttributesService()).thenReturn(attributesService); when(attributesService.find(any(), eq(getEntityId()), eq(SERVER_SCOPE), anyCollection())) .thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); @@ -142,9 +136,6 @@ public abstract class AbstractAttributeNodeTest { void failedChainUsedIfCustomerCannotBeFound(User user) { msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(null)); - node.onMsg(ctx, msg); verify(ctx).tellNext(msg, FAILURE); assertTrue(msg.getMetaData().getData().isEmpty()); @@ -158,41 +149,29 @@ public abstract class AbstractAttributeNodeTest { void usersCustomerAttributesFetched(User user) { msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(user)); - entityAttributeFetched(getEntityId()); } void assetsCustomerAttributesFetched(Asset asset) { msg = TbMsg.newMsg("ASSET", asset.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - when(ctx.getAssetService()).thenReturn(assetService); - when(assetService.findAssetByIdAsync(any(), eq(asset.getId()))).thenReturn(Futures.immediateFuture(asset)); - entityAttributeFetched(getEntityId()); } void deviceCustomerAttributesFetched(Device device) { - msg = TbMsg.newMsg("USER", device.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getDeviceService()).thenReturn(deviceService); - when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); + msg = TbMsg.newMsg("DEVICE", device.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); entityAttributeFetched(getEntityId()); } void deviceCustomerTelemetryFetched(Device device) throws TbNodeException { ObjectMapper mapper = JacksonUtil.OBJECT_MAPPER; - TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(getTbNodeConfigFotTelemetry())); + TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(getTbNodeConfigForTelemetry())); TbEntityGetAttrNode node = getEmptyNode(); node.init(null, nodeConfiguration); - msg = TbMsg.newMsg("USER", device.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getDeviceService()).thenReturn(deviceService); - when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); + msg = TbMsg.newMsg("DEVICE", device.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); List timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); @@ -217,11 +196,40 @@ public abstract class AbstractAttributeNodeTest { assertEquals(msg.getMetaData().getValue("result"), "high"); } - protected abstract TbEntityGetAttrNode getEmptyNode(); + TbGetEntityAttrNodeConfiguration getTbNodeConfig() { + return getConfig(false); + } + + TbGetEntityAttrNodeConfiguration getTbNodeConfigForTelemetry() { + return getConfig(true); + } - abstract T getTbNodeConfig(); + @NotNull + private TbGetEntityAttrNodeConfiguration getConfig(boolean isTelemetry) { + TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); + Map conf = new HashMap<>(); + conf.put(keyAttrConf, valueAttrConf); + config.setAttrMapping(conf); + config.setTelemetry(isTelemetry); + return config; + } - abstract T getTbNodeConfigFotTelemetry(); + protected abstract TbEntityGetAttrNode getEmptyNode(); abstract EntityId getEntityId(); + + void mockFindDevice(Device device) { + when(ctx.getDeviceService()).thenReturn(deviceService); + when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); + } + + void mockFindAsset(Asset asset) { + when(ctx.getAssetService()).thenReturn(assetService); + when(assetService.findAssetByIdAsync(any(), eq(asset.getId()))).thenReturn(Futures.immediateFuture(asset)); + } + + void mockFindUser(User user) { + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(user)); + } } 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 3ed2e35702..707a204334 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 @@ -16,6 +16,7 @@ package org.thingsboard.rule.engine.metadata; import com.datastax.oss.driver.api.core.uuid.Uuids; +import com.google.common.util.concurrent.Futures; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -29,10 +30,12 @@ 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 java.util.HashMap; -import java.util.Map; import java.util.UUID; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.when; + @RunWith(MockitoJUnitRunner.class) public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest { User user = new User(); @@ -57,26 +60,6 @@ public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest { return new TbGetCustomerAttributeNode(); } - @Override - T getTbNodeConfig() { - TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); - Map conf = new HashMap<>(); - conf.put(keyAttrConf, valueAttrConf); - config.setAttrMapping(conf); - config.setTelemetry(false); - return (T) config; - } - - @Override - T getTbNodeConfigFotTelemetry() { - TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); - Map conf = new HashMap<>(); - conf.put(keyAttrConf, valueAttrConf); - config.setAttrMapping(conf); - config.setTelemetry(true); - return (T) config; - } - @Override EntityId getEntityId() { return customerId; @@ -84,16 +67,20 @@ public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest { @Test public void errorThrownIfCannotLoadAttributes() { + mockFindUser(user); errorThrownIfCannotLoadAttributes(user); } @Test public void errorThrownIfCannotLoadAttributesAsync() { + mockFindUser(user); errorThrownIfCannotLoadAttributesAsync(user); } @Test public void failedChainUsedIfCustomerCannotBeFound() { + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(null)); failedChainUsedIfCustomerCannotBeFound(user); } @@ -104,21 +91,25 @@ public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest { @Test public void usersCustomerAttributesFetched() { + mockFindUser(user); usersCustomerAttributesFetched(user); } @Test public void assetsCustomerAttributesFetched() { + mockFindAsset(asset); assetsCustomerAttributesFetched(asset); } @Test public void deviceCustomerAttributesFetched() { + mockFindDevice(device); deviceCustomerAttributesFetched(device); } @Test public void deviceCustomerTelemetryFetched() throws TbNodeException { + mockFindDevice(device); deviceCustomerTelemetryFetched(device); } } 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 d2067ca893..4f125a2466 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 @@ -16,6 +16,7 @@ package org.thingsboard.rule.engine.metadata; import com.google.common.util.concurrent.Futures; +import org.jetbrains.annotations.NotNull; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -40,7 +41,7 @@ import java.util.UUID; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.when; -@RunWith(MockitoJUnitRunner.Silent.class) +@RunWith(MockitoJUnitRunner.class) public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest { User user = new User(); Asset asset = new Asset(); @@ -74,25 +75,24 @@ public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest { } @Override - T getTbNodeConfig() { - TbGetRelatedAttrNodeConfiguration config = new TbGetRelatedAttrNodeConfiguration(); - config = config.defaultConfiguration(); - Map conf = new HashMap<>(); - conf.put(keyAttrConf, valueAttrConf); - config.setAttrMapping(conf); - config.setTelemetry(false); - return (T) config; + TbGetEntityAttrNodeConfiguration getTbNodeConfig() { + return getConfig(false); } @Override - T getTbNodeConfigFotTelemetry() { + TbGetEntityAttrNodeConfiguration getTbNodeConfigForTelemetry() { + return getConfig(true); + } + + @NotNull + private TbGetEntityAttrNodeConfiguration getConfig(boolean isTelemetry) { TbGetRelatedAttrNodeConfiguration config = new TbGetRelatedAttrNodeConfiguration(); config = config.defaultConfiguration(); Map conf = new HashMap<>(); conf.put(keyAttrConf, valueAttrConf); config.setAttrMapping(conf); - config.setTelemetry(true); - return (T) config; + config.setTelemetry(isTelemetry); + return config; } @Override 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 e297722928..2a5798d1e2 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 @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.metadata; +import com.google.common.util.concurrent.Futures; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -28,10 +29,12 @@ 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 java.util.HashMap; -import java.util.Map; import java.util.UUID; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.when; + @RunWith(MockitoJUnitRunner.class) public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest { @@ -57,26 +60,6 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest { return new TbGetTenantAttributeNode(); } - @Override - T getTbNodeConfig() { - TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); - Map conf = new HashMap<>(); - conf.put(keyAttrConf, valueAttrConf); - config.setAttrMapping(conf); - config.setTelemetry(false); - return (T) config; - } - - @Override - T getTbNodeConfigFotTelemetry() { - TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); - Map conf = new HashMap<>(); - conf.put(keyAttrConf, valueAttrConf); - config.setAttrMapping(conf); - config.setTelemetry(true); - return (T) config; - } - @Override EntityId getEntityId() { return tenantId; @@ -84,16 +67,20 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest { @Test public void errorThrownIfCannotLoadAttributes() { + mockFindUser(user); errorThrownIfCannotLoadAttributes(user); } @Test public void errorThrownIfCannotLoadAttributesAsync() { + mockFindUser(user); errorThrownIfCannotLoadAttributesAsync(user); } @Test public void failedChainUsedIfCustomerCannotBeFound() { + when(ctx.getUserService()).thenReturn(userService); + when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(null)); failedChainUsedIfCustomerCannotBeFound(user); } @@ -104,21 +91,25 @@ public class TbGetTenantAttributeNodeTest extends AbstractAttributeNodeTest { @Test public void usersCustomerAttributesFetched() { + mockFindUser(user); usersCustomerAttributesFetched(user); } @Test public void assetsCustomerAttributesFetched() { + mockFindAsset(asset); assetsCustomerAttributesFetched(asset); } @Test public void deviceCustomerAttributesFetched() { + mockFindDevice(device); deviceCustomerAttributesFetched(device); } @Test public void deviceCustomerTelemetryFetched() throws TbNodeException { + mockFindDevice(device); deviceCustomerTelemetryFetched(device); } } From 551f67ec02de132df028a282c89e989b5baf1a70 Mon Sep 17 00:00:00 2001 From: van-vanich Date: Wed, 22 Dec 2021 11:14:04 +0200 Subject: [PATCH 7/8] add this feature for target attribute: use ${metadataKey} for value from metadata, $[messageKey] for value from the message body --- .../rule/engine/metadata/TbEntityGetAttrNode.java | 8 ++++++-- .../rule/engine/metadata/AbstractAttributeNodeTest.java | 7 ++++--- 2 files changed, 10 insertions(+), 5 deletions(-) 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 index 747549cd77..13451a3b34 100644 --- 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 @@ -68,8 +68,11 @@ public abstract class TbEntityGetAttrNode implements TbNode Map mappingsMap = new HashMap<>(); config.getAttrMapping().forEach((key, value) -> { - String processPattern = TbNodeUtils.processPattern(key, msg); - mappingsMap.put(processPattern, value); + log.info("key = {}, value = {}", key, value); + log.info(msg.toString()); + String processPatternKey = TbNodeUtils.processPattern(key, msg); + String processPatternValue = TbNodeUtils.processPattern(value, msg); + mappingsMap.put(processPatternKey, processPatternValue); }); List keys = List.copyOf(mappingsMap.keySet()); @@ -94,6 +97,7 @@ public abstract class TbEntityGetAttrNode implements TbNode private void putAttributesAndTell(TbContext ctx, TbMsg msg, List attributes, Map map) { attributes.forEach(r -> { String attrName = map.get(r.getKey()); + log.info("attrName = {}, value = {}", attrName, r.getValueAsString()); msg.getMetaData().putValue(attrName, r.getValueAsString()); }); ctx.tellSuccess(msg); 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/AbstractAttributeNodeTest.java index 005b40b4e7..a0f2879851 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/AbstractAttributeNodeTest.java @@ -72,7 +72,7 @@ public abstract class AbstractAttributeNodeTest { final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); final String keyAttrConf = "${word}"; - final String valueAttrConf = "result"; + final String valueAttrConf = "${result}"; @Mock TbContext ctx; @Mock @@ -95,6 +95,7 @@ public abstract class AbstractAttributeNodeTest { metaData = new HashMap<>(); metaData.putIfAbsent("word", "temperature"); + metaData.putIfAbsent("result", "answer"); this.node = node; this.node.init(null, nodeConfiguration); @@ -181,7 +182,7 @@ public abstract class AbstractAttributeNodeTest { node.onMsg(ctx, msg); verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("result"), "highest"); + assertEquals(msg.getMetaData().getValue("answer"), "highest"); } void entityAttributeFetched(EntityId entityId) { @@ -193,7 +194,7 @@ public abstract class AbstractAttributeNodeTest { node.onMsg(ctx, msg); verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("result"), "high"); + assertEquals(msg.getMetaData().getValue("answer"), "high"); } TbGetEntityAttrNodeConfiguration getTbNodeConfig() { From 167a40d4169548cdf2da853db170c359974fb837 Mon Sep 17 00:00:00 2001 From: van-vanich Date: Wed, 22 Dec 2021 11:30:18 +0200 Subject: [PATCH 8/8] remove unnecessary code --- .../thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java | 3 --- 1 file changed, 3 deletions(-) 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 index 13451a3b34..c2e19c6826 100644 --- 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 @@ -68,8 +68,6 @@ public abstract class TbEntityGetAttrNode implements TbNode Map mappingsMap = new HashMap<>(); config.getAttrMapping().forEach((key, value) -> { - log.info("key = {}, value = {}", key, value); - log.info(msg.toString()); String processPatternKey = TbNodeUtils.processPattern(key, msg); String processPatternValue = TbNodeUtils.processPattern(value, msg); mappingsMap.put(processPatternKey, processPatternValue); @@ -97,7 +95,6 @@ public abstract class TbEntityGetAttrNode implements TbNode private void putAttributesAndTell(TbContext ctx, TbMsg msg, List attributes, Map map) { attributes.forEach(r -> { String attrName = map.get(r.getKey()); - log.info("attrName = {}, value = {}", attrName, r.getValueAsString()); msg.getMetaData().putValue(attrName, r.getValueAsString()); }); ctx.tellSuccess(msg);