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