Browse Source

rule chain. implemented save/update/delete @Transactional. Added sync DAO methods to run in transaction. Moved loadRuleChainMetaData outside of saveRuleChainMetaData transaction

pull/4460/head
Sergey Matvienko 5 years ago
committed by Andrew Shvayka
parent
commit
8b653d7065
  1. 3
      application/src/main/java/org/thingsboard/server/controller/RuleChainController.java
  2. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java
  3. 50
      dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java
  4. 8
      dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java
  5. 54
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java
  6. 22
      dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java
  7. 6
      dao/src/test/java/org/thingsboard/server/dao/service/BaseRuleChainServiceTest.java

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

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

50
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<EntityRelation> inboundRelations = new ArrayList<>();
for (RelationTypeGroup typeGroup : RelationTypeGroup.values()) {
inboundRelations.addAll(relationDao.findAllByTo(tenantId, entityId, typeGroup));
}
List<EntityRelation> 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<ListenableFuture<List<EntityRelation>>> inboundRelationsList = new ArrayList<>();
for (RelationTypeGroup typeGroup : RelationTypeGroup.values()) {
inboundRelationsList.add(relationDao.findAllByTo(tenantId, entityId, typeGroup));
inboundRelationsList.add(relationDao.findAllByToAsync(tenantId, entityId, typeGroup));
}
ListenableFuture<List<List<EntityRelation>>> inboundRelations = Futures.allAsList(inboundRelationsList);
List<ListenableFuture<List<EntityRelation>>> outboundRelationsList = new ArrayList<>();
for (RelationTypeGroup typeGroup : RelationTypeGroup.values()) {
outboundRelationsList.add(relationDao.findAllByFrom(tenantId, entityId, typeGroup));
outboundRelationsList.add(relationDao.findAllByFromAsync(tenantId, entityId, typeGroup));
}
ListenableFuture<List<List<EntityRelation>>> 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<Object> 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<List<EntityRelation>> relationsFuture = relationDao.findAllByFrom(tenantId, from, typeGroup);
ListenableFuture<List<EntityRelation>> relationsFuture = relationDao.findAllByFromAsync(tenantId, from, typeGroup);
Futures.addCallback(relationsFuture,
new FutureCallback<List<EntityRelation>>() {
@Override
@ -327,7 +353,7 @@ public class BaseRelationService implements RelationService {
log.trace("Executing findInfoByFrom [{}][{}]", from, typeGroup);
validate(from);
validateTypeGroup(typeGroup);
ListenableFuture<List<EntityRelation>> relations = relationDao.findAllByFrom(tenantId, from, typeGroup);
ListenableFuture<List<EntityRelation>> relations = relationDao.findAllByFromAsync(tenantId, from, typeGroup);
return Futures.transformAsync(relations,
relations1 -> {
List<ListenableFuture<EntityRelationInfo>> 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<List<EntityRelation>> relationsFuture = relationDao.findAllByTo(tenantId, to, typeGroup);
ListenableFuture<List<EntityRelation>> relationsFuture = relationDao.findAllByToAsync(tenantId, to, typeGroup);
Futures.addCallback(relationsFuture,
new FutureCallback<List<EntityRelation>>() {
@Override
@ -409,7 +435,7 @@ public class BaseRelationService implements RelationService {
log.trace("Executing findInfoByTo [{}][{}]", to, typeGroup);
validate(to);
validateTypeGroup(typeGroup);
ListenableFuture<List<EntityRelation>> relations = relationDao.findAllByTo(tenantId, to, typeGroup);
ListenableFuture<List<EntityRelation>> relations = relationDao.findAllByToAsync(tenantId, to, typeGroup);
return Futures.transformAsync(relations,
relations1 -> {
List<ListenableFuture<EntityRelationInfo>> futures = new ArrayList<>();

8
dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java

@ -31,11 +31,15 @@ import java.util.List;
*/
public interface RelationDao {
ListenableFuture<List<EntityRelation>> findAllByFrom(TenantId tenantId, EntityId from, RelationTypeGroup typeGroup);
ListenableFuture<List<EntityRelation>> findAllByFromAsync(TenantId tenantId, EntityId from, RelationTypeGroup typeGroup);
List<EntityRelation> findAllByFrom(TenantId tenantId, EntityId from, RelationTypeGroup typeGroup);
ListenableFuture<List<EntityRelation>> findAllByFromAndType(TenantId tenantId, EntityId from, String relationType, RelationTypeGroup typeGroup);
ListenableFuture<List<EntityRelation>> findAllByTo(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup);
ListenableFuture<List<EntityRelation>> findAllByToAsync(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup);
List<EntityRelation> findAllByTo(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup);
ListenableFuture<List<EntityRelation>> findAllByToAndType(TenantId tenantId, EntityId to, String relationType, RelationTypeGroup typeGroup);

54
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<NodeConnectionInfo> 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<RuleChain> ruleChainValidator =
new DataValidator<RuleChain>() {

22
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<List<EntityRelation>> findAllByFrom(TenantId tenantId, EntityId from, RelationTypeGroup typeGroup) {
return service.submit(() -> DaoUtil.convertDataList(
public ListenableFuture<List<EntityRelation>> findAllByFromAsync(TenantId tenantId, EntityId from, RelationTypeGroup typeGroup) {
return service.submit(() -> findAllByFrom(tenantId, from, typeGroup));
}
@Override
public List<EntityRelation> 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<List<EntityRelation>> findAllByTo(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup) {
return service.submit(() -> DaoUtil.convertDataList(
public ListenableFuture<List<EntityRelation>> findAllByToAsync(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup) {
return service.submit(() -> findAllByTo(tenantId, to, typeGroup));
}
@Override
public List<EntityRelation> findAllByTo(TenantId tenantId, EntityId to, RelationTypeGroup typeGroup) {
return DaoUtil.convertDataList(
relationRepository.findAllByToIdAndToTypeAndRelationTypeGroup(
to.getId(),
to.getEntityType().name(),
typeGroup.name())));
typeGroup.name()));
}
@Override

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

Loading…
Cancel
Save