From 4956c0c573fbeba976f5fe5c6d9c138358439bed Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Tue, 13 Mar 2018 13:12:00 +0200 Subject: [PATCH] Rule Chain service implementation. --- .../data/relation/RelationTypeGroup.java | 4 +- .../common/data/rule/RuleChainMetaData.java | 75 ++++++ .../server/dao/rule/BaseRuleChainService.java | 225 +++++++++++++++++- .../server/dao/rule/RuleChainService.java | 14 ++ .../server/dao/tenant/TenantServiceImpl.java | 5 + 5 files changed, 316 insertions(+), 7 deletions(-) create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChainMetaData.java diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/relation/RelationTypeGroup.java b/common/data/src/main/java/org/thingsboard/server/common/data/relation/RelationTypeGroup.java index 90e0253370..55990552b0 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/relation/RelationTypeGroup.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/relation/RelationTypeGroup.java @@ -19,6 +19,8 @@ public enum RelationTypeGroup { COMMON, ALARM, - DASHBOARD + DASHBOARD, + RULE_CHAIN, + RULE_NODE } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChainMetaData.java b/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChainMetaData.java new file mode 100644 index 0000000000..af141d6142 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChainMetaData.java @@ -0,0 +1,75 @@ +/** + * Copyright © 2016-2018 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.common.data.rule; + +import lombok.Data; +import org.thingsboard.server.common.data.id.RuleChainId; + +import java.util.ArrayList; +import java.util.List; + +/** + * Created by igor on 3/13/18. + */ +@Data +public class RuleChainMetaData { + + private RuleChainId ruleChainId; + + private Integer firstNodeIndex; + + private List nodes; + + private List connections; + + private List ruleChainConnections; + + public void addConnectionInfo(int fromIndex, int toIndex, String type) { + NodeConnectionInfo connectionInfo = new NodeConnectionInfo(); + connectionInfo.setFromIndex(fromIndex); + connectionInfo.setToIndex(toIndex); + connectionInfo.setType(type); + if (connections == null) { + connections = new ArrayList<>(); + } + connections.add(connectionInfo); + } + public void addRuleChainConnectionInfo(int fromIndex, RuleChainId targetRuleChainId, String type) { + RuleChainConnectionInfo connectionInfo = new RuleChainConnectionInfo(); + connectionInfo.setFromIndex(fromIndex); + connectionInfo.setTargetRuleChainId(targetRuleChainId); + connectionInfo.setType(type); + if (ruleChainConnections == null) { + ruleChainConnections = new ArrayList<>(); + } + ruleChainConnections.add(connectionInfo); + } + + @Data + public class NodeConnectionInfo { + private int fromIndex; + private int toIndex; + private String type; + } + + @Data + public class RuleChainConnectionInfo { + private int fromIndex; + private RuleChainId targetRuleChainId; + private String type; + } + +} 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 eb6be43fcf..073ccb912e 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 @@ -20,19 +20,33 @@ import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.Tenant; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.RuleChainId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; +import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleChain; +import org.thingsboard.server.common.data.rule.RuleChainMetaData; +import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.PaginatedRemover; import org.thingsboard.server.dao.service.Validator; +import org.thingsboard.server.dao.tenant.TenantDao; +import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +import java.util.Map; +import java.util.concurrent.ExecutionException; +import java.util.stream.Collectors; /** * Created by igor on 3/12/18. @@ -46,6 +60,12 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC @Autowired private RuleChainDao ruleChainDao; + @Autowired + private RuleNodeDao ruleNodeDao; + + @Autowired + private TenantDao tenantDao; + @Override public RuleChain saveRuleChain(RuleChain ruleChain) { ruleChainValidator.validate(ruleChain); @@ -53,7 +73,144 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC log.trace("Save system rule chain with predefined id {}", SYSTEM_TENANT); ruleChain.setTenantId(SYSTEM_TENANT); } - return ruleChainDao.save(ruleChain); + RuleChain savedRuleChain = ruleChainDao.save(ruleChain); + if (ruleChain.isRoot() && ruleChain.getTenantId() != null && ruleChain.getId() == null) { + try { + createRelation(new EntityRelation(savedRuleChain.getTenantId(), savedRuleChain.getId(), + EntityRelation.CONTAINS_TYPE, RelationTypeGroup.RULE_CHAIN)); + } catch (ExecutionException | InterruptedException e) { + log.warn("[{}] Failed to create tenant to root rule chain relation. from: [{}], to: [{}]", + savedRuleChain.getTenantId(), savedRuleChain.getId()); + throw new RuntimeException(e); + } + } + return savedRuleChain; + } + + @Override + public RuleChainMetaData saveRuleChainMetaData(RuleChainMetaData ruleChainMetaData) { + Validator.validateId(ruleChainMetaData.getRuleChainId(), "Incorrect rule chain id."); + RuleChain ruleChain = findRuleChainById(ruleChainMetaData.getRuleChainId()); + if (ruleChain == null) { + return null; + } + + List nodes = ruleChainMetaData.getNodes(); + List toAdd = new ArrayList<>(); + List toUpdate = new ArrayList<>(); + List toDelete = new ArrayList<>(); + + Map ruleNodeIndexMap = new HashMap<>(); + if (nodes != null) { + for (RuleNode node : nodes) { + if (node.getId() != null) { + ruleNodeIndexMap.put(node.getId(), nodes.indexOf(node)); + } else { + toAdd.add(node); + } + } + } + + List existingRuleNodes = getRuleChainNodes(ruleChainMetaData.getRuleChainId()); + for (RuleNode existingNode : existingRuleNodes) { + deleteEntityRelations(existingNode.getId()); + Integer index = ruleNodeIndexMap.get(existingNode.getId()); + if (index != null) { + toUpdate.add(ruleChainMetaData.getNodes().get(index)); + } else { + toDelete.add(existingNode); + } + } + for (RuleNode node : toAdd) { + RuleNode savedNode = ruleNodeDao.save(node); + try { + createRelation(new EntityRelation(ruleChainMetaData.getRuleChainId(), savedNode.getId(), + EntityRelation.CONTAINS_TYPE, RelationTypeGroup.RULE_CHAIN)); + } catch (ExecutionException | InterruptedException e) { + log.warn("[{}] Failed to create rule chain to rule node relation. from: [{}], to: [{}]", + ruleChainMetaData.getRuleChainId(), savedNode.getId()); + throw new RuntimeException(e); + } + int index = nodes.indexOf(node); + nodes.set(index, savedNode); + ruleNodeIndexMap.put(savedNode.getId(), index); + } + for (RuleNode node: toDelete) { + deleteRuleNode(node.getId()); + } + RuleNodeId firstRuleNodeId = null; + if (ruleChainMetaData.getFirstNodeIndex() != null) { + firstRuleNodeId = nodes.get(ruleChainMetaData.getFirstNodeIndex()).getId(); + } + if ((ruleChain.getFirstRuleNodeId() != null && !ruleChain.getFirstRuleNodeId().equals(firstRuleNodeId)) + || (ruleChain.getFirstRuleNodeId() == null && firstRuleNodeId != null)) { + ruleChain.setFirstRuleNodeId(firstRuleNodeId); + ruleChainDao.save(ruleChain); + } + if (ruleChainMetaData.getConnections() != null) { + for (RuleChainMetaData.NodeConnectionInfo nodeConnection : ruleChainMetaData.getConnections()) { + EntityId from = nodes.get(nodeConnection.getFromIndex()).getId(); + EntityId to = nodes.get(nodeConnection.getToIndex()).getId(); + String type = nodeConnection.getType(); + try { + createRelation(new EntityRelation(from, to, type, RelationTypeGroup.RULE_NODE)); + } catch (ExecutionException | InterruptedException e) { + log.warn("[{}] Failed to create rule node relation. from: [{}], to: [{}]", from, to); + throw new RuntimeException(e); + } + } + } + if (ruleChainMetaData.getRuleChainConnections() != null) { + for (RuleChainMetaData.RuleChainConnectionInfo nodeToRuleChainConnection : ruleChainMetaData.getRuleChainConnections()) { + EntityId from = nodes.get(nodeToRuleChainConnection.getFromIndex()).getId(); + EntityId to = nodeToRuleChainConnection.getTargetRuleChainId(); + String type = nodeToRuleChainConnection.getType(); + try { + createRelation(new EntityRelation(from, to, type, RelationTypeGroup.RULE_NODE)); + } catch (ExecutionException | InterruptedException e) { + log.warn("[{}] Failed to create rule node to rule chain relation. from: [{}], to: [{}]", from, to); + throw new RuntimeException(e); + } + } + } + + return loadRuleChainMetaData(ruleChainMetaData.getRuleChainId()); + } + + @Override + public RuleChainMetaData loadRuleChainMetaData(RuleChainId ruleChainId) { + Validator.validateId(ruleChainId, "Incorrect rule chain id."); + RuleChain ruleChain = findRuleChainById(ruleChainId); + if (ruleChain == null) { + return null; + } + RuleChainMetaData ruleChainMetaData = new RuleChainMetaData(); + ruleChainMetaData.setRuleChainId(ruleChainId); + List ruleNodes = getRuleChainNodes(ruleChainId); + Map ruleNodeIndexMap = new HashMap<>(); + for (RuleNode node : ruleNodes) { + ruleNodeIndexMap.put(node.getId(), ruleNodes.indexOf(node)); + } + ruleChainMetaData.setNodes(ruleNodes); + if (ruleChain.getFirstRuleNodeId() != null) { + ruleChainMetaData.setFirstNodeIndex(ruleNodeIndexMap.get(ruleChain.getFirstRuleNodeId())); + } + for (RuleNode node : ruleNodes) { + int fromIndex = ruleNodeIndexMap.get(node.getId()); + List nodeRelations = getRuleNodeRelations(node.getId()); + for (EntityRelation nodeRelation : nodeRelations) { + String type = nodeRelation.getType(); + if (nodeRelation.getTo().getEntityType() == EntityType.RULE_NODE) { + RuleNodeId toNodeId = new RuleNodeId(nodeRelation.getTo().getId()); + int toIndex = ruleNodeIndexMap.get(toNodeId); + ruleChainMetaData.addConnectionInfo(fromIndex, toIndex, type); + } else if (nodeRelation.getTo().getEntityType() == EntityType.RULE_CHAIN) { + RuleChainId targetRuleChainId = new RuleChainId(nodeRelation.getTo().getId()); + ruleChainMetaData.addRuleChainConnectionInfo(fromIndex, targetRuleChainId, type); + } + } + } + return ruleChainMetaData; } @Override @@ -62,6 +219,33 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC return ruleChainDao.findById(ruleChainId.getId()); } + @Override + public RuleChain getRootTenantRuleChain(TenantId tenantId) { + Validator.validateId(tenantId, "Incorrect tenant id for search request."); + List relations = relationService.findByFrom(tenantId, RelationTypeGroup.RULE_CHAIN); + if (relations != null && !relations.isEmpty()) { + EntityRelation relation = relations.get(0); + RuleChainId ruleChainId = new RuleChainId(relation.getTo().getId()); + return findRuleChainById(ruleChainId); + } else { + return null; + } + } + + @Override + public List getRuleChainNodes(RuleChainId ruleChainId) { + Validator.validateId(ruleChainId, "Incorrect rule chain id for search request."); + List relations = getRuleChainToNodeRelations(ruleChainId); + List ruleNodes = relations.stream().map(relation -> ruleNodeDao.findById(relation.getTo().getId())).collect(Collectors.toList()); + return ruleNodes; + } + + @Override + public List getRuleNodeRelations(RuleNodeId ruleNodeId) { + Validator.validateId(ruleNodeId, "Incorrect rule node id for search request."); + return relationService.findByFrom(ruleNodeId, RelationTypeGroup.RULE_NODE); + } + @Override public TextPageData findSystemRuleChains(TextPageLink pageLink) { Validator.validatePageLink(pageLink, "Incorrect PageLink object for search system rule chain request."); @@ -92,16 +276,33 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC checkRuleNodesAndDelete(ruleChainId); } + @Override + public void deleteRuleChainsByTenantId(TenantId tenantId) { + Validator.validateId(tenantId, "Incorrect tenant id for delete rule chains request."); + tenantRuleChainsRemover.removeEntities(tenantId); + } + private void checkRuleNodesAndDelete(RuleChainId ruleChainId) { - //TODO: + List nodeRelations = getRuleChainToNodeRelations(ruleChainId); + for (EntityRelation relation : nodeRelations) { + deleteRuleNode(relation.getTo()); + } deleteEntityRelations(ruleChainId); ruleChainDao.removeById(ruleChainId.getId()); } - @Override - public void deleteRuleChainsByTenantId(TenantId tenantId) { - Validator.validateId(tenantId, "Incorrect tenant id for delete rule chains request."); - tenantRuleChainsRemover.removeEntities(tenantId); + private List getRuleChainToNodeRelations(RuleChainId ruleChainId) { + return relationService.findByFrom(ruleChainId, RelationTypeGroup.RULE_CHAIN); + } + + private void deleteRuleNode(EntityId entityId) { + deleteEntityRelations(entityId); + ruleNodeDao.removeById(entityId.getId()); + } + + private void createRelation(EntityRelation relation) throws ExecutionException, InterruptedException { + log.debug("Creating relation: {}", relation); + relationService.saveRelationAsync(relation).get(); } private DataValidator ruleChainValidator = @@ -111,6 +312,18 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC if (StringUtils.isEmpty(ruleChain.getName())) { throw new DataValidationException("Rule chain name should be specified!."); } + if (ruleChain.getTenantId() != null && !ruleChain.getTenantId().isNullUid()) { + Tenant tenant = tenantDao.findById(ruleChain.getTenantId().getId()); + if (tenant == null) { + throw new DataValidationException("Rule chain is referencing to non-existent tenant!"); + } + if (ruleChain.isRoot()) { + RuleChain rootRuleChain = getRootTenantRuleChain(ruleChain.getTenantId()); + if (ruleChain.getId() == null || !ruleChain.getId().equals(rootRuleChain.getId())) { + throw new DataValidationException("Another root rule chain is present in scope of current tenant!"); + } + } + } } }; diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java index 00dc97315e..c6a2941938 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java @@ -18,11 +18,15 @@ package org.thingsboard.server.dao.rule; import org.thingsboard.server.common.data.id.PluginId; import org.thingsboard.server.common.data.id.RuleChainId; +import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.plugin.PluginMetaData; +import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.rule.RuleChain; +import org.thingsboard.server.common.data.rule.RuleChainMetaData; +import org.thingsboard.server.common.data.rule.RuleNode; import java.util.List; @@ -33,8 +37,18 @@ public interface RuleChainService { RuleChain saveRuleChain(RuleChain ruleChain); + RuleChainMetaData saveRuleChainMetaData(RuleChainMetaData ruleChainMetaData); + + RuleChainMetaData loadRuleChainMetaData(RuleChainId ruleChainId); + RuleChain findRuleChainById(RuleChainId ruleChainId); + RuleChain getRootTenantRuleChain(TenantId tenantId); + + List getRuleChainNodes(RuleChainId ruleChainId); + + List getRuleNodeRelations(RuleNodeId ruleNodeId); + TextPageData findSystemRuleChains(TextPageLink pageLink); TextPageData findTenantRuleChains(TenantId tenantId, TextPageLink pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java index 24abe52802..3a607b623e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java @@ -31,6 +31,7 @@ import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.plugin.PluginService; +import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.rule.RuleService; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.PaginatedRemover; @@ -76,6 +77,9 @@ public class TenantServiceImpl extends AbstractEntityService implements TenantSe @Autowired private PluginService pluginService; + @Autowired + private RuleChainService ruleChainService; + @Override public Tenant findTenantById(TenantId tenantId) { log.trace("Executing findTenantById [{}]", tenantId); @@ -108,6 +112,7 @@ public class TenantServiceImpl extends AbstractEntityService implements TenantSe assetService.deleteAssetsByTenantId(tenantId); deviceService.deleteDevicesByTenantId(tenantId); userService.deleteTenantAdmins(tenantId); + ruleChainService.deleteRuleChainsByTenantId(tenantId); ruleService.deleteRulesByTenantId(tenantId); pluginService.deletePluginsByTenantId(tenantId); tenantDao.removeById(tenantId.getId());