From 9085b42d4304c000f5a66cfab275899458300ce1 Mon Sep 17 00:00:00 2001 From: Nikita Mazurenko Date: Fri, 21 Nov 2025 14:26:47 +0200 Subject: [PATCH 1/3] Add handling of ENTITY_DELETED_RPC_MESSAGE during edge user message processing --- .../rpc/processor/user/BaseUserProcessor.java | 10 +++++++- .../rpc/processor/user/UserEdgeProcessor.java | 14 ++++++++++- .../thingsboard/server/edge/UserEdgeTest.java | 25 +++++++++++++++++++ 3 files changed, 47 insertions(+), 2 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java index c836c14233..14e4da53c0 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java @@ -78,6 +78,15 @@ public abstract class BaseUserProcessor extends BaseEdgeProcessor { return Pair.of(isCreated, userEmailUpdated); } + protected User deleteUser(TenantId tenantId, UserId userId) { + User userById = edgeCtx.getUserService().findUserById(tenantId, userId); + if (userById == null) { + throw new IllegalArgumentException(String.format("[%s] Failed to find User with id [%s]", tenantId, userId)); + } + edgeCtx.getUserService().deleteUser(tenantId, userById); + return userById; + } + protected void updateUserCredentials(TenantId tenantId, UserCredentialsUpdateMsg updateMsg) { UserCredentials userCredentialsFromUpdateMsg = JacksonUtil.fromString(updateMsg.getEntity(), UserCredentials.class, true); if (userCredentialsFromUpdateMsg == null) { @@ -117,7 +126,6 @@ public abstract class BaseUserProcessor extends BaseEdgeProcessor { tenantId, user.getName(), updateMsg, e); throw new RuntimeException(e); } - } protected abstract void setCustomerId(TenantId tenantId, CustomerId customerId, User user, UserUpdateMsg userUpdateMsg); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java index e7813caae9..cfe888b1a5 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java @@ -61,6 +61,10 @@ public class UserEdgeProcessor extends BaseUserProcessor implements UserProcesso saveOrUpdateUser(tenantId, userId, userUpdateMsg, edge); yield Futures.immediateFuture(null); } + case ENTITY_DELETED_RPC_MESSAGE -> { + deleteUserAndPushEntityDeletedEventToRuleEngine(tenantId, userId, edge); + yield Futures.immediateFuture(null); + } default -> handleUnsupportedMsgType(userUpdateMsg.getMsgType()); }; } catch (DataValidationException e) { @@ -116,7 +120,15 @@ public class UserEdgeProcessor extends BaseUserProcessor implements UserProcesso } } - @Override + private void deleteUserAndPushEntityDeletedEventToRuleEngine(TenantId tenantId, UserId userId, Edge edge) { + User removedUser = deleteUser(tenantId, userId); + CustomerId userCustomerId = removedUser.getCustomerId(); + String userAsString = JacksonUtil.toString(removedUser); + TbMsgMetaData msgMetaData = getEdgeActionTbMsgMetaData(edge, userCustomerId); + pushEntityEventToRuleEngine(tenantId, userId, userCustomerId, TbMsgType.ENTITY_DELETED, userAsString, msgMetaData); + } + + @Override public DownlinkMsg convertEdgeEventToDownlink(EdgeEvent edgeEvent, EdgeVersion edgeVersion) { UserId userId = new UserId(edgeEvent.getEntityId()); switch (edgeEvent.getAction()) { diff --git a/application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java index 8b5261150b..f590b21b84 100644 --- a/application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java @@ -148,6 +148,31 @@ public class UserEdgeTest extends AbstractEdgeTest { testAutoGeneratedCodeByProtobuf(userCredentialsUpdateMsg); } + @Test + public void testSendUserDeleteFromEdgeToCloud() throws Exception { + // create customer + Customer savedCustomer = createAndAssignCustomerToEdge(); + + // create user + User customerUser = buildUser(Authority.CUSTOMER_USER, savedCustomer.getId()); + User savedCustomerUser = createAndVerifyUserOnEdge(customerUser); + + // simulate user removal event from edge to cloud + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); + UserUpdateMsg.newBuilder().setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE) + .setIdMSB(savedCustomerUser.getUuidId().getMostSignificantBits()) + .setIdLSB(savedCustomerUser.getUuidId().getLeastSignificantBits()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + + // expect edge message sent & cloud message response + edgeImitator.expectResponsesAmount(1); + edgeImitator.expectMessageAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + Assert.assertTrue(edgeImitator.waitForResponses()); + Assert.assertTrue(edgeImitator.waitForMessages()); + } + private Customer createAndAssignCustomerToEdge() throws Exception { edgeImitator.expectMessageAmount(1); Customer customer = new Customer(); From 12ed2a56be5e728a6e8e3f127db40c534fa595c5 Mon Sep 17 00:00:00 2001 From: Nikita Mazurenko Date: Fri, 21 Nov 2025 15:40:37 +0200 Subject: [PATCH 2/3] Remove exception when user doesn't exist during edge user delete event processing --- .../service/edge/rpc/processor/user/BaseUserProcessor.java | 3 ++- .../service/edge/rpc/processor/user/UserEdgeProcessor.java | 3 +++ 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java index 14e4da53c0..897e49a7cd 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java @@ -81,7 +81,8 @@ public abstract class BaseUserProcessor extends BaseEdgeProcessor { protected User deleteUser(TenantId tenantId, UserId userId) { User userById = edgeCtx.getUserService().findUserById(tenantId, userId); if (userById == null) { - throw new IllegalArgumentException(String.format("[%s] Failed to find User with id [%s]", tenantId, userId)); + log.trace("[{}] User with id {} does not exist", tenantId, userId); + return null; } edgeCtx.getUserService().deleteUser(tenantId, userById); return userById; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java index cfe888b1a5..7905e802bd 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java @@ -122,6 +122,9 @@ public class UserEdgeProcessor extends BaseUserProcessor implements UserProcesso private void deleteUserAndPushEntityDeletedEventToRuleEngine(TenantId tenantId, UserId userId, Edge edge) { User removedUser = deleteUser(tenantId, userId); + if (removedUser == null) { + return; + } CustomerId userCustomerId = removedUser.getCustomerId(); String userAsString = JacksonUtil.toString(removedUser); TbMsgMetaData msgMetaData = getEdgeActionTbMsgMetaData(edge, userCustomerId); From bc1af0ba3bc290ecd7bf5f95a01782a977be5171 Mon Sep 17 00:00:00 2001 From: Nikita Mazurenko Date: Fri, 21 Nov 2025 17:04:42 +0200 Subject: [PATCH 3/3] Fix testSendUserDeleteFromEdgeToCloud --- .../rpc/processor/user/UserEdgeProcessor.java | 2 +- .../thingsboard/server/edge/UserEdgeTest.java | 28 +++++++++++++++---- 2 files changed, 23 insertions(+), 7 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java index 7905e802bd..968d573618 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java @@ -131,7 +131,7 @@ public class UserEdgeProcessor extends BaseUserProcessor implements UserProcesso pushEntityEventToRuleEngine(tenantId, userId, userCustomerId, TbMsgType.ENTITY_DELETED, userAsString, msgMetaData); } - @Override + @Override public DownlinkMsg convertEdgeEventToDownlink(EdgeEvent edgeEvent, EdgeVersion edgeVersion) { UserId userId = new UserId(edgeEvent.getEntityId()); switch (edgeEvent.getAction()) { diff --git a/application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java index f590b21b84..6bf51eb159 100644 --- a/application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java @@ -20,8 +20,11 @@ import org.junit.Assert; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder; +import org.springframework.test.web.servlet.ResultMatcher; +import org.testcontainers.shaded.org.awaitility.Awaitility; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Customer; +import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.edge.Edge; @@ -41,6 +44,7 @@ import org.thingsboard.server.service.security.model.ChangePasswordRequest; import java.util.Optional; import java.util.UUID; +import java.util.concurrent.TimeUnit; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.thingsboard.server.dao.user.UserServiceImpl.DEFAULT_TOKEN_LENGTH; @@ -158,19 +162,31 @@ public class UserEdgeTest extends AbstractEdgeTest { User savedCustomerUser = createAndVerifyUserOnEdge(customerUser); // simulate user removal event from edge to cloud - UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); - UserUpdateMsg.newBuilder().setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE) + UserUpdateMsg.Builder userUpdateMsg = UserUpdateMsg.newBuilder().setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE) .setIdMSB(savedCustomerUser.getUuidId().getMostSignificantBits()) .setIdLSB(savedCustomerUser.getUuidId().getLeastSignificantBits()); - testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + UplinkMsg uplink = UplinkMsg.newBuilder() + .setUplinkMsgId(EdgeUtils.nextPositiveInt()) + .addUserUpdateMsg(userUpdateMsg).build(); + + testAutoGeneratedCodeByProtobuf(userUpdateMsg); // expect edge message sent & cloud message response edgeImitator.expectResponsesAmount(1); - edgeImitator.expectMessageAmount(1); - edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + edgeImitator.sendUplinkMsg(uplink); Assert.assertTrue(edgeImitator.waitForResponses()); - Assert.assertTrue(edgeImitator.waitForMessages()); + + loginTenantAdmin(); + Awaitility.await().atMost(10, TimeUnit.SECONDS) + .until(() -> { + try { + doGet("/api/user/" + savedCustomerUser.getId(), User.class, status().isNotFound()); + return true; + } catch (Throwable ex) { + return false; + } + }); } private Customer createAndAssignCustomerToEdge() throws Exception {