Browse Source

Merge pull request #53 from BohdanSmetanyuk/testing

Testing
pull/2436/head
VoBa 6 years ago
committed by GitHub
parent
commit
7bc7e1dbb8
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 663
      application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java
  2. 73
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java
  3. 133
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java

663
application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java

@ -15,7 +15,13 @@
*/ */
package org.thingsboard.server.edge; package org.thingsboard.server.edge;
import com.datastax.driver.core.utils.UUIDs;
import com.fasterxml.jackson.core.type.TypeReference; 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 lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.junit.After; import org.junit.After;
import org.junit.Assert; import org.junit.Assert;
@ -23,7 +29,7 @@ import org.junit.Before;
import org.junit.Test; import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.Dashboard;
import org.thingsboard.server.common.data.DashboardInfo; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.Tenant;
@ -33,7 +39,11 @@ import org.thingsboard.server.common.data.alarm.AlarmInfo;
import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmSeverity;
import org.thingsboard.server.common.data.alarm.AlarmStatus; import org.thingsboard.server.common.data.alarm.AlarmStatus;
import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.page.TimePageData; import org.thingsboard.server.common.data.page.TimePageData;
@ -42,17 +52,26 @@ import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChain;
import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.common.data.rule.RuleChainType;
import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.common.transport.adaptor.JsonConverter;
import org.thingsboard.server.controller.AbstractControllerTest; import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.edge.EdgeEventService;;
import org.thingsboard.server.edge.imitator.EdgeImitator; import org.thingsboard.server.edge.imitator.EdgeImitator;
import org.thingsboard.server.gen.edge.AlarmUpdateMsg;
import org.thingsboard.server.gen.edge.AssetUpdateMsg;
import org.thingsboard.server.gen.edge.DashboardUpdateMsg;
import org.thingsboard.server.gen.edge.DeviceUpdateMsg;
import org.thingsboard.server.gen.edge.EdgeConfiguration; import org.thingsboard.server.gen.edge.EdgeConfiguration;
import org.thingsboard.server.gen.edge.EntityDataProto;
import org.thingsboard.server.gen.edge.RelationUpdateMsg;
import org.thingsboard.server.gen.edge.RuleChainUpdateMsg;
import org.thingsboard.server.gen.edge.UpdateMsgType;
import org.thingsboard.server.gen.edge.UplinkMsg;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Set; import java.util.Optional;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@ -67,6 +86,9 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
private EdgeImitator edgeImitator; private EdgeImitator edgeImitator;
private Edge edge; private Edge edge;
@Autowired
private EdgeEventService edgeEventService;
@Before @Before
public void beforeTest() throws Exception { public void beforeTest() throws Exception {
loginSysAdmin(); loginSysAdmin();
@ -89,7 +111,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
edgeImitator = new EdgeImitator("localhost", 7070, edge.getRoutingKey(), edge.getSecret()); edgeImitator = new EdgeImitator("localhost", 7070, edge.getRoutingKey(), edge.getSecret());
// should be 3, but 3 events from sync service + 3 from controller. will be fixed in next releases // should be 3, but 3 events from sync service + 3 from controller. will be fixed in next releases
edgeImitator.getStorage().expectMessageAmount(6); edgeImitator.expectMessageAmount(6);
edgeImitator.connect(); edgeImitator.connect();
} }
@ -114,109 +136,149 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
testDashboards(); testDashboards();
testRelations(); testRelations();
testAlarms(); testAlarms();
testTimeseries();
testAttributes();
testSendMessagesToCloud();
} }
private void testReceivedInitialData() throws Exception { private void testReceivedInitialData() throws Exception {
log.info("Checking received data"); log.info("Checking received data");
edgeImitator.getStorage().waitForMessages(); edgeImitator.waitForMessages();
EdgeConfiguration configuration = edgeImitator.getStorage().getConfiguration(); EdgeConfiguration configuration = edgeImitator.getConfiguration();
Assert.assertNotNull(configuration); Assert.assertNotNull(configuration);
Map<UUID, EntityType> entities = edgeImitator.getStorage().getEntities(); Optional<DeviceUpdateMsg> optionalMsg1 = edgeImitator.findMessageByType(DeviceUpdateMsg.class);
Assert.assertFalse(entities.isEmpty()); Assert.assertTrue(optionalMsg1.isPresent());
DeviceUpdateMsg deviceUpdateMsg = optionalMsg1.get();
Set<UUID> devices = edgeImitator.getStorage().getEntitiesByType(EntityType.DEVICE); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceUpdateMsg.getMsgType());
Assert.assertEquals(1, devices.size()); UUID deviceUUID = new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB());
TimePageData<Device> pageDataDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", Device device = doGet("/api/device/" + deviceUUID.toString(), Device.class);
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)); Assert.assertNotNull(device);
for (Device device: pageDataDevices.getData()) { List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
Assert.assertTrue(devices.contains(device.getUuidId())); new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
} Assert.assertTrue(edgeDevices.contains(device));
Set<UUID> assets = edgeImitator.getStorage().getEntitiesByType(EntityType.ASSET); Optional<AssetUpdateMsg> optionalMsg2 = edgeImitator.findMessageByType(AssetUpdateMsg.class);
Assert.assertEquals(1, assets.size()); Assert.assertTrue(optionalMsg2.isPresent());
TimePageData<Asset> pageDataAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", AssetUpdateMsg assetUpdateMsg = optionalMsg2.get();
new TypeReference<TimePageData<Asset>>() {}, new TextPageLink(100)); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetUpdateMsg.getMsgType());
for (Asset asset: pageDataAssets.getData()) { UUID assetUUID = new UUID(assetUpdateMsg.getIdMSB(), assetUpdateMsg.getIdLSB());
Assert.assertTrue(assets.contains(asset.getUuidId())); 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<TimePageData<Asset>>() {}, new TextPageLink(100)).getData();
Assert.assertTrue(edgeAssets.contains(asset));
Optional<RuleChainUpdateMsg> optionalMsg3 = edgeImitator.findMessageByType(RuleChainUpdateMsg.class);
Assert.assertTrue(optionalMsg3.isPresent());
RuleChainUpdateMsg ruleChainUpdateMsg = optionalMsg3.get();
Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, ruleChainUpdateMsg.getMsgType());
UUID ruleChainUUID = new UUID(ruleChainUpdateMsg.getIdMSB(), ruleChainUpdateMsg.getIdLSB());
RuleChain ruleChain = doGet("/api/ruleChain/" + ruleChainUUID.toString(), RuleChain.class);
Assert.assertNotNull(ruleChain);
List<RuleChain> edgeRuleChains = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/ruleChains?",
new TypeReference<TimePageData<RuleChain>>() {}, new TextPageLink(100)).getData();
Assert.assertTrue(edgeRuleChains.contains(ruleChain));
Set<UUID> ruleChains = edgeImitator.getStorage().getEntitiesByType(EntityType.RULE_CHAIN);
Assert.assertEquals(1, ruleChains.size());
TimePageData<RuleChain> pageDataRuleChains = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/ruleChains?",
new TypeReference<TimePageData<RuleChain>>() {}, new TextPageLink(100));
for (RuleChain ruleChain: pageDataRuleChains.getData()) {
Assert.assertTrue(ruleChains.contains(ruleChain.getUuidId()));
}
log.info("Received data checked"); log.info("Received data checked");
} }
private void testDevices() throws Exception { private void testDevices() throws Exception {
log.info("Testing devices"); log.info("Testing devices");
Device device = new Device(); Device device = new Device();
device.setName("Edge Device 2"); device.setName("Edge Device 2");
device.setType("test"); device.setType("test");
Device savedDevice = doPost("/api/device", device, Device.class); Device savedDevice = doPost("/api/device", device, Device.class);
edgeImitator.getStorage().expectMessageAmount(1);
edgeImitator.expectMessageAmount(1);
doPost("/api/edge/" + edge.getId().getId().toString() doPost("/api/edge/" + edge.getId().getId().toString()
+ "/device/" + savedDevice.getId().getId().toString(), Device.class); + "/device/" + savedDevice.getId().getId().toString(), Device.class);
edgeImitator.waitForMessages();
TimePageData<Device> pageDataDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)); AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(pageDataDevices.getData().contains(savedDevice)); Assert.assertTrue(latestMessage instanceof DeviceUpdateMsg);
edgeImitator.getStorage().waitForMessages(); DeviceUpdateMsg deviceUpdateMsg = (DeviceUpdateMsg) latestMessage;
Set<UUID> devices = edgeImitator.getStorage().getEntitiesByType(EntityType.DEVICE); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, deviceUpdateMsg.getMsgType());
Assert.assertEquals(2, devices.size()); Assert.assertEquals(deviceUpdateMsg.getIdMSB(), savedDevice.getUuidId().getMostSignificantBits());
Assert.assertTrue(devices.contains(savedDevice.getUuidId())); Assert.assertEquals(deviceUpdateMsg.getIdLSB(), savedDevice.getUuidId().getLeastSignificantBits());
Assert.assertEquals(deviceUpdateMsg.getName(), savedDevice.getName());
edgeImitator.getStorage().expectMessageAmount(1); Assert.assertEquals(deviceUpdateMsg.getType(), savedDevice.getType());
edgeImitator.expectMessageAmount(1);
doDelete("/api/edge/" + edge.getId().getId().toString() doDelete("/api/edge/" + edge.getId().getId().toString()
+ "/device/" + savedDevice.getId().getId().toString(), Device.class); + "/device/" + savedDevice.getId().getId().toString(), Device.class);
pageDataDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", edgeImitator.waitForMessages();
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100));
Assert.assertFalse(pageDataDevices.getData().contains(savedDevice));
edgeImitator.getStorage().waitForMessages();
devices = edgeImitator.getStorage().getEntitiesByType(EntityType.DEVICE);
Assert.assertEquals(1, devices.size());
Assert.assertFalse(devices.contains(savedDevice.getUuidId()));
latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof DeviceUpdateMsg);
deviceUpdateMsg = (DeviceUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, deviceUpdateMsg.getMsgType());
Assert.assertEquals(deviceUpdateMsg.getIdMSB(), savedDevice.getUuidId().getMostSignificantBits());
Assert.assertEquals(deviceUpdateMsg.getIdLSB(), savedDevice.getUuidId().getLeastSignificantBits());
edgeImitator.expectMessageAmount(1);
doDelete("/api/device/" + savedDevice.getId().getId().toString()) doDelete("/api/device/" + savedDevice.getId().getId().toString())
.andExpect(status().isOk()); .andExpect(status().isOk());
edgeImitator.waitForMessages();
latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof DeviceUpdateMsg);
deviceUpdateMsg = (DeviceUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, deviceUpdateMsg.getMsgType());
Assert.assertEquals(deviceUpdateMsg.getIdMSB(), savedDevice.getUuidId().getMostSignificantBits());
Assert.assertEquals(deviceUpdateMsg.getIdLSB(), savedDevice.getUuidId().getLeastSignificantBits());
log.info("Devices tested successfully"); log.info("Devices tested successfully");
} }
private void testAssets() throws Exception { private void testAssets() throws Exception {
log.info("Testing assets"); log.info("Testing assets");
Asset asset = new Asset(); Asset asset = new Asset();
asset.setName("Edge Asset 2"); asset.setName("Edge Asset 2");
asset.setType("test"); asset.setType("test");
Asset savedAsset = doPost("/api/asset", asset, Asset.class); Asset savedAsset = doPost("/api/asset", asset, Asset.class);
edgeImitator.getStorage().expectMessageAmount(1);
edgeImitator.expectMessageAmount(1);
doPost("/api/edge/" + edge.getId().getId().toString() doPost("/api/edge/" + edge.getId().getId().toString()
+ "/asset/" + savedAsset.getId().getId().toString(), Asset.class); + "/asset/" + savedAsset.getId().getId().toString(), Asset.class);
edgeImitator.waitForMessages();
TimePageData<Asset> pageDataAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?",
new TypeReference<TimePageData<Asset>>() {}, new TextPageLink(100)); AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(pageDataAssets.getData().contains(savedAsset)); Assert.assertTrue(latestMessage instanceof AssetUpdateMsg);
edgeImitator.getStorage().waitForMessages(); AssetUpdateMsg assetUpdateMsg = (AssetUpdateMsg) latestMessage;
Set<UUID> assets = edgeImitator.getStorage().getEntitiesByType(EntityType.ASSET); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, assetUpdateMsg.getMsgType());
Assert.assertEquals(2, assets.size()); Assert.assertEquals(assetUpdateMsg.getIdMSB(), savedAsset.getUuidId().getMostSignificantBits());
Assert.assertTrue(assets.contains(savedAsset.getUuidId())); Assert.assertEquals(assetUpdateMsg.getIdLSB(), savedAsset.getUuidId().getLeastSignificantBits());
Assert.assertEquals(assetUpdateMsg.getName(), savedAsset.getName());
edgeImitator.getStorage().expectMessageAmount(1); Assert.assertEquals(assetUpdateMsg.getType(), savedAsset.getType());
edgeImitator.expectMessageAmount(1);
doDelete("/api/edge/" + edge.getId().getId().toString() doDelete("/api/edge/" + edge.getId().getId().toString()
+ "/asset/" + savedAsset.getId().getId().toString(), Asset.class); + "/asset/" + savedAsset.getId().getId().toString(), Asset.class);
pageDataAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?", edgeImitator.waitForMessages();
new TypeReference<TimePageData<Asset>>() {}, new TextPageLink(100));
Assert.assertFalse(pageDataAssets.getData().contains(savedAsset)); latestMessage = edgeImitator.getLatestMessage();
edgeImitator.getStorage().waitForMessages(); Assert.assertTrue(latestMessage instanceof AssetUpdateMsg);
assets = edgeImitator.getStorage().getEntitiesByType(EntityType.ASSET); assetUpdateMsg = (AssetUpdateMsg) latestMessage;
Assert.assertEquals(1, assets.size()); Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, assetUpdateMsg.getMsgType());
Assert.assertFalse(assets.contains(savedAsset.getUuidId())); Assert.assertEquals(assetUpdateMsg.getIdMSB(), savedAsset.getUuidId().getMostSignificantBits());
Assert.assertEquals(assetUpdateMsg.getIdLSB(), savedAsset.getUuidId().getLeastSignificantBits());
edgeImitator.expectMessageAmount(1);
doDelete("/api/asset/" + savedAsset.getId().getId().toString()) doDelete("/api/asset/" + savedAsset.getId().getId().toString())
.andExpect(status().isOk()); .andExpect(status().isOk());
edgeImitator.waitForMessages();
latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof AssetUpdateMsg);
assetUpdateMsg = (AssetUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, assetUpdateMsg.getMsgType());
Assert.assertEquals(assetUpdateMsg.getIdMSB(), savedAsset.getUuidId().getMostSignificantBits());
Assert.assertEquals(assetUpdateMsg.getIdLSB(), savedAsset.getUuidId().getLeastSignificantBits());
log.info("Assets tested successfully"); log.info("Assets tested successfully");
} }
@ -226,33 +288,45 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
ruleChain.setName("Edge Test Rule Chain"); ruleChain.setName("Edge Test Rule Chain");
ruleChain.setType(RuleChainType.EDGE); ruleChain.setType(RuleChainType.EDGE);
RuleChain savedRuleChain = doPost("/api/ruleChain", ruleChain, RuleChain.class); RuleChain savedRuleChain = doPost("/api/ruleChain", ruleChain, RuleChain.class);
edgeImitator.getStorage().expectMessageAmount(1);
edgeImitator.expectMessageAmount(1);
doPost("/api/edge/" + edge.getId().getId().toString() doPost("/api/edge/" + edge.getId().getId().toString()
+ "/ruleChain/" + savedRuleChain.getId().getId().toString(), RuleChain.class); + "/ruleChain/" + savedRuleChain.getId().getId().toString(), RuleChain.class);
edgeImitator.waitForMessages();
TimePageData<RuleChain> pageDataRuleChain = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/ruleChains?", AbstractMessage latestMessage = edgeImitator.getLatestMessage();
new TypeReference<TimePageData<RuleChain>>() {}, new TextPageLink(100)); Assert.assertTrue(latestMessage instanceof RuleChainUpdateMsg);
Assert.assertTrue(pageDataRuleChain.getData().contains(savedRuleChain)); RuleChainUpdateMsg ruleChainUpdateMsg = (RuleChainUpdateMsg) latestMessage;
edgeImitator.getStorage().waitForMessages(); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, ruleChainUpdateMsg.getMsgType());
Set<UUID> ruleChains = edgeImitator.getStorage().getEntitiesByType(EntityType.RULE_CHAIN); Assert.assertEquals(ruleChainUpdateMsg.getIdMSB(), savedRuleChain.getUuidId().getMostSignificantBits());
Assert.assertEquals(2, ruleChains.size()); Assert.assertEquals(ruleChainUpdateMsg.getIdLSB(), savedRuleChain.getUuidId().getLeastSignificantBits());
Assert.assertTrue(ruleChains.contains(savedRuleChain.getUuidId())); Assert.assertEquals(ruleChainUpdateMsg.getName(), savedRuleChain.getName());
edgeImitator.getStorage().expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
doDelete("/api/edge/" + edge.getId().getId().toString() doDelete("/api/edge/" + edge.getId().getId().toString()
+ "/ruleChain/" + savedRuleChain.getId().getId().toString(), RuleChain.class); + "/ruleChain/" + savedRuleChain.getId().getId().toString(), RuleChain.class);
pageDataRuleChain = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/ruleChains?", edgeImitator.waitForMessages();
new TypeReference<TimePageData<RuleChain>>() {}, new TextPageLink(100));
Assert.assertFalse(pageDataRuleChain.getData().contains(savedRuleChain));
edgeImitator.getStorage().waitForMessages();
ruleChains = edgeImitator.getStorage().getEntitiesByType(EntityType.RULE_CHAIN);
Assert.assertEquals(1, ruleChains.size());
Assert.assertFalse(ruleChains.contains(savedRuleChain.getUuidId()));
latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof RuleChainUpdateMsg);
ruleChainUpdateMsg = (RuleChainUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, ruleChainUpdateMsg.getMsgType());
Assert.assertEquals(ruleChainUpdateMsg.getIdMSB(), savedRuleChain.getUuidId().getMostSignificantBits());
Assert.assertEquals(ruleChainUpdateMsg.getIdLSB(), savedRuleChain.getUuidId().getLeastSignificantBits());
edgeImitator.expectMessageAmount(1);
doDelete("/api/ruleChain/" + savedRuleChain.getId().getId().toString()) doDelete("/api/ruleChain/" + savedRuleChain.getId().getId().toString())
.andExpect(status().isOk()); .andExpect(status().isOk());
log.info("RuleChains tested successfully"); edgeImitator.waitForMessages();
latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof RuleChainUpdateMsg);
ruleChainUpdateMsg = (RuleChainUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, ruleChainUpdateMsg.getMsgType());
Assert.assertEquals(ruleChainUpdateMsg.getIdMSB(), savedRuleChain.getUuidId().getMostSignificantBits());
Assert.assertEquals(ruleChainUpdateMsg.getIdLSB(), savedRuleChain.getUuidId().getLeastSignificantBits());
log.info("RuleChains tested successfully");
} }
private void testDashboards() throws Exception { private void testDashboards() throws Exception {
@ -260,33 +334,45 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
Dashboard dashboard = new Dashboard(); Dashboard dashboard = new Dashboard();
dashboard.setTitle("Edge Test Dashboard"); dashboard.setTitle("Edge Test Dashboard");
Dashboard savedDashboard = doPost("/api/dashboard", dashboard, Dashboard.class); Dashboard savedDashboard = doPost("/api/dashboard", dashboard, Dashboard.class);
edgeImitator.getStorage().expectMessageAmount(1);
edgeImitator.expectMessageAmount(1);
doPost("/api/edge/" + edge.getId().getId().toString() doPost("/api/edge/" + edge.getId().getId().toString()
+ "/dashboard/" + savedDashboard.getId().getId().toString(), Dashboard.class); + "/dashboard/" + savedDashboard.getId().getId().toString(), Dashboard.class);
edgeImitator.waitForMessages();
TimePageData<DashboardInfo> pageDataDashboard = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/dashboards?", AbstractMessage latestMessage = edgeImitator.getLatestMessage();
new TypeReference<TimePageData<DashboardInfo>>() {}, new TextPageLink(100)); Assert.assertTrue(latestMessage instanceof DashboardUpdateMsg);
Assert.assertTrue(pageDataDashboard.getData().stream().allMatch(dashboardInfo -> dashboardInfo.getUuidId().equals(savedDashboard.getUuidId()))); DashboardUpdateMsg dashboardUpdateMsg = (DashboardUpdateMsg) latestMessage;
edgeImitator.getStorage().waitForMessages(); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, dashboardUpdateMsg.getMsgType());
Set<UUID> dashboards = edgeImitator.getStorage().getEntitiesByType(EntityType.DASHBOARD); Assert.assertEquals(dashboardUpdateMsg.getIdMSB(), savedDashboard.getUuidId().getMostSignificantBits());
Assert.assertEquals(1, dashboards.size()); Assert.assertEquals(dashboardUpdateMsg.getIdLSB(), savedDashboard.getUuidId().getLeastSignificantBits());
Assert.assertTrue(dashboards.contains(savedDashboard.getUuidId())); Assert.assertEquals(dashboardUpdateMsg.getTitle(), savedDashboard.getName());
edgeImitator.getStorage().expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
doDelete("/api/edge/" + edge.getId().getId().toString() doDelete("/api/edge/" + edge.getId().getId().toString()
+ "/dashboard/" + savedDashboard.getId().getId().toString(), Dashboard.class); + "/dashboard/" + savedDashboard.getId().getId().toString(), Dashboard.class);
pageDataDashboard = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/dashboards?", edgeImitator.waitForMessages();
new TypeReference<TimePageData<DashboardInfo>>() {}, new TextPageLink(100));
Assert.assertFalse(pageDataDashboard.getData().stream().anyMatch(dashboardInfo -> dashboardInfo.getUuidId().equals(savedDashboard.getUuidId()))); latestMessage = edgeImitator.getLatestMessage();
edgeImitator.getStorage().waitForMessages(); Assert.assertTrue(latestMessage instanceof DashboardUpdateMsg);
dashboards = edgeImitator.getStorage().getEntitiesByType(EntityType.DASHBOARD); dashboardUpdateMsg = (DashboardUpdateMsg) latestMessage;
Assert.assertEquals(0, dashboards.size()); Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, dashboardUpdateMsg.getMsgType());
Assert.assertFalse(dashboards.contains(savedDashboard.getUuidId())); Assert.assertEquals(dashboardUpdateMsg.getIdMSB(), savedDashboard.getUuidId().getMostSignificantBits());
Assert.assertEquals(dashboardUpdateMsg.getIdLSB(), savedDashboard.getUuidId().getLeastSignificantBits());
edgeImitator.expectMessageAmount(1);
doDelete("/api/dashboard/" + savedDashboard.getId().getId().toString()) doDelete("/api/dashboard/" + savedDashboard.getId().getId().toString())
.andExpect(status().isOk()); .andExpect(status().isOk());
log.info("Dashboards tested successfully"); edgeImitator.waitForMessages();
latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof DashboardUpdateMsg);
dashboardUpdateMsg = (DashboardUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, dashboardUpdateMsg.getMsgType());
Assert.assertEquals(dashboardUpdateMsg.getIdMSB(), savedDashboard.getUuidId().getMostSignificantBits());
Assert.assertEquals(dashboardUpdateMsg.getIdLSB(), savedDashboard.getUuidId().getLeastSignificantBits());
log.info("Dashboards tested successfully");
} }
private void testRelations() throws Exception { private void testRelations() throws Exception {
@ -308,14 +394,24 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
relation.setFrom(device.getId()); relation.setFrom(device.getId());
relation.setTo(asset.getId()); relation.setTo(asset.getId());
relation.setTypeGroup(RelationTypeGroup.COMMON); relation.setTypeGroup(RelationTypeGroup.COMMON);
edgeImitator.getStorage().expectMessageAmount(1);
doPost("/api/relation", relation);
edgeImitator.getStorage().waitForMessages(); edgeImitator.expectMessageAmount(1);
List<EntityRelation> relations = edgeImitator.getStorage().getRelations(); doPost("/api/relation", relation);
Assert.assertEquals(1, relations.size()); edgeImitator.waitForMessages();
Assert.assertTrue(relations.contains(relation));
edgeImitator.getStorage().expectMessageAmount(1); AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof RelationUpdateMsg);
RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType());
Assert.assertEquals(relationUpdateMsg.getType(), relation.getType());
Assert.assertEquals(relationUpdateMsg.getFromIdMSB(), relation.getFrom().getId().getMostSignificantBits());
Assert.assertEquals(relationUpdateMsg.getFromIdLSB(), relation.getFrom().getId().getLeastSignificantBits());
Assert.assertEquals(relationUpdateMsg.getToEntityType(), relation.getTo().getEntityType().name());Assert.assertEquals(relationUpdateMsg.getFromIdMSB(), relation.getFrom().getId().getMostSignificantBits());
Assert.assertEquals(relationUpdateMsg.getToIdLSB(), relation.getTo().getId().getLeastSignificantBits());
Assert.assertEquals(relationUpdateMsg.getToEntityType(), relation.getTo().getEntityType().name());
Assert.assertEquals(relationUpdateMsg.getTypeGroup(), relation.getTypeGroup().name());
edgeImitator.expectMessageAmount(1);
doDelete("/api/relation?" + doDelete("/api/relation?" +
"fromId=" + relation.getFrom().getId().toString() + "fromId=" + relation.getFrom().getId().toString() +
"&fromType=" + relation.getFrom().getEntityType().name() + "&fromType=" + relation.getFrom().getEntityType().name() +
@ -324,15 +420,23 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
"&toId=" + relation.getTo().getId().toString() + "&toId=" + relation.getTo().getId().toString() +
"&toType=" + relation.getTo().getEntityType().name()) "&toType=" + relation.getTo().getEntityType().name())
.andExpect(status().isOk()); .andExpect(status().isOk());
edgeImitator.waitForMessages();
latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof RelationUpdateMsg);
relationUpdateMsg = (RelationUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, relationUpdateMsg.getMsgType());
Assert.assertEquals(relationUpdateMsg.getType(), relation.getType());
Assert.assertEquals(relationUpdateMsg.getFromIdMSB(), relation.getFrom().getId().getMostSignificantBits());
Assert.assertEquals(relationUpdateMsg.getFromIdLSB(), relation.getFrom().getId().getLeastSignificantBits());
Assert.assertEquals(relationUpdateMsg.getToEntityType(), relation.getTo().getEntityType().name());Assert.assertEquals(relationUpdateMsg.getFromIdMSB(), relation.getFrom().getId().getMostSignificantBits());
Assert.assertEquals(relationUpdateMsg.getToIdLSB(), relation.getTo().getId().getLeastSignificantBits());
Assert.assertEquals(relationUpdateMsg.getToEntityType(), relation.getTo().getEntityType().name());
Assert.assertEquals(relationUpdateMsg.getTypeGroup(), relation.getTypeGroup().name());
edgeImitator.getStorage().waitForMessages();
relations = edgeImitator.getStorage().getRelations();
Assert.assertEquals(0, relations.size());
Assert.assertFalse(relations.contains(relation));
log.info("Relations tested successfully"); log.info("Relations tested successfully");
} }
private void testAlarms() throws Exception { private void testAlarms() throws Exception {
log.info("Testing Alarms"); log.info("Testing Alarms");
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
@ -347,35 +451,313 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
alarm.setType("alarm"); alarm.setType("alarm");
alarm.setSeverity(AlarmSeverity.CRITICAL); alarm.setSeverity(AlarmSeverity.CRITICAL);
edgeImitator.getStorage().expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
Alarm savedAlarm = doPost("/api/alarm", alarm, Alarm.class); Alarm savedAlarm = doPost("/api/alarm", alarm, Alarm.class);
AlarmInfo alarmInfo = doGet("/api/alarm/info/" + savedAlarm.getId().getId().toString(), AlarmInfo.class); edgeImitator.waitForMessages();
edgeImitator.getStorage().waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertEquals(1, edgeImitator.getStorage().getAlarms().size()); Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg);
Assert.assertTrue(edgeImitator.getStorage().getAlarms().containsKey(alarmInfo.getType())); AlarmUpdateMsg alarmUpdateMsg = (AlarmUpdateMsg) latestMessage;
Assert.assertEquals(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()), alarmInfo.getStatus()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, alarmUpdateMsg.getMsgType());
edgeImitator.getStorage().expectMessageAmount(1); Assert.assertEquals(alarmUpdateMsg.getType(), savedAlarm.getType());
Assert.assertEquals(alarmUpdateMsg.getName(), savedAlarm.getName());
Assert.assertEquals(alarmUpdateMsg.getOriginatorName(), device.getName());
Assert.assertEquals(alarmUpdateMsg.getStatus(), savedAlarm.getStatus().name());
Assert.assertEquals(alarmUpdateMsg.getSeverity(), savedAlarm.getSeverity().name());
edgeImitator.expectMessageAmount(1);
doPost("/api/alarm/" + savedAlarm.getId().getId().toString() + "/ack"); doPost("/api/alarm/" + savedAlarm.getId().getId().toString() + "/ack");
edgeImitator.waitForMessages();
edgeImitator.getStorage().waitForMessages();
alarmInfo = doGet("/api/alarm/info/" + savedAlarm.getId().getId().toString(), AlarmInfo.class); latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()).isAck()); Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg);
Assert.assertEquals(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()), alarmInfo.getStatus()); alarmUpdateMsg = (AlarmUpdateMsg) latestMessage;
edgeImitator.getStorage().expectMessageAmount(1); Assert.assertEquals(UpdateMsgType.ALARM_ACK_RPC_MESSAGE, alarmUpdateMsg.getMsgType());
Assert.assertEquals(alarmUpdateMsg.getType(), savedAlarm.getType());
Assert.assertEquals(alarmUpdateMsg.getName(), savedAlarm.getName());
Assert.assertEquals(alarmUpdateMsg.getOriginatorName(), device.getName());
Assert.assertEquals(alarmUpdateMsg.getStatus(), AlarmStatus.ACTIVE_ACK.name());
edgeImitator.expectMessageAmount(1);
doPost("/api/alarm/" + savedAlarm.getId().getId().toString() + "/clear"); doPost("/api/alarm/" + savedAlarm.getId().getId().toString() + "/clear");
edgeImitator.waitForMessages();
edgeImitator.getStorage().waitForMessages(); latestMessage = edgeImitator.getLatestMessage();
alarmInfo = doGet("/api/alarm/info/" + savedAlarm.getId().getId().toString(), AlarmInfo.class); Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg);
Assert.assertTrue(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()).isAck()); alarmUpdateMsg = (AlarmUpdateMsg) latestMessage;
Assert.assertTrue(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()).isCleared()); Assert.assertEquals(UpdateMsgType.ALARM_CLEAR_RPC_MESSAGE, alarmUpdateMsg.getMsgType());
Assert.assertEquals(edgeImitator.getStorage().getAlarms().get(alarmInfo.getType()), alarmInfo.getStatus()); Assert.assertEquals(alarmUpdateMsg.getType(), savedAlarm.getType());
Assert.assertEquals(alarmUpdateMsg.getName(), savedAlarm.getName());
Assert.assertEquals(alarmUpdateMsg.getOriginatorName(), device.getName());
Assert.assertEquals(alarmUpdateMsg.getStatus(), AlarmStatus.CLEARED_ACK.name());
doDelete("/api/alarm/" + savedAlarm.getId().getId().toString()) doDelete("/api/alarm/" + savedAlarm.getId().getId().toString())
.andExpect(status().isOk()); .andExpect(status().isOk());
log.info("Alarms tested successfully"); log.info("Alarms tested successfully");
} }
private void testTimeseries() throws Exception {
log.info("Testing timeseries");
ObjectMapper mapper = new ObjectMapper();
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Assert.assertEquals(1, edgeDevices.size());
Device device = edgeDevices.get(0);
Assert.assertEquals("Edge Device 1", device.getName());
String timeseriesData = "{\"data\":{\"temperature\":25},\"ts\":" + System.currentTimeMillis() + "}";
JsonNode timeseriesEntityData = mapper.readTree(timeseriesData);
EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), ActionType.TIMESERIES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData);
edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent1);
edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof EntityDataProto);
EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage;
Assert.assertEquals(latestEntityDataMsg.getEntityIdMSB(), device.getUuidId().getMostSignificantBits());
Assert.assertEquals(latestEntityDataMsg.getEntityIdLSB(), device.getUuidId().getLeastSignificantBits());
Assert.assertEquals(latestEntityDataMsg.getEntityType(), device.getId().getEntityType().name());
Assert.assertTrue(latestEntityDataMsg.hasPostTelemetryMsg());
TransportProtos.PostTelemetryMsg postTelemetryMsg = latestEntityDataMsg.getPostTelemetryMsg();
Assert.assertEquals(1, postTelemetryMsg.getTsKvListCount());
TransportProtos.TsKvListProto tsKvListProto = postTelemetryMsg.getTsKvList(0);
Assert.assertEquals(timeseriesEntityData.get("ts").asLong(), tsKvListProto.getTs());
Assert.assertEquals(1, tsKvListProto.getKvCount());
TransportProtos.KeyValueProto keyValueProto = tsKvListProto.getKv(0);
Assert.assertEquals("temperature", keyValueProto.getKey());
Assert.assertEquals(25, keyValueProto.getLongV());
log.info("Timeseries tested successfully");
}
private void testAttributes() throws Exception {
log.info("Testing attributes");
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Assert.assertEquals(1, edgeDevices.size());
Device device = edgeDevices.get(0);
Assert.assertEquals("Edge Device 1", device.getName());
String attributesData = "{\"scope\":\"SERVER_SCOPE\",\"kv\":{\"key\":\"value\"}}";
JsonNode attributesEntityData = mapper.readTree(attributesData);
EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), ActionType.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(latestEntityDataMsg.getEntityIdMSB(), device.getUuidId().getMostSignificantBits());
Assert.assertEquals(latestEntityDataMsg.getEntityIdLSB(), device.getUuidId().getLeastSignificantBits());
Assert.assertEquals(latestEntityDataMsg.getEntityType(), device.getId().getEntityType().name());
Assert.assertEquals(latestEntityDataMsg.getPostAttributeScope(), attributesEntityData.get("scope").asText());
Assert.assertTrue(latestEntityDataMsg.hasAttributesUpdatedMsg());
TransportProtos.PostAttributeMsg attributesUpdatedMsg = latestEntityDataMsg.getAttributesUpdatedMsg();
Assert.assertEquals(1, attributesUpdatedMsg.getKvCount());
TransportProtos.KeyValueProto keyValueProto = attributesUpdatedMsg.getKv(0);
Assert.assertEquals("key", keyValueProto.getKey());
Assert.assertEquals("value", keyValueProto.getStringV());
((ObjectNode) attributesEntityData).put("isPostAttributes", true);
EdgeEvent edgeEvent2 = constructEdgeEvent(tenantId, edge.getId(), ActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData);
edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent2);
edgeImitator.waitForMessages();
latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof EntityDataProto);
latestEntityDataMsg = (EntityDataProto) latestMessage;
Assert.assertEquals(latestEntityDataMsg.getEntityIdMSB(), device.getUuidId().getMostSignificantBits());
Assert.assertEquals(latestEntityDataMsg.getEntityIdLSB(), device.getUuidId().getLeastSignificantBits());
Assert.assertEquals(latestEntityDataMsg.getEntityType(), device.getId().getEntityType().name());
Assert.assertEquals(latestEntityDataMsg.getPostAttributeScope(), attributesEntityData.get("scope").asText());
Assert.assertTrue(latestEntityDataMsg.hasPostAttributesMsg());
attributesUpdatedMsg = latestEntityDataMsg.getPostAttributesMsg();
Assert.assertEquals(1, attributesUpdatedMsg.getKvCount());
keyValueProto = attributesUpdatedMsg.getKv(0);
Assert.assertEquals("key", keyValueProto.getKey());
Assert.assertEquals("value", keyValueProto.getStringV());
log.info("Attributes tested successfully");
}
private void testSendMessagesToCloud() throws Exception {
log.info("Sending messages to cloud");
sendDevice();
sendAlarm();
sendTelemetry();
sendRelation();
sendDeleteDeviceOnEdge();
log.info("Messages were sent successfully");
}
private void sendDevice() throws Exception {
UUID uuid = UUIDs.timeBased();
UplinkMsg.Builder builder = UplinkMsg.newBuilder();
DeviceUpdateMsg.Builder deviceUpdateMsgBuilder = DeviceUpdateMsg.newBuilder();
deviceUpdateMsgBuilder.setIdMSB(uuid.getMostSignificantBits());
deviceUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits());
deviceUpdateMsgBuilder.setName("Edge Device 2");
deviceUpdateMsgBuilder.setType("test");
deviceUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE);
builder.addDeviceUpdateMsg(deviceUpdateMsgBuilder.build());
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(builder.build());
edgeImitator.waitForResponses();
Device device = doGet("/api/device/" + uuid.toString(), Device.class);
Assert.assertNotNull(device);
Assert.assertEquals("Edge Device 2", device.getName());
}
private void sendAlarm() throws Exception {
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Optional<Device> foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 2")).findAny();
Assert.assertTrue(foundDevice.isPresent());
Device device = foundDevice.get();
UplinkMsg.Builder builder = UplinkMsg.newBuilder();
AlarmUpdateMsg.Builder alarmUpdateMgBuilder = AlarmUpdateMsg.newBuilder();
alarmUpdateMgBuilder.setName("alarm from edge");
alarmUpdateMgBuilder.setStatus(AlarmStatus.ACTIVE_UNACK.name());
alarmUpdateMgBuilder.setSeverity(AlarmSeverity.CRITICAL.name());
alarmUpdateMgBuilder.setOriginatorName(device.getName());
alarmUpdateMgBuilder.setOriginatorType(EntityType.DEVICE.name());
builder.addAlarmUpdateMsg(alarmUpdateMgBuilder.build());
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(builder.build());
edgeImitator.waitForResponses();
List<AlarmInfo> alarms = doGetTypedWithPageLink("/api/alarm/{entityType}/{entityId}?",
new TypeReference<TimePageData<AlarmInfo>>() {},
new TextPageLink(100), device.getId().getEntityType().name(), device.getId().getId().toString())
.getData();
Optional<AlarmInfo> foundAlarm = alarms.stream().filter(alarm -> alarm.getType().equals("alarm from edge")).findAny();
Assert.assertTrue(foundAlarm.isPresent());
AlarmInfo alarmInfo = foundAlarm.get();
Assert.assertEquals(device.getId(), alarmInfo.getOriginator());
Assert.assertEquals(AlarmStatus.ACTIVE_UNACK, alarmInfo.getStatus());
Assert.assertEquals(AlarmSeverity.CRITICAL, alarmInfo.getSeverity());
}
private void sendRelation() throws Exception {
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Optional<Device> foundDevice1 = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 1")).findAny();
Assert.assertTrue(foundDevice1.isPresent());
Device device1 = foundDevice1.get();
Optional<Device> foundDevice2 = edgeDevices.stream().filter(device2 -> device2.getName().equals("Edge Device 2")).findAny();
Assert.assertTrue(foundDevice2.isPresent());
Device device2 = foundDevice2.get();
UplinkMsg.Builder builder = UplinkMsg.newBuilder();
RelationUpdateMsg.Builder relationUpdateMsgBuilder = RelationUpdateMsg.newBuilder();
relationUpdateMsgBuilder.setType("test");
relationUpdateMsgBuilder.setTypeGroup(RelationTypeGroup.COMMON.name());
relationUpdateMsgBuilder.setToIdMSB(device1.getId().getId().getMostSignificantBits());
relationUpdateMsgBuilder.setToIdLSB(device1.getId().getId().getLeastSignificantBits());
relationUpdateMsgBuilder.setToEntityType(device1.getId().getEntityType().name());
relationUpdateMsgBuilder.setFromIdMSB(device2.getId().getId().getMostSignificantBits());
relationUpdateMsgBuilder.setFromIdLSB(device2.getId().getId().getLeastSignificantBits());
relationUpdateMsgBuilder.setFromEntityType(device2.getId().getEntityType().name());
relationUpdateMsgBuilder.setAdditionalInfo("{}");
builder.addRelationUpdateMsg(relationUpdateMsgBuilder.build());
UplinkMsg msg = builder.build();
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(msg);
edgeImitator.waitForResponses();
EntityRelation relation = doGet("/api/relation?" +
"&fromId=" + device2.getId().getId().toString() +
"&fromType=" + device2.getId().getEntityType().name() +
"&relationType=" + "test" +
"&relationTypeGroup=" + RelationTypeGroup.COMMON.name() +
"&toId=" + device1.getId().getId().toString() +
"&toType=" + device1.getId().getEntityType().name(), EntityRelation.class);
Assert.assertNotNull(relation);
}
private void sendTelemetry() throws Exception {
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Optional<Device> foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 2")).findAny();
Assert.assertTrue(foundDevice.isPresent());
Device device = foundDevice.get();
edgeImitator.expectResponsesAmount(2);
JsonObject data = new JsonObject();
String timeseriesKey = "key";
String timeseriesValue = "25";
data.addProperty(timeseriesKey, timeseriesValue);
UplinkMsg.Builder builder1 = UplinkMsg.newBuilder();
EntityDataProto.Builder entityDataBuilder = EntityDataProto.newBuilder();
entityDataBuilder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(data, System.currentTimeMillis()));
entityDataBuilder.setEntityType(device.getId().getEntityType().name());
entityDataBuilder.setEntityIdMSB(device.getUuidId().getMostSignificantBits());
entityDataBuilder.setEntityIdLSB(device.getUuidId().getLeastSignificantBits());
builder1.addEntityData(entityDataBuilder.build());
edgeImitator.sendUplinkMsg(builder1.build());
JsonObject attributesData = new JsonObject();
String attributesKey = "test_attr";
String attributesValue = "test_value";
attributesData.addProperty(attributesKey, attributesValue);
UplinkMsg.Builder builder2 = UplinkMsg.newBuilder();
EntityDataProto.Builder entityDataBuilder2 = EntityDataProto.newBuilder();
entityDataBuilder2.setEntityType(device.getId().getEntityType().name());
entityDataBuilder2.setEntityIdMSB(device.getId().getId().getMostSignificantBits());
entityDataBuilder2.setEntityIdLSB(device.getId().getId().getLeastSignificantBits());
entityDataBuilder2.setAttributesUpdatedMsg(JsonConverter.convertToAttributesProto(attributesData));
entityDataBuilder2.setPostAttributeScope(DataConstants.SERVER_SCOPE);
builder2.addEntityData(entityDataBuilder2.build());
edgeImitator.sendUplinkMsg(builder2.build());
edgeImitator.waitForResponses();
Thread.sleep(1000);
Map<String, List<Map<String, String>>> timeseries = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/values/timeseries?keys=" + timeseriesKey, Map.class);
Assert.assertTrue(timeseries.containsKey(timeseriesKey));
Assert.assertEquals(1, timeseries.get(timeseriesKey).size());
Assert.assertEquals(timeseriesValue, timeseries.get(timeseriesKey).get(0).get("value"));
List<Map<String, String>> attributes = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getId() + "/values/attributes/" + DataConstants.SERVER_SCOPE, List.class);
Assert.assertEquals(1, attributes.size());
Assert.assertEquals(attributes.get(0).get("key"), attributesKey);
Assert.assertEquals(attributes.get(0).get("value"), attributesValue);
}
private void sendDeleteDeviceOnEdge() throws Exception {
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Optional<Device> foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 2")).findAny();
Assert.assertTrue(foundDevice.isPresent());
Device device = foundDevice.get();
UplinkMsg.Builder builder = UplinkMsg.newBuilder();
DeviceUpdateMsg.Builder deviceDeleteMsgBuilder = DeviceUpdateMsg.newBuilder();
deviceDeleteMsgBuilder.setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE);
deviceDeleteMsgBuilder.setIdMSB(device.getId().getId().getMostSignificantBits());
deviceDeleteMsgBuilder.setIdLSB(device.getId().getId().getLeastSignificantBits());
builder.addDeviceUpdateMsg(deviceDeleteMsgBuilder.build());
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(builder.build());
edgeImitator.waitForResponses();
device = doGet("/api/device/" + device.getId().getId().toString(), Device.class);
Assert.assertNotNull(device);
edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {
}, new TextPageLink(100)).getData();
Assert.assertFalse(edgeDevices.contains(device));
}
private void installation() throws Exception { private void installation() throws Exception {
edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class); edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class);
@ -413,4 +795,15 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
doDelete("/api/edge/" + edge.getId().getId().toString()) doDelete("/api/edge/" + edge.getId().getId().toString())
.andExpect(status().isOk()); .andExpect(status().isOk());
} }
private EdgeEvent constructEdgeEvent(TenantId tenantId, EdgeId edgeId, ActionType edgeEventAction, UUID entityId, EdgeEventType edgeEventType, JsonNode entityBody) {
EdgeEvent edgeEvent = new EdgeEvent();
edgeEvent.setEdgeId(edgeId);
edgeEvent.setTenantId(tenantId);
edgeEvent.setAction(edgeEventAction.name());
edgeEvent.setEntityId(entityId);
edgeEvent.setType(edgeEventType);
edgeEvent.setBody(entityBody);
return edgeEvent;
}
} }

73
application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java

@ -19,12 +19,12 @@ import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.MoreExecutors;
import com.google.protobuf.AbstractMessage;
import lombok.Getter; import lombok.Getter;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.Nullable; import org.checkerframework.checker.nullness.qual.Nullable;
import org.thingsboard.edge.rpc.EdgeGrpcClient; import org.thingsboard.edge.rpc.EdgeGrpcClient;
import org.thingsboard.edge.rpc.EdgeRpcClient; import org.thingsboard.edge.rpc.EdgeRpcClient;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.gen.edge.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.AlarmUpdateMsg;
import org.thingsboard.server.gen.edge.AssetUpdateMsg; import org.thingsboard.server.gen.edge.AssetUpdateMsg;
import org.thingsboard.server.gen.edge.DashboardUpdateMsg; import org.thingsboard.server.gen.edge.DashboardUpdateMsg;
@ -32,14 +32,18 @@ import org.thingsboard.server.gen.edge.DeviceUpdateMsg;
import org.thingsboard.server.gen.edge.DownlinkMsg; import org.thingsboard.server.gen.edge.DownlinkMsg;
import org.thingsboard.server.gen.edge.DownlinkResponseMsg; import org.thingsboard.server.gen.edge.DownlinkResponseMsg;
import org.thingsboard.server.gen.edge.EdgeConfiguration; import org.thingsboard.server.gen.edge.EdgeConfiguration;
import org.thingsboard.server.gen.edge.EntityDataProto;
import org.thingsboard.server.gen.edge.RelationUpdateMsg; import org.thingsboard.server.gen.edge.RelationUpdateMsg;
import org.thingsboard.server.gen.edge.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.RuleChainUpdateMsg;
import org.thingsboard.server.gen.edge.UplinkMsg;
import org.thingsboard.server.gen.edge.UplinkResponseMsg; import org.thingsboard.server.gen.edge.UplinkResponseMsg;
import java.lang.reflect.Field; import java.lang.reflect.Field;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.UUID; import java.util.Optional;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
@Slf4j @Slf4j
public class EdgeImitator { public class EdgeImitator {
@ -49,17 +53,24 @@ public class EdgeImitator {
private EdgeRpcClient edgeRpcClient; private EdgeRpcClient edgeRpcClient;
@Getter private CountDownLatch messagesLatch;
private EdgeStorage storage; private CountDownLatch responsesLatch;
@Getter
private EdgeConfiguration configuration;
@Getter
private List<AbstractMessage> downlinkMsgs;
public EdgeImitator(String host, int port, String routingKey, String routingSecret) throws NoSuchFieldException, IllegalAccessException { public EdgeImitator(String host, int port, String routingKey, String routingSecret) throws NoSuchFieldException, IllegalAccessException {
edgeRpcClient = new EdgeGrpcClient(); edgeRpcClient = new EdgeGrpcClient();
storage = new EdgeStorage(); messagesLatch = new CountDownLatch(0);
responsesLatch = new CountDownLatch(0);
downlinkMsgs = new ArrayList<>();
this.routingKey = routingKey; this.routingKey = routingKey;
this.routingSecret = routingSecret; this.routingSecret = routingSecret;
setEdgeCredentials("rpcHost", host); setEdgeCredentials("rpcHost", host);
setEdgeCredentials("rpcPort", port); setEdgeCredentials("rpcPort", port);
setEdgeCredentials("keepAliveTimeSec", 300);
} }
private void setEdgeCredentials(String fieldName, Object value) throws NoSuchFieldException, IllegalAccessException { private void setEdgeCredentials(String fieldName, Object value) throws NoSuchFieldException, IllegalAccessException {
@ -83,12 +94,17 @@ public class EdgeImitator {
edgeRpcClient.disconnect(false); edgeRpcClient.disconnect(false);
} }
public void sendUplinkMsg(UplinkMsg uplinkMsg) {
edgeRpcClient.sendUplinkMsg(uplinkMsg);
}
private void onUplinkResponse(UplinkResponseMsg msg) { private void onUplinkResponse(UplinkResponseMsg msg) {
log.info("onUplinkResponse: {}", msg); log.info("onUplinkResponse: {}", msg);
responsesLatch.countDown();
} }
private void onEdgeUpdate(EdgeConfiguration edgeConfiguration) { private void onEdgeUpdate(EdgeConfiguration edgeConfiguration) {
storage.setConfiguration(edgeConfiguration); this.configuration = edgeConfiguration;
} }
private void onDownlink(DownlinkMsg downlinkMsg) { private void onDownlink(DownlinkMsg downlinkMsg) {
@ -116,35 +132,68 @@ public class EdgeImitator {
List<ListenableFuture<Void>> result = new ArrayList<>(); List<ListenableFuture<Void>> result = new ArrayList<>();
if (downlinkMsg.getDeviceUpdateMsgList() != null && !downlinkMsg.getDeviceUpdateMsgList().isEmpty()) { if (downlinkMsg.getDeviceUpdateMsgList() != null && !downlinkMsg.getDeviceUpdateMsgList().isEmpty()) {
for (DeviceUpdateMsg deviceUpdateMsg: downlinkMsg.getDeviceUpdateMsgList()) { for (DeviceUpdateMsg deviceUpdateMsg: downlinkMsg.getDeviceUpdateMsgList()) {
result.add(storage.processEntity(deviceUpdateMsg.getMsgType(), EntityType.DEVICE, new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB()))); saveDownlinkMsg(deviceUpdateMsg);
} }
} }
if (downlinkMsg.getAssetUpdateMsgList() != null && !downlinkMsg.getAssetUpdateMsgList().isEmpty()) { if (downlinkMsg.getAssetUpdateMsgList() != null && !downlinkMsg.getAssetUpdateMsgList().isEmpty()) {
for (AssetUpdateMsg assetUpdateMsg: downlinkMsg.getAssetUpdateMsgList()) { for (AssetUpdateMsg assetUpdateMsg: downlinkMsg.getAssetUpdateMsgList()) {
result.add(storage.processEntity(assetUpdateMsg.getMsgType(), EntityType.ASSET, new UUID(assetUpdateMsg.getIdMSB(), assetUpdateMsg.getIdLSB()))); saveDownlinkMsg(assetUpdateMsg);
} }
} }
if (downlinkMsg.getRuleChainUpdateMsgList() != null && !downlinkMsg.getRuleChainUpdateMsgList().isEmpty()) { if (downlinkMsg.getRuleChainUpdateMsgList() != null && !downlinkMsg.getRuleChainUpdateMsgList().isEmpty()) {
for (RuleChainUpdateMsg ruleChainUpdateMsg: downlinkMsg.getRuleChainUpdateMsgList()) { for (RuleChainUpdateMsg ruleChainUpdateMsg: downlinkMsg.getRuleChainUpdateMsgList()) {
result.add(storage.processEntity(ruleChainUpdateMsg.getMsgType(), EntityType.RULE_CHAIN, new UUID(ruleChainUpdateMsg.getIdMSB(), ruleChainUpdateMsg.getIdLSB()))); saveDownlinkMsg(ruleChainUpdateMsg);
} }
} }
if (downlinkMsg.getDashboardUpdateMsgList() != null && !downlinkMsg.getDashboardUpdateMsgList().isEmpty()) { if (downlinkMsg.getDashboardUpdateMsgList() != null && !downlinkMsg.getDashboardUpdateMsgList().isEmpty()) {
for (DashboardUpdateMsg dashboardUpdateMsg: downlinkMsg.getDashboardUpdateMsgList()) { for (DashboardUpdateMsg dashboardUpdateMsg: downlinkMsg.getDashboardUpdateMsgList()) {
result.add(storage.processEntity(dashboardUpdateMsg.getMsgType(), EntityType.DASHBOARD, new UUID(dashboardUpdateMsg.getIdMSB(), dashboardUpdateMsg.getIdLSB()))); saveDownlinkMsg(dashboardUpdateMsg);
} }
} }
if (downlinkMsg.getRelationUpdateMsgList() != null && !downlinkMsg.getRelationUpdateMsgList().isEmpty()) { if (downlinkMsg.getRelationUpdateMsgList() != null && !downlinkMsg.getRelationUpdateMsgList().isEmpty()) {
for (RelationUpdateMsg relationUpdateMsg: downlinkMsg.getRelationUpdateMsgList()) { for (RelationUpdateMsg relationUpdateMsg: downlinkMsg.getRelationUpdateMsgList()) {
result.add(storage.processRelation(relationUpdateMsg)); saveDownlinkMsg(relationUpdateMsg);
} }
} }
if (downlinkMsg.getAlarmUpdateMsgList() != null && !downlinkMsg.getAlarmUpdateMsgList().isEmpty()) { if (downlinkMsg.getAlarmUpdateMsgList() != null && !downlinkMsg.getAlarmUpdateMsgList().isEmpty()) {
for (AlarmUpdateMsg alarmUpdateMsg: downlinkMsg.getAlarmUpdateMsgList()) { for (AlarmUpdateMsg alarmUpdateMsg: downlinkMsg.getAlarmUpdateMsgList()) {
result.add(storage.processAlarm(alarmUpdateMsg)); saveDownlinkMsg(alarmUpdateMsg);
}
}
if (downlinkMsg.getEntityDataList() != null && !downlinkMsg.getEntityDataList().isEmpty()) {
for (EntityDataProto entityData: downlinkMsg.getEntityDataList()) {
saveDownlinkMsg(entityData);
} }
} }
return Futures.allAsList(result); return Futures.allAsList(result);
} }
private ListenableFuture<Void> saveDownlinkMsg(AbstractMessage message) {
downlinkMsgs.add(message);
messagesLatch.countDown();
return Futures.immediateFuture(null);
}
public void waitForMessages() throws InterruptedException {
messagesLatch.await(5, TimeUnit.SECONDS);
}
public void expectMessageAmount(int messageAmount) {
messagesLatch = new CountDownLatch(messageAmount);
}
public void waitForResponses() throws InterruptedException { responsesLatch.await(5, TimeUnit.SECONDS); }
public void expectResponsesAmount(int messageAmount) {
responsesLatch = new CountDownLatch(messageAmount);
}
public <T> Optional<T> findMessageByType(Class<T> tClass) {
return (Optional<T>) downlinkMsgs.stream().filter(downlinkMsg -> downlinkMsg.getClass().isAssignableFrom(tClass)).findAny();
}
public AbstractMessage getLatestMessage() {
return downlinkMsgs.get(downlinkMsgs.size() - 1);
}
} }

133
application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java

@ -1,133 +0,0 @@
/**
* Copyright © 2016-2020 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.edge.imitator;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.Getter;
import lombok.Setter;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.alarm.AlarmStatus;
import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.gen.edge.AlarmUpdateMsg;
import org.thingsboard.server.gen.edge.EdgeConfiguration;
import org.thingsboard.server.gen.edge.RelationUpdateMsg;
import org.thingsboard.server.gen.edge.UpdateMsgType;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
@Slf4j
@Getter
@Setter
public class EdgeStorage {
private EdgeConfiguration configuration;
private CountDownLatch latch;
private Map<UUID, EntityType> entities;
private Map<String, AlarmStatus> alarms;
private List<EntityRelation> relations;
public EdgeStorage() {
latch = new CountDownLatch(0);
entities = new HashMap<>();
alarms = new HashMap<>();
relations = new ArrayList<>();
}
public ListenableFuture<Void> processEntity(UpdateMsgType msgType, EntityType type, UUID uuid) {
switch (msgType) {
case ENTITY_CREATED_RPC_MESSAGE:
case ENTITY_UPDATED_RPC_MESSAGE:
entities.put(uuid, type);
latch.countDown();
break;
case ENTITY_DELETED_RPC_MESSAGE:
if (entities.remove(uuid) != null) {
latch.countDown();
}
break;
}
return Futures.immediateFuture(null);
}
public ListenableFuture<Void> processRelation(RelationUpdateMsg relationMsg) {
boolean result = false;
EntityRelation relation = new EntityRelation();
relation.setType(relationMsg.getType());
relation.setTypeGroup(RelationTypeGroup.valueOf(relationMsg.getTypeGroup()));
relation.setTo(EntityIdFactory.getByTypeAndUuid(relationMsg.getToEntityType(), new UUID(relationMsg.getToIdMSB(), relationMsg.getToIdLSB())));
relation.setFrom(EntityIdFactory.getByTypeAndUuid(relationMsg.getFromEntityType(), new UUID(relationMsg.getFromIdMSB(), relationMsg.getFromIdLSB())));
switch (relationMsg.getMsgType()) {
case ENTITY_CREATED_RPC_MESSAGE:
case ENTITY_UPDATED_RPC_MESSAGE:
result = relations.add(relation);
break;
case ENTITY_DELETED_RPC_MESSAGE:
result = relations.remove(relation);
break;
}
if (result) {
latch.countDown();
}
return Futures.immediateFuture(null);
}
public ListenableFuture<Void> processAlarm(AlarmUpdateMsg alarmMsg) {
switch (alarmMsg.getMsgType()) {
case ENTITY_CREATED_RPC_MESSAGE:
case ENTITY_UPDATED_RPC_MESSAGE:
case ALARM_ACK_RPC_MESSAGE:
case ALARM_CLEAR_RPC_MESSAGE:
alarms.put(alarmMsg.getType(), AlarmStatus.valueOf(alarmMsg.getStatus()));
latch.countDown();
break;
case ENTITY_DELETED_RPC_MESSAGE:
if (alarms.remove(alarmMsg.getName()) != null) {
latch.countDown();
}
break;
}
return Futures.immediateFuture(null);
}
public Set<UUID> getEntitiesByType(EntityType type) {
return entities.entrySet().stream()
.filter(entry -> entry.getValue().equals(type))
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)).keySet();
}
public void waitForMessages() throws InterruptedException {
latch.await(5, TimeUnit.SECONDS);
}
public void expectMessageAmount(int messageAmount) {
latch = new CountDownLatch(messageAmount);
}
}
Loading…
Cancel
Save