From 52c8e1b48f258f0f3060c5b5e244fd9fd609845a Mon Sep 17 00:00:00 2001 From: Bohdan Smetaniuk Date: Mon, 28 Sep 2020 17:02:14 +0300 Subject: [PATCH 1/4] fixed attributes saving --- .../thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java index f0f5d2418f..5da3df68ae 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java @@ -66,7 +66,7 @@ public class TbMsgAttributesNode implements TbNode { if (StringUtils.isEmpty(msg.getMetaData().getValue(SCOPE))) { msg.getMetaData().putValue(SCOPE, config.getScope()); } - ctx.getTelemetryService().saveAndNotify(ctx.getTenantId(), msg.getOriginator(), config.getScope(), + ctx.getTelemetryService().saveAndNotify(ctx.getTenantId(), msg.getOriginator(), msg.getMetaData().getValue(SCOPE), new ArrayList<>(attributes), new TelemetryNodeCallback(ctx, msg)); } From 9e215251584092941187f587a68a064b48e95cce Mon Sep 17 00:00:00 2001 From: Bohdan Smetaniuk Date: Tue, 29 Sep 2020 13:43:35 +0300 Subject: [PATCH 2/4] code coverage for processors and telemetry constructor --- .../thingsboard/server/edge/BaseEdgeTest.java | 284 +++++++++++++++++- .../server/edge/imitator/EdgeImitator.java | 25 ++ .../server/edge/imitator/EdgeStorage.java | 29 +- .../engine/telemetry/TbMsgAttributesNode.java | 2 +- 4 files changed, 326 insertions(+), 14 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java index be26e64f1a..084def021c 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java @@ -15,7 +15,11 @@ */ package org.thingsboard.server.edge; +import com.datastax.driver.core.utils.UUIDs; import com.fasterxml.jackson.core.type.TypeReference; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.gson.JsonObject; import lombok.extern.slf4j.Slf4j; import org.junit.After; import org.junit.Assert; @@ -24,6 +28,7 @@ import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.DashboardInfo; +import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.Tenant; @@ -33,7 +38,11 @@ import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmStatus; import org.thingsboard.server.common.data.asset.Asset; +import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.edge.EdgeEventType; +import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.page.TimePageData; @@ -42,17 +51,24 @@ import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.common.data.security.Authority; +import org.thingsboard.server.common.transport.adaptor.JsonConverter; import org.thingsboard.server.controller.AbstractControllerTest; -import org.thingsboard.server.dao.rule.RuleChainService; +import org.thingsboard.server.dao.edge.EdgeEventService;; import org.thingsboard.server.edge.imitator.EdgeImitator; +import org.thingsboard.server.gen.edge.AlarmUpdateMsg; +import org.thingsboard.server.gen.edge.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.EdgeConfiguration; +import org.thingsboard.server.gen.edge.EntityDataProto; +import org.thingsboard.server.gen.edge.RelationUpdateMsg; +import org.thingsboard.server.gen.edge.UpdateMsgType; +import org.thingsboard.server.gen.edge.UplinkMsg; +import org.thingsboard.server.gen.transport.TransportProtos; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.Set; import java.util.UUID; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.TimeUnit; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @@ -67,6 +83,9 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { private EdgeImitator edgeImitator; private Edge edge; + @Autowired + private EdgeEventService edgeEventService; + @Before public void beforeTest() throws Exception { loginSysAdmin(); @@ -114,6 +133,9 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { testDashboards(); testRelations(); testAlarms(); + testTimeseries(); + testAttributes(); + testSendMessagesToCloud(); } private void testReceivedInitialData() throws Exception { @@ -376,6 +398,251 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { log.info("Alarms tested successfully"); } + private void testTimeseries() throws Exception { + log.info("Testing timeseries"); + ObjectMapper mapper = new ObjectMapper(); + List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() {}, new TextPageLink(100)).getData(); + Assert.assertEquals(1, edgeDevices.size()); + Device device = edgeDevices.get(0); + Assert.assertEquals("Edge Device 1", device.getName()); + + String timeseriesData = "{\"data\":{\"temperature\":25},\"ts\":" + System.currentTimeMillis() + "}"; + JsonNode timeseriesEntityData = mapper.readTree(timeseriesData); + EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), ActionType.TIMESERIES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData); + edgeImitator.getStorage().expectMessageAmount(1); + edgeEventService.saveAsync(edgeEvent1); + edgeImitator.getStorage().waitForMessages(); + + EntityDataProto latestEntityDataMsg = edgeImitator.getStorage().getLatestEntityDataMsg(); + Assert.assertNotNull(latestEntityDataMsg); + UUID uuid = new UUID(latestEntityDataMsg.getEntityIdMSB(), latestEntityDataMsg.getEntityIdLSB()); + Assert.assertEquals(device.getId().getId(), uuid); + Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); + Assert.assertTrue(latestEntityDataMsg.hasPostTelemetryMsg()); + + TransportProtos.PostTelemetryMsg postTelemetryMsg = latestEntityDataMsg.getPostTelemetryMsg(); + Assert.assertEquals(1, postTelemetryMsg.getTsKvListCount()); + TransportProtos.TsKvListProto tsKvListProto = postTelemetryMsg.getTsKvList(0); + Assert.assertEquals(timeseriesEntityData.get("ts").asLong(), tsKvListProto.getTs()); + Assert.assertEquals(1, tsKvListProto.getKvCount()); + TransportProtos.KeyValueProto keyValueProto = tsKvListProto.getKv(0); + Assert.assertEquals("temperature", keyValueProto.getKey()); + Assert.assertEquals(25, keyValueProto.getLongV()); + edgeImitator.getStorage().setLatestEntityDataMsg(null); + log.info("Timeseries tested successfully"); + } + + private void testAttributes() throws Exception { + log.info("Testing attributes"); + List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() {}, new TextPageLink(100)).getData(); + Assert.assertEquals(1, edgeDevices.size()); + Device device = edgeDevices.get(0); + Assert.assertEquals("Edge Device 1", device.getName()); + + String attributesData = "{\"scope\":\"SERVER_SCOPE\",\"kv\":{\"test\":\"test\"}}"; + JsonNode attributesEntityData = mapper.readTree(attributesData); + EdgeEvent edgeEvent2 = constructEdgeEvent(tenantId, edge.getId(), ActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); + edgeImitator.getStorage().expectMessageAmount(1); + edgeEventService.saveAsync(edgeEvent2); + edgeImitator.getStorage().waitForMessages(); + + EntityDataProto latestEntityDataMsg = edgeImitator.getStorage().getLatestEntityDataMsg(); + Assert.assertNotNull(latestEntityDataMsg); + UUID uuid = new UUID(latestEntityDataMsg.getEntityIdMSB(), latestEntityDataMsg.getEntityIdLSB()); + Assert.assertEquals(device.getId().getId(), uuid); + Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); + Assert.assertEquals(attributesEntityData.get("scope").asText(), latestEntityDataMsg.getPostAttributeScope()); + Assert.assertTrue(latestEntityDataMsg.hasPostAttributesMsg()); + + TransportProtos.PostAttributeMsg postAttributeMsg = latestEntityDataMsg.getPostAttributesMsg(); + Assert.assertEquals(1, postAttributeMsg.getKvCount()); + TransportProtos.KeyValueProto keyValueProto = postAttributeMsg.getKv(0); + Assert.assertEquals("test", keyValueProto.getKey()); + Assert.assertEquals("test", keyValueProto.getStringV()); + edgeImitator.getStorage().setLatestEntityDataMsg(null); + log.info("Attributes tested successfully"); + } + + private void testSendMessagesToCloud() throws Exception { + log.info("Sending messages to cloud"); + sendDevice(); + sendAlarm(); + sendTelemetry(); + sendRelation(); + sendDeleteDeviceOnEdge(); + log.info("Messages were sent successfully"); + } + + private void sendDevice() throws Exception { + UUID uuid = UUIDs.timeBased(); + + UplinkMsg.Builder builder = UplinkMsg.newBuilder(); + DeviceUpdateMsg.Builder deviceUpdateMsgBuilder = DeviceUpdateMsg.newBuilder(); + deviceUpdateMsgBuilder.setIdMSB(uuid.getMostSignificantBits()); + deviceUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits()); + deviceUpdateMsgBuilder.setName("Edge Device 2"); + deviceUpdateMsgBuilder.setType("test"); + deviceUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE); + builder.addDeviceUpdateMsg(deviceUpdateMsgBuilder.build()); + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(builder.build()); + edgeImitator.waitForResponses(); + + Device device = doGet("/api/device/" + uuid.toString(), Device.class); + Assert.assertNotNull(device); + Assert.assertEquals("Edge Device 2", device.getName()); + } + + private void sendAlarm() throws Exception { + List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() {}, new TextPageLink(100)).getData(); + Optional foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 2")).findAny(); + Assert.assertTrue(foundDevice.isPresent()); + Device device = foundDevice.get(); + + UplinkMsg.Builder builder = UplinkMsg.newBuilder(); + AlarmUpdateMsg.Builder alarmUpdateMgBuilder = AlarmUpdateMsg.newBuilder(); + alarmUpdateMgBuilder.setName("alarm from edge"); + alarmUpdateMgBuilder.setStatus(AlarmStatus.ACTIVE_UNACK.name()); + alarmUpdateMgBuilder.setSeverity(AlarmSeverity.CRITICAL.name()); + alarmUpdateMgBuilder.setOriginatorName(device.getName()); + alarmUpdateMgBuilder.setOriginatorType(EntityType.DEVICE.name()); + builder.addAlarmUpdateMsg(alarmUpdateMgBuilder.build()); + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(builder.build()); + edgeImitator.waitForResponses(); + + + List alarms = doGetTypedWithPageLink("/api/alarm/{entityType}/{entityId}?", + new TypeReference>() {}, + new TextPageLink(100), device.getId().getEntityType().name(), device.getId().getId().toString()) + .getData(); + + for (AlarmInfo alarmInfo: alarms) { + log.info(String.valueOf(alarmInfo)); + } + + Optional foundAlarm = alarms.stream().filter(alarm -> alarm.getType().equals("alarm from edge")).findAny(); + Assert.assertTrue(foundAlarm.isPresent()); + AlarmInfo alarmInfo = foundAlarm.get(); + Assert.assertEquals(device.getId(), alarmInfo.getOriginator()); + Assert.assertEquals(AlarmStatus.ACTIVE_UNACK, alarmInfo.getStatus()); + Assert.assertEquals(AlarmSeverity.CRITICAL, alarmInfo.getSeverity()); + } + + private void sendRelation() throws Exception { + List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() {}, new TextPageLink(100)).getData(); + Optional foundDevice1 = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 1")).findAny(); + Assert.assertTrue(foundDevice1.isPresent()); + Device device1 = foundDevice1.get(); + Optional foundDevice2 = edgeDevices.stream().filter(device2 -> device2.getName().equals("Edge Device 2")).findAny(); + Assert.assertTrue(foundDevice2.isPresent()); + Device device2 = foundDevice2.get(); + + UplinkMsg.Builder builder = UplinkMsg.newBuilder(); + RelationUpdateMsg.Builder relationUpdateMsgBuilder = RelationUpdateMsg.newBuilder(); + relationUpdateMsgBuilder.setType("test"); + relationUpdateMsgBuilder.setTypeGroup(RelationTypeGroup.COMMON.name()); + relationUpdateMsgBuilder.setToIdMSB(device1.getId().getId().getMostSignificantBits()); + relationUpdateMsgBuilder.setToIdLSB(device1.getId().getId().getLeastSignificantBits()); + relationUpdateMsgBuilder.setToEntityType(device1.getId().getEntityType().name()); + relationUpdateMsgBuilder.setFromIdMSB(device2.getId().getId().getMostSignificantBits()); + relationUpdateMsgBuilder.setFromIdLSB(device2.getId().getId().getLeastSignificantBits()); + relationUpdateMsgBuilder.setFromEntityType(device2.getId().getEntityType().name()); + relationUpdateMsgBuilder.setAdditionalInfo("{}"); + builder.addRelationUpdateMsg(relationUpdateMsgBuilder.build()); + UplinkMsg msg = builder.build(); + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(msg); + edgeImitator.waitForResponses(); + + EntityRelation relation = doGet("/api/relation?" + + "&fromId=" + device2.getId().getId().toString() + + "&fromType=" + device2.getId().getEntityType().name() + + "&relationType=" + "test" + + "&relationTypeGroup=" + RelationTypeGroup.COMMON.name() + + "&toId=" + device1.getId().getId().toString() + + "&toType=" + device1.getId().getEntityType().name(), EntityRelation.class); + Assert.assertNotNull(relation); + } + + private void sendTelemetry() throws Exception { + List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() {}, new TextPageLink(100)).getData(); + Optional foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 2")).findAny(); + Assert.assertTrue(foundDevice.isPresent()); + Device device = foundDevice.get(); + + edgeImitator.expectResponsesAmount(2); + + JsonObject data = new JsonObject(); + String timeseriesKey = "key"; + String timeseriesValue = "25"; + data.addProperty(timeseriesKey, timeseriesValue); + UplinkMsg.Builder builder1 = UplinkMsg.newBuilder(); + EntityDataProto.Builder entityDataBuilder = EntityDataProto.newBuilder(); + entityDataBuilder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(data, System.currentTimeMillis())); + entityDataBuilder.setEntityType(device.getId().getEntityType().name()); + entityDataBuilder.setEntityIdMSB(device.getUuidId().getMostSignificantBits()); + entityDataBuilder.setEntityIdLSB(device.getUuidId().getLeastSignificantBits()); + builder1.addEntityData(entityDataBuilder.build()); + edgeImitator.sendUplinkMsg(builder1.build()); + + JsonObject attributesData = new JsonObject(); + String attributesKey = "test_attr"; + String attributesValue = "test_value"; + attributesData.addProperty(attributesKey, attributesValue); + UplinkMsg.Builder builder2 = UplinkMsg.newBuilder(); + EntityDataProto.Builder entityDataBuilder2 = EntityDataProto.newBuilder(); + entityDataBuilder2.setEntityType(device.getId().getEntityType().name()); + entityDataBuilder2.setEntityIdMSB(device.getId().getId().getMostSignificantBits()); + entityDataBuilder2.setEntityIdLSB(device.getId().getId().getLeastSignificantBits()); + entityDataBuilder2.setPostAttributesMsg(JsonConverter.convertToAttributesProto(attributesData)); + entityDataBuilder2.setPostAttributeScope(DataConstants.SERVER_SCOPE); + builder2.addEntityData(entityDataBuilder2.build()); + edgeImitator.sendUplinkMsg(builder2.build()); + + edgeImitator.waitForResponses(); + + Thread.sleep(1000); + Map>> timeseries = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/values/timeseries?keys=" + timeseriesKey, Map.class); + Assert.assertTrue(timeseries.containsKey(timeseriesKey)); + Assert.assertEquals(1, timeseries.get(timeseriesKey).size()); + Assert.assertEquals(timeseriesValue, timeseries.get(timeseriesKey).get(0).get("value")); + + List> attributes = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getId() + "/values/attributes/" + DataConstants.SERVER_SCOPE, List.class); + Assert.assertEquals(1, attributes.size()); + Assert.assertEquals(attributes.get(0).get("key"), attributesKey); + Assert.assertEquals(attributes.get(0).get("value"), attributesValue); + + } + + private void sendDeleteDeviceOnEdge() throws Exception { + List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() {}, new TextPageLink(100)).getData(); + Optional foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 2")).findAny(); + Assert.assertTrue(foundDevice.isPresent()); + Device device = foundDevice.get(); + UplinkMsg.Builder builder = UplinkMsg.newBuilder(); + DeviceUpdateMsg.Builder deviceDeleteMsgBuilder = DeviceUpdateMsg.newBuilder(); + deviceDeleteMsgBuilder.setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE); + deviceDeleteMsgBuilder.setIdMSB(device.getId().getId().getMostSignificantBits()); + deviceDeleteMsgBuilder.setIdLSB(device.getId().getId().getLeastSignificantBits()); + builder.addDeviceUpdateMsg(deviceDeleteMsgBuilder.build()); + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(builder.build()); + edgeImitator.waitForResponses(); + device = doGet("/api/device/" + device.getId().getId().toString(), Device.class); + Assert.assertNotNull(device); + edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() { + }, new TextPageLink(100)).getData(); + Assert.assertFalse(edgeDevices.contains(device)); + } + private void installation() throws Exception { edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class); @@ -413,4 +680,15 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { doDelete("/api/edge/" + edge.getId().getId().toString()) .andExpect(status().isOk()); } + + private EdgeEvent constructEdgeEvent(TenantId tenantId, EdgeId edgeId, ActionType edgeEventAction, UUID entityId, EdgeEventType edgeEventType, JsonNode entityBody) { + EdgeEvent edgeEvent = new EdgeEvent(); + edgeEvent.setEdgeId(edgeId); + edgeEvent.setTenantId(tenantId); + edgeEvent.setAction(edgeEventAction.name()); + edgeEvent.setEntityId(entityId); + edgeEvent.setType(edgeEventType); + edgeEvent.setBody(entityBody); + return edgeEvent; + } } diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java index e1b54b6f84..90ec80838f 100644 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java +++ b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java @@ -32,14 +32,18 @@ import org.thingsboard.server.gen.edge.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.DownlinkMsg; import org.thingsboard.server.gen.edge.DownlinkResponseMsg; import org.thingsboard.server.gen.edge.EdgeConfiguration; +import org.thingsboard.server.gen.edge.EntityDataProto; import org.thingsboard.server.gen.edge.RelationUpdateMsg; import org.thingsboard.server.gen.edge.RuleChainUpdateMsg; +import org.thingsboard.server.gen.edge.UplinkMsg; import org.thingsboard.server.gen.edge.UplinkResponseMsg; import java.lang.reflect.Field; import java.util.ArrayList; import java.util.List; import java.util.UUID; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; @Slf4j public class EdgeImitator { @@ -49,6 +53,8 @@ public class EdgeImitator { private EdgeRpcClient edgeRpcClient; + private CountDownLatch responsesLatch; + @Getter private EdgeStorage storage; @@ -56,10 +62,12 @@ public class EdgeImitator { public EdgeImitator(String host, int port, String routingKey, String routingSecret) throws NoSuchFieldException, IllegalAccessException { edgeRpcClient = new EdgeGrpcClient(); storage = new EdgeStorage(); + responsesLatch = new CountDownLatch(0); this.routingKey = routingKey; this.routingSecret = routingSecret; setEdgeCredentials("rpcHost", host); setEdgeCredentials("rpcPort", port); + setEdgeCredentials("keepAliveTimeSec", 300); } private void setEdgeCredentials(String fieldName, Object value) throws NoSuchFieldException, IllegalAccessException { @@ -83,8 +91,13 @@ public class EdgeImitator { edgeRpcClient.disconnect(false); } + public void sendUplinkMsg(UplinkMsg uplinkMsg) { + edgeRpcClient.sendUplinkMsg(uplinkMsg); + } + private void onUplinkResponse(UplinkResponseMsg msg) { log.info("onUplinkResponse: {}", msg); + responsesLatch.countDown(); } private void onEdgeUpdate(EdgeConfiguration edgeConfiguration) { @@ -144,7 +157,19 @@ public class EdgeImitator { result.add(storage.processAlarm(alarmUpdateMsg)); } } + if (downlinkMsg.getEntityDataList() != null && !downlinkMsg.getEntityDataList().isEmpty()) { + for (EntityDataProto entityData: downlinkMsg.getEntityDataList()) { + result.add(storage.processEntityData(entityData)); + } + } return Futures.allAsList(result); } + public void waitForResponses() throws InterruptedException { responsesLatch.await(5, TimeUnit.SECONDS); + } + + public void expectResponsesAmount(int messageAmount) { + responsesLatch = new CountDownLatch(messageAmount); + } + } diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java index cadd580f58..a8301fda33 100644 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java +++ b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java @@ -27,6 +27,7 @@ import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.gen.edge.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.EdgeConfiguration; +import org.thingsboard.server.gen.edge.EntityDataProto; import org.thingsboard.server.gen.edge.RelationUpdateMsg; import org.thingsboard.server.gen.edge.UpdateMsgType; @@ -47,17 +48,20 @@ public class EdgeStorage { private EdgeConfiguration configuration; - private CountDownLatch latch; + private CountDownLatch messagesLatch; private Map entities; private Map alarms; private List relations; + private EntityDataProto latestEntityDataMsg; + public EdgeStorage() { - latch = new CountDownLatch(0); + messagesLatch = new CountDownLatch(0); entities = new HashMap<>(); alarms = new HashMap<>(); relations = new ArrayList<>(); + latestEntityDataMsg = null; } public ListenableFuture processEntity(UpdateMsgType msgType, EntityType type, UUID uuid) { @@ -65,11 +69,11 @@ public class EdgeStorage { case ENTITY_CREATED_RPC_MESSAGE: case ENTITY_UPDATED_RPC_MESSAGE: entities.put(uuid, type); - latch.countDown(); + messagesLatch.countDown(); break; case ENTITY_DELETED_RPC_MESSAGE: if (entities.remove(uuid) != null) { - latch.countDown(); + messagesLatch.countDown(); } break; } @@ -93,7 +97,7 @@ public class EdgeStorage { break; } if (result) { - latch.countDown(); + messagesLatch.countDown(); } return Futures.immediateFuture(null); } @@ -105,17 +109,23 @@ public class EdgeStorage { case ALARM_ACK_RPC_MESSAGE: case ALARM_CLEAR_RPC_MESSAGE: alarms.put(alarmMsg.getType(), AlarmStatus.valueOf(alarmMsg.getStatus())); - latch.countDown(); + messagesLatch.countDown(); break; case ENTITY_DELETED_RPC_MESSAGE: if (alarms.remove(alarmMsg.getName()) != null) { - latch.countDown(); + messagesLatch.countDown(); } break; } return Futures.immediateFuture(null); } + public ListenableFuture processEntityData(EntityDataProto entityData) { + latestEntityDataMsg = entityData; + messagesLatch.countDown(); + return Futures.immediateFuture(null); + } + public Set getEntitiesByType(EntityType type) { return entities.entrySet().stream() .filter(entry -> entry.getValue().equals(type)) @@ -123,11 +133,10 @@ public class EdgeStorage { } public void waitForMessages() throws InterruptedException { - latch.await(5, TimeUnit.SECONDS); + messagesLatch.await(5, TimeUnit.SECONDS); } public void expectMessageAmount(int messageAmount) { - latch = new CountDownLatch(messageAmount); + messagesLatch = new CountDownLatch(messageAmount); } - } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java index 5da3df68ae..f0f5d2418f 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java @@ -66,7 +66,7 @@ public class TbMsgAttributesNode implements TbNode { if (StringUtils.isEmpty(msg.getMetaData().getValue(SCOPE))) { msg.getMetaData().putValue(SCOPE, config.getScope()); } - ctx.getTelemetryService().saveAndNotify(ctx.getTenantId(), msg.getOriginator(), msg.getMetaData().getValue(SCOPE), + ctx.getTelemetryService().saveAndNotify(ctx.getTenantId(), msg.getOriginator(), config.getScope(), new ArrayList<>(attributes), new TelemetryNodeCallback(ctx, msg)); } From 773b856c2547924cc76fb61fec38ecc43056fd21 Mon Sep 17 00:00:00 2001 From: Bohdan Smetaniuk Date: Fri, 2 Oct 2020 15:34:54 +0300 Subject: [PATCH 3/4] tested attributes updated/post attributes separation --- .../thingsboard/server/edge/BaseEdgeTest.java | 37 +++++++++++++++---- 1 file changed, 30 insertions(+), 7 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java index 084def021c..db7bb48f75 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java @@ -19,6 +19,7 @@ import com.datastax.driver.core.utils.UUIDs; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.gson.JsonObject; import lombok.extern.slf4j.Slf4j; import org.junit.After; @@ -441,11 +442,11 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Device device = edgeDevices.get(0); Assert.assertEquals("Edge Device 1", device.getName()); - String attributesData = "{\"scope\":\"SERVER_SCOPE\",\"kv\":{\"test\":\"test\"}}"; + String attributesData = "{\"scope\":\"SERVER_SCOPE\",\"kv\":{\"key\":\"value\"}}"; JsonNode attributesEntityData = mapper.readTree(attributesData); - EdgeEvent edgeEvent2 = constructEdgeEvent(tenantId, edge.getId(), ActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); + EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), ActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); edgeImitator.getStorage().expectMessageAmount(1); - edgeEventService.saveAsync(edgeEvent2); + edgeEventService.saveAsync(edgeEvent1); edgeImitator.getStorage().waitForMessages(); EntityDataProto latestEntityDataMsg = edgeImitator.getStorage().getLatestEntityDataMsg(); @@ -454,14 +455,36 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertEquals(device.getId().getId(), uuid); Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); Assert.assertEquals(attributesEntityData.get("scope").asText(), latestEntityDataMsg.getPostAttributeScope()); + Assert.assertTrue(latestEntityDataMsg.hasAttributesUpdatedMsg()); + + TransportProtos.PostAttributeMsg attributesUpdatedMsg = latestEntityDataMsg.getAttributesUpdatedMsg(); + Assert.assertEquals(1, attributesUpdatedMsg.getKvCount()); + TransportProtos.KeyValueProto keyValueProto = attributesUpdatedMsg.getKv(0); + Assert.assertEquals("key", keyValueProto.getKey()); + Assert.assertEquals("value", keyValueProto.getStringV()); + edgeImitator.getStorage().setLatestEntityDataMsg(null); + + ((ObjectNode) attributesEntityData).put("isPostAttributes", true); + EdgeEvent edgeEvent2 = constructEdgeEvent(tenantId, edge.getId(), ActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); + edgeImitator.getStorage().expectMessageAmount(1); + edgeEventService.saveAsync(edgeEvent2); + edgeImitator.getStorage().waitForMessages(); + + latestEntityDataMsg = edgeImitator.getStorage().getLatestEntityDataMsg(); + Assert.assertNotNull(latestEntityDataMsg); + uuid = new UUID(latestEntityDataMsg.getEntityIdMSB(), latestEntityDataMsg.getEntityIdLSB()); + Assert.assertEquals(device.getId().getId(), uuid); + Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); + Assert.assertEquals(attributesEntityData.get("scope").asText(), latestEntityDataMsg.getPostAttributeScope()); Assert.assertTrue(latestEntityDataMsg.hasPostAttributesMsg()); TransportProtos.PostAttributeMsg postAttributeMsg = latestEntityDataMsg.getPostAttributesMsg(); Assert.assertEquals(1, postAttributeMsg.getKvCount()); - TransportProtos.KeyValueProto keyValueProto = postAttributeMsg.getKv(0); - Assert.assertEquals("test", keyValueProto.getKey()); - Assert.assertEquals("test", keyValueProto.getStringV()); + keyValueProto = postAttributeMsg.getKv(0); + Assert.assertEquals("key", keyValueProto.getKey()); + Assert.assertEquals("value", keyValueProto.getStringV()); edgeImitator.getStorage().setLatestEntityDataMsg(null); + log.info("Attributes tested successfully"); } @@ -600,7 +623,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { entityDataBuilder2.setEntityType(device.getId().getEntityType().name()); entityDataBuilder2.setEntityIdMSB(device.getId().getId().getMostSignificantBits()); entityDataBuilder2.setEntityIdLSB(device.getId().getId().getLeastSignificantBits()); - entityDataBuilder2.setPostAttributesMsg(JsonConverter.convertToAttributesProto(attributesData)); + entityDataBuilder2.setAttributesUpdatedMsg(JsonConverter.convertToAttributesProto(attributesData)); entityDataBuilder2.setPostAttributeScope(DataConstants.SERVER_SCOPE); builder2.addEntityData(entityDataBuilder2.build()); edgeImitator.sendUplinkMsg(builder2.build()); From e09b97100b965343d977300b56794252be96cfd3 Mon Sep 17 00:00:00 2001 From: Bohdan Smetaniuk Date: Fri, 2 Oct 2020 19:29:35 +0300 Subject: [PATCH 4/4] tests refactored --- .../thingsboard/server/edge/BaseEdgeTest.java | 430 +++++++++++------- .../server/edge/imitator/EdgeImitator.java | 52 ++- .../server/edge/imitator/EdgeStorage.java | 142 ------ 3 files changed, 299 insertions(+), 325 deletions(-) delete mode 100644 application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java index db7bb48f75..c5db7a0e97 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java @@ -21,6 +21,7 @@ import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.gson.JsonObject; +import com.google.protobuf.AbstractMessage; import lombok.extern.slf4j.Slf4j; import org.junit.After; import org.junit.Assert; @@ -28,7 +29,6 @@ import org.junit.Before; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.server.common.data.Dashboard; -import org.thingsboard.server.common.data.DashboardInfo; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntityType; @@ -57,10 +57,13 @@ import org.thingsboard.server.controller.AbstractControllerTest; import org.thingsboard.server.dao.edge.EdgeEventService;; import org.thingsboard.server.edge.imitator.EdgeImitator; import org.thingsboard.server.gen.edge.AlarmUpdateMsg; +import org.thingsboard.server.gen.edge.AssetUpdateMsg; +import org.thingsboard.server.gen.edge.DashboardUpdateMsg; import org.thingsboard.server.gen.edge.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.EdgeConfiguration; import org.thingsboard.server.gen.edge.EntityDataProto; import org.thingsboard.server.gen.edge.RelationUpdateMsg; +import org.thingsboard.server.gen.edge.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.UpdateMsgType; import org.thingsboard.server.gen.edge.UplinkMsg; import org.thingsboard.server.gen.transport.TransportProtos; @@ -68,7 +71,6 @@ import org.thingsboard.server.gen.transport.TransportProtos; import java.util.List; import java.util.Map; import java.util.Optional; -import java.util.Set; import java.util.UUID; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @@ -109,7 +111,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { edgeImitator = new EdgeImitator("localhost", 7070, edge.getRoutingKey(), edge.getSecret()); // should be 3, but 3 events from sync service + 3 from controller. will be fixed in next releases - edgeImitator.getStorage().expectMessageAmount(6); + edgeImitator.expectMessageAmount(6); edgeImitator.connect(); } @@ -141,105 +143,142 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { private void testReceivedInitialData() throws Exception { log.info("Checking received data"); - edgeImitator.getStorage().waitForMessages(); + edgeImitator.waitForMessages(); - EdgeConfiguration configuration = edgeImitator.getStorage().getConfiguration(); + EdgeConfiguration configuration = edgeImitator.getConfiguration(); Assert.assertNotNull(configuration); - Map entities = edgeImitator.getStorage().getEntities(); - Assert.assertFalse(entities.isEmpty()); - - Set devices = edgeImitator.getStorage().getEntitiesByType(EntityType.DEVICE); - Assert.assertEquals(1, devices.size()); - TimePageData pageDataDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", - new TypeReference>() {}, new TextPageLink(100)); - for (Device device: pageDataDevices.getData()) { - Assert.assertTrue(devices.contains(device.getUuidId())); - } - - Set assets = edgeImitator.getStorage().getEntitiesByType(EntityType.ASSET); - Assert.assertEquals(1, assets.size()); - TimePageData pageDataAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", - new TypeReference>() {}, new TextPageLink(100)); - for (Asset asset: pageDataAssets.getData()) { - Assert.assertTrue(assets.contains(asset.getUuidId())); - } + Optional optionalMsg1 = edgeImitator.findMessageByType(DeviceUpdateMsg.class); + Assert.assertTrue(optionalMsg1.isPresent()); + DeviceUpdateMsg deviceUpdateMsg = optionalMsg1.get(); + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceUpdateMsg.getMsgType()); + UUID deviceUUID = new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB()); + Device device = doGet("/api/device/" + deviceUUID.toString(), Device.class); + Assert.assertNotNull(device); + List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() {}, new TextPageLink(100)).getData(); + Assert.assertTrue(edgeDevices.contains(device)); + + Optional optionalMsg2 = edgeImitator.findMessageByType(AssetUpdateMsg.class); + Assert.assertTrue(optionalMsg2.isPresent()); + AssetUpdateMsg assetUpdateMsg = optionalMsg2.get(); + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetUpdateMsg.getMsgType()); + UUID assetUUID = new UUID(assetUpdateMsg.getIdMSB(), assetUpdateMsg.getIdLSB()); + Asset asset = doGet("/api/asset/" + assetUUID.toString(), Asset.class); + Assert.assertNotNull(asset); + List edgeAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", + new TypeReference>() {}, new TextPageLink(100)).getData(); + Assert.assertTrue(edgeAssets.contains(asset)); + + Optional optionalMsg3 = edgeImitator.findMessageByType(RuleChainUpdateMsg.class); + Assert.assertTrue(optionalMsg3.isPresent()); + RuleChainUpdateMsg ruleChainUpdateMsg = optionalMsg3.get(); + Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, ruleChainUpdateMsg.getMsgType()); + UUID ruleChainUUID = new UUID(ruleChainUpdateMsg.getIdMSB(), ruleChainUpdateMsg.getIdLSB()); + RuleChain ruleChain = doGet("/api/ruleChain/" + ruleChainUUID.toString(), RuleChain.class); + Assert.assertNotNull(ruleChain); + List edgeRuleChains = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/ruleChains?", + new TypeReference>() {}, new TextPageLink(100)).getData(); + Assert.assertTrue(edgeRuleChains.contains(ruleChain)); - Set ruleChains = edgeImitator.getStorage().getEntitiesByType(EntityType.RULE_CHAIN); - Assert.assertEquals(1, ruleChains.size()); - TimePageData pageDataRuleChains = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/ruleChains?", - new TypeReference>() {}, new TextPageLink(100)); - for (RuleChain ruleChain: pageDataRuleChains.getData()) { - Assert.assertTrue(ruleChains.contains(ruleChain.getUuidId())); - } log.info("Received data checked"); } - private void testDevices() throws Exception { + private void testDevices() throws Exception { log.info("Testing devices"); + Device device = new Device(); device.setName("Edge Device 2"); device.setType("test"); Device savedDevice = doPost("/api/device", device, Device.class); - edgeImitator.getStorage().expectMessageAmount(1); + + edgeImitator.expectMessageAmount(1); doPost("/api/edge/" + edge.getId().getId().toString() + "/device/" + savedDevice.getId().getId().toString(), Device.class); - - TimePageData pageDataDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", - new TypeReference>() {}, new TextPageLink(100)); - Assert.assertTrue(pageDataDevices.getData().contains(savedDevice)); - edgeImitator.getStorage().waitForMessages(); - Set devices = edgeImitator.getStorage().getEntitiesByType(EntityType.DEVICE); - Assert.assertEquals(2, devices.size()); - Assert.assertTrue(devices.contains(savedDevice.getUuidId())); - - edgeImitator.getStorage().expectMessageAmount(1); + edgeImitator.waitForMessages(); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof DeviceUpdateMsg); + DeviceUpdateMsg deviceUpdateMsg = (DeviceUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceUpdateMsg.getMsgType()); + Assert.assertEquals(deviceUpdateMsg.getIdMSB(), savedDevice.getUuidId().getMostSignificantBits()); + Assert.assertEquals(deviceUpdateMsg.getIdLSB(), savedDevice.getUuidId().getLeastSignificantBits()); + Assert.assertEquals(deviceUpdateMsg.getName(), savedDevice.getName()); + Assert.assertEquals(deviceUpdateMsg.getType(), savedDevice.getType()); + + edgeImitator.expectMessageAmount(1); doDelete("/api/edge/" + edge.getId().getId().toString() + "/device/" + savedDevice.getId().getId().toString(), Device.class); - pageDataDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", - new TypeReference>() {}, new TextPageLink(100)); - Assert.assertFalse(pageDataDevices.getData().contains(savedDevice)); - edgeImitator.getStorage().waitForMessages(); - devices = edgeImitator.getStorage().getEntitiesByType(EntityType.DEVICE); - Assert.assertEquals(1, devices.size()); - Assert.assertFalse(devices.contains(savedDevice.getUuidId())); + edgeImitator.waitForMessages(); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof DeviceUpdateMsg); + deviceUpdateMsg = (DeviceUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, deviceUpdateMsg.getMsgType()); + Assert.assertEquals(deviceUpdateMsg.getIdMSB(), savedDevice.getUuidId().getMostSignificantBits()); + Assert.assertEquals(deviceUpdateMsg.getIdLSB(), savedDevice.getUuidId().getLeastSignificantBits()); + + edgeImitator.expectMessageAmount(1); doDelete("/api/device/" + savedDevice.getId().getId().toString()) .andExpect(status().isOk()); + edgeImitator.waitForMessages(); + + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof DeviceUpdateMsg); + deviceUpdateMsg = (DeviceUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, deviceUpdateMsg.getMsgType()); + Assert.assertEquals(deviceUpdateMsg.getIdMSB(), savedDevice.getUuidId().getMostSignificantBits()); + Assert.assertEquals(deviceUpdateMsg.getIdLSB(), savedDevice.getUuidId().getLeastSignificantBits()); + log.info("Devices tested successfully"); } + private void testAssets() throws Exception { log.info("Testing assets"); Asset asset = new Asset(); asset.setName("Edge Asset 2"); asset.setType("test"); Asset savedAsset = doPost("/api/asset", asset, Asset.class); - edgeImitator.getStorage().expectMessageAmount(1); + + edgeImitator.expectMessageAmount(1); doPost("/api/edge/" + edge.getId().getId().toString() + "/asset/" + savedAsset.getId().getId().toString(), Asset.class); - - TimePageData pageDataAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", - new TypeReference>() {}, new TextPageLink(100)); - Assert.assertTrue(pageDataAssets.getData().contains(savedAsset)); - edgeImitator.getStorage().waitForMessages(); - Set assets = edgeImitator.getStorage().getEntitiesByType(EntityType.ASSET); - Assert.assertEquals(2, assets.size()); - Assert.assertTrue(assets.contains(savedAsset.getUuidId())); - - edgeImitator.getStorage().expectMessageAmount(1); + edgeImitator.waitForMessages(); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof AssetUpdateMsg); + AssetUpdateMsg assetUpdateMsg = (AssetUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetUpdateMsg.getMsgType()); + Assert.assertEquals(assetUpdateMsg.getIdMSB(), savedAsset.getUuidId().getMostSignificantBits()); + Assert.assertEquals(assetUpdateMsg.getIdLSB(), savedAsset.getUuidId().getLeastSignificantBits()); + Assert.assertEquals(assetUpdateMsg.getName(), savedAsset.getName()); + Assert.assertEquals(assetUpdateMsg.getType(), savedAsset.getType()); + + edgeImitator.expectMessageAmount(1); doDelete("/api/edge/" + edge.getId().getId().toString() + "/asset/" + savedAsset.getId().getId().toString(), Asset.class); - pageDataAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", - new TypeReference>() {}, new TextPageLink(100)); - Assert.assertFalse(pageDataAssets.getData().contains(savedAsset)); - edgeImitator.getStorage().waitForMessages(); - assets = edgeImitator.getStorage().getEntitiesByType(EntityType.ASSET); - Assert.assertEquals(1, assets.size()); - Assert.assertFalse(assets.contains(savedAsset.getUuidId())); + edgeImitator.waitForMessages(); + + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof AssetUpdateMsg); + assetUpdateMsg = (AssetUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, assetUpdateMsg.getMsgType()); + Assert.assertEquals(assetUpdateMsg.getIdMSB(), savedAsset.getUuidId().getMostSignificantBits()); + Assert.assertEquals(assetUpdateMsg.getIdLSB(), savedAsset.getUuidId().getLeastSignificantBits()); + edgeImitator.expectMessageAmount(1); doDelete("/api/asset/" + savedAsset.getId().getId().toString()) .andExpect(status().isOk()); + edgeImitator.waitForMessages(); + + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof AssetUpdateMsg); + assetUpdateMsg = (AssetUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, assetUpdateMsg.getMsgType()); + Assert.assertEquals(assetUpdateMsg.getIdMSB(), savedAsset.getUuidId().getMostSignificantBits()); + Assert.assertEquals(assetUpdateMsg.getIdLSB(), savedAsset.getUuidId().getLeastSignificantBits()); + log.info("Assets tested successfully"); } @@ -249,33 +288,45 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { ruleChain.setName("Edge Test Rule Chain"); ruleChain.setType(RuleChainType.EDGE); RuleChain savedRuleChain = doPost("/api/ruleChain", ruleChain, RuleChain.class); - edgeImitator.getStorage().expectMessageAmount(1); + + edgeImitator.expectMessageAmount(1); doPost("/api/edge/" + edge.getId().getId().toString() + "/ruleChain/" + savedRuleChain.getId().getId().toString(), RuleChain.class); + edgeImitator.waitForMessages(); - TimePageData pageDataRuleChain = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/ruleChains?", - new TypeReference>() {}, new TextPageLink(100)); - Assert.assertTrue(pageDataRuleChain.getData().contains(savedRuleChain)); - edgeImitator.getStorage().waitForMessages(); - Set ruleChains = edgeImitator.getStorage().getEntitiesByType(EntityType.RULE_CHAIN); - Assert.assertEquals(2, ruleChains.size()); - Assert.assertTrue(ruleChains.contains(savedRuleChain.getUuidId())); + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof RuleChainUpdateMsg); + RuleChainUpdateMsg ruleChainUpdateMsg = (RuleChainUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, ruleChainUpdateMsg.getMsgType()); + Assert.assertEquals(ruleChainUpdateMsg.getIdMSB(), savedRuleChain.getUuidId().getMostSignificantBits()); + Assert.assertEquals(ruleChainUpdateMsg.getIdLSB(), savedRuleChain.getUuidId().getLeastSignificantBits()); + Assert.assertEquals(ruleChainUpdateMsg.getName(), savedRuleChain.getName()); - edgeImitator.getStorage().expectMessageAmount(1); + edgeImitator.expectMessageAmount(1); doDelete("/api/edge/" + edge.getId().getId().toString() + "/ruleChain/" + savedRuleChain.getId().getId().toString(), RuleChain.class); - pageDataRuleChain = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/ruleChains?", - new TypeReference>() {}, new TextPageLink(100)); - Assert.assertFalse(pageDataRuleChain.getData().contains(savedRuleChain)); - edgeImitator.getStorage().waitForMessages(); - ruleChains = edgeImitator.getStorage().getEntitiesByType(EntityType.RULE_CHAIN); - Assert.assertEquals(1, ruleChains.size()); - Assert.assertFalse(ruleChains.contains(savedRuleChain.getUuidId())); + edgeImitator.waitForMessages(); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof RuleChainUpdateMsg); + ruleChainUpdateMsg = (RuleChainUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, ruleChainUpdateMsg.getMsgType()); + Assert.assertEquals(ruleChainUpdateMsg.getIdMSB(), savedRuleChain.getUuidId().getMostSignificantBits()); + Assert.assertEquals(ruleChainUpdateMsg.getIdLSB(), savedRuleChain.getUuidId().getLeastSignificantBits()); + + edgeImitator.expectMessageAmount(1); doDelete("/api/ruleChain/" + savedRuleChain.getId().getId().toString()) .andExpect(status().isOk()); - log.info("RuleChains tested successfully"); + edgeImitator.waitForMessages(); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof RuleChainUpdateMsg); + ruleChainUpdateMsg = (RuleChainUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, ruleChainUpdateMsg.getMsgType()); + Assert.assertEquals(ruleChainUpdateMsg.getIdMSB(), savedRuleChain.getUuidId().getMostSignificantBits()); + Assert.assertEquals(ruleChainUpdateMsg.getIdLSB(), savedRuleChain.getUuidId().getLeastSignificantBits()); + + log.info("RuleChains tested successfully"); } private void testDashboards() throws Exception { @@ -283,33 +334,45 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Dashboard dashboard = new Dashboard(); dashboard.setTitle("Edge Test Dashboard"); Dashboard savedDashboard = doPost("/api/dashboard", dashboard, Dashboard.class); - edgeImitator.getStorage().expectMessageAmount(1); + + edgeImitator.expectMessageAmount(1); doPost("/api/edge/" + edge.getId().getId().toString() + "/dashboard/" + savedDashboard.getId().getId().toString(), Dashboard.class); + edgeImitator.waitForMessages(); - TimePageData pageDataDashboard = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/dashboards?", - new TypeReference>() {}, new TextPageLink(100)); - Assert.assertTrue(pageDataDashboard.getData().stream().allMatch(dashboardInfo -> dashboardInfo.getUuidId().equals(savedDashboard.getUuidId()))); - edgeImitator.getStorage().waitForMessages(); - Set dashboards = edgeImitator.getStorage().getEntitiesByType(EntityType.DASHBOARD); - Assert.assertEquals(1, dashboards.size()); - Assert.assertTrue(dashboards.contains(savedDashboard.getUuidId())); + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof DashboardUpdateMsg); + DashboardUpdateMsg dashboardUpdateMsg = (DashboardUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, dashboardUpdateMsg.getMsgType()); + Assert.assertEquals(dashboardUpdateMsg.getIdMSB(), savedDashboard.getUuidId().getMostSignificantBits()); + Assert.assertEquals(dashboardUpdateMsg.getIdLSB(), savedDashboard.getUuidId().getLeastSignificantBits()); + Assert.assertEquals(dashboardUpdateMsg.getTitle(), savedDashboard.getName()); - edgeImitator.getStorage().expectMessageAmount(1); + edgeImitator.expectMessageAmount(1); doDelete("/api/edge/" + edge.getId().getId().toString() + "/dashboard/" + savedDashboard.getId().getId().toString(), Dashboard.class); - pageDataDashboard = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/dashboards?", - new TypeReference>() {}, new TextPageLink(100)); - Assert.assertFalse(pageDataDashboard.getData().stream().anyMatch(dashboardInfo -> dashboardInfo.getUuidId().equals(savedDashboard.getUuidId()))); - edgeImitator.getStorage().waitForMessages(); - dashboards = edgeImitator.getStorage().getEntitiesByType(EntityType.DASHBOARD); - Assert.assertEquals(0, dashboards.size()); - Assert.assertFalse(dashboards.contains(savedDashboard.getUuidId())); + edgeImitator.waitForMessages(); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof DashboardUpdateMsg); + dashboardUpdateMsg = (DashboardUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, dashboardUpdateMsg.getMsgType()); + Assert.assertEquals(dashboardUpdateMsg.getIdMSB(), savedDashboard.getUuidId().getMostSignificantBits()); + Assert.assertEquals(dashboardUpdateMsg.getIdLSB(), savedDashboard.getUuidId().getLeastSignificantBits()); + + edgeImitator.expectMessageAmount(1); doDelete("/api/dashboard/" + savedDashboard.getId().getId().toString()) .andExpect(status().isOk()); - log.info("Dashboards tested successfully"); + edgeImitator.waitForMessages(); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof DashboardUpdateMsg); + dashboardUpdateMsg = (DashboardUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, dashboardUpdateMsg.getMsgType()); + Assert.assertEquals(dashboardUpdateMsg.getIdMSB(), savedDashboard.getUuidId().getMostSignificantBits()); + Assert.assertEquals(dashboardUpdateMsg.getIdLSB(), savedDashboard.getUuidId().getLeastSignificantBits()); + + log.info("Dashboards tested successfully"); } private void testRelations() throws Exception { @@ -331,14 +394,24 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { relation.setFrom(device.getId()); relation.setTo(asset.getId()); relation.setTypeGroup(RelationTypeGroup.COMMON); - edgeImitator.getStorage().expectMessageAmount(1); - doPost("/api/relation", relation); - edgeImitator.getStorage().waitForMessages(); - List relations = edgeImitator.getStorage().getRelations(); - Assert.assertEquals(1, relations.size()); - Assert.assertTrue(relations.contains(relation)); - edgeImitator.getStorage().expectMessageAmount(1); + edgeImitator.expectMessageAmount(1); + doPost("/api/relation", relation); + edgeImitator.waitForMessages(); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof RelationUpdateMsg); + RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); + Assert.assertEquals(relationUpdateMsg.getType(), relation.getType()); + Assert.assertEquals(relationUpdateMsg.getFromIdMSB(), relation.getFrom().getId().getMostSignificantBits()); + Assert.assertEquals(relationUpdateMsg.getFromIdLSB(), relation.getFrom().getId().getLeastSignificantBits()); + Assert.assertEquals(relationUpdateMsg.getToEntityType(), relation.getTo().getEntityType().name());Assert.assertEquals(relationUpdateMsg.getFromIdMSB(), relation.getFrom().getId().getMostSignificantBits()); + Assert.assertEquals(relationUpdateMsg.getToIdLSB(), relation.getTo().getId().getLeastSignificantBits()); + Assert.assertEquals(relationUpdateMsg.getToEntityType(), relation.getTo().getEntityType().name()); + Assert.assertEquals(relationUpdateMsg.getTypeGroup(), relation.getTypeGroup().name()); + + edgeImitator.expectMessageAmount(1); doDelete("/api/relation?" + "fromId=" + relation.getFrom().getId().toString() + "&fromType=" + relation.getFrom().getEntityType().name() + @@ -347,15 +420,23 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { "&toId=" + relation.getTo().getId().toString() + "&toType=" + relation.getTo().getEntityType().name()) .andExpect(status().isOk()); + edgeImitator.waitForMessages(); + + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof RelationUpdateMsg); + relationUpdateMsg = (RelationUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); + Assert.assertEquals(relationUpdateMsg.getType(), relation.getType()); + Assert.assertEquals(relationUpdateMsg.getFromIdMSB(), relation.getFrom().getId().getMostSignificantBits()); + Assert.assertEquals(relationUpdateMsg.getFromIdLSB(), relation.getFrom().getId().getLeastSignificantBits()); + Assert.assertEquals(relationUpdateMsg.getToEntityType(), relation.getTo().getEntityType().name());Assert.assertEquals(relationUpdateMsg.getFromIdMSB(), relation.getFrom().getId().getMostSignificantBits()); + Assert.assertEquals(relationUpdateMsg.getToIdLSB(), relation.getTo().getId().getLeastSignificantBits()); + Assert.assertEquals(relationUpdateMsg.getToEntityType(), relation.getTo().getEntityType().name()); + Assert.assertEquals(relationUpdateMsg.getTypeGroup(), relation.getTypeGroup().name()); - edgeImitator.getStorage().waitForMessages(); - relations = edgeImitator.getStorage().getRelations(); - Assert.assertEquals(0, relations.size()); - Assert.assertFalse(relations.contains(relation)); log.info("Relations tested successfully"); } - private void testAlarms() throws Exception { log.info("Testing Alarms"); List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", @@ -370,29 +451,45 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { alarm.setType("alarm"); alarm.setSeverity(AlarmSeverity.CRITICAL); - edgeImitator.getStorage().expectMessageAmount(1); + edgeImitator.expectMessageAmount(1); Alarm savedAlarm = doPost("/api/alarm", alarm, Alarm.class); - AlarmInfo alarmInfo = doGet("/api/alarm/info/" + savedAlarm.getId().getId().toString(), AlarmInfo.class); - edgeImitator.getStorage().waitForMessages(); - - Assert.assertEquals(1, edgeImitator.getStorage().getAlarms().size()); - Assert.assertTrue(edgeImitator.getStorage().getAlarms().containsKey(alarmInfo.getType())); - Assert.assertEquals(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()), alarmInfo.getStatus()); - edgeImitator.getStorage().expectMessageAmount(1); + edgeImitator.waitForMessages(); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg); + AlarmUpdateMsg alarmUpdateMsg = (AlarmUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, alarmUpdateMsg.getMsgType()); + Assert.assertEquals(alarmUpdateMsg.getType(), savedAlarm.getType()); + Assert.assertEquals(alarmUpdateMsg.getName(), savedAlarm.getName()); + Assert.assertEquals(alarmUpdateMsg.getOriginatorName(), device.getName()); + Assert.assertEquals(alarmUpdateMsg.getStatus(), savedAlarm.getStatus().name()); + Assert.assertEquals(alarmUpdateMsg.getSeverity(), savedAlarm.getSeverity().name()); + + edgeImitator.expectMessageAmount(1); doPost("/api/alarm/" + savedAlarm.getId().getId().toString() + "/ack"); - - edgeImitator.getStorage().waitForMessages(); - alarmInfo = doGet("/api/alarm/info/" + savedAlarm.getId().getId().toString(), AlarmInfo.class); - Assert.assertTrue(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()).isAck()); - Assert.assertEquals(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()), alarmInfo.getStatus()); - edgeImitator.getStorage().expectMessageAmount(1); + edgeImitator.waitForMessages(); + + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg); + alarmUpdateMsg = (AlarmUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ALARM_ACK_RPC_MESSAGE, alarmUpdateMsg.getMsgType()); + Assert.assertEquals(alarmUpdateMsg.getType(), savedAlarm.getType()); + Assert.assertEquals(alarmUpdateMsg.getName(), savedAlarm.getName()); + Assert.assertEquals(alarmUpdateMsg.getOriginatorName(), device.getName()); + Assert.assertEquals(alarmUpdateMsg.getStatus(), AlarmStatus.ACTIVE_ACK.name()); + + edgeImitator.expectMessageAmount(1); doPost("/api/alarm/" + savedAlarm.getId().getId().toString() + "/clear"); + edgeImitator.waitForMessages(); - edgeImitator.getStorage().waitForMessages(); - alarmInfo = doGet("/api/alarm/info/" + savedAlarm.getId().getId().toString(), AlarmInfo.class); - Assert.assertTrue(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()).isAck()); - Assert.assertTrue(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()).isCleared()); - Assert.assertEquals(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()), alarmInfo.getStatus()); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg); + alarmUpdateMsg = (AlarmUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ALARM_CLEAR_RPC_MESSAGE, alarmUpdateMsg.getMsgType()); + Assert.assertEquals(alarmUpdateMsg.getType(), savedAlarm.getType()); + Assert.assertEquals(alarmUpdateMsg.getName(), savedAlarm.getName()); + Assert.assertEquals(alarmUpdateMsg.getOriginatorName(), device.getName()); + Assert.assertEquals(alarmUpdateMsg.getStatus(), AlarmStatus.CLEARED_ACK.name()); doDelete("/api/alarm/" + savedAlarm.getId().getId().toString()) .andExpect(status().isOk()); @@ -411,15 +508,16 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { String timeseriesData = "{\"data\":{\"temperature\":25},\"ts\":" + System.currentTimeMillis() + "}"; JsonNode timeseriesEntityData = mapper.readTree(timeseriesData); EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), ActionType.TIMESERIES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData); - edgeImitator.getStorage().expectMessageAmount(1); + edgeImitator.expectMessageAmount(1); edgeEventService.saveAsync(edgeEvent1); - edgeImitator.getStorage().waitForMessages(); - - EntityDataProto latestEntityDataMsg = edgeImitator.getStorage().getLatestEntityDataMsg(); - Assert.assertNotNull(latestEntityDataMsg); - UUID uuid = new UUID(latestEntityDataMsg.getEntityIdMSB(), latestEntityDataMsg.getEntityIdLSB()); - Assert.assertEquals(device.getId().getId(), uuid); - Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); + edgeImitator.waitForMessages(); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof EntityDataProto); + EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage; + Assert.assertEquals(latestEntityDataMsg.getEntityIdMSB(), device.getUuidId().getMostSignificantBits()); + Assert.assertEquals(latestEntityDataMsg.getEntityIdLSB(), device.getUuidId().getLeastSignificantBits()); + Assert.assertEquals(latestEntityDataMsg.getEntityType(), device.getId().getEntityType().name()); Assert.assertTrue(latestEntityDataMsg.hasPostTelemetryMsg()); TransportProtos.PostTelemetryMsg postTelemetryMsg = latestEntityDataMsg.getPostTelemetryMsg(); @@ -430,7 +528,6 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { TransportProtos.KeyValueProto keyValueProto = tsKvListProto.getKv(0); Assert.assertEquals("temperature", keyValueProto.getKey()); Assert.assertEquals(25, keyValueProto.getLongV()); - edgeImitator.getStorage().setLatestEntityDataMsg(null); log.info("Timeseries tested successfully"); } @@ -445,16 +542,17 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { String attributesData = "{\"scope\":\"SERVER_SCOPE\",\"kv\":{\"key\":\"value\"}}"; JsonNode attributesEntityData = mapper.readTree(attributesData); EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), ActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); - edgeImitator.getStorage().expectMessageAmount(1); + edgeImitator.expectMessageAmount(1); edgeEventService.saveAsync(edgeEvent1); - edgeImitator.getStorage().waitForMessages(); - - EntityDataProto latestEntityDataMsg = edgeImitator.getStorage().getLatestEntityDataMsg(); - Assert.assertNotNull(latestEntityDataMsg); - UUID uuid = new UUID(latestEntityDataMsg.getEntityIdMSB(), latestEntityDataMsg.getEntityIdLSB()); - Assert.assertEquals(device.getId().getId(), uuid); - Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); - Assert.assertEquals(attributesEntityData.get("scope").asText(), latestEntityDataMsg.getPostAttributeScope()); + edgeImitator.waitForMessages(); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof EntityDataProto); + EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage; + Assert.assertEquals(latestEntityDataMsg.getEntityIdMSB(), device.getUuidId().getMostSignificantBits()); + Assert.assertEquals(latestEntityDataMsg.getEntityIdLSB(), device.getUuidId().getLeastSignificantBits()); + Assert.assertEquals(latestEntityDataMsg.getEntityType(), device.getId().getEntityType().name()); + Assert.assertEquals(latestEntityDataMsg.getPostAttributeScope(), attributesEntityData.get("scope").asText()); Assert.assertTrue(latestEntityDataMsg.hasAttributesUpdatedMsg()); TransportProtos.PostAttributeMsg attributesUpdatedMsg = latestEntityDataMsg.getAttributesUpdatedMsg(); @@ -462,28 +560,27 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { TransportProtos.KeyValueProto keyValueProto = attributesUpdatedMsg.getKv(0); Assert.assertEquals("key", keyValueProto.getKey()); Assert.assertEquals("value", keyValueProto.getStringV()); - edgeImitator.getStorage().setLatestEntityDataMsg(null); ((ObjectNode) attributesEntityData).put("isPostAttributes", true); EdgeEvent edgeEvent2 = constructEdgeEvent(tenantId, edge.getId(), ActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); - edgeImitator.getStorage().expectMessageAmount(1); + edgeImitator.expectMessageAmount(1); edgeEventService.saveAsync(edgeEvent2); - edgeImitator.getStorage().waitForMessages(); - - latestEntityDataMsg = edgeImitator.getStorage().getLatestEntityDataMsg(); - Assert.assertNotNull(latestEntityDataMsg); - uuid = new UUID(latestEntityDataMsg.getEntityIdMSB(), latestEntityDataMsg.getEntityIdLSB()); - Assert.assertEquals(device.getId().getId(), uuid); - Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); - Assert.assertEquals(attributesEntityData.get("scope").asText(), latestEntityDataMsg.getPostAttributeScope()); + edgeImitator.waitForMessages(); + + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof EntityDataProto); + latestEntityDataMsg = (EntityDataProto) latestMessage; + Assert.assertEquals(latestEntityDataMsg.getEntityIdMSB(), device.getUuidId().getMostSignificantBits()); + Assert.assertEquals(latestEntityDataMsg.getEntityIdLSB(), device.getUuidId().getLeastSignificantBits()); + Assert.assertEquals(latestEntityDataMsg.getEntityType(), device.getId().getEntityType().name()); + Assert.assertEquals(latestEntityDataMsg.getPostAttributeScope(), attributesEntityData.get("scope").asText()); Assert.assertTrue(latestEntityDataMsg.hasPostAttributesMsg()); - TransportProtos.PostAttributeMsg postAttributeMsg = latestEntityDataMsg.getPostAttributesMsg(); - Assert.assertEquals(1, postAttributeMsg.getKvCount()); - keyValueProto = postAttributeMsg.getKv(0); + attributesUpdatedMsg = latestEntityDataMsg.getPostAttributesMsg(); + Assert.assertEquals(1, attributesUpdatedMsg.getKvCount()); + keyValueProto = attributesUpdatedMsg.getKv(0); Assert.assertEquals("key", keyValueProto.getKey()); Assert.assertEquals("value", keyValueProto.getStringV()); - edgeImitator.getStorage().setLatestEntityDataMsg(null); log.info("Attributes tested successfully"); } @@ -542,11 +639,6 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { new TypeReference>() {}, new TextPageLink(100), device.getId().getEntityType().name(), device.getId().getId().toString()) .getData(); - - for (AlarmInfo alarmInfo: alarms) { - log.info(String.valueOf(alarmInfo)); - } - Optional foundAlarm = alarms.stream().filter(alarm -> alarm.getType().equals("alarm from edge")).findAny(); Assert.assertTrue(foundAlarm.isPresent()); AlarmInfo alarmInfo = foundAlarm.get(); diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java index 90ec80838f..4cdc7ba84c 100644 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java +++ b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java @@ -19,12 +19,12 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; +import com.google.protobuf.AbstractMessage; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.checkerframework.checker.nullness.qual.Nullable; import org.thingsboard.edge.rpc.EdgeGrpcClient; import org.thingsboard.edge.rpc.EdgeRpcClient; -import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.gen.edge.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.AssetUpdateMsg; import org.thingsboard.server.gen.edge.DashboardUpdateMsg; @@ -41,7 +41,7 @@ import org.thingsboard.server.gen.edge.UplinkResponseMsg; import java.lang.reflect.Field; import java.util.ArrayList; import java.util.List; -import java.util.UUID; +import java.util.Optional; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -53,16 +53,19 @@ public class EdgeImitator { private EdgeRpcClient edgeRpcClient; + private CountDownLatch messagesLatch; private CountDownLatch responsesLatch; @Getter - private EdgeStorage storage; - + private EdgeConfiguration configuration; + @Getter + private List downlinkMsgs; public EdgeImitator(String host, int port, String routingKey, String routingSecret) throws NoSuchFieldException, IllegalAccessException { edgeRpcClient = new EdgeGrpcClient(); - storage = new EdgeStorage(); + messagesLatch = new CountDownLatch(0); responsesLatch = new CountDownLatch(0); + downlinkMsgs = new ArrayList<>(); this.routingKey = routingKey; this.routingSecret = routingSecret; setEdgeCredentials("rpcHost", host); @@ -101,7 +104,7 @@ public class EdgeImitator { } private void onEdgeUpdate(EdgeConfiguration edgeConfiguration) { - storage.setConfiguration(edgeConfiguration); + this.configuration = edgeConfiguration; } private void onDownlink(DownlinkMsg downlinkMsg) { @@ -129,47 +132,68 @@ public class EdgeImitator { List> result = new ArrayList<>(); if (downlinkMsg.getDeviceUpdateMsgList() != null && !downlinkMsg.getDeviceUpdateMsgList().isEmpty()) { for (DeviceUpdateMsg deviceUpdateMsg: downlinkMsg.getDeviceUpdateMsgList()) { - result.add(storage.processEntity(deviceUpdateMsg.getMsgType(), EntityType.DEVICE, new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB()))); + saveDownlinkMsg(deviceUpdateMsg); } } if (downlinkMsg.getAssetUpdateMsgList() != null && !downlinkMsg.getAssetUpdateMsgList().isEmpty()) { for (AssetUpdateMsg assetUpdateMsg: downlinkMsg.getAssetUpdateMsgList()) { - result.add(storage.processEntity(assetUpdateMsg.getMsgType(), EntityType.ASSET, new UUID(assetUpdateMsg.getIdMSB(), assetUpdateMsg.getIdLSB()))); + saveDownlinkMsg(assetUpdateMsg); } } if (downlinkMsg.getRuleChainUpdateMsgList() != null && !downlinkMsg.getRuleChainUpdateMsgList().isEmpty()) { for (RuleChainUpdateMsg ruleChainUpdateMsg: downlinkMsg.getRuleChainUpdateMsgList()) { - result.add(storage.processEntity(ruleChainUpdateMsg.getMsgType(), EntityType.RULE_CHAIN, new UUID(ruleChainUpdateMsg.getIdMSB(), ruleChainUpdateMsg.getIdLSB()))); + saveDownlinkMsg(ruleChainUpdateMsg); } } if (downlinkMsg.getDashboardUpdateMsgList() != null && !downlinkMsg.getDashboardUpdateMsgList().isEmpty()) { for (DashboardUpdateMsg dashboardUpdateMsg: downlinkMsg.getDashboardUpdateMsgList()) { - result.add(storage.processEntity(dashboardUpdateMsg.getMsgType(), EntityType.DASHBOARD, new UUID(dashboardUpdateMsg.getIdMSB(), dashboardUpdateMsg.getIdLSB()))); + saveDownlinkMsg(dashboardUpdateMsg); } } if (downlinkMsg.getRelationUpdateMsgList() != null && !downlinkMsg.getRelationUpdateMsgList().isEmpty()) { for (RelationUpdateMsg relationUpdateMsg: downlinkMsg.getRelationUpdateMsgList()) { - result.add(storage.processRelation(relationUpdateMsg)); + saveDownlinkMsg(relationUpdateMsg); } } if (downlinkMsg.getAlarmUpdateMsgList() != null && !downlinkMsg.getAlarmUpdateMsgList().isEmpty()) { for (AlarmUpdateMsg alarmUpdateMsg: downlinkMsg.getAlarmUpdateMsgList()) { - result.add(storage.processAlarm(alarmUpdateMsg)); + saveDownlinkMsg(alarmUpdateMsg); } } if (downlinkMsg.getEntityDataList() != null && !downlinkMsg.getEntityDataList().isEmpty()) { for (EntityDataProto entityData: downlinkMsg.getEntityDataList()) { - result.add(storage.processEntityData(entityData)); + saveDownlinkMsg(entityData); } } return Futures.allAsList(result); } - public void waitForResponses() throws InterruptedException { responsesLatch.await(5, TimeUnit.SECONDS); + private ListenableFuture saveDownlinkMsg(AbstractMessage message) { + downlinkMsgs.add(message); + messagesLatch.countDown(); + return Futures.immediateFuture(null); + } + + public void waitForMessages() throws InterruptedException { + messagesLatch.await(5, TimeUnit.SECONDS); + } + + public void expectMessageAmount(int messageAmount) { + messagesLatch = new CountDownLatch(messageAmount); } + public void waitForResponses() throws InterruptedException { responsesLatch.await(5, TimeUnit.SECONDS); } + public void expectResponsesAmount(int messageAmount) { responsesLatch = new CountDownLatch(messageAmount); } + public Optional findMessageByType(Class tClass) { + return (Optional) downlinkMsgs.stream().filter(downlinkMsg -> downlinkMsg.getClass().isAssignableFrom(tClass)).findAny(); + } + + public AbstractMessage getLatestMessage() { + return downlinkMsgs.get(downlinkMsgs.size() - 1); + } + } diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java deleted file mode 100644 index a8301fda33..0000000000 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java +++ /dev/null @@ -1,142 +0,0 @@ -/** - * Copyright © 2016-2020 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.server.edge.imitator; - -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; -import lombok.Getter; -import lombok.Setter; -import lombok.extern.slf4j.Slf4j; -import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.alarm.AlarmStatus; -import org.thingsboard.server.common.data.id.EntityIdFactory; -import org.thingsboard.server.common.data.relation.EntityRelation; -import org.thingsboard.server.common.data.relation.RelationTypeGroup; -import org.thingsboard.server.gen.edge.AlarmUpdateMsg; -import org.thingsboard.server.gen.edge.EdgeConfiguration; -import org.thingsboard.server.gen.edge.EntityDataProto; -import org.thingsboard.server.gen.edge.RelationUpdateMsg; -import org.thingsboard.server.gen.edge.UpdateMsgType; - -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.Set; -import java.util.UUID; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.TimeUnit; -import java.util.stream.Collectors; - -@Slf4j -@Getter -@Setter -public class EdgeStorage { - - private EdgeConfiguration configuration; - - private CountDownLatch messagesLatch; - - private Map entities; - private Map alarms; - private List relations; - - private EntityDataProto latestEntityDataMsg; - - public EdgeStorage() { - messagesLatch = new CountDownLatch(0); - entities = new HashMap<>(); - alarms = new HashMap<>(); - relations = new ArrayList<>(); - latestEntityDataMsg = null; - } - - public ListenableFuture processEntity(UpdateMsgType msgType, EntityType type, UUID uuid) { - switch (msgType) { - case ENTITY_CREATED_RPC_MESSAGE: - case ENTITY_UPDATED_RPC_MESSAGE: - entities.put(uuid, type); - messagesLatch.countDown(); - break; - case ENTITY_DELETED_RPC_MESSAGE: - if (entities.remove(uuid) != null) { - messagesLatch.countDown(); - } - break; - } - return Futures.immediateFuture(null); - } - - public ListenableFuture processRelation(RelationUpdateMsg relationMsg) { - boolean result = false; - EntityRelation relation = new EntityRelation(); - relation.setType(relationMsg.getType()); - relation.setTypeGroup(RelationTypeGroup.valueOf(relationMsg.getTypeGroup())); - relation.setTo(EntityIdFactory.getByTypeAndUuid(relationMsg.getToEntityType(), new UUID(relationMsg.getToIdMSB(), relationMsg.getToIdLSB()))); - relation.setFrom(EntityIdFactory.getByTypeAndUuid(relationMsg.getFromEntityType(), new UUID(relationMsg.getFromIdMSB(), relationMsg.getFromIdLSB()))); - switch (relationMsg.getMsgType()) { - case ENTITY_CREATED_RPC_MESSAGE: - case ENTITY_UPDATED_RPC_MESSAGE: - result = relations.add(relation); - break; - case ENTITY_DELETED_RPC_MESSAGE: - result = relations.remove(relation); - break; - } - if (result) { - messagesLatch.countDown(); - } - return Futures.immediateFuture(null); - } - - public ListenableFuture processAlarm(AlarmUpdateMsg alarmMsg) { - switch (alarmMsg.getMsgType()) { - case ENTITY_CREATED_RPC_MESSAGE: - case ENTITY_UPDATED_RPC_MESSAGE: - case ALARM_ACK_RPC_MESSAGE: - case ALARM_CLEAR_RPC_MESSAGE: - alarms.put(alarmMsg.getType(), AlarmStatus.valueOf(alarmMsg.getStatus())); - messagesLatch.countDown(); - break; - case ENTITY_DELETED_RPC_MESSAGE: - if (alarms.remove(alarmMsg.getName()) != null) { - messagesLatch.countDown(); - } - break; - } - return Futures.immediateFuture(null); - } - - public ListenableFuture processEntityData(EntityDataProto entityData) { - latestEntityDataMsg = entityData; - messagesLatch.countDown(); - return Futures.immediateFuture(null); - } - - public Set getEntitiesByType(EntityType type) { - return entities.entrySet().stream() - .filter(entry -> entry.getValue().equals(type)) - .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)).keySet(); - } - - public void waitForMessages() throws InterruptedException { - messagesLatch.await(5, TimeUnit.SECONDS); - } - - public void expectMessageAmount(int messageAmount) { - messagesLatch = new CountDownLatch(messageAmount); - } -}