From 436e19d6524532765edb063f83115c9861fcdba8 Mon Sep 17 00:00:00 2001 From: van-vanich Date: Tue, 23 Nov 2021 16:14:22 +0200 Subject: [PATCH] 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); } }