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 a7a1282dc0..a50c2b4bc1 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java @@ -166,7 +166,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { // sleep 1 seconds to avoid CREDENTIALS updated message for the user // user credentials is going to be stored and updated event pushed to edge notification service // while service will be processing this event edge could be already added and additional message will be pushed - Thread.sleep(1000); + Thread.sleep(500); installation(); @@ -177,6 +177,22 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { testReceivedInitialData(); } + private void installation() throws Exception { + edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class); + + DeviceProfile deviceProfile = this.createDeviceProfile(CUSTOM_DEVICE_PROFILE_NAME, null); + extendDeviceProfileData(deviceProfile); + doPost("/api/deviceProfile", deviceProfile, DeviceProfile.class); + + Device savedDevice = saveDevice("Edge Device 1", CUSTOM_DEVICE_PROFILE_NAME); + doPost("/api/edge/" + edge.getId().getId().toString() + + "/device/" + savedDevice.getId().getId().toString(), Device.class); + + Asset savedAsset = saveAsset("Edge Asset 1"); + doPost("/api/edge/" + edge.getId().getId().toString() + + "/asset/" + savedAsset.getId().getId().toString(), Asset.class); + } + @After public void afterTest() throws Exception { try { @@ -190,7 +206,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { } @Test - public void generalTest() throws Exception { + public void test() throws Exception { testDevices(); testAssets(); @@ -216,179 +232,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { testSendMessagesToCloud(); testRpcCall(); - } - - @Test - public void testTimeseriesWithFailures() throws Exception { - log.info("Testing timeseries with failures"); - - int numberOfTimeseriesToSend = 1000; - - edgeImitator.setRandomFailuresOnTimeseriesDownlink(true); - // imitator will generate failure in 5% of cases - edgeImitator.setFailureProbability(5.0); - - edgeImitator.expectMessageAmount(numberOfTimeseriesToSend); - Device device = findDeviceByName("Edge Device 1"); - for (int idx = 1; idx <= numberOfTimeseriesToSend; idx++) { - String timeseriesData = "{\"data\":{\"idx\":" + idx + "},\"ts\":" + System.currentTimeMillis() + "}"; - JsonNode timeseriesEntityData = mapper.readTree(timeseriesData); - EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.TIMESERIES_UPDATED, - device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData); - edgeEventService.saveAsync(edgeEvent); - clusterService.onEdgeEventUpdate(tenantId, edge.getId()); - } - Assert.assertTrue(edgeImitator.waitForMessages(60)); - - List allTelemetryMsgs = edgeImitator.findAllMessagesByType(EntityDataProto.class); - Assert.assertEquals(numberOfTimeseriesToSend, allTelemetryMsgs.size()); - - for (int idx = 1; idx <= numberOfTimeseriesToSend; idx++) { - Assert.assertTrue(isIdxExistsInTheDownlinkList(idx, allTelemetryMsgs)); - } - - edgeImitator.setRandomFailuresOnTimeseriesDownlink(false); - log.info("Timeseries with failures tested successfully"); - } - - private boolean isIdxExistsInTheDownlinkList(int idx, List allTelemetryMsgs) { - for (EntityDataProto proto : allTelemetryMsgs) { - TransportProtos.PostTelemetryMsg postTelemetryMsg = proto.getPostTelemetryMsg(); - Assert.assertEquals(1, postTelemetryMsg.getTsKvListCount()); - TransportProtos.TsKvListProto tsKvListProto = postTelemetryMsg.getTsKvList(0); - Assert.assertEquals(1, tsKvListProto.getKvCount()); - TransportProtos.KeyValueProto keyValueProto = tsKvListProto.getKv(0); - Assert.assertEquals("idx", keyValueProto.getKey()); - if (keyValueProto.getLongV() == idx) { - return true; - } - } - return false; - } - - private Device findDeviceByName(String deviceName) throws Exception { - List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", - new TypeReference>() { - }, new PageLink(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 PageLink(100)).getData(); - - Assert.assertEquals(1, edgeAssets.size()); - Asset asset = edgeAssets.get(0); - Assert.assertEquals(assetName, asset.getName()); - return asset; - } - - private Device saveDevice(String deviceName, String type) throws Exception { - Device device = new Device(); - device.setName(deviceName); - device.setType(type); - 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"); - - ObjectNode body = mapper.createObjectNode(); - body.put("requestId", new Random().nextInt()); - body.put("requestUUID", Uuids.timeBased().toString()); - body.put("oneway", false); - body.put("expirationTime", System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10)); - body.put("method", "test_method"); - body.put("params", "{\"param1\":\"value1\"}"); - - EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body); - edgeImitator.expectMessageAmount(1); - edgeEventService.saveAsync(edgeEvent); - clusterService.onEdgeEventUpdate(tenantId, edge.getId()); - Assert.assertTrue(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 { - log.info("Checking received data"); - Assert.assertTrue(edgeImitator.waitForMessages()); - - EdgeConfiguration configuration = edgeImitator.getConfiguration(); - Assert.assertNotNull(configuration); - - testAutoGeneratedCodeByProtobuf(configuration); - - Optional deviceUpdateMsgOpt = edgeImitator.findMessageByType(DeviceUpdateMsg.class); - Assert.assertTrue(deviceUpdateMsgOpt.isPresent()); - DeviceUpdateMsg deviceUpdateMsg = deviceUpdateMsgOpt.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 PageLink(100)).getData(); - Assert.assertTrue(edgeDevices.contains(device)); - - List deviceProfileUpdateMsgList = edgeImitator.findAllMessagesByType(DeviceProfileUpdateMsg.class); - Assert.assertEquals(3, deviceProfileUpdateMsgList.size()); - Optional deviceProfileUpdateMsgOpt = - deviceProfileUpdateMsgList.stream().filter(dfum -> CUSTOM_DEVICE_PROFILE_NAME.equals(dfum.getName())).findAny(); - Assert.assertTrue(deviceProfileUpdateMsgOpt.isPresent()); - DeviceProfileUpdateMsg deviceProfileUpdateMsg = deviceProfileUpdateMsgOpt.get(); - Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); - UUID deviceProfileUUID = new UUID(deviceProfileUpdateMsg.getIdMSB(), deviceProfileUpdateMsg.getIdLSB()); - DeviceProfile deviceProfile = doGet("/api/deviceProfile/" + deviceProfileUUID.toString(), DeviceProfile.class); - Assert.assertNotNull(deviceProfile); - Assert.assertNotNull(deviceProfile.getProfileData()); - Assert.assertNotNull(deviceProfile.getProfileData().getAlarms()); - Assert.assertNotNull(deviceProfile.getProfileData().getAlarms().get(0).getClearRule()); - - testAutoGeneratedCodeByProtobuf(deviceProfileUpdateMsg); - - Optional assetUpdateMsgOpt = edgeImitator.findMessageByType(AssetUpdateMsg.class); - Assert.assertTrue(assetUpdateMsgOpt.isPresent()); - AssetUpdateMsg assetUpdateMsg = assetUpdateMsgOpt.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 PageLink(100)).getData(); - Assert.assertTrue(edgeAssets.contains(asset)); - - testAutoGeneratedCodeByProtobuf(assetUpdateMsg); - - Optional ruleChainUpdateMsgOpt = edgeImitator.findMessageByType(RuleChainUpdateMsg.class); - Assert.assertTrue(ruleChainUpdateMsgOpt.isPresent()); - RuleChainUpdateMsg ruleChainUpdateMsg = ruleChainUpdateMsgOpt.get(); - Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_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 PageLink(100)).getData(); - Assert.assertTrue(edgeRuleChains.contains(ruleChain)); - - testAutoGeneratedCodeByProtobuf(ruleChainUpdateMsg); - - log.info("Received data checked"); + testTimeseriesWithFailures(); } private void testDevices() throws Exception { @@ -529,23 +374,20 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { log.info("Testing RuleChains"); // 1 - edgeImitator.expectMessageAmount(1); + edgeImitator.expectMessageAmount(2); RuleChain ruleChain = new RuleChain(); ruleChain.setName("Edge Test Rule Chain"); ruleChain.setType(RuleChainType.EDGE); RuleChain savedRuleChain = doPost("/api/ruleChain", ruleChain, RuleChain.class); - createRuleChainMetadata(savedRuleChain); - // sleep 1 seconds to avoid ENTITY_UPDATED_RPC_MESSAGE for the rule chain - // rule chain metadata is going to be stored and updated event pushed to edge notification service - // while service will be processing this event assignment rule chain to edge will be completed if bad timing - Thread.sleep(1000); doPost("/api/edge/" + edge.getId().getId().toString() + "/ruleChain/" + savedRuleChain.getId().getId().toString(), RuleChain.class); + createRuleChainMetadata(savedRuleChain); Assert.assertTrue(edgeImitator.waitForMessages()); - AbstractMessage latestMessage = edgeImitator.getLatestMessage(); - Assert.assertTrue(latestMessage instanceof RuleChainUpdateMsg); - RuleChainUpdateMsg ruleChainUpdateMsg = (RuleChainUpdateMsg) latestMessage; - Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, ruleChainUpdateMsg.getMsgType()); + Optional ruleChainUpdateMsgOpt = edgeImitator.findMessageByType(RuleChainUpdateMsg.class); + Assert.assertTrue(ruleChainUpdateMsgOpt.isPresent()); + RuleChainUpdateMsg ruleChainUpdateMsg = ruleChainUpdateMsgOpt.get(); + Assert.assertTrue(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE.equals(ruleChainUpdateMsg.getMsgType()) || + UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE.equals(ruleChainUpdateMsg.getMsgType())); Assert.assertEquals(ruleChainUpdateMsg.getIdMSB(), savedRuleChain.getUuidId().getMostSignificantBits()); Assert.assertEquals(ruleChainUpdateMsg.getIdLSB(), savedRuleChain.getUuidId().getLeastSignificantBits()); Assert.assertEquals(ruleChainUpdateMsg.getName(), savedRuleChain.getName()); @@ -558,9 +400,9 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { doDelete("/api/edge/" + edge.getId().getId().toString() + "/ruleChain/" + savedRuleChain.getId().getId().toString(), RuleChain.class); Assert.assertTrue(edgeImitator.waitForMessages()); - latestMessage = edgeImitator.getLatestMessage(); - Assert.assertTrue(latestMessage instanceof RuleChainUpdateMsg); - ruleChainUpdateMsg = (RuleChainUpdateMsg) latestMessage; + ruleChainUpdateMsgOpt = edgeImitator.findMessageByType(RuleChainUpdateMsg.class); + Assert.assertTrue(ruleChainUpdateMsgOpt.isPresent()); + ruleChainUpdateMsg = ruleChainUpdateMsgOpt.get(); Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, ruleChainUpdateMsg.getMsgType()); Assert.assertEquals(ruleChainUpdateMsg.getIdMSB(), savedRuleChain.getUuidId().getMostSignificantBits()); Assert.assertEquals(ruleChainUpdateMsg.getIdLSB(), savedRuleChain.getUuidId().getLeastSignificantBits()); @@ -574,67 +416,6 @@ 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()); - Assert.assertTrue(edgeImitator.waitForResponses()); - Assert.assertTrue(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"); @@ -1037,12 +818,12 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { 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); + 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(edgeEvent); + edgeEventService.saveAsync(edgeEvent1); clusterService.onEdgeEventUpdate(tenantId, edge.getId()); Assert.assertTrue(edgeImitator.waitForMessages()); @@ -1052,15 +833,14 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { 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()); - Assert.assertTrue(latestEntityDataMsg.hasAttributeDeleteMsg()); - - 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)); + 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 testPostAttributesMsg(Device device) throws JsonProcessingException, InterruptedException { @@ -1088,29 +868,192 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertEquals("value2", keyValueProto.getStringV()); } - 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); + 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); clusterService.onEdgeEventUpdate(tenantId, edge.getId()); Assert.assertTrue(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()); + 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.assertTrue(latestEntityDataMsg.hasAttributeDeleteMsg()); + + 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 Device findDeviceByName(String deviceName) throws Exception { + List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() { + }, new PageLink(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 PageLink(100)).getData(); + + Assert.assertEquals(1, edgeAssets.size()); + Asset asset = edgeAssets.get(0); + Assert.assertEquals(assetName, asset.getName()); + return asset; + } + + private Device saveDevice(String deviceName, String type) throws Exception { + Device device = new Device(); + device.setName(deviceName); + device.setType(type); + 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 testReceivedInitialData() throws Exception { + log.info("Checking received data"); + Assert.assertTrue(edgeImitator.waitForMessages()); + + EdgeConfiguration configuration = edgeImitator.getConfiguration(); + Assert.assertNotNull(configuration); + + testAutoGeneratedCodeByProtobuf(configuration); + + Optional deviceUpdateMsgOpt = edgeImitator.findMessageByType(DeviceUpdateMsg.class); + Assert.assertTrue(deviceUpdateMsgOpt.isPresent()); + DeviceUpdateMsg deviceUpdateMsg = deviceUpdateMsgOpt.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 PageLink(100)).getData(); + Assert.assertTrue(edgeDevices.contains(device)); + + List deviceProfileUpdateMsgList = edgeImitator.findAllMessagesByType(DeviceProfileUpdateMsg.class); + Assert.assertEquals(3, deviceProfileUpdateMsgList.size()); + Optional deviceProfileUpdateMsgOpt = + deviceProfileUpdateMsgList.stream().filter(dfum -> CUSTOM_DEVICE_PROFILE_NAME.equals(dfum.getName())).findAny(); + Assert.assertTrue(deviceProfileUpdateMsgOpt.isPresent()); + DeviceProfileUpdateMsg deviceProfileUpdateMsg = deviceProfileUpdateMsgOpt.get(); + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceProfileUpdateMsg.getMsgType()); + UUID deviceProfileUUID = new UUID(deviceProfileUpdateMsg.getIdMSB(), deviceProfileUpdateMsg.getIdLSB()); + DeviceProfile deviceProfile = doGet("/api/deviceProfile/" + deviceProfileUUID.toString(), DeviceProfile.class); + Assert.assertNotNull(deviceProfile); + Assert.assertNotNull(deviceProfile.getProfileData()); + Assert.assertNotNull(deviceProfile.getProfileData().getAlarms()); + Assert.assertNotNull(deviceProfile.getProfileData().getAlarms().get(0).getClearRule()); + + testAutoGeneratedCodeByProtobuf(deviceProfileUpdateMsg); + + Optional assetUpdateMsgOpt = edgeImitator.findMessageByType(AssetUpdateMsg.class); + Assert.assertTrue(assetUpdateMsgOpt.isPresent()); + AssetUpdateMsg assetUpdateMsg = assetUpdateMsgOpt.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 PageLink(100)).getData(); + Assert.assertTrue(edgeAssets.contains(asset)); + + testAutoGeneratedCodeByProtobuf(assetUpdateMsg); + + Optional ruleChainUpdateMsgOpt = edgeImitator.findMessageByType(RuleChainUpdateMsg.class); + Assert.assertTrue(ruleChainUpdateMsgOpt.isPresent()); + RuleChainUpdateMsg ruleChainUpdateMsg = ruleChainUpdateMsgOpt.get(); + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_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 PageLink(100)).getData(); + Assert.assertTrue(edgeRuleChains.contains(ruleChain)); + + testAutoGeneratedCodeByProtobuf(ruleChainUpdateMsg); + + log.info("Received data checked"); + } + + 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()); + Assert.assertTrue(edgeImitator.waitForResponses()); + Assert.assertTrue(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"); - 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()); + ruleChainMetaData.addRuleChainConnectionInfo(2, edge.getRootRuleChainId(), "success", mapper.createObjectNode()); + + doPost("/api/ruleChain/metadata", ruleChainMetaData, RuleChainMetaData.class); } private void testSendMessagesToCloud() throws Exception { @@ -1284,46 +1227,6 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertEquals(AlarmSeverity.CRITICAL, alarmInfo.getSeverity()); } - private void sendRelation() throws Exception { - List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", - new TypeReference>() {}, new PageLink(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 uplinkMsgBuilder = 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("{}"); - testAutoGeneratedCodeByProtobuf(relationUpdateMsgBuilder); - uplinkMsgBuilder.addRelationUpdateMsg(relationUpdateMsgBuilder.build()); - - testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); - - edgeImitator.expectResponsesAmount(1); - edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); - Assert.assertTrue(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 PageLink(100)).getData(); @@ -1368,9 +1271,15 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { edgeImitator.sendUplinkMsg(uplinkMsgBuilder2.build()); Assert.assertTrue(edgeImitator.waitForResponses()); - // Wait before device attributes saved to database before requesting them from controller - Thread.sleep(1000); - Map>> timeseries = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/values/timeseries?keys=" + timeseriesKey, new TypeReference<>() {}); + int attempt = 0; + Map>> timeseries; + do { + timeseries = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/values/timeseries?keys=" + timeseriesKey, + new TypeReference<>() {}); + // Wait before device attributes saved to database before requesting them from controller + Thread.sleep(100); + attempt++; + } while (!timeseries.containsKey(timeseriesKey) || attempt < 10); Assert.assertTrue(timeseries.containsKey(timeseriesKey)); Assert.assertEquals(1, timeseries.get(timeseriesKey).size()); Assert.assertEquals(timeseriesValue, timeseries.get(timeseriesKey).get(0).get("value")); @@ -1382,6 +1291,69 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { } + private void sendRelation() throws Exception { + List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() {}, new PageLink(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 uplinkMsgBuilder = 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("{}"); + testAutoGeneratedCodeByProtobuf(relationUpdateMsgBuilder); + uplinkMsgBuilder.addRelationUpdateMsg(relationUpdateMsgBuilder.build()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + Assert.assertTrue(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 sendDeleteDeviceOnEdge() throws Exception { + Device device = findDeviceByName("Edge Device 2"); + 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()); + testAutoGeneratedCodeByProtobuf(deviceDeleteMsgBuilder); + + upLinkMsgBuilder.addDeviceUpdateMsg(deviceDeleteMsgBuilder.build()); + testAutoGeneratedCodeByProtobuf(upLinkMsgBuilder); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(upLinkMsgBuilder.build()); + Assert.assertTrue(edgeImitator.waitForResponses()); + device = doGet("/api/device/" + device.getId().getId().toString(), Device.class); + Assert.assertNotNull(device); + List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", + new TypeReference>() { + }, new PageLink(100)).getData(); + Assert.assertFalse(edgeDevices.contains(device)); + } + private void sendRuleChainMetadataRequest() throws Exception { RuleChainId edgeRootRuleChainId = edge.getRootRuleChainId(); @@ -1463,25 +1435,6 @@ 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()); - Assert.assertTrue(edgeImitator.waitForResponses()); - } - private void sendDeviceRpcResponse() throws Exception { Device device = findDeviceByName("Edge Device 1"); @@ -1507,6 +1460,25 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertTrue(edgeImitator.waitForResponses()); } + 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()); + Assert.assertTrue(edgeImitator.waitForResponses()); + } + private void sendAttributesRequest() throws Exception { Device device = findDeviceByName("Edge Device 1"); sendAttributesRequest(device, DataConstants.SERVER_SCOPE, "{\"key1\":\"value1\"}", "key1", "value1"); @@ -1520,7 +1492,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { attributesData); // Wait before device attributes saved to database before requesting them from edge - Thread.sleep(1000); + // queue used to save attributes to database + Thread.sleep(500); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); AttributesRequestMsg.Builder attributesRequestMsgBuilder = AttributesRequestMsg.newBuilder(); @@ -1554,43 +1527,75 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertEquals(expectedValue, keyValueProto.getStringV()); } - private void sendDeleteDeviceOnEdge() throws Exception { - Device device = findDeviceByName("Edge Device 2"); - 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()); - testAutoGeneratedCodeByProtobuf(deviceDeleteMsgBuilder); + private void testRpcCall() throws Exception { + Device device = findDeviceByName("Edge Device 1"); - upLinkMsgBuilder.addDeviceUpdateMsg(deviceDeleteMsgBuilder.build()); - testAutoGeneratedCodeByProtobuf(upLinkMsgBuilder); + ObjectNode body = mapper.createObjectNode(); + body.put("requestId", new Random().nextInt()); + body.put("requestUUID", Uuids.timeBased().toString()); + body.put("oneway", false); + body.put("expirationTime", System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10)); + body.put("method", "test_method"); + body.put("params", "{\"param1\":\"value1\"}"); - edgeImitator.expectResponsesAmount(1); - edgeImitator.sendUplinkMsg(upLinkMsgBuilder.build()); - Assert.assertTrue(edgeImitator.waitForResponses()); - device = doGet("/api/device/" + device.getId().getId().toString(), Device.class); - Assert.assertNotNull(device); - List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", - new TypeReference>() { - }, new PageLink(100)).getData(); - Assert.assertFalse(edgeDevices.contains(device)); + EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body); + edgeImitator.expectMessageAmount(1); + edgeEventService.saveAsync(edgeEvent); + clusterService.onEdgeEventUpdate(tenantId, edge.getId()); + Assert.assertTrue(edgeImitator.waitForMessages()); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof DeviceRpcCallMsg); + DeviceRpcCallMsg latestDeviceRpcCallMsg = (DeviceRpcCallMsg) latestMessage; + Assert.assertEquals("test_method", latestDeviceRpcCallMsg.getRequestMsg().getMethod()); } - private void installation() throws Exception { - edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class); + private void testTimeseriesWithFailures() throws Exception { + log.info("Testing timeseries with failures"); - DeviceProfile deviceProfile = this.createDeviceProfile(CUSTOM_DEVICE_PROFILE_NAME, null); - extendDeviceProfileData(deviceProfile); - doPost("/api/deviceProfile", deviceProfile, DeviceProfile.class); + int numberOfTimeseriesToSend = 1000; - Device savedDevice = saveDevice("Edge Device 1", CUSTOM_DEVICE_PROFILE_NAME); - doPost("/api/edge/" + edge.getId().getId().toString() - + "/device/" + savedDevice.getId().getId().toString(), Device.class); + edgeImitator.setRandomFailuresOnTimeseriesDownlink(true); + // imitator will generate failure in 5% of cases + edgeImitator.setFailureProbability(5.0); - Asset savedAsset = saveAsset("Edge Asset 1"); - doPost("/api/edge/" + edge.getId().getId().toString() - + "/asset/" + savedAsset.getId().getId().toString(), Asset.class); + edgeImitator.expectMessageAmount(numberOfTimeseriesToSend); + Device device = findDeviceByName("Edge Device 1"); + for (int idx = 1; idx <= numberOfTimeseriesToSend; idx++) { + String timeseriesData = "{\"data\":{\"idx\":" + idx + "},\"ts\":" + System.currentTimeMillis() + "}"; + JsonNode timeseriesEntityData = mapper.readTree(timeseriesData); + EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.TIMESERIES_UPDATED, + device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData); + edgeEventService.saveAsync(edgeEvent); + clusterService.onEdgeEventUpdate(tenantId, edge.getId()); + } + + Assert.assertTrue(edgeImitator.waitForMessages(60)); + + List allTelemetryMsgs = edgeImitator.findAllMessagesByType(EntityDataProto.class); + Assert.assertEquals(numberOfTimeseriesToSend, allTelemetryMsgs.size()); + + for (int idx = 1; idx <= numberOfTimeseriesToSend; idx++) { + Assert.assertTrue(isIdxExistsInTheDownlinkList(idx, allTelemetryMsgs)); + } + + edgeImitator.setRandomFailuresOnTimeseriesDownlink(false); + log.info("Timeseries with failures tested successfully"); + } + + private boolean isIdxExistsInTheDownlinkList(int idx, List allTelemetryMsgs) { + for (EntityDataProto proto : allTelemetryMsgs) { + TransportProtos.PostTelemetryMsg postTelemetryMsg = proto.getPostTelemetryMsg(); + Assert.assertEquals(1, postTelemetryMsg.getTsKvListCount()); + TransportProtos.TsKvListProto tsKvListProto = postTelemetryMsg.getTsKvList(0); + Assert.assertEquals(1, tsKvListProto.getKvCount()); + TransportProtos.KeyValueProto keyValueProto = tsKvListProto.getKv(0); + Assert.assertEquals("idx", keyValueProto.getKey()); + if (keyValueProto.getLongV() == idx) { + return true; + } + } + return false; } private void extendDeviceProfileData(DeviceProfile deviceProfile) { diff --git a/application/src/test/resources/application-test.properties b/application/src/test/resources/application-test.properties index 02d43723ed..40ef8f06c3 100644 --- a/application/src/test/resources/application-test.properties +++ b/application/src/test/resources/application-test.properties @@ -1,4 +1,6 @@ transport.lwm2m.security.key_store=lwm2m/credentials/serverKeyStore.jks transport.lwm2m.security.key_store_password=server edges.enabled=true +edges.storage.no_read_records_sleep=500 +edges.storage.sleep_between_batches=500 transport.lwm2m.bootstrap.security.alias=server \ No newline at end of file