From 9e215251584092941187f587a68a064b48e95cce Mon Sep 17 00:00:00 2001 From: Bohdan Smetaniuk Date: Tue, 29 Sep 2020 13:43:35 +0300 Subject: [PATCH] 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)); }