Browse Source

Merge pull request #86 from MazurenkoNick/support-user-sync

Handle ENTITY_DELETED_RPC_MESSAGE during edge user message processing
pull/14352/head
Volodymyr Babak 10 months ago
committed by GitHub
parent
commit
1ba2cfc93a
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 11
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java
  2. 15
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java
  3. 41
      application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java

11
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java

@ -78,6 +78,16 @@ 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) {
log.trace("[{}] User with id {} does not exist", tenantId, userId);
return null;
}
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 +127,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);

15
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,6 +120,17 @@ 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);
pushEntityEventToRuleEngine(tenantId, userId, userCustomerId, TbMsgType.ENTITY_DELETED, userAsString, msgMetaData);
}
@Override
public DownlinkMsg convertEdgeEventToDownlink(EdgeEvent edgeEvent, EdgeVersion edgeVersion) {
UserId userId = new UserId(edgeEvent.getEntityId());

41
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;
@ -148,6 +152,43 @@ 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
UserUpdateMsg.Builder userUpdateMsg = UserUpdateMsg.newBuilder().setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE)
.setIdMSB(savedCustomerUser.getUuidId().getMostSignificantBits())
.setIdLSB(savedCustomerUser.getUuidId().getLeastSignificantBits());
UplinkMsg uplink = UplinkMsg.newBuilder()
.setUplinkMsgId(EdgeUtils.nextPositiveInt())
.addUserUpdateMsg(userUpdateMsg).build();
testAutoGeneratedCodeByProtobuf(userUpdateMsg);
// expect edge message sent & cloud message response
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(uplink);
Assert.assertTrue(edgeImitator.waitForResponses());
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 {
edgeImitator.expectMessageAmount(1);
Customer customer = new Customer();

Loading…
Cancel
Save