diff --git a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java index 198c433416..934b8616de 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java @@ -40,8 +40,19 @@ import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.rpc.processor.AlarmEdgeProcessor; -import org.thingsboard.server.service.edge.rpc.processor.EntityEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.AssetEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.CustomerEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.DeviceEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.DeviceProfileEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.EdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.EntityViewEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.OtaPackageEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.QueueEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.RelationEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.RuleChainEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.UserEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.WidgetBundleEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.WidgetTypeEdgeProcessor; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; @@ -64,7 +75,43 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { private TbClusterService clusterService; @Autowired - private EntityEdgeProcessor entityProcessor; + private EdgeProcessor edgeProcessor; + + @Autowired + private AssetEdgeProcessor assetProcessor; + + @Autowired + private DeviceEdgeProcessor deviceProcessor; + + @Autowired + private EntityViewEdgeProcessor entityViewProcessor; + + @Autowired + private DeviceEdgeProcessor dashboardProcessor; + + @Autowired + private RuleChainEdgeProcessor ruleChainProcessor; + + @Autowired + private UserEdgeProcessor userProcessor; + + @Autowired + private CustomerEdgeProcessor customerProcessor; + + @Autowired + private DeviceProfileEdgeProcessor deviceProfileProcessor; + + @Autowired + private OtaPackageEdgeProcessor otaPackageProcessor; + + @Autowired + private WidgetBundleEdgeProcessor widgetBundleProcessor; + + @Autowired + private WidgetTypeEdgeProcessor widgetTypeProcessor; + + @Autowired + private QueueEdgeProcessor queueEdgeProcessor; @Autowired private AlarmEdgeProcessor alarmProcessor; @@ -120,21 +167,43 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { ListenableFuture future; switch (type) { case EDGE: + future = edgeProcessor.processEdgeNotification(tenantId, edgeNotificationMsg); + break; case ASSET: + future = assetProcessor.processAssetNotification(tenantId, edgeNotificationMsg); + break; case DEVICE: + future = deviceProcessor.processDeviceNotification(tenantId, edgeNotificationMsg); + break; case ENTITY_VIEW: + future = entityViewProcessor.processEntityViewNotification(tenantId, edgeNotificationMsg); + break; case DASHBOARD: + future = dashboardProcessor.processDashboardNotification(tenantId, edgeNotificationMsg); + break; case RULE_CHAIN: - future = entityProcessor.processEntityNotification(tenantId, edgeNotificationMsg); + future = ruleChainProcessor.processRuleChainNotification(tenantId, edgeNotificationMsg); break; case USER: + future = userProcessor.processUserNotification(tenantId, edgeNotificationMsg); + break; case CUSTOMER: + future = customerProcessor.processCustomerNotification(tenantId, edgeNotificationMsg); + break; case DEVICE_PROFILE: + future = deviceProfileProcessor.processDeviceProfileNotification(tenantId, edgeNotificationMsg); + break; case OTA_PACKAGE: + future = otaPackageProcessor.processOtaPackageNotification(tenantId, edgeNotificationMsg); + break; case WIDGETS_BUNDLE: + future = widgetBundleProcessor.processWidgetsBundleNotification(tenantId, edgeNotificationMsg); + break; case WIDGET_TYPE: + future = widgetTypeProcessor.processWidgetTypeNotification(tenantId, edgeNotificationMsg); + break; case QUEUE: - future = entityProcessor.processEntityNotificationForAllEdges(tenantId, edgeNotificationMsg); + future = queueEdgeProcessor.processQueueNotification(tenantId, edgeNotificationMsg); break; case ALARM: future = alarmProcessor.processAlarmNotification(tenantId, edgeNotificationMsg); 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 089cb6d794..0a771ce7e4 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 @@ -45,7 +45,6 @@ import org.thingsboard.server.service.edge.rpc.processor.DashboardEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.DeviceEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.DeviceProfileEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.EdgeProcessor; -import org.thingsboard.server.service.edge.rpc.processor.EntityEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.EntityViewEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.OtaPackageEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.QueueEdgeProcessor; @@ -125,9 +124,6 @@ public class EdgeContextComponent { @Autowired private DeviceEdgeProcessor deviceProcessor; - @Autowired - private EntityEdgeProcessor entityProcessor; - @Autowired private AssetEdgeProcessor assetProcessor; 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 8d4aa06082..39bcd22743 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 @@ -28,7 +28,6 @@ import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEvent; -import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; @@ -60,7 +59,6 @@ import org.thingsboard.server.gen.edge.v1.RequestMsgType; import org.thingsboard.server.gen.edge.v1.ResponseMsg; import org.thingsboard.server.gen.edge.v1.RuleChainMetadataRequestMsg; import org.thingsboard.server.gen.edge.v1.SyncCompletedMsg; -import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.edge.v1.UplinkMsg; import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg; import org.thingsboard.server.gen.edge.v1.UserCredentialsRequestMsg; @@ -428,7 +426,6 @@ public final class EdgeGrpcSession implements Closeable { TimeUnit.MILLISECONDS) ); } - } private DownlinkMsg convertToDownlinkMsg(EdgeEvent edgeEvent) { @@ -448,23 +445,17 @@ public final class EdgeGrpcSession implements Closeable { case RELATION_DELETED: case ASSIGNED_TO_CUSTOMER: case UNASSIGNED_FROM_CUSTOMER: - downlinkMsg = processEntityMessage(edgeEvent); + case CREDENTIALS_REQUEST: + case ENTITY_MERGE_REQUEST: + case RPC_CALL: + downlinkMsg = convertEntityEventToDownlink(edgeEvent); log.trace("[{}][{}] entity message processed [{}]", edgeEvent.getTenantId(), this.sessionId, downlinkMsg); break; case ATTRIBUTES_UPDATED: case POST_ATTRIBUTES: case ATTRIBUTES_DELETED: case TIMESERIES_UPDATED: - downlinkMsg = ctx.getTelemetryProcessor().processTelemetryMessageToEdge(edgeEvent); - break; - case CREDENTIALS_REQUEST: - downlinkMsg = ctx.getEntityProcessor().processCredentialsRequestMessageToEdge(edgeEvent); - break; - case ENTITY_MERGE_REQUEST: - downlinkMsg = ctx.getEntityProcessor().processEntityMergeRequestMessageToEdge(edge, edgeEvent); - break; - case RPC_CALL: - downlinkMsg = ctx.getDeviceProcessor().processRpcCallMsgToEdge(edgeEvent); + downlinkMsg = ctx.getTelemetryProcessor().convertTelemetryEventToDownlink(edgeEvent); break; default: log.warn("[{}][{}] Unsupported action type [{}]", edge.getTenantId(), this.sessionId, edgeEvent.getAction()); @@ -504,43 +495,43 @@ public final class EdgeGrpcSession implements Closeable { return ctx.getAttributesService().save(edge.getTenantId(), edge.getId(), DataConstants.SERVER_SCOPE, attributes); } - private DownlinkMsg processEntityMessage(EdgeEvent edgeEvent) { - log.trace("Executing processEntityMessage, edgeEvent [{}], action [{}]", edgeEvent, edgeEvent.getAction()); + private DownlinkMsg convertEntityEventToDownlink(EdgeEvent edgeEvent) { + log.trace("Executing convertEntityEventToDownlink, edgeEvent [{}], action [{}]", edgeEvent, edgeEvent.getAction()); switch (edgeEvent.getType()) { case EDGE: - return ctx.getEdgeProcessor().processEdgeToEdge(edgeEvent); + return ctx.getEdgeProcessor().convertEdgeEventToDownlink(edgeEvent); case DEVICE: - return ctx.getDeviceProcessor().processDeviceToEdge(edgeEvent); + return ctx.getDeviceProcessor().convertDeviceEventToDownlink(edgeEvent); case DEVICE_PROFILE: - return ctx.getDeviceProfileProcessor().processDeviceProfileToEdge(edgeEvent); + return ctx.getDeviceProfileProcessor().convertDeviceProfileEventToDownlink(edgeEvent); case ASSET: - return ctx.getAssetProcessor().processAssetToEdge(edgeEvent); + return ctx.getAssetProcessor().convertAssetEventToDownlink(edgeEvent); case ENTITY_VIEW: - return ctx.getEntityViewProcessor().processEntityViewToEdge(edgeEvent); + return ctx.getEntityViewProcessor().convertEntityViewEventToDownlink(edgeEvent); case DASHBOARD: - return ctx.getDashboardProcessor().processDashboardToEdge(edgeEvent); + return ctx.getDashboardProcessor().convertDashboardEventToDownlink(edgeEvent); case CUSTOMER: - return ctx.getCustomerProcessor().processCustomerToEdge(edgeEvent); + return ctx.getCustomerProcessor().convertCustomerEventToDownlink(edgeEvent); case RULE_CHAIN: - return ctx.getRuleChainProcessor().processRuleChainToEdge(edge, edgeEvent); + return ctx.getRuleChainProcessor().convertRuleChainEventToDownlink(edge, edgeEvent); case RULE_CHAIN_METADATA: - return ctx.getRuleChainProcessor().processRuleChainMetadataToEdge(edgeEvent, this.edgeVersion); + return ctx.getRuleChainProcessor().convertRuleChainMetadataEventToDownlink(edgeEvent, this.edgeVersion); case ALARM: - return ctx.getAlarmProcessor().processAlarmToEdge(edgeEvent); + return ctx.getAlarmProcessor().convertAlarmEventToDownlink(edgeEvent); case USER: - return ctx.getUserProcessor().processUserToEdge(edgeEvent); + return ctx.getUserProcessor().convertUserEventToDownlink(edgeEvent); case RELATION: - return ctx.getRelationProcessor().processRelationToEdge(edgeEvent); + return ctx.getRelationProcessor().convertRelationEventToDownlink(edgeEvent); case WIDGETS_BUNDLE: - return ctx.getWidgetBundleProcessor().processWidgetsBundleToEdge(edgeEvent); + return ctx.getWidgetBundleProcessor().convertWidgetsBundleEventToDownlink(edgeEvent); case WIDGET_TYPE: - return ctx.getWidgetTypeProcessor().processWidgetTypeToEdge(edgeEvent); + return ctx.getWidgetTypeProcessor().convertWidgetTypeEventToDownlink(edgeEvent); case ADMIN_SETTINGS: - return ctx.getAdminSettingsProcessor().processAdminSettingsToEdge(edgeEvent); + return ctx.getAdminSettingsProcessor().convertAdminSettingsEventToDownlink(edgeEvent); case OTA_PACKAGE: - return ctx.getOtaPackageEdgeProcessor().processOtaPackageToEdge(edgeEvent); + return ctx.getOtaPackageEdgeProcessor().convertOtaPackageEventToDownlink(edgeEvent); case QUEUE: - return ctx.getQueueEdgeProcessor().processQueueToEdge(edgeEvent); + return ctx.getQueueEdgeProcessor().convertQueueEventToDownlink(edgeEvent); default: log.warn("Unsupported edge event type [{}]", edgeEvent); return null; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AdminSettingsEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AdminSettingsEdgeProcessor.java index db30ee2036..c3d3df4492 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AdminSettingsEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AdminSettingsEdgeProcessor.java @@ -30,7 +30,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class AdminSettingsEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processAdminSettingsToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertAdminSettingsEventToDownlink(EdgeEvent edgeEvent) { AdminSettings adminSettings = JacksonUtil.OBJECT_MAPPER.convertValue(edgeEvent.getBody(), AdminSettings.class); AdminSettingsUpdateMsg adminSettingsUpdateMsg = adminSettingsMsgConstructor.constructAdminSettingsUpdateMsg(adminSettings); return DownlinkMsg.newBuilder() diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java index 999f8a0f86..b5d256bfa8 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java @@ -116,7 +116,7 @@ public class AlarmEdgeProcessor extends BaseEdgeProcessor { } } - public DownlinkMsg processAlarmToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertAlarmEventToDownlink(EdgeEvent edgeEvent) { AlarmId alarmId = new AlarmId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AssetEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AssetEdgeProcessor.java index 51ce3a214a..908d9e571b 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AssetEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AssetEdgeProcessor.java @@ -15,15 +15,18 @@ */ package org.thingsboard.server.service.edge.rpc.processor; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.AssetId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.gen.edge.v1.AssetUpdateMsg; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @Component @@ -31,7 +34,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class AssetEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processAssetToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertAssetEventToDownlink(EdgeEvent edgeEvent) { AssetId assetId = new AssetId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -63,4 +66,8 @@ public class AssetEdgeProcessor extends BaseEdgeProcessor { } return downlinkMsg; } + + public ListenableFuture processAssetNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotification(tenantId, edgeNotificationMsg); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java index ded5a0322a..9ff48204b1 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java @@ -30,9 +30,13 @@ import org.thingsboard.server.common.data.edge.EdgeEventActionType; 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.EntityIdFactory; +import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.common.data.rule.RuleChain; +import org.thingsboard.server.common.data.rule.RuleChainConnectionInfo; import org.thingsboard.server.dao.alarm.AlarmService; import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.dao.attributes.AttributesService; @@ -54,6 +58,7 @@ 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.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.service.edge.rpc.constructor.AdminSettingsMsgConstructor; @@ -79,6 +84,7 @@ import org.thingsboard.server.service.state.DeviceStateService; import java.util.ArrayList; import java.util.List; +import java.util.UUID; @Slf4j public abstract class BaseEdgeProcessor { @@ -292,4 +298,111 @@ public abstract class BaseEdgeProcessor { throw new RuntimeException("Unsupported actionType [" + actionType + "]"); } } + + protected ListenableFuture processEntityNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()); + EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); + EntityId entityId = EntityIdFactory.getByEdgeEventTypeAndUuid(type, + new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); + EdgeId edgeId = safeGetEdgeId(edgeNotificationMsg); + switch (actionType) { + case ADDED: + case UPDATED: + case CREDENTIALS_UPDATED: + case ASSIGNED_TO_CUSTOMER: + case UNASSIGNED_FROM_CUSTOMER: + case DELETED: + if (edgeId != null) { + return saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null); + } else { + return pushNotificationToAllRelatedEdges(tenantId, entityId, type, actionType); + } + case ASSIGNED_TO_EDGE: + case UNASSIGNED_FROM_EDGE: + ListenableFuture future = saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null); + return Futures.transformAsync(future, unused -> { + if (type.equals(EdgeEventType.RULE_CHAIN)) { + return updateDependentRuleChains(tenantId, new RuleChainId(entityId.getId()), edgeId); + } else { + return Futures.immediateFuture(null); + } + }, dbCallbackExecutorService); + default: + return Futures.immediateFuture(null); + } + } + + private EdgeId safeGetEdgeId(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + if (edgeNotificationMsg.getEdgeIdMSB() != 0 && edgeNotificationMsg.getEdgeIdLSB() != 0) { + return new EdgeId(new UUID(edgeNotificationMsg.getEdgeIdMSB(), edgeNotificationMsg.getEdgeIdLSB())); + } else { + return null; + } + } + + private ListenableFuture pushNotificationToAllRelatedEdges(TenantId tenantId, EntityId entityId, EdgeEventType type, EdgeEventActionType actionType) { + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); + PageData pageData; + List> futures = new ArrayList<>(); + do { + pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink); + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + for (EdgeId relatedEdgeId : pageData.getData()) { + futures.add(saveEdgeEvent(tenantId, relatedEdgeId, type, actionType, entityId, null)); + } + if (pageData.hasNext()) { + pageLink = pageLink.nextPageLink(); + } + } + } while (pageData != null && pageData.hasNext()); + return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); + } + + private ListenableFuture updateDependentRuleChains(TenantId tenantId, RuleChainId processingRuleChainId, EdgeId edgeId) { + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); + PageData pageData; + List> futures = new ArrayList<>(); + do { + pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink); + if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { + for (RuleChain ruleChain : pageData.getData()) { + if (!ruleChain.getId().equals(processingRuleChainId)) { + List connectionInfos = + ruleChainService.loadRuleChainMetaData(ruleChain.getTenantId(), ruleChain.getId()).getRuleChainConnections(); + if (connectionInfos != null && !connectionInfos.isEmpty()) { + for (RuleChainConnectionInfo connectionInfo : connectionInfos) { + if (connectionInfo.getTargetRuleChainId().equals(processingRuleChainId)) { + futures.add(saveEdgeEvent(tenantId, + edgeId, + EdgeEventType.RULE_CHAIN_METADATA, + EdgeEventActionType.UPDATED, + ruleChain.getId(), + null)); + } + } + } + } + } + if (pageData.hasNext()) { + pageLink = pageLink.nextPageLink(); + } + } + } while (pageData != null && pageData.hasNext()); + return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); + } + + protected ListenableFuture processEntityNotificationForAllEdges(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()); + EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); + EntityId entityId = EntityIdFactory.getByEdgeEventTypeAndUuid(type, new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); + switch (actionType) { + case ADDED: + case UPDATED: + case DELETED: + case CREDENTIALS_UPDATED: // used by USER entity + return processActionForAllEdges(tenantId, type, actionType, entityId); + default: + return Futures.immediateFuture(null); + } + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/CustomerEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/CustomerEdgeProcessor.java index 0e2f7f57d4..ec32203d19 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/CustomerEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/CustomerEdgeProcessor.java @@ -15,15 +15,18 @@ */ package org.thingsboard.server.service.edge.rpc.processor; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.gen.edge.v1.CustomerUpdateMsg; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @Component @@ -31,7 +34,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class CustomerEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processCustomerToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertCustomerEventToDownlink(EdgeEvent edgeEvent) { CustomerId customerId = new CustomerId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -60,4 +63,7 @@ public class CustomerEdgeProcessor extends BaseEdgeProcessor { return downlinkMsg; } + public ListenableFuture processCustomerNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotificationForAllEdges(tenantId, edgeNotificationMsg); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DashboardEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DashboardEdgeProcessor.java index aec157fbfe..f8ddab678d 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DashboardEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DashboardEdgeProcessor.java @@ -31,7 +31,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class DashboardEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processDashboardToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertDashboardEventToDownlink(EdgeEvent edgeEvent) { DashboardId dashboardId = new DashboardId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java index 49db071f7b..21df8d53c7 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java @@ -52,11 +52,13 @@ import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.dao.exception.DataValidationException; +import org.thingsboard.server.gen.edge.v1.DeviceCredentialsRequestMsg; import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.v1.DeviceRpcCallMsg; import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.TbQueueCallback; import org.thingsboard.server.queue.TbQueueMsgMetadata; import org.thingsboard.server.queue.util.DataDecodingEncodingService; @@ -356,7 +358,7 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { return futureToSet; } - public DownlinkMsg processDeviceToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertDeviceEventToDownlink(EdgeEvent edgeEvent) { DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -396,12 +398,18 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { .build(); } break; + case RPC_CALL: + return convertRpcCallEventToDownlink(edgeEvent); + case CREDENTIALS_REQUEST: + return convertCredentialsRequestEventToDownlink(edgeEvent); + case ENTITY_MERGE_REQUEST: + return convertEntityMergeRequestEventToDownlink(edgeEvent); } return downlinkMsg; } - public DownlinkMsg processRpcCallMsgToEdge(EdgeEvent edgeEvent) { - log.trace("Executing processRpcCall, edgeEvent [{}]", edgeEvent); + private DownlinkMsg convertRpcCallEventToDownlink(EdgeEvent edgeEvent) { + log.trace("Executing convertRpcCallEventToDownlink, edgeEvent [{}]", edgeEvent); DeviceRpcCallMsg deviceRpcCallMsg = deviceMsgConstructor.constructDeviceRpcCallMsg(edgeEvent.getEntityId(), edgeEvent.getBody()); return DownlinkMsg.newBuilder() @@ -409,4 +417,39 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { .addDeviceRpcCallMsg(deviceRpcCallMsg) .build(); } + + private DownlinkMsg convertCredentialsRequestEventToDownlink(EdgeEvent edgeEvent) { + DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); + DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = DeviceCredentialsRequestMsg.newBuilder() + .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) + .build(); + DownlinkMsg.Builder builder = DownlinkMsg.newBuilder() + .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) + .addDeviceCredentialsRequestMsg(deviceCredentialsRequestMsg); + return builder.build(); + } + + public DownlinkMsg convertEntityMergeRequestEventToDownlink(EdgeEvent edgeEvent) { + DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); + Device device = deviceService.findDeviceById(edgeEvent.getTenantId(), deviceId); + String conflictName = null; + if(edgeEvent.getBody() != null) { + conflictName = edgeEvent.getBody().get("conflictName").asText(); + } + DeviceUpdateMsg deviceUpdateMsg = deviceMsgConstructor + .constructDeviceUpdatedMsg(UpdateMsgType.ENTITY_MERGE_RPC_MESSAGE, device, conflictName); + return DownlinkMsg.newBuilder() + .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) + .addDeviceUpdateMsg(deviceUpdateMsg) + .build(); + } + + public ListenableFuture processDeviceNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotification(tenantId, edgeNotificationMsg); + } + + public ListenableFuture processDashboardNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processDeviceNotification(tenantId, edgeNotificationMsg); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProfileEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProfileEdgeProcessor.java index 46732d4a03..d5d36ee640 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProfileEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceProfileEdgeProcessor.java @@ -15,15 +15,18 @@ */ package org.thingsboard.server.service.edge.rpc.processor; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.DeviceProfileId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.gen.edge.v1.DeviceProfileUpdateMsg; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @Component @@ -31,7 +34,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class DeviceProfileEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processDeviceProfileToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertDeviceProfileEventToDownlink(EdgeEvent edgeEvent) { DeviceProfileId deviceProfileId = new DeviceProfileId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -60,4 +63,7 @@ public class DeviceProfileEdgeProcessor extends BaseEdgeProcessor { return downlinkMsg; } + public ListenableFuture processDeviceProfileNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotificationForAllEdges(tenantId, edgeNotificationMsg); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeProcessor.java index d696baa027..0aaedfe2f9 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EdgeProcessor.java @@ -15,14 +15,17 @@ */ package org.thingsboard.server.service.edge.rpc.processor; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.EdgeConfiguration; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @Component @@ -30,7 +33,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class EdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processEdgeToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertEdgeEventToDownlink(EdgeEvent edgeEvent) { EdgeId edgeId = new EdgeId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -49,4 +52,8 @@ public class EdgeProcessor extends BaseEdgeProcessor { } return downlinkMsg; } + + public ListenableFuture processEdgeNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotification(tenantId, edgeNotificationMsg); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityEdgeProcessor.java deleted file mode 100644 index 530b1d1cb5..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityEdgeProcessor.java +++ /dev/null @@ -1,198 +0,0 @@ -/** - * Copyright © 2016-2022 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.common.util.JacksonUtil; -import org.thingsboard.server.common.data.Device; -import org.thingsboard.server.common.data.EdgeUtils; -import org.thingsboard.server.common.data.edge.Edge; -import org.thingsboard.server.common.data.edge.EdgeEvent; -import org.thingsboard.server.common.data.edge.EdgeEventActionType; -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.EntityIdFactory; -import org.thingsboard.server.common.data.id.RuleChainId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.page.PageData; -import org.thingsboard.server.common.data.page.PageLink; -import org.thingsboard.server.common.data.rule.RuleChain; -import org.thingsboard.server.common.data.rule.RuleChainConnectionInfo; -import org.thingsboard.server.gen.edge.v1.DeviceCredentialsRequestMsg; -import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg; -import org.thingsboard.server.gen.edge.v1.DownlinkMsg; -import org.thingsboard.server.gen.edge.v1.UpdateMsgType; -import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.util.TbCoreComponent; - -import java.util.ArrayList; -import java.util.List; -import java.util.UUID; - -@Component -@Slf4j -@TbCoreComponent -public class EntityEdgeProcessor extends BaseEdgeProcessor { - - public DownlinkMsg processEntityMergeRequestMessageToEdge(Edge edge, EdgeEvent edgeEvent) { - DownlinkMsg downlinkMsg = null; - if (EdgeEventType.DEVICE.equals(edgeEvent.getType())) { - DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); - Device device = deviceService.findDeviceById(edge.getTenantId(), deviceId); - String conflictName = null; - if(edgeEvent.getBody() != null) { - conflictName = edgeEvent.getBody().get("conflictName").asText(); - } - DeviceUpdateMsg deviceUpdateMsg = deviceMsgConstructor - .constructDeviceUpdatedMsg(UpdateMsgType.ENTITY_MERGE_RPC_MESSAGE, device, conflictName); - downlinkMsg = DownlinkMsg.newBuilder() - .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) - .addDeviceUpdateMsg(deviceUpdateMsg) - .build(); - } - return downlinkMsg; - } - - public DownlinkMsg processCredentialsRequestMessageToEdge(EdgeEvent edgeEvent) { - DownlinkMsg downlinkMsg = null; - if (EdgeEventType.DEVICE.equals(edgeEvent.getType())) { - DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); - DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = DeviceCredentialsRequestMsg.newBuilder() - .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) - .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) - .build(); - DownlinkMsg.Builder builder = DownlinkMsg.newBuilder() - .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) - .addDeviceCredentialsRequestMsg(deviceCredentialsRequestMsg); - downlinkMsg = builder.build(); - } - return downlinkMsg; - } - - public ListenableFuture processEntityNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { - EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()); - EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); - EntityId entityId = EntityIdFactory.getByEdgeEventTypeAndUuid(type, - new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); - EdgeId edgeId = safeGetEdgeId(edgeNotificationMsg); - switch (actionType) { - case ADDED: - case UPDATED: - case CREDENTIALS_UPDATED: - case ASSIGNED_TO_CUSTOMER: - case UNASSIGNED_FROM_CUSTOMER: - case DELETED: - if (edgeId != null) { - return saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null); - } else { - return pushNotificationToAllRelatedEdges(tenantId, entityId, type, actionType); - } - case ASSIGNED_TO_EDGE: - case UNASSIGNED_FROM_EDGE: - ListenableFuture future = saveEdgeEvent(tenantId, edgeId, type, actionType, entityId, null); - return Futures.transformAsync(future, unused -> { - if (type.equals(EdgeEventType.RULE_CHAIN)) { - return updateDependentRuleChains(tenantId, new RuleChainId(entityId.getId()), edgeId); - } else { - return Futures.immediateFuture(null); - } - }, dbCallbackExecutorService); - default: - return Futures.immediateFuture(null); - } - } - - private EdgeId safeGetEdgeId(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { - if (edgeNotificationMsg.getEdgeIdMSB() != 0 && edgeNotificationMsg.getEdgeIdLSB() != 0) { - return new EdgeId(new UUID(edgeNotificationMsg.getEdgeIdMSB(), edgeNotificationMsg.getEdgeIdLSB())); - } else { - return null; - } - } - - private ListenableFuture pushNotificationToAllRelatedEdges(TenantId tenantId, EntityId entityId, EdgeEventType type, EdgeEventActionType actionType) { - PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); - PageData pageData; - List> futures = new ArrayList<>(); - do { - pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink); - if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { - for (EdgeId relatedEdgeId : pageData.getData()) { - futures.add(saveEdgeEvent(tenantId, relatedEdgeId, type, actionType, entityId, null)); - } - if (pageData.hasNext()) { - pageLink = pageLink.nextPageLink(); - } - } - } while (pageData != null && pageData.hasNext()); - return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); - } - - private ListenableFuture updateDependentRuleChains(TenantId tenantId, RuleChainId processingRuleChainId, EdgeId edgeId) { - PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); - PageData pageData; - List> futures = new ArrayList<>(); - do { - pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink); - if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { - for (RuleChain ruleChain : pageData.getData()) { - if (!ruleChain.getId().equals(processingRuleChainId)) { - List connectionInfos = - ruleChainService.loadRuleChainMetaData(ruleChain.getTenantId(), ruleChain.getId()).getRuleChainConnections(); - if (connectionInfos != null && !connectionInfos.isEmpty()) { - for (RuleChainConnectionInfo connectionInfo : connectionInfos) { - if (connectionInfo.getTargetRuleChainId().equals(processingRuleChainId)) { - futures.add(saveEdgeEvent(tenantId, - edgeId, - EdgeEventType.RULE_CHAIN_METADATA, - EdgeEventActionType.UPDATED, - ruleChain.getId(), - null)); - } - } - } - } - } - if (pageData.hasNext()) { - pageLink = pageLink.nextPageLink(); - } - } - } while (pageData != null && pageData.hasNext()); - return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); - } - - public ListenableFuture processEntityNotificationForAllEdges(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { - EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()); - EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); - EntityId entityId = EntityIdFactory.getByEdgeEventTypeAndUuid(type, new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); - switch (actionType) { - case ADDED: - case UPDATED: - case DELETED: - case CREDENTIALS_UPDATED: // used by USER entity - return processActionForAllEdges(tenantId, type, actionType, entityId); - default: - return Futures.immediateFuture(null); - } - } -} - diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityViewEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityViewEdgeProcessor.java index 9207595750..b89ce15a5d 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityViewEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/EntityViewEdgeProcessor.java @@ -15,15 +15,18 @@ */ package org.thingsboard.server.service.edge.rpc.processor; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.EntityViewId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @Component @@ -31,7 +34,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class EntityViewEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processEntityViewToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertEntityViewEventToDownlink(EdgeEvent edgeEvent) { EntityViewId entityViewId = new EntityViewId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -63,4 +66,8 @@ public class EntityViewEdgeProcessor extends BaseEdgeProcessor { } return downlinkMsg; } + + public ListenableFuture processEntityViewNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotification(tenantId, edgeNotificationMsg); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/OtaPackageEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/OtaPackageEdgeProcessor.java index 0047eb39dc..37ae5142af 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/OtaPackageEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/OtaPackageEdgeProcessor.java @@ -15,15 +15,18 @@ */ package org.thingsboard.server.service.edge.rpc.processor; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.OtaPackage; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.OtaPackageId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.OtaPackageUpdateMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @Component @@ -31,7 +34,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class OtaPackageEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processOtaPackageToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertOtaPackageEventToDownlink(EdgeEvent edgeEvent) { OtaPackageId otaPackageId = new OtaPackageId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -60,4 +63,7 @@ public class OtaPackageEdgeProcessor extends BaseEdgeProcessor { return downlinkMsg; } + public ListenableFuture processOtaPackageNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotificationForAllEdges(tenantId, edgeNotificationMsg); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/QueueEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/QueueEdgeProcessor.java index e5881f7b27..6a7e66a58a 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/QueueEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/QueueEdgeProcessor.java @@ -15,15 +15,18 @@ */ package org.thingsboard.server.service.edge.rpc.processor; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.QueueId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.QueueUpdateMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @Component @@ -31,7 +34,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class QueueEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processQueueToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertQueueEventToDownlink(EdgeEvent edgeEvent) { QueueId queueId = new QueueId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -60,4 +63,7 @@ public class QueueEdgeProcessor extends BaseEdgeProcessor { return downlinkMsg; } + public ListenableFuture processQueueNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotificationForAllEdges(tenantId, edgeNotificationMsg); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java index 1c4db0d151..34e4d1a24f 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RelationEdgeProcessor.java @@ -118,7 +118,7 @@ public class RelationEdgeProcessor extends BaseEdgeProcessor { } } - public DownlinkMsg processRelationToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertRelationEventToDownlink(EdgeEvent edgeEvent) { EntityRelation entityRelation = JacksonUtil.OBJECT_MAPPER.convertValue(edgeEvent.getBody(), EntityRelation.class); UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); RelationUpdateMsg relationUpdateMsg = relationMsgConstructor.constructRelationUpdatedMsg(msgType, entityRelation); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RuleChainEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RuleChainEdgeProcessor.java index f55c72adbc..5704e7d30d 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RuleChainEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/RuleChainEdgeProcessor.java @@ -15,12 +15,14 @@ */ package org.thingsboard.server.service.edge.rpc.processor; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.RuleChainId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; @@ -28,6 +30,7 @@ import org.thingsboard.server.gen.edge.v1.EdgeVersion; import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg; import org.thingsboard.server.gen.edge.v1.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @Component @@ -35,7 +38,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class RuleChainEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processRuleChainToEdge(Edge edge, EdgeEvent edgeEvent) { + public DownlinkMsg convertRuleChainEventToDownlink(Edge edge, EdgeEvent edgeEvent) { RuleChainId ruleChainId = new RuleChainId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -64,7 +67,7 @@ public class RuleChainEdgeProcessor extends BaseEdgeProcessor { return downlinkMsg; } - public DownlinkMsg processRuleChainMetadataToEdge(EdgeEvent edgeEvent, EdgeVersion edgeVersion) { + public DownlinkMsg convertRuleChainMetadataEventToDownlink(EdgeEvent edgeEvent, EdgeVersion edgeVersion) { RuleChainId ruleChainId = new RuleChainId(edgeEvent.getEntityId()); RuleChain ruleChain = ruleChainService.findRuleChainById(edgeEvent.getTenantId(), ruleChainId); DownlinkMsg downlinkMsg = null; @@ -82,4 +85,8 @@ public class RuleChainEdgeProcessor extends BaseEdgeProcessor { } return downlinkMsg; } + + public ListenableFuture processRuleChainNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotification(tenantId, edgeNotificationMsg); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryEdgeProcessor.java index 9c75373cad..57e460ac2f 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/TelemetryEdgeProcessor.java @@ -304,7 +304,7 @@ public class TelemetryEdgeProcessor extends BaseEdgeProcessor { } } - public DownlinkMsg processTelemetryMessageToEdge(EdgeEvent edgeEvent) throws JsonProcessingException { + public DownlinkMsg convertTelemetryEventToDownlink(EdgeEvent edgeEvent) throws JsonProcessingException { EntityId entityId; switch (edgeEvent.getType()) { case DEVICE: diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/UserEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/UserEdgeProcessor.java index b771ba96f0..390fafe464 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/UserEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/UserEdgeProcessor.java @@ -15,16 +15,19 @@ */ package org.thingsboard.server.service.edge.rpc.processor; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.security.UserCredentials; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.edge.v1.UserCredentialsUpdateMsg; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @Component @@ -32,7 +35,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class UserEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processUserToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertUserEventToDownlink(EdgeEvent edgeEvent) { UserId userId = new UserId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -67,4 +70,7 @@ public class UserEdgeProcessor extends BaseEdgeProcessor { return downlinkMsg; } + public ListenableFuture processUserNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotificationForAllEdges(tenantId, edgeNotificationMsg); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/WidgetBundleEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/WidgetBundleEdgeProcessor.java index 0652d44740..b74e3c1cec 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/WidgetBundleEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/WidgetBundleEdgeProcessor.java @@ -15,15 +15,18 @@ */ package org.thingsboard.server.service.edge.rpc.processor; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.WidgetsBundleId; import org.thingsboard.server.common.data.widget.WidgetsBundle; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.edge.v1.WidgetsBundleUpdateMsg; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @Component @@ -31,7 +34,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class WidgetBundleEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processWidgetsBundleToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertWidgetsBundleEventToDownlink(EdgeEvent edgeEvent) { WidgetsBundleId widgetsBundleId = new WidgetsBundleId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -59,4 +62,8 @@ public class WidgetBundleEdgeProcessor extends BaseEdgeProcessor { } return downlinkMsg; } + + public ListenableFuture processWidgetsBundleNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotificationForAllEdges(tenantId, edgeNotificationMsg); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/WidgetTypeEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/WidgetTypeEdgeProcessor.java index cce4dc75b7..10807b7b5c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/WidgetTypeEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/WidgetTypeEdgeProcessor.java @@ -15,15 +15,18 @@ */ package org.thingsboard.server.service.edge.rpc.processor; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.WidgetTypeId; import org.thingsboard.server.common.data.widget.WidgetTypeDetails; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.edge.v1.WidgetTypeUpdateMsg; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @Component @@ -31,7 +34,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; @TbCoreComponent public class WidgetTypeEdgeProcessor extends BaseEdgeProcessor { - public DownlinkMsg processWidgetTypeToEdge(EdgeEvent edgeEvent) { + public DownlinkMsg convertWidgetTypeEventToDownlink(EdgeEvent edgeEvent) { WidgetTypeId widgetTypeId = new WidgetTypeId(edgeEvent.getEntityId()); DownlinkMsg downlinkMsg = null; switch (edgeEvent.getAction()) { @@ -60,4 +63,7 @@ public class WidgetTypeEdgeProcessor extends BaseEdgeProcessor { return downlinkMsg; } + public ListenableFuture processWidgetTypeNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { + return processEntityNotificationForAllEdges(tenantId, edgeNotificationMsg); + } }