Browse Source

added lock on entity creation to fix race condition on multiple entity creation

pull/14206/head
dashevchenko 11 months ago
parent
commit
c01302fe6b
  1. 4
      dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java
  2. 2
      dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java
  3. 4
      dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java
  4. 4
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java
  5. 4
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java
  6. 23
      dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java
  7. 4
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java
  8. 4
      dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java
  9. 6
      dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java
  10. 57
      dao/src/test/java/org/thingsboard/server/dao/service/AssetServiceTest.java
  11. 31
      dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java

4
dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java

@ -148,6 +148,10 @@ public class BaseAssetService extends AbstractCachedEntityService<AssetCacheKey,
@Override @Override
public Asset saveAsset(Asset asset, boolean doValidate) { public Asset saveAsset(Asset asset, boolean doValidate) {
return saveLimitedEntity(asset, () -> doSaveAsset(asset, doValidate));
}
private Asset doSaveAsset(Asset asset, boolean doValidate) {
log.trace("Executing saveAsset [{}]", asset); log.trace("Executing saveAsset [{}]", asset);
Asset oldAsset = null; Asset oldAsset = null;
if (doValidate) { if (doValidate) {

2
dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java

@ -139,7 +139,7 @@ public class CustomerServiceImpl extends AbstractCachedEntityService<CustomerCac
@Override @Override
@Transactional @Transactional
public Customer saveCustomer(Customer customer) { public Customer saveCustomer(Customer customer) {
return saveCustomer(customer, true); return saveLimitedEntity(customer, () -> saveCustomer(customer, true));
} }
private Customer saveCustomer(Customer customer, boolean doValidate) { private Customer saveCustomer(Customer customer, boolean doValidate) {

4
dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java

@ -157,6 +157,10 @@ public class DashboardServiceImpl extends AbstractEntityService implements Dashb
@Override @Override
public Dashboard saveDashboard(Dashboard dashboard, boolean doValidate) { public Dashboard saveDashboard(Dashboard dashboard, boolean doValidate) {
return saveLimitedEntity(dashboard, () -> doSaveDashboard(dashboard, doValidate));
}
private Dashboard doSaveDashboard(Dashboard dashboard, boolean doValidate) {
log.trace("Executing saveDashboard [{}]", dashboard); log.trace("Executing saveDashboard [{}]", dashboard);
if (doValidate) { if (doValidate) {
dashboardValidator.validate(dashboard, DashboardInfo::getTenantId); dashboardValidator.validate(dashboard, DashboardInfo::getTenantId);

4
dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java

@ -210,6 +210,10 @@ public class DeviceServiceImpl extends CachedVersionedEntityService<DeviceCacheK
} }
private Device saveDeviceWithoutCredentials(Device device, boolean doValidate) { private Device saveDeviceWithoutCredentials(Device device, boolean doValidate) {
return saveLimitedEntity(device, () -> doSaveDeviceWithoutCredentials(device, doValidate));
}
private Device doSaveDeviceWithoutCredentials(Device device, boolean doValidate) {
log.trace("Executing saveDevice [{}]", device); log.trace("Executing saveDevice [{}]", device);
Device oldDevice = null; Device oldDevice = null;
if (doValidate) { if (doValidate) {

4
dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java

@ -201,6 +201,10 @@ public class EdgeServiceImpl extends AbstractCachedEntityService<EdgeCacheKey, E
@Override @Override
public Edge saveEdge(Edge edge) { public Edge saveEdge(Edge edge) {
return saveLimitedEntity(edge, () -> doSaveEdge(edge));
}
private Edge doSaveEdge(Edge edge) {
log.trace("Executing saveEdge [{}]", edge); log.trace("Executing saveEdge [{}]", edge);
Edge oldEdge = edgeValidator.validate(edge, Edge::getTenantId); Edge oldEdge = edgeValidator.validate(edge, Edge::getTenantId);
EdgeCacheEvictEvent evictEvent = new EdgeCacheEvictEvent(edge.getTenantId(), edge.getName(), oldEdge != null ? oldEdge.getName() : null); EdgeCacheEvictEvent evictEvent = new EdgeCacheEvictEvent(edge.getTenantId(), edge.getName(), oldEdge != null ? oldEdge.getName() : null);

23
dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java

@ -21,13 +21,16 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisher;
import org.springframework.context.annotation.Lazy; import org.springframework.context.annotation.Lazy;
import org.springframework.util.ConcurrentReferenceHashMap;
import org.thingsboard.common.util.DebugModeUtil; import org.thingsboard.common.util.DebugModeUtil;
import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.HasDebugSettings; import org.thingsboard.server.common.data.HasDebugSettings;
import org.thingsboard.server.common.data.HasTenantId;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.debug.DebugSettings; import org.thingsboard.server.common.data.debug.DebugSettings;
import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.HasId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.relation.RelationTypeGroup;
@ -44,7 +47,10 @@ import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Optional; import java.util.Optional;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Supplier;
@Slf4j @Slf4j
public abstract class AbstractEntityService { public abstract class AbstractEntityService {
@ -52,6 +58,8 @@ public abstract class AbstractEntityService {
public static final String INCORRECT_EDGE_ID = "Incorrect edgeId "; public static final String INCORRECT_EDGE_ID = "Incorrect edgeId ";
public static final String INCORRECT_PAGE_LINK = "Incorrect page link "; public static final String INCORRECT_PAGE_LINK = "Incorrect page link ";
private final ConcurrentMap<TenantId, ReentrantLock> entityCreationLocks = new ConcurrentReferenceHashMap<>(16);
@Autowired @Autowired
protected ApplicationEventPublisher eventPublisher; protected ApplicationEventPublisher eventPublisher;
@ -86,6 +94,21 @@ public abstract class AbstractEntityService {
@Value("${debug.settings.default_duration:15}") @Value("${debug.settings.default_duration:15}")
private int defaultDebugDurationMinutes; private int defaultDebugDurationMinutes;
protected <E extends HasId & HasTenantId> E saveLimitedEntity(E entity, Supplier<E> saveFunction) {
log.debug("Creating limited entity: {}", entity);
if (entity.getId() == null) {
ReentrantLock lock = entityCreationLocks.computeIfAbsent(entity.getTenantId(), id -> new ReentrantLock());
lock.lock();
try {
return saveFunction.get();
} finally {
lock.unlock();
}
} else {
return saveFunction.get();
}
}
protected void createRelation(TenantId tenantId, EntityRelation relation) { protected void createRelation(TenantId tenantId, EntityRelation relation) {
log.debug("Creating relation: {}", relation); log.debug("Creating relation: {}", relation);
relationService.saveRelation(tenantId, relation); relationService.saveRelation(tenantId, relation);

4
dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java

@ -125,6 +125,10 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
@Override @Override
@Transactional @Transactional
public RuleChain saveRuleChain(RuleChain ruleChain, boolean publishSaveEvent, boolean doValidate) { public RuleChain saveRuleChain(RuleChain ruleChain, boolean publishSaveEvent, boolean doValidate) {
return saveLimitedEntity(ruleChain, () -> doSaveRuleChain(ruleChain, publishSaveEvent, true));
}
private RuleChain doSaveRuleChain(RuleChain ruleChain, boolean publishSaveEvent, boolean doValidate) {
log.trace("Executing doSaveRuleChain [{}]", ruleChain); log.trace("Executing doSaveRuleChain [{}]", ruleChain);
if (doValidate) { if (doValidate) {
ruleChainValidator.validate(ruleChain, RuleChain::getTenantId); ruleChainValidator.validate(ruleChain, RuleChain::getTenantId);

4
dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java

@ -159,6 +159,10 @@ public class UserServiceImpl extends AbstractCachedEntityService<UserCacheKey, U
@Override @Override
@Transactional @Transactional
public User saveUser(TenantId tenantId, User user) { public User saveUser(TenantId tenantId, User user) {
return saveLimitedEntity(user, () -> doSaveUser(tenantId, user));
}
private User doSaveUser(TenantId tenantId, User user) {
log.trace("Executing saveUser [{}]", user); log.trace("Executing saveUser [{}]", user);
User oldUser = userValidator.validate(user, User::getTenantId); User oldUser = userValidator.validate(user, User::getTenantId);
if (!userLoginCaseSensitive) { if (!userLoginCaseSensitive) {

6
dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java

@ -48,6 +48,7 @@ import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.HasId; import org.thingsboard.server.common.data.id.HasId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.oauth2.MapperType; import org.thingsboard.server.common.data.oauth2.MapperType;
import org.thingsboard.server.common.data.oauth2.OAuth2Client; import org.thingsboard.server.common.data.oauth2.OAuth2Client;
import org.thingsboard.server.common.data.oauth2.OAuth2CustomMapperConfig; import org.thingsboard.server.common.data.oauth2.OAuth2CustomMapperConfig;
@ -185,8 +186,13 @@ public abstract class AbstractServiceTest {
} }
public Tenant createTenant() { public Tenant createTenant() {
return createTenant(null);
}
public Tenant createTenant(TenantProfileId tenantProfileId) {
Tenant tenant = new Tenant(); Tenant tenant = new Tenant();
tenant.setTitle("My tenant " + UUID.randomUUID()); tenant.setTitle("My tenant " + UUID.randomUUID());
tenant.setTenantProfileId(tenantProfileId);
Tenant savedTenant = tenantService.saveTenant(tenant); Tenant savedTenant = tenantService.saveTenant(tenant);
assertNotNull(savedTenant); assertNotNull(savedTenant);
return savedTenant; return savedTenant;

57
dao/src/test/java/org/thingsboard/server/dao/service/AssetServiceTest.java

@ -16,17 +16,25 @@
package org.thingsboard.server.dao.service; package org.thingsboard.server.dao.service;
import com.datastax.oss.driver.api.core.uuid.Uuids; import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.google.common.util.concurrent.ListeningExecutorService;
import com.google.common.util.concurrent.MoreExecutors;
import org.junit.After;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Before;
import org.junit.Test; import org.junit.Test;
import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Assertions;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.transaction.PlatformTransactionManager; import org.springframework.transaction.PlatformTransactionManager;
import org.springframework.transaction.TransactionStatus; import org.springframework.transaction.TransactionStatus;
import org.springframework.transaction.support.DefaultTransactionDefinition; import org.springframework.transaction.support.DefaultTransactionDefinition;
import org.testcontainers.shaded.org.apache.commons.lang3.RandomStringUtils;
import org.testcontainers.shaded.org.awaitility.Awaitility;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntitySubtype;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.asset.AssetInfo; import org.thingsboard.server.common.data.asset.AssetInfo;
import org.thingsboard.server.common.data.asset.AssetProfile; import org.thingsboard.server.common.data.asset.AssetProfile;
@ -44,6 +52,8 @@ import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.data.tenant.profile.TenantProfileData;
import org.thingsboard.server.dao.asset.AssetDao; import org.thingsboard.server.dao.asset.AssetDao;
import org.thingsboard.server.dao.asset.AssetProfileService; import org.thingsboard.server.dao.asset.AssetProfileService;
import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.dao.asset.AssetService;
@ -51,11 +61,14 @@ import org.thingsboard.server.dao.cf.CalculatedFieldService;
import org.thingsboard.server.dao.customer.CustomerService; import org.thingsboard.server.dao.customer.CustomerService;
import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.dao.tenant.TenantProfileService;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collections; import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID; import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID;
@ -72,6 +85,8 @@ public class AssetServiceTest extends AbstractServiceTest {
@Autowired @Autowired
RelationService relationService; RelationService relationService;
@Autowired @Autowired
TenantProfileService tenantProfileService;
@Autowired
private AssetProfileService assetProfileService; private AssetProfileService assetProfileService;
@Autowired @Autowired
private CalculatedFieldService calculatedFieldService; private CalculatedFieldService calculatedFieldService;
@ -79,6 +94,18 @@ public class AssetServiceTest extends AbstractServiceTest {
private PlatformTransactionManager platformTransactionManager; private PlatformTransactionManager platformTransactionManager;
private IdComparator<Asset> idComparator = new IdComparator<>(); private IdComparator<Asset> idComparator = new IdComparator<>();
ListeningExecutorService executor;
private TenantId anotherTenantId;
@Before
public void before() {
executor = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(10, ThingsBoardThreadFactory.forName(getClass().getSimpleName() + "-test-scope")));
}
@After
public void after() {
executor.shutdownNow();
}
@Test @Test
public void testSaveAsset() { public void testSaveAsset() {
@ -105,6 +132,36 @@ public class AssetServiceTest extends AbstractServiceTest {
assetService.deleteAsset(tenantId, savedAsset.getId()); assetService.deleteAsset(tenantId, savedAsset.getId());
} }
@Test
public void testAssetLimitOnTenantProfileLevel() {
TenantProfile tenantProfile = new TenantProfile();
tenantProfile.setName("Test profile");
tenantProfile.setDescription("Test");
TenantProfileData profileData = new TenantProfileData();
profileData.setConfiguration(DefaultTenantProfileConfiguration.builder().maxAssets(5l).build());
tenantProfile.setProfileData(profileData);
tenantProfile.setDefault(false);
tenantProfile.setIsolatedTbRuleEngine(false);
tenantProfile = tenantProfileService.saveTenantProfile(anotherTenantId, tenantProfile);
anotherTenantId = createTenant(tenantProfile.getId()).getId();
for (int i = 0; i < 20; i++) {
executor.submit(() -> {
Asset asset = new Asset();
asset.setTenantId(anotherTenantId);
asset.setName(RandomStringUtils.randomAlphabetic(10));
asset.setType("default");
assetService.saveAsset(asset);
});
}
Awaitility.await().atMost(10, TimeUnit.SECONDS).until(() -> {
long countByTenantId = assetService.countByTenantId(anotherTenantId);
return countByTenantId == 5;
});
}
@Test @Test
public void testShouldNotPutInCacheRolledbackAssetProfile() { public void testShouldNotPutInCacheRolledbackAssetProfile() {
AssetProfile assetProfile = new AssetProfile(); AssetProfile assetProfile = new AssetProfile();

31
dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java

@ -16,6 +16,8 @@
package org.thingsboard.server.dao.service; package org.thingsboard.server.dao.service;
import com.datastax.oss.driver.api.core.uuid.Uuids; import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.google.common.util.concurrent.ListeningExecutorService;
import com.google.common.util.concurrent.MoreExecutors;
import org.junit.After; import org.junit.After;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Before; import org.junit.Before;
@ -27,6 +29,8 @@ import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.transaction.PlatformTransactionManager; import org.springframework.transaction.PlatformTransactionManager;
import org.springframework.transaction.TransactionStatus; import org.springframework.transaction.TransactionStatus;
import org.springframework.transaction.support.DefaultTransactionDefinition; import org.springframework.transaction.support.DefaultTransactionDefinition;
import org.testcontainers.shaded.org.awaitility.Awaitility;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceInfo; import org.thingsboard.server.common.data.DeviceInfo;
@ -74,6 +78,8 @@ import java.util.ArrayList;
import java.util.Collections; import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.any;
@ -105,10 +111,12 @@ public class DeviceServiceTest extends AbstractServiceTest {
private IdComparator<Device> idComparator = new IdComparator<>(); private IdComparator<Device> idComparator = new IdComparator<>();
private TenantId anotherTenantId; private TenantId anotherTenantId;
private ListeningExecutorService executor;
@Before @Before
public void before() { public void before() {
anotherTenantId = createTenant().getId(); anotherTenantId = createTenant().getId();
executor = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(10, ThingsBoardThreadFactory.forName(getClass().getSimpleName() + "-test-scope")));
} }
@After @After
@ -118,6 +126,7 @@ public class DeviceServiceTest extends AbstractServiceTest {
tenantProfileService.deleteTenantProfiles(tenantId); tenantProfileService.deleteTenantProfiles(tenantId);
tenantProfileService.deleteTenantProfiles(anotherTenantId); tenantProfileService.deleteTenantProfiles(anotherTenantId);
executor.shutdownNow();
} }
@Test @Test
@ -136,6 +145,28 @@ public class DeviceServiceTest extends AbstractServiceTest {
deleteDevice(tenantId, device); deleteDevice(tenantId, device);
} }
@Test
public void testDeviceLimitOnTenantProfileLevel() {
TenantProfile defaultTenantProfile = tenantProfileService.findDefaultTenantProfile(tenantId);
defaultTenantProfile.getProfileData().setConfiguration(DefaultTenantProfileConfiguration.builder().maxDevices(5l).build());
tenantProfileService.saveTenantProfile(tenantId, defaultTenantProfile);
for (int i = 0; i < 20; i++) {
executor.submit(() -> {
Device device = new Device();
device.setTenantId(tenantId);
device.setName(StringUtils.randomAlphabetic(10));
device.setType("default");
deviceService.saveDevice(device, true);
});
}
Awaitility.await().atMost(10, TimeUnit.SECONDS).until(() -> {
long countByTenantId = deviceService.countByTenantId(tenantId);
return countByTenantId == 5;
});
}
@Test @Test
public void testSaveDevicesWithMaxDeviceOutOfLimit() { public void testSaveDevicesWithMaxDeviceOutOfLimit() {
TenantProfile defaultTenantProfile = tenantProfileService.findDefaultTenantProfile(tenantId); TenantProfile defaultTenantProfile = tenantProfileService.findDefaultTenantProfile(tenantId);

Loading…
Cancel
Save