diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 35f045d66b..1dea607dc4 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -335,6 +335,7 @@ public final class EdgeGrpcSession implements Closeable { downlinkMsg = processEntityMessage(edgeEvent, edgeEvent.getAction()); break; case ATTRIBUTES_UPDATED: + case POST_ATTRIBUTES: case ATTRIBUTES_DELETED: case TIMESERIES_UPDATED: downlinkMsg = processTelemetryMessage(edgeEvent); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java index fb005d7c3a..63f0a55ea3 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java @@ -89,9 +89,11 @@ public class DeviceMsgConstructor { .setRequestIdMSB(request.getRequestUUID().getMostSignificantBits()) .setRequestIdLSB(request.getRequestUUID().getLeastSignificantBits()) .setExpirationTime(request.getExpirationTime()) - .setOriginServiceId(request.getOriginServiceId()) .setOneway(request.isOneway()) .setRequestMsg(requestBuilder.build()); + if (request.getOriginServiceId() != null) { + builder.setOriginServiceId(request.getOriginServiceId()); + } return builder.build(); } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java index 19fa3f9036..469207a0cc 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java @@ -60,15 +60,21 @@ public class EntityDataMsgConstructor { case ATTRIBUTES_UPDATED: try { JsonObject data = entityData.getAsJsonObject(); - TransportProtos.PostAttributeMsg postAttributeMsg = JsonConverter.convertToAttributesProto(data.getAsJsonObject("kv")); - if (data.has("isPostAttributes") && data.getAsJsonPrimitive("isPostAttributes").getAsBoolean()) { - builder.setPostAttributesMsg(postAttributeMsg); - } else { - builder.setAttributesUpdatedMsg(postAttributeMsg); - } + TransportProtos.PostAttributeMsg attributesUpdatedMsg = JsonConverter.convertToAttributesProto(data.getAsJsonObject("kv")); + builder.setAttributesUpdatedMsg(attributesUpdatedMsg); + builder.setPostAttributeScope(data.getAsJsonPrimitive("scope").getAsString()); + } catch (Exception e) { + log.warn("[{}] Can't convert to AttributesUpdatedMsg proto, entityData [{}]", entityId, entityData, e); + } + break; + case POST_ATTRIBUTES: + try { + JsonObject data = entityData.getAsJsonObject(); + TransportProtos.PostAttributeMsg postAttributesMsg = JsonConverter.convertToAttributesProto(data.getAsJsonObject("kv")); + builder.setPostAttributesMsg(postAttributesMsg); builder.setPostAttributeScope(data.getAsJsonPrimitive("scope").getAsString()); } catch (Exception e) { - log.warn("[{}] Can't convert to attributes proto, entityData [{}]", entityId, entityData, e); + log.warn("[{}] Can't convert to PostAttributesMsg, entityData [{}]", entityId, entityData, e); } break; case ATTRIBUTES_DELETED: 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 ee5161869a..5095d32566 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java @@ -16,17 +16,22 @@ package org.thingsboard.server.edge; import com.datastax.driver.core.utils.UUIDs; +import com.fasterxml.jackson.core.JsonProcessingException; 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 com.google.protobuf.AbstractMessage; +import com.google.protobuf.InvalidProtocolBufferException; +import com.google.protobuf.MessageLite; import lombok.extern.slf4j.Slf4j; import org.junit.After; import org.junit.Assert; import org.junit.Before; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; +import org.thingsboard.rule.engine.api.RuleEngineDeviceRpcRequest; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.DataConstants; @@ -45,6 +50,8 @@ import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; @@ -53,9 +60,12 @@ import org.thingsboard.server.common.data.page.TimePageData; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleChain; +import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.rule.RuleChainType; +import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.DeviceCredentials; +import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.data.widget.WidgetType; import org.thingsboard.server.common.data.widget.WidgetsBundle; import org.thingsboard.server.common.transport.adaptor.JsonConverter; @@ -65,15 +75,20 @@ import org.thingsboard.server.dao.util.mapping.JacksonUtil; 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.AttributeDeleteMsg; +import org.thingsboard.server.gen.edge.AttributesRequestMsg; import org.thingsboard.server.gen.edge.CustomerUpdateMsg; import org.thingsboard.server.gen.edge.DashboardUpdateMsg; import org.thingsboard.server.gen.edge.DeviceCredentialsRequestMsg; import org.thingsboard.server.gen.edge.DeviceCredentialsUpdateMsg; +import org.thingsboard.server.gen.edge.DeviceRpcCallMsg; 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.EntityViewUpdateMsg; +import org.thingsboard.server.gen.edge.RelationRequestMsg; import org.thingsboard.server.gen.edge.RelationUpdateMsg; +import org.thingsboard.server.gen.edge.RpcResponseMsg; import org.thingsboard.server.gen.edge.RuleChainMetadataRequestMsg; import org.thingsboard.server.gen.edge.RuleChainMetadataUpdateMsg; import org.thingsboard.server.gen.edge.RuleChainUpdateMsg; @@ -85,16 +100,16 @@ import org.thingsboard.server.gen.edge.WidgetTypeUpdateMsg; import org.thingsboard.server.gen.edge.WidgetsBundleUpdateMsg; import org.thingsboard.server.gen.transport.TransportProtos; +import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.Random; import java.util.UUID; +import java.util.concurrent.TimeUnit; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; -; - - @Slf4j abstract public class BaseEdgeTest extends AbstractControllerTest { @@ -137,7 +152,6 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { @After public void afterTest() throws Exception { edgeImitator.disconnect(); - uninstallation(); loginSysAdmin(); @@ -161,6 +175,69 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { testTimeseries(); testAttributes(); testSendMessagesToCloud(); + testRpcCall(); + } + + private Device findDeviceByName(String deviceName) throws Exception { + List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() {}, new TextPageLink(100)).getData(); + Optional foundDevice = edgeDevices.stream().filter(d -> d.getName().equals(deviceName)).findAny(); + Assert.assertTrue(foundDevice.isPresent()); + Device device = foundDevice.get(); + Assert.assertEquals(deviceName, device.getName()); + return device; + } + + private Asset findAssetByName(String assetName) throws Exception { + List edgeAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", + new TypeReference>() {}, new TextPageLink(100)).getData(); + + Assert.assertEquals(1, edgeAssets.size()); + Asset asset = edgeAssets.get(0); + Assert.assertEquals(assetName, asset.getName()); + return asset; + } + + private Device saveDevice(String deviceName) throws Exception { + Device device = new Device(); + device.setName(deviceName); + device.setType("test"); + return doPost("/api/device", device, Device.class); + } + + private Asset saveAsset(String assetName) throws Exception { + Asset asset = new Asset(); + asset.setName(assetName); + asset.setType("test"); + return doPost("/api/asset", asset, Asset.class); + } + + private void testRpcCall() throws Exception { + Device device = findDeviceByName("Edge Device 1"); + + RuleEngineDeviceRpcRequest request = RuleEngineDeviceRpcRequest.builder() + .oneway(true) + .method("test_method") + .body("{\"param1\":\"value1\"}") + .tenantId(device.getTenantId()) + .deviceId(device.getId()) + .requestId(new Random().nextInt()) + .requestUUID(UUIDs.timeBased()) + .originServiceId("originServiceId") + .expirationTime(System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10)) + .restApiCall(true) + .build(); + + JsonNode body = mapper.valueToTree(request); + EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body); + edgeImitator.expectMessageAmount(1); + edgeEventService.saveAsync(edgeEvent); + edgeImitator.waitForMessages(); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof DeviceRpcCallMsg); + DeviceRpcCallMsg latestDeviceRpcCallMsg = (DeviceRpcCallMsg) latestMessage; + Assert.assertEquals("test_method", latestDeviceRpcCallMsg.getRequestMsg().getMethod()); } private void testReceivedInitialData() throws Exception { @@ -170,6 +247,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { EdgeConfiguration configuration = edgeImitator.getConfiguration(); Assert.assertNotNull(configuration); + testAutoGeneratedCodeByProtobuf(configuration); + UserId userId = edgeImitator.getUserId(); Assert.assertNotNull(userId); @@ -195,6 +274,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { new TypeReference>() {}, new TextPageLink(100)).getData(); Assert.assertTrue(edgeAssets.contains(asset)); + testAutoGeneratedCodeByProtobuf(assetUpdateMsg); + Optional optionalMsg3 = edgeImitator.findMessageByType(RuleChainUpdateMsg.class); Assert.assertTrue(optionalMsg3.isPresent()); RuleChainUpdateMsg ruleChainUpdateMsg = optionalMsg3.get(); @@ -206,16 +287,15 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { new TypeReference>() {}, new TextPageLink(100)).getData(); Assert.assertTrue(edgeRuleChains.contains(ruleChain)); + testAutoGeneratedCodeByProtobuf(ruleChainUpdateMsg); + log.info("Received data checked"); } 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); + Device savedDevice = saveDevice("Edge Device 2"); edgeImitator.expectMessageAmount(1); doPost("/api/edge/" + edge.getId().getId().toString() @@ -261,10 +341,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { 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); + Asset savedAsset = saveAsset("Edge Asset 2"); edgeImitator.expectMessageAmount(1); doPost("/api/edge/" + edge.getId().getId().toString() @@ -314,6 +391,11 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { ruleChain.setType(RuleChainType.EDGE); RuleChain savedRuleChain = doPost("/api/ruleChain", ruleChain, RuleChain.class); + createRuleChainMetadata(savedRuleChain); + + // Wait before rule chain metadata saved to database before rule chain is assigned to edge + Thread.sleep(1000); + edgeImitator.expectMessageAmount(1); doPost("/api/edge/" + edge.getId().getId().toString() + "/ruleChain/" + savedRuleChain.getId().getId().toString(), RuleChain.class); @@ -327,6 +409,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertEquals(ruleChainUpdateMsg.getIdLSB(), savedRuleChain.getUuidId().getLeastSignificantBits()); Assert.assertEquals(ruleChainUpdateMsg.getName(), savedRuleChain.getName()); + testRuleChainMetadataRequestMsg(savedRuleChain.getId()); + edgeImitator.expectMessageAmount(1); doDelete("/api/edge/" + edge.getId().getId().toString() + "/ruleChain/" + savedRuleChain.getId().getId().toString(), RuleChain.class); @@ -354,6 +438,67 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { log.info("RuleChains tested successfully"); } + private void testRuleChainMetadataRequestMsg(RuleChainId ruleChainId) throws Exception { + RuleChainMetadataRequestMsg.Builder ruleChainMetadataRequestMsgBuilder = RuleChainMetadataRequestMsg.newBuilder() + .setRuleChainIdMSB(ruleChainId.getId().getMostSignificantBits()) + .setRuleChainIdLSB(ruleChainId.getId().getLeastSignificantBits()); + testAutoGeneratedCodeByProtobuf(ruleChainMetadataRequestMsgBuilder); + + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder() + .addRuleChainMetadataRequestMsg(ruleChainMetadataRequestMsgBuilder.build()); + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.expectMessageAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + edgeImitator.waitForResponses(); + edgeImitator.waitForMessages(); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof RuleChainMetadataUpdateMsg); + RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = (RuleChainMetadataUpdateMsg) latestMessage; + RuleChainId receivedRuleChainId = + new RuleChainId(new UUID(ruleChainMetadataUpdateMsg.getRuleChainIdMSB(), ruleChainMetadataUpdateMsg.getRuleChainIdLSB())); + Assert.assertEquals(ruleChainId, receivedRuleChainId); + } + + private void createRuleChainMetadata(RuleChain ruleChain) throws Exception { + RuleChainMetaData ruleChainMetaData = new RuleChainMetaData(); + ruleChainMetaData.setRuleChainId(ruleChain.getId()); + + ObjectMapper mapper = new ObjectMapper(); + + RuleNode ruleNode1 = new RuleNode(); + ruleNode1.setName("name1"); + ruleNode1.setType("type1"); + ruleNode1.setConfiguration(mapper.readTree("\"key1\": \"val1\"")); + + RuleNode ruleNode2 = new RuleNode(); + ruleNode2.setName("name2"); + ruleNode2.setType("type2"); + ruleNode2.setConfiguration(mapper.readTree("\"key2\": \"val2\"")); + + RuleNode ruleNode3 = new RuleNode(); + ruleNode3.setName("name3"); + ruleNode3.setType("type3"); + ruleNode3.setConfiguration(mapper.readTree("\"key3\": \"val3\"")); + + List ruleNodes = new ArrayList<>(); + ruleNodes.add(ruleNode1); + ruleNodes.add(ruleNode2); + ruleNodes.add(ruleNode3); + ruleChainMetaData.setFirstNodeIndex(0); + ruleChainMetaData.setNodes(ruleNodes); + + ruleChainMetaData.addConnectionInfo(0,1,"success"); + ruleChainMetaData.addConnectionInfo(0,2,"fail"); + ruleChainMetaData.addConnectionInfo(1,2,"success"); + + ruleChainMetaData.addRuleChainConnectionInfo(2, edge.getRootRuleChainId(), "success", mapper.createObjectNode()); + + doPost("/api/ruleChain/metadata", ruleChainMetaData, RuleChainMetaData.class); + } + private void testDashboards() throws Exception { log.info("Testing Dashboards"); Dashboard dashboard = new Dashboard(); @@ -373,6 +518,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertEquals(dashboardUpdateMsg.getIdLSB(), savedDashboard.getUuidId().getLeastSignificantBits()); Assert.assertEquals(dashboardUpdateMsg.getTitle(), savedDashboard.getName()); + testAutoGeneratedCodeByProtobuf(dashboardUpdateMsg); + edgeImitator.expectMessageAmount(1); savedDashboard.setTitle("Updated Edge Test Dashboard"); doPost("/api/dashboard", savedDashboard, Dashboard.class); @@ -413,18 +560,10 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { private void testRelations() throws Exception { log.info("Testing Relations"); - List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", - new TypeReference>() {}, new TextPageLink(100)).getData(); - List edgeAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", - new TypeReference>() {}, new TextPageLink(100)).getData(); - - Assert.assertEquals(1, edgeDevices.size()); - Assert.assertEquals(1, edgeAssets.size()); - Device device = edgeDevices.get(0); - Asset asset = edgeAssets.get(0); - Assert.assertEquals("Edge Device 1", device.getName()); - Assert.assertEquals("Edge Asset 1", asset.getName()); + Device device = findDeviceByName("Edge Device 1"); + Asset asset = findAssetByName("Edge Asset 1"); + EntityRelation relation = new EntityRelation(); relation.setType("test"); relation.setFrom(device.getId()); @@ -476,11 +615,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { private void testAlarms() throws Exception { log.info("Testing Alarms"); - 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()); + Device device = findDeviceByName("Edge Device 1"); Alarm alarm = new Alarm(); alarm.setOriginator(device.getId()); @@ -535,11 +670,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { private void testEntityView() throws Exception { log.info("Testing EntityView"); - 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()); + Device device = findDeviceByName("Edge Device 1"); EntityView entityView = new EntityView(); entityView.setName("Edge EntityView 1"); @@ -611,6 +742,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertEquals(customerUpdateMsg.getIdLSB(), savedCustomer.getUuidId().getLeastSignificantBits()); Assert.assertEquals(customerUpdateMsg.getTitle(), savedCustomer.getTitle()); + testAutoGeneratedCodeByProtobuf(customerUpdateMsg); + edgeImitator.expectMessageAmount(1); doDelete("/api/customer/edge/" + edge.getId().getId().toString(), Edge.class); edgeImitator.waitForMessages(); @@ -656,6 +789,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertEquals(widgetsBundleUpdateMsg.getAlias(), savedWidgetsBundle.getAlias()); Assert.assertEquals(widgetsBundleUpdateMsg.getTitle(), savedWidgetsBundle.getTitle()); + testAutoGeneratedCodeByProtobuf(widgetsBundleUpdateMsg); + WidgetType widgetType = new WidgetType(); widgetType.setName("Test Widget Type"); widgetType.setBundleAlias(savedWidgetsBundle.getAlias()); @@ -706,11 +841,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { private void testTimeseries() throws Exception { log.info("Testing timeseries"); - 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()); + Device device = findDeviceByName("Edge Device 1"); String timeseriesData = "{\"data\":{\"temperature\":25},\"ts\":" + System.currentTimeMillis() + "}"; JsonNode timeseriesEntityData = mapper.readTree(timeseriesData); @@ -740,61 +871,92 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { 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()); + Device device = findDeviceByName("Edge Device 1"); - String attributesData = "{\"scope\":\"SERVER_SCOPE\",\"kv\":{\"key\":\"value\"}}"; - JsonNode attributesEntityData = mapper.readTree(attributesData); - EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); + testAttributesUpdatedMsg(device); + testPostAttributesMsg(device); + testAttributesDeleteMsg(device); + + log.info("Attributes tested successfully"); + } + + private void testAttributesDeleteMsg(Device device) throws JsonProcessingException, InterruptedException { + String deleteAttributesData = "{\"scope\":\"SERVER_SCOPE\",\"keys\":[\"key1\",\"key2\"]}"; + JsonNode deleteAttributesEntityData = mapper.readTree(deleteAttributesData); + EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_DELETED, device.getId().getId(), EdgeEventType.DEVICE, deleteAttributesEntityData); edgeImitator.expectMessageAmount(1); - edgeEventService.saveAsync(edgeEvent1); + edgeEventService.saveAsync(edgeEvent); 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()); + Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB()); + Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB()); + Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); - 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()); + Assert.assertTrue(latestEntityDataMsg.hasAttributeDeleteMsg()); - ((ObjectNode) attributesEntityData).put("isPostAttributes", true); - EdgeEvent edgeEvent2 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); + AttributeDeleteMsg attributeDeleteMsg = latestEntityDataMsg.getAttributeDeleteMsg(); + Assert.assertEquals(attributeDeleteMsg.getScope(), deleteAttributesEntityData.get("scope").asText()); + + Assert.assertEquals(2, attributeDeleteMsg.getAttributeNamesCount()); + Assert.assertEquals("key1", attributeDeleteMsg.getAttributeNames(0)); + Assert.assertEquals("key2", attributeDeleteMsg.getAttributeNames(1)); + } + + private void testPostAttributesMsg(Device device) throws JsonProcessingException, InterruptedException { + String postAttributesData = "{\"scope\":\"SERVER_SCOPE\",\"kv\":{\"key2\":\"value2\"}}"; + JsonNode postAttributesEntityData = mapper.readTree(postAttributesData); + EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.POST_ATTRIBUTES, device.getId().getId(), EdgeEventType.DEVICE, postAttributesEntityData); edgeImitator.expectMessageAmount(1); - edgeEventService.saveAsync(edgeEvent2); + edgeEventService.saveAsync(edgeEvent); edgeImitator.waitForMessages(); - latestMessage = edgeImitator.getLatestMessage(); + AbstractMessage 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()); + EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage; + Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB()); + Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB()); + Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); + Assert.assertEquals("SERVER_SCOPE", latestEntityDataMsg.getPostAttributeScope()); Assert.assertTrue(latestEntityDataMsg.hasPostAttributesMsg()); - attributesUpdatedMsg = latestEntityDataMsg.getPostAttributesMsg(); - Assert.assertEquals(1, attributesUpdatedMsg.getKvCount()); - keyValueProto = attributesUpdatedMsg.getKv(0); - Assert.assertEquals("key", keyValueProto.getKey()); - Assert.assertEquals("value", keyValueProto.getStringV()); + TransportProtos.PostAttributeMsg postAttributesMsg = latestEntityDataMsg.getPostAttributesMsg(); + Assert.assertEquals(1, postAttributesMsg.getKvCount()); + TransportProtos.KeyValueProto keyValueProto = postAttributesMsg.getKv(0); + Assert.assertEquals("key2", keyValueProto.getKey()); + Assert.assertEquals("value2", keyValueProto.getStringV()); + } - log.info("Attributes tested successfully"); + private void testAttributesUpdatedMsg(Device device) throws JsonProcessingException, InterruptedException { + String attributesData = "{\"scope\":\"SERVER_SCOPE\",\"kv\":{\"key1\":\"value1\"}}"; + JsonNode attributesEntityData = mapper.readTree(attributesData); + EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); + edgeImitator.expectMessageAmount(1); + edgeEventService.saveAsync(edgeEvent1); + edgeImitator.waitForMessages(); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof EntityDataProto); + EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage; + Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB()); + Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB()); + Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); + Assert.assertEquals("SERVER_SCOPE", latestEntityDataMsg.getPostAttributeScope()); + Assert.assertTrue(latestEntityDataMsg.hasAttributesUpdatedMsg()); + + TransportProtos.PostAttributeMsg attributesUpdatedMsg = latestEntityDataMsg.getAttributesUpdatedMsg(); + Assert.assertEquals(1, attributesUpdatedMsg.getKvCount()); + TransportProtos.KeyValueProto keyValueProto = attributesUpdatedMsg.getKv(0); + Assert.assertEquals("key1", keyValueProto.getKey()); + Assert.assertEquals("value1", keyValueProto.getStringV()); } private void testSendMessagesToCloud() throws Exception { log.info("Sending messages to cloud"); sendDevice(); + sendRelationRequest(); sendAlarm(); sendTelemetry(); sendRelation(); @@ -802,22 +964,30 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { sendRuleChainMetadataRequest(); sendUserCredentialsRequest(); sendDeviceCredentialsRequest(); + sendDeviceRpcResponse(); + sendDeviceCredentialsUpdate(); + sendAttributesRequest(); log.info("Messages were sent successfully"); } private void sendDevice() throws Exception { UUID uuid = UUIDs.timeBased(); - UplinkMsg.Builder builder = UplinkMsg.newBuilder(); + UplinkMsg.Builder uplinkMsgBuilder = 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()); + testAutoGeneratedCodeByProtobuf(deviceUpdateMsgBuilder); + uplinkMsgBuilder.addDeviceUpdateMsg(deviceUpdateMsgBuilder.build()); + edgeImitator.expectResponsesAmount(1); - edgeImitator.sendUplinkMsg(builder.build()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); edgeImitator.waitForResponses(); Device device = doGet("/api/device/" + uuid.toString(), Device.class); @@ -825,23 +995,70 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertEquals("Edge Device 2", device.getName()); } + private void sendRelationRequest() throws Exception { + Device device = findDeviceByName("Edge Device 1"); + Asset asset = findAssetByName("Edge Asset 1"); + + EntityRelation relation = new EntityRelation(); + relation.setType("test"); + relation.setFrom(device.getId()); + relation.setTo(asset.getId()); + relation.setTypeGroup(RelationTypeGroup.COMMON); + + edgeImitator.expectMessageAmount(1); + doPost("/api/relation", relation); + edgeImitator.waitForMessages(); + + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); + RelationRequestMsg.Builder relationRequestMsgBuilder = RelationRequestMsg.newBuilder(); + relationRequestMsgBuilder.setEntityIdMSB(device.getId().getId().getMostSignificantBits()); + relationRequestMsgBuilder.setEntityIdLSB(device.getId().getId().getLeastSignificantBits()); + relationRequestMsgBuilder.setEntityType(device.getId().getEntityType().name()); + testAutoGeneratedCodeByProtobuf(relationRequestMsgBuilder); + + uplinkMsgBuilder.addRelationRequestMsg(relationRequestMsgBuilder.build()); + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.expectMessageAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + edgeImitator.waitForResponses(); + 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(relation.getType(), relationUpdateMsg.getType()); + + UUID fromUUID = new UUID(relationUpdateMsg.getFromIdMSB(), relationUpdateMsg.getFromIdLSB()); + EntityId fromEntityId = EntityIdFactory.getByTypeAndUuid(relationUpdateMsg.getFromEntityType(), fromUUID); + Assert.assertEquals(relation.getFrom(), fromEntityId); + + UUID toUUID = new UUID(relationUpdateMsg.getToIdMSB(), relationUpdateMsg.getToIdLSB()); + EntityId toEntityId = EntityIdFactory.getByTypeAndUuid(relationUpdateMsg.getToEntityType(), toUUID); + Assert.assertEquals(relation.getTo(), toEntityId); + + Assert.assertEquals(relation.getTypeGroup().name(), relationUpdateMsg.getTypeGroup()); + } + 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(); + Device device = findDeviceByName("Edge Device 2"); - UplinkMsg.Builder builder = UplinkMsg.newBuilder(); + UplinkMsg.Builder uplinkMsgBuilder = 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()); + testAutoGeneratedCodeByProtobuf(alarmUpdateMgBuilder); + uplinkMsgBuilder.addAlarmUpdateMsg(alarmUpdateMgBuilder.build()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + edgeImitator.expectResponsesAmount(1); - edgeImitator.sendUplinkMsg(builder.build()); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); edgeImitator.waitForResponses(); @@ -867,7 +1084,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertTrue(foundDevice2.isPresent()); Device device2 = foundDevice2.get(); - UplinkMsg.Builder builder = UplinkMsg.newBuilder(); + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); RelationUpdateMsg.Builder relationUpdateMsgBuilder = RelationUpdateMsg.newBuilder(); relationUpdateMsgBuilder.setType("test"); relationUpdateMsgBuilder.setTypeGroup(RelationTypeGroup.COMMON.name()); @@ -878,10 +1095,13 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { relationUpdateMsgBuilder.setFromIdLSB(device2.getId().getId().getLeastSignificantBits()); relationUpdateMsgBuilder.setFromEntityType(device2.getId().getEntityType().name()); relationUpdateMsgBuilder.setAdditionalInfo("{}"); - builder.addRelationUpdateMsg(relationUpdateMsgBuilder.build()); - UplinkMsg msg = builder.build(); + testAutoGeneratedCodeByProtobuf(relationUpdateMsgBuilder); + uplinkMsgBuilder.addRelationUpdateMsg(relationUpdateMsgBuilder.build()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + edgeImitator.expectResponsesAmount(1); - edgeImitator.sendUplinkMsg(msg); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); edgeImitator.waitForResponses(); EntityRelation relation = doGet("/api/relation?" + @@ -907,28 +1127,35 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { String timeseriesKey = "key"; String timeseriesValue = "25"; data.addProperty(timeseriesKey, timeseriesValue); - UplinkMsg.Builder builder1 = UplinkMsg.newBuilder(); + UplinkMsg.Builder uplinkMsgBuilder1 = 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()); + testAutoGeneratedCodeByProtobuf(entityDataBuilder); + uplinkMsgBuilder1.addEntityData(entityDataBuilder.build()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder1.build()); JsonObject attributesData = new JsonObject(); String attributesKey = "test_attr"; String attributesValue = "test_value"; attributesData.addProperty(attributesKey, attributesValue); - UplinkMsg.Builder builder2 = UplinkMsg.newBuilder(); + UplinkMsg.Builder uplinkMsgBuilder2 = 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.setAttributesUpdatedMsg(JsonConverter.convertToAttributesProto(attributesData)); entityDataBuilder2.setPostAttributeScope(DataConstants.SERVER_SCOPE); - builder2.addEntityData(entityDataBuilder2.build()); - edgeImitator.sendUplinkMsg(builder2.build()); + testAutoGeneratedCodeByProtobuf(entityDataBuilder2); + + uplinkMsgBuilder2.addEntityData(entityDataBuilder2.build()); + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder2); + + edgeImitator.sendUplinkMsg(uplinkMsgBuilder2.build()); edgeImitator.waitForResponses(); Thread.sleep(1000); @@ -947,14 +1174,18 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { private void sendRuleChainMetadataRequest() throws Exception { RuleChainId edgeRootRuleChainId = edge.getRootRuleChainId(); - UplinkMsg.Builder builder = UplinkMsg.newBuilder(); + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); RuleChainMetadataRequestMsg.Builder ruleChainMetadataRequestMsgBuilder = RuleChainMetadataRequestMsg.newBuilder(); ruleChainMetadataRequestMsgBuilder.setRuleChainIdMSB(edgeRootRuleChainId.getId().getMostSignificantBits()); ruleChainMetadataRequestMsgBuilder.setRuleChainIdLSB(edgeRootRuleChainId.getId().getLeastSignificantBits()); - builder.addRuleChainMetadataRequestMsg(ruleChainMetadataRequestMsgBuilder.build()); + testAutoGeneratedCodeByProtobuf(ruleChainMetadataRequestMsgBuilder); + uplinkMsgBuilder.addRuleChainMetadataRequestMsg(ruleChainMetadataRequestMsgBuilder.build()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + edgeImitator.expectResponsesAmount(1); edgeImitator.expectMessageAmount(1); - edgeImitator.sendUplinkMsg(builder.build()); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); edgeImitator.waitForResponses(); edgeImitator.waitForMessages(); @@ -963,19 +1194,25 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = (RuleChainMetadataUpdateMsg) latestMessage; Assert.assertEquals(ruleChainMetadataUpdateMsg.getRuleChainIdMSB(), edgeRootRuleChainId.getId().getMostSignificantBits()); Assert.assertEquals(ruleChainMetadataUpdateMsg.getRuleChainIdLSB(), edgeRootRuleChainId.getId().getLeastSignificantBits()); + + testAutoGeneratedCodeByProtobuf(ruleChainMetadataUpdateMsg); } private void sendUserCredentialsRequest() throws Exception { UserId userId = edgeImitator.getUserId(); - UplinkMsg.Builder builder = UplinkMsg.newBuilder(); + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); UserCredentialsRequestMsg.Builder userCredentialsRequestMsgBuilder = UserCredentialsRequestMsg.newBuilder(); userCredentialsRequestMsgBuilder.setUserIdMSB(userId.getId().getMostSignificantBits()); userCredentialsRequestMsgBuilder.setUserIdLSB(userId.getId().getLeastSignificantBits()); - builder.addUserCredentialsRequestMsg(userCredentialsRequestMsgBuilder.build()); + testAutoGeneratedCodeByProtobuf(userCredentialsRequestMsgBuilder); + uplinkMsgBuilder.addUserCredentialsRequestMsg(userCredentialsRequestMsgBuilder.build()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + edgeImitator.expectResponsesAmount(1); edgeImitator.expectMessageAmount(1); - edgeImitator.sendUplinkMsg(builder.build()); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); edgeImitator.waitForResponses(); edgeImitator.waitForMessages(); @@ -984,26 +1221,27 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { UserCredentialsUpdateMsg userCredentialsUpdateMsg = (UserCredentialsUpdateMsg) latestMessage; Assert.assertEquals(userCredentialsUpdateMsg.getUserIdMSB(), userId.getId().getMostSignificantBits()); Assert.assertEquals(userCredentialsUpdateMsg.getUserIdLSB(), userId.getId().getLeastSignificantBits()); + + testAutoGeneratedCodeByProtobuf(userCredentialsUpdateMsg); } private void sendDeviceCredentialsRequest() 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 1")).findAny(); - Assert.assertTrue(foundDevice.isPresent()); - Device device = foundDevice.get(); + Device device = findDeviceByName("Edge Device 1"); DeviceCredentials deviceCredentials = doGet("/api/device/" + device.getId().getId().toString() + "/credentials", DeviceCredentials.class); - UplinkMsg.Builder builder = UplinkMsg.newBuilder(); + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); DeviceCredentialsRequestMsg.Builder deviceCredentialsRequestMsgBuilder = DeviceCredentialsRequestMsg.newBuilder(); deviceCredentialsRequestMsgBuilder.setDeviceIdMSB(device.getUuidId().getMostSignificantBits()); deviceCredentialsRequestMsgBuilder.setDeviceIdLSB(device.getUuidId().getLeastSignificantBits()); - builder.addDeviceCredentialsRequestMsg(deviceCredentialsRequestMsgBuilder.build()); + testAutoGeneratedCodeByProtobuf(deviceCredentialsRequestMsgBuilder); + uplinkMsgBuilder.addDeviceCredentialsRequestMsg(deviceCredentialsRequestMsgBuilder.build()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); edgeImitator.expectResponsesAmount(1); edgeImitator.expectMessageAmount(1); - edgeImitator.sendUplinkMsg(builder.build()); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); edgeImitator.waitForResponses(); edgeImitator.waitForMessages(); @@ -1016,20 +1254,108 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertEquals(deviceCredentialsUpdateMsg.getCredentialsId(), deviceCredentials.getCredentialsId()); } + private void sendDeviceCredentialsUpdate() throws Exception { + Device device = findDeviceByName("Edge Device 1"); + + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); + DeviceCredentialsUpdateMsg.Builder deviceCredentialsUpdateMsgBuilder = DeviceCredentialsUpdateMsg.newBuilder(); + deviceCredentialsUpdateMsgBuilder.setDeviceIdMSB(device.getUuidId().getMostSignificantBits()); + deviceCredentialsUpdateMsgBuilder.setDeviceIdLSB(device.getUuidId().getLeastSignificantBits()); + deviceCredentialsUpdateMsgBuilder.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN.name()); + deviceCredentialsUpdateMsgBuilder.setCredentialsId("NEW_TOKEN"); + testAutoGeneratedCodeByProtobuf(deviceCredentialsUpdateMsgBuilder); + uplinkMsgBuilder.addDeviceCredentialsUpdateMsg(deviceCredentialsUpdateMsgBuilder.build()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + edgeImitator.waitForResponses(); + } + + private void sendDeviceRpcResponse() throws Exception { + Device device = findDeviceByName("Edge Device 1"); + + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); + DeviceRpcCallMsg.Builder deviceRpcCallResponseBuilder = DeviceRpcCallMsg.newBuilder(); + deviceRpcCallResponseBuilder.setDeviceIdMSB(device.getUuidId().getMostSignificantBits()); + deviceRpcCallResponseBuilder.setDeviceIdLSB(device.getUuidId().getLeastSignificantBits()); + deviceRpcCallResponseBuilder.setOneway(true); + deviceRpcCallResponseBuilder.setOriginServiceId("originServiceId"); + deviceRpcCallResponseBuilder.setExpirationTime(System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10)); + RpcResponseMsg.Builder responseBuilder = + RpcResponseMsg.newBuilder().setResponse("{}"); + testAutoGeneratedCodeByProtobuf(responseBuilder); + + deviceRpcCallResponseBuilder.setResponseMsg(responseBuilder.build()); + testAutoGeneratedCodeByProtobuf(deviceRpcCallResponseBuilder); + + uplinkMsgBuilder.addDeviceRpcCallMsg(deviceRpcCallResponseBuilder.build()); + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + edgeImitator.waitForResponses(); + } + + private void sendAttributesRequest() throws Exception { + Device device = findDeviceByName("Edge Device 1"); + + String attributesDataStr = "{\"key1\":\"value1\"}"; + JsonNode attributesData = mapper.readTree(attributesDataStr); + + doPost("/api/plugins/telemetry/DEVICE/" + device.getId().getId().toString() + "/attributes/" + DataConstants.SERVER_SCOPE, + attributesData); + + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); + AttributesRequestMsg.Builder attributesRequestMsgBuilder = AttributesRequestMsg.newBuilder(); + attributesRequestMsgBuilder.setEntityIdMSB(device.getUuidId().getMostSignificantBits()); + attributesRequestMsgBuilder.setEntityIdLSB(device.getUuidId().getLeastSignificantBits()); + attributesRequestMsgBuilder.setEntityType(EntityType.DEVICE.name()); + testAutoGeneratedCodeByProtobuf(attributesRequestMsgBuilder); + uplinkMsgBuilder.addAttributesRequestMsg(attributesRequestMsgBuilder.build()); + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.expectMessageAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + edgeImitator.waitForResponses(); + edgeImitator.waitForMessages(); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof EntityDataProto); + EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage; + Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB()); + Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB()); + Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType()); + Assert.assertEquals("SERVER_SCOPE", latestEntityDataMsg.getPostAttributeScope()); + Assert.assertTrue(latestEntityDataMsg.hasAttributesUpdatedMsg()); + + TransportProtos.PostAttributeMsg attributesUpdatedMsg = latestEntityDataMsg.getAttributesUpdatedMsg(); + Assert.assertEquals(1, attributesUpdatedMsg.getKvCount()); + TransportProtos.KeyValueProto keyValueProto = attributesUpdatedMsg.getKv(0); + Assert.assertEquals("key1", keyValueProto.getKey()); + Assert.assertEquals("value1", keyValueProto.getStringV()); + } + 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(); + UplinkMsg.Builder upLinkMsgBuilder = 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()); + testAutoGeneratedCodeByProtobuf(deviceDeleteMsgBuilder); + + upLinkMsgBuilder.addDeviceUpdateMsg(deviceDeleteMsgBuilder.build()); + testAutoGeneratedCodeByProtobuf(upLinkMsgBuilder); + edgeImitator.expectResponsesAmount(1); - edgeImitator.sendUplinkMsg(builder.build()); + edgeImitator.sendUplinkMsg(upLinkMsgBuilder.build()); edgeImitator.waitForResponses(); device = doGet("/api/device/" + device.getId().getId().toString(), Device.class); Assert.assertNotNull(device); @@ -1042,41 +1368,15 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { private void installation() throws Exception { edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class); - Device device = new Device(); - device.setName("Edge Device 1"); - device.setType("test"); - Device savedDevice = doPost("/api/device", device, Device.class); + Device savedDevice = saveDevice("Edge Device 1"); doPost("/api/edge/" + edge.getId().getId().toString() + "/device/" + savedDevice.getId().getId().toString(), Device.class); - Asset asset = new Asset(); - asset.setName("Edge Asset 1"); - asset.setType("test"); - Asset savedAsset = doPost("/api/asset", asset, Asset.class); + Asset savedAsset = saveAsset("Edge Asset 1"); doPost("/api/edge/" + edge.getId().getId().toString() + "/asset/" + savedAsset.getId().getId().toString(), Asset.class); } - private void uninstallation() throws Exception { - - TimePageData pageDataDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", - new TypeReference>() {}, new TextPageLink(100)); - for (Device device: pageDataDevices.getData()) { - doDelete("/api/device/" + device.getId().getId().toString()) - .andExpect(status().isOk()); - } - - TimePageData pageDataAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", - new TypeReference>() {}, new TextPageLink(100)); - for (Asset asset: pageDataAssets.getData()) { - doDelete("/api/asset/" + asset.getId().getId().toString()) - .andExpect(status().isOk()); - } - - doDelete("/api/edge/" + edge.getId().getId().toString()) - .andExpect(status().isOk()); - } - private EdgeEvent constructEdgeEvent(TenantId tenantId, EdgeId edgeId, EdgeEventActionType edgeEventAction, UUID entityId, EdgeEventType edgeEventType, JsonNode entityBody) { EdgeEvent edgeEvent = new EdgeEvent(); edgeEvent.setEdgeId(edgeId); @@ -1087,4 +1387,19 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { edgeEvent.setBody(entityBody); return edgeEvent; } + + private void testAutoGeneratedCodeByProtobuf(MessageLite.Builder builder) throws InvalidProtocolBufferException { + MessageLite source = builder.build(); + + testAutoGeneratedCodeByProtobuf(source); + + MessageLite target = source.getParserForType().parseFrom(source.toByteArray()); + builder.clear().mergeFrom(target); + } + + private void testAutoGeneratedCodeByProtobuf(MessageLite source) throws InvalidProtocolBufferException { + MessageLite target = source.getParserForType().parseFrom(source.toByteArray()); + Assert.assertEquals(source, target); + Assert.assertEquals(source.hashCode(), target.hashCode()); + } } 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 3eca5438dd..c2122a9c83 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 @@ -30,7 +30,9 @@ import org.thingsboard.server.gen.edge.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.AssetUpdateMsg; import org.thingsboard.server.gen.edge.CustomerUpdateMsg; import org.thingsboard.server.gen.edge.DashboardUpdateMsg; +import org.thingsboard.server.gen.edge.DeviceCredentialsRequestMsg; import org.thingsboard.server.gen.edge.DeviceCredentialsUpdateMsg; +import org.thingsboard.server.gen.edge.DeviceRpcCallMsg; import org.thingsboard.server.gen.edge.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.DownlinkMsg; import org.thingsboard.server.gen.edge.DownlinkResponseMsg; @@ -224,6 +226,16 @@ public class EdgeImitator { result.add(saveDownlinkMsg(userCredentialsUpdateMsg)); } } + if (downlinkMsg.getDeviceRpcCallMsgList() != null && !downlinkMsg.getDeviceRpcCallMsgList().isEmpty()) { + for (DeviceRpcCallMsg deviceRpcCallMsg: downlinkMsg.getDeviceRpcCallMsgList()) { + result.add(saveDownlinkMsg(deviceRpcCallMsg)); + } + } + if (downlinkMsg.getDeviceCredentialsRequestMsgList() != null && !downlinkMsg.getDeviceCredentialsRequestMsgList().isEmpty()) { + for (DeviceCredentialsRequestMsg deviceCredentialsRequestMsg: downlinkMsg.getDeviceCredentialsRequestMsgList()) { + result.add(saveDownlinkMsg(deviceCredentialsRequestMsg)); + } + } return Futures.allAsList(result); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java index 2422aec651..b8b386ac1a 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java @@ -19,6 +19,7 @@ public enum EdgeEventActionType { ADDED, DELETED, UPDATED, + POST_ATTRIBUTES, ATTRIBUTES_UPDATED, ATTRIBUTES_DELETED, TIMESERIES_UPDATED, diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java index 5bd33a056a..6072e0cae9 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java @@ -152,11 +152,9 @@ public class TbMsgPushToEdgeNode implements TbNode { JsonNode dataJson = json.readTree(msg.getData()); switch (actionType) { case ATTRIBUTES_UPDATED: + case POST_ATTRIBUTES: entityBody.put("kv", dataJson); entityBody.put("scope", metadata.get("scope")); - if (SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType)) { - entityBody.put("isPostAttributes", true); - } break; case ATTRIBUTES_DELETED: List keys = json.treeToValue(dataJson.get("attributes"), List.class); @@ -192,9 +190,10 @@ public class TbMsgPushToEdgeNode implements TbNode { EdgeEventActionType actionType; if (SessionMsgType.POST_TELEMETRY_REQUEST.name().equals(msgType)) { actionType = EdgeEventActionType.TIMESERIES_UPDATED; - } else if (SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType) - || DataConstants.ATTRIBUTES_UPDATED.equals(msgType)) { + } else if (DataConstants.ATTRIBUTES_UPDATED.equals(msgType)) { actionType = EdgeEventActionType.ATTRIBUTES_UPDATED; + } else if (SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType)) { + actionType = EdgeEventActionType.POST_ATTRIBUTES; } else { actionType = EdgeEventActionType.ATTRIBUTES_DELETED; } diff --git a/ui/src/app/event/event-table.directive.js b/ui/src/app/event/event-table.directive.js index 93a425a200..762d4f308d 100644 --- a/ui/src/app/event/event-table.directive.js +++ b/ui/src/app/event/event-table.directive.js @@ -220,39 +220,37 @@ export default function EventTableDirective($compile, $templateCache, $rootScope } scope.subscriptionId = null; - scope.queueStartTs; + scope.queueStartTs = 0; scope.loadEdgeInfo = function() { - attributeService.getEntityAttributesValues(scope.entityType, scope.entityId, types.attributesScope.server.value, - types.edgeAttributeKeys.queueStartTs, {}) + attributeService.getEntityAttributesValues( + scope.entityType, + scope.entityId, + types.attributesScope.server.value, + types.edgeAttributeKeys.queueStartTs, + {}) .then(function success(attributes) { - scope.onUpdate(attributes); + scope.onEdgeAttributesUpdate(attributes); }); scope.checkSubscription(); - - attributeService.getEntityAttributes(scope.entityType, scope.entityId, types.attributesScope.server.value, {order: '', limit: 1, page: 1, search: ''}, - function (attributes) { - if (attributes && attributes.data) { - scope.onUpdate(attributes.data); - } - }); } - scope.onUpdate = function(attributes) { - let edge = attributes.reduce(function (map, attribute) { + scope.onEdgeAttributesUpdate = function(attributes) { + let edgeAttributes = attributes.reduce(function (map, attribute) { map[attribute.key] = attribute; return map; }, {}); - if (edge.queueStartTs) { - scope.queueStartTs = edge.queueStartTs.lastUpdateTs; + if (edgeAttributes.queueStartTs) { + scope.queueStartTs = edgeAttributes.queueStartTs.lastUpdateTs; } } scope.checkSubscription = function() { var newSubscriptionId = null; if (scope.entityId && scope.entityType && types.attributesScope.server.value) { - newSubscriptionId = attributeService.subscribeForEntityAttributes(scope.entityType, scope.entityId, types.attributesScope.server.value); + newSubscriptionId = + attributeService.subscribeForEntityAttributes(scope.entityType, scope.entityId, types.attributesScope.server.value); } if (scope.subscriptionId && scope.subscriptionId != newSubscriptionId) { attributeService.unsubscribeForEntityAttributes(scope.subscriptionId);