diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index 5c5eea32fd..0f469b15b8 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -73,6 +73,7 @@ import org.thingsboard.server.service.edge.rpc.processor.resource.ResourceEdgePr import org.thingsboard.server.service.edge.rpc.processor.rule.RuleChainEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.rule.RuleChainMetadataEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.telemetry.TelemetryEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.user.UserProcessor; import org.thingsboard.server.service.edge.rpc.sync.EdgeRequestsService; import org.thingsboard.server.service.executors.GrpcCallbackExecutorService; @@ -261,6 +262,9 @@ public class EdgeContextComponent { @Autowired private CalculatedFieldProcessor calculatedFieldProcessor; + @Autowired + private UserProcessor userProcessor; + public EdgeProcessor getProcessor(EdgeEventType edgeEventType) { EdgeProcessor processor = processorMap.get(edgeEventType); if (processor == null) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index a86409c6af..393d1ba3da 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -83,6 +83,8 @@ import org.thingsboard.server.gen.edge.v1.SyncCompletedMsg; import org.thingsboard.server.gen.edge.v1.UplinkMsg; import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg; import org.thingsboard.server.gen.edge.v1.UserCredentialsRequestMsg; +import org.thingsboard.server.gen.edge.v1.UserCredentialsUpdateMsg; +import org.thingsboard.server.gen.edge.v1.UserUpdateMsg; import org.thingsboard.server.gen.edge.v1.WidgetBundleTypesRequestMsg; import org.thingsboard.server.service.edge.EdgeContextComponent; import org.thingsboard.server.service.edge.EdgeMsgConstructorUtils; @@ -100,6 +102,7 @@ import java.util.UUID; import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.function.BiConsumer; @@ -121,6 +124,7 @@ public abstract class EdgeGrpcSession implements Closeable { private final EdgeSessionState sessionState = new EdgeSessionState(); private final ReentrantLock downlinkMsgLock = new ReentrantLock(); + private final Lock sequenceDependencyLock = new ReentrantLock(); protected EdgeContextComponent ctx; protected Edge edge; @@ -934,6 +938,26 @@ public abstract class EdgeGrpcSession implements Closeable { result.add(ctx.getCalculatedFieldProcessor().processCalculatedFieldMsgFromEdge(edge.getTenantId(), edge, calculatedFieldUpdateMsg)); } } + if (uplinkMsg.getUserUpdateMsgCount() > 0) { + for (UserUpdateMsg userUpdateMsg : uplinkMsg.getUserUpdateMsgList()) { + sequenceDependencyLock.lock(); + try { + result.add(ctx.getUserProcessor().processUserMsgFromEdge(edge.getTenantId(), edge, userUpdateMsg)); + } finally { + sequenceDependencyLock.unlock(); + } + } + } + if (uplinkMsg.getUserCredentialsUpdateMsgCount() > 0) { + for (UserCredentialsUpdateMsg userCredentialsUpdateMsg : uplinkMsg.getUserCredentialsUpdateMsgList()) { + sequenceDependencyLock.lock(); + try { + result.add(ctx.getUserProcessor().processUserCredentialsMsgFromEdge(edge.getTenantId(), edge, userCredentialsUpdateMsg)); + } finally { + sequenceDependencyLock.unlock(); + } + } + } } catch (Exception e) { String failureMsg = String.format("Can't process uplink msg [%s] from edge", uplinkMsg); log.trace("[{}][{}] Can't process uplink msg [{}]", tenantId, edge.getId(), uplinkMsg, e); 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 new file mode 100644 index 0000000000..92d5664a4f --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java @@ -0,0 +1,130 @@ +/** + * Copyright © 2016-2025 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.service.edge.rpc.processor.user; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.data.util.Pair; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.UserId; +import org.thingsboard.server.common.data.security.UserCredentials; +import org.thingsboard.server.dao.service.DataValidator; +import org.thingsboard.server.gen.edge.v1.UserCredentialsUpdateMsg; +import org.thingsboard.server.gen.edge.v1.UserUpdateMsg; +import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; + +@Slf4j +public abstract class BaseUserProcessor extends BaseEdgeProcessor { + + @Autowired + private DataValidator userValidator; + + protected Pair saveOrUpdateUser(TenantId tenantId, UserId userId, UserUpdateMsg userUpdateMsg) { + boolean isCreated = false; + boolean userEmailUpdated = false; + + try { + User user = JacksonUtil.fromString(userUpdateMsg.getEntity(), User.class, true); + if (user == null) { + throw new RuntimeException("[{" + tenantId + "}] userUpdateMsg {" + userUpdateMsg + "} cannot be converted to user"); + } + + User userById = edgeCtx.getUserService().findUserById(tenantId, userId); + if (userById == null) { + isCreated = true; + user.setId(null); + } else { + user.setId(userId); + } + + String userEmail = user.getEmail(); + User existing = edgeCtx.getUserService().findUserByTenantIdAndEmail(tenantId, user.getEmail()); + + if (existing != null && !existing.getId().equals(user.getId())) { + String[] splitEmail = userEmail.split("@"); + userEmail = splitEmail[0] + "_" + StringUtils.randomAlphanumeric(15) + "@" + splitEmail[1]; + log.warn("[{}] User with email {} already exists. Renaming User email to {}", + tenantId, user.getEmail(), userEmail); + userEmailUpdated = true; + } + user.setEmail(userEmail); + setCustomerId(tenantId, isCreated ? null : userById.getCustomerId(), user, userUpdateMsg); + + userValidator.validate(user, User::getTenantId); + + if (isCreated) { + user.setId(userId); + } + + edgeCtx.getUserService().saveUser(tenantId, user, false); + } catch (Exception e) { + log.error("[{}] Failed to process user update msg [{}]", tenantId, userUpdateMsg, e); + throw e; + } + + return Pair.of(isCreated, userEmailUpdated); + } + + protected void updateUserCredentials(TenantId tenantId, UserCredentialsUpdateMsg updateMsg) { + UserCredentials userCredentialsFromUpdateMsg = JacksonUtil.fromString(updateMsg.getEntity(), UserCredentials.class, true); + if (userCredentialsFromUpdateMsg == null) { + throw new RuntimeException(String.format("[%s] Failed to parse UserCredentials from updateMsg: %s", tenantId, updateMsg)); + } + + User user = edgeCtx.getUserService().findUserById(tenantId, userCredentialsFromUpdateMsg.getUserId()); + if (user == null) { + log.warn("[{}] Can't find user by id [{}] skipping credentials update. UserCredentialsUpdateMsg [{}]", + tenantId, userCredentialsFromUpdateMsg.getUserId(), updateMsg); + return; + } + + log.debug("[{}] Updating user credentials for user [{}]. New credentials Id [{}], enabled [{}]", + tenantId, user.getName(), userCredentialsFromUpdateMsg.getId(), userCredentialsFromUpdateMsg.isEnabled()); + + try { + UserCredentials existing = edgeCtx.getUserService().findUserCredentialsByUserId(tenantId, user.getId()); + boolean created = existing == null; + + UserCredentials updated = created ? new UserCredentials() : existing; + updated.setId(userCredentialsFromUpdateMsg.getId()); + updated.setUserId(user.getId()); + updated.setEnabled(userCredentialsFromUpdateMsg.isEnabled()); + updated.setActivateToken(userCredentialsFromUpdateMsg.getActivateToken()); + updated.setAdditionalInfo(userCredentialsFromUpdateMsg.getAdditionalInfo()); + updated.setPassword(userCredentialsFromUpdateMsg.getPassword()); + updated.setResetToken(userCredentialsFromUpdateMsg.getResetToken()); + + + if (created) { + edgeCtx.getUserService().saveUserCredentials(tenantId, updated, false); + } else { + edgeCtx.getUserService().replaceUserCredentials(tenantId, updated, existing.getId(), false); + } + } catch (Exception e) { + log.error("[{}] Can't update user credentials for user [{}], userCredentialsUpdateMsg [{}]", + 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 fdd03d63f3..e0bd07c6c6 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 @@ -15,26 +15,106 @@ */ package org.thingsboard.server.service.edge.rpc.processor.user; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; +import org.springframework.data.util.Pair; import org.springframework.stereotype.Component; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; +import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.security.UserCredentials; +import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.EdgeVersion; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.edge.v1.UserCredentialsUpdateMsg; +import org.thingsboard.server.gen.edge.v1.UserUpdateMsg; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.EdgeMsgConstructorUtils; -import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; + +import java.util.UUID; @Slf4j @Component @TbCoreComponent -public class UserEdgeProcessor extends BaseEdgeProcessor { +public class UserEdgeProcessor extends BaseUserProcessor implements UserProcessor { + + @Override + public ListenableFuture processUserMsgFromEdge(TenantId tenantId, Edge edge, UserUpdateMsg userUpdateMsg) { + log.trace("[{}] executing processUserMsgFromEdge [{}] from edge [{}]", tenantId, userUpdateMsg, edge.getId()); + UserId userId = new UserId(new UUID(userUpdateMsg.getIdMSB(), userUpdateMsg.getIdLSB())); + try { + edgeSynchronizationManager.getEdgeId().set(edge.getId()); + + return switch (userUpdateMsg.getMsgType()) { + case ENTITY_CREATED_RPC_MESSAGE, ENTITY_UPDATED_RPC_MESSAGE -> { + saveOrUpdateUser(tenantId, userId, userUpdateMsg, edge); + yield Futures.immediateFuture(null); + } + default -> handleUnsupportedMsgType(userUpdateMsg.getMsgType()); + }; + } catch (DataValidationException e) { + if (e.getMessage().contains("limit reached")) { + log.warn("[{}] Number of allowed users violated {}", tenantId, userUpdateMsg, e); + return Futures.immediateFuture(null); + } else { + return Futures.immediateFailedFuture(e); + } + } finally { + edgeSynchronizationManager.getEdgeId().remove(); + } + } + + @Override + public ListenableFuture processUserCredentialsMsgFromEdge(TenantId tenantId, Edge edge, UserCredentialsUpdateMsg userCredentialsUpdateMsg) { + log.debug("[{}] Executing processUserCredentialsMsgFromEdge, userCredentialsUpdateMsg [{}]", tenantId, userCredentialsUpdateMsg); + try { + edgeSynchronizationManager.getEdgeId().set(edge.getId()); + + super.updateUserCredentials(tenantId, userCredentialsUpdateMsg); + } finally { + edgeSynchronizationManager.getEdgeId().remove(); + } + return Futures.immediateFuture(null); + } + + private void saveOrUpdateUser(TenantId tenantId, UserId userId, UserUpdateMsg userUpdateMsg, Edge edge) { + Pair resultPair = super.saveOrUpdateUser(tenantId, userId, userUpdateMsg); + boolean isCreated = resultPair.getFirst(); + if (isCreated) { + createRelationFromEdge(tenantId, edge.getId(), userId); + pushUserCreatedEventToRuleEngine(tenantId, edge, userId); + } + + Boolean userEmailUpdated = resultPair.getSecond(); + + if (userEmailUpdated) { + saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.USER, EdgeEventActionType.UPDATED, userId, null); + } + } + + private void pushUserCreatedEventToRuleEngine(TenantId tenantId, Edge edge, UserId userId) { + try { + User user = edgeCtx.getUserService().findUserById(tenantId, userId); + if (user != null) { + String userAsString = JacksonUtil.toString(user); + TbMsgMetaData msgMetaData = getEdgeActionTbMsgMetaData(edge, user.getCustomerId()); + pushEntityEventToRuleEngine(tenantId, userId, user.getCustomerId(), TbMsgType.ENTITY_CREATED, userAsString, msgMetaData); + } + } catch (Exception e) { + log.warn("[{}][{}] Failed to push user action to rule engine: {}", tenantId, userId, TbMsgType.ENTITY_CREATED.name(), e); + } + } @Override public DownlinkMsg convertEdgeEventToDownlink(EdgeEvent edgeEvent, EdgeVersion edgeVersion) { @@ -48,7 +128,7 @@ public class UserEdgeProcessor extends BaseEdgeProcessor { .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) .addUserUpdateMsg(EdgeMsgConstructorUtils.constructUserUpdatedMsg(msgType, user)); UserCredentials userCredentialsByUserId = edgeCtx.getUserService().findUserCredentialsByUserId(edgeEvent.getTenantId(), userId); - if (userCredentialsByUserId != null && userCredentialsByUserId.isEnabled()) { + if (userCredentialsByUserId != null) { builder.addUserCredentialsUpdateMsg(EdgeMsgConstructorUtils.constructUserCredentialsUpdatedMsg(userCredentialsByUserId)); } return builder.build(); @@ -62,11 +142,10 @@ public class UserEdgeProcessor extends BaseEdgeProcessor { } case CREDENTIALS_UPDATED -> { UserCredentials userCredentialsByUserId = edgeCtx.getUserService().findUserCredentialsByUserId(edgeEvent.getTenantId(), userId); - if (userCredentialsByUserId != null && userCredentialsByUserId.isEnabled()) { - UserCredentialsUpdateMsg userCredentialsUpdateMsg = EdgeMsgConstructorUtils.constructUserCredentialsUpdatedMsg(userCredentialsByUserId); + if (userCredentialsByUserId != null) { return DownlinkMsg.newBuilder() .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) - .addUserCredentialsUpdateMsg(userCredentialsUpdateMsg) + .addUserCredentialsUpdateMsg(EdgeMsgConstructorUtils.constructUserCredentialsUpdatedMsg(userCredentialsByUserId)) .build(); } } @@ -79,4 +158,10 @@ public class UserEdgeProcessor extends BaseEdgeProcessor { return EdgeEventType.USER; } + @Override + protected void setCustomerId(TenantId tenantId, CustomerId customerId, User user, UserUpdateMsg userUpdateMsg) { + CustomerId customerUUID = user.getCustomerId() != null ? user.getCustomerId() : customerId; + user.setCustomerId(customerUUID); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserProcessor.java new file mode 100644 index 0000000000..dd9659c1f4 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserProcessor.java @@ -0,0 +1,31 @@ +/** + * Copyright © 2016-2025 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.service.edge.rpc.processor.user; + +import com.google.common.util.concurrent.ListenableFuture; +import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.gen.edge.v1.UserCredentialsUpdateMsg; +import org.thingsboard.server.gen.edge.v1.UserUpdateMsg; +import org.thingsboard.server.service.edge.rpc.processor.EdgeProcessor; + +public interface UserProcessor extends EdgeProcessor { + + ListenableFuture processUserMsgFromEdge(TenantId tenantId, Edge edge, UserUpdateMsg userUpdateMsg); + + ListenableFuture processUserCredentialsMsgFromEdge(TenantId tenantId, Edge edge, UserCredentialsUpdateMsg userCredentialsUpdateMsg); + +} 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 05debdc015..84442e2d0a 100644 --- a/application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.edge; +import com.fasterxml.jackson.databind.JsonNode; import com.google.protobuf.AbstractMessage; import org.junit.Assert; import org.junit.Test; @@ -22,8 +23,12 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Customer; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.UserCredentialsId; +import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.UserCredentials; import org.thingsboard.server.dao.service.DaoSqlTest; @@ -32,11 +37,14 @@ import org.thingsboard.server.gen.edge.v1.UplinkMsg; import org.thingsboard.server.gen.edge.v1.UserCredentialsRequestMsg; import org.thingsboard.server.gen.edge.v1.UserCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.v1.UserUpdateMsg; +import org.thingsboard.server.service.edge.EdgeMsgConstructorUtils; import org.thingsboard.server.service.security.model.ChangePasswordRequest; import java.util.Optional; +import java.util.UUID; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; +import static org.thingsboard.server.dao.user.UserServiceImpl.DEFAULT_TOKEN_LENGTH; @DaoSqlTest public class UserEdgeTest extends AbstractEdgeTest { @@ -44,46 +52,35 @@ public class UserEdgeTest extends AbstractEdgeTest { @Autowired private BCryptPasswordEncoder passwordEncoder; + private static final String DEFAULT_FIRST_NAME = "Boris"; + private static final String UPDATED_LAST_NAME = "Borisov"; + @Test public void testCreateUpdateDeleteTenantUser() throws Exception { // create user edgeImitator.expectMessageAmount(3); - User newTenantAdmin = new User(); - newTenantAdmin.setAuthority(Authority.TENANT_ADMIN); - newTenantAdmin.setTenantId(tenantId); - newTenantAdmin.setEmail("tenantAdmin@thingsboard.org"); - newTenantAdmin.setFirstName("Boris"); - newTenantAdmin.setLastName("Johnson"); + User newTenantAdmin = buildUser(Authority.TENANT_ADMIN, null, "tenantAdmin@thingsboard.org", DEFAULT_FIRST_NAME, "Johnson"); User savedTenantAdmin = createUser(newTenantAdmin, "tenant"); Assert.assertTrue(edgeImitator.waitForMessages()); // wait 3 messages - x1 user update msg and x2 user credentials update msgs (create + authenticate user) Assert.assertEquals(1, edgeImitator.findAllMessagesByType(UserUpdateMsg.class).size()); Assert.assertEquals(2, edgeImitator.findAllMessagesByType(UserCredentialsUpdateMsg.class).size()); - Optional userUpdateMsgOpt = edgeImitator.findMessageByType(UserUpdateMsg.class); - Assert.assertTrue(userUpdateMsgOpt.isPresent()); - UserUpdateMsg userUpdateMsg = userUpdateMsgOpt.get(); + + UserUpdateMsg userUpdateMsg = getLatestUserUpdateMsg(); User userMsg = JacksonUtil.fromString(userUpdateMsg.getEntity(), User.class, true); Assert.assertNotNull(userMsg); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, userUpdateMsg.getMsgType()); - Assert.assertEquals(savedTenantAdmin.getId(), userMsg.getId()); - Assert.assertEquals(savedTenantAdmin.getAuthority(), userMsg.getAuthority()); - Assert.assertEquals(savedTenantAdmin.getEmail(), userMsg.getEmail()); - Assert.assertEquals(savedTenantAdmin.getFirstName(), userMsg.getFirstName()); - Assert.assertEquals(savedTenantAdmin.getLastName(), userMsg.getLastName()); - Optional userCredentialsUpdateMsgOpt = edgeImitator.findMessageByType(UserCredentialsUpdateMsg.class); - Assert.assertTrue(userCredentialsUpdateMsgOpt.isPresent()); // update user edgeImitator.expectMessageAmount(2); - savedTenantAdmin.setLastName("Borisov"); + savedTenantAdmin.setLastName(UPDATED_LAST_NAME); savedTenantAdmin = doPost("/api/user", savedTenantAdmin, User.class); Assert.assertTrue(edgeImitator.waitForMessages()); - userUpdateMsgOpt = edgeImitator.findMessageByType(UserUpdateMsg.class); - Assert.assertTrue(userUpdateMsgOpt.isPresent()); - userUpdateMsg = userUpdateMsgOpt.get(); - userMsg = JacksonUtil.fromString(userUpdateMsg.getEntity(), User.class, true); - Assert.assertNotNull(userMsg); + + userUpdateMsg = getLatestUserUpdateMsg(); + User userFromMsg = JacksonUtil.fromString(userUpdateMsg.getEntity(), User.class, true); + Assert.assertNotNull(userFromMsg); Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, userUpdateMsg.getMsgType()); - Assert.assertEquals(savedTenantAdmin.getLastName(), userMsg.getLastName()); + Assert.assertEquals(UPDATED_LAST_NAME, userFromMsg.getLastName()); // update user credentials login(savedTenantAdmin.getEmail(), "tenant"); @@ -94,6 +91,7 @@ public class UserEdgeTest extends AbstractEdgeTest { changePasswordRequest.setNewPassword("newTenant"); doPost("/api/auth/changePassword", changePasswordRequest); Assert.assertTrue(edgeImitator.waitForMessages()); + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); Assert.assertTrue(latestMessage instanceof UserCredentialsUpdateMsg); UserCredentialsUpdateMsg userCredentialsUpdateMsg = (UserCredentialsUpdateMsg) latestMessage; @@ -109,6 +107,7 @@ public class UserEdgeTest extends AbstractEdgeTest { doDelete("/api/user/" + savedTenantAdmin.getUuidId()) .andExpect(status().isOk()); Assert.assertTrue(edgeImitator.waitForMessages()); + latestMessage = edgeImitator.getLatestMessage(); Assert.assertTrue(latestMessage instanceof UserUpdateMsg); userUpdateMsg = (UserUpdateMsg) latestMessage; @@ -120,27 +119,11 @@ public class UserEdgeTest extends AbstractEdgeTest { @Test public void testCreateUpdateDeleteCustomerUser() throws Exception { // create customer - edgeImitator.expectMessageAmount(1); - Customer customer = new Customer(); - customer.setTitle("Edge Customer"); - Customer savedCustomer = doPost("/api/customer", customer, Customer.class); - Assert.assertFalse(edgeImitator.waitForMessages(5)); - - // assign edge to customer - edgeImitator.expectMessageAmount(2); - doPost("/api/customer/" + savedCustomer.getUuidId() - + "/edge/" + edge.getUuidId(), Edge.class); - Assert.assertTrue(edgeImitator.waitForMessages()); + Customer savedCustomer = createAndAssignCustomerToEdge("Edge Customer"); // create user edgeImitator.expectMessageAmount(3); - User customerUser = new User(); - customerUser.setAuthority(Authority.CUSTOMER_USER); - customerUser.setTenantId(tenantId); - customerUser.setCustomerId(savedCustomer.getId()); - customerUser.setEmail("customerUser@thingsboard.org"); - customerUser.setFirstName("John"); - customerUser.setLastName("Edwards"); + User customerUser = buildUser(Authority.CUSTOMER_USER, savedCustomer.getId(), "customerUser@thingsboard.org", "John", "Edwards"); User savedCustomerUser = createUser(customerUser, "customer"); Assert.assertTrue(edgeImitator.waitForMessages()); // wait 3 messages - x1 user update msg and x2 user credentials update msgs (create + authenticate user) Assert.assertEquals(1, edgeImitator.findAllMessagesByType(UserUpdateMsg.class).size()); @@ -205,32 +188,57 @@ public class UserEdgeTest extends AbstractEdgeTest { } @Test - public void testSendUserCredentialsRequestToCloud() throws Exception { - UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); - UserCredentialsRequestMsg.Builder userCredentialsRequestMsgBuilder = UserCredentialsRequestMsg.newBuilder(); - userCredentialsRequestMsgBuilder.setUserIdMSB(tenantAdminUserId.getId().getMostSignificantBits()); - userCredentialsRequestMsgBuilder.setUserIdLSB(tenantAdminUserId.getId().getLeastSignificantBits()); - testAutoGeneratedCodeByProtobuf(userCredentialsRequestMsgBuilder); - uplinkMsgBuilder.addUserCredentialsRequestMsg(userCredentialsRequestMsgBuilder.build()); + public void testSendUserToCloudFromEdge() throws Exception { + // create customer + Customer savedCustomer = createAndAssignCustomerToEdge("Edge Customer"); - testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + // create user + User customerUser = buildUser(Authority.CUSTOMER_USER, savedCustomer.getId(), "customerUser@thingsboard.org", DEFAULT_FIRST_NAME, "Johnson"); + + UUID uuid = UUID.randomUUID(); + customerUser.setId(new UserId(uuid)); + UUID userCredentialsUuid = UUID.randomUUID(); + UplinkMsg uplinkMsg = constructUserUplinkMsg(customerUser, UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, userCredentialsUuid); edgeImitator.expectResponsesAmount(1); - edgeImitator.expectMessageAmount(1); - edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + edgeImitator.sendUplinkMsg(uplinkMsg); Assert.assertTrue(edgeImitator.waitForResponses()); - Assert.assertTrue(edgeImitator.waitForMessages()); - AbstractMessage latestMessage = edgeImitator.getLatestMessage(); - Assert.assertTrue(latestMessage instanceof UserCredentialsUpdateMsg); - UserCredentialsUpdateMsg userCredentialsUpdateMsg = (UserCredentialsUpdateMsg) latestMessage; - UserCredentials userCredentialsMsg = JacksonUtil.fromString(userCredentialsUpdateMsg.getEntity(), UserCredentials.class, true); - Assert.assertNotNull(userCredentialsMsg); - Assert.assertEquals(tenantAdminUserId, userCredentialsMsg.getUserId()); + User userFromCloud = doGet("/api/user/" + uuid, User.class); + Assert.assertNotNull(userFromCloud); + Assert.assertEquals(customerUser.getEmail(), userFromCloud.getEmail()); + //check user with existing email + User userWithExistingEmail = buildUser(Authority.CUSTOMER_USER, savedCustomer.getId(), "customerUser@thingsboard.org", DEFAULT_FIRST_NAME, "Johnson"); + + UUID uuidForExistingEmail = UUID.randomUUID(); + userWithExistingEmail.setId(new UserId(uuidForExistingEmail)); + UUID userCredentialsUuidForExistingEmail = UUID.randomUUID(); + UplinkMsg uplinkMsgForExistingEmail = constructUserUplinkMsg(userWithExistingEmail, UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, userCredentialsUuidForExistingEmail); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgForExistingEmail); + Assert.assertTrue(edgeImitator.waitForResponses()); + + User userFromCloudWithExistingEmail = doGet("/api/user/" + uuidForExistingEmail, User.class); + Assert.assertNotNull(userFromCloudWithExistingEmail); + Assert.assertNotEquals(userWithExistingEmail.getEmail(), userFromCloudWithExistingEmail.getEmail()); + + assertUserCredentialsFlags(userFromCloud, false, false); + + UplinkMsg enabledCredentialsUplinkMsg = constructUserCredentialsUplinkMsg(customerUser.getId(), "password", true, userCredentialsUuid); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(enabledCredentialsUplinkMsg); + Assert.assertTrue(edgeImitator.waitForResponses()); + + User cloudUserWithCredentials = doGet("/api/user/" + uuid, User.class); + Assert.assertNotNull(cloudUserWithCredentials); + + assertUserCredentialsFlags(cloudUserWithCredentials, true, true); } @Test - public void sendUserCredentialsRequest() throws Exception { + public void testSendUserCredentialsRequest() throws Exception { UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); UserCredentialsRequestMsg.Builder userCredentialsRequestMsgBuilder = UserCredentialsRequestMsg.newBuilder(); userCredentialsRequestMsgBuilder.setUserIdMSB(tenantAdminUserId.getId().getMostSignificantBits()); @@ -256,4 +264,71 @@ public class UserEdgeTest extends AbstractEdgeTest { testAutoGeneratedCodeByProtobuf(userCredentialsUpdateMsg); } + private User buildUser(Authority authority, CustomerId customerId, String email, String firstName, String lastName) { + User customerUser = new User(); + customerUser.setAuthority(authority); + customerUser.setTenantId(tenantId); + customerUser.setCustomerId(customerId); + customerUser.setEmail(email); + customerUser.setFirstName(firstName); + customerUser.setLastName(lastName); + return customerUser; + } + + private Customer createAndAssignCustomerToEdge(String title) throws Exception { + edgeImitator.expectMessageAmount(1); + Customer customer = new Customer(); + customer.setTitle(title); + Customer savedCustomer = doPost("/api/customer", customer, Customer.class); + Assert.assertFalse(edgeImitator.waitForMessages(5)); + + edgeImitator.expectMessageAmount(2); + doPost("/api/customer/" + savedCustomer.getUuidId() + "/edge/" + edge.getUuidId(), Edge.class); + Assert.assertTrue(edgeImitator.waitForMessages()); + + return savedCustomer; + } + + private UplinkMsg constructUserUplinkMsg(User user, UpdateMsgType msgType, UUID userCredentialsUuid) { + UserUpdateMsg userUpdateMsg = EdgeMsgConstructorUtils.constructUserUpdatedMsg(msgType, user); + + UserCredentials userCredentials = new UserCredentials(); + userCredentials.setId(new UserCredentialsId(userCredentialsUuid)); + userCredentials.setUserId(user.getId()); + userCredentials.setEnabled(false); + userCredentials.setAdditionalInfo(JacksonUtil.newObjectNode()); + userCredentials.setActivateToken(StringUtils.randomAlphanumeric(DEFAULT_TOKEN_LENGTH)); + UserCredentialsUpdateMsg userCredentialsMsg = EdgeMsgConstructorUtils.constructUserCredentialsUpdatedMsg(userCredentials); + + return UplinkMsg.newBuilder() + .addUserUpdateMsg(userUpdateMsg) + .addUserCredentialsUpdateMsg(userCredentialsMsg) + .build(); + } + + private UplinkMsg constructUserCredentialsUplinkMsg(UserId userId, String password, boolean enabled, UUID userCredentialsUuid) { + UserCredentials userCredentials = new UserCredentials(); + userCredentials.setId(new UserCredentialsId(userCredentialsUuid)); + userCredentials.setUserId(userId); + userCredentials.setEnabled(enabled); + userCredentials.setPassword(password); + UserCredentialsUpdateMsg credsMsg = EdgeMsgConstructorUtils.constructUserCredentialsUpdatedMsg(userCredentials); + return UplinkMsg.newBuilder() + .addUserCredentialsUpdateMsg(credsMsg) + .build(); + } + + private void assertUserCredentialsFlags(User user, boolean enabled, boolean activated) { + JsonNode info = user.getAdditionalInfo(); + Assert.assertNotNull(info); + Assert.assertEquals(enabled, info.get("userCredentialsEnabled").asBoolean()); + Assert.assertEquals(activated, info.get("userActivated").asBoolean()); + } + + private UserUpdateMsg getLatestUserUpdateMsg() { + Optional opt = edgeImitator.findMessageByType(UserUpdateMsg.class); + Assert.assertTrue(opt.isPresent()); + return opt.get(); + } + } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/user/UserService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/user/UserService.java index c016631064..1c8943ad7a 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/user/UserService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/user/UserService.java @@ -45,6 +45,8 @@ public interface UserService extends EntityDaoService { User saveUser(TenantId tenantId, User user); + User saveUser(TenantId tenantId, User user, boolean doValidate); + UserCredentials findUserCredentialsByUserId(TenantId tenantId, UserId userId); UserCredentials findUserCredentialsByActivateToken(TenantId tenantId, String activateToken); @@ -53,6 +55,8 @@ public interface UserService extends EntityDaoService { UserCredentials saveUserCredentials(TenantId tenantId, UserCredentials userCredentials); + UserCredentials saveUserCredentials(TenantId tenantId, UserCredentials userCredentials, boolean doValidate); + UserCredentials activateUserCredentials(TenantId tenantId, String activateToken, String password); UserCredentials requestPasswordReset(TenantId tenantId, String email); @@ -67,6 +71,9 @@ public interface UserService extends EntityDaoService { UserCredentials replaceUserCredentials(TenantId tenantId, UserCredentials userCredentials); + UserCredentials replaceUserCredentials(TenantId tenantId, UserCredentials userCredentials, + UserCredentialsId oldUserCredentialsId, boolean doValidate); + void deleteUser(TenantId tenantId, User user); PageData findUsersByTenantId(TenantId tenantId, PageLink pageLink); diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index dbda462a99..7eddb326ef 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -441,6 +441,8 @@ message UplinkMsg { repeated RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = 24; repeated CalculatedFieldUpdateMsg calculatedFieldUpdateMsg = 25; repeated CalculatedFieldRequestMsg calculatedFieldRequestMsg = 26; + repeated UserUpdateMsg userUpdateMsg = 27; + repeated UserCredentialsUpdateMsg userCredentialsUpdateMsg = 28; } message UplinkResponseMsg { diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/UserDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/UserDataValidator.java index 777ce4dfc8..64d0fc21b8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/UserDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/UserDataValidator.java @@ -136,4 +136,5 @@ public class UserDataValidator extends DataValidator { } } } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java index 5c94ba1891..8c229c815d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java @@ -70,6 +70,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Optional; +import java.util.UUID; import java.util.concurrent.TimeUnit; import static org.thingsboard.server.common.data.StringUtils.generateSafeToken; @@ -84,7 +85,7 @@ public class UserServiceImpl extends AbstractCachedEntityService tenantId); + if (doValidate) { + userCredentialsValidator.validate(userCredentials, data -> tenantId); + } UserCredentials result = userCredentialsDao.save(tenantId, userCredentials); eventPublisher.publishEvent(ActionEntityEvent.builder() .tenantId(tenantId) @@ -304,19 +324,44 @@ public class UserServiceImpl extends AbstractCachedEntityService tenantId); - userCredentialsDao.removeById(tenantId, userCredentials.getUuidId()); - userCredentials.setId(null); - if (userCredentials.getPassword() != null) { - updatePasswordHistory(userCredentials); + return replaceUserCredentialsInternal(tenantId, userCredentials, userCredentials.getUuidId(), true); + } + + @Override + public UserCredentials replaceUserCredentials(TenantId tenantId, UserCredentials userCredentials, + UserCredentialsId oldUserCredentialsId, boolean doValidate) { + return replaceUserCredentialsInternal(tenantId, userCredentials, oldUserCredentialsId.getId(), doValidate); + } + + private UserCredentials replaceUserCredentialsInternal(TenantId tenantId, UserCredentials userCredentials, + UUID oldCredentialsUuid, boolean doValidate) { + log.trace("[{}] Replacing user credentials for user [{}], old credentials ID [{}]", + tenantId, userCredentials.getUserId(), oldCredentialsUuid); + + if (doValidate) { + userCredentialsValidator.validate(userCredentials, data -> tenantId); + } + + try { + userCredentialsDao.removeById(tenantId, oldCredentialsUuid); + + if (userCredentials.getPassword() != null) { + updatePasswordHistory(userCredentials); + } + + UserCredentials savedCredentials = userCredentialsDao.save(tenantId, userCredentials); + + eventPublisher.publishEvent(ActionEntityEvent.builder() + .tenantId(tenantId) + .entityId(userCredentials.getUserId()) + .actionType(ActionType.CREDENTIALS_UPDATED) + .build()); + + return savedCredentials; + } catch (Exception e) { + log.error("[{}] Failed to replace user credentials for user [{}]", tenantId, userCredentials.getUserId(), e); + throw new RuntimeException("Failed to replace user credentials", e); } - UserCredentials result = userCredentialsDao.save(tenantId, userCredentials); - eventPublisher.publishEvent(ActionEntityEvent.builder() - .tenantId(tenantId) - .entityId(userCredentials.getUserId()) - .actionType(ActionType.CREDENTIALS_UPDATED).build()); - return result; } @Override