diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java index 48425d58db..d69b9aec0b 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java @@ -64,6 +64,9 @@ public class DeviceActor extends ContextAwareActor { case DEVICE_ATTRIBUTES_UPDATE_TO_DEVICE_ACTOR_MSG: processor.processAttributesUpdate((DeviceAttributesEventNotificationMsg) msg); break; + case DEVICE_DELETE_TO_DEVICE_ACTOR_MSG: + ctx.stop(ctx.getSelf()); + break; case DEVICE_CREDENTIALS_UPDATE_TO_DEVICE_ACTOR_MSG: processor.processCredentialsUpdate(msg); break; diff --git a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java index c9f5448df3..86f338e89d 100644 --- a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java @@ -53,10 +53,13 @@ import org.thingsboard.server.common.msg.queue.PartitionChangeMsg; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; import org.thingsboard.server.common.msg.queue.RuleEngineException; import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.rule.engine.DeviceDeleteMsg; import org.thingsboard.server.service.edge.rpc.EdgeRpcService; import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper; +import java.util.HashSet; import java.util.List; +import java.util.Set; @Slf4j public class TenantActor extends RuleChainManagerActor { @@ -65,8 +68,11 @@ public class TenantActor extends RuleChainManagerActor { private boolean isCore; private ApiUsageState apiUsageState; + private Set deletedDevices; + private TenantActor(ActorSystemContext systemContext, TenantId tenantId) { super(systemContext, tenantId); + this.deletedDevices = new HashSet<>(); } boolean cantFindTenant = false; @@ -221,6 +227,10 @@ public class TenantActor extends RuleChainManagerActor { if (!isCore) { log.warn("RECEIVED INVALID MESSAGE: {}", msg); } + if (deletedDevices.contains(msg.getDeviceId())) { + log.debug("RECEIVED MESSAGE FOR DELETED DEVICE: {}", msg); + return; + } TbActorRef deviceActor = getOrCreateDeviceActor(msg.getDeviceId()); if (priority) { deviceActor.tellWithHighPriority(msg); @@ -240,7 +250,8 @@ public class TenantActor extends RuleChainManagerActor { log.info("[{}] Received API state update. Going to ENABLE Rule Engine execution.", tenantId); initRuleChains(); } - } else if (msg.getEntityId().getEntityType() == EntityType.EDGE) { + } + if (msg.getEntityId().getEntityType() == EntityType.EDGE) { EdgeId edgeId = new EdgeId(msg.getEntityId().getId()); EdgeRpcService edgeRpcService = systemContext.getEdgeRpcService(); if (msg.getEvent() == ComponentLifecycleEvent.DELETED) { @@ -249,7 +260,13 @@ public class TenantActor extends RuleChainManagerActor { Edge edge = systemContext.getEdgeService().findEdgeById(tenantId, edgeId); edgeRpcService.updateEdge(tenantId, edge); } - } else if (isRuleEngine) { + } + if (msg.getEntityId().getEntityType() == EntityType.DEVICE && ComponentLifecycleEvent.DELETED == msg.getEvent()) { + DeviceId deviceId = (DeviceId) msg.getEntityId(); + onToDeviceActorMsg(new DeviceDeleteMsg(tenantId, deviceId), true); + deletedDevices.add(deviceId); + } + if (isRuleEngine) { TbActorRef target = getEntityActorRef(msg.getEntityId()); if (target != null) { if (msg.getEntityId().getEntityType() == EntityType.RULE_CHAIN) { diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java index 72e8f44020..8682d18a62 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java @@ -112,7 +112,7 @@ public class DefaultTbNotificationEntityService implements TbNotificationEntityS public void notifyDeleteDevice(TenantId tenantId, DeviceId deviceId, CustomerId customerId, Device device, User user, Object... additionalInfo) { gatewayNotificationsService.onDeviceDeleted(device); - tbClusterService.onDeviceDeleted(device, null); + tbClusterService.onDeviceDeleted(tenantId, device, null); logEntityAction(tenantId, deviceId, device, customerId, ActionType.DELETED, user, additionalInfo); } @@ -126,6 +126,7 @@ public class DefaultTbNotificationEntityService implements TbNotificationEntityS @Override public void notifyAssignDeviceToTenant(TenantId tenantId, TenantId newTenantId, DeviceId deviceId, CustomerId customerId, Device device, Tenant tenant, User user, Object... additionalInfo) { + tbClusterService.onDeviceAssignedToTenant(tenantId, device); logEntityAction(tenantId, deviceId, device, customerId, ActionType.ASSIGNED_TO_TENANT, user, additionalInfo); pushAssignedFromNotification(tenant, newTenantId, device); } 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 f96e49930c..af2a54a7ac 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 @@ -308,10 +308,17 @@ public class DefaultTbClusterService implements TbClusterService { } @Override - 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); + public void onDeviceDeleted(TenantId tenantId, Device device, TbQueueCallback callback) { + DeviceId deviceId = device.getId(); + broadcastEntityDeleteToTransport(tenantId, deviceId, device.getName(), callback); + sendDeviceStateServiceEvent(tenantId, deviceId, false, false, true); + broadcastEntityStateChangeEvent(tenantId, deviceId, ComponentLifecycleEvent.DELETED); + } + + @Override + public void onDeviceAssignedToTenant(TenantId oldTenantId, Device device) { + onDeviceDeleted(oldTenantId, device, null); + sendDeviceStateServiceEvent(device.getTenantId(), device.getId(), true, false, false); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/queue/ProtoUtils.java b/application/src/main/java/org/thingsboard/server/service/queue/ProtoUtils.java index 36038568bc..946993151d 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/ProtoUtils.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/ProtoUtils.java @@ -46,6 +46,7 @@ import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequestActorMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceCredentialsUpdateNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceDeleteMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.gen.transport.TransportProtos; @@ -384,6 +385,21 @@ public class ProtoUtils { ); } + private static TransportProtos.DeviceDeleteMsgProto toProto(DeviceDeleteMsg msg) { + return TransportProtos.DeviceDeleteMsgProto.newBuilder() + .setTenantIdMSB(msg.getTenantId().getId().getMostSignificantBits()) + .setTenantIdLSB(msg.getTenantId().getId().getLeastSignificantBits()) + .setDeviceIdMSB(msg.getDeviceId().getId().getMostSignificantBits()) + .setDeviceIdLSB(msg.getDeviceId().getId().getLeastSignificantBits()) + .build(); + } + + private static DeviceDeleteMsg fromProto(TransportProtos.DeviceDeleteMsgProto proto) { + return new DeviceDeleteMsg( + TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())), + new DeviceId(new UUID(proto.getDeviceIdMSB(), proto.getDeviceIdLSB()))); + } + public static TransportProtos.ToDeviceActorNotificationMsgProto toProto(ToDeviceActorNotificationMsg msg) { if (msg instanceof DeviceEdgeUpdateMsg) { DeviceEdgeUpdateMsg updateMsg = (DeviceEdgeUpdateMsg) msg; @@ -413,6 +429,10 @@ public class ProtoUtils { RemoveRpcActorMsg updateMsg = (RemoveRpcActorMsg) msg; TransportProtos.RemoveRpcActorMsgProto proto = toProto(updateMsg); return TransportProtos.ToDeviceActorNotificationMsgProto.newBuilder().setRemoveRpcActorMsg(proto).build(); + } else if (msg instanceof DeviceDeleteMsg) { + DeviceDeleteMsg updateMsg = (DeviceDeleteMsg) msg; + TransportProtos.DeviceDeleteMsgProto proto = toProto(updateMsg); + return TransportProtos.ToDeviceActorNotificationMsgProto.newBuilder().setDeviceDeleteMsg(proto).build(); } return null; } @@ -432,6 +452,8 @@ public class ProtoUtils { return fromProto(proto.getFromDeviceRpcResponseMsg()); } else if (proto.hasRemoveRpcActorMsg()) { return fromProto(proto.getRemoveRpcActorMsg()); + } else if (proto.hasDeviceDeleteMsg()) { + return fromProto(proto.getDeviceDeleteMsg()); } return null; } diff --git a/application/src/test/java/org/thingsboard/server/controller/DeviceControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/DeviceControllerTest.java index 778a11a552..cefe24e935 100644 --- a/application/src/test/java/org/thingsboard/server/controller/DeviceControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/DeviceControllerTest.java @@ -51,7 +51,6 @@ import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmSeverity; -import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.id.CustomerId; @@ -74,6 +73,7 @@ import org.thingsboard.server.dao.exception.DeviceCredentialsValidationException import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.service.gateway_device.GatewayNotificationsService; +import org.thingsboard.server.service.state.DeviceStateService; import java.util.ArrayList; import java.util.List; @@ -82,6 +82,8 @@ import java.util.concurrent.TimeUnit; import static org.assertj.core.api.Assertions.assertThat; import static org.hamcrest.Matchers.containsString; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.argThat; import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @@ -106,6 +108,9 @@ public class DeviceControllerTest extends AbstractControllerTest { @SpyBean private GatewayNotificationsService gatewayNotificationsService; + @SpyBean + private DeviceStateService deviceStateService; + @Autowired private DeviceDao deviceDao; @@ -1340,6 +1345,28 @@ public class DeviceControllerTest extends AbstractControllerTest { ActionType.ASSIGNED_TO_TENANT, savedDifferentTenant.getId().getId().toString(), savedDifferentTenant.getTitle()); testNotificationUpdateGatewayNever(); + Mockito.verify(deviceStateService, times(1)).onQueueMsg( + argThat(proto -> + proto.getTenantIdMSB() == savedTenant.getUuidId().getMostSignificantBits() && + proto.getTenantIdLSB() == savedTenant.getUuidId().getLeastSignificantBits() && + proto.getDeviceIdMSB() == savedDevice.getUuidId().getMostSignificantBits() && + proto.getDeviceIdLSB() == savedDevice.getUuidId().getLeastSignificantBits() && + proto.getDeleted() + ), + any() + ); + + Mockito.verify(deviceStateService, times(1)).onQueueMsg( + argThat(proto -> + proto.getTenantIdMSB() == savedDifferentTenant.getUuidId().getMostSignificantBits() && + proto.getTenantIdLSB() == savedDifferentTenant.getUuidId().getLeastSignificantBits() && + proto.getDeviceIdMSB() == savedDevice.getUuidId().getMostSignificantBits() && + proto.getDeviceIdLSB() == savedDevice.getUuidId().getLeastSignificantBits() && + proto.getAdded() + ), + any() + ); + login("tenant9@thingsboard.org", "testPassword1"); Device foundDevice1 = doGet("/api/device/" + assignedDevice.getId().getId(), Device.class); diff --git a/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java b/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java index 5b1eb1ae4f..5c3b9076ff 100644 --- a/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java +++ b/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java @@ -83,7 +83,9 @@ public interface TbClusterService extends TbQueueClusterService { void onDeviceUpdated(Device device, Device old); - void onDeviceDeleted(Device device, TbQueueCallback callback); + void onDeviceDeleted(TenantId tenantId, Device device, TbQueueCallback callback); + + void onDeviceAssignedToTenant(TenantId oldTenantId, Device device); void onResourceChange(TbResource resource, TbQueueCallback callback); diff --git a/common/cluster-api/src/main/proto/queue.proto b/common/cluster-api/src/main/proto/queue.proto index 0cb91faee0..ffae1048dd 100644 --- a/common/cluster-api/src/main/proto/queue.proto +++ b/common/cluster-api/src/main/proto/queue.proto @@ -1010,6 +1010,13 @@ message RemoveRpcActorMsgProto { int64 deviceIdLSB = 6; } +message DeviceDeleteMsgProto { + int64 tenantIdMSB = 1; + int64 tenantIdLSB = 2; + int64 deviceIdMSB = 3; + int64 deviceIdLSB = 4; +} + message ToDeviceActorNotificationMsgProto { DeviceEdgeUpdateMsgProto deviceEdgeUpdateMsg = 1; DeviceNameOrTypeUpdateMsgProto deviceNameOrTypeMsg = 2; @@ -1018,6 +1025,7 @@ message ToDeviceActorNotificationMsgProto { ToDeviceRpcRequestActorMsgProto toDeviceRpcRequestMsg = 5; FromDeviceRpcResponseActorMsgProto fromDeviceRpcResponseMsg = 6; RemoveRpcActorMsgProto removeRpcActorMsg = 7; + DeviceDeleteMsgProto deviceDeleteMsg = 8; } /** diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java b/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java index 0a39d8d51b..c3d3623d8e 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java @@ -93,6 +93,8 @@ public enum MsgType { DEVICE_NAME_OR_TYPE_UPDATE_TO_DEVICE_ACTOR_MSG, + DEVICE_DELETE_TO_DEVICE_ACTOR_MSG, + DEVICE_EDGE_UPDATE_TO_DEVICE_ACTOR_MSG, DEVICE_RPC_REQUEST_TO_DEVICE_ACTOR_MSG, diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceDeleteMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceDeleteMsg.java new file mode 100644 index 0000000000..1bdab7409b --- /dev/null +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceDeleteMsg.java @@ -0,0 +1,36 @@ +/** + * Copyright © 2016-2023 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.msg.rule.engine; + +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 +public class DeviceDeleteMsg implements ToDeviceActorNotificationMsg { + + private static final long serialVersionUID = 4679029228395462172L; + + private final TenantId tenantId; + private final DeviceId deviceId; + + @Override + public MsgType getMsgType() { + return MsgType.DEVICE_DELETE_TO_DEVICE_ACTOR_MSG; + } +}