diff --git a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java index 85ec0950e3..d5fcc2706f 100644 --- a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java +++ b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java @@ -244,7 +244,8 @@ public class RuleChainController extends BaseController { } RuleChain ruleChain = checkRuleChain(ruleChainMetaData.getRuleChainId(), Operation.WRITE); - RuleChainMetaData savedRuleChainMetaData = checkNotNull(ruleChainService.saveRuleChainMetaData(tenantId, ruleChainMetaData)); + checkNotNull(ruleChainService.saveRuleChainMetaData(tenantId, ruleChainMetaData) ? true : null); + RuleChainMetaData savedRuleChainMetaData = checkNotNull(ruleChainService.loadRuleChainMetaData(tenantId, ruleChainMetaData.getRuleChainId())); if (RuleChainType.CORE.equals(ruleChain.getType())) { tbClusterService.onEntityStateChange(ruleChain.getTenantId(), ruleChain.getId(), ComponentLifecycleEvent.UPDATED); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java index 086d954e27..5740abb354 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java @@ -43,7 +43,7 @@ public interface RuleChainService { boolean setRootRuleChain(TenantId tenantId, RuleChainId ruleChainId); - RuleChainMetaData saveRuleChainMetaData(TenantId tenantId, RuleChainMetaData ruleChainMetaData); + boolean saveRuleChainMetaData(TenantId tenantId, RuleChainMetaData ruleChainMetaData); RuleChainMetaData loadRuleChainMetaData(TenantId tenantId, RuleChainId ruleChainId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java b/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java index 44b328fd9f..15ddb9a390 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java @@ -179,11 +179,28 @@ public class BaseRelationService implements RelationService { @Override public void deleteEntityRelations(TenantId tenantId, EntityId entityId) { - try { - deleteEntityRelationsAsync(tenantId, entityId).get(); - } catch (InterruptedException | ExecutionException e) { - throw new RuntimeException(e); + log.trace("Executing deleteEntityRelations [{}]", entityId); + validate(entityId); + final Cache cache = cacheManager.getCache(RELATIONS_CACHE); + List inboundRelations = new ArrayList<>(); + for (RelationTypeGroup typeGroup : RelationTypeGroup.values()) { + inboundRelations.addAll(relationDao.findAllByTo(tenantId, entityId, typeGroup)); + } + + List outboundRelations = new ArrayList<>(); + for (RelationTypeGroup typeGroup : RelationTypeGroup.values()) { + outboundRelations.addAll(relationDao.findAllByFrom(tenantId, entityId, typeGroup)); + } + + for (EntityRelation relation : inboundRelations){ + delete(tenantId, cache, relation, true); } + + for (EntityRelation relation : outboundRelations){ + delete(tenantId, cache, relation, false); + } + + relationDao.deleteOutboundRelations(tenantId, entityId); } @Override @@ -193,14 +210,14 @@ public class BaseRelationService implements RelationService { validate(entityId); List>> inboundRelationsList = new ArrayList<>(); for (RelationTypeGroup typeGroup : RelationTypeGroup.values()) { - inboundRelationsList.add(relationDao.findAllByTo(tenantId, entityId, typeGroup)); + inboundRelationsList.add(relationDao.findAllByToAsync(tenantId, entityId, typeGroup)); } ListenableFuture>> inboundRelations = Futures.allAsList(inboundRelationsList); List>> outboundRelationsList = new ArrayList<>(); for (RelationTypeGroup typeGroup : RelationTypeGroup.values()) { - outboundRelationsList.add(relationDao.findAllByFrom(tenantId, entityId, typeGroup)); + outboundRelationsList.add(relationDao.findAllByFromAsync(tenantId, entityId, typeGroup)); } ListenableFuture>> outboundRelations = Futures.allAsList(outboundRelationsList); @@ -242,6 +259,15 @@ public class BaseRelationService implements RelationService { } } + boolean delete(TenantId tenantId, Cache cache, EntityRelation relation, boolean deleteFromDb) { + cacheEviction(relation, cache); + if (deleteFromDb) { + return relationDao.deleteRelation(tenantId, relation); + } else { + return false; + } + } + private void cacheEviction(EntityRelation relation, Cache cache) { List fromToTypeAndTypeGroup = new ArrayList<>(); fromToTypeAndTypeGroup.add(relation.getFrom()); @@ -283,7 +309,7 @@ public class BaseRelationService implements RelationService { validate(from); validateTypeGroup(typeGroup); try { - return relationDao.findAllByFrom(tenantId, from, typeGroup).get(); + return relationDao.findAllByFromAsync(tenantId, from, typeGroup).get(); } catch (InterruptedException | ExecutionException e) { throw new RuntimeException(e); } @@ -306,7 +332,7 @@ public class BaseRelationService implements RelationService { if (fromCache != null) { return Futures.immediateFuture(fromCache); } else { - ListenableFuture> relationsFuture = relationDao.findAllByFrom(tenantId, from, typeGroup); + ListenableFuture> relationsFuture = relationDao.findAllByFromAsync(tenantId, from, typeGroup); Futures.addCallback(relationsFuture, new FutureCallback>() { @Override @@ -327,7 +353,7 @@ public class BaseRelationService implements RelationService { log.trace("Executing findInfoByFrom [{}][{}]", from, typeGroup); validate(from); validateTypeGroup(typeGroup); - ListenableFuture> relations = relationDao.findAllByFrom(tenantId, from, typeGroup); + ListenableFuture> relations = relationDao.findAllByFromAsync(tenantId, from, typeGroup); return Futures.transformAsync(relations, relations1 -> { List> futures = new ArrayList<>(); @@ -365,7 +391,7 @@ public class BaseRelationService implements RelationService { validate(to); validateTypeGroup(typeGroup); try { - return relationDao.findAllByTo(tenantId, to, typeGroup).get(); + return relationDao.findAllByToAsync(tenantId, to, typeGroup).get(); } catch (InterruptedException | ExecutionException e) { throw new RuntimeException(e); } @@ -388,7 +414,7 @@ public class BaseRelationService implements RelationService { if (fromCache != null) { return Futures.immediateFuture(fromCache); } else { - ListenableFuture> relationsFuture = relationDao.findAllByTo(tenantId, to, typeGroup); + ListenableFuture> relationsFuture = relationDao.findAllByToAsync(tenantId, to, typeGroup); Futures.addCallback(relationsFuture, new FutureCallback>() { @Override @@ -409,7 +435,7 @@ public class BaseRelationService implements RelationService { log.trace("Executing findInfoByTo [{}][{}]", to, typeGroup); validate(to); validateTypeGroup(typeGroup); - ListenableFuture> relations = relationDao.findAllByTo(tenantId, to, typeGroup); + ListenableFuture> relations = relationDao.findAllByToAsync(tenantId, to, typeGroup); return Futures.transformAsync(relations, relations1 -> { List> futures = new ArrayList<>(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java b/dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java index a8c4d9bbc2..d60f176832 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java @@ -31,11 +31,15 @@ import java.util.List; */ public interface RelationDao { - ListenableFuture> findAllByFrom(TenantId tenantId, EntityId from, RelationTypeGroup typeGroup); + ListenableFuture> findAllByFromAsync(TenantId tenantId, EntityId from, RelationTypeGroup typeGroup); + + List findAllByFrom(TenantId tenantId, EntityId from, RelationTypeGroup typeGroup); ListenableFuture> findAllByFromAndType(TenantId tenantId, EntityId from, String relationType, RelationTypeGroup typeGroup); - ListenableFuture> findAllByTo(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup); + ListenableFuture> findAllByToAsync(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup); + + List findAllByTo(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup); ListenableFuture> findAllByToAndType(TenantId tenantId, EntityId to, String relationType, RelationTypeGroup typeGroup); diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java index b85a257ed2..baa3ae8272 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java @@ -26,6 +26,7 @@ import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; import org.thingsboard.server.common.data.BaseData; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.Tenant; @@ -96,23 +97,19 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC private TbTenantProfileCache tenantProfileCache; @Override + @Transactional public RuleChain saveRuleChain(RuleChain ruleChain) { ruleChainValidator.validate(ruleChain, RuleChain::getTenantId); RuleChain savedRuleChain = ruleChainDao.save(ruleChain.getTenantId(), ruleChain); if (ruleChain.isRoot() && ruleChain.getId() == null) { - try { - createRelation(ruleChain.getTenantId(), new EntityRelation(savedRuleChain.getTenantId(), savedRuleChain.getId(), - EntityRelation.CONTAINS_TYPE, RelationTypeGroup.RULE_CHAIN)); - } catch (Exception e) { - log.warn("[{}] Failed to create tenant to root rule chain relation. from: [{}], to: [{}]", - savedRuleChain.getTenantId(), savedRuleChain.getTenantId(), savedRuleChain.getId(), e); - throw new RuntimeException(e); - } + createRelation(ruleChain.getTenantId(), new EntityRelation(savedRuleChain.getTenantId(), savedRuleChain.getId(), + EntityRelation.CONTAINS_TYPE, RelationTypeGroup.RULE_CHAIN)); } return savedRuleChain; } @Override + @Transactional public boolean setRootRuleChain(TenantId tenantId, RuleChainId ruleChainId) { RuleChain ruleChain = ruleChainDao.findById(tenantId, ruleChainId.getId()); if (!ruleChain.isRoot()) { @@ -145,11 +142,12 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC } @Override - public RuleChainMetaData saveRuleChainMetaData(TenantId tenantId, RuleChainMetaData ruleChainMetaData) { + @Transactional + public boolean saveRuleChainMetaData(TenantId tenantId, RuleChainMetaData ruleChainMetaData) { Validator.validateId(ruleChainMetaData.getRuleChainId(), "Incorrect rule chain id."); RuleChain ruleChain = findRuleChainById(tenantId, ruleChainMetaData.getRuleChainId()); if (ruleChain == null) { - return null; + return false; } if (CollectionUtils.isNotEmpty(ruleChainMetaData.getConnections())) { @@ -184,14 +182,8 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC for (RuleNode node : toAddOrUpdate) { node.setRuleChainId(ruleChain.getId()); RuleNode savedNode = ruleNodeDao.save(tenantId, node); - try { - createRelation(tenantId, new EntityRelation(ruleChainMetaData.getRuleChainId(), savedNode.getId(), - EntityRelation.CONTAINS_TYPE, RelationTypeGroup.RULE_CHAIN)); - } catch (Exception e) { - log.warn("[{}] Failed to create rule chain to rule node relation. from: [{}], to: [{}]", - ruleChainMetaData.getRuleChainId(), savedNode.getId()); - throw new RuntimeException(e); - } + createRelation(tenantId, new EntityRelation(ruleChainMetaData.getRuleChainId(), savedNode.getId(), + EntityRelation.CONTAINS_TYPE, RelationTypeGroup.RULE_CHAIN)); int index = nodes.indexOf(node); nodes.set(index, savedNode); ruleNodeIndexMap.put(savedNode.getId(), index); @@ -213,12 +205,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC EntityId from = nodes.get(nodeConnection.getFromIndex()).getId(); EntityId to = nodes.get(nodeConnection.getToIndex()).getId(); String type = nodeConnection.getType(); - try { - createRelation(tenantId, new EntityRelation(from, to, type, RelationTypeGroup.RULE_NODE)); - } catch (Exception e) { - log.warn("[{}] Failed to create rule node relation. from: [{}], to: [{}]", from, to); - throw new RuntimeException(e); - } + createRelation(tenantId, new EntityRelation(from, to, type, RelationTypeGroup.RULE_NODE)); } } if (ruleChainMetaData.getRuleChainConnections() != null) { @@ -226,16 +213,11 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC EntityId from = nodes.get(nodeToRuleChainConnection.getFromIndex()).getId(); EntityId to = nodeToRuleChainConnection.getTargetRuleChainId(); String type = nodeToRuleChainConnection.getType(); - try { - createRelation(tenantId, new EntityRelation(from, to, type, RelationTypeGroup.RULE_NODE, nodeToRuleChainConnection.getAdditionalInfo())); - } catch (Exception e) { - log.warn("[{}] Failed to create rule node to rule chain relation. from: [{}], to: [{}]", from, to); - throw new RuntimeException(e); - } + createRelation(tenantId, new EntityRelation(from, to, type, RelationTypeGroup.RULE_NODE, nodeToRuleChainConnection.getAdditionalInfo())); } } - return loadRuleChainMetaData(tenantId, ruleChainMetaData.getRuleChainId()); + return true; } private void validateCircles(List connectionInfos) { @@ -410,6 +392,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC } @Override + @Transactional public void deleteRuleChainById(TenantId tenantId, RuleChainId ruleChainId) { Validator.validateId(ruleChainId, "Incorrect rule chain id for delete request."); RuleChain ruleChain = ruleChainDao.findById(tenantId, ruleChainId.getId()); @@ -716,6 +699,15 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC ruleNodeDao.removeById(tenantId, entityId.getId()); } + private void createRelation(TenantId tenantId, EntityRelation relation) { + log.debug("Creating relation: {}", relation); + relationService.saveRelation(tenantId, relation); + } + + private void deleteRelation(TenantId tenantId, EntityRelation relation) { + log.debug("Deleting relation: {}", relation); + relationService.deleteRelation(tenantId, relation); + } private DataValidator ruleChainValidator = new DataValidator() { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java index 89b63b87cf..5e10cabda2 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java @@ -49,12 +49,17 @@ public class JpaRelationDao extends JpaAbstractDaoListeningExecutorService imple private RelationInsertRepository relationInsertRepository; @Override - public ListenableFuture> findAllByFrom(TenantId tenantId, EntityId from, RelationTypeGroup typeGroup) { - return service.submit(() -> DaoUtil.convertDataList( + public ListenableFuture> findAllByFromAsync(TenantId tenantId, EntityId from, RelationTypeGroup typeGroup) { + return service.submit(() -> findAllByFrom(tenantId, from, typeGroup)); + } + + @Override + public List findAllByFrom(TenantId tenantId, EntityId from, RelationTypeGroup typeGroup) { + return DaoUtil.convertDataList( relationRepository.findAllByFromIdAndFromTypeAndRelationTypeGroup( from.getId(), from.getEntityType().name(), - typeGroup.name()))); + typeGroup.name())); } @Override @@ -68,12 +73,17 @@ public class JpaRelationDao extends JpaAbstractDaoListeningExecutorService imple } @Override - public ListenableFuture> findAllByTo(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup) { - return service.submit(() -> DaoUtil.convertDataList( + public ListenableFuture> findAllByToAsync(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup) { + return service.submit(() -> findAllByTo(tenantId, to, typeGroup)); + } + + @Override + public List findAllByTo(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup) { + return DaoUtil.convertDataList( relationRepository.findAllByToIdAndToTypeAndRelationTypeGroup( to.getId(), to.getEntityType().name(), - typeGroup.name()))); + typeGroup.name())); } @Override diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseRuleChainServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseRuleChainServiceTest.java index e7cd349e73..3f4b103ce2 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseRuleChainServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseRuleChainServiceTest.java @@ -292,7 +292,8 @@ public abstract class BaseRuleChainServiceTest extends AbstractServiceTest { ruleNodes.set(name3Index, ruleNode4); - RuleChainMetaData updatedRuleChainMetaData = ruleChainService.saveRuleChainMetaData(tenantId, savedRuleChainMetaData); + Assert.assertTrue(ruleChainService.saveRuleChainMetaData(tenantId, savedRuleChainMetaData)); + RuleChainMetaData updatedRuleChainMetaData = ruleChainService.loadRuleChainMetaData(tenantId, savedRuleChainMetaData.getRuleChainId()); Assert.assertEquals(3, updatedRuleChainMetaData.getNodes().size()); Assert.assertEquals(3, updatedRuleChainMetaData.getConnections().size()); @@ -404,7 +405,8 @@ public abstract class BaseRuleChainServiceTest extends AbstractServiceTest { ruleChainMetaData.addConnectionInfo(0,2,"fail"); ruleChainMetaData.addConnectionInfo(1,2,"success"); - return ruleChainService.saveRuleChainMetaData(tenantId, ruleChainMetaData); + Assert.assertTrue(ruleChainService.saveRuleChainMetaData(tenantId, ruleChainMetaData)); + return ruleChainService.loadRuleChainMetaData(tenantId, ruleChainMetaData.getRuleChainId()); } private RuleChainMetaData createRuleChainMetadataWithCirclingRelation() throws Exception {