|
|
|
@ -177,22 +177,6 @@ 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 { |
|
|
|
@ -229,11 +213,121 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { |
|
|
|
|
|
|
|
testAttributes(); |
|
|
|
|
|
|
|
testSendMessagesToCloud(); |
|
|
|
|
|
|
|
testRpcCall(); |
|
|
|
|
|
|
|
testTimeseriesWithFailures(); |
|
|
|
|
|
|
|
testSendMessagesToCloud(); |
|
|
|
} |
|
|
|
|
|
|
|
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); |
|
|
|
} |
|
|
|
|
|
|
|
private void extendDeviceProfileData(DeviceProfile deviceProfile) { |
|
|
|
DeviceProfileData profileData = deviceProfile.getProfileData(); |
|
|
|
List<DeviceProfileAlarm> alarms = new ArrayList<>(); |
|
|
|
DeviceProfileAlarm deviceProfileAlarm = new DeviceProfileAlarm(); |
|
|
|
deviceProfileAlarm.setAlarmType("High Temperature"); |
|
|
|
AlarmRule alarmRule = new AlarmRule(); |
|
|
|
alarmRule.setAlarmDetails("Alarm Details"); |
|
|
|
AlarmCondition alarmCondition = new AlarmCondition(); |
|
|
|
alarmCondition.setSpec(new SimpleAlarmConditionSpec()); |
|
|
|
List<AlarmConditionFilter> condition = new ArrayList<>(); |
|
|
|
AlarmConditionFilter alarmConditionFilter = new AlarmConditionFilter(); |
|
|
|
alarmConditionFilter.setKey(new AlarmConditionFilterKey(AlarmConditionKeyType.ATTRIBUTE, "temperature")); |
|
|
|
NumericFilterPredicate predicate = new NumericFilterPredicate(); |
|
|
|
predicate.setOperation(NumericFilterPredicate.NumericOperation.GREATER); |
|
|
|
predicate.setValue(new FilterPredicateValue<>(55.0)); |
|
|
|
alarmConditionFilter.setPredicate(predicate); |
|
|
|
alarmConditionFilter.setValueType(EntityKeyValueType.NUMERIC); |
|
|
|
condition.add(alarmConditionFilter); |
|
|
|
alarmCondition.setCondition(condition); |
|
|
|
alarmRule.setCondition(alarmCondition); |
|
|
|
deviceProfileAlarm.setClearRule(alarmRule); |
|
|
|
TreeMap<AlarmSeverity, AlarmRule> createRules = new TreeMap<>(); |
|
|
|
createRules.put(AlarmSeverity.CRITICAL, alarmRule); |
|
|
|
deviceProfileAlarm.setCreateRules(createRules); |
|
|
|
alarms.add(deviceProfileAlarm); |
|
|
|
profileData.setAlarms(alarms); |
|
|
|
profileData.setProvisionConfiguration(new AllowCreateNewDevicesDeviceProfileProvisionConfiguration("123")); |
|
|
|
} |
|
|
|
|
|
|
|
private void testReceivedInitialData() throws Exception { |
|
|
|
log.info("Checking received data"); |
|
|
|
Assert.assertTrue(edgeImitator.waitForMessages()); |
|
|
|
|
|
|
|
EdgeConfiguration configuration = edgeImitator.getConfiguration(); |
|
|
|
Assert.assertNotNull(configuration); |
|
|
|
|
|
|
|
testAutoGeneratedCodeByProtobuf(configuration); |
|
|
|
|
|
|
|
Optional<DeviceUpdateMsg> 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<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", |
|
|
|
new TypeReference<PageData<Device>>() {}, new PageLink(100)).getData(); |
|
|
|
Assert.assertTrue(edgeDevices.contains(device)); |
|
|
|
|
|
|
|
List<DeviceProfileUpdateMsg> deviceProfileUpdateMsgList = edgeImitator.findAllMessagesByType(DeviceProfileUpdateMsg.class); |
|
|
|
Assert.assertEquals(3, deviceProfileUpdateMsgList.size()); |
|
|
|
Optional<DeviceProfileUpdateMsg> 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<AssetUpdateMsg> 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<Asset> edgeAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", |
|
|
|
new TypeReference<PageData<Asset>>() {}, new PageLink(100)).getData(); |
|
|
|
Assert.assertTrue(edgeAssets.contains(asset)); |
|
|
|
|
|
|
|
testAutoGeneratedCodeByProtobuf(assetUpdateMsg); |
|
|
|
|
|
|
|
Optional<RuleChainUpdateMsg> 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<RuleChain> edgeRuleChains = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/ruleChains?", |
|
|
|
new TypeReference<PageData<RuleChain>>() {}, new PageLink(100)).getData(); |
|
|
|
Assert.assertTrue(edgeRuleChains.contains(ruleChain)); |
|
|
|
|
|
|
|
testAutoGeneratedCodeByProtobuf(ruleChainUpdateMsg); |
|
|
|
|
|
|
|
log.info("Received data checked"); |
|
|
|
} |
|
|
|
|
|
|
|
private void testDevices() throws Exception { |
|
|
|
@ -416,6 +510,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()); |
|
|
|
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<RuleNode> 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"); |
|
|
|
|
|
|
|
@ -894,166 +1049,75 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { |
|
|
|
Assert.assertEquals("key2", attributeDeleteMsg.getAttributeNames(1)); |
|
|
|
} |
|
|
|
|
|
|
|
private Device findDeviceByName(String deviceName) throws Exception { |
|
|
|
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", |
|
|
|
new TypeReference<PageData<Device>>() { |
|
|
|
}, new PageLink(100)).getData(); |
|
|
|
Optional<Device> 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<Asset> edgeAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", |
|
|
|
new TypeReference<PageData<Asset>>() { |
|
|
|
}, 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<DeviceUpdateMsg> 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<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", |
|
|
|
new TypeReference<PageData<Device>>() {}, new PageLink(100)).getData(); |
|
|
|
Assert.assertTrue(edgeDevices.contains(device)); |
|
|
|
|
|
|
|
List<DeviceProfileUpdateMsg> deviceProfileUpdateMsgList = edgeImitator.findAllMessagesByType(DeviceProfileUpdateMsg.class); |
|
|
|
Assert.assertEquals(3, deviceProfileUpdateMsgList.size()); |
|
|
|
Optional<DeviceProfileUpdateMsg> 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<AssetUpdateMsg> 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<Asset> edgeAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", |
|
|
|
new TypeReference<PageData<Asset>>() {}, new PageLink(100)).getData(); |
|
|
|
Assert.assertTrue(edgeAssets.contains(asset)); |
|
|
|
|
|
|
|
testAutoGeneratedCodeByProtobuf(assetUpdateMsg); |
|
|
|
|
|
|
|
Optional<RuleChainUpdateMsg> 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<RuleChain> edgeRuleChains = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/ruleChains?", |
|
|
|
new TypeReference<PageData<RuleChain>>() {}, 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); |
|
|
|
private void testRpcCall() throws Exception { |
|
|
|
Device device = findDeviceByName("Edge Device 1"); |
|
|
|
|
|
|
|
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder() |
|
|
|
.addRuleChainMetadataRequestMsg(ruleChainMetadataRequestMsgBuilder.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); |
|
|
|
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body); |
|
|
|
edgeImitator.expectMessageAmount(1); |
|
|
|
edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); |
|
|
|
Assert.assertTrue(edgeImitator.waitForResponses()); |
|
|
|
edgeEventService.saveAsync(edgeEvent); |
|
|
|
clusterService.onEdgeEventUpdate(tenantId, edge.getId()); |
|
|
|
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); |
|
|
|
Assert.assertTrue(latestMessage instanceof DeviceRpcCallMsg); |
|
|
|
DeviceRpcCallMsg latestDeviceRpcCallMsg = (DeviceRpcCallMsg) latestMessage; |
|
|
|
Assert.assertEquals("test_method", latestDeviceRpcCallMsg.getRequestMsg().getMethod()); |
|
|
|
} |
|
|
|
|
|
|
|
private void createRuleChainMetadata(RuleChain ruleChain) throws Exception { |
|
|
|
RuleChainMetaData ruleChainMetaData = new RuleChainMetaData(); |
|
|
|
ruleChainMetaData.setRuleChainId(ruleChain.getId()); |
|
|
|
private void testTimeseriesWithFailures() throws Exception { |
|
|
|
log.info("Testing timeseries with failures"); |
|
|
|
|
|
|
|
ObjectMapper mapper = new ObjectMapper(); |
|
|
|
int numberOfTimeseriesToSend = 1000; |
|
|
|
|
|
|
|
RuleNode ruleNode1 = new RuleNode(); |
|
|
|
ruleNode1.setName("name1"); |
|
|
|
ruleNode1.setType("type1"); |
|
|
|
ruleNode1.setConfiguration(mapper.readTree("\"key1\": \"val1\"")); |
|
|
|
edgeImitator.setRandomFailuresOnTimeseriesDownlink(true); |
|
|
|
// imitator will generate failure in 5% of cases
|
|
|
|
edgeImitator.setFailureProbability(5.0); |
|
|
|
|
|
|
|
RuleNode ruleNode2 = new RuleNode(); |
|
|
|
ruleNode2.setName("name2"); |
|
|
|
ruleNode2.setType("type2"); |
|
|
|
ruleNode2.setConfiguration(mapper.readTree("\"key2\": \"val2\"")); |
|
|
|
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()); |
|
|
|
} |
|
|
|
|
|
|
|
RuleNode ruleNode3 = new RuleNode(); |
|
|
|
ruleNode3.setName("name3"); |
|
|
|
ruleNode3.setType("type3"); |
|
|
|
ruleNode3.setConfiguration(mapper.readTree("\"key3\": \"val3\"")); |
|
|
|
Assert.assertTrue(edgeImitator.waitForMessages(60)); |
|
|
|
|
|
|
|
List<RuleNode> ruleNodes = new ArrayList<>(); |
|
|
|
ruleNodes.add(ruleNode1); |
|
|
|
ruleNodes.add(ruleNode2); |
|
|
|
ruleNodes.add(ruleNode3); |
|
|
|
ruleChainMetaData.setFirstNodeIndex(0); |
|
|
|
ruleChainMetaData.setNodes(ruleNodes); |
|
|
|
List<EntityDataProto> allTelemetryMsgs = edgeImitator.findAllMessagesByType(EntityDataProto.class); |
|
|
|
Assert.assertEquals(numberOfTimeseriesToSend, allTelemetryMsgs.size()); |
|
|
|
|
|
|
|
ruleChainMetaData.addConnectionInfo(0, 1, "success"); |
|
|
|
ruleChainMetaData.addConnectionInfo(0, 2, "fail"); |
|
|
|
ruleChainMetaData.addConnectionInfo(1, 2, "success"); |
|
|
|
for (int idx = 1; idx <= numberOfTimeseriesToSend; idx++) { |
|
|
|
Assert.assertTrue(isIdxExistsInTheDownlinkList(idx, allTelemetryMsgs)); |
|
|
|
} |
|
|
|
|
|
|
|
ruleChainMetaData.addRuleChainConnectionInfo(2, edge.getRootRuleChainId(), "success", mapper.createObjectNode()); |
|
|
|
edgeImitator.setRandomFailuresOnTimeseriesDownlink(false); |
|
|
|
log.info("Timeseries with failures tested successfully"); |
|
|
|
} |
|
|
|
|
|
|
|
doPost("/api/ruleChain/metadata", ruleChainMetaData, RuleChainMetaData.class); |
|
|
|
private boolean isIdxExistsInTheDownlinkList(int idx, List<EntityDataProto> 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 testSendMessagesToCloud() throws Exception { |
|
|
|
@ -1527,104 +1591,42 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { |
|
|
|
Assert.assertEquals(expectedValue, keyValueProto.getStringV()); |
|
|
|
} |
|
|
|
|
|
|
|
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()); |
|
|
|
// Utility methods
|
|
|
|
|
|
|
|
AbstractMessage latestMessage = edgeImitator.getLatestMessage(); |
|
|
|
Assert.assertTrue(latestMessage instanceof DeviceRpcCallMsg); |
|
|
|
DeviceRpcCallMsg latestDeviceRpcCallMsg = (DeviceRpcCallMsg) latestMessage; |
|
|
|
Assert.assertEquals("test_method", latestDeviceRpcCallMsg.getRequestMsg().getMethod()); |
|
|
|
private Device findDeviceByName(String deviceName) throws Exception { |
|
|
|
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", |
|
|
|
new TypeReference<PageData<Device>>() { |
|
|
|
}, new PageLink(100)).getData(); |
|
|
|
Optional<Device> 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 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<EntityDataProto> allTelemetryMsgs = edgeImitator.findAllMessagesByType(EntityDataProto.class); |
|
|
|
Assert.assertEquals(numberOfTimeseriesToSend, allTelemetryMsgs.size()); |
|
|
|
|
|
|
|
for (int idx = 1; idx <= numberOfTimeseriesToSend; idx++) { |
|
|
|
Assert.assertTrue(isIdxExistsInTheDownlinkList(idx, allTelemetryMsgs)); |
|
|
|
} |
|
|
|
private Asset findAssetByName(String assetName) throws Exception { |
|
|
|
List<Asset> edgeAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", |
|
|
|
new TypeReference<PageData<Asset>>() { |
|
|
|
}, new PageLink(100)).getData(); |
|
|
|
|
|
|
|
edgeImitator.setRandomFailuresOnTimeseriesDownlink(false); |
|
|
|
log.info("Timeseries with failures tested successfully"); |
|
|
|
Assert.assertEquals(1, edgeAssets.size()); |
|
|
|
Asset asset = edgeAssets.get(0); |
|
|
|
Assert.assertEquals(assetName, asset.getName()); |
|
|
|
return asset; |
|
|
|
} |
|
|
|
|
|
|
|
private boolean isIdxExistsInTheDownlinkList(int idx, List<EntityDataProto> 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 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 void extendDeviceProfileData(DeviceProfile deviceProfile) { |
|
|
|
DeviceProfileData profileData = deviceProfile.getProfileData(); |
|
|
|
List<DeviceProfileAlarm> alarms = new ArrayList<>(); |
|
|
|
DeviceProfileAlarm deviceProfileAlarm = new DeviceProfileAlarm(); |
|
|
|
deviceProfileAlarm.setAlarmType("High Temperature"); |
|
|
|
AlarmRule alarmRule = new AlarmRule(); |
|
|
|
alarmRule.setAlarmDetails("Alarm Details"); |
|
|
|
AlarmCondition alarmCondition = new AlarmCondition(); |
|
|
|
alarmCondition.setSpec(new SimpleAlarmConditionSpec()); |
|
|
|
List<AlarmConditionFilter> condition = new ArrayList<>(); |
|
|
|
AlarmConditionFilter alarmConditionFilter = new AlarmConditionFilter(); |
|
|
|
alarmConditionFilter.setKey(new AlarmConditionFilterKey(AlarmConditionKeyType.ATTRIBUTE, "temperature")); |
|
|
|
NumericFilterPredicate predicate = new NumericFilterPredicate(); |
|
|
|
predicate.setOperation(NumericFilterPredicate.NumericOperation.GREATER); |
|
|
|
predicate.setValue(new FilterPredicateValue<>(55.0)); |
|
|
|
alarmConditionFilter.setPredicate(predicate); |
|
|
|
alarmConditionFilter.setValueType(EntityKeyValueType.NUMERIC); |
|
|
|
condition.add(alarmConditionFilter); |
|
|
|
alarmCondition.setCondition(condition); |
|
|
|
alarmRule.setCondition(alarmCondition); |
|
|
|
deviceProfileAlarm.setClearRule(alarmRule); |
|
|
|
TreeMap<AlarmSeverity, AlarmRule> createRules = new TreeMap<>(); |
|
|
|
createRules.put(AlarmSeverity.CRITICAL, alarmRule); |
|
|
|
deviceProfileAlarm.setCreateRules(createRules); |
|
|
|
alarms.add(deviceProfileAlarm); |
|
|
|
profileData.setAlarms(alarms); |
|
|
|
profileData.setProvisionConfiguration(new AllowCreateNewDevicesDeviceProfileProvisionConfiguration("123")); |
|
|
|
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 EdgeEvent constructEdgeEvent(TenantId tenantId, EdgeId edgeId, EdgeEventActionType edgeEventAction, UUID entityId, EdgeEventType edgeEventType, JsonNode entityBody) { |
|
|
|
|