Browse Source

Merge remote-tracking branch 'origin/feature/edge' into feature/edge

pull/2436/head
Artem Babak 6 years ago
parent
commit
94437d84fb
  1. 1
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  2. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java
  3. 20
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java
  4. 591
      application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java
  5. 12
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java
  6. 1
      common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java
  7. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java
  8. 30
      ui/src/app/event/event-table.directive.js

1
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java

@ -335,6 +335,7 @@ public final class EdgeGrpcSession implements Closeable {
downlinkMsg = processEntityMessage(edgeEvent, edgeEvent.getAction()); downlinkMsg = processEntityMessage(edgeEvent, edgeEvent.getAction());
break; break;
case ATTRIBUTES_UPDATED: case ATTRIBUTES_UPDATED:
case POST_ATTRIBUTES:
case ATTRIBUTES_DELETED: case ATTRIBUTES_DELETED:
case TIMESERIES_UPDATED: case TIMESERIES_UPDATED:
downlinkMsg = processTelemetryMessage(edgeEvent); downlinkMsg = processTelemetryMessage(edgeEvent);

4
application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java

@ -89,9 +89,11 @@ public class DeviceMsgConstructor {
.setRequestIdMSB(request.getRequestUUID().getMostSignificantBits()) .setRequestIdMSB(request.getRequestUUID().getMostSignificantBits())
.setRequestIdLSB(request.getRequestUUID().getLeastSignificantBits()) .setRequestIdLSB(request.getRequestUUID().getLeastSignificantBits())
.setExpirationTime(request.getExpirationTime()) .setExpirationTime(request.getExpirationTime())
.setOriginServiceId(request.getOriginServiceId())
.setOneway(request.isOneway()) .setOneway(request.isOneway())
.setRequestMsg(requestBuilder.build()); .setRequestMsg(requestBuilder.build());
if (request.getOriginServiceId() != null) {
builder.setOriginServiceId(request.getOriginServiceId());
}
return builder.build(); return builder.build();
} }
} }

20
application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java

@ -60,15 +60,21 @@ public class EntityDataMsgConstructor {
case ATTRIBUTES_UPDATED: case ATTRIBUTES_UPDATED:
try { try {
JsonObject data = entityData.getAsJsonObject(); JsonObject data = entityData.getAsJsonObject();
TransportProtos.PostAttributeMsg postAttributeMsg = JsonConverter.convertToAttributesProto(data.getAsJsonObject("kv")); TransportProtos.PostAttributeMsg attributesUpdatedMsg = JsonConverter.convertToAttributesProto(data.getAsJsonObject("kv"));
if (data.has("isPostAttributes") && data.getAsJsonPrimitive("isPostAttributes").getAsBoolean()) { builder.setAttributesUpdatedMsg(attributesUpdatedMsg);
builder.setPostAttributesMsg(postAttributeMsg); builder.setPostAttributeScope(data.getAsJsonPrimitive("scope").getAsString());
} else { } catch (Exception e) {
builder.setAttributesUpdatedMsg(postAttributeMsg); log.warn("[{}] Can't convert to AttributesUpdatedMsg proto, entityData [{}]", entityId, entityData, e);
} }
break;
case POST_ATTRIBUTES:
try {
JsonObject data = entityData.getAsJsonObject();
TransportProtos.PostAttributeMsg postAttributesMsg = JsonConverter.convertToAttributesProto(data.getAsJsonObject("kv"));
builder.setPostAttributesMsg(postAttributesMsg);
builder.setPostAttributeScope(data.getAsJsonPrimitive("scope").getAsString()); builder.setPostAttributeScope(data.getAsJsonPrimitive("scope").getAsString());
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}] Can't convert to attributes proto, entityData [{}]", entityId, entityData, e); log.warn("[{}] Can't convert to PostAttributesMsg, entityData [{}]", entityId, entityData, e);
} }
break; break;
case ATTRIBUTES_DELETED: case ATTRIBUTES_DELETED:

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

@ -16,17 +16,22 @@
package org.thingsboard.server.edge; package org.thingsboard.server.edge;
import com.datastax.driver.core.utils.UUIDs; import com.datastax.driver.core.utils.UUIDs;
import com.fasterxml.jackson.core.JsonProcessingException;
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.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode; import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.gson.JsonObject; import com.google.gson.JsonObject;
import com.google.protobuf.AbstractMessage; import com.google.protobuf.AbstractMessage;
import com.google.protobuf.InvalidProtocolBufferException;
import com.google.protobuf.MessageLite;
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;
import org.junit.Before; 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.rule.engine.api.RuleEngineDeviceRpcRequest;
import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.Dashboard;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
@ -45,6 +50,8 @@ import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.id.UserId;
@ -53,9 +60,12 @@ import org.thingsboard.server.common.data.page.TimePageData;
import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup; 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.RuleChainMetaData;
import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.common.data.rule.RuleChainType;
import org.thingsboard.server.common.data.rule.RuleNode;
import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType;
import org.thingsboard.server.common.data.widget.WidgetType; import org.thingsboard.server.common.data.widget.WidgetType;
import org.thingsboard.server.common.data.widget.WidgetsBundle; import org.thingsboard.server.common.data.widget.WidgetsBundle;
import org.thingsboard.server.common.transport.adaptor.JsonConverter; import org.thingsboard.server.common.transport.adaptor.JsonConverter;
@ -65,15 +75,20 @@ import org.thingsboard.server.dao.util.mapping.JacksonUtil;
import org.thingsboard.server.edge.imitator.EdgeImitator; import org.thingsboard.server.edge.imitator.EdgeImitator;
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.AttributeDeleteMsg;
import org.thingsboard.server.gen.edge.AttributesRequestMsg;
import org.thingsboard.server.gen.edge.CustomerUpdateMsg; import org.thingsboard.server.gen.edge.CustomerUpdateMsg;
import org.thingsboard.server.gen.edge.DashboardUpdateMsg; import org.thingsboard.server.gen.edge.DashboardUpdateMsg;
import org.thingsboard.server.gen.edge.DeviceCredentialsRequestMsg; import org.thingsboard.server.gen.edge.DeviceCredentialsRequestMsg;
import org.thingsboard.server.gen.edge.DeviceCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.DeviceCredentialsUpdateMsg;
import org.thingsboard.server.gen.edge.DeviceRpcCallMsg;
import org.thingsboard.server.gen.edge.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.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.EntityDataProto;
import org.thingsboard.server.gen.edge.EntityViewUpdateMsg; import org.thingsboard.server.gen.edge.EntityViewUpdateMsg;
import org.thingsboard.server.gen.edge.RelationRequestMsg;
import org.thingsboard.server.gen.edge.RelationUpdateMsg; import org.thingsboard.server.gen.edge.RelationUpdateMsg;
import org.thingsboard.server.gen.edge.RpcResponseMsg;
import org.thingsboard.server.gen.edge.RuleChainMetadataRequestMsg; import org.thingsboard.server.gen.edge.RuleChainMetadataRequestMsg;
import org.thingsboard.server.gen.edge.RuleChainMetadataUpdateMsg; import org.thingsboard.server.gen.edge.RuleChainMetadataUpdateMsg;
import org.thingsboard.server.gen.edge.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.RuleChainUpdateMsg;
@ -85,16 +100,16 @@ import org.thingsboard.server.gen.edge.WidgetTypeUpdateMsg;
import org.thingsboard.server.gen.edge.WidgetsBundleUpdateMsg; import org.thingsboard.server.gen.edge.WidgetsBundleUpdateMsg;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Optional; import java.util.Optional;
import java.util.Random;
import java.util.UUID; import java.util.UUID;
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;
;
@Slf4j @Slf4j
abstract public class BaseEdgeTest extends AbstractControllerTest { abstract public class BaseEdgeTest extends AbstractControllerTest {
@ -137,7 +152,6 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
@After @After
public void afterTest() throws Exception { public void afterTest() throws Exception {
edgeImitator.disconnect(); edgeImitator.disconnect();
uninstallation();
loginSysAdmin(); loginSysAdmin();
@ -161,6 +175,69 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
testTimeseries(); testTimeseries();
testAttributes(); testAttributes();
testSendMessagesToCloud(); testSendMessagesToCloud();
testRpcCall();
}
private Device findDeviceByName(String deviceName) 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(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<TimePageData<Asset>>() {}, new TextPageLink(100)).getData();
Assert.assertEquals(1, edgeAssets.size());
Asset asset = edgeAssets.get(0);
Assert.assertEquals(assetName, asset.getName());
return asset;
}
private Device saveDevice(String deviceName) throws Exception {
Device device = new Device();
device.setName(deviceName);
device.setType("test");
return doPost("/api/device", device, Device.class);
}
private Asset saveAsset(String assetName) throws Exception {
Asset asset = new Asset();
asset.setName(assetName);
asset.setType("test");
return doPost("/api/asset", asset, Asset.class);
}
private void testRpcCall() throws Exception {
Device device = findDeviceByName("Edge Device 1");
RuleEngineDeviceRpcRequest request = RuleEngineDeviceRpcRequest.builder()
.oneway(true)
.method("test_method")
.body("{\"param1\":\"value1\"}")
.tenantId(device.getTenantId())
.deviceId(device.getId())
.requestId(new Random().nextInt())
.requestUUID(UUIDs.timeBased())
.originServiceId("originServiceId")
.expirationTime(System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10))
.restApiCall(true)
.build();
JsonNode body = mapper.valueToTree(request);
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body);
edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent);
edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof DeviceRpcCallMsg);
DeviceRpcCallMsg latestDeviceRpcCallMsg = (DeviceRpcCallMsg) latestMessage;
Assert.assertEquals("test_method", latestDeviceRpcCallMsg.getRequestMsg().getMethod());
} }
private void testReceivedInitialData() throws Exception { private void testReceivedInitialData() throws Exception {
@ -170,6 +247,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
EdgeConfiguration configuration = edgeImitator.getConfiguration(); EdgeConfiguration configuration = edgeImitator.getConfiguration();
Assert.assertNotNull(configuration); Assert.assertNotNull(configuration);
testAutoGeneratedCodeByProtobuf(configuration);
UserId userId = edgeImitator.getUserId(); UserId userId = edgeImitator.getUserId();
Assert.assertNotNull(userId); Assert.assertNotNull(userId);
@ -195,6 +274,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
new TypeReference<TimePageData<Asset>>() {}, new TextPageLink(100)).getData(); new TypeReference<TimePageData<Asset>>() {}, new TextPageLink(100)).getData();
Assert.assertTrue(edgeAssets.contains(asset)); Assert.assertTrue(edgeAssets.contains(asset));
testAutoGeneratedCodeByProtobuf(assetUpdateMsg);
Optional<RuleChainUpdateMsg> optionalMsg3 = edgeImitator.findMessageByType(RuleChainUpdateMsg.class); Optional<RuleChainUpdateMsg> optionalMsg3 = edgeImitator.findMessageByType(RuleChainUpdateMsg.class);
Assert.assertTrue(optionalMsg3.isPresent()); Assert.assertTrue(optionalMsg3.isPresent());
RuleChainUpdateMsg ruleChainUpdateMsg = optionalMsg3.get(); RuleChainUpdateMsg ruleChainUpdateMsg = optionalMsg3.get();
@ -206,16 +287,15 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
new TypeReference<TimePageData<RuleChain>>() {}, new TextPageLink(100)).getData(); new TypeReference<TimePageData<RuleChain>>() {}, new TextPageLink(100)).getData();
Assert.assertTrue(edgeRuleChains.contains(ruleChain)); Assert.assertTrue(edgeRuleChains.contains(ruleChain));
testAutoGeneratedCodeByProtobuf(ruleChainUpdateMsg);
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 savedDevice = saveDevice("Edge Device 2");
device.setName("Edge Device 2");
device.setType("test");
Device savedDevice = doPost("/api/device", device, Device.class);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
doPost("/api/edge/" + edge.getId().getId().toString() doPost("/api/edge/" + edge.getId().getId().toString()
@ -261,10 +341,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
private void testAssets() throws Exception { private void testAssets() throws Exception {
log.info("Testing assets"); log.info("Testing assets");
Asset asset = new Asset(); Asset savedAsset = saveAsset("Edge Asset 2");
asset.setName("Edge Asset 2");
asset.setType("test");
Asset savedAsset = doPost("/api/asset", asset, Asset.class);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
doPost("/api/edge/" + edge.getId().getId().toString() doPost("/api/edge/" + edge.getId().getId().toString()
@ -314,6 +391,11 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
ruleChain.setType(RuleChainType.EDGE); ruleChain.setType(RuleChainType.EDGE);
RuleChain savedRuleChain = doPost("/api/ruleChain", ruleChain, RuleChain.class); RuleChain savedRuleChain = doPost("/api/ruleChain", ruleChain, RuleChain.class);
createRuleChainMetadata(savedRuleChain);
// Wait before rule chain metadata saved to database before rule chain is assigned to edge
Thread.sleep(1000);
edgeImitator.expectMessageAmount(1); 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);
@ -327,6 +409,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
Assert.assertEquals(ruleChainUpdateMsg.getIdLSB(), savedRuleChain.getUuidId().getLeastSignificantBits()); Assert.assertEquals(ruleChainUpdateMsg.getIdLSB(), savedRuleChain.getUuidId().getLeastSignificantBits());
Assert.assertEquals(ruleChainUpdateMsg.getName(), savedRuleChain.getName()); Assert.assertEquals(ruleChainUpdateMsg.getName(), savedRuleChain.getName());
testRuleChainMetadataRequestMsg(savedRuleChain.getId());
edgeImitator.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);
@ -354,6 +438,67 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
log.info("RuleChains tested successfully"); log.info("RuleChains tested successfully");
} }
private void testRuleChainMetadataRequestMsg(RuleChainId ruleChainId) throws Exception {
RuleChainMetadataRequestMsg.Builder ruleChainMetadataRequestMsgBuilder = RuleChainMetadataRequestMsg.newBuilder()
.setRuleChainIdMSB(ruleChainId.getId().getMostSignificantBits())
.setRuleChainIdLSB(ruleChainId.getId().getLeastSignificantBits());
testAutoGeneratedCodeByProtobuf(ruleChainMetadataRequestMsgBuilder);
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder()
.addRuleChainMetadataRequestMsg(ruleChainMetadataRequestMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1);
edgeImitator.expectMessageAmount(1);
edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
edgeImitator.waitForResponses();
edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof RuleChainMetadataUpdateMsg);
RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = (RuleChainMetadataUpdateMsg) latestMessage;
RuleChainId receivedRuleChainId =
new RuleChainId(new UUID(ruleChainMetadataUpdateMsg.getRuleChainIdMSB(), ruleChainMetadataUpdateMsg.getRuleChainIdLSB()));
Assert.assertEquals(ruleChainId, receivedRuleChainId);
}
private void createRuleChainMetadata(RuleChain ruleChain) throws Exception {
RuleChainMetaData ruleChainMetaData = new RuleChainMetaData();
ruleChainMetaData.setRuleChainId(ruleChain.getId());
ObjectMapper mapper = new ObjectMapper();
RuleNode ruleNode1 = new RuleNode();
ruleNode1.setName("name1");
ruleNode1.setType("type1");
ruleNode1.setConfiguration(mapper.readTree("\"key1\": \"val1\""));
RuleNode ruleNode2 = new RuleNode();
ruleNode2.setName("name2");
ruleNode2.setType("type2");
ruleNode2.setConfiguration(mapper.readTree("\"key2\": \"val2\""));
RuleNode ruleNode3 = new RuleNode();
ruleNode3.setName("name3");
ruleNode3.setType("type3");
ruleNode3.setConfiguration(mapper.readTree("\"key3\": \"val3\""));
List<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 { private void testDashboards() throws Exception {
log.info("Testing Dashboards"); log.info("Testing Dashboards");
Dashboard dashboard = new Dashboard(); Dashboard dashboard = new Dashboard();
@ -373,6 +518,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
Assert.assertEquals(dashboardUpdateMsg.getIdLSB(), savedDashboard.getUuidId().getLeastSignificantBits()); Assert.assertEquals(dashboardUpdateMsg.getIdLSB(), savedDashboard.getUuidId().getLeastSignificantBits());
Assert.assertEquals(dashboardUpdateMsg.getTitle(), savedDashboard.getName()); Assert.assertEquals(dashboardUpdateMsg.getTitle(), savedDashboard.getName());
testAutoGeneratedCodeByProtobuf(dashboardUpdateMsg);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
savedDashboard.setTitle("Updated Edge Test Dashboard"); savedDashboard.setTitle("Updated Edge Test Dashboard");
doPost("/api/dashboard", savedDashboard, Dashboard.class); doPost("/api/dashboard", savedDashboard, Dashboard.class);
@ -413,18 +560,10 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
private void testRelations() throws Exception { private void testRelations() throws Exception {
log.info("Testing Relations"); log.info("Testing Relations");
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
List<Asset> edgeAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?",
new TypeReference<TimePageData<Asset>>() {}, new TextPageLink(100)).getData();
Assert.assertEquals(1, edgeDevices.size());
Assert.assertEquals(1, edgeAssets.size());
Device device = edgeDevices.get(0);
Asset asset = edgeAssets.get(0);
Assert.assertEquals("Edge Device 1", device.getName());
Assert.assertEquals("Edge Asset 1", asset.getName());
Device device = findDeviceByName("Edge Device 1");
Asset asset = findAssetByName("Edge Asset 1");
EntityRelation relation = new EntityRelation(); EntityRelation relation = new EntityRelation();
relation.setType("test"); relation.setType("test");
relation.setFrom(device.getId()); relation.setFrom(device.getId());
@ -476,11 +615,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
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?", Device device = findDeviceByName("Edge Device 1");
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());
Alarm alarm = new Alarm(); Alarm alarm = new Alarm();
alarm.setOriginator(device.getId()); alarm.setOriginator(device.getId());
@ -535,11 +670,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
private void testEntityView() throws Exception { private void testEntityView() throws Exception {
log.info("Testing EntityView"); log.info("Testing EntityView");
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", Device device = findDeviceByName("Edge Device 1");
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());
EntityView entityView = new EntityView(); EntityView entityView = new EntityView();
entityView.setName("Edge EntityView 1"); entityView.setName("Edge EntityView 1");
@ -611,6 +742,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
Assert.assertEquals(customerUpdateMsg.getIdLSB(), savedCustomer.getUuidId().getLeastSignificantBits()); Assert.assertEquals(customerUpdateMsg.getIdLSB(), savedCustomer.getUuidId().getLeastSignificantBits());
Assert.assertEquals(customerUpdateMsg.getTitle(), savedCustomer.getTitle()); Assert.assertEquals(customerUpdateMsg.getTitle(), savedCustomer.getTitle());
testAutoGeneratedCodeByProtobuf(customerUpdateMsg);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
doDelete("/api/customer/edge/" + edge.getId().getId().toString(), Edge.class); doDelete("/api/customer/edge/" + edge.getId().getId().toString(), Edge.class);
edgeImitator.waitForMessages(); edgeImitator.waitForMessages();
@ -656,6 +789,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
Assert.assertEquals(widgetsBundleUpdateMsg.getAlias(), savedWidgetsBundle.getAlias()); Assert.assertEquals(widgetsBundleUpdateMsg.getAlias(), savedWidgetsBundle.getAlias());
Assert.assertEquals(widgetsBundleUpdateMsg.getTitle(), savedWidgetsBundle.getTitle()); Assert.assertEquals(widgetsBundleUpdateMsg.getTitle(), savedWidgetsBundle.getTitle());
testAutoGeneratedCodeByProtobuf(widgetsBundleUpdateMsg);
WidgetType widgetType = new WidgetType(); WidgetType widgetType = new WidgetType();
widgetType.setName("Test Widget Type"); widgetType.setName("Test Widget Type");
widgetType.setBundleAlias(savedWidgetsBundle.getAlias()); widgetType.setBundleAlias(savedWidgetsBundle.getAlias());
@ -706,11 +841,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
private void testTimeseries() throws Exception { private void testTimeseries() throws Exception {
log.info("Testing timeseries"); log.info("Testing timeseries");
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", Device device = findDeviceByName("Edge Device 1");
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() + "}"; String timeseriesData = "{\"data\":{\"temperature\":25},\"ts\":" + System.currentTimeMillis() + "}";
JsonNode timeseriesEntityData = mapper.readTree(timeseriesData); JsonNode timeseriesEntityData = mapper.readTree(timeseriesData);
@ -740,61 +871,92 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
private void testAttributes() throws Exception { private void testAttributes() throws Exception {
log.info("Testing attributes"); log.info("Testing attributes");
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", Device device = findDeviceByName("Edge Device 1");
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\"}}"; testAttributesUpdatedMsg(device);
JsonNode attributesEntityData = mapper.readTree(attributesData); testPostAttributesMsg(device);
EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); testAttributesDeleteMsg(device);
log.info("Attributes tested successfully");
}
private void testAttributesDeleteMsg(Device device) throws JsonProcessingException, InterruptedException {
String deleteAttributesData = "{\"scope\":\"SERVER_SCOPE\",\"keys\":[\"key1\",\"key2\"]}";
JsonNode deleteAttributesEntityData = mapper.readTree(deleteAttributesData);
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_DELETED, device.getId().getId(), EdgeEventType.DEVICE, deleteAttributesEntityData);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent1); edgeEventService.saveAsync(edgeEvent);
edgeImitator.waitForMessages(); edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage(); AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof EntityDataProto); Assert.assertTrue(latestMessage instanceof EntityDataProto);
EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage; EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage;
Assert.assertEquals(latestEntityDataMsg.getEntityIdMSB(), device.getUuidId().getMostSignificantBits()); Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB());
Assert.assertEquals(latestEntityDataMsg.getEntityIdLSB(), device.getUuidId().getLeastSignificantBits()); Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB());
Assert.assertEquals(latestEntityDataMsg.getEntityType(), device.getId().getEntityType().name()); Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType());
Assert.assertEquals(latestEntityDataMsg.getPostAttributeScope(), attributesEntityData.get("scope").asText());
Assert.assertTrue(latestEntityDataMsg.hasAttributesUpdatedMsg());
TransportProtos.PostAttributeMsg attributesUpdatedMsg = latestEntityDataMsg.getAttributesUpdatedMsg(); Assert.assertTrue(latestEntityDataMsg.hasAttributeDeleteMsg());
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); AttributeDeleteMsg attributeDeleteMsg = latestEntityDataMsg.getAttributeDeleteMsg();
EdgeEvent edgeEvent2 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData); Assert.assertEquals(attributeDeleteMsg.getScope(), deleteAttributesEntityData.get("scope").asText());
Assert.assertEquals(2, attributeDeleteMsg.getAttributeNamesCount());
Assert.assertEquals("key1", attributeDeleteMsg.getAttributeNames(0));
Assert.assertEquals("key2", attributeDeleteMsg.getAttributeNames(1));
}
private void testPostAttributesMsg(Device device) throws JsonProcessingException, InterruptedException {
String postAttributesData = "{\"scope\":\"SERVER_SCOPE\",\"kv\":{\"key2\":\"value2\"}}";
JsonNode postAttributesEntityData = mapper.readTree(postAttributesData);
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.POST_ATTRIBUTES, device.getId().getId(), EdgeEventType.DEVICE, postAttributesEntityData);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent2); edgeEventService.saveAsync(edgeEvent);
edgeImitator.waitForMessages(); edgeImitator.waitForMessages();
latestMessage = edgeImitator.getLatestMessage(); AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof EntityDataProto); Assert.assertTrue(latestMessage instanceof EntityDataProto);
latestEntityDataMsg = (EntityDataProto) latestMessage; EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage;
Assert.assertEquals(latestEntityDataMsg.getEntityIdMSB(), device.getUuidId().getMostSignificantBits()); Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB());
Assert.assertEquals(latestEntityDataMsg.getEntityIdLSB(), device.getUuidId().getLeastSignificantBits()); Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB());
Assert.assertEquals(latestEntityDataMsg.getEntityType(), device.getId().getEntityType().name()); Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType());
Assert.assertEquals(latestEntityDataMsg.getPostAttributeScope(), attributesEntityData.get("scope").asText()); Assert.assertEquals("SERVER_SCOPE", latestEntityDataMsg.getPostAttributeScope());
Assert.assertTrue(latestEntityDataMsg.hasPostAttributesMsg()); Assert.assertTrue(latestEntityDataMsg.hasPostAttributesMsg());
attributesUpdatedMsg = latestEntityDataMsg.getPostAttributesMsg(); TransportProtos.PostAttributeMsg postAttributesMsg = latestEntityDataMsg.getPostAttributesMsg();
Assert.assertEquals(1, attributesUpdatedMsg.getKvCount()); Assert.assertEquals(1, postAttributesMsg.getKvCount());
keyValueProto = attributesUpdatedMsg.getKv(0); TransportProtos.KeyValueProto keyValueProto = postAttributesMsg.getKv(0);
Assert.assertEquals("key", keyValueProto.getKey()); Assert.assertEquals("key2", keyValueProto.getKey());
Assert.assertEquals("value", keyValueProto.getStringV()); Assert.assertEquals("value2", keyValueProto.getStringV());
}
log.info("Attributes tested successfully"); private void testAttributesUpdatedMsg(Device device) throws JsonProcessingException, InterruptedException {
String attributesData = "{\"scope\":\"SERVER_SCOPE\",\"kv\":{\"key1\":\"value1\"}}";
JsonNode attributesEntityData = mapper.readTree(attributesData);
EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData);
edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent1);
edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof EntityDataProto);
EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage;
Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB());
Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB());
Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType());
Assert.assertEquals("SERVER_SCOPE", latestEntityDataMsg.getPostAttributeScope());
Assert.assertTrue(latestEntityDataMsg.hasAttributesUpdatedMsg());
TransportProtos.PostAttributeMsg attributesUpdatedMsg = latestEntityDataMsg.getAttributesUpdatedMsg();
Assert.assertEquals(1, attributesUpdatedMsg.getKvCount());
TransportProtos.KeyValueProto keyValueProto = attributesUpdatedMsg.getKv(0);
Assert.assertEquals("key1", keyValueProto.getKey());
Assert.assertEquals("value1", keyValueProto.getStringV());
} }
private void testSendMessagesToCloud() throws Exception { private void testSendMessagesToCloud() throws Exception {
log.info("Sending messages to cloud"); log.info("Sending messages to cloud");
sendDevice(); sendDevice();
sendRelationRequest();
sendAlarm(); sendAlarm();
sendTelemetry(); sendTelemetry();
sendRelation(); sendRelation();
@ -802,22 +964,30 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
sendRuleChainMetadataRequest(); sendRuleChainMetadataRequest();
sendUserCredentialsRequest(); sendUserCredentialsRequest();
sendDeviceCredentialsRequest(); sendDeviceCredentialsRequest();
sendDeviceRpcResponse();
sendDeviceCredentialsUpdate();
sendAttributesRequest();
log.info("Messages were sent successfully"); log.info("Messages were sent successfully");
} }
private void sendDevice() throws Exception { private void sendDevice() throws Exception {
UUID uuid = UUIDs.timeBased(); UUID uuid = UUIDs.timeBased();
UplinkMsg.Builder builder = UplinkMsg.newBuilder(); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
DeviceUpdateMsg.Builder deviceUpdateMsgBuilder = DeviceUpdateMsg.newBuilder(); DeviceUpdateMsg.Builder deviceUpdateMsgBuilder = DeviceUpdateMsg.newBuilder();
deviceUpdateMsgBuilder.setIdMSB(uuid.getMostSignificantBits()); deviceUpdateMsgBuilder.setIdMSB(uuid.getMostSignificantBits());
deviceUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits()); deviceUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits());
deviceUpdateMsgBuilder.setName("Edge Device 2"); deviceUpdateMsgBuilder.setName("Edge Device 2");
deviceUpdateMsgBuilder.setType("test"); deviceUpdateMsgBuilder.setType("test");
deviceUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE); deviceUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE);
builder.addDeviceUpdateMsg(deviceUpdateMsgBuilder.build()); testAutoGeneratedCodeByProtobuf(deviceUpdateMsgBuilder);
uplinkMsgBuilder.addDeviceUpdateMsg(deviceUpdateMsgBuilder.build());
edgeImitator.expectResponsesAmount(1); edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(builder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
edgeImitator.waitForResponses(); edgeImitator.waitForResponses();
Device device = doGet("/api/device/" + uuid.toString(), Device.class); Device device = doGet("/api/device/" + uuid.toString(), Device.class);
@ -825,23 +995,70 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
Assert.assertEquals("Edge Device 2", device.getName()); Assert.assertEquals("Edge Device 2", device.getName());
} }
private void sendRelationRequest() throws Exception {
Device device = findDeviceByName("Edge Device 1");
Asset asset = findAssetByName("Edge Asset 1");
EntityRelation relation = new EntityRelation();
relation.setType("test");
relation.setFrom(device.getId());
relation.setTo(asset.getId());
relation.setTypeGroup(RelationTypeGroup.COMMON);
edgeImitator.expectMessageAmount(1);
doPost("/api/relation", relation);
edgeImitator.waitForMessages();
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
RelationRequestMsg.Builder relationRequestMsgBuilder = RelationRequestMsg.newBuilder();
relationRequestMsgBuilder.setEntityIdMSB(device.getId().getId().getMostSignificantBits());
relationRequestMsgBuilder.setEntityIdLSB(device.getId().getId().getLeastSignificantBits());
relationRequestMsgBuilder.setEntityType(device.getId().getEntityType().name());
testAutoGeneratedCodeByProtobuf(relationRequestMsgBuilder);
uplinkMsgBuilder.addRelationRequestMsg(relationRequestMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1);
edgeImitator.expectMessageAmount(1);
edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
edgeImitator.waitForResponses();
edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof RelationUpdateMsg);
RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType());
Assert.assertEquals(relation.getType(), relationUpdateMsg.getType());
UUID fromUUID = new UUID(relationUpdateMsg.getFromIdMSB(), relationUpdateMsg.getFromIdLSB());
EntityId fromEntityId = EntityIdFactory.getByTypeAndUuid(relationUpdateMsg.getFromEntityType(), fromUUID);
Assert.assertEquals(relation.getFrom(), fromEntityId);
UUID toUUID = new UUID(relationUpdateMsg.getToIdMSB(), relationUpdateMsg.getToIdLSB());
EntityId toEntityId = EntityIdFactory.getByTypeAndUuid(relationUpdateMsg.getToEntityType(), toUUID);
Assert.assertEquals(relation.getTo(), toEntityId);
Assert.assertEquals(relation.getTypeGroup().name(), relationUpdateMsg.getTypeGroup());
}
private void sendAlarm() throws Exception { private void sendAlarm() throws Exception {
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", Device device = findDeviceByName("Edge Device 2");
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(); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
AlarmUpdateMsg.Builder alarmUpdateMgBuilder = AlarmUpdateMsg.newBuilder(); AlarmUpdateMsg.Builder alarmUpdateMgBuilder = AlarmUpdateMsg.newBuilder();
alarmUpdateMgBuilder.setName("alarm from edge"); alarmUpdateMgBuilder.setName("alarm from edge");
alarmUpdateMgBuilder.setStatus(AlarmStatus.ACTIVE_UNACK.name()); alarmUpdateMgBuilder.setStatus(AlarmStatus.ACTIVE_UNACK.name());
alarmUpdateMgBuilder.setSeverity(AlarmSeverity.CRITICAL.name()); alarmUpdateMgBuilder.setSeverity(AlarmSeverity.CRITICAL.name());
alarmUpdateMgBuilder.setOriginatorName(device.getName()); alarmUpdateMgBuilder.setOriginatorName(device.getName());
alarmUpdateMgBuilder.setOriginatorType(EntityType.DEVICE.name()); alarmUpdateMgBuilder.setOriginatorType(EntityType.DEVICE.name());
builder.addAlarmUpdateMsg(alarmUpdateMgBuilder.build()); testAutoGeneratedCodeByProtobuf(alarmUpdateMgBuilder);
uplinkMsgBuilder.addAlarmUpdateMsg(alarmUpdateMgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1); edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(builder.build()); edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
edgeImitator.waitForResponses(); edgeImitator.waitForResponses();
@ -867,7 +1084,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
Assert.assertTrue(foundDevice2.isPresent()); Assert.assertTrue(foundDevice2.isPresent());
Device device2 = foundDevice2.get(); Device device2 = foundDevice2.get();
UplinkMsg.Builder builder = UplinkMsg.newBuilder(); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
RelationUpdateMsg.Builder relationUpdateMsgBuilder = RelationUpdateMsg.newBuilder(); RelationUpdateMsg.Builder relationUpdateMsgBuilder = RelationUpdateMsg.newBuilder();
relationUpdateMsgBuilder.setType("test"); relationUpdateMsgBuilder.setType("test");
relationUpdateMsgBuilder.setTypeGroup(RelationTypeGroup.COMMON.name()); relationUpdateMsgBuilder.setTypeGroup(RelationTypeGroup.COMMON.name());
@ -878,10 +1095,13 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
relationUpdateMsgBuilder.setFromIdLSB(device2.getId().getId().getLeastSignificantBits()); relationUpdateMsgBuilder.setFromIdLSB(device2.getId().getId().getLeastSignificantBits());
relationUpdateMsgBuilder.setFromEntityType(device2.getId().getEntityType().name()); relationUpdateMsgBuilder.setFromEntityType(device2.getId().getEntityType().name());
relationUpdateMsgBuilder.setAdditionalInfo("{}"); relationUpdateMsgBuilder.setAdditionalInfo("{}");
builder.addRelationUpdateMsg(relationUpdateMsgBuilder.build()); testAutoGeneratedCodeByProtobuf(relationUpdateMsgBuilder);
UplinkMsg msg = builder.build(); uplinkMsgBuilder.addRelationUpdateMsg(relationUpdateMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1); edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(msg); edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
edgeImitator.waitForResponses(); edgeImitator.waitForResponses();
EntityRelation relation = doGet("/api/relation?" + EntityRelation relation = doGet("/api/relation?" +
@ -907,28 +1127,35 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
String timeseriesKey = "key"; String timeseriesKey = "key";
String timeseriesValue = "25"; String timeseriesValue = "25";
data.addProperty(timeseriesKey, timeseriesValue); data.addProperty(timeseriesKey, timeseriesValue);
UplinkMsg.Builder builder1 = UplinkMsg.newBuilder(); UplinkMsg.Builder uplinkMsgBuilder1 = UplinkMsg.newBuilder();
EntityDataProto.Builder entityDataBuilder = EntityDataProto.newBuilder(); EntityDataProto.Builder entityDataBuilder = EntityDataProto.newBuilder();
entityDataBuilder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(data, System.currentTimeMillis())); entityDataBuilder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(data, System.currentTimeMillis()));
entityDataBuilder.setEntityType(device.getId().getEntityType().name()); entityDataBuilder.setEntityType(device.getId().getEntityType().name());
entityDataBuilder.setEntityIdMSB(device.getUuidId().getMostSignificantBits()); entityDataBuilder.setEntityIdMSB(device.getUuidId().getMostSignificantBits());
entityDataBuilder.setEntityIdLSB(device.getUuidId().getLeastSignificantBits()); entityDataBuilder.setEntityIdLSB(device.getUuidId().getLeastSignificantBits());
builder1.addEntityData(entityDataBuilder.build()); testAutoGeneratedCodeByProtobuf(entityDataBuilder);
edgeImitator.sendUplinkMsg(builder1.build()); uplinkMsgBuilder1.addEntityData(entityDataBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder1);
edgeImitator.sendUplinkMsg(uplinkMsgBuilder1.build());
JsonObject attributesData = new JsonObject(); JsonObject attributesData = new JsonObject();
String attributesKey = "test_attr"; String attributesKey = "test_attr";
String attributesValue = "test_value"; String attributesValue = "test_value";
attributesData.addProperty(attributesKey, attributesValue); attributesData.addProperty(attributesKey, attributesValue);
UplinkMsg.Builder builder2 = UplinkMsg.newBuilder(); UplinkMsg.Builder uplinkMsgBuilder2 = UplinkMsg.newBuilder();
EntityDataProto.Builder entityDataBuilder2 = EntityDataProto.newBuilder(); EntityDataProto.Builder entityDataBuilder2 = EntityDataProto.newBuilder();
entityDataBuilder2.setEntityType(device.getId().getEntityType().name()); entityDataBuilder2.setEntityType(device.getId().getEntityType().name());
entityDataBuilder2.setEntityIdMSB(device.getId().getId().getMostSignificantBits()); entityDataBuilder2.setEntityIdMSB(device.getId().getId().getMostSignificantBits());
entityDataBuilder2.setEntityIdLSB(device.getId().getId().getLeastSignificantBits()); entityDataBuilder2.setEntityIdLSB(device.getId().getId().getLeastSignificantBits());
entityDataBuilder2.setAttributesUpdatedMsg(JsonConverter.convertToAttributesProto(attributesData)); entityDataBuilder2.setAttributesUpdatedMsg(JsonConverter.convertToAttributesProto(attributesData));
entityDataBuilder2.setPostAttributeScope(DataConstants.SERVER_SCOPE); entityDataBuilder2.setPostAttributeScope(DataConstants.SERVER_SCOPE);
builder2.addEntityData(entityDataBuilder2.build()); testAutoGeneratedCodeByProtobuf(entityDataBuilder2);
edgeImitator.sendUplinkMsg(builder2.build());
uplinkMsgBuilder2.addEntityData(entityDataBuilder2.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder2);
edgeImitator.sendUplinkMsg(uplinkMsgBuilder2.build());
edgeImitator.waitForResponses(); edgeImitator.waitForResponses();
Thread.sleep(1000); Thread.sleep(1000);
@ -947,14 +1174,18 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
private void sendRuleChainMetadataRequest() throws Exception { private void sendRuleChainMetadataRequest() throws Exception {
RuleChainId edgeRootRuleChainId = edge.getRootRuleChainId(); RuleChainId edgeRootRuleChainId = edge.getRootRuleChainId();
UplinkMsg.Builder builder = UplinkMsg.newBuilder(); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
RuleChainMetadataRequestMsg.Builder ruleChainMetadataRequestMsgBuilder = RuleChainMetadataRequestMsg.newBuilder(); RuleChainMetadataRequestMsg.Builder ruleChainMetadataRequestMsgBuilder = RuleChainMetadataRequestMsg.newBuilder();
ruleChainMetadataRequestMsgBuilder.setRuleChainIdMSB(edgeRootRuleChainId.getId().getMostSignificantBits()); ruleChainMetadataRequestMsgBuilder.setRuleChainIdMSB(edgeRootRuleChainId.getId().getMostSignificantBits());
ruleChainMetadataRequestMsgBuilder.setRuleChainIdLSB(edgeRootRuleChainId.getId().getLeastSignificantBits()); ruleChainMetadataRequestMsgBuilder.setRuleChainIdLSB(edgeRootRuleChainId.getId().getLeastSignificantBits());
builder.addRuleChainMetadataRequestMsg(ruleChainMetadataRequestMsgBuilder.build()); testAutoGeneratedCodeByProtobuf(ruleChainMetadataRequestMsgBuilder);
uplinkMsgBuilder.addRuleChainMetadataRequestMsg(ruleChainMetadataRequestMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1); edgeImitator.expectResponsesAmount(1);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
edgeImitator.sendUplinkMsg(builder.build()); edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
edgeImitator.waitForResponses(); edgeImitator.waitForResponses();
edgeImitator.waitForMessages(); edgeImitator.waitForMessages();
@ -963,19 +1194,25 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = (RuleChainMetadataUpdateMsg) latestMessage; RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = (RuleChainMetadataUpdateMsg) latestMessage;
Assert.assertEquals(ruleChainMetadataUpdateMsg.getRuleChainIdMSB(), edgeRootRuleChainId.getId().getMostSignificantBits()); Assert.assertEquals(ruleChainMetadataUpdateMsg.getRuleChainIdMSB(), edgeRootRuleChainId.getId().getMostSignificantBits());
Assert.assertEquals(ruleChainMetadataUpdateMsg.getRuleChainIdLSB(), edgeRootRuleChainId.getId().getLeastSignificantBits()); Assert.assertEquals(ruleChainMetadataUpdateMsg.getRuleChainIdLSB(), edgeRootRuleChainId.getId().getLeastSignificantBits());
testAutoGeneratedCodeByProtobuf(ruleChainMetadataUpdateMsg);
} }
private void sendUserCredentialsRequest() throws Exception { private void sendUserCredentialsRequest() throws Exception {
UserId userId = edgeImitator.getUserId(); UserId userId = edgeImitator.getUserId();
UplinkMsg.Builder builder = UplinkMsg.newBuilder(); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
UserCredentialsRequestMsg.Builder userCredentialsRequestMsgBuilder = UserCredentialsRequestMsg.newBuilder(); UserCredentialsRequestMsg.Builder userCredentialsRequestMsgBuilder = UserCredentialsRequestMsg.newBuilder();
userCredentialsRequestMsgBuilder.setUserIdMSB(userId.getId().getMostSignificantBits()); userCredentialsRequestMsgBuilder.setUserIdMSB(userId.getId().getMostSignificantBits());
userCredentialsRequestMsgBuilder.setUserIdLSB(userId.getId().getLeastSignificantBits()); userCredentialsRequestMsgBuilder.setUserIdLSB(userId.getId().getLeastSignificantBits());
builder.addUserCredentialsRequestMsg(userCredentialsRequestMsgBuilder.build()); testAutoGeneratedCodeByProtobuf(userCredentialsRequestMsgBuilder);
uplinkMsgBuilder.addUserCredentialsRequestMsg(userCredentialsRequestMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1); edgeImitator.expectResponsesAmount(1);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
edgeImitator.sendUplinkMsg(builder.build()); edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
edgeImitator.waitForResponses(); edgeImitator.waitForResponses();
edgeImitator.waitForMessages(); edgeImitator.waitForMessages();
@ -984,26 +1221,27 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
UserCredentialsUpdateMsg userCredentialsUpdateMsg = (UserCredentialsUpdateMsg) latestMessage; UserCredentialsUpdateMsg userCredentialsUpdateMsg = (UserCredentialsUpdateMsg) latestMessage;
Assert.assertEquals(userCredentialsUpdateMsg.getUserIdMSB(), userId.getId().getMostSignificantBits()); Assert.assertEquals(userCredentialsUpdateMsg.getUserIdMSB(), userId.getId().getMostSignificantBits());
Assert.assertEquals(userCredentialsUpdateMsg.getUserIdLSB(), userId.getId().getLeastSignificantBits()); Assert.assertEquals(userCredentialsUpdateMsg.getUserIdLSB(), userId.getId().getLeastSignificantBits());
testAutoGeneratedCodeByProtobuf(userCredentialsUpdateMsg);
} }
private void sendDeviceCredentialsRequest() throws Exception { private void sendDeviceCredentialsRequest() throws Exception {
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", Device device = findDeviceByName("Edge Device 1");
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Optional<Device> foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 1")).findAny();
Assert.assertTrue(foundDevice.isPresent());
Device device = foundDevice.get();
DeviceCredentials deviceCredentials = doGet("/api/device/" + device.getId().getId().toString() + "/credentials", DeviceCredentials.class); DeviceCredentials deviceCredentials = doGet("/api/device/" + device.getId().getId().toString() + "/credentials", DeviceCredentials.class);
UplinkMsg.Builder builder = UplinkMsg.newBuilder(); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
DeviceCredentialsRequestMsg.Builder deviceCredentialsRequestMsgBuilder = DeviceCredentialsRequestMsg.newBuilder(); DeviceCredentialsRequestMsg.Builder deviceCredentialsRequestMsgBuilder = DeviceCredentialsRequestMsg.newBuilder();
deviceCredentialsRequestMsgBuilder.setDeviceIdMSB(device.getUuidId().getMostSignificantBits()); deviceCredentialsRequestMsgBuilder.setDeviceIdMSB(device.getUuidId().getMostSignificantBits());
deviceCredentialsRequestMsgBuilder.setDeviceIdLSB(device.getUuidId().getLeastSignificantBits()); deviceCredentialsRequestMsgBuilder.setDeviceIdLSB(device.getUuidId().getLeastSignificantBits());
builder.addDeviceCredentialsRequestMsg(deviceCredentialsRequestMsgBuilder.build()); testAutoGeneratedCodeByProtobuf(deviceCredentialsRequestMsgBuilder);
uplinkMsgBuilder.addDeviceCredentialsRequestMsg(deviceCredentialsRequestMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1); edgeImitator.expectResponsesAmount(1);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
edgeImitator.sendUplinkMsg(builder.build()); edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
edgeImitator.waitForResponses(); edgeImitator.waitForResponses();
edgeImitator.waitForMessages(); edgeImitator.waitForMessages();
@ -1016,20 +1254,108 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
Assert.assertEquals(deviceCredentialsUpdateMsg.getCredentialsId(), deviceCredentials.getCredentialsId()); Assert.assertEquals(deviceCredentialsUpdateMsg.getCredentialsId(), deviceCredentials.getCredentialsId());
} }
private void sendDeviceCredentialsUpdate() throws Exception {
Device device = findDeviceByName("Edge Device 1");
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
DeviceCredentialsUpdateMsg.Builder deviceCredentialsUpdateMsgBuilder = DeviceCredentialsUpdateMsg.newBuilder();
deviceCredentialsUpdateMsgBuilder.setDeviceIdMSB(device.getUuidId().getMostSignificantBits());
deviceCredentialsUpdateMsgBuilder.setDeviceIdLSB(device.getUuidId().getLeastSignificantBits());
deviceCredentialsUpdateMsgBuilder.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN.name());
deviceCredentialsUpdateMsgBuilder.setCredentialsId("NEW_TOKEN");
testAutoGeneratedCodeByProtobuf(deviceCredentialsUpdateMsgBuilder);
uplinkMsgBuilder.addDeviceCredentialsUpdateMsg(deviceCredentialsUpdateMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
edgeImitator.waitForResponses();
}
private void sendDeviceRpcResponse() throws Exception {
Device device = findDeviceByName("Edge Device 1");
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
DeviceRpcCallMsg.Builder deviceRpcCallResponseBuilder = DeviceRpcCallMsg.newBuilder();
deviceRpcCallResponseBuilder.setDeviceIdMSB(device.getUuidId().getMostSignificantBits());
deviceRpcCallResponseBuilder.setDeviceIdLSB(device.getUuidId().getLeastSignificantBits());
deviceRpcCallResponseBuilder.setOneway(true);
deviceRpcCallResponseBuilder.setOriginServiceId("originServiceId");
deviceRpcCallResponseBuilder.setExpirationTime(System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10));
RpcResponseMsg.Builder responseBuilder =
RpcResponseMsg.newBuilder().setResponse("{}");
testAutoGeneratedCodeByProtobuf(responseBuilder);
deviceRpcCallResponseBuilder.setResponseMsg(responseBuilder.build());
testAutoGeneratedCodeByProtobuf(deviceRpcCallResponseBuilder);
uplinkMsgBuilder.addDeviceRpcCallMsg(deviceRpcCallResponseBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
edgeImitator.waitForResponses();
}
private void sendAttributesRequest() throws Exception {
Device device = findDeviceByName("Edge Device 1");
String attributesDataStr = "{\"key1\":\"value1\"}";
JsonNode attributesData = mapper.readTree(attributesDataStr);
doPost("/api/plugins/telemetry/DEVICE/" + device.getId().getId().toString() + "/attributes/" + DataConstants.SERVER_SCOPE,
attributesData);
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
AttributesRequestMsg.Builder attributesRequestMsgBuilder = AttributesRequestMsg.newBuilder();
attributesRequestMsgBuilder.setEntityIdMSB(device.getUuidId().getMostSignificantBits());
attributesRequestMsgBuilder.setEntityIdLSB(device.getUuidId().getLeastSignificantBits());
attributesRequestMsgBuilder.setEntityType(EntityType.DEVICE.name());
testAutoGeneratedCodeByProtobuf(attributesRequestMsgBuilder);
uplinkMsgBuilder.addAttributesRequestMsg(attributesRequestMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder);
edgeImitator.expectResponsesAmount(1);
edgeImitator.expectMessageAmount(1);
edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build());
edgeImitator.waitForResponses();
edgeImitator.waitForMessages();
AbstractMessage latestMessage = edgeImitator.getLatestMessage();
Assert.assertTrue(latestMessage instanceof EntityDataProto);
EntityDataProto latestEntityDataMsg = (EntityDataProto) latestMessage;
Assert.assertEquals(device.getUuidId().getMostSignificantBits(), latestEntityDataMsg.getEntityIdMSB());
Assert.assertEquals(device.getUuidId().getLeastSignificantBits(), latestEntityDataMsg.getEntityIdLSB());
Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType());
Assert.assertEquals("SERVER_SCOPE", latestEntityDataMsg.getPostAttributeScope());
Assert.assertTrue(latestEntityDataMsg.hasAttributesUpdatedMsg());
TransportProtos.PostAttributeMsg attributesUpdatedMsg = latestEntityDataMsg.getAttributesUpdatedMsg();
Assert.assertEquals(1, attributesUpdatedMsg.getKvCount());
TransportProtos.KeyValueProto keyValueProto = attributesUpdatedMsg.getKv(0);
Assert.assertEquals("key1", keyValueProto.getKey());
Assert.assertEquals("value1", keyValueProto.getStringV());
}
private void sendDeleteDeviceOnEdge() throws Exception { private void sendDeleteDeviceOnEdge() throws Exception {
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData(); new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Optional<Device> foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 2")).findAny(); Optional<Device> foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 2")).findAny();
Assert.assertTrue(foundDevice.isPresent()); Assert.assertTrue(foundDevice.isPresent());
Device device = foundDevice.get(); Device device = foundDevice.get();
UplinkMsg.Builder builder = UplinkMsg.newBuilder(); UplinkMsg.Builder upLinkMsgBuilder = UplinkMsg.newBuilder();
DeviceUpdateMsg.Builder deviceDeleteMsgBuilder = DeviceUpdateMsg.newBuilder(); DeviceUpdateMsg.Builder deviceDeleteMsgBuilder = DeviceUpdateMsg.newBuilder();
deviceDeleteMsgBuilder.setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE); deviceDeleteMsgBuilder.setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE);
deviceDeleteMsgBuilder.setIdMSB(device.getId().getId().getMostSignificantBits()); deviceDeleteMsgBuilder.setIdMSB(device.getId().getId().getMostSignificantBits());
deviceDeleteMsgBuilder.setIdLSB(device.getId().getId().getLeastSignificantBits()); deviceDeleteMsgBuilder.setIdLSB(device.getId().getId().getLeastSignificantBits());
builder.addDeviceUpdateMsg(deviceDeleteMsgBuilder.build()); testAutoGeneratedCodeByProtobuf(deviceDeleteMsgBuilder);
upLinkMsgBuilder.addDeviceUpdateMsg(deviceDeleteMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(upLinkMsgBuilder);
edgeImitator.expectResponsesAmount(1); edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(builder.build()); edgeImitator.sendUplinkMsg(upLinkMsgBuilder.build());
edgeImitator.waitForResponses(); edgeImitator.waitForResponses();
device = doGet("/api/device/" + device.getId().getId().toString(), Device.class); device = doGet("/api/device/" + device.getId().getId().toString(), Device.class);
Assert.assertNotNull(device); Assert.assertNotNull(device);
@ -1042,41 +1368,15 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
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);
Device device = new Device(); Device savedDevice = saveDevice("Edge Device 1");
device.setName("Edge Device 1");
device.setType("test");
Device savedDevice = doPost("/api/device", device, Device.class);
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);
Asset asset = new Asset(); Asset savedAsset = saveAsset("Edge Asset 1");
asset.setName("Edge Asset 1");
asset.setType("test");
Asset savedAsset = doPost("/api/asset", asset, Asset.class);
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);
} }
private void uninstallation() throws Exception {
TimePageData<Device> pageDataDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100));
for (Device device: pageDataDevices.getData()) {
doDelete("/api/device/" + device.getId().getId().toString())
.andExpect(status().isOk());
}
TimePageData<Asset> pageDataAssets = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/assets?",
new TypeReference<TimePageData<Asset>>() {}, new TextPageLink(100));
for (Asset asset: pageDataAssets.getData()) {
doDelete("/api/asset/" + asset.getId().getId().toString())
.andExpect(status().isOk());
}
doDelete("/api/edge/" + edge.getId().getId().toString())
.andExpect(status().isOk());
}
private EdgeEvent constructEdgeEvent(TenantId tenantId, EdgeId edgeId, EdgeEventActionType edgeEventAction, UUID entityId, EdgeEventType edgeEventType, JsonNode entityBody) { private EdgeEvent constructEdgeEvent(TenantId tenantId, EdgeId edgeId, EdgeEventActionType edgeEventAction, UUID entityId, EdgeEventType edgeEventType, JsonNode entityBody) {
EdgeEvent edgeEvent = new EdgeEvent(); EdgeEvent edgeEvent = new EdgeEvent();
edgeEvent.setEdgeId(edgeId); edgeEvent.setEdgeId(edgeId);
@ -1087,4 +1387,19 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
edgeEvent.setBody(entityBody); edgeEvent.setBody(entityBody);
return edgeEvent; return edgeEvent;
} }
private void testAutoGeneratedCodeByProtobuf(MessageLite.Builder builder) throws InvalidProtocolBufferException {
MessageLite source = builder.build();
testAutoGeneratedCodeByProtobuf(source);
MessageLite target = source.getParserForType().parseFrom(source.toByteArray());
builder.clear().mergeFrom(target);
}
private void testAutoGeneratedCodeByProtobuf(MessageLite source) throws InvalidProtocolBufferException {
MessageLite target = source.getParserForType().parseFrom(source.toByteArray());
Assert.assertEquals(source, target);
Assert.assertEquals(source.hashCode(), target.hashCode());
}
} }

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

@ -30,7 +30,9 @@ import org.thingsboard.server.gen.edge.AlarmUpdateMsg;
import org.thingsboard.server.gen.edge.AssetUpdateMsg; import org.thingsboard.server.gen.edge.AssetUpdateMsg;
import org.thingsboard.server.gen.edge.CustomerUpdateMsg; import org.thingsboard.server.gen.edge.CustomerUpdateMsg;
import org.thingsboard.server.gen.edge.DashboardUpdateMsg; import org.thingsboard.server.gen.edge.DashboardUpdateMsg;
import org.thingsboard.server.gen.edge.DeviceCredentialsRequestMsg;
import org.thingsboard.server.gen.edge.DeviceCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.DeviceCredentialsUpdateMsg;
import org.thingsboard.server.gen.edge.DeviceRpcCallMsg;
import org.thingsboard.server.gen.edge.DeviceUpdateMsg; 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;
@ -224,6 +226,16 @@ public class EdgeImitator {
result.add(saveDownlinkMsg(userCredentialsUpdateMsg)); result.add(saveDownlinkMsg(userCredentialsUpdateMsg));
} }
} }
if (downlinkMsg.getDeviceRpcCallMsgList() != null && !downlinkMsg.getDeviceRpcCallMsgList().isEmpty()) {
for (DeviceRpcCallMsg deviceRpcCallMsg: downlinkMsg.getDeviceRpcCallMsgList()) {
result.add(saveDownlinkMsg(deviceRpcCallMsg));
}
}
if (downlinkMsg.getDeviceCredentialsRequestMsgList() != null && !downlinkMsg.getDeviceCredentialsRequestMsgList().isEmpty()) {
for (DeviceCredentialsRequestMsg deviceCredentialsRequestMsg: downlinkMsg.getDeviceCredentialsRequestMsgList()) {
result.add(saveDownlinkMsg(deviceCredentialsRequestMsg));
}
}
return Futures.allAsList(result); return Futures.allAsList(result);
} }

1
common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventActionType.java

@ -19,6 +19,7 @@ public enum EdgeEventActionType {
ADDED, ADDED,
DELETED, DELETED,
UPDATED, UPDATED,
POST_ATTRIBUTES,
ATTRIBUTES_UPDATED, ATTRIBUTES_UPDATED,
ATTRIBUTES_DELETED, ATTRIBUTES_DELETED,
TIMESERIES_UPDATED, TIMESERIES_UPDATED,

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java

@ -152,11 +152,9 @@ public class TbMsgPushToEdgeNode implements TbNode {
JsonNode dataJson = json.readTree(msg.getData()); JsonNode dataJson = json.readTree(msg.getData());
switch (actionType) { switch (actionType) {
case ATTRIBUTES_UPDATED: case ATTRIBUTES_UPDATED:
case POST_ATTRIBUTES:
entityBody.put("kv", dataJson); entityBody.put("kv", dataJson);
entityBody.put("scope", metadata.get("scope")); entityBody.put("scope", metadata.get("scope"));
if (SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType)) {
entityBody.put("isPostAttributes", true);
}
break; break;
case ATTRIBUTES_DELETED: case ATTRIBUTES_DELETED:
List<String> keys = json.treeToValue(dataJson.get("attributes"), List.class); List<String> keys = json.treeToValue(dataJson.get("attributes"), List.class);
@ -192,9 +190,10 @@ public class TbMsgPushToEdgeNode implements TbNode {
EdgeEventActionType actionType; EdgeEventActionType actionType;
if (SessionMsgType.POST_TELEMETRY_REQUEST.name().equals(msgType)) { if (SessionMsgType.POST_TELEMETRY_REQUEST.name().equals(msgType)) {
actionType = EdgeEventActionType.TIMESERIES_UPDATED; actionType = EdgeEventActionType.TIMESERIES_UPDATED;
} else if (SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType) } else if (DataConstants.ATTRIBUTES_UPDATED.equals(msgType)) {
|| DataConstants.ATTRIBUTES_UPDATED.equals(msgType)) {
actionType = EdgeEventActionType.ATTRIBUTES_UPDATED; actionType = EdgeEventActionType.ATTRIBUTES_UPDATED;
} else if (SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType)) {
actionType = EdgeEventActionType.POST_ATTRIBUTES;
} else { } else {
actionType = EdgeEventActionType.ATTRIBUTES_DELETED; actionType = EdgeEventActionType.ATTRIBUTES_DELETED;
} }

30
ui/src/app/event/event-table.directive.js

@ -220,39 +220,37 @@ export default function EventTableDirective($compile, $templateCache, $rootScope
} }
scope.subscriptionId = null; scope.subscriptionId = null;
scope.queueStartTs; scope.queueStartTs = 0;
scope.loadEdgeInfo = function() { scope.loadEdgeInfo = function() {
attributeService.getEntityAttributesValues(scope.entityType, scope.entityId, types.attributesScope.server.value, attributeService.getEntityAttributesValues(
types.edgeAttributeKeys.queueStartTs, {}) scope.entityType,
scope.entityId,
types.attributesScope.server.value,
types.edgeAttributeKeys.queueStartTs,
{})
.then(function success(attributes) { .then(function success(attributes) {
scope.onUpdate(attributes); scope.onEdgeAttributesUpdate(attributes);
}); });
scope.checkSubscription(); scope.checkSubscription();
attributeService.getEntityAttributes(scope.entityType, scope.entityId, types.attributesScope.server.value, {order: '', limit: 1, page: 1, search: ''},
function (attributes) {
if (attributes && attributes.data) {
scope.onUpdate(attributes.data);
}
});
} }
scope.onUpdate = function(attributes) { scope.onEdgeAttributesUpdate = function(attributes) {
let edge = attributes.reduce(function (map, attribute) { let edgeAttributes = attributes.reduce(function (map, attribute) {
map[attribute.key] = attribute; map[attribute.key] = attribute;
return map; return map;
}, {}); }, {});
if (edge.queueStartTs) { if (edgeAttributes.queueStartTs) {
scope.queueStartTs = edge.queueStartTs.lastUpdateTs; scope.queueStartTs = edgeAttributes.queueStartTs.lastUpdateTs;
} }
} }
scope.checkSubscription = function() { scope.checkSubscription = function() {
var newSubscriptionId = null; var newSubscriptionId = null;
if (scope.entityId && scope.entityType && types.attributesScope.server.value) { if (scope.entityId && scope.entityType && types.attributesScope.server.value) {
newSubscriptionId = attributeService.subscribeForEntityAttributes(scope.entityType, scope.entityId, types.attributesScope.server.value); newSubscriptionId =
attributeService.subscribeForEntityAttributes(scope.entityType, scope.entityId, types.attributesScope.server.value);
} }
if (scope.subscriptionId && scope.subscriptionId != newSubscriptionId) { if (scope.subscriptionId && scope.subscriptionId != newSubscriptionId) {
attributeService.unsubscribeForEntityAttributes(scope.subscriptionId); attributeService.unsubscribeForEntityAttributes(scope.subscriptionId);

Loading…
Cancel
Save