Browse Source

refactoring and improve code

pull/5550/head
van-vanich 5 years ago
parent
commit
436e19d652
  1. 31
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java
  2. 64
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/AbstractAttributeNodeTest.java
  3. 35
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java
  4. 24
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNodeTest.java
  5. 35
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetTenantAttributeNodeTest.java

31
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<T extends EntityId> implements TbNode
return;
}
List<String> 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<String, String> mappingsMap = new HashMap<>();
config.getAttrMapping().forEach((key, value) -> {
String processPattern = TbNodeUtils.processPattern(key, msg);
mappingsMap.put(processPattern, value);
});
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, List<String> attrProcess) {
ListenableFuture<List<AttributeKvEntry>> latest = ctx.getAttributesService().find(ctx.getTenantId(), entityId, SERVER_SCOPE, attrProcess);
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, List<String> latestProcess) {
ListenableFuture<List<TsKvEntry>> latest = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), entityId, latestProcess);
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) {
Map<String, String> 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<? extends KvEntry> attributes, Map<String, String> 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);

64
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<TsKvEntry> 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> T getTbNodeConfig();
@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;
}
abstract <T> 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));
}
}

35
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> T getTbNodeConfig() {
TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration();
Map<String, String> conf = new HashMap<>();
conf.put(keyAttrConf, valueAttrConf);
config.setAttrMapping(conf);
config.setTelemetry(false);
return (T) config;
}
@Override
<T> T getTbNodeConfigFotTelemetry() {
TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration();
Map<String, String> 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);
}
}

24
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> T getTbNodeConfig() {
TbGetRelatedAttrNodeConfiguration config = new TbGetRelatedAttrNodeConfiguration();
config = config.defaultConfiguration();
Map<String, String> conf = new HashMap<>();
conf.put(keyAttrConf, valueAttrConf);
config.setAttrMapping(conf);
config.setTelemetry(false);
return (T) config;
TbGetEntityAttrNodeConfiguration getTbNodeConfig() {
return getConfig(false);
}
@Override
<T> T getTbNodeConfigFotTelemetry() {
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(true);
return (T) config;
config.setTelemetry(isTelemetry);
return config;
}
@Override

35
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> T getTbNodeConfig() {
TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration();
Map<String, String> conf = new HashMap<>();
conf.put(keyAttrConf, valueAttrConf);
config.setAttrMapping(conf);
config.setTelemetry(false);
return (T) config;
}
@Override
<T> T getTbNodeConfigFotTelemetry() {
TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration();
Map<String, String> 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);
}
}

Loading…
Cancel
Save