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..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 @@ -30,12 +30,13 @@ import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.msg.TbMsg; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.stream.Collectors; import static org.thingsboard.common.util.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; -import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; @Slf4j @@ -65,27 +66,35 @@ public abstract class TbEntityGetAttrNode implements TbNode return; } - withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId) : getAttributesAsync(ctx, entityId), - attributes -> putAttributesAndTell(ctx, msg, attributes), + Map mappingsMap = new HashMap<>(); + config.getAttrMapping().forEach((key, value) -> { + String processPatternKey = TbNodeUtils.processPattern(key, msg); + String processPatternValue = TbNodeUtils.processPattern(value, msg); + mappingsMap.put(processPatternKey, processPatternValue); + }); + + List keys = List.copyOf(mappingsMap.keySet()); + withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId, keys) : getAttributesAsync(ctx, entityId, keys), + attributes -> putAttributesAndTell(ctx, msg, attributes, mappingsMap), t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } - private ListenableFuture> getAttributesAsync(TbContext ctx, EntityId entityId) { - ListenableFuture> latest = ctx.getAttributesService().find(ctx.getTenantId(), entityId, SERVER_SCOPE, config.getAttrMapping().keySet()); + 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) { - ListenableFuture> latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, config.getAttrMapping().keySet()); + 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) { + private void putAttributesAndTell(TbContext ctx, TbMsg msg, List attributes, Map map) { attributes.forEach(r -> { - String attrName = config.getAttrMapping().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 new file mode 100644 index 0000000000..a0f2879851 --- /dev/null +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/AbstractAttributeNodeTest.java @@ -0,0 +1,236 @@ +/** + * 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.jetbrains.annotations.NotNull; +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.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"); + metaData.putIfAbsent("result", "answer"); + + 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.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.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); + + 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); + + entityAttributeFetched(getEntityId()); + } + + void assetsCustomerAttributesFetched(Asset asset) { + msg = TbMsg.newMsg("ASSET", asset.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + entityAttributeFetched(getEntityId()); + } + + void deviceCustomerAttributesFetched(Device 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(getTbNodeConfigForTelemetry())); + + TbEntityGetAttrNode node = getEmptyNode(); + node.init(null, nodeConfiguration); + + msg = TbMsg.newMsg("DEVICE", device.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); + + 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("answer"), "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("answer"), "high"); + } + + TbGetEntityAttrNodeConfiguration getTbNodeConfig() { + return getConfig(false); + } + + TbGetEntityAttrNodeConfiguration getTbNodeConfigForTelemetry() { + return getConfig(true); + } + + @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; + } + + 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 ca88290808..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,258 +16,100 @@ 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.Collections; -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.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 TbGetCustomerAttributeNodeTest { +public class TbGetCustomerAttributeNodeTest extends AbstractAttributeNodeTest { + User user = new User(); + Asset asset = new Asset(); + Device device = new Device(); - 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; + @Before + public void initDataForTests() throws TbNodeException { + init(new TbGetCustomerAttributeNode()); + user.setCustomerId(customerId); + user.setId(new UserId(UUID.randomUUID())); - private TbMsg msg; + asset.setCustomerId(customerId); + asset.setId(new AssetId(UUID.randomUUID())); - private RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); - private RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); + device.setCustomerId(customerId); + device.setId(new DeviceId(Uuids.timeBased())); + } - @Before - public void init() throws TbNodeException { - TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); - Map attrMapping = new HashMap<>(); - attrMapping.putIfAbsent("temperature", "tempo"); - config.setAttrMapping(attrMapping); - config.setTelemetry(false); - ObjectMapper mapper = new ObjectMapper(); - TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + @Override + protected TbEntityGetAttrNode getEmptyNode() { + return new TbGetCustomerAttributeNode(); + } - 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), eq(Collections.singleton("temperature")))) - .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()); + mockFindUser(user); + 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), eq(Collections.singleton("temperature")))) - .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()); + mockFindUser(user); + 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()); + when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(null)); + failedChainUsedIfCustomerCannotBeFound(user); } @Test public void customerAttributeAddedInMetadata() { - CustomerId customerId = new CustomerId(Uuids.timeBased()); - msg = TbMsg.newMsg( "CUSTOMER", customerId, new TbMsgMetaData(), 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(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getUserService()).thenReturn(userService); - when(userService.findUserByIdAsync(any(), eq(userId))).thenReturn(Futures.immediateFuture(user)); - - entityAttributeFetched(customerId); + mockFindUser(user); + 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(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getAssetService()).thenReturn(assetService); - when(assetService.findAssetByIdAsync(any(), eq(assetId))).thenReturn(Futures.immediateFuture(asset)); - - entityAttributeFetched(customerId); + mockFindAsset(asset); + 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(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); - - when(ctx.getDeviceService()).thenReturn(deviceService); - when(deviceService.findDeviceByIdAsync(any(), eq(deviceId))).thenReturn(Futures.immediateFuture(device)); - - entityAttributeFetched(customerId); + mockFindDevice(device); + deviceCustomerAttributesFetched(device); } @Test public void deviceCustomerTelemetryFetched() throws TbNodeException { - TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); - Map attrMapping = new HashMap<>(); - attrMapping.putIfAbsent("temperature", "tempo"); - config.setAttrMapping(attrMapping); - 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(), 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("temperature")))) - .thenReturn(Futures.immediateFuture(timeseries)); - - node.onMsg(ctx, msg); - verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("tempo"), "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")))) - .thenReturn(Futures.immediateFuture(attributes)); - - node.onMsg(ctx, msg); - verify(ctx).tellSuccess(msg); - assertEquals(msg.getMetaData().getValue("tempo"), "high"); + 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 new file mode 100644 index 0000000000..4f125a2466 --- /dev/null +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java @@ -0,0 +1,166 @@ +/** + * 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.google.common.util.concurrent.Futures; +import org.jetbrains.annotations.NotNull; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.junit.MockitoJUnitRunner; +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.DeviceId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.UserId; +import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.dao.relation.RelationService; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.UUID; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.when; + +@RunWith(MockitoJUnitRunner.class) +public class TbGetRelatedAttributeNodeTest extends AbstractAttributeNodeTest { + User user = new User(); + Asset asset = new Asset(); + Device device = new Device(); + @Mock + private RelationService relationService; + private EntityRelation entityRelation; + + @Before + 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 + TbGetEntityAttrNodeConfiguration getTbNodeConfig() { + return getConfig(false); + } + + @Override + 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(isTelemetry); + return config; + } + + @Override + EntityId getEntityId() { + return customerId; + } + + @Test + public void errorThrownIfCannotLoadAttributes() { + entityRelation.setFrom(user.getId()); + entityRelation.setTo(customerId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + errorThrownIfCannotLoadAttributes(user); + } + + @Test + public void errorThrownIfCannotLoadAttributesAsync() { + entityRelation.setFrom(user.getId()); + entityRelation.setTo(customerId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + errorThrownIfCannotLoadAttributesAsync(user); + } + + @Test + public void failedChainUsedIfCustomerCannotBeFound() { + entityRelation.setFrom(customerId); + entityRelation.setTo(null); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + failedChainUsedIfCustomerCannotBeFound(user); + } + + @Test + public void customerAttributeAddedInMetadata() { + entityRelation.setFrom(customerId); + entityRelation.setTo(customerId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + entityAttributeAddedInMetadata(customerId, "CUSTOMER"); + } + + @Test + public void usersCustomerAttributesFetched() { + entityRelation.setFrom(user.getId()); + entityRelation.setTo(customerId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + usersCustomerAttributesFetched(user); + } + + @Test + public void assetsCustomerAttributesFetched() { + entityRelation.setFrom(asset.getId()); + entityRelation.setTo(customerId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + assetsCustomerAttributesFetched(asset); + } + + @Test + public void deviceCustomerAttributesFetched() { + entityRelation.setFrom(device.getId()); + entityRelation.setTo(customerId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + deviceCustomerAttributesFetched(device); + } + + @Test + public void deviceCustomerTelemetryFetched() throws TbNodeException { + entityRelation.setFrom(device.getId()); + entityRelation.setTo(customerId); + when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); + deviceCustomerTelemetryFetched(device); + } +} 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..2a5798d1e2 --- /dev/null +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java @@ -0,0 +1,115 @@ +/** + * 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.google.common.util.concurrent.Futures; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.junit.MockitoJUnitRunner; +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.DeviceId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.UserId; + +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 { + + User user = new User(); + Asset asset = new Asset(); + Device device = new Device(); + + @Before + 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 + EntityId getEntityId() { + return tenantId; + } + + @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); + } + + @Test + public void customerAttributeAddedInMetadata() { + entityAttributeAddedInMetadata(tenantId, "TENANT"); + } + + @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); + } +}