Browse Source

Added User sync for Edge

pull/14040/head
Yevhenii 10 months ago
parent
commit
fcc2a917fd
  1. 4
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  2. 24
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  3. 130
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/BaseUserProcessor.java
  4. 97
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserEdgeProcessor.java
  5. 31
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/user/UserProcessor.java
  6. 193
      application/src/test/java/org/thingsboard/server/edge/UserEdgeTest.java
  7. 7
      common/dao-api/src/main/java/org/thingsboard/server/dao/user/UserService.java
  8. 2
      common/edge-api/src/main/proto/edge.proto
  9. 1
      dao/src/main/java/org/thingsboard/server/dao/service/validator/UserDataValidator.java
  10. 77
      dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java

4
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) {

24
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);

130
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<User> userValidator;
protected Pair<Boolean, Boolean> 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);
}

97
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<Void> 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<Void> 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<Boolean, Boolean> 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);
}
}

31
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<Void> processUserMsgFromEdge(TenantId tenantId, Edge edge, UserUpdateMsg userUpdateMsg);
ListenableFuture<Void> processUserCredentialsMsgFromEdge(TenantId tenantId, Edge edge, UserCredentialsUpdateMsg userCredentialsUpdateMsg);
}

193
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<UserUpdateMsg> 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<UserCredentialsUpdateMsg> 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<UserUpdateMsg> opt = edgeImitator.findMessageByType(UserUpdateMsg.class);
Assert.assertTrue(opt.isPresent());
return opt.get();
}
}

7
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<User> findUsersByTenantId(TenantId tenantId, PageLink pageLink);

2
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 {

1
dao/src/main/java/org/thingsboard/server/dao/service/validator/UserDataValidator.java

@ -136,4 +136,5 @@ public class UserDataValidator extends DataValidator<User> {
}
}
}
}

77
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<UserCacheKey, U
public static final String USER_PASSWORD_HISTORY = "userPasswordHistory";
private static final int DEFAULT_TOKEN_LENGTH = 30;
public static final int DEFAULT_TOKEN_LENGTH = 30;
public static final String INCORRECT_USER_ID = "Incorrect userId ";
public static final String INCORRECT_TENANT_ID = "Incorrect tenantId ";
@ -159,8 +160,20 @@ public class UserServiceImpl extends AbstractCachedEntityService<UserCacheKey, U
@Override
@Transactional
public User saveUser(TenantId tenantId, User user) {
return saveUser(tenantId, user, true);
}
@Override
@Transactional
public User saveUser(TenantId tenantId, User user, boolean doValidate) {
boolean isCreatedOnCloud = doValidate;
log.trace("Executing saveUser [{}]", user);
User oldUser = userValidator.validate(user, User::getTenantId);
User oldUser = null;
if (doValidate) {
oldUser = userValidator.validate(user, User::getTenantId);
} else if (user.getId() != null) {
oldUser = findUserById(user.getTenantId(), user.getId());
}
if (!userLoginCaseSensitive) {
user.setEmail(user.getEmail().toLowerCase());
}
@ -169,7 +182,7 @@ public class UserServiceImpl extends AbstractCachedEntityService<UserCacheKey, U
try {
savedUser = userDao.saveAndFlush(user.getTenantId(), user);
publishEvictEvent(evictEvent);
if (user.getId() == null) {
if (user.getId() == null && isCreatedOnCloud) {
countService.publishCountEntityEvictEvent(savedUser.getTenantId(), EntityType.USER);
UserCredentials userCredentials = new UserCredentials();
userCredentials.setEnabled(false);
@ -215,8 +228,15 @@ public class UserServiceImpl extends AbstractCachedEntityService<UserCacheKey, U
@Override
public UserCredentials saveUserCredentials(TenantId tenantId, UserCredentials userCredentials) {
return saveUserCredentials(tenantId, userCredentials, true);
}
@Override
public UserCredentials saveUserCredentials(TenantId tenantId, UserCredentials userCredentials, boolean doValidate) {
log.trace("Executing saveUserCredentials [{}]", userCredentials);
userCredentialsValidator.validate(userCredentials, data -> 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<UserCacheKey, U
@Override
public UserCredentials replaceUserCredentials(TenantId tenantId, UserCredentials userCredentials) {
log.trace("Executing replaceUserCredentials [{}]", userCredentials);
userCredentialsValidator.validate(userCredentials, data -> 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

Loading…
Cancel
Save