From 450026d53bc6c8f69f51a2660c8890f127ac3d78 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Fri, 4 Sep 2020 16:44:31 +0300 Subject: [PATCH] Added processor. Renaming constructors --- .../service/edge/EdgeContextComponent.java | 73 ++- .../service/edge/rpc/EdgeGrpcSession.java | 585 +++--------------- ....java => AdminSettingsMsgConstructor.java} | 2 +- ...structor.java => AlarmMsgConstructor.java} | 2 +- ...structor.java => AssetMsgConstructor.java} | 2 +- ...uctor.java => CustomerMsgConstructor.java} | 2 +- ...ctor.java => DashboardMsgConstructor.java} | 2 +- ...tructor.java => DeviceMsgConstructor.java} | 2 +- ...tor.java => EntityViewMsgConstructor.java} | 2 +- ...uctor.java => RelationMsgConstructor.java} | 2 +- ...ctor.java => RuleChainMsgConstructor.java} | 2 +- ...nstructor.java => UserMsgConstructor.java} | 2 +- ...tor.java => WidgetTypeMsgConstructor.java} | 2 +- ....java => WidgetsBundleMsgConstructor.java} | 2 +- .../edge/rpc/processor/AlarmProcessor.java | 96 +++ .../edge/rpc/processor/BaseProcessor.java | 112 ++++ .../edge/rpc/processor/DeviceProcessor.java | 216 +++++++ .../edge/rpc/processor/RelationProcessor.java | 102 +++ .../rpc/processor/TelemetryProcessor.java | 202 ++++++ .../server/common/data/audit/ActionType.java | 3 +- 20 files changed, 856 insertions(+), 557 deletions(-) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{AdminSettingsUpdateMsgConstructor.java => AdminSettingsMsgConstructor.java} (96%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{AlarmUpdateMsgConstructor.java => AlarmMsgConstructor.java} (98%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{AssetUpdateMsgConstructor.java => AssetMsgConstructor.java} (98%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{CustomerUpdateMsgConstructor.java => CustomerMsgConstructor.java} (98%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{DashboardUpdateMsgConstructor.java => DashboardMsgConstructor.java} (98%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{DeviceUpdateMsgConstructor.java => DeviceMsgConstructor.java} (98%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{EntityViewUpdateMsgConstructor.java => EntityViewMsgConstructor.java} (98%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{RelationUpdateMsgConstructor.java => RelationMsgConstructor.java} (97%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{RuleChainUpdateMsgConstructor.java => RuleChainMsgConstructor.java} (99%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{UserUpdateMsgConstructor.java => UserMsgConstructor.java} (98%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{WidgetTypeUpdateMsgConstructor.java => WidgetTypeMsgConstructor.java} (98%) rename application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/{WidgetsBundleUpdateMsgConstructor.java => WidgetsBundleMsgConstructor.java} (97%) create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmProcessor.java create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseProcessor.java create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProcessor.java create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationProcessor.java create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryProcessor.java 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 f5a97e2326..c90cc99965 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 @@ -36,23 +36,26 @@ import org.thingsboard.server.dao.user.UserService; import org.thingsboard.server.dao.widget.WidgetTypeService; import org.thingsboard.server.dao.widget.WidgetsBundleService; import org.thingsboard.server.queue.discovery.PartitionService; -import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.rpc.EdgeEventStorageSettings; -import org.thingsboard.server.service.edge.rpc.constructor.AdminSettingsUpdateMsgConstructor; -import org.thingsboard.server.service.edge.rpc.constructor.AlarmUpdateMsgConstructor; -import org.thingsboard.server.service.edge.rpc.constructor.AssetUpdateMsgConstructor; -import org.thingsboard.server.service.edge.rpc.constructor.CustomerUpdateMsgConstructor; -import org.thingsboard.server.service.edge.rpc.constructor.DashboardUpdateMsgConstructor; -import org.thingsboard.server.service.edge.rpc.constructor.DeviceUpdateMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.AdminSettingsMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.AlarmMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.AssetMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.CustomerMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.DashboardMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.DeviceMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.EntityDataMsgConstructor; -import org.thingsboard.server.service.edge.rpc.constructor.EntityViewUpdateMsgConstructor; -import org.thingsboard.server.service.edge.rpc.constructor.RelationUpdateMsgConstructor; -import org.thingsboard.server.service.edge.rpc.constructor.RuleChainUpdateMsgConstructor; -import org.thingsboard.server.service.edge.rpc.constructor.UserUpdateMsgConstructor; -import org.thingsboard.server.service.edge.rpc.constructor.WidgetTypeUpdateMsgConstructor; -import org.thingsboard.server.service.edge.rpc.constructor.WidgetsBundleUpdateMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.EntityViewMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.RelationMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.RuleChainMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.UserMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.WidgetTypeMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.WidgetsBundleMsgConstructor; import org.thingsboard.server.service.edge.rpc.init.SyncEdgeService; +import org.thingsboard.server.service.edge.rpc.processor.AlarmProcessor; +import org.thingsboard.server.service.edge.rpc.processor.DeviceProcessor; +import org.thingsboard.server.service.edge.rpc.processor.RelationProcessor; +import org.thingsboard.server.service.edge.rpc.processor.TelemetryProcessor; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.queue.TbClusterService; import org.thingsboard.server.service.state.DeviceStateService; @@ -69,10 +72,6 @@ public class EdgeContextComponent { @Autowired private PartitionService partitionService; - @Autowired - @Lazy - private TbQueueProducerProvider producerProvider; - @Lazy @Autowired private EdgeNotificationService edgeNotificationService; @@ -147,56 +146,72 @@ public class EdgeContextComponent { @Lazy @Autowired - private RuleChainUpdateMsgConstructor ruleChainUpdateMsgConstructor; + private RuleChainMsgConstructor ruleChainMsgConstructor; @Lazy @Autowired - private AlarmUpdateMsgConstructor alarmUpdateMsgConstructor; + private AlarmMsgConstructor alarmMsgConstructor; @Lazy @Autowired - private DeviceUpdateMsgConstructor deviceUpdateMsgConstructor; + private DeviceMsgConstructor deviceMsgConstructor; @Lazy @Autowired - private AssetUpdateMsgConstructor assetUpdateMsgConstructor; + private AssetMsgConstructor assetMsgConstructor; @Lazy @Autowired - private EntityViewUpdateMsgConstructor entityViewUpdateMsgConstructor; + private EntityViewMsgConstructor entityViewMsgConstructor; @Lazy @Autowired - private DashboardUpdateMsgConstructor dashboardUpdateMsgConstructor; + private DashboardMsgConstructor dashboardMsgConstructor; @Lazy @Autowired - private CustomerUpdateMsgConstructor customerUpdateMsgConstructor; + private CustomerMsgConstructor customerMsgConstructor; @Lazy @Autowired - private UserUpdateMsgConstructor userUpdateMsgConstructor; + private UserMsgConstructor userMsgConstructor; @Lazy @Autowired - private RelationUpdateMsgConstructor relationUpdateMsgConstructor; + private RelationMsgConstructor relationMsgConstructor; @Lazy @Autowired - private WidgetsBundleUpdateMsgConstructor widgetsBundleUpdateMsgConstructor; + private WidgetsBundleMsgConstructor widgetsBundleMsgConstructor; @Lazy @Autowired - private WidgetTypeUpdateMsgConstructor widgetTypeUpdateMsgConstructor; + private WidgetTypeMsgConstructor widgetTypeMsgConstructor; @Lazy @Autowired - private AdminSettingsUpdateMsgConstructor adminSettingsUpdateMsgConstructor; + private AdminSettingsMsgConstructor adminSettingsMsgConstructor; @Lazy @Autowired private EntityDataMsgConstructor entityDataMsgConstructor; + @Lazy + @Autowired + private AlarmProcessor alarmProcessor; + + @Lazy + @Autowired + private DeviceProcessor deviceProcessor; + + @Lazy + @Autowired + private RelationProcessor relationProcessor; + + @Lazy + @Autowired + private TelemetryProcessor telemetryProcessor; + @Lazy @Autowired private EdgeEventStorageSettings edgeEventStorageSettings; 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 2232c43286..078978435d 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 @@ -18,39 +18,29 @@ package org.thingsboard.server.service.edge.rpc; import com.datastax.driver.core.utils.UUIDs; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; -import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; -import com.google.common.util.concurrent.SettableFuture; -import com.google.gson.Gson; import com.google.gson.JsonElement; -import com.google.gson.JsonObject; import io.grpc.stub.StreamObserver; import lombok.Data; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang.RandomStringUtils; import org.checkerframework.checker.nullness.qual.Nullable; -import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg; import org.thingsboard.server.common.data.AdminSettings; 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.HasCustomerId; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.alarm.Alarm; -import org.thingsboard.server.common.data.alarm.AlarmSeverity; -import org.thingsboard.server.common.data.alarm.AlarmStatus; 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.EdgeEvent; -import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; -import org.thingsboard.server.common.data.exception.ThingsboardException; +import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.CustomerId; @@ -58,38 +48,28 @@ import org.thingsboard.server.common.data.id.DashboardId; import org.thingsboard.server.common.data.id.DeviceId; 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.EntityViewId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.id.WidgetTypeId; import org.thingsboard.server.common.data.id.WidgetsBundleId; -import org.thingsboard.server.common.data.kv.AttributeKey; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.page.TimePageData; import org.thingsboard.server.common.data.page.TimePageLink; 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.data.security.UserCredentials; import org.thingsboard.server.common.data.widget.WidgetType; import org.thingsboard.server.common.data.widget.WidgetsBundle; -import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.common.msg.TbMsgMetaData; -import org.thingsboard.server.common.msg.queue.ServiceType; -import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; -import org.thingsboard.server.common.msg.session.SessionMsgType; import org.thingsboard.server.common.transport.util.JsonUtils; import org.thingsboard.server.gen.edge.AdminSettingsUpdateMsg; import org.thingsboard.server.gen.edge.AlarmUpdateMsg; 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.ConnectRequestMsg; import org.thingsboard.server.gen.edge.ConnectResponseCode; @@ -119,20 +99,13 @@ import org.thingsboard.server.gen.edge.UserCredentialsRequestMsg; import org.thingsboard.server.gen.edge.UserCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.WidgetTypeUpdateMsg; import org.thingsboard.server.gen.edge.WidgetsBundleUpdateMsg; -import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.TbQueueCallback; -import org.thingsboard.server.queue.TbQueueMsgMetadata; -import org.thingsboard.server.queue.TbQueueProducer; -import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.service.edge.EdgeContextComponent; import java.io.Closeable; import java.util.ArrayList; import java.util.Collections; -import java.util.HashSet; import java.util.List; import java.util.Optional; -import java.util.Set; import java.util.UUID; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; @@ -145,12 +118,8 @@ import java.util.function.Consumer; @Data public final class EdgeGrpcSession implements Closeable { - private static final ReentrantLock deviceCreationLock = new ReentrantLock(); - private static final ReentrantLock responseMsgLock = new ReentrantLock(); - private final Gson gson = new Gson(); - private static final String QUEUE_START_TS_ATTR_KEY = "queueStartTs"; private final UUID sessionId; @@ -166,8 +135,6 @@ public final class EdgeGrpcSession implements Closeable { private CountDownLatch latch; - private TbQueueProducer> ruleEngineMsgProducer; - EdgeGrpcSession(EdgeContextComponent ctx, StreamObserver outputStream, BiConsumer sessionOpenListener, Consumer sessionCloseListener, ObjectMapper mapper) { this.sessionId = UUID.randomUUID(); @@ -176,7 +143,6 @@ public final class EdgeGrpcSession implements Closeable { this.sessionOpenListener = sessionOpenListener; this.sessionCloseListener = sessionCloseListener; this.mapper = mapper; - this.ruleEngineMsgProducer = ctx.getProducerProvider().getRuleEngineMsgProducer(); initInputStream(); } @@ -361,6 +327,12 @@ public final class EdgeGrpcSession implements Closeable { case TIMESERIES_UPDATED: downlinkMsg = processTelemetryMessage(edgeEvent); break; + case CREDENTIALS_REQUEST: + downlinkMsg = processCredentialsRequestMessage(edgeEvent); + break; + case ENTITY_EXISTS_REQUEST: + downlinkMsg = processEntityExistsRequestMessage(edgeEvent); + break; } if (downlinkMsg != null) { result.add(downlinkMsg); @@ -372,6 +344,36 @@ public final class EdgeGrpcSession implements Closeable { return result; } + private DownlinkMsg processEntityExistsRequestMessage(EdgeEvent edgeEvent) { + DownlinkMsg downlinkMsg = null; + if (EdgeEventType.DEVICE.equals(edgeEvent.getEdgeEventType())) { + DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); + Device device = ctx.getDeviceService().findDeviceById(edge.getTenantId(), deviceId); + CustomerId customerId = getCustomerIdIfEdgeAssignedToCustomer(device); + DeviceUpdateMsg d = ctx.getDeviceMsgConstructor().constructDeviceUpdatedMsg(UpdateMsgType.DEVICE_CONFLICT_RPC_MESSAGE, device, customerId); + downlinkMsg = DownlinkMsg.newBuilder() + .addAllDeviceUpdateMsg(Collections.singletonList(d)) + .build(); + } + return downlinkMsg; + } + + private DownlinkMsg processCredentialsRequestMessage(EdgeEvent edgeEvent) { + DownlinkMsg downlinkMsg = null; + if (EdgeEventType.DEVICE.equals(edgeEvent.getEdgeEventType())) { + DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); + DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = DeviceCredentialsRequestMsg.newBuilder() + .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) + .build(); + DownlinkMsg.Builder builder = DownlinkMsg.newBuilder() + .addAllDeviceCredentialsRequestMsg(Collections.singletonList(deviceCredentialsRequestMsg)); + downlinkMsg = builder.build(); + } + return downlinkMsg; + } + + private ListenableFuture getQueueStartTs() { ListenableFuture> future = ctx.getAttributesService().find(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, QUEUE_START_TS_ATTR_KEY); @@ -479,7 +481,7 @@ public final class EdgeGrpcSession implements Closeable { if (device != null) { CustomerId customerId = getCustomerIdIfEdgeAssignedToCustomer(device); DeviceUpdateMsg deviceUpdateMsg = - ctx.getDeviceUpdateMsgConstructor().constructDeviceUpdatedMsg(msgType, device, customerId); + ctx.getDeviceMsgConstructor().constructDeviceUpdatedMsg(msgType, device, customerId); downlinkMsg = DownlinkMsg.newBuilder() .addAllDeviceUpdateMsg(Collections.singletonList(deviceUpdateMsg)) .build(); @@ -488,7 +490,7 @@ public final class EdgeGrpcSession implements Closeable { case DELETED: case UNASSIGNED_FROM_EDGE: DeviceUpdateMsg deviceUpdateMsg = - ctx.getDeviceUpdateMsgConstructor().constructDeviceDeleteMsg(deviceId); + ctx.getDeviceMsgConstructor().constructDeviceDeleteMsg(deviceId); downlinkMsg = DownlinkMsg.newBuilder() .addAllDeviceUpdateMsg(Collections.singletonList(deviceUpdateMsg)) .build(); @@ -497,7 +499,7 @@ public final class EdgeGrpcSession implements Closeable { DeviceCredentials deviceCredentials = ctx.getDeviceCredentialsService().findDeviceCredentialsByDeviceId(edge.getTenantId(), deviceId); if (deviceCredentials != null) { DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = - ctx.getDeviceUpdateMsgConstructor().constructDeviceCredentialsUpdatedMsg(deviceCredentials); + ctx.getDeviceMsgConstructor().constructDeviceCredentialsUpdatedMsg(deviceCredentials); downlinkMsg = DownlinkMsg.newBuilder() .addAllDeviceCredentialsUpdateMsg(Collections.singletonList(deviceCredentialsUpdateMsg)) .build(); @@ -520,7 +522,7 @@ public final class EdgeGrpcSession implements Closeable { if (asset != null) { CustomerId customerId = getCustomerIdIfEdgeAssignedToCustomer(asset); AssetUpdateMsg assetUpdateMsg = - ctx.getAssetUpdateMsgConstructor().constructAssetUpdatedMsg(msgType, asset, customerId); + ctx.getAssetMsgConstructor().constructAssetUpdatedMsg(msgType, asset, customerId); downlinkMsg = DownlinkMsg.newBuilder() .addAllAssetUpdateMsg(Collections.singletonList(assetUpdateMsg)) .build(); @@ -529,7 +531,7 @@ public final class EdgeGrpcSession implements Closeable { case DELETED: case UNASSIGNED_FROM_EDGE: AssetUpdateMsg assetUpdateMsg = - ctx.getAssetUpdateMsgConstructor().constructAssetDeleteMsg(assetId); + ctx.getAssetMsgConstructor().constructAssetDeleteMsg(assetId); downlinkMsg = DownlinkMsg.newBuilder() .addAllAssetUpdateMsg(Collections.singletonList(assetUpdateMsg)) .build(); @@ -551,7 +553,7 @@ public final class EdgeGrpcSession implements Closeable { if (entityView != null) { CustomerId customerId = getCustomerIdIfEdgeAssignedToCustomer(entityView); EntityViewUpdateMsg entityViewUpdateMsg = - ctx.getEntityViewUpdateMsgConstructor().constructEntityViewUpdatedMsg(msgType, entityView, customerId); + ctx.getEntityViewMsgConstructor().constructEntityViewUpdatedMsg(msgType, entityView, customerId); downlinkMsg = DownlinkMsg.newBuilder() .addAllEntityViewUpdateMsg(Collections.singletonList(entityViewUpdateMsg)) .build(); @@ -560,7 +562,7 @@ public final class EdgeGrpcSession implements Closeable { case DELETED: case UNASSIGNED_FROM_EDGE: EntityViewUpdateMsg entityViewUpdateMsg = - ctx.getEntityViewUpdateMsgConstructor().constructEntityViewDeleteMsg(entityViewId); + ctx.getEntityViewMsgConstructor().constructEntityViewDeleteMsg(entityViewId); downlinkMsg = DownlinkMsg.newBuilder() .addAllEntityViewUpdateMsg(Collections.singletonList(entityViewUpdateMsg)) .build(); @@ -585,7 +587,7 @@ public final class EdgeGrpcSession implements Closeable { customerId = edge.getCustomerId(); } DashboardUpdateMsg dashboardUpdateMsg = - ctx.getDashboardUpdateMsgConstructor().constructDashboardUpdatedMsg(msgType, dashboard, customerId); + ctx.getDashboardMsgConstructor().constructDashboardUpdatedMsg(msgType, dashboard, customerId); downlinkMsg = DownlinkMsg.newBuilder() .addAllDashboardUpdateMsg(Collections.singletonList(dashboardUpdateMsg)) .build(); @@ -594,7 +596,7 @@ public final class EdgeGrpcSession implements Closeable { case DELETED: case UNASSIGNED_FROM_EDGE: DashboardUpdateMsg dashboardUpdateMsg = - ctx.getDashboardUpdateMsgConstructor().constructDashboardDeleteMsg(dashboardId); + ctx.getDashboardMsgConstructor().constructDashboardDeleteMsg(dashboardId); downlinkMsg = DownlinkMsg.newBuilder() .addAllDashboardUpdateMsg(Collections.singletonList(dashboardUpdateMsg)) .build(); @@ -612,7 +614,7 @@ public final class EdgeGrpcSession implements Closeable { Customer customer = ctx.getCustomerService().findCustomerById(edgeEvent.getTenantId(), customerId); if (customer != null) { CustomerUpdateMsg customerUpdateMsg = - ctx.getCustomerUpdateMsgConstructor().constructCustomerUpdatedMsg(msgType, customer); + ctx.getCustomerMsgConstructor().constructCustomerUpdatedMsg(msgType, customer); downlinkMsg = DownlinkMsg.newBuilder() .addAllCustomerUpdateMsg(Collections.singletonList(customerUpdateMsg)) .build(); @@ -620,7 +622,7 @@ public final class EdgeGrpcSession implements Closeable { break; case DELETED: CustomerUpdateMsg customerUpdateMsg = - ctx.getCustomerUpdateMsgConstructor().constructCustomerDeleteMsg(customerId); + ctx.getCustomerMsgConstructor().constructCustomerDeleteMsg(customerId); downlinkMsg = DownlinkMsg.newBuilder() .addAllCustomerUpdateMsg(Collections.singletonList(customerUpdateMsg)) .build(); @@ -639,7 +641,7 @@ public final class EdgeGrpcSession implements Closeable { RuleChain ruleChain = ctx.getRuleChainService().findRuleChainById(edgeEvent.getTenantId(), ruleChainId); if (ruleChain != null) { RuleChainUpdateMsg ruleChainUpdateMsg = - ctx.getRuleChainUpdateMsgConstructor().constructRuleChainUpdatedMsg(edge.getRootRuleChainId(), msgType, ruleChain); + ctx.getRuleChainMsgConstructor().constructRuleChainUpdatedMsg(edge.getRootRuleChainId(), msgType, ruleChain); downlinkMsg = DownlinkMsg.newBuilder() .addAllRuleChainUpdateMsg(Collections.singletonList(ruleChainUpdateMsg)) .build(); @@ -648,7 +650,7 @@ public final class EdgeGrpcSession implements Closeable { case DELETED: case UNASSIGNED_FROM_EDGE: downlinkMsg = DownlinkMsg.newBuilder() - .addAllRuleChainUpdateMsg(Collections.singletonList(ctx.getRuleChainUpdateMsgConstructor().constructRuleChainDeleteMsg(ruleChainId))) + .addAllRuleChainUpdateMsg(Collections.singletonList(ctx.getRuleChainMsgConstructor().constructRuleChainDeleteMsg(ruleChainId))) .build(); break; } @@ -662,7 +664,7 @@ public final class EdgeGrpcSession implements Closeable { if (ruleChain != null) { RuleChainMetaData ruleChainMetaData = ctx.getRuleChainService().loadRuleChainMetaData(edgeEvent.getTenantId(), ruleChainId); RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = - ctx.getRuleChainUpdateMsgConstructor().constructRuleChainMetadataUpdatedMsg(msgType, ruleChainMetaData); + ctx.getRuleChainMsgConstructor().constructRuleChainMetadataUpdatedMsg(msgType, ruleChainMetaData); if (ruleChainMetadataUpdateMsg != null) { downlinkMsg = DownlinkMsg.newBuilder() .addAllRuleChainMetadataUpdateMsg(Collections.singletonList(ruleChainMetadataUpdateMsg)) @@ -682,20 +684,20 @@ public final class EdgeGrpcSession implements Closeable { if (user != null) { CustomerId customerId = getCustomerIdIfEdgeAssignedToCustomer(user); downlinkMsg = DownlinkMsg.newBuilder() - .addAllUserUpdateMsg(Collections.singletonList(ctx.getUserUpdateMsgConstructor().constructUserUpdatedMsg(msgType, user, customerId))) + .addAllUserUpdateMsg(Collections.singletonList(ctx.getUserMsgConstructor().constructUserUpdatedMsg(msgType, user, customerId))) .build(); } break; case DELETED: downlinkMsg = DownlinkMsg.newBuilder() - .addAllUserUpdateMsg(Collections.singletonList(ctx.getUserUpdateMsgConstructor().constructUserDeleteMsg(userId))) + .addAllUserUpdateMsg(Collections.singletonList(ctx.getUserMsgConstructor().constructUserDeleteMsg(userId))) .build(); break; case CREDENTIALS_UPDATED: UserCredentials userCredentialsByUserId = ctx.getUserService().findUserCredentialsByUserId(edge.getTenantId(), userId); if (userCredentialsByUserId != null && userCredentialsByUserId.isEnabled()) { UserCredentialsUpdateMsg userCredentialsUpdateMsg = - ctx.getUserUpdateMsgConstructor().constructUserCredentialsUpdatedMsg(userCredentialsByUserId); + ctx.getUserMsgConstructor().constructUserCredentialsUpdatedMsg(userCredentialsByUserId); downlinkMsg = DownlinkMsg.newBuilder() .addAllUserCredentialsUpdateMsg(Collections.singletonList(userCredentialsUpdateMsg)) .build(); @@ -714,7 +716,7 @@ public final class EdgeGrpcSession implements Closeable { private DownlinkMsg processRelation(EdgeEvent edgeEvent, UpdateMsgType msgType) { EntityRelation entityRelation = mapper.convertValue(edgeEvent.getEntityBody(), EntityRelation.class); - RelationUpdateMsg r = ctx.getRelationUpdateMsgConstructor().constructRelationUpdatedMsg(msgType, entityRelation); + RelationUpdateMsg r = ctx.getRelationMsgConstructor().constructRelationUpdatedMsg(msgType, entityRelation); return DownlinkMsg.newBuilder() .addAllRelationUpdateMsg(Collections.singletonList(r)) .build(); @@ -727,7 +729,7 @@ public final class EdgeGrpcSession implements Closeable { Alarm alarm = ctx.getAlarmService().findAlarmByIdAsync(edgeEvent.getTenantId(), alarmId).get(); if (alarm != null) { downlinkMsg = DownlinkMsg.newBuilder() - .addAllAlarmUpdateMsg(Collections.singletonList(ctx.getAlarmUpdateMsgConstructor().constructAlarmUpdatedMsg(edge.getTenantId(), msgType, alarm))) + .addAllAlarmUpdateMsg(Collections.singletonList(ctx.getAlarmMsgConstructor().constructAlarmUpdatedMsg(edge.getTenantId(), msgType, alarm))) .build(); } } catch (Exception e) { @@ -745,7 +747,7 @@ public final class EdgeGrpcSession implements Closeable { WidgetsBundle widgetsBundle = ctx.getWidgetsBundleService().findWidgetsBundleById(edgeEvent.getTenantId(), widgetsBundleId); if (widgetsBundle != null) { WidgetsBundleUpdateMsg widgetsBundleUpdateMsg = - ctx.getWidgetsBundleUpdateMsgConstructor().constructWidgetsBundleUpdateMsg(msgType, widgetsBundle); + ctx.getWidgetsBundleMsgConstructor().constructWidgetsBundleUpdateMsg(msgType, widgetsBundle); downlinkMsg = DownlinkMsg.newBuilder() .addAllWidgetsBundleUpdateMsg(Collections.singletonList(widgetsBundleUpdateMsg)) .build(); @@ -753,7 +755,7 @@ public final class EdgeGrpcSession implements Closeable { break; case DELETED: WidgetsBundleUpdateMsg widgetsBundleUpdateMsg = - ctx.getWidgetsBundleUpdateMsgConstructor().constructWidgetsBundleDeleteMsg(widgetsBundleId); + ctx.getWidgetsBundleMsgConstructor().constructWidgetsBundleDeleteMsg(widgetsBundleId); downlinkMsg = DownlinkMsg.newBuilder() .addAllWidgetsBundleUpdateMsg(Collections.singletonList(widgetsBundleUpdateMsg)) .build(); @@ -771,7 +773,7 @@ public final class EdgeGrpcSession implements Closeable { WidgetType widgetType = ctx.getWidgetTypeService().findWidgetTypeById(edgeEvent.getTenantId(), widgetTypeId); if (widgetType != null) { WidgetTypeUpdateMsg widgetTypeUpdateMsg = - ctx.getWidgetTypeUpdateMsgConstructor().constructWidgetTypeUpdateMsg(msgType, widgetType); + ctx.getWidgetTypeMsgConstructor().constructWidgetTypeUpdateMsg(msgType, widgetType); downlinkMsg = DownlinkMsg.newBuilder() .addAllWidgetTypeUpdateMsg(Collections.singletonList(widgetTypeUpdateMsg)) .build(); @@ -779,10 +781,10 @@ public final class EdgeGrpcSession implements Closeable { break; case DELETED: WidgetTypeUpdateMsg widgetTypeUpdateMsg = - ctx.getWidgetTypeUpdateMsgConstructor().constructWidgetTypeDeleteMsg(widgetTypeId); - downlinkMsg = DownlinkMsg.newBuilder() - .addAllWidgetTypeUpdateMsg(Collections.singletonList(widgetTypeUpdateMsg)) - .build(); + ctx.getWidgetTypeMsgConstructor().constructWidgetTypeDeleteMsg(widgetTypeId); + downlinkMsg = DownlinkMsg.newBuilder() + .addAllWidgetTypeUpdateMsg(Collections.singletonList(widgetTypeUpdateMsg)) + .build(); break; } return downlinkMsg; @@ -790,7 +792,7 @@ public final class EdgeGrpcSession implements Closeable { private DownlinkMsg processAdminSettings(EdgeEvent edgeEvent) { AdminSettings adminSettings = mapper.convertValue(edgeEvent.getEntityBody(), AdminSettings.class); - AdminSettingsUpdateMsg t = ctx.getAdminSettingsUpdateMsgConstructor().constructAdminSettingsUpdateMsg(adminSettings); + AdminSettingsUpdateMsg t = ctx.getAdminSettingsMsgConstructor().constructAdminSettingsUpdateMsg(adminSettings); return DownlinkMsg.newBuilder() .addAllAdminSettingsUpdateMsg(Collections.singletonList(t)) .build(); @@ -832,42 +834,28 @@ public final class EdgeGrpcSession implements Closeable { try { if (uplinkMsg.getEntityDataList() != null && !uplinkMsg.getEntityDataList().isEmpty()) { for (EntityDataProto entityData : uplinkMsg.getEntityDataList()) { - EntityId entityId = constructEntityId(entityData); - if ((entityData.hasPostAttributesMsg() || entityData.hasPostTelemetryMsg()) && entityId != null) { - TbMsgMetaData metaData = constructBaseMsgMetadata(entityId); - metaData.putValue(DataConstants.MSG_SOURCE_KEY, DataConstants.EDGE_MSG_SOURCE); - if (entityData.hasPostAttributesMsg()) { - metaData.putValue("scope", entityData.getPostAttributeScope()); - result.add(processPostAttributes(entityId, entityData.getPostAttributesMsg(), metaData)); - } - if (entityData.hasPostTelemetryMsg()) { - result.add(processPostTelemetry(entityId, entityData.getPostTelemetryMsg(), metaData)); - } - } - if (entityData.hasAttributeDeleteMsg()) { - result.add(processAttributeDeleteMsg(entityId, entityData.getAttributeDeleteMsg(), entityData.getEntityType())); - } + result.addAll(ctx.getTelemetryProcessor().onTelemetryUpdate(edge.getTenantId(), entityData)); } } if (uplinkMsg.getDeviceUpdateMsgList() != null && !uplinkMsg.getDeviceUpdateMsgList().isEmpty()) { for (DeviceUpdateMsg deviceUpdateMsg : uplinkMsg.getDeviceUpdateMsgList()) { - result.add(onDeviceUpdate(deviceUpdateMsg)); + result.add(ctx.getDeviceProcessor().onDeviceUpdate(edge.getTenantId(), edge, deviceUpdateMsg)); } } if (uplinkMsg.getDeviceCredentialsUpdateMsgList() != null && !uplinkMsg.getDeviceCredentialsUpdateMsgList().isEmpty()) { for (DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg : uplinkMsg.getDeviceCredentialsUpdateMsgList()) { - result.add(onDeviceCredentialsUpdate(deviceCredentialsUpdateMsg)); + result.add(ctx.getDeviceProcessor().onDeviceCredentialsUpdate(edge.getTenantId(), deviceCredentialsUpdateMsg)); } } if (uplinkMsg.getAlarmUpdateMsgList() != null && !uplinkMsg.getAlarmUpdateMsgList().isEmpty()) { for (AlarmUpdateMsg alarmUpdateMsg : uplinkMsg.getAlarmUpdateMsgList()) { - result.add(onAlarmUpdate(alarmUpdateMsg)); + result.add(ctx.getAlarmProcessor().onAlarmUpdate(edge.getTenantId(), alarmUpdateMsg)); } } if (uplinkMsg.getRelationUpdateMsgList() != null && !uplinkMsg.getRelationUpdateMsgList().isEmpty()) { - for (RelationUpdateMsg relationUpdateMsg: uplinkMsg.getRelationUpdateMsgList()) { - result.add(onRelationUpdate(relationUpdateMsg)); + for (RelationUpdateMsg relationUpdateMsg : uplinkMsg.getRelationUpdateMsgList()) { + result.add(ctx.getRelationProcessor().onRelationUpdate(edge.getTenantId(), relationUpdateMsg)); } } if (uplinkMsg.getRuleChainMetadataRequestMsgList() != null && !uplinkMsg.getRuleChainMetadataRequestMsgList().isEmpty()) { @@ -901,430 +889,6 @@ public final class EdgeGrpcSession implements Closeable { return Futures.allAsList(result); } - private TbMsgMetaData constructBaseMsgMetadata(EntityId entityId) { - TbMsgMetaData metaData = new TbMsgMetaData(); - switch (entityId.getEntityType()) { - case DEVICE: - Device device = ctx.getDeviceService().findDeviceById(edge.getTenantId(), new DeviceId(entityId.getId())); - if (device != null) { - metaData.putValue("deviceName", device.getName()); - metaData.putValue("deviceType", device.getType()); - } - break; - case ASSET: - Asset asset = ctx.getAssetService().findAssetById(edge.getTenantId(), new AssetId(entityId.getId())); - if (asset != null) { - metaData.putValue("assetName", asset.getName()); - metaData.putValue("assetType", asset.getType()); - } - break; - case ENTITY_VIEW: - EntityView entityView = ctx.getEntityViewService().findEntityViewById(edge.getTenantId(), new EntityViewId(entityId.getId())); - if (entityView != null) { - metaData.putValue("entityViewName", entityView.getName()); - metaData.putValue("entityViewType", entityView.getType()); - } - break; - default: - log.debug("Using empty metadata for entityId [{}]", entityId); - break; - } - return metaData; - } - - private EntityId constructEntityId(EntityDataProto entityData) { - EntityType entityType = EntityType.valueOf(entityData.getEntityType()); - switch (entityType) { - case DEVICE: - return new DeviceId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); - case ASSET: - return new AssetId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); - case ENTITY_VIEW: - return new EntityViewId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); - case DASHBOARD: - return new DashboardId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); - case TENANT: - return new TenantId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); - case CUSTOMER: - return new CustomerId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); - case USER: - return new UserId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); - default: - log.warn("Unsupported entity type [{}] during construct of entity id. EntityDataProto [{}]", entityData.getEntityType(), entityData); - return null; - } - } - - private ListenableFuture processPostTelemetry(EntityId entityId, TransportProtos.PostTelemetryMsg msg, TbMsgMetaData metaData) { - SettableFuture futureToSet = SettableFuture.create(); - for (TransportProtos.TsKvListProto tsKv : msg.getTsKvListList()) { - JsonObject json = JsonUtils.getJsonObject(tsKv.getKvList()); - metaData.putValue("ts", tsKv.getTs() + ""); - TbMsg tbMsg = TbMsg.newMsg(SessionMsgType.POST_TELEMETRY_REQUEST.name(), entityId, metaData, gson.toJson(json)); - ctx.getTbClusterService().pushMsgToRuleEngine(edge.getTenantId(), tbMsg.getOriginator(), tbMsg, new TbQueueCallback() { - @Override - public void onSuccess(TbQueueMsgMetadata metadata) { - futureToSet.set(null); - } - - @Override - public void onFailure(Throwable t) { - log.error("Can't process post telemetry [{}]", msg, t); - futureToSet.setException(t); - } - }); - } - return futureToSet; - } - - private ListenableFuture processPostAttributes(EntityId entityId, TransportProtos.PostAttributeMsg msg, TbMsgMetaData metaData) { - SettableFuture futureToSet = SettableFuture.create(); - JsonObject json = JsonUtils.getJsonObject(msg.getKvList()); - TbMsg tbMsg = TbMsg.newMsg(SessionMsgType.POST_ATTRIBUTES_REQUEST.name(), entityId, metaData, gson.toJson(json)); - ctx.getTbClusterService().pushMsgToRuleEngine(edge.getTenantId(), tbMsg.getOriginator(), tbMsg, new TbQueueCallback() { - @Override - public void onSuccess(TbQueueMsgMetadata metadata) { - futureToSet.set(null); - } - - @Override - public void onFailure(Throwable t) { - log.error("Can't process post attributes [{}]", msg, t); - futureToSet.setException(t); - } - }); - return futureToSet; - } - - private ListenableFuture processAttributeDeleteMsg(EntityId entityId, AttributeDeleteMsg attributeDeleteMsg, String entityType) { - SettableFuture futureToSet = SettableFuture.create(); - String scope = attributeDeleteMsg.getScope(); - List attributeNames = attributeDeleteMsg.getAttributeNamesList(); - ctx.getAttributesService().removeAll(edge.getTenantId(), entityId, scope, attributeNames); - if (EntityType.DEVICE.name().equals(entityType)) { - Set attributeKeys = new HashSet<>(); - for (String attributeName : attributeNames) { - attributeKeys.add(new AttributeKey(scope, attributeName)); - } - ctx.getTbClusterService().pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete( - edge.getTenantId(), (DeviceId) entityId, attributeKeys), new TbQueueCallback() { - @Override - public void onSuccess(TbQueueMsgMetadata metadata) { - futureToSet.set(null); - } - - @Override - public void onFailure(Throwable t) { - log.error("Can't process attribute delete msg [{}]", attributeDeleteMsg, t); - futureToSet.setException(t); - } - }); - } - return futureToSet; - } - - private ListenableFuture onDeviceUpdate(DeviceUpdateMsg 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)) { - CustomerId customerId = getCustomerIdIfEdgeAssignedToCustomer(device); - DeviceUpdateMsg d = ctx.getDeviceUpdateMsgConstructor().constructDeviceUpdatedMsg(UpdateMsgType.DEVICE_CONFLICT_RPC_MESSAGE, device, customerId); - DownlinkMsg downlinkMsg = DownlinkMsg.newBuilder() - .addAllDeviceUpdateMsg(Collections.singletonList(d)) - .build(); - sendResponseMsg(ResponseMsg.newBuilder() - .setDownlinkMsg(downlinkMsg) - .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 = createDevice(deviceUpdateMsg); - CustomerId customerId = getCustomerIdIfEdgeAssignedToCustomer(device); - DeviceUpdateMsg d = ctx.getDeviceUpdateMsgConstructor().constructDeviceUpdatedMsg(UpdateMsgType.DEVICE_CONFLICT_RPC_MESSAGE, device, customerId); - DownlinkMsg downlinkMsg = DownlinkMsg.newBuilder() - .addAllDeviceUpdateMsg(Collections.singletonList(d)) - .build(); - sendResponseMsg(ResponseMsg.newBuilder() - .setDownlinkMsg(downlinkMsg) - .build()); - } else { - device = createDevice(deviceUpdateMsg); - } - } - // TODO: voba - assign device only in case device is not assigned yet. Missing functionality to check this relation prior assignment - ctx.getDeviceService().assignDeviceToEdge(edge.getTenantId(), device.getId(), edge.getId()); - 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 {}", deviceUpdateMsg.getMsgType()); - return Futures.immediateFailedFuture(new RuntimeException("Unsupported msg type " + deviceUpdateMsg.getMsgType())); - } - return Futures.immediateFuture(null); - } - - 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); - - requestDeviceCredentialsFromEdge(device); - } - - private ListenableFuture onDeviceCredentialsUpdate(DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg) { - log.debug("Executing onDeviceCredentialsUpdate, deviceCredentialsUpdateMsg [{}]", deviceCredentialsUpdateMsg); - DeviceId deviceId = new DeviceId(new UUID(deviceCredentialsUpdateMsg.getDeviceIdMSB(), deviceCredentialsUpdateMsg.getDeviceIdLSB())); - ListenableFuture deviceFuture = ctx.getDeviceService().findDeviceByIdAsync(edge.getTenantId(), deviceId); - return Futures.transform(deviceFuture, device -> { - if (device != null) { - log.debug("Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]", - device.getName(), deviceCredentialsUpdateMsg.getCredentialsId(), deviceCredentialsUpdateMsg.getCredentialsValue()); - try { - DeviceCredentials deviceCredentials = ctx.getDeviceCredentialsService().findDeviceCredentialsByDeviceId(edge.getTenantId(), device.getId()); - deviceCredentials.setCredentialsType(DeviceCredentialsType.valueOf(deviceCredentialsUpdateMsg.getCredentialsType())); - deviceCredentials.setCredentialsId(deviceCredentialsUpdateMsg.getCredentialsId()); - deviceCredentials.setCredentialsValue(deviceCredentialsUpdateMsg.getCredentialsValue()); - ctx.getDeviceCredentialsService().updateDeviceCredentials(edge.getTenantId(), deviceCredentials); - } catch (Exception e) { - log.error("Can't update device credentials for device [{}], deviceCredentialsUpdateMsg [{}]", device.getName(), deviceCredentialsUpdateMsg, e); - throw new RuntimeException(e); - } - } - return null; - }, ctx.getDbCallbackExecutor()); - } - - private void requestDeviceCredentialsFromEdge(Device device) { - log.debug("Executing requestDeviceCredentialsFromEdge device [{}]", device); - - DownlinkMsg downlinkMsg = constructDeviceCredentialsRequestMsg(device.getId()); - sendResponseMsg(ResponseMsg.newBuilder() - .setDownlinkMsg(downlinkMsg) - .build()); - } - - private DownlinkMsg constructDeviceCredentialsRequestMsg(DeviceId deviceId) { - DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = DeviceCredentialsRequestMsg.newBuilder() - .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) - .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) - .build(); - DownlinkMsg.Builder builder = DownlinkMsg.newBuilder() - .addAllDeviceCredentialsRequestMsg(Collections.singletonList(deviceCredentialsRequestMsg)); - return builder.build(); - } - - 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); - createDeviceCredentials(device); - createRelationFromEdge(device.getId()); - ctx.getDeviceStateService().onDeviceAdded(device); - pushDeviceCreatedEventToRuleEngine(device); - requestDeviceCredentialsFromEdge(device); - } finally { - deviceCreationLock.unlock(); - } - return device; - } - - private void createDeviceCredentials(Device device) { - DeviceCredentials deviceCredentials = new DeviceCredentials(); - deviceCredentials.setDeviceId(device.getId()); - deviceCredentials.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN); - deviceCredentials.setCredentialsId(RandomStringUtils.randomAlphanumeric(20)); - ctx.getDeviceCredentialsService().createDeviceCredentials(device.getTenantId(), deviceCredentials); - } - - private void pushDeviceCreatedEventToRuleEngine(Device device) { - try { - ObjectNode entityNode = mapper.valueToTree(device); - TbMsg tbMsg = TbMsg.newMsg(DataConstants.ENTITY_CREATED, device.getId(), getActionTbMsgMetaData(device.getCustomerId()), mapper.writeValueAsString(entityNode)); - sendToRuleEngine(edge.getTenantId(), tbMsg, new TbQueueCallback() { - @Override - public void onSuccess(TbQueueMsgMetadata metadata) { - // TODO: voba - handle success - log.debug("Successfully send ENTITY_CREATED EVENT to rule engine [{}]", device); - } - - @Override - public void onFailure(Throwable t) { - // TODO: voba - handle failure - log.debug("Failed to send ENTITY_CREATED EVENT to rule engine [{}]", device, t); - } - }); - } catch (JsonProcessingException | IllegalArgumentException e) { - log.warn("[{}] Failed to push device action to rule engine: {}", device.getId(), DataConstants.ENTITY_CREATED, e); - } - } - - protected void sendToRuleEngine(TenantId tenantId, TbMsg tbMsg, TbQueueCallback callback) { - TopicPartitionInfo tpi = ctx.getPartitionService().resolve(ServiceType.TB_RULE_ENGINE, tenantId, tbMsg.getOriginator()); - TransportProtos.ToRuleEngineMsg msg = TransportProtos.ToRuleEngineMsg.newBuilder().setTbMsg(TbMsg.toByteString(tbMsg)) - .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) - .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()).build(); - ruleEngineMsgProducer.send(tpi, new TbProtoQueueMsg<>(tbMsg.getId(), msg), callback); - } - - private TbMsgMetaData getActionTbMsgMetaData(CustomerId customerId) { - TbMsgMetaData metaData = getTbMsgMetaData(edge); - if (customerId != null && !customerId.isNullUid()) { - metaData.putValue("customerId", customerId.toString()); - } - return metaData; - } - - private TbMsgMetaData getTbMsgMetaData(Edge edge) { - TbMsgMetaData metaData = new TbMsgMetaData(); - metaData.putValue("edgeId", edge.getId().toString()); - metaData.putValue("edgeName", edge.getName()); - return metaData; - } - - private EntityId getAlarmOriginator(String entityName, org.thingsboard.server.common.data.EntityType entityType) { - switch (entityType) { - case DEVICE: - return ctx.getDeviceService().findDeviceByTenantIdAndName(edge.getTenantId(), entityName).getId(); - case ASSET: - return ctx.getAssetService().findAssetByTenantIdAndName(edge.getTenantId(), entityName).getId(); - case ENTITY_VIEW: - return ctx.getEntityViewService().findEntityViewByTenantIdAndName(edge.getTenantId(), entityName).getId(); - default: - return null; - } - } - - private ListenableFuture onAlarmUpdate(AlarmUpdateMsg alarmUpdateMsg) { - EntityId originatorId = getAlarmOriginator(alarmUpdateMsg.getOriginatorName(), org.thingsboard.server.common.data.EntityType.valueOf(alarmUpdateMsg.getOriginatorType())); - if (originatorId == null) { - return Futures.immediateFuture(null); - } - try { - Alarm existentAlarm = ctx.getAlarmService().findLatestByOriginatorAndType(edge.getTenantId(), originatorId, alarmUpdateMsg.getType()).get(); - switch (alarmUpdateMsg.getMsgType()) { - case ENTITY_CREATED_RPC_MESSAGE: - case ENTITY_UPDATED_RPC_MESSAGE: - if (existentAlarm == null || existentAlarm.getStatus().isCleared()) { - existentAlarm = new Alarm(); - existentAlarm.setTenantId(edge.getTenantId()); - existentAlarm.setType(alarmUpdateMsg.getName()); - existentAlarm.setOriginator(originatorId); - existentAlarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity())); - existentAlarm.setStartTs(alarmUpdateMsg.getStartTs()); - existentAlarm.setClearTs(alarmUpdateMsg.getClearTs()); - existentAlarm.setPropagate(alarmUpdateMsg.getPropagate()); - } - existentAlarm.setStatus(AlarmStatus.valueOf(alarmUpdateMsg.getStatus())); - existentAlarm.setAckTs(alarmUpdateMsg.getAckTs()); - existentAlarm.setEndTs(alarmUpdateMsg.getEndTs()); - existentAlarm.setDetails(mapper.readTree(alarmUpdateMsg.getDetails())); - ctx.getAlarmService().createOrUpdateAlarm(existentAlarm); - break; - case ALARM_ACK_RPC_MESSAGE: - if (existentAlarm != null) { - ctx.getAlarmService().ackAlarm(edge.getTenantId(), existentAlarm.getId(), alarmUpdateMsg.getAckTs()); - } - break; - case ALARM_CLEAR_RPC_MESSAGE: - if (existentAlarm != null) { - ctx.getAlarmService().clearAlarm(edge.getTenantId(), existentAlarm.getId(), mapper.readTree(alarmUpdateMsg.getDetails()), alarmUpdateMsg.getAckTs()); - } - break; - case ENTITY_DELETED_RPC_MESSAGE: - if (existentAlarm != null) { - ctx.getAlarmService().deleteAlarm(edge.getTenantId(), existentAlarm.getId()); - } - break; - } - return Futures.immediateFuture(null); - } catch (Exception e) { - log.error("Failed to process alarm update msg [{}]", alarmUpdateMsg, e); - return Futures.immediateFailedFuture(new RuntimeException("Failed to process alarm update msg", e)); - } - } - - private ListenableFuture onRelationUpdate(RelationUpdateMsg relationUpdateMsg) { - log.info("onRelationUpdate {}", relationUpdateMsg); - try { - EntityRelation entityRelation = new EntityRelation(); - - UUID fromUUID = new UUID(relationUpdateMsg.getFromIdMSB(), relationUpdateMsg.getFromIdLSB()); - EntityId fromId = EntityIdFactory.getByTypeAndUuid(EntityType.valueOf(relationUpdateMsg.getFromEntityType()), fromUUID); - entityRelation.setFrom(fromId); - - UUID toUUID = new UUID(relationUpdateMsg.getToIdMSB(), relationUpdateMsg.getToIdLSB()); - EntityId toId = EntityIdFactory.getByTypeAndUuid(EntityType.valueOf(relationUpdateMsg.getToEntityType()), toUUID); - entityRelation.setTo(toId); - - entityRelation.setType(relationUpdateMsg.getType()); - entityRelation.setTypeGroup(RelationTypeGroup.valueOf(relationUpdateMsg.getTypeGroup())); - entityRelation.setAdditionalInfo(mapper.readTree(relationUpdateMsg.getAdditionalInfo())); - switch (relationUpdateMsg.getMsgType()) { - case ENTITY_CREATED_RPC_MESSAGE: - case ENTITY_UPDATED_RPC_MESSAGE: - if (isEntityExists(edge.getTenantId(), entityRelation.getTo()) - && isEntityExists(edge.getTenantId(), entityRelation.getFrom())) { - ctx.getRelationService().saveRelationAsync(edge.getTenantId(), entityRelation); - } - break; - case ENTITY_DELETED_RPC_MESSAGE: - ctx.getRelationService().deleteRelation(edge.getTenantId(), entityRelation); - break; - case UNRECOGNIZED: - log.error("Unsupported msg type"); - } - return Futures.immediateFuture(null); - } catch (Exception e) { - log.error("Failed to process relation update msg [{}]", relationUpdateMsg, e); - return Futures.immediateFailedFuture(new RuntimeException("Failed to process relation update msg", e)); - } - } - - private boolean isEntityExists(TenantId tenantId, EntityId entityId) throws ThingsboardException { - switch (entityId.getEntityType()) { - case DEVICE: - return ctx.getDeviceService().findDeviceById(tenantId, new DeviceId(entityId.getId())) != null; - case ASSET: - return ctx.getAssetService().findAssetById(tenantId, new AssetId(entityId.getId())) != null; - case ENTITY_VIEW: - return ctx.getEntityViewService().findEntityViewById(tenantId, new EntityViewId(entityId.getId())) != null; - case CUSTOMER: - return ctx.getCustomerService().findCustomerById(tenantId, new CustomerId(entityId.getId())) != null; - case USER: - return ctx.getUserService().findUserById(tenantId, new UserId(entityId.getId())) != null; - case DASHBOARD: - return ctx.getDashboardService().findDashboardById(tenantId, new DashboardId(entityId.getId())) != null; - default: - throw new ThingsboardException("Unsupported entity type " + entityId.getEntityType(), ThingsboardErrorCode.INVALID_ARGUMENTS); - } - } - private ConnectResponseMsg processConnect(ConnectRequestMsg request) { Optional optional = ctx.getEdgeService().findEdgeByRoutingKey(TenantId.SYS_TENANT_ID, request.getEdgeRoutingKey()); if (optional.isPresent()) { @@ -1355,15 +919,6 @@ public final class EdgeGrpcSession implements Closeable { .setConfiguration(EdgeConfiguration.getDefaultInstance()).build(); } - private void createRelationFromEdge(EntityId entityId) { - EntityRelation relation = new EntityRelation(); - relation.setFrom(edge.getId()); - relation.setTo(entityId); - relation.setTypeGroup(RelationTypeGroup.COMMON); - relation.setType(EntityRelation.EDGE_TYPE); - ctx.getRelationService().saveRelation(edge.getTenantId(), relation); - } - private EdgeConfiguration constructEdgeConfigProto(Edge edge) throws JsonProcessingException { return EdgeConfiguration.newBuilder() .setEdgeIdMSB(edge.getId().getId().getMostSignificantBits()) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AdminSettingsUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AdminSettingsMsgConstructor.java similarity index 96% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AdminSettingsUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AdminSettingsMsgConstructor.java index 31d6683d89..6a6b67e2e6 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AdminSettingsUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AdminSettingsMsgConstructor.java @@ -23,7 +23,7 @@ import org.thingsboard.server.gen.edge.AdminSettingsUpdateMsg; @Slf4j @Component -public class AdminSettingsUpdateMsgConstructor { +public class AdminSettingsMsgConstructor { public AdminSettingsUpdateMsg constructAdminSettingsUpdateMsg(AdminSettings adminSettings) { AdminSettingsUpdateMsg.Builder builder = AdminSettingsUpdateMsg.newBuilder() diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AlarmUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AlarmMsgConstructor.java similarity index 98% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AlarmUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AlarmMsgConstructor.java index 91f3570c4e..3414727575 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AlarmUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AlarmMsgConstructor.java @@ -34,7 +34,7 @@ import org.thingsboard.server.gen.edge.UpdateMsgType; @Component @Slf4j -public class AlarmUpdateMsgConstructor { +public class AlarmMsgConstructor { @Autowired private DeviceService deviceService; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AssetUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AssetMsgConstructor.java similarity index 98% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AssetUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AssetMsgConstructor.java index ddcb531583..694efcd081 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AssetUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/AssetMsgConstructor.java @@ -27,7 +27,7 @@ import org.thingsboard.server.gen.edge.UpdateMsgType; @Component @Slf4j -public class AssetUpdateMsgConstructor { +public class AssetMsgConstructor { public AssetUpdateMsg constructAssetUpdatedMsg(UpdateMsgType msgType, Asset asset, CustomerId customerId) { AssetUpdateMsg.Builder builder = AssetUpdateMsg.newBuilder() diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/CustomerUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/CustomerMsgConstructor.java similarity index 98% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/CustomerUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/CustomerMsgConstructor.java index a5f54fb443..94687b88d4 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/CustomerUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/CustomerMsgConstructor.java @@ -25,7 +25,7 @@ import org.thingsboard.server.gen.edge.UpdateMsgType; @Component @Slf4j -public class CustomerUpdateMsgConstructor { +public class CustomerMsgConstructor { public CustomerUpdateMsg constructCustomerUpdatedMsg(UpdateMsgType msgType, Customer customer) { CustomerUpdateMsg.Builder builder = CustomerUpdateMsg.newBuilder() diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DashboardUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DashboardMsgConstructor.java similarity index 98% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DashboardUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DashboardMsgConstructor.java index 69b9d8dbd2..a41b987c1c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DashboardUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DashboardMsgConstructor.java @@ -27,7 +27,7 @@ import org.thingsboard.server.gen.edge.UpdateMsgType; @Component @Slf4j -public class DashboardUpdateMsgConstructor { +public class DashboardMsgConstructor { public DashboardUpdateMsg constructDashboardUpdatedMsg(UpdateMsgType msgType, Dashboard dashboard, CustomerId customerId) { DashboardUpdateMsg.Builder builder = DashboardUpdateMsg.newBuilder() diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java similarity index 98% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java index 885511b3db..ebc180d0f9 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java @@ -28,7 +28,7 @@ import org.thingsboard.server.gen.edge.UpdateMsgType; @Component @Slf4j -public class DeviceUpdateMsgConstructor { +public class DeviceMsgConstructor { public DeviceUpdateMsg constructDeviceUpdatedMsg(UpdateMsgType msgType, Device device, CustomerId customerId) { DeviceUpdateMsg.Builder builder = DeviceUpdateMsg.newBuilder() diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityViewUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityViewMsgConstructor.java similarity index 98% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityViewUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityViewMsgConstructor.java index af61e9e4ad..c19a84cfcb 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityViewUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityViewMsgConstructor.java @@ -27,7 +27,7 @@ import org.thingsboard.server.gen.edge.UpdateMsgType; @Component @Slf4j -public class EntityViewUpdateMsgConstructor { +public class EntityViewMsgConstructor { public EntityViewUpdateMsg constructEntityViewUpdatedMsg(UpdateMsgType msgType, EntityView entityView, CustomerId customerId) { EdgeEntityType entityType; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RelationUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RelationMsgConstructor.java similarity index 97% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RelationUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RelationMsgConstructor.java index e09b9f2e8b..5a1f7ce39d 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RelationUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RelationMsgConstructor.java @@ -24,7 +24,7 @@ import org.thingsboard.server.gen.edge.UpdateMsgType; @Component @Slf4j -public class RelationUpdateMsgConstructor { +public class RelationMsgConstructor { public RelationUpdateMsg constructRelationUpdatedMsg(UpdateMsgType msgType, EntityRelation entityRelation) { RelationUpdateMsg.Builder builder = RelationUpdateMsg.newBuilder() diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainMsgConstructor.java similarity index 99% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainMsgConstructor.java index 1d81763a1f..22017c56c2 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/RuleChainMsgConstructor.java @@ -38,7 +38,7 @@ import java.util.List; @Component @Slf4j -public class RuleChainUpdateMsgConstructor { +public class RuleChainMsgConstructor { private static final ObjectMapper objectMapper = new ObjectMapper(); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/UserUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/UserMsgConstructor.java similarity index 98% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/UserUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/UserMsgConstructor.java index 4240d21292..b5b6b080f7 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/UserUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/UserMsgConstructor.java @@ -31,7 +31,7 @@ import java.util.UUID; @Component @Slf4j -public class UserUpdateMsgConstructor { +public class UserMsgConstructor { public UserUpdateMsg constructUserUpdatedMsg(UpdateMsgType msgType, User user, CustomerId customerId) { UserUpdateMsg.Builder builder = UserUpdateMsg.newBuilder() diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetTypeUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetTypeMsgConstructor.java similarity index 98% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetTypeUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetTypeMsgConstructor.java index 784bfd6fdc..e01b2afa66 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetTypeUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetTypeMsgConstructor.java @@ -26,7 +26,7 @@ import org.thingsboard.server.gen.edge.WidgetTypeUpdateMsg; @Component @Slf4j -public class WidgetTypeUpdateMsgConstructor { +public class WidgetTypeMsgConstructor { public WidgetTypeUpdateMsg constructWidgetTypeUpdateMsg(UpdateMsgType msgType, WidgetType widgetType) { WidgetTypeUpdateMsg.Builder builder = WidgetTypeUpdateMsg.newBuilder() diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetsBundleUpdateMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetsBundleMsgConstructor.java similarity index 97% rename from application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetsBundleUpdateMsgConstructor.java rename to application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetsBundleMsgConstructor.java index 4d26e0be59..13dd78980c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetsBundleUpdateMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/WidgetsBundleMsgConstructor.java @@ -26,7 +26,7 @@ import org.thingsboard.server.gen.edge.WidgetsBundleUpdateMsg; @Component @Slf4j -public class WidgetsBundleUpdateMsgConstructor { +public class WidgetsBundleMsgConstructor { public WidgetsBundleUpdateMsg constructWidgetsBundleUpdateMsg(UpdateMsgType msgType, WidgetsBundle widgetsBundle) { WidgetsBundleUpdateMsg.Builder builder = WidgetsBundleUpdateMsg.newBuilder() diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmProcessor.java new file mode 100644 index 0000000000..257bb746c5 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmProcessor.java @@ -0,0 +1,96 @@ +/** + * 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.service.edge.rpc.processor; + +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.alarm.Alarm; +import org.thingsboard.server.common.data.alarm.AlarmSeverity; +import org.thingsboard.server.common.data.alarm.AlarmStatus; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.gen.edge.AlarmUpdateMsg; + +@Component +@Slf4j +public class AlarmProcessor extends BaseProcessor { + + public ListenableFuture onAlarmUpdate(TenantId tenantId, AlarmUpdateMsg alarmUpdateMsg) { + EntityId originatorId = getAlarmOriginator(tenantId, alarmUpdateMsg.getOriginatorName(), + EntityType.valueOf(alarmUpdateMsg.getOriginatorType())); + if (originatorId == null) { + return Futures.immediateFuture(null); + } + try { + Alarm existentAlarm = alarmService.findLatestByOriginatorAndType(tenantId, originatorId, alarmUpdateMsg.getType()).get(); + switch (alarmUpdateMsg.getMsgType()) { + case ENTITY_CREATED_RPC_MESSAGE: + case ENTITY_UPDATED_RPC_MESSAGE: + if (existentAlarm == null || existentAlarm.getStatus().isCleared()) { + existentAlarm = new Alarm(); + existentAlarm.setTenantId(tenantId); + existentAlarm.setType(alarmUpdateMsg.getName()); + existentAlarm.setOriginator(originatorId); + existentAlarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity())); + existentAlarm.setStartTs(alarmUpdateMsg.getStartTs()); + existentAlarm.setClearTs(alarmUpdateMsg.getClearTs()); + existentAlarm.setPropagate(alarmUpdateMsg.getPropagate()); + } + existentAlarm.setStatus(AlarmStatus.valueOf(alarmUpdateMsg.getStatus())); + existentAlarm.setAckTs(alarmUpdateMsg.getAckTs()); + existentAlarm.setEndTs(alarmUpdateMsg.getEndTs()); + existentAlarm.setDetails(mapper.readTree(alarmUpdateMsg.getDetails())); + alarmService.createOrUpdateAlarm(existentAlarm); + break; + case ALARM_ACK_RPC_MESSAGE: + if (existentAlarm != null) { + alarmService.ackAlarm(tenantId, existentAlarm.getId(), alarmUpdateMsg.getAckTs()); + } + break; + case ALARM_CLEAR_RPC_MESSAGE: + if (existentAlarm != null) { + alarmService.clearAlarm(tenantId, existentAlarm.getId(), mapper.readTree(alarmUpdateMsg.getDetails()), alarmUpdateMsg.getAckTs()); + } + break; + case ENTITY_DELETED_RPC_MESSAGE: + if (existentAlarm != null) { + alarmService.deleteAlarm(tenantId, existentAlarm.getId()); + } + break; + } + return Futures.immediateFuture(null); + } catch (Exception e) { + log.error("Failed to process alarm update msg [{}]", alarmUpdateMsg, e); + return Futures.immediateFailedFuture(new RuntimeException("Failed to process alarm update msg", e)); + } + } + + private EntityId getAlarmOriginator(TenantId tenantId, String entityName, EntityType entityType) { + switch (entityType) { + case DEVICE: + return deviceService.findDeviceByTenantIdAndName(tenantId, entityName).getId(); + case ASSET: + return assetService.findAssetByTenantIdAndName(tenantId, entityName).getId(); + case ENTITY_VIEW: + return entityViewService.findEntityViewByTenantIdAndName(tenantId, entityName).getId(); + default: + return null; + } + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseProcessor.java new file mode 100644 index 0000000000..ee5c545999 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseProcessor.java @@ -0,0 +1,112 @@ +/** + * 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.service.edge.rpc.processor; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.common.util.concurrent.ListenableFuture; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.thingsboard.server.common.data.audit.ActionType; +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.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.dao.alarm.AlarmService; +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.EdgeEventService; +import org.thingsboard.server.dao.entityview.EntityViewService; +import org.thingsboard.server.dao.relation.RelationService; +import org.thingsboard.server.dao.user.UserService; +import org.thingsboard.server.service.executors.DbCallbackExecutorService; +import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.service.state.DeviceStateService; + +@Slf4j +public abstract class BaseProcessor { + + protected static final ObjectMapper mapper = new ObjectMapper(); + + @Autowired + protected AlarmService alarmService; + + @Autowired + protected DeviceService deviceService; + + @Autowired + protected DashboardService dashboardService; + + @Autowired + protected AssetService assetService; + + @Autowired + protected EntityViewService entityViewService; + + @Autowired + protected CustomerService customerService; + + @Autowired + protected UserService userService; + + @Autowired + protected RelationService relationService; + + @Autowired + protected DeviceCredentialsService deviceCredentialsService; + + @Autowired + protected AttributesService attributesService; + + @Autowired + protected TbClusterService tbClusterService; + + @Autowired + protected DeviceStateService deviceStateService; + + @Autowired + protected EdgeEventService edgeEventService; + + @Autowired + protected DbCallbackExecutorService dbCallbackExecutorService; + + protected ListenableFuture saveEdgeEvent(TenantId tenantId, + EdgeId edgeId, + EdgeEventType edgeEventType, + ActionType edgeEventAction, + EntityId entityId, + JsonNode entityBody) { + log.debug("Pushing event to edge queue. tenantId [{}], edgeId [{}], edgeEventType[{}], " + + "edgeEventAction [{}], entityId [{}], entityBody [{}]", + tenantId, edgeId, edgeEventType, edgeEventAction, entityId, entityBody); + + EdgeEvent edgeEvent = new EdgeEvent(); + edgeEvent.setTenantId(tenantId); + edgeEvent.setEdgeId(edgeId); + edgeEvent.setEdgeEventType(edgeEventType); + edgeEvent.setEdgeEventAction(edgeEventAction.name()); + if (entityId != null) { + edgeEvent.setEntityId(entityId.getId()); + } + edgeEvent.setEntityBody(entityBody); + return edgeEventService.saveAsync(edgeEvent); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProcessor.java new file mode 100644 index 0000000000..db122e31f9 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProcessor.java @@ -0,0 +1,216 @@ +/** + * 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.service.edge.rpc.processor; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang.RandomStringUtils; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.audit.ActionType; +import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.edge.EdgeEventType; +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.common.data.relation.RelationTypeGroup; +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; +import org.thingsboard.server.gen.edge.DeviceCredentialsUpdateMsg; +import org.thingsboard.server.gen.edge.DeviceUpdateMsg; +import org.thingsboard.server.queue.TbQueueCallback; +import org.thingsboard.server.queue.TbQueueMsgMetadata; + +import java.util.UUID; +import java.util.concurrent.locks.ReentrantLock; + +@Component +@Slf4j +public class DeviceProcessor extends BaseProcessor { + + private static final ReentrantLock deviceCreationLock = new ReentrantLock(); + + public ListenableFuture onDeviceUpdate(TenantId tenantId, Edge edge, DeviceUpdateMsg 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 = deviceService.findDeviceByTenantIdAndName(tenantId, deviceName); + if (device != null) { + // device with this name already exists on the cloud - update ID on the edge + if (!device.getId().equals(edgeDeviceId)) { + saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, ActionType.ENTITY_EXISTS_REQUEST, device.getId(), null); + } + } else { + Device deviceById = deviceService.findDeviceById(edge.getTenantId(), edgeDeviceId); + if (deviceById != null) { + // this ID already used by other device - create new device and update ID on the edge + device = createDevice(tenantId, edge, deviceUpdateMsg); + saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, ActionType.ENTITY_EXISTS_REQUEST, device.getId(), null); + } else { + device = createDevice(tenantId, edge, deviceUpdateMsg); + } + } + // TODO: voba - assign device only in case device is not assigned yet. Missing functionality to check this relation prior assignment + deviceService.assignDeviceToEdge(edge.getTenantId(), device.getId(), edge.getId()); + break; + case ENTITY_UPDATED_RPC_MESSAGE: + updateDevice(tenantId, edge, deviceUpdateMsg); + break; + case ENTITY_DELETED_RPC_MESSAGE: + Device deviceToDelete = deviceService.findDeviceById(tenantId, edgeDeviceId); + if (deviceToDelete != null) { + deviceService.unassignDeviceFromEdge(tenantId, edgeDeviceId, edge.getId()); + } + break; + case UNRECOGNIZED: + log.error("Unsupported msg type {}", deviceUpdateMsg.getMsgType()); + return Futures.immediateFailedFuture(new RuntimeException("Unsupported msg type " + deviceUpdateMsg.getMsgType())); + } + return Futures.immediateFuture(null); + } + + + public ListenableFuture onDeviceCredentialsUpdate(TenantId tenantId, DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg) { + log.debug("Executing onDeviceCredentialsUpdate, deviceCredentialsUpdateMsg [{}]", deviceCredentialsUpdateMsg); + DeviceId deviceId = new DeviceId(new UUID(deviceCredentialsUpdateMsg.getDeviceIdMSB(), deviceCredentialsUpdateMsg.getDeviceIdLSB())); + ListenableFuture deviceFuture = deviceService.findDeviceByIdAsync(tenantId, deviceId); + return Futures.transform(deviceFuture, device -> { + if (device != null) { + log.debug("Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]", + device.getName(), deviceCredentialsUpdateMsg.getCredentialsId(), deviceCredentialsUpdateMsg.getCredentialsValue()); + try { + DeviceCredentials deviceCredentials = deviceCredentialsService.findDeviceCredentialsByDeviceId(tenantId, device.getId()); + deviceCredentials.setCredentialsType(DeviceCredentialsType.valueOf(deviceCredentialsUpdateMsg.getCredentialsType())); + deviceCredentials.setCredentialsId(deviceCredentialsUpdateMsg.getCredentialsId()); + deviceCredentials.setCredentialsValue(deviceCredentialsUpdateMsg.getCredentialsValue()); + deviceCredentialsService.updateDeviceCredentials(tenantId, deviceCredentials); + } catch (Exception e) { + log.error("Can't update device credentials for device [{}], deviceCredentialsUpdateMsg [{}]", device.getName(), deviceCredentialsUpdateMsg, e); + throw new RuntimeException(e); + } + } + return null; + }, dbCallbackExecutorService); + } + + + private void updateDevice(TenantId tenantId, Edge edge, DeviceUpdateMsg deviceUpdateMsg) { + DeviceId deviceId = new DeviceId(new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB())); + Device device = deviceService.findDeviceById(tenantId, deviceId); + device.setName(deviceUpdateMsg.getName()); + device.setType(deviceUpdateMsg.getType()); + device.setLabel(deviceUpdateMsg.getLabel()); + deviceService.saveDevice(device); + + saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, ActionType.CREDENTIALS_REQUEST, deviceId, null); + } + + private Device createDevice(TenantId tenantId, Edge edge, 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 = deviceService.saveDevice(device); + createDeviceCredentials(device); + createRelationFromEdge(tenantId, edge.getId(), device.getId()); + deviceStateService.onDeviceAdded(device); + pushDeviceCreatedEventToRuleEngine(tenantId, edge, device); + + saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, ActionType.CREDENTIALS_REQUEST, deviceId, null); + } finally { + deviceCreationLock.unlock(); + } + return device; + } + + private void createRelationFromEdge(TenantId tenantId, EdgeId edgeId, EntityId entityId) { + EntityRelation relation = new EntityRelation(); + relation.setFrom(edgeId); + relation.setTo(entityId); + relation.setTypeGroup(RelationTypeGroup.COMMON); + relation.setType(EntityRelation.EDGE_TYPE); + relationService.saveRelation(tenantId, relation); + } + + private void createDeviceCredentials(Device device) { + DeviceCredentials deviceCredentials = new DeviceCredentials(); + deviceCredentials.setDeviceId(device.getId()); + deviceCredentials.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN); + deviceCredentials.setCredentialsId(RandomStringUtils.randomAlphanumeric(20)); + deviceCredentialsService.createDeviceCredentials(device.getTenantId(), deviceCredentials); + } + + private void pushDeviceCreatedEventToRuleEngine(TenantId tenantId, Edge edge, Device device) { + try { + DeviceId deviceId = device.getId(); + ObjectNode entityNode = mapper.valueToTree(device); + TbMsg tbMsg = TbMsg.newMsg(DataConstants.ENTITY_CREATED, deviceId, + getActionTbMsgMetaData(edge, device.getCustomerId()), TbMsgDataType.JSON, mapper.writeValueAsString(entityNode)); + tbClusterService.pushMsgToRuleEngine(tenantId, deviceId, tbMsg, new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + // TODO: voba - handle success + log.debug("Successfully send ENTITY_CREATED EVENT to rule engine [{}]", device); + } + + @Override + public void onFailure(Throwable t) { + // TODO: voba - handle failure + log.debug("Failed to send ENTITY_CREATED EVENT to rule engine [{}]", device, t); + } + + ; + }); + } catch (JsonProcessingException | IllegalArgumentException e) { + log.warn("[{}] Failed to push device action to rule engine: {}", device.getId(), DataConstants.ENTITY_CREATED, e); + } + } + + + private TbMsgMetaData getActionTbMsgMetaData(Edge edge, CustomerId customerId) { + TbMsgMetaData metaData = getTbMsgMetaData(edge); + if (customerId != null && !customerId.isNullUid()) { + metaData.putValue("customerId", customerId.toString()); + } + return metaData; + } + + + private TbMsgMetaData getTbMsgMetaData(Edge edge) { + TbMsgMetaData metaData = new TbMsgMetaData(); + metaData.putValue("edgeId", edge.getId().toString()); + metaData.putValue("edgeName", edge.getName()); + return metaData; + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationProcessor.java new file mode 100644 index 0000000000..1ef4044d75 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationProcessor.java @@ -0,0 +1,102 @@ +/** + * 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.service.edge.rpc.processor; + +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; +import org.thingsboard.server.common.data.exception.ThingsboardException; +import org.thingsboard.server.common.data.id.AssetId; +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.DashboardId; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.EntityIdFactory; +import org.thingsboard.server.common.data.id.EntityViewId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.UserId; +import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.common.data.relation.RelationTypeGroup; +import org.thingsboard.server.gen.edge.RelationUpdateMsg; + +import java.util.UUID; + +@Component +@Slf4j +public class RelationProcessor extends BaseProcessor { + + public ListenableFuture onRelationUpdate(TenantId tenantId, RelationUpdateMsg relationUpdateMsg) { + log.info("onRelationUpdate {}", relationUpdateMsg); + try { + EntityRelation entityRelation = new EntityRelation(); + + UUID fromUUID = new UUID(relationUpdateMsg.getFromIdMSB(), relationUpdateMsg.getFromIdLSB()); + EntityId fromId = EntityIdFactory.getByTypeAndUuid(EntityType.valueOf(relationUpdateMsg.getFromEntityType()), fromUUID); + entityRelation.setFrom(fromId); + + UUID toUUID = new UUID(relationUpdateMsg.getToIdMSB(), relationUpdateMsg.getToIdLSB()); + EntityId toId = EntityIdFactory.getByTypeAndUuid(EntityType.valueOf(relationUpdateMsg.getToEntityType()), toUUID); + entityRelation.setTo(toId); + + entityRelation.setType(relationUpdateMsg.getType()); + entityRelation.setTypeGroup(RelationTypeGroup.valueOf(relationUpdateMsg.getTypeGroup())); + entityRelation.setAdditionalInfo(mapper.readTree(relationUpdateMsg.getAdditionalInfo())); + switch (relationUpdateMsg.getMsgType()) { + case ENTITY_CREATED_RPC_MESSAGE: + case ENTITY_UPDATED_RPC_MESSAGE: + if (isEntityExists(tenantId, entityRelation.getTo()) + && isEntityExists(tenantId, entityRelation.getFrom())) { + relationService.saveRelationAsync(tenantId, entityRelation); + } + break; + case ENTITY_DELETED_RPC_MESSAGE: + relationService.deleteRelation(tenantId, entityRelation); + break; + case UNRECOGNIZED: + log.error("Unsupported msg type"); + } + return Futures.immediateFuture(null); + } catch (Exception e) { + log.error("Failed to process relation update msg [{}]", relationUpdateMsg, e); + return Futures.immediateFailedFuture(new RuntimeException("Failed to process relation update msg", e)); + } + } + + + private boolean isEntityExists(TenantId tenantId, EntityId entityId) throws ThingsboardException { + switch (entityId.getEntityType()) { + case DEVICE: + return deviceService.findDeviceById(tenantId, new DeviceId(entityId.getId())) != null; + case ASSET: + return assetService.findAssetById(tenantId, new AssetId(entityId.getId())) != null; + case ENTITY_VIEW: + return entityViewService.findEntityViewById(tenantId, new EntityViewId(entityId.getId())) != null; + case CUSTOMER: + return customerService.findCustomerById(tenantId, new CustomerId(entityId.getId())) != null; + case USER: + return userService.findUserById(tenantId, new UserId(entityId.getId())) != null; + case DASHBOARD: + return dashboardService.findDashboardById(tenantId, new DashboardId(entityId.getId())) != null; + default: + throw new ThingsboardException("Unsupported entity type " + entityId.getEntityType(), ThingsboardErrorCode.INVALID_ARGUMENTS); + } + } + + +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryProcessor.java new file mode 100644 index 0000000000..506d9f6ba7 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryProcessor.java @@ -0,0 +1,202 @@ +/** + * 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.service.edge.rpc.processor; + +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.SettableFuture; +import com.google.gson.Gson; +import com.google.gson.JsonObject; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; +import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg; +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.asset.Asset; +import org.thingsboard.server.common.data.id.AssetId; +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.DashboardId; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.EntityViewId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.UserId; +import org.thingsboard.server.common.data.kv.AttributeKey; +import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.common.msg.session.SessionMsgType; +import org.thingsboard.server.common.transport.util.JsonUtils; +import org.thingsboard.server.gen.edge.AttributeDeleteMsg; +import org.thingsboard.server.gen.edge.EntityDataProto; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.queue.TbQueueCallback; +import org.thingsboard.server.queue.TbQueueMsgMetadata; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.UUID; + +@Component +@Slf4j +public class TelemetryProcessor extends BaseProcessor { + + private final Gson gson = new Gson(); + + public List> onTelemetryUpdate(TenantId tenantId, EntityDataProto entityData) { + List> result = new ArrayList<>(); + EntityId entityId = constructEntityId(entityData); + if ((entityData.hasPostAttributesMsg() || entityData.hasPostTelemetryMsg()) && entityId != null) { + TbMsgMetaData metaData = constructBaseMsgMetadata(tenantId, entityId); + metaData.putValue(DataConstants.MSG_SOURCE_KEY, DataConstants.EDGE_MSG_SOURCE); + if (entityData.hasPostAttributesMsg()) { + metaData.putValue("scope", entityData.getPostAttributeScope()); + result.add(processPostAttributes(tenantId, entityId, entityData.getPostAttributesMsg(), metaData)); + } + if (entityData.hasPostTelemetryMsg()) { + result.add(processPostTelemetry(tenantId, entityId, entityData.getPostTelemetryMsg(), metaData)); + } + } + if (entityData.hasAttributeDeleteMsg()) { + result.add(processAttributeDeleteMsg(tenantId, entityId, entityData.getAttributeDeleteMsg(), entityData.getEntityType())); + } + return result; + } + + private TbMsgMetaData constructBaseMsgMetadata(TenantId tenantId, EntityId entityId) { + TbMsgMetaData metaData = new TbMsgMetaData(); + switch (entityId.getEntityType()) { + case DEVICE: + Device device = deviceService.findDeviceById(tenantId, new DeviceId(entityId.getId())); + if (device != null) { + metaData.putValue("deviceName", device.getName()); + metaData.putValue("deviceType", device.getType()); + } + break; + case ASSET: + Asset asset = assetService.findAssetById(tenantId, new AssetId(entityId.getId())); + if (asset != null) { + metaData.putValue("assetName", asset.getName()); + metaData.putValue("assetType", asset.getType()); + } + break; + case ENTITY_VIEW: + EntityView entityView = entityViewService.findEntityViewById(tenantId, new EntityViewId(entityId.getId())); + if (entityView != null) { + metaData.putValue("entityViewName", entityView.getName()); + metaData.putValue("entityViewType", entityView.getType()); + } + break; + default: + log.debug("Using empty metadata for entityId [{}]", entityId); + break; + } + return metaData; + } + + private ListenableFuture processPostTelemetry(TenantId tenantId, EntityId entityId, TransportProtos.PostTelemetryMsg msg, TbMsgMetaData metaData) { + SettableFuture futureToSet = SettableFuture.create(); + for (TransportProtos.TsKvListProto tsKv : msg.getTsKvListList()) { + JsonObject json = JsonUtils.getJsonObject(tsKv.getKvList()); + metaData.putValue("ts", tsKv.getTs() + ""); + TbMsg tbMsg = TbMsg.newMsg(SessionMsgType.POST_TELEMETRY_REQUEST.name(), entityId, metaData, gson.toJson(json)); + tbClusterService.pushMsgToRuleEngine(tenantId, tbMsg.getOriginator(), tbMsg, new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + futureToSet.set(null); + } + + @Override + public void onFailure(Throwable t) { + log.error("Can't process post telemetry [{}]", msg, t); + futureToSet.setException(t); + } + }); + } + return futureToSet; + } + + private ListenableFuture processPostAttributes(TenantId tenantId, EntityId entityId, TransportProtos.PostAttributeMsg msg, TbMsgMetaData metaData) { + SettableFuture futureToSet = SettableFuture.create(); + JsonObject json = JsonUtils.getJsonObject(msg.getKvList()); + TbMsg tbMsg = TbMsg.newMsg(SessionMsgType.POST_ATTRIBUTES_REQUEST.name(), entityId, metaData, gson.toJson(json)); + tbClusterService.pushMsgToRuleEngine(tenantId, tbMsg.getOriginator(), tbMsg, new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + futureToSet.set(null); + } + + @Override + public void onFailure(Throwable t) { + log.error("Can't process post attributes [{}]", msg, t); + futureToSet.setException(t); + } + }); + return futureToSet; + } + + private ListenableFuture processAttributeDeleteMsg(TenantId tenantId, EntityId entityId, AttributeDeleteMsg attributeDeleteMsg, String entityType) { + SettableFuture futureToSet = SettableFuture.create(); + String scope = attributeDeleteMsg.getScope(); + List attributeNames = attributeDeleteMsg.getAttributeNamesList(); + attributesService.removeAll(tenantId, entityId, scope, attributeNames); + if (EntityType.DEVICE.name().equals(entityType)) { + Set attributeKeys = new HashSet<>(); + for (String attributeName : attributeNames) { + attributeKeys.add(new AttributeKey(scope, attributeName)); + } + tbClusterService.pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete( + tenantId, (DeviceId) entityId, attributeKeys), new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + futureToSet.set(null); + } + + @Override + public void onFailure(Throwable t) { + log.error("Can't process attribute delete msg [{}]", attributeDeleteMsg, t); + futureToSet.setException(t); + } + }); + } + return futureToSet; + } + + private EntityId constructEntityId(EntityDataProto entityData) { + EntityType entityType = EntityType.valueOf(entityData.getEntityType()); + switch (entityType) { + case DEVICE: + return new DeviceId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); + case ASSET: + return new AssetId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); + case ENTITY_VIEW: + return new EntityViewId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); + case DASHBOARD: + return new DashboardId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); + case TENANT: + return new TenantId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); + case CUSTOMER: + return new CustomerId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); + case USER: + return new UserId(new UUID(entityData.getEntityIdMSB(), entityData.getEntityIdLSB())); + default: + log.warn("Unsupported entity type [{}] during construct of entity id. EntityDataProto [{}]", entityData.getEntityType(), entityData); + return null; + } + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/audit/ActionType.java b/common/data/src/main/java/org/thingsboard/server/common/data/audit/ActionType.java index 84273148c5..158bf62e4c 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/audit/ActionType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/audit/ActionType.java @@ -44,7 +44,8 @@ public enum ActionType { LOCKOUT(false), ASSIGNED_TO_EDGE(false), // log edge name UNASSIGNED_FROM_EDGE(false), // log edge name - CREDENTIALS_REQUEST(false); // request credentials from edge + CREDENTIALS_REQUEST(false), // request credentials from edge + ENTITY_EXISTS_REQUEST(false); // request to recreate entity on edge private final boolean isRead;