Browse Source

Merge pull request #5550 from van-vanich/upgrade_enrichment_attributes_rule_node

[3.3.3] Upgrade enrichment attributes rule node
pull/5823/head
Andrew Shvayka 5 years ago
committed by GitHub
parent
commit
530765487c
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 27
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java
  2. 236
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/AbstractAttributeNodeTest.java
  3. 234
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java
  4. 166
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java
  5. 115
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java

27
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<T extends EntityId> implements TbNode
return;
}
withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId) : getAttributesAsync(ctx, entityId),
attributes -> putAttributesAndTell(ctx, msg, attributes),
Map<String, String> 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<String> 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<List<KvEntry>> getAttributesAsync(TbContext ctx, EntityId entityId) {
ListenableFuture<List<AttributeKvEntry>> latest = ctx.getAttributesService().find(ctx.getTenantId(), entityId, SERVER_SCOPE, config.getAttrMapping().keySet());
private ListenableFuture<List<KvEntry>> getAttributesAsync(TbContext ctx, EntityId entityId, List<String> attrKeys) {
ListenableFuture<List<AttributeKvEntry>> 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<List<KvEntry>> getLatestTelemetry(TbContext ctx, EntityId entityId) {
ListenableFuture<List<TsKvEntry>> latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, config.getAttrMapping().keySet());
private ListenableFuture<List<KvEntry>> getLatestTelemetry(TbContext ctx, EntityId entityId, List<String> timeseriesKeys) {
ListenableFuture<List<TsKvEntry>> 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<? extends KvEntry> attributes) {
private void putAttributesAndTell(TbContext ctx, TbMsg msg, List<? extends KvEntry> attributes, Map<String, String> 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);

236
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<String, String> 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<Throwable> 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<Throwable> 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<TsKvEntry> 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<AttributeKvEntry> 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<String, String> 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));
}
}

234
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<String, String> 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<Throwable> 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<Throwable> 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<String, String> 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<TsKvEntry> 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<AttributeKvEntry> 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);
}
}

166
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<String, String> 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);
}
}

115
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);
}
}
Loading…
Cancel
Save