From 492a14221b40aafc8af0523b3c0b6e304dcb775c Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 10 Jun 2020 19:07:16 +0300 Subject: [PATCH] Handling device creation from edge to cloud --- .../service/edge/EdgeContextComponent.java | 5 + .../service/edge/rpc/EdgeGrpcSession.java | 179 +++++++++++------- common/edge-api/src/main/proto/edge.proto | 13 +- 3 files changed, 118 insertions(+), 79 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index 1e51106717..f936ee82a5 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -25,6 +25,7 @@ import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.customer.CustomerService; import org.thingsboard.server.dao.dashboard.DashboardService; +import org.thingsboard.server.dao.device.DeviceCredentialsService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entityview.EntityViewService; @@ -57,6 +58,10 @@ public class EdgeContextComponent { @Autowired private DeviceService deviceService; + @Lazy + @Autowired + private DeviceCredentialsService deviceCredentialsService; + @Lazy @Autowired private EntityViewService entityViewService; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index ea1aad462c..e6b68f486a 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -31,6 +31,7 @@ import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.Event; import org.thingsboard.server.common.data.User; @@ -56,6 +57,8 @@ import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; +import org.thingsboard.server.common.data.security.DeviceCredentials; +import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgMetaData; @@ -99,7 +102,7 @@ import static org.thingsboard.server.gen.edge.UpdateMsgType.ENTITY_CREATED_RPC_M @Data public final class EdgeGrpcSession implements Closeable { - private static final ReentrantLock entityCreationLock = new ReentrantLock(); + private static final ReentrantLock deviceCreationLock = new ReentrantLock(); private static final String QUEUE_START_TS_ATTR_KEY = "queueStartTs"; @@ -518,32 +521,15 @@ public final class EdgeGrpcSession implements Closeable { for (EntityDataProto entityData : uplinkMsg.getEntityDataList()) { TbMsg tbMsg = null; TbMsg originalTbMsg = TbMsg.fromBytes(entityData.getTbMsg().toByteArray(), TbMsgCallback.EMPTY); - switch (originalTbMsg.getOriginator().getEntityType()) { - case DEVICE: - String deviceName = entityData.getEntityName(); - String deviceType = entityData.getEntityType(); - Device device = getOrCreateDevice(deviceName, deviceType); - if (device != null) { - tbMsg = TbMsg.newMsg(originalTbMsg.getType(), device.getId(), originalTbMsg.getMetaData().copy(), - originalTbMsg.getDataType(), originalTbMsg.getData()); - } - break; - case ASSET: - String assetName = entityData.getEntityName(); - Asset asset = ctx.getAssetService().findAssetByTenantIdAndName(edge.getTenantId(), assetName); - if (asset != null) { - tbMsg = TbMsg.newMsg(originalTbMsg.getType(), asset.getId(), originalTbMsg.getMetaData().copy(), - originalTbMsg.getDataType(), originalTbMsg.getData()); - } - break; - case ENTITY_VIEW: - String entityViewName = entityData.getEntityName(); - EntityView entityView = ctx.getEntityViewService().findEntityViewByTenantIdAndName(edge.getTenantId(), entityViewName); - if (entityView != null) { - tbMsg = TbMsg.newMsg(originalTbMsg.getType(), entityView.getId(), originalTbMsg.getMetaData().copy(), - originalTbMsg.getDataType(), originalTbMsg.getData()); - } - break; + if (originalTbMsg.getOriginator().getEntityType() == EntityType.DEVICE) { + String deviceName = entityData.getEntityName(); + Device device = ctx.getDeviceService().findDeviceByTenantIdAndName(edge.getTenantId(), deviceName); + if (device != null) { + tbMsg = TbMsg.newMsg(originalTbMsg.getType(), device.getId(), originalTbMsg.getMetaData().copy(), + originalTbMsg.getDataType(), originalTbMsg.getData()); + } + } else { + tbMsg = originalTbMsg; } if (tbMsg != null) { ctx.getTbClusterService().pushMsgToRuleEngine(edge.getTenantId(), tbMsg.getOriginator(), tbMsg, null); @@ -552,19 +538,7 @@ public final class EdgeGrpcSession implements Closeable { } if (uplinkMsg.getDeviceUpdateMsgList() != null && !uplinkMsg.getDeviceUpdateMsgList().isEmpty()) { for (DeviceUpdateMsg deviceUpdateMsg : uplinkMsg.getDeviceUpdateMsgList()) { - String deviceName = deviceUpdateMsg.getName(); - String deviceType = deviceUpdateMsg.getType(); - switch (deviceUpdateMsg.getMsgType()) { - case ENTITY_CREATED_RPC_MESSAGE: - getOrCreateDevice(deviceName, deviceType); - break; - case ENTITY_DELETED_RPC_MESSAGE: - Device device = ctx.getDeviceService().findDeviceByTenantIdAndName(edge.getTenantId(), deviceName); - if (device != null) { - ctx.getDeviceService().unassignDeviceFromEdge(edge.getTenantId(), device.getId(), edge.getId()); - } - break; - } + onDeviceUpdate(deviceUpdateMsg); } } if (uplinkMsg.getAlarmUpdateMsgList() != null && !uplinkMsg.getAlarmUpdateMsgList().isEmpty()) { @@ -584,6 +558,101 @@ public final class EdgeGrpcSession implements Closeable { return UplinkResponseMsg.newBuilder().setSuccess(true).build(); } + private void onDeviceUpdate(DeviceUpdateMsg deviceUpdateMsg) { + log.info("onDeviceUpdate {}", deviceUpdateMsg); + DeviceId edgeDeviceId = new DeviceId(new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB())); + switch (deviceUpdateMsg.getMsgType()) { + case ENTITY_CREATED_RPC_MESSAGE: + String deviceName = deviceUpdateMsg.getName(); + Device device = ctx.getDeviceService().findDeviceByTenantIdAndName(edge.getTenantId(), deviceName); + if (device != null) { + // device with this name already exists on the cloud - update ID on the edge + if (!device.getId().equals(edgeDeviceId)) { + EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() + .setDeviceUpdateMsg(ctx.getDeviceUpdateMsgConstructor().constructDeviceUpdatedMsg(UpdateMsgType.DEVICE_CONFLICT_RPC_MESSAGE, device)) + .build(); + outputStream.onNext(ResponseMsg.newBuilder() + .setEntityUpdateMsg(entityUpdateMsg) + .build()); + } + } else { + Device deviceById = ctx.getDeviceService().findDeviceById(edge.getTenantId(), edgeDeviceId); + if (deviceById != null) { + // this ID already used by other device - create new device and update ID on the edge + Device savedDevice = createDevice(deviceUpdateMsg); + EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg.newBuilder() + .setDeviceUpdateMsg(ctx.getDeviceUpdateMsgConstructor().constructDeviceUpdatedMsg(UpdateMsgType.DEVICE_CONFLICT_RPC_MESSAGE, savedDevice)) + .build(); + outputStream.onNext(ResponseMsg.newBuilder() + .setEntityUpdateMsg(entityUpdateMsg) + .build()); + } else { + createDevice(deviceUpdateMsg); + } + } + break; + case ENTITY_UPDATED_RPC_MESSAGE: + updateDevice(deviceUpdateMsg); + break; + case ENTITY_DELETED_RPC_MESSAGE: + Device deviceToDelete = ctx.getDeviceService().findDeviceById(edge.getTenantId(), edgeDeviceId); + if (deviceToDelete != null) { + ctx.getDeviceService().unassignDeviceFromEdge(edge.getTenantId(), edgeDeviceId, edge.getId()); + } + break; + case UNRECOGNIZED: + log.error("Unsupported msg type"); + } + } + + private void updateDevice(DeviceUpdateMsg deviceUpdateMsg) { + DeviceId deviceId = new DeviceId(new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB())); + Device device = ctx.getDeviceService().findDeviceById(edge.getTenantId(), deviceId); + device.setName(deviceUpdateMsg.getName()); + device.setType(deviceUpdateMsg.getType()); + device.setLabel(deviceUpdateMsg.getLabel()); + device = ctx.getDeviceService().saveDevice(device); + updateDeviceCredentials(deviceUpdateMsg, device); + } + + private void updateDeviceCredentials(DeviceUpdateMsg deviceUpdateMsg, Device device) { + log.debug("Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]", + device.getName(), deviceUpdateMsg.getCredentialsId(), deviceUpdateMsg.getCredentialsValue()); + + DeviceCredentials deviceCredentials = ctx.getDeviceCredentialsService().findDeviceCredentialsByDeviceId(edge.getTenantId(), device.getId()); + deviceCredentials.setCredentialsType(DeviceCredentialsType.valueOf(deviceUpdateMsg.getCredentialsType())); + deviceCredentials.setCredentialsId(deviceUpdateMsg.getCredentialsId()); + deviceCredentials.setCredentialsValue(deviceUpdateMsg.getCredentialsValue()); + ctx.getDeviceCredentialsService().updateDeviceCredentials(edge.getTenantId(), deviceCredentials); + log.debug("Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]", + device.getName(), deviceUpdateMsg.getCredentialsId(), deviceUpdateMsg.getCredentialsValue()); + + } + + private Device createDevice(DeviceUpdateMsg deviceUpdateMsg) { + Device device; + try { + deviceCreationLock.lock(); + DeviceId deviceId = new DeviceId(new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB())); + device = new Device(); + device.setTenantId(edge.getTenantId()); + device.setCustomerId(edge.getCustomerId()); + device.setId(deviceId); + device.setName(deviceUpdateMsg.getName()); + device.setType(deviceUpdateMsg.getType()); + device.setLabel(deviceUpdateMsg.getLabel()); + device = ctx.getDeviceService().saveDevice(device); + device = ctx.getDeviceService().assignDeviceToEdge(edge.getTenantId(), device.getId(), edge.getId()); + createRelationFromEdge(device.getId()); + ctx.getRelationService().saveRelationAsync(TenantId.SYS_TENANT_ID, new EntityRelation(edge.getId(), device.getId(), "Created")); + ctx.getDeviceStateService().onDeviceAdded(device); + updateDeviceCredentials(deviceUpdateMsg, device); + } finally { + deviceCreationLock.unlock(); + } + return device; + } + private EntityId getAlarmOriginator(String entityName, org.thingsboard.server.common.data.EntityType entityType) { switch (entityType) { case DEVICE: @@ -643,36 +712,6 @@ public final class EdgeGrpcSession implements Closeable { } } - private Device getOrCreateDevice(String deviceName, String deviceType) { - Device device = ctx.getDeviceService().findDeviceByTenantIdAndName(edge.getTenantId(), deviceName); - if (device == null) { - entityCreationLock.lock(); - try { - return processGetOrCreateDevice(deviceName, deviceType); - } finally { - entityCreationLock.unlock(); - } - } - return device; - } - - private Device processGetOrCreateDevice(String deviceName, String deviceType) { - Device device = ctx.getDeviceService().findDeviceByTenantIdAndName(edge.getTenantId(), deviceName); - if (device == null) { - device = new Device(); - device.setName(deviceName); - device.setType(deviceType); - device.setTenantId(edge.getTenantId()); - device.setCustomerId(edge.getCustomerId()); - device = ctx.getDeviceService().saveDevice(device); - device = ctx.getDeviceService().assignDeviceToEdge(edge.getTenantId(), device.getId(), edge.getId()); - createRelationFromEdge(device.getId()); - ctx.getRelationService().saveRelationAsync(TenantId.SYS_TENANT_ID, new EntityRelation(edge.getId(), device.getId(), "Created")); - ctx.getDeviceStateService().onDeviceAdded(device); - } - return device; - } - private ConnectResponseMsg processConnect(ConnectRequestMsg request) { Optional optional = ctx.getEdgeService().findEdgeByRoutingKey(TenantId.SYS_TENANT_ID, request.getEdgeRoutingKey()); if (optional.isPresent()) { diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index ca32322977..9c38c71baa 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -94,14 +94,14 @@ enum UpdateMsgType { ALARM_ACK_RPC_MESSAGE = 3; ALARM_CLEAR_RPC_MESSAGE = 4; RULE_CHAIN_CUSTOM_MESSAGE = 5; + DEVICE_CONFLICT_RPC_MESSAGE = 6; } message EntityDataProto { string entityName = 1; - string entityType = 2; - int64 entityIdMSB = 3; - int64 entityIdLSB = 4; - bytes tbMsg = 5; + int64 entityIdMSB = 2; + int64 entityIdLSB = 3; + bytes tbMsg = 4; } message RuleChainUpdateMsg { @@ -156,7 +156,6 @@ message DashboardUpdateMsg { int64 idLSB = 3; string title = 4; string configuration = 5; - string groupName = 6; } message DeviceUpdateMsg { @@ -169,7 +168,6 @@ message DeviceUpdateMsg { string credentialsType = 7; string credentialsId = 8; string credentialsValue = 9; - string groupName = 10; } message AssetUpdateMsg { @@ -179,7 +177,6 @@ message AssetUpdateMsg { string name = 4; string type = 5; string label = 6; - string groupName = 7; } message EntityViewUpdateMsg { @@ -191,7 +188,6 @@ message EntityViewUpdateMsg { int64 entityIdMSB = 6; int64 entityIdLSB = 7; EdgeEntityType entityType = 8; - string groupName = 9; } message AlarmUpdateMsg { @@ -237,7 +233,6 @@ message UserUpdateMsg { string additionalInfo = 8; bool enabled = 9; string password = 10; - string groupName = 11; } message RuleChainMetadataRequestMsg {