From 2078c71d32a6e8a2b2e3b691b3a61258cba6ef2d Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 6 Aug 2021 14:29:54 +0300 Subject: [PATCH] Added module cluster-api and used it in rule engine and other services --- application/pom.xml | 4 + .../server/actors/ActorSystemContext.java | 2 +- .../device/DeviceActorMessageProcessor.java | 6 +- .../actors/ruleChain/DefaultTbContext.java | 6 + .../RuleChainActorMessageProcessor.java | 2 +- .../controller/AbstractRpcController.java | 15 +- .../server/controller/BaseController.java | 41 +---- .../server/controller/CustomerController.java | 2 +- .../server/controller/DeviceController.java | 21 +-- .../controller/DeviceProfileController.java | 4 +- .../server/controller/EdgeController.java | 12 +- .../server/controller/Lwm2mController.java | 15 +- .../controller/RuleChainController.java | 16 +- .../server/controller/TenantController.java | 4 +- .../controller/TenantProfileController.java | 3 +- .../action/RuleEngineEntityActionService.java | 2 +- .../DefaultTbApiUsageStateService.java | 2 +- .../device/ClaimDevicesServiceImpl.java | 13 +- .../device/DeviceProvisionServiceImpl.java | 7 +- .../edge/DefaultEdgeNotificationService.java | 2 +- .../edge/rpc/processor/BaseEdgeProcessor.java | 2 +- .../rpc/processor/DeviceEdgeProcessor.java | 11 +- .../rpc/sync/DefaultEdgeRequestsService.java | 2 +- .../DefaultSystemDataLoaderService.java | 1 + .../ota/DefaultOtaPackageStateService.java | 8 +- .../queue/DefaultTbClusterService.java | 122 +++++++++++---- .../queue/DefaultTbCoreConsumerService.java | 4 +- .../DefaultTbRuleEngineConsumerService.java | 4 +- .../rpc/DefaultTbCoreDeviceRpcService.java | 5 +- .../rpc/DefaultTbRuleEngineRpcService.java | 5 +- .../rpc/FromDeviceRpcResponseActorMsg.java | 3 +- .../service/rpc/TbCoreDeviceRpcService.java | 1 + .../server/service/rpc/TbRpcService.java | 2 +- .../rpc/TbRuleEngineDeviceRpcService.java | 1 + .../rpc/ToDeviceRpcRequestActorMsg.java | 2 +- .../oauth2/AbstractOAuth2ClientMapper.java | 4 +- .../state/DefaultDeviceStateService.java | 49 ++---- .../service/state/DeviceStateService.java | 6 - .../DefaultSubscriptionManagerService.java | 2 +- .../DefaultTbLocalSubscriptionService.java | 3 +- .../AbstractSubscriptionService.java | 4 +- .../DefaultAlarmSubscriptionService.java | 2 +- .../DefaultTelemetrySubscriptionService.java | 2 +- .../transport/DefaultTransportApiService.java | 10 +- .../thingsboard/server/edge/BaseEdgeTest.java | 4 +- .../state/DefaultDeviceStateServiceTest.java | 2 +- common/cluster-api/pom.xml | 144 ++++++++++++++++++ .../server/cluster}/TbClusterService.java | 15 +- .../server/queue/TbQueueAdmin.java | 0 .../server/queue/TbQueueCallback.java | 0 .../server/queue/TbQueueConsumer.java | 0 .../server/queue/TbQueueHandler.java | 0 .../thingsboard/server/queue/TbQueueMsg.java | 0 .../server/queue/TbQueueMsgDecoder.java | 0 .../server/queue/TbQueueMsgHeaders.java | 0 .../server/queue/TbQueueMsgMetadata.java | 0 .../server/queue/TbQueueProducer.java | 0 .../server/queue/TbQueueRequestTemplate.java | 0 .../server/queue/TbQueueResponseTemplate.java | 0 .../src/main/proto/jsinvoke.proto | 0 .../src/main/proto/queue.proto | 0 .../server/dao/device/DeviceService.java | 4 +- .../server/common/data/rpc}/RpcError.java | 2 +- .../msg/ToDeviceActorNotificationMsg.java | 3 +- .../msg}/rpc/FromDeviceRpcResponse.java | 4 +- common/pom.xml | 1 + common/queue/pom.xml | 13 +- .../mqtt/session/GatewaySessionHandler.java | 2 +- pom.xml | 5 + rule-engine/rule-engine-api/pom.xml | 5 + .../api/RuleEngineDeviceRpcResponse.java | 1 + .../rule/engine/api/TbContext.java | 3 + .../DeviceAttributesEventNotificationMsg.java | 1 + ...eviceCredentialsUpdateNotificationMsg.java | 7 +- .../engine/api/msg/DeviceEdgeUpdateMsg.java | 1 + .../api/msg/DeviceNameOrTypeUpdateMsg.java | 1 + .../action/TbAbstractRelationActionNode.java | 5 +- 77 files changed, 394 insertions(+), 263 deletions(-) create mode 100644 common/cluster-api/pom.xml rename {application/src/main/java/org/thingsboard/server/service/queue => common/cluster-api/src/main/java/org/thingsboard/server/cluster}/TbClusterService.java (83%) rename common/{queue => cluster-api}/src/main/java/org/thingsboard/server/queue/TbQueueAdmin.java (100%) rename common/{queue => cluster-api}/src/main/java/org/thingsboard/server/queue/TbQueueCallback.java (100%) rename common/{queue => cluster-api}/src/main/java/org/thingsboard/server/queue/TbQueueConsumer.java (100%) rename common/{queue => cluster-api}/src/main/java/org/thingsboard/server/queue/TbQueueHandler.java (100%) rename common/{queue => cluster-api}/src/main/java/org/thingsboard/server/queue/TbQueueMsg.java (100%) rename common/{queue => cluster-api}/src/main/java/org/thingsboard/server/queue/TbQueueMsgDecoder.java (100%) rename common/{queue => cluster-api}/src/main/java/org/thingsboard/server/queue/TbQueueMsgHeaders.java (100%) rename common/{queue => cluster-api}/src/main/java/org/thingsboard/server/queue/TbQueueMsgMetadata.java (100%) rename common/{queue => cluster-api}/src/main/java/org/thingsboard/server/queue/TbQueueProducer.java (100%) rename common/{queue => cluster-api}/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java (100%) rename common/{queue => cluster-api}/src/main/java/org/thingsboard/server/queue/TbQueueResponseTemplate.java (100%) rename common/{queue => cluster-api}/src/main/proto/jsinvoke.proto (100%) rename common/{queue => cluster-api}/src/main/proto/queue.proto (100%) rename {rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api => common/data/src/main/java/org/thingsboard/server/common/data/rpc}/RpcError.java (93%) rename {rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api => common/message/src/main/java/org/thingsboard/server/common}/msg/ToDeviceActorNotificationMsg.java (90%) rename {application/src/main/java/org/thingsboard/server/service => common/message/src/main/java/org/thingsboard/server/common/msg}/rpc/FromDeviceRpcResponse.java (92%) diff --git a/application/pom.xml b/application/pom.xml index 83a9eefdc3..7aa0c54594 100644 --- a/application/pom.xml +++ b/application/pom.xml @@ -65,6 +65,10 @@ org.thingsboard.rule-engine rule-engine-api + + org.thingsboard.common + cluster-api + org.thingsboard.rule-engine rule-engine-components diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 7ad3d912ac..5516762783 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -81,7 +81,7 @@ import org.thingsboard.server.service.executors.ExternalCallExecutorService; import org.thingsboard.server.service.executors.SharedEventLoopGroupService; import org.thingsboard.server.service.mail.MailExecutorService; import org.thingsboard.server.service.profile.TbDeviceProfileCache; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService; import org.thingsboard.server.service.rpc.TbRpcService; import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService; diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index ac0a59e0cf..137b8b8863 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -26,7 +26,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections.CollectionUtils; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.LinkedHashMapRemoveEldest; -import org.thingsboard.rule.engine.api.RpcError; +import org.thingsboard.server.common.data.rpc.RpcError; import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg; import org.thingsboard.rule.engine.api.msg.DeviceCredentialsUpdateNotificationMsg; import org.thingsboard.rule.engine.api.msg.DeviceEdgeUpdateMsg; @@ -86,7 +86,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto; import org.thingsboard.server.gen.transport.TransportProtos.TransportToDeviceActorMsg; import org.thingsboard.server.gen.transport.TransportProtos.TsKvProto; -import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; +import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.service.rpc.FromDeviceRpcResponseActorMsg; import org.thingsboard.server.service.rpc.ToDeviceRpcRequestActorMsg; import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper; @@ -585,7 +585,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { log.debug("[{}] Received duplicate session open event [{}]", deviceId, sessionId); return; } - log.info("[{}] Processing new session [{}]. Current sessions size {}", deviceId, sessionId, sessions.size()); + log.debug("[{}] Processing new session [{}]. Current sessions size {}", deviceId, sessionId, sessions.size()); sessions.put(sessionId, new SessionInfoMetaData(new SessionInfo(SessionType.ASYNC, sessionInfo.getNodeId()))); if (sessions.size() == 1) { diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java index 9a1afb9ff8..740ff0102d 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java @@ -33,6 +33,7 @@ import org.thingsboard.rule.engine.api.TbRelationTypes; import org.thingsboard.rule.engine.api.sms.SmsSenderFactory; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.TbActorRef; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; @@ -453,6 +454,11 @@ class DefaultTbContext implements TbContext { return mainCtx.getDeviceService(); } + @Override + public TbClusterService getClusterService() { + return mainCtx.getClusterService(); + } + @Override public DashboardService getDashboardService() { return mainCtx.getDashboardService(); diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java index 10d6272beb..cb916f9baa 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java @@ -49,7 +49,7 @@ import org.thingsboard.server.queue.TbQueueCallback; import org.thingsboard.server.queue.common.MultipleTbQueueTbMsgCallbackWrapper; import org.thingsboard.server.queue.common.TbQueueTbMsgCallbackWrapper; import org.thingsboard.server.queue.usagestats.TbApiUsageClient; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import java.util.ArrayList; import java.util.Collections; diff --git a/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java b/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java index f7aa23a37b..ba14a353d1 100644 --- a/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java +++ b/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java @@ -22,33 +22,22 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.http.HttpStatus; import org.springframework.http.ResponseEntity; -import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.util.StringUtils; -import org.springframework.web.bind.annotation.PathVariable; -import org.springframework.web.bind.annotation.RequestMapping; -import org.springframework.web.bind.annotation.RequestMethod; -import org.springframework.web.bind.annotation.RequestParam; -import org.springframework.web.bind.annotation.ResponseBody; import org.springframework.web.context.request.async.DeferredResult; import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.rule.engine.api.RpcError; +import org.thingsboard.server.common.data.rpc.RpcError; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.RpcId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UUIDBased; -import org.thingsboard.server.common.data.page.PageData; -import org.thingsboard.server.common.data.page.PageLink; -import org.thingsboard.server.common.data.rpc.Rpc; -import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.data.rpc.ToDeviceRpcRequestBody; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; +import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.service.rpc.LocalRequestMetaData; import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService; import org.thingsboard.server.service.security.AccessValidator; diff --git a/application/src/main/java/org/thingsboard/server/controller/BaseController.java b/application/src/main/java/org/thingsboard/server/controller/BaseController.java index 09281a0ec0..ca727f4b4c 100644 --- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java +++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java @@ -33,7 +33,6 @@ import org.thingsboard.server.common.data.DashboardInfo; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceInfo; import org.thingsboard.server.common.data.DeviceProfile; -import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.EntityViewInfo; @@ -118,7 +117,6 @@ 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.exception.ThingsboardErrorResponseHandler; -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.queue.util.TbCoreComponent; @@ -129,7 +127,7 @@ import org.thingsboard.server.service.edge.rpc.EdgeRpcService; import org.thingsboard.server.service.lwm2m.LwM2MServerSecurityInfoRepository; import org.thingsboard.server.service.ota.OtaPackageStateService; import org.thingsboard.server.service.profile.TbDeviceProfileCache; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.resource.TbResourceService; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.permission.AccessControlService; @@ -888,42 +886,7 @@ public abstract class BaseController { } private void sendNotificationMsgToEdgeService(TenantId tenantId, EdgeId edgeId, EntityId entityId, String body, EdgeEventType type, EdgeEventActionType action) { - if (!edgesEnabled) { - return; - } - if (type == null) { - if (entityId != null) { - type = EdgeUtils.getEdgeEventTypeByEntityType(entityId.getEntityType()); - } else { - log.trace("[{}] entity id and type are null. Ignoring this notification", tenantId); - return; - } - if (type == null) { - log.trace("[{}] edge event type is null. Ignoring this notification [{}]", tenantId, entityId); - return; - } - } - TransportProtos.EdgeNotificationMsgProto.Builder builder = TransportProtos.EdgeNotificationMsgProto.newBuilder(); - builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); - builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); - builder.setType(type.name()); - builder.setAction(action.name()); - if (entityId != null) { - builder.setEntityIdMSB(entityId.getId().getMostSignificantBits()); - builder.setEntityIdLSB(entityId.getId().getLeastSignificantBits()); - builder.setEntityType(entityId.getEntityType().name()); - } - if (edgeId != null) { - builder.setEdgeIdMSB(edgeId.getId().getMostSignificantBits()); - builder.setEdgeIdLSB(edgeId.getId().getLeastSignificantBits()); - } - if (body != null) { - builder.setBody(body); - } - TransportProtos.EdgeNotificationMsgProto msg = builder.build(); - log.trace("[{}] sending notification to edge service {}", tenantId.getId(), msg); - tbClusterService.pushMsgToCore(tenantId, entityId != null ? entityId : tenantId, - TransportProtos.ToCoreMsg.newBuilder().setEdgeNotificationMsg(msg).build(), null); + tbClusterService.sendNotificationMsgToEdgeService(tenantId, edgeId, entityId, body, type, action); } protected List findRelatedEdgeIds(TenantId tenantId, EntityId entityId) { diff --git a/application/src/main/java/org/thingsboard/server/controller/CustomerController.java b/application/src/main/java/org/thingsboard/server/controller/CustomerController.java index 9eb1dc044a..3e58e0656e 100644 --- a/application/src/main/java/org/thingsboard/server/controller/CustomerController.java +++ b/application/src/main/java/org/thingsboard/server/controller/CustomerController.java @@ -149,7 +149,7 @@ public class CustomerController extends BaseController { ActionType.DELETED, null, strCustomerId); sendDeleteNotificationMsg(getTenantId(), customerId, relatedEdgeIds); - tbClusterService.onEntityStateChange(getTenantId(), customerId, ComponentLifecycleEvent.DELETED); + tbClusterService.broadcastEntityStateChangeEvent(getTenantId(), customerId, ComponentLifecycleEvent.DELETED); } catch (Exception e) { logEntityAction(emptyId(EntityType.CUSTOMER), diff --git a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java index c261d50f23..f66ba01e42 100644 --- a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java +++ b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java @@ -33,7 +33,6 @@ import org.springframework.web.bind.annotation.RestController; import org.springframework.web.context.request.async.DeferredResult; import org.thingsboard.rule.engine.api.msg.DeviceCredentialsUpdateNotificationMsg; import org.thingsboard.rule.engine.api.msg.DeviceEdgeUpdateMsg; -import org.thingsboard.rule.engine.api.msg.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.common.data.ClaimRequest; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.DataConstants; @@ -136,27 +135,12 @@ public class DeviceController extends BaseController { Device savedDevice = checkNotNull(deviceService.saveDeviceWithAccessToken(device, accessToken)); - tbClusterService.onDeviceChange(savedDevice, null); - tbClusterService.pushMsgToCore(new DeviceNameOrTypeUpdateMsg(savedDevice.getTenantId(), - savedDevice.getId(), savedDevice.getName(), savedDevice.getType()), null); - tbClusterService.onEntityStateChange(savedDevice.getTenantId(), savedDevice.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); - - if (!created) { - sendEntityNotificationMsg(savedDevice.getTenantId(), savedDevice.getId(), EdgeEventActionType.UPDATED); - } + tbClusterService.onDeviceUpdated(savedDevice, oldDevice); logEntityAction(savedDevice.getId(), savedDevice, savedDevice.getCustomerId(), created ? ActionType.ADDED : ActionType.UPDATED, null); - if (device.getId() == null) { - deviceStateService.onDeviceAdded(savedDevice); - } else { - deviceStateService.onDeviceUpdated(savedDevice); - } - - otaPackageStateService.update(savedDevice, oldDevice); - return savedDevice; } catch (Exception e) { logEntityAction(emptyId(EntityType.DEVICE), device, @@ -180,15 +164,12 @@ public class DeviceController extends BaseController { deviceService.deleteDevice(getCurrentUser().getTenantId(), deviceId); tbClusterService.onDeviceDeleted(device, null); - tbClusterService.onEntityStateChange(device.getTenantId(), deviceId, ComponentLifecycleEvent.DELETED); logEntityAction(deviceId, device, device.getCustomerId(), ActionType.DELETED, null, strDeviceId); sendDeleteNotificationMsg(getTenantId(), deviceId, relatedEdgeIds); - - deviceStateService.onDeviceDeleted(device); } catch (Exception e) { logEntityAction(emptyId(EntityType.DEVICE), null, diff --git a/application/src/main/java/org/thingsboard/server/controller/DeviceProfileController.java b/application/src/main/java/org/thingsboard/server/controller/DeviceProfileController.java index ceee45147b..3a68551773 100644 --- a/application/src/main/java/org/thingsboard/server/controller/DeviceProfileController.java +++ b/application/src/main/java/org/thingsboard/server/controller/DeviceProfileController.java @@ -161,7 +161,7 @@ public class DeviceProfileController extends BaseController { DeviceProfile savedDeviceProfile = checkNotNull(deviceProfileService.saveDeviceProfile(deviceProfile)); tbClusterService.onDeviceProfileChange(savedDeviceProfile, null); - tbClusterService.onEntityStateChange(deviceProfile.getTenantId(), savedDeviceProfile.getId(), + tbClusterService.broadcastEntityStateChangeEvent(deviceProfile.getTenantId(), savedDeviceProfile.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); logEntityAction(savedDeviceProfile.getId(), savedDeviceProfile, @@ -191,7 +191,7 @@ public class DeviceProfileController extends BaseController { deviceProfileService.deleteDeviceProfile(getTenantId(), deviceProfileId); tbClusterService.onDeviceProfileDelete(deviceProfile, null); - tbClusterService.onEntityStateChange(deviceProfile.getTenantId(), deviceProfile.getId(), ComponentLifecycleEvent.DELETED); + tbClusterService.broadcastEntityStateChangeEvent(deviceProfile.getTenantId(), deviceProfile.getId(), ComponentLifecycleEvent.DELETED); logEntityAction(deviceProfileId, deviceProfile, null, diff --git a/application/src/main/java/org/thingsboard/server/controller/EdgeController.java b/application/src/main/java/org/thingsboard/server/controller/EdgeController.java index 426c899588..5184591fe5 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EdgeController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EdgeController.java @@ -135,7 +135,7 @@ public class EdgeController extends BaseController { edgeService.assignDefaultRuleChainsToEdge(tenantId, savedEdge.getId()); } - tbClusterService.onEntityStateChange(savedEdge.getTenantId(), savedEdge.getId(), + tbClusterService.broadcastEntityStateChangeEvent(savedEdge.getTenantId(), savedEdge.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); logEntityAction(savedEdge.getId(), savedEdge, null, created ? ActionType.ADDED : ActionType.UPDATED, null); @@ -157,7 +157,7 @@ public class EdgeController extends BaseController { Edge edge = checkEdgeId(edgeId, Operation.DELETE); edgeService.deleteEdge(getTenantId(), edgeId); - tbClusterService.onEntityStateChange(getTenantId(), edgeId, + tbClusterService.broadcastEntityStateChangeEvent(getTenantId(), edgeId, ComponentLifecycleEvent.DELETED); logEntityAction(edgeId, edge, @@ -208,7 +208,7 @@ public class EdgeController extends BaseController { Edge savedEdge = checkNotNull(edgeService.assignEdgeToCustomer(getCurrentUser().getTenantId(), edgeId, customerId)); - tbClusterService.onEntityStateChange(getTenantId(), edgeId, + tbClusterService.broadcastEntityStateChangeEvent(getTenantId(), edgeId, ComponentLifecycleEvent.UPDATED); logEntityAction(edgeId, savedEdge, @@ -242,7 +242,7 @@ public class EdgeController extends BaseController { Edge savedEdge = checkNotNull(edgeService.unassignEdgeFromCustomer(getCurrentUser().getTenantId(), edgeId)); - tbClusterService.onEntityStateChange(getTenantId(), edgeId, + tbClusterService.broadcastEntityStateChangeEvent(getTenantId(), edgeId, ComponentLifecycleEvent.UPDATED); logEntityAction(edgeId, edge, @@ -272,7 +272,7 @@ public class EdgeController extends BaseController { Customer publicCustomer = customerService.findOrCreatePublicCustomer(edge.getTenantId()); Edge savedEdge = checkNotNull(edgeService.assignEdgeToCustomer(getCurrentUser().getTenantId(), edgeId, publicCustomer.getId())); - tbClusterService.onEntityStateChange(getTenantId(), edgeId, + tbClusterService.broadcastEntityStateChangeEvent(getTenantId(), edgeId, ComponentLifecycleEvent.UPDATED); logEntityAction(edgeId, savedEdge, @@ -364,7 +364,7 @@ public class EdgeController extends BaseController { Edge updatedEdge = edgeNotificationService.setEdgeRootRuleChain(getTenantId(), edge, ruleChainId); - tbClusterService.onEntityStateChange(updatedEdge.getTenantId(), updatedEdge.getId(), ComponentLifecycleEvent.UPDATED); + tbClusterService.broadcastEntityStateChangeEvent(updatedEdge.getTenantId(), updatedEdge.getId(), ComponentLifecycleEvent.UPDATED); logEntityAction(updatedEdge.getId(), updatedEdge, null, ActionType.UPDATED, null); diff --git a/application/src/main/java/org/thingsboard/server/controller/Lwm2mController.java b/application/src/main/java/org/thingsboard/server/controller/Lwm2mController.java index d94d26fc87..64537d9f5f 100644 --- a/application/src/main/java/org/thingsboard/server/controller/Lwm2mController.java +++ b/application/src/main/java/org/thingsboard/server/controller/Lwm2mController.java @@ -24,13 +24,11 @@ import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RequestMethod; import org.springframework.web.bind.annotation.ResponseBody; import org.springframework.web.bind.annotation.RestController; -import org.thingsboard.rule.engine.api.msg.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.lwm2m.ServerSecurityConfig; -import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.security.permission.Resource; @@ -66,22 +64,11 @@ public class Lwm2mController extends BaseController { checkEntity(device.getId(), device, Resource.DEVICE); Device savedDevice = deviceService.saveDeviceWithCredentials(device, credentials); checkNotNull(savedDevice); - - tbClusterService.onDeviceChange(savedDevice, null); - tbClusterService.pushMsgToCore(new DeviceNameOrTypeUpdateMsg(savedDevice.getTenantId(), - savedDevice.getId(), savedDevice.getName(), savedDevice.getType()), null); - tbClusterService.onEntityStateChange(savedDevice.getTenantId(), savedDevice.getId(), - device.getId() == null ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); - + tbClusterService.onDeviceUpdated(savedDevice, device); logEntityAction(savedDevice.getId(), savedDevice, savedDevice.getCustomerId(), device.getId() == null ? ActionType.ADDED : ActionType.UPDATED, null); - if (device.getId() == null) { - deviceStateService.onDeviceAdded(savedDevice); - } else { - deviceStateService.onDeviceUpdated(savedDevice); - } return savedDevice; } catch (Exception e) { logEntityAction(emptyId(EntityType.DEVICE), device, diff --git a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java index 7b6821d63d..a71fb70e0a 100644 --- a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java +++ b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java @@ -149,7 +149,7 @@ public class RuleChainController extends BaseController { RuleChain savedRuleChain = checkNotNull(ruleChainService.saveRuleChain(ruleChain)); if (RuleChainType.CORE.equals(savedRuleChain.getType())) { - tbClusterService.onEntityStateChange(ruleChain.getTenantId(), savedRuleChain.getId(), + tbClusterService.broadcastEntityStateChangeEvent(ruleChain.getTenantId(), savedRuleChain.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); } @@ -183,7 +183,7 @@ public class RuleChainController extends BaseController { RuleChain savedRuleChain = installScripts.createDefaultRuleChain(getCurrentUser().getTenantId(), request.getName()); - tbClusterService.onEntityStateChange(savedRuleChain.getTenantId(), savedRuleChain.getId(), ComponentLifecycleEvent.CREATED); + tbClusterService.broadcastEntityStateChangeEvent(savedRuleChain.getTenantId(), savedRuleChain.getId(), ComponentLifecycleEvent.CREATED); logEntityAction(savedRuleChain.getId(), savedRuleChain, null, ActionType.ADDED, null); @@ -210,7 +210,7 @@ public class RuleChainController extends BaseController { if (previousRootRuleChain != null) { previousRootRuleChain = ruleChainService.findRuleChainById(getTenantId(), previousRootRuleChain.getId()); - tbClusterService.onEntityStateChange(previousRootRuleChain.getTenantId(), previousRootRuleChain.getId(), + tbClusterService.broadcastEntityStateChangeEvent(previousRootRuleChain.getTenantId(), previousRootRuleChain.getId(), ComponentLifecycleEvent.UPDATED); logEntityAction(previousRootRuleChain.getId(), previousRootRuleChain, @@ -218,7 +218,7 @@ public class RuleChainController extends BaseController { } ruleChain = ruleChainService.findRuleChainById(getTenantId(), ruleChainId); - tbClusterService.onEntityStateChange(ruleChain.getTenantId(), ruleChain.getId(), + tbClusterService.broadcastEntityStateChangeEvent(ruleChain.getTenantId(), ruleChain.getId(), ComponentLifecycleEvent.UPDATED); logEntityAction(ruleChain.getId(), ruleChain, @@ -254,7 +254,7 @@ public class RuleChainController extends BaseController { RuleChainMetaData savedRuleChainMetaData = checkNotNull(ruleChainService.loadRuleChainMetaData(tenantId, ruleChainMetaData.getRuleChainId())); if (RuleChainType.CORE.equals(ruleChain.getType())) { - tbClusterService.onEntityStateChange(ruleChain.getTenantId(), ruleChain.getId(), ComponentLifecycleEvent.UPDATED); + tbClusterService.broadcastEntityStateChangeEvent(ruleChain.getTenantId(), ruleChain.getId(), ComponentLifecycleEvent.UPDATED); } logEntityAction(ruleChain.getId(), ruleChain, @@ -323,9 +323,9 @@ public class RuleChainController extends BaseController { if (RuleChainType.CORE.equals(ruleChain.getType())) { referencingRuleChainIds.forEach(referencingRuleChainId -> - tbClusterService.onEntityStateChange(ruleChain.getTenantId(), referencingRuleChainId, ComponentLifecycleEvent.UPDATED)); + tbClusterService.broadcastEntityStateChangeEvent(ruleChain.getTenantId(), referencingRuleChainId, ComponentLifecycleEvent.UPDATED)); - tbClusterService.onEntityStateChange(ruleChain.getTenantId(), ruleChain.getId(), ComponentLifecycleEvent.DELETED); + tbClusterService.broadcastEntityStateChangeEvent(ruleChain.getTenantId(), ruleChain.getId(), ComponentLifecycleEvent.DELETED); } logEntityAction(ruleChainId, ruleChain, @@ -456,7 +456,7 @@ public class RuleChainController extends BaseController { List importResults = ruleChainService.importTenantRuleChains(tenantId, ruleChainData, RuleChainType.CORE, overwrite); if (!CollectionUtils.isEmpty(importResults)) { for (RuleChainImportResult importResult : importResults) { - tbClusterService.onEntityStateChange(importResult.getTenantId(), importResult.getRuleChainId(), importResult.getLifecycleEvent()); + tbClusterService.broadcastEntityStateChangeEvent(importResult.getTenantId(), importResult.getRuleChainId(), importResult.getLifecycleEvent()); } } } catch (Exception e) { diff --git a/application/src/main/java/org/thingsboard/server/controller/TenantController.java b/application/src/main/java/org/thingsboard/server/controller/TenantController.java index 04d37371e3..f8ac4acf49 100644 --- a/application/src/main/java/org/thingsboard/server/controller/TenantController.java +++ b/application/src/main/java/org/thingsboard/server/controller/TenantController.java @@ -99,7 +99,7 @@ public class TenantController extends BaseController { } tenantProfileCache.evict(tenant.getId()); tbClusterService.onTenantChange(tenant, null); - tbClusterService.onEntityStateChange(tenant.getId(), tenant.getId(), + tbClusterService.broadcastEntityStateChangeEvent(tenant.getId(), tenant.getId(), newTenant ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); return tenant; } catch (Exception e) { @@ -118,7 +118,7 @@ public class TenantController extends BaseController { tenantService.deleteTenant(tenantId); tenantProfileCache.evict(tenantId); tbClusterService.onTenantDelete(tenant, null); - tbClusterService.onEntityStateChange(tenantId, tenantId, ComponentLifecycleEvent.DELETED); + tbClusterService.broadcastEntityStateChangeEvent(tenantId, tenantId, ComponentLifecycleEvent.DELETED); } catch (Exception e) { throw handleException(e); } diff --git a/application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java b/application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java index fc0a1b0291..e201c8f815 100644 --- a/application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java +++ b/application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java @@ -34,7 +34,6 @@ import org.thingsboard.server.common.data.id.TenantProfileId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; -import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Resource; @@ -98,7 +97,7 @@ public class TenantProfileController extends BaseController { tenantProfile = checkNotNull(tenantProfileService.saveTenantProfile(getTenantId(), tenantProfile)); tenantProfileCache.put(tenantProfile); tbClusterService.onTenantProfileChange(tenantProfile, null); - tbClusterService.onEntityStateChange(TenantId.SYS_TENANT_ID, tenantProfile.getId(), + tbClusterService.broadcastEntityStateChangeEvent(TenantId.SYS_TENANT_ID, tenantProfile.getId(), newTenantProfile ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); return tenantProfile; } catch (Exception e) { diff --git a/application/src/main/java/org/thingsboard/server/service/action/RuleEngineEntityActionService.java b/application/src/main/java/org/thingsboard/server/service/action/RuleEngineEntityActionService.java index f1320d1a7a..ef7d077edb 100644 --- a/application/src/main/java/org/thingsboard/server/service/action/RuleEngineEntityActionService.java +++ b/application/src/main/java/org/thingsboard/server/service/action/RuleEngineEntityActionService.java @@ -39,7 +39,7 @@ 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.queue.util.TbCoreComponent; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import java.util.List; import java.util.Map; diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java index 123537194a..47f698baff 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java @@ -61,7 +61,7 @@ import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.scheduler.SchedulerComponent; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.telemetry.InternalTelemetryService; import javax.annotation.PostConstruct; diff --git a/application/src/main/java/org/thingsboard/server/service/device/ClaimDevicesServiceImpl.java b/application/src/main/java/org/thingsboard/server/service/device/ClaimDevicesServiceImpl.java index 7b3c98c445..dc00c63fd0 100644 --- a/application/src/main/java/org/thingsboard/server/service/device/ClaimDevicesServiceImpl.java +++ b/application/src/main/java/org/thingsboard/server/service/device/ClaimDevicesServiceImpl.java @@ -29,6 +29,7 @@ import org.springframework.cache.CacheManager; import org.springframework.stereotype.Service; import org.springframework.util.StringUtils; import org.thingsboard.rule.engine.api.RuleEngineTelemetryService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; @@ -69,6 +70,8 @@ public class ClaimDevicesServiceImpl implements ClaimDevicesService { private static final String CLAIM_DATA_ATTRIBUTE_NAME = "claimingData"; private static final ObjectMapper mapper = new ObjectMapper(); + @Autowired + private TbClusterService clusterService; @Autowired private DeviceService deviceService; @Autowired @@ -155,6 +158,7 @@ public class ClaimDevicesServiceImpl implements ClaimDevicesService { if (device.getCustomerId().getId().equals(ModelConstants.NULL_UUID)) { device.setCustomerId(customerId); Device savedDevice = deviceService.saveDevice(device); + clusterService.onDeviceUpdated(savedDevice, device); return Futures.transform(removeClaimingSavedData(cache, claimData, device), result -> new ClaimResult(savedDevice, ClaimResponse.SUCCESS), MoreExecutors.directExecutor()); } return Futures.transform(removeClaimingSavedData(cache, claimData, device), result -> new ClaimResult(null, ClaimResponse.CLAIMED), MoreExecutors.directExecutor()); @@ -179,13 +183,14 @@ public class ClaimDevicesServiceImpl implements ClaimDevicesService { cacheEviction(device.getId()); Customer unassignedCustomer = customerService.findCustomerById(tenantId, device.getCustomerId()); device.setCustomerId(null); - deviceService.saveDevice(device); + Device savedDevice = deviceService.saveDevice(device); + clusterService.onDeviceUpdated(savedDevice, device); if (isAllowedClaimingByDefault) { return Futures.immediateFuture(new ReclaimResult(unassignedCustomer)); } SettableFuture result = SettableFuture.create(); telemetryService.saveAndNotify( - tenantId, device.getId(), DataConstants.SERVER_SCOPE, Collections.singletonList( + tenantId, savedDevice.getId(), DataConstants.SERVER_SCOPE, Collections.singletonList( new BaseAttributeKvEntry(new BooleanDataEntry(CLAIM_ATTRIBUTE_NAME, true), System.currentTimeMillis()) ), new FutureCallback<>() { @@ -198,7 +203,7 @@ public class ClaimDevicesServiceImpl implements ClaimDevicesService { public void onFailure(Throwable t) { result.setException(t); } - }); + }); return result; } cacheEviction(device.getId()); @@ -238,7 +243,7 @@ public class ClaimDevicesServiceImpl implements ClaimDevicesService { public void onFailure(Throwable t) { result.setException(t); } - }); + }); return result; } diff --git a/application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java b/application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java index 36f4ea06af..8b8cfe3131 100644 --- a/application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java +++ b/application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java @@ -24,6 +24,7 @@ import org.apache.commons.lang3.RandomStringUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import org.springframework.util.StringUtils; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; @@ -77,6 +78,9 @@ public class DeviceProvisionServiceImpl implements DeviceProvisionService { private static final String DEVICE_PROVISION_STATE = "provisionState"; private static final String PROVISIONED_STATE = "provisioned"; + @Autowired + TbClusterService clusterService; + @Autowired DeviceDao deviceDao; @@ -190,8 +194,7 @@ public class DeviceProvisionServiceImpl implements DeviceProvisionService { provisionRequest.setDeviceName(newDeviceName); } Device savedDevice = deviceService.saveDevice(provisionRequest, profile); - - deviceStateService.onDeviceAdded(savedDevice); + clusterService.onDeviceUpdated(savedDevice, null); saveProvisionStateAttribute(savedDevice).get(); pushDeviceCreatedEventToRuleEngine(savedDevice); notify(savedDevice, provisionRequest, DataConstants.PROVISION_SUCCESS, true); 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 8a1d70751b..cb24a57fb2 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 @@ -37,7 +37,7 @@ import org.thingsboard.server.service.edge.rpc.processor.CustomerEdgeProcessor; 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.RelationEdgeProcessor; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; 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 e213313228..c5cbc39b46 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 @@ -62,7 +62,7 @@ import org.thingsboard.server.service.edge.rpc.constructor.WidgetTypeMsgConstruc import org.thingsboard.server.service.edge.rpc.constructor.WidgetsBundleMsgConstructor; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.profile.TbDeviceProfileCache; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.state.DeviceStateService; @Slf4j 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 c036eebc16..2f24c1a219 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 @@ -18,7 +18,6 @@ package org.thingsboard.server.service.edge.rpc.processor; import com.datastax.oss.driver.api.core.uuid.Uuids; import com.fasterxml.jackson.core.JsonProcessingException; 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.SettableFuture; @@ -27,7 +26,7 @@ import org.apache.commons.lang3.RandomStringUtils; import org.apache.commons.lang3.StringUtils; import org.springframework.stereotype.Component; import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.rule.engine.api.RpcError; +import org.thingsboard.server.common.data.rpc.RpcError; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; @@ -60,7 +59,7 @@ import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.queue.TbQueueCallback; import org.thingsboard.server.queue.TbQueueMsgMetadata; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; +import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.service.rpc.FromDeviceRpcResponseActorMsg; import java.util.UUID; @@ -176,7 +175,8 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { deviceUpdateMsg.getDeviceProfileIdLSB().getValue())); device.setDeviceProfileId(deviceProfileId); } - deviceService.saveDevice(device); + Device savedDevice = deviceService.saveDevice(device); + tbClusterService.onDeviceUpdated(savedDevice, device); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_REQUEST, deviceId, null); } else { log.warn("[{}] can't find device [{}], edge [{}]", tenantId, deviceUpdateMsg, edge.getId()); @@ -215,14 +215,13 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { device.setDeviceProfileId(deviceProfileId); } Device savedDevice = deviceService.saveDevice(device, false); + tbClusterService.onDeviceUpdated(savedDevice, device); if (created) { DeviceCredentials deviceCredentials = new DeviceCredentials(); deviceCredentials.setDeviceId(new DeviceId(savedDevice.getUuidId())); deviceCredentials.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN); deviceCredentials.setCredentialsId(org.apache.commons.lang3.RandomStringUtils.randomAlphanumeric(20)); deviceCredentialsService.createDeviceCredentials(device.getTenantId(), deviceCredentials); - - deviceStateService.onDeviceAdded(savedDevice); } createRelationFromEdge(tenantId, edge.getId(), device.getId()); pushDeviceCreatedEventToRuleEngine(tenantId, edge, device); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java index d7971e78be..a1f8fc89ec 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java @@ -73,7 +73,7 @@ import org.thingsboard.server.gen.edge.v1.UserCredentialsRequestMsg; import org.thingsboard.server.gen.edge.v1.WidgetBundleTypesRequestMsg; import org.thingsboard.server.service.edge.rpc.EdgeEventUtils; import org.thingsboard.server.service.executors.DbCallbackExecutorService; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import java.util.ArrayList; import java.util.HashMap; diff --git a/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java b/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java index 7042bb6de1..ab1d2bbc61 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java @@ -481,6 +481,7 @@ public class DefaultSystemDataLoaderService implements SystemDataLoaderService { device.setAdditionalInfo(additionalInfo); } device = deviceService.saveDevice(device); + //TODO: No access to cluster service, so we should manually update the status of device. DeviceCredentials deviceCredentials = deviceCredentialsService.findDeviceCredentialsByDeviceId(TenantId.SYS_TENANT_ID, device.getId()); deviceCredentials.setCredentialsId(accessToken); deviceCredentialsService.updateDeviceCredentials(TenantId.SYS_TENANT_ID, deviceCredentials); diff --git a/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java b/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java index 4e5a104577..c4e780af36 100644 --- a/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java @@ -17,6 +17,7 @@ package org.thingsboard.server.service.ota; import com.google.common.util.concurrent.FutureCallback; import lombok.extern.slf4j.Slf4j; +import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.thingsboard.rule.engine.api.RuleEngineTelemetryService; import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg; @@ -49,7 +50,7 @@ import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import javax.annotation.Nullable; import java.util.ArrayList; @@ -87,10 +88,11 @@ public class DefaultOtaPackageStateService implements OtaPackageStateService { private final RuleEngineTelemetryService telemetryService; private final TbQueueProducer> otaPackageStateMsgProducer; - public DefaultOtaPackageStateService(TbClusterService tbClusterService, OtaPackageService otaPackageService, + public DefaultOtaPackageStateService(@Lazy TbClusterService tbClusterService, + OtaPackageService otaPackageService, DeviceService deviceService, DeviceProfileService deviceProfileService, - RuleEngineTelemetryService telemetryService, + @Lazy RuleEngineTelemetryService telemetryService, TbCoreQueueFactory coreQueueFactory) { this.tbClusterService = tbClusterService; this.otaPackageService = otaPackageService; diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java index fff648401b..0829583fbc 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java @@ -16,11 +16,17 @@ package org.thingsboard.server.service.queue; import com.google.protobuf.ByteString; +import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; -import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg; +import org.thingsboard.rule.engine.api.msg.DeviceNameOrTypeUpdateMsg; +import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.common.data.EdgeUtils; +import org.thingsboard.server.common.data.edge.EdgeEventActionType; +import org.thingsboard.server.common.data.edge.EdgeEventType; +import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg; import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; @@ -29,7 +35,6 @@ import org.thingsboard.server.common.data.HasName; import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; -import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EdgeId; @@ -56,8 +61,9 @@ import org.thingsboard.server.queue.common.MultipleTbQueueCallbackWrapper; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; +import org.thingsboard.server.service.ota.OtaPackageStateService; import org.thingsboard.server.service.profile.TbDeviceProfileCache; -import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; +import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import java.util.HashSet; import java.util.Set; @@ -66,10 +72,13 @@ import java.util.concurrent.atomic.AtomicInteger; @Service @Slf4j +@RequiredArgsConstructor public class DefaultTbClusterService implements TbClusterService { @Value("${cluster.stats.enabled:false}") private boolean statsEnabled; + @Value("${edges.enabled}") + protected boolean edgesEnabled; private final AtomicInteger toCoreMsgs = new AtomicInteger(0); private final AtomicInteger toCoreNfs = new AtomicInteger(0); @@ -81,13 +90,7 @@ public class DefaultTbClusterService implements TbClusterService { private final PartitionService partitionService; private final DataDecodingEncodingService encodingService; private final TbDeviceProfileCache deviceProfileCache; - - public DefaultTbClusterService(TbQueueProducerProvider producerProvider, PartitionService partitionService, DataDecodingEncodingService encodingService, TbDeviceProfileCache deviceProfileCache) { - this.producerProvider = producerProvider; - this.partitionService = partitionService; - this.encodingService = encodingService; - this.deviceProfileCache = deviceProfileCache; - } + private final OtaPackageStateService otaPackageStateService; @Override public void pushMsgToCore(TenantId tenantId, EntityId entityId, ToCoreMsg msg, TbQueueCallback callback) { @@ -200,55 +203,52 @@ public class DefaultTbClusterService implements TbClusterService { } @Override - public void onEntityStateChange(TenantId tenantId, EntityId entityId, ComponentLifecycleEvent state) { + public void broadcastEntityStateChangeEvent(TenantId tenantId, EntityId entityId, ComponentLifecycleEvent state) { log.trace("[{}] Processing {} state change event: {}", tenantId, entityId.getEntityType(), state); broadcast(new ComponentLifecycleMsg(tenantId, entityId, state)); } @Override public void onDeviceProfileChange(DeviceProfile deviceProfile, TbQueueCallback callback) { - onEntityChange(deviceProfile.getTenantId(), deviceProfile.getId(), deviceProfile, callback); + broadcastEntityChangeToTransport(deviceProfile.getTenantId(), deviceProfile.getId(), deviceProfile, callback); } @Override public void onTenantProfileChange(TenantProfile tenantProfile, TbQueueCallback callback) { - onEntityChange(TenantId.SYS_TENANT_ID, tenantProfile.getId(), tenantProfile, callback); + broadcastEntityChangeToTransport(TenantId.SYS_TENANT_ID, tenantProfile.getId(), tenantProfile, callback); } @Override public void onTenantChange(Tenant tenant, TbQueueCallback callback) { - onEntityChange(TenantId.SYS_TENANT_ID, tenant.getId(), tenant, callback); + broadcastEntityChangeToTransport(TenantId.SYS_TENANT_ID, tenant.getId(), tenant, callback); } @Override public void onApiStateChange(ApiUsageState apiUsageState, TbQueueCallback callback) { - onEntityChange(apiUsageState.getTenantId(), apiUsageState.getId(), apiUsageState, callback); + broadcastEntityChangeToTransport(apiUsageState.getTenantId(), apiUsageState.getId(), apiUsageState, callback); broadcast(new ComponentLifecycleMsg(apiUsageState.getTenantId(), apiUsageState.getId(), ComponentLifecycleEvent.UPDATED)); } @Override public void onDeviceProfileDelete(DeviceProfile entity, TbQueueCallback callback) { - onEntityDelete(entity.getTenantId(), entity.getId(), entity.getName(), callback); + broadcastEntityDeleteToTransport(entity.getTenantId(), entity.getId(), entity.getName(), callback); } @Override public void onTenantProfileDelete(TenantProfile entity, TbQueueCallback callback) { - onEntityDelete(TenantId.SYS_TENANT_ID, entity.getId(), entity.getName(), callback); + broadcastEntityDeleteToTransport(TenantId.SYS_TENANT_ID, entity.getId(), entity.getName(), callback); } @Override public void onTenantDelete(Tenant entity, TbQueueCallback callback) { - onEntityDelete(TenantId.SYS_TENANT_ID, entity.getId(), entity.getName(), callback); + broadcastEntityDeleteToTransport(TenantId.SYS_TENANT_ID, entity.getId(), entity.getName(), callback); } @Override - public void onDeviceChange(Device entity, TbQueueCallback callback) { - onEntityChange(entity.getTenantId(), entity.getId(), entity, callback); - } - - @Override - public void onDeviceDeleted(Device entity, TbQueueCallback callback) { - onEntityDelete(entity.getTenantId(), entity.getId(), entity.getName(), callback); + public void onDeviceDeleted(Device device, TbQueueCallback callback) { + broadcastEntityDeleteToTransport(device.getTenantId(), device.getId(), device.getName(), callback); + sendDeviceStateServiceEvent(device.getTenantId(), device.getId(), false, false, true); + broadcastEntityStateChangeEvent(device.getTenantId(), device.getId(), ComponentLifecycleEvent.DELETED); } @Override @@ -278,7 +278,7 @@ public class DefaultTbClusterService implements TbClusterService { broadcast(transportMsg, callback); } - public void onEntityChange(TenantId tenantId, EntityId entityid, T entity, TbQueueCallback callback) { + public void broadcastEntityChangeToTransport(TenantId tenantId, EntityId entityid, T entity, TbQueueCallback callback) { String entityName = (entity instanceof HasName) ? ((HasName) entity).getName() : entity.getClass().getName(); log.trace("[{}][{}][{}] Processing [{}] change event", tenantId, entityid.getEntityType(), entityid.getId(), entityName); TransportProtos.EntityUpdateMsg entityUpdateMsg = TransportProtos.EntityUpdateMsg.newBuilder() @@ -288,7 +288,7 @@ public class DefaultTbClusterService implements TbClusterService { broadcast(transportMsg, callback); } - private void onEntityDelete(TenantId tenantId, EntityId entityId, String name, TbQueueCallback callback) { + private void broadcastEntityDeleteToTransport(TenantId tenantId, EntityId entityId, String name, TbQueueCallback callback) { log.trace("[{}][{}][{}] Processing [{}] delete event", tenantId, entityId.getEntityType(), entityId.getId(), name); TransportProtos.EntityDeleteMsg entityDeleteMsg = TransportProtos.EntityDeleteMsg.newBuilder() .setEntityType(entityId.getEntityType().name()) @@ -369,4 +369,72 @@ public class DefaultTbClusterService implements TbClusterService { } } } + + private void sendDeviceStateServiceEvent(TenantId tenantId, DeviceId deviceId, boolean added, boolean updated, boolean deleted) { + TransportProtos.DeviceStateServiceMsgProto.Builder builder = TransportProtos.DeviceStateServiceMsgProto.newBuilder(); + builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); + builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); + builder.setDeviceIdMSB(deviceId.getId().getMostSignificantBits()); + builder.setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()); + builder.setAdded(added); + builder.setUpdated(updated); + builder.setDeleted(deleted); + TransportProtos.DeviceStateServiceMsgProto msg = builder.build(); + pushMsgToCore(tenantId, deviceId, TransportProtos.ToCoreMsg.newBuilder().setDeviceStateServiceMsg(msg).build(), null); + } + + @Override + public void onDeviceUpdated(Device device, Device old) { + var created = old == null; + broadcastEntityChangeToTransport(device.getTenantId(), device.getId(), device, null); + if (old != null && (!device.getName().equals(old.getName()) || !device.getType().equals(old.getType()))) { + pushMsgToCore(new DeviceNameOrTypeUpdateMsg(device.getTenantId(), device.getId(), device.getName(), device.getType()), null); + } + broadcastEntityStateChangeEvent(device.getTenantId(), device.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); + sendDeviceStateServiceEvent(device.getTenantId(), device.getId(), created, !created, false); + otaPackageStateService.update(device, old); + if (!created) { + sendNotificationMsgToEdgeService(device.getTenantId(), null, device.getId(), null, null, EdgeEventActionType.UPDATED); + } + } + + @Override + public void sendNotificationMsgToEdgeService(TenantId tenantId, EdgeId edgeId, EntityId entityId, String body, EdgeEventType type, EdgeEventActionType action) { + if (!edgesEnabled) { + return; + } + if (type == null) { + if (entityId != null) { + type = EdgeUtils.getEdgeEventTypeByEntityType(entityId.getEntityType()); + } else { + log.trace("[{}] entity id and type are null. Ignoring this notification", tenantId); + return; + } + if (type == null) { + log.trace("[{}] edge event type is null. Ignoring this notification [{}]", tenantId, entityId); + return; + } + } + TransportProtos.EdgeNotificationMsgProto.Builder builder = TransportProtos.EdgeNotificationMsgProto.newBuilder(); + builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); + builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); + builder.setType(type.name()); + builder.setAction(action.name()); + if (entityId != null) { + builder.setEntityIdMSB(entityId.getId().getMostSignificantBits()); + builder.setEntityIdLSB(entityId.getId().getLeastSignificantBits()); + builder.setEntityType(entityId.getEntityType().name()); + } + if (edgeId != null) { + builder.setEdgeIdMSB(edgeId.getId().getMostSignificantBits()); + builder.setEdgeIdLSB(edgeId.getId().getLeastSignificantBits()); + } + if (body != null) { + builder.setBody(body); + } + TransportProtos.EdgeNotificationMsgProto msg = builder.build(); + log.trace("[{}] sending notification to edge service {}", tenantId.getId(), msg); + pushMsgToCore(tenantId, entityId != null ? entityId : tenantId, TransportProtos.ToCoreMsg.newBuilder().setEdgeNotificationMsg(msg).build(), null); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index ceb364f921..ad7a0340cf 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -26,7 +26,7 @@ import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardThreadFactory; -import org.thingsboard.rule.engine.api.RpcError; +import org.thingsboard.server.common.data.rpc.RpcError; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.id.TenantId; @@ -64,7 +64,7 @@ import org.thingsboard.server.service.ota.OtaPackageStateService; import org.thingsboard.server.service.profile.TbDeviceProfileCache; import org.thingsboard.server.service.queue.processing.AbstractConsumerService; import org.thingsboard.server.service.queue.processing.IdMsgPair; -import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; +import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService; import org.thingsboard.server.service.rpc.ToDeviceRpcRequestActorMsg; import org.thingsboard.server.service.state.DeviceStateService; diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java index 2edaceee8e..4b1392e14a 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java @@ -21,7 +21,7 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardThreadFactory; -import org.thingsboard.rule.engine.api.RpcError; +import org.thingsboard.server.common.data.rpc.RpcError; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.TbMsg; @@ -55,7 +55,7 @@ import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingStr import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingStrategyFactory; import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategy; import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategyFactory; -import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; +import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService; import org.thingsboard.server.service.stats.RuleEngineStatisticsService; diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java index 9fbbb142cf..036beaea48 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java @@ -22,18 +22,19 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardThreadFactory; -import org.thingsboard.rule.engine.api.RpcError; +import org.thingsboard.server.common.data.rpc.RpcError; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; 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.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.security.model.SecurityUser; import javax.annotation.PostConstruct; diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java index 505c8db7d2..230f6e759e 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java @@ -19,18 +19,19 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardThreadFactory; -import org.thingsboard.rule.engine.api.RpcError; +import org.thingsboard.server.common.data.rpc.RpcError; import org.thingsboard.rule.engine.api.RuleEngineDeviceRpcRequest; import org.thingsboard.rule.engine.api.RuleEngineDeviceRpcResponse; import org.thingsboard.server.common.data.rpc.ToDeviceRpcRequestBody; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.util.TbRuleEngineComponent; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/FromDeviceRpcResponseActorMsg.java b/application/src/main/java/org/thingsboard/server/service/rpc/FromDeviceRpcResponseActorMsg.java index 9efe9a3b1f..1316c08a45 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/FromDeviceRpcResponseActorMsg.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/FromDeviceRpcResponseActorMsg.java @@ -18,10 +18,11 @@ package org.thingsboard.server.service.rpc; import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.ToString; -import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg; +import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.MsgType; +import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; @ToString @RequiredArgsConstructor diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/TbCoreDeviceRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/TbCoreDeviceRpcService.java index c3c6f9b888..db22e1ea64 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/TbCoreDeviceRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/TbCoreDeviceRpcService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.rpc; +import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; import org.thingsboard.server.service.security.model.SecurityUser; diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java index 456ea7d6d6..5a7f3de1f0 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java @@ -31,7 +31,7 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.dao.rpc.RpcService; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; @TbCoreComponent @Service diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/TbRuleEngineDeviceRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/TbRuleEngineDeviceRpcService.java index 3b920ac2eb..34197797e6 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/TbRuleEngineDeviceRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/TbRuleEngineDeviceRpcService.java @@ -16,6 +16,7 @@ package org.thingsboard.server.service.rpc; import org.thingsboard.rule.engine.api.RuleEngineRpcService; +import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; /** * Created by ashvayka on 16.04.18. diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/ToDeviceRpcRequestActorMsg.java b/application/src/main/java/org/thingsboard/server/service/rpc/ToDeviceRpcRequestActorMsg.java index 64a08b2ef4..7575e8b434 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/ToDeviceRpcRequestActorMsg.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/ToDeviceRpcRequestActorMsg.java @@ -18,7 +18,7 @@ package org.thingsboard.server.service.rpc; import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.ToString; -import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg; +import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.MsgType; diff --git a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java index f06dfc0711..4d4780f7ca 100644 --- a/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java +++ b/application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/AbstractOAuth2ClientMapper.java @@ -47,7 +47,7 @@ import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.user.UserService; import org.thingsboard.server.service.install.InstallScripts; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.model.UserPrincipal; @@ -180,7 +180,7 @@ public abstract class AbstractOAuth2ClientMapper { installScripts.createDefaultEdgeRuleChains(tenant.getId()); tenantProfileCache.evict(tenant.getId()); tbClusterService.onTenantChange(tenant, null); - tbClusterService.onEntityStateChange(tenant.getId(), tenant.getId(), + tbClusterService.broadcastEntityStateChangeEvent(tenant.getId(), tenant.getId(), ComponentLifecycleEvent.CREATED); } else { tenant = tenants.get(0); diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index 26b2fe5324..95bba6165f 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -22,7 +22,6 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningScheduledExecutorService; import com.google.common.util.concurrent.MoreExecutors; -import com.google.common.util.concurrent.SettableFuture; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -59,7 +58,7 @@ import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import javax.annotation.Nonnull; @@ -71,7 +70,6 @@ import java.util.Arrays; import java.util.Collections; import java.util.HashSet; import java.util.List; -import java.util.Optional; import java.util.Queue; import java.util.Random; import java.util.Set; @@ -175,21 +173,6 @@ public class DefaultDeviceStateService extends TbApplicationEventListener Function, DeviceStateData> extractDeviceStateData(Device device) { - return new Function, DeviceStateData>() { + return new Function<>() { @Nonnull @Override public DeviceStateData apply(@Nullable List data) { @@ -669,9 +640,9 @@ public class DefaultDeviceStateService extends TbApplicationEventListener(deviceId, key, value)); + new TelemetrySaveCallback<>(deviceId, key, value)); } else { - tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, DataConstants.SERVER_SCOPE, key, value, new AttributeSaveCallback<>(deviceId, key, value)); + tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, DataConstants.SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value)); } } @@ -680,18 +651,18 @@ public class DefaultDeviceStateService extends TbApplicationEventListener(deviceId, key, value)); + new TelemetrySaveCallback<>(deviceId, key, value)); } else { - tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, DataConstants.SERVER_SCOPE, key, value, new AttributeSaveCallback<>(deviceId, key, value)); + tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, DataConstants.SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value)); } } - private static class AttributeSaveCallback implements FutureCallback { + private static class TelemetrySaveCallback implements FutureCallback { private final DeviceId deviceId; private final String key; private final Object value; - AttributeSaveCallback(DeviceId deviceId, String key, Object value) { + TelemetrySaveCallback(DeviceId deviceId, String key, Object value) { this.deviceId = deviceId; this.key = key; this.value = value; diff --git a/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java index a717769310..025a29dddb 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DeviceStateService.java @@ -28,12 +28,6 @@ import org.thingsboard.server.common.msg.queue.TbCallback; */ public interface DeviceStateService extends ApplicationListener { - void onDeviceAdded(Device device); - - void onDeviceUpdated(Device device); - - void onDeviceDeleted(Device device); - void onDeviceConnect(TenantId tenantId, DeviceId deviceId); void onDeviceActivity(TenantId tenantId, DeviceId deviceId, long lastReportedActivityTime); diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java index 1ca4a41602..15cf6f1a17 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java @@ -53,7 +53,7 @@ import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.state.DefaultDeviceStateService; import org.thingsboard.server.service.state.DeviceStateService; import org.thingsboard.server.service.telemetry.sub.AlarmSubscriptionUpdate; diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java index 0f1765f316..7db0478e03 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java @@ -30,7 +30,7 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.telemetry.sub.AlarmSubscriptionUpdate; import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; @@ -42,7 +42,6 @@ import java.util.Map; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; @Slf4j @TbCoreComponent diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java index 23abba15f3..bbd3811075 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/AbstractSubscriptionService.java @@ -20,15 +20,13 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.context.ApplicationListener; -import org.springframework.context.event.EventListener; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbApplicationEventListener; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.subscription.SubscriptionManagerService; import javax.annotation.Nullable; diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java index d85ac8d736..14d264fb17 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultAlarmSubscriptionService.java @@ -46,7 +46,7 @@ import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.usagestats.TbApiUsageClient; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.subscription.SubscriptionManagerService; import org.thingsboard.server.service.subscription.TbSubscriptionUtils; diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java index 33173283b3..a3771988cf 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java @@ -45,7 +45,7 @@ import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.usagestats.TbApiUsageClient; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.subscription.TbSubscriptionUtils; import javax.annotation.Nullable; diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index bb34b14e19..55a04b108b 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java @@ -95,7 +95,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.profile.TbDeviceProfileCache; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.resource.TbResourceService; import org.thingsboard.server.service.state.DeviceStateService; @@ -265,9 +265,11 @@ public class DefaultTransportApiService implements TransportApiService { device.setCustomerId(gateway.getCustomerId()); DeviceProfile deviceProfile = deviceProfileCache.findOrCreateDeviceProfile(gateway.getTenantId(), requestMsg.getDeviceType()); device.setDeviceProfileId(deviceProfile.getId()); - device = deviceService.saveDevice(device); + Device savedDevice = deviceService.saveDevice(device); + tbClusterService.onDeviceUpdated(savedDevice, device); + device = savedDevice; + relationService.saveRelationAsync(TenantId.SYS_TENANT_ID, new EntityRelation(gateway.getId(), device.getId(), "Created")); - deviceStateService.onDeviceAdded(device); TbMsgMetaData metaData = new TbMsgMetaData(); CustomerId customerId = gateway.getCustomerId(); @@ -581,7 +583,7 @@ public class DefaultTransportApiService implements TransportApiService { device.setName(deviceName); device.setType("LwM2M"); device = deviceService.saveDevice(device); - deviceStateService.onDeviceAdded(device); + tbClusterService.onDeviceUpdated(device, null); } TransportProtos.LwM2MRegistrationResponseMsg registrationResponseMsg = TransportProtos.LwM2MRegistrationResponseMsg.newBuilder() diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java index d2267facfb..81042a47fb 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java @@ -28,7 +28,6 @@ import com.google.protobuf.MessageLite; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.RandomStringUtils; import org.awaitility.Awaitility; -import org.hamcrest.Matchers; import org.junit.After; import org.junit.Assert; import org.junit.Before; @@ -116,8 +115,7 @@ import org.thingsboard.server.gen.edge.v1.UserUpdateMsg; import org.thingsboard.server.gen.edge.v1.WidgetTypeUpdateMsg; import org.thingsboard.server.gen.edge.v1.WidgetsBundleUpdateMsg; import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.service.edge.rpc.EdgeProtoUtils; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import java.util.ArrayList; import java.util.List; diff --git a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java index 98118409bc..ed32020fb7 100644 --- a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java @@ -27,7 +27,7 @@ import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.queue.discovery.PartitionService; -import org.thingsboard.server.service.queue.TbClusterService; +import org.thingsboard.server.cluster.TbClusterService; import static org.hamcrest.CoreMatchers.is; import static org.hamcrest.MatcherAssert.assertThat; diff --git a/common/cluster-api/pom.xml b/common/cluster-api/pom.xml new file mode 100644 index 0000000000..dfe2657e4c --- /dev/null +++ b/common/cluster-api/pom.xml @@ -0,0 +1,144 @@ + + + 4.0.0 + + org.thingsboard + 3.3.0-SNAPSHOT + common + + org.thingsboard.common + cluster-api + jar + + Thingsboard Server Common Cluster API + https://thingsboard.io + + + UTF-8 + ${basedir}/../.. + + + + + org.thingsboard.common + data + + + org.thingsboard.common + message + + + org.thingsboard.common + stats + + + com.google.guava + guava + + + javax.annotation + javax.annotation-api + + + com.github.fge + json-schema-validator + + + org.slf4j + slf4j-api + + + org.slf4j + log4j-over-slf4j + + + ch.qos.logback + logback-core + + + ch.qos.logback + logback-classic + + + com.fasterxml.jackson.core + jackson-databind + + + org.springframework.boot + spring-boot-autoconfigure + provided + + + com.datastax.oss + java-driver-core + provided + + + io.dropwizard.metrics + metrics-jmx + provided + + + org.apache.commons + commons-lang3 + provided + + + junit + junit + test + + + org.mockito + mockito-core + test + + + + + + + org.xolstice.maven.plugins + protobuf-maven-plugin + + + org.apache.maven.plugins + maven-source-plugin + + + attach-sources + + jar + + + + + + org.apache.maven.plugins + maven-deploy-plugin + + false + + + + + + + diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java b/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java similarity index 83% rename from application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java rename to common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java index c5848bed58..fce6945f98 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/TbClusterService.java +++ b/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java @@ -13,26 +13,29 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.queue; +package org.thingsboard.server.cluster; -import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg; +import org.thingsboard.server.common.data.edge.EdgeEventActionType; +import org.thingsboard.server.common.data.edge.EdgeEventType; +import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg; import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; +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.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg; import org.thingsboard.server.queue.TbQueueCallback; -import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; import java.util.UUID; @@ -54,7 +57,7 @@ public interface TbClusterService { void pushNotificationToTransport(String targetServiceId, ToTransportMsg response, TbQueueCallback callback); - void onEntityStateChange(TenantId tenantId, EntityId entityId, ComponentLifecycleEvent state); + void broadcastEntityStateChangeEvent(TenantId tenantId, EntityId entityId, ComponentLifecycleEvent state); void onDeviceProfileChange(DeviceProfile deviceProfile, TbQueueCallback callback); @@ -70,7 +73,7 @@ public interface TbClusterService { void onApiStateChange(ApiUsageState apiUsageState, TbQueueCallback callback); - void onDeviceChange(Device device, TbQueueCallback callback); + void onDeviceUpdated(Device device, Device old); void onDeviceDeleted(Device device, TbQueueCallback callback); @@ -79,4 +82,6 @@ public interface TbClusterService { void onResourceDeleted(TbResource resource, TbQueueCallback callback); void onEdgeEventUpdate(TenantId tenantId, EdgeId edgeId); + + void sendNotificationMsgToEdgeService(TenantId tenantId, EdgeId edgeId, EntityId entityId, String body, EdgeEventType type, EdgeEventActionType action); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueAdmin.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueAdmin.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/TbQueueAdmin.java rename to common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueAdmin.java diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueCallback.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueCallback.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/TbQueueCallback.java rename to common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueCallback.java diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueConsumer.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueConsumer.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/TbQueueConsumer.java rename to common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueConsumer.java diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueHandler.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueHandler.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/TbQueueHandler.java rename to common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueHandler.java diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueMsg.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueMsg.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/TbQueueMsg.java rename to common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueMsg.java diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueMsgDecoder.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueMsgDecoder.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/TbQueueMsgDecoder.java rename to common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueMsgDecoder.java diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueMsgHeaders.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueMsgHeaders.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/TbQueueMsgHeaders.java rename to common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueMsgHeaders.java diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueMsgMetadata.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueMsgMetadata.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/TbQueueMsgMetadata.java rename to common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueMsgMetadata.java diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueProducer.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueProducer.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/TbQueueProducer.java rename to common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueProducer.java diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java rename to common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueResponseTemplate.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueResponseTemplate.java similarity index 100% rename from common/queue/src/main/java/org/thingsboard/server/queue/TbQueueResponseTemplate.java rename to common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueResponseTemplate.java diff --git a/common/queue/src/main/proto/jsinvoke.proto b/common/cluster-api/src/main/proto/jsinvoke.proto similarity index 100% rename from common/queue/src/main/proto/jsinvoke.proto rename to common/cluster-api/src/main/proto/jsinvoke.proto diff --git a/common/queue/src/main/proto/queue.proto b/common/cluster-api/src/main/proto/queue.proto similarity index 100% rename from common/queue/src/main/proto/queue.proto rename to common/cluster-api/src/main/proto/queue.proto diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java index 162959fa5a..82f7a4449d 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceService.java @@ -54,6 +54,8 @@ public interface DeviceService { Device saveDeviceWithCredentials(Device device, DeviceCredentials deviceCredentials); + Device saveDevice(ProvisionRequest provisionRequest, DeviceProfile profile); + Device assignDeviceToCustomer(TenantId tenantId, DeviceId deviceId, CustomerId customerId); Device unassignDeviceFromCustomer(TenantId tenantId, DeviceId deviceId); @@ -98,8 +100,6 @@ public interface DeviceService { Device assignDeviceToTenant(TenantId tenantId, Device device); - Device saveDevice(ProvisionRequest provisionRequest, DeviceProfile profile); - PageData findDevicesIdsByDeviceProfileTransportType(DeviceTransportType transportType, PageLink pageLink); Device assignDeviceToEdge(TenantId tenantId, DeviceId deviceId, EdgeId edgeId); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RpcError.java b/common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcError.java similarity index 93% rename from rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RpcError.java rename to common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcError.java index 05c27ab654..072a3e63a0 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RpcError.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcError.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.rule.engine.api; +package org.thingsboard.server.common.data.rpc; /** * @author Andrew Shvayka diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/ToDeviceActorNotificationMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/ToDeviceActorNotificationMsg.java similarity index 90% rename from rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/ToDeviceActorNotificationMsg.java rename to common/message/src/main/java/org/thingsboard/server/common/msg/ToDeviceActorNotificationMsg.java index 2f3486b5d4..036ab1b2c5 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/ToDeviceActorNotificationMsg.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/ToDeviceActorNotificationMsg.java @@ -13,9 +13,8 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.rule.engine.api.msg; +package org.thingsboard.server.common.msg; -import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.aware.DeviceAwareMsg; import org.thingsboard.server.common.msg.aware.TenantAwareMsg; diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/FromDeviceRpcResponse.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/FromDeviceRpcResponse.java similarity index 92% rename from application/src/main/java/org/thingsboard/server/service/rpc/FromDeviceRpcResponse.java rename to common/message/src/main/java/org/thingsboard/server/common/msg/rpc/FromDeviceRpcResponse.java index b7ce551c0e..06a6cb1d18 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/FromDeviceRpcResponse.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/FromDeviceRpcResponse.java @@ -13,12 +13,12 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.rpc; +package org.thingsboard.server.common.msg.rpc; import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.ToString; -import org.thingsboard.rule.engine.api.RpcError; +import org.thingsboard.server.common.data.rpc.RpcError; import java.io.Serializable; import java.util.Optional; diff --git a/common/pom.xml b/common/pom.xml index 892f4fb582..eea3c0e617 100644 --- a/common/pom.xml +++ b/common/pom.xml @@ -41,6 +41,7 @@ queue transport dao-api + cluster-api stats cache coap-server diff --git a/common/queue/pom.xml b/common/queue/pom.xml index 2b1d87fc7b..92e480c660 100644 --- a/common/queue/pom.xml +++ b/common/queue/pom.xml @@ -52,6 +52,10 @@ org.thingsboard.common stats + + org.thingsboard.common + cluster-api + org.apache.kafka kafka-clients @@ -137,13 +141,4 @@ - - - - org.xolstice.maven.plugins - protobuf-maven-plugin - - - - diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java index 0a547eb11a..7b45b5f9d9 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java @@ -341,7 +341,7 @@ public class GatewaySessionHandler { for (Map.Entry deviceEntry : jsonObj.entrySet()) { String deviceName = deviceEntry.getKey(); Futures.addCallback(checkDeviceConnected(deviceName), - new FutureCallback() { + new FutureCallback<>() { @Override public void onSuccess(@Nullable GatewayDeviceSessionCtx deviceCtx) { if (!deviceEntry.getValue().isJsonArray()) { diff --git a/pom.xml b/pom.xml index 59c160c6a9..b6d460acd5 100755 --- a/pom.xml +++ b/pom.xml @@ -880,6 +880,11 @@ dao-api ${project.version} + + org.thingsboard.common + cluster-api + ${project.version} + org.thingsboard.rule-engine rule-engine-api diff --git a/rule-engine/rule-engine-api/pom.xml b/rule-engine/rule-engine-api/pom.xml index bb348793c1..e0db697309 100644 --- a/rule-engine/rule-engine-api/pom.xml +++ b/rule-engine/rule-engine-api/pom.xml @@ -48,6 +48,11 @@ dao-api provided + + org.thingsboard.common + cluster-api + provided + org.thingsboard.common util diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcResponse.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcResponse.java index 9c43696650..c55e2ccbcf 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcResponse.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcResponse.java @@ -18,6 +18,7 @@ package org.thingsboard.rule.engine.api; import lombok.Builder; import lombok.Data; import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.rpc.RpcError; import java.util.Optional; diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java index 62ed9b0fa0..76b079c3ee 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java @@ -18,6 +18,7 @@ package org.thingsboard.rule.engine.api; import io.netty.channel.EventLoopGroup; import org.thingsboard.common.util.ListeningExecutor; import org.thingsboard.rule.engine.api.sms.SmsSenderFactory; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; @@ -188,6 +189,8 @@ public interface TbContext { DeviceService getDeviceService(); + TbClusterService getClusterService(); + DashboardService getDashboardService(); RuleEngineAlarmService getAlarmService(); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceAttributesEventNotificationMsg.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceAttributesEventNotificationMsg.java index a00b2372a2..dfbf72997b 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceAttributesEventNotificationMsg.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceAttributesEventNotificationMsg.java @@ -23,6 +23,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.AttributeKey; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.msg.MsgType; +import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg; import java.util.List; import java.util.Set; diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceCredentialsUpdateNotificationMsg.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceCredentialsUpdateNotificationMsg.java index 4b039dbf53..17413c2b0f 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceCredentialsUpdateNotificationMsg.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceCredentialsUpdateNotificationMsg.java @@ -16,16 +16,11 @@ package org.thingsboard.rule.engine.api.msg; import lombok.Data; -import lombok.Getter; -import lombok.ToString; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.kv.AttributeKey; import org.thingsboard.server.common.data.security.DeviceCredentials; -import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.msg.MsgType; - -import java.util.Set; +import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg; /** * @author Andrew Shvayka diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceEdgeUpdateMsg.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceEdgeUpdateMsg.java index 3c7fc7a68c..00a69091b2 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceEdgeUpdateMsg.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceEdgeUpdateMsg.java @@ -21,6 +21,7 @@ import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.MsgType; +import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg; @Data @AllArgsConstructor diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceNameOrTypeUpdateMsg.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceNameOrTypeUpdateMsg.java index 0418986998..47ffcc59ba 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceNameOrTypeUpdateMsg.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/msg/DeviceNameOrTypeUpdateMsg.java @@ -20,6 +20,7 @@ import lombok.Data; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.MsgType; +import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg; @Data @AllArgsConstructor diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractRelationActionNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractRelationActionNode.java index 861e2a3f0d..33c637f9d9 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractRelationActionNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractRelationActionNode.java @@ -191,6 +191,7 @@ public abstract class TbAbstractRelationActionNode log.trace("Pushed Device Created message: {}", savedDevice), throwable -> log.warn("Failed to push Device Created message: {}", savedDevice, throwable)); @@ -259,10 +260,10 @@ public abstract class TbAbstractRelationActionNode