From 6ab3714ac51706c684ca60e35665fdd53f901d82 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 27 May 2020 09:22:25 +0300 Subject: [PATCH] Added default rule chain for server side --- .../server/controller/EdgeController.java | 4 +- .../controller/RuleChainController.java | 202 +++++------------- .../CassandraDatabaseUpgradeService.java | 8 - .../install/SqlDatabaseUpgradeService.java | 6 - .../server/dao/edge/EdgeService.java | 6 +- .../server/dao/rule/RuleChainService.java | 9 +- .../server/common/data/rule/RuleChain.java | 27 --- .../server/dao/edge/CassandraEdgeDao.java | 29 +++ .../thingsboard/server/dao/edge/EdgeDao.java | 11 + .../server/dao/edge/EdgeServiceImpl.java | 39 +++- .../server/dao/model/ModelConstants.java | 1 - .../dao/model/nosql/RuleChainEntity.java | 30 --- .../server/dao/model/sql/RuleChainEntity.java | 28 --- .../server/dao/rule/BaseRuleChainService.java | 133 ++++++------ .../dao/rule/CassandraRuleChainDao.java | 15 ++ .../server/dao/rule/RuleChainDao.java | 8 + .../server/dao/sql/edge/JpaEdgeDao.java | 27 +++ .../server/dao/sql/rule/JpaRuleChainDao.java | 15 ++ .../resources/sql/schema-entities-hsql.sql | 4 +- .../main/resources/sql/schema-entities.sql | 4 +- ui/src/app/api/rule-chain.service.js | 109 ++++------ ui/src/app/rulechain/rulechains.controller.js | 27 ++- 22 files changed, 327 insertions(+), 415 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/EdgeController.java b/application/src/main/java/org/thingsboard/server/controller/EdgeController.java index bab992f456..bb9106719d 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EdgeController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EdgeController.java @@ -96,7 +96,7 @@ public class EdgeController extends BaseController { if (created) { ruleChainService.assignRuleChainToEdge(tenantId, defaultRootEdgeRuleChain.getId(), result.getId()); - edgeService.setRootRuleChain(tenantId, result, defaultRootEdgeRuleChain.getId()); + edgeService.setEdgeRootRuleChain(tenantId, result, defaultRootEdgeRuleChain.getId()); } logEntityAction(result.getId(), result, null, created ? ActionType.ADDED : ActionType.UPDATED, null); @@ -281,7 +281,7 @@ public class EdgeController extends BaseController { accessControlService.checkPermission(getCurrentUser(), Resource.EDGE, Operation.WRITE, edge.getId(), edge); - Edge updatedEdge = edgeService.setRootRuleChain(getTenantId(), edge, ruleChainId); + Edge updatedEdge = edgeService.setEdgeRootRuleChain(getTenantId(), edge, ruleChainId); logEntityAction(updatedEdge.getId(), updatedEdge, null, ActionType.UPDATED, null); 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 f1486a8c19..3d6b1f3349 100644 --- a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java +++ b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.controller; -import com.datastax.driver.core.utils.UUIDs; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; @@ -40,7 +39,6 @@ import org.thingsboard.server.actors.tenant.DebugTbRateLimits; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.Event; -import org.thingsboard.server.common.data.ShortEdgeInfo; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.exception.ThingsboardException; @@ -62,13 +60,11 @@ import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.dao.event.EventService; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.service.script.JsInvokeService; import org.thingsboard.server.service.script.RuleNodeJsScriptEngine; import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Resource; -import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Set; @@ -450,189 +446,75 @@ public class RuleChainController extends BaseController { } } - @PreAuthorize("hasAuthority('TENANT_ADMIN')") - @RequestMapping(value = "/ruleChain/{ruleChainId}/edges", method = RequestMethod.POST) + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") + @RequestMapping(value = "/edge/{edgeId}/ruleChains", params = {"limit"}, method = RequestMethod.GET) @ResponseBody - public RuleChain updateRuleChainEdges(@PathVariable(RULE_CHAIN_ID) String strRuleChainId, - @RequestBody String[] strEdgeIds) throws ThingsboardException { - checkParameter(RULE_CHAIN_ID, strRuleChainId); + public TimePageData getEdgeRuleChains( + @PathVariable("edgeId") String strEdgeId, + @RequestParam int limit, + @RequestParam(required = false) Long startTime, + @RequestParam(required = false) Long endTime, + @RequestParam(required = false, defaultValue = "false") boolean ascOrder, + @RequestParam(required = false) String offset) throws ThingsboardException { + checkParameter("edgeId", strEdgeId); try { - RuleChainId ruleChainId = new RuleChainId(toUUID(strRuleChainId)); - RuleChain ruleChain = checkRuleChain(ruleChainId, Operation.ASSIGN_TO_EDGE); - - Set edgeIds = new HashSet<>(); - if (strEdgeIds != null) { - for (String strEdgeId : strEdgeIds) { - edgeIds.add(new EdgeId(toUUID(strEdgeId))); - } - } - - Set addedEdgeIds = new HashSet<>(); - Set removedEdgeIds = new HashSet<>(); - for (EdgeId edgeId : edgeIds) { - if (!ruleChain.isAssignedToEdge(edgeId)) { - addedEdgeIds.add(edgeId); - } - } - - Set assignedEdges = ruleChain.getAssignedEdges(); - if (assignedEdges != null) { - for (ShortEdgeInfo edgeInfo : assignedEdges) { - if (!edgeIds.contains(edgeInfo.getEdgeId())) { - removedEdgeIds.add(edgeInfo.getEdgeId()); - } - } - } - - if (addedEdgeIds.isEmpty() && removedEdgeIds.isEmpty()) { - return ruleChain; - } else { - RuleChain savedRuleChain = null; - for (EdgeId edgeId : addedEdgeIds) { - savedRuleChain = checkNotNull(ruleChainService.assignRuleChainToEdge(getCurrentUser().getTenantId(), ruleChainId, edgeId)); - ShortEdgeInfo edgeInfo = savedRuleChain.getAssignedEdgeInfo(edgeId); - logEntityAction(ruleChainId, savedRuleChain, - null, - ActionType.ASSIGNED_TO_EDGE, null, strRuleChainId, edgeId.toString(), edgeInfo.getTitle()); - } - for (EdgeId edgeId : removedEdgeIds) { - ShortEdgeInfo edgeInfo = ruleChain.getAssignedEdgeInfo(edgeId); - savedRuleChain = checkNotNull(ruleChainService.unassignRuleChainFromEdge(getCurrentUser().getTenantId(), ruleChainId, edgeId, false)); - logEntityAction(ruleChainId, ruleChain, - null, - ActionType.UNASSIGNED_FROM_EDGE, null, strRuleChainId, edgeId.toString(), edgeInfo.getTitle()); - - } - return savedRuleChain; - } + TenantId tenantId = getCurrentUser().getTenantId(); + EdgeId edgeId = new EdgeId(toUUID(strEdgeId)); + checkEdgeId(edgeId, Operation.READ); + TimePageLink pageLink = createPageLink(limit, startTime, endTime, ascOrder, offset); + return checkNotNull(ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink).get()); } catch (Exception e) { - - logEntityAction(emptyId(EntityType.RULE_CHAIN), null, - null, - ActionType.ASSIGNED_TO_EDGE, e, strRuleChainId); - throw handleException(e); } } - @PreAuthorize("hasAuthority('TENANT_ADMIN')") - @RequestMapping(value = "/ruleChain/{ruleChainId}/edges/add", method = RequestMethod.POST) + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") + @RequestMapping(value = "/ruleChain/{ruleChainId}/defaultRootEdge", method = RequestMethod.POST) @ResponseBody - public RuleChain addRuleChainEdges(@PathVariable(RULE_CHAIN_ID) String strRuleChainId, - @RequestBody String[] strEdgeIds) throws ThingsboardException { + public RuleChain setDefaultRootEdgeRuleChain(@PathVariable(RULE_CHAIN_ID) String strRuleChainId) throws ThingsboardException { checkParameter(RULE_CHAIN_ID, strRuleChainId); try { RuleChainId ruleChainId = new RuleChainId(toUUID(strRuleChainId)); - RuleChain ruleChain = checkRuleChain(ruleChainId, Operation.ASSIGN_TO_EDGE); - - Set edgeIds = new HashSet<>(); - if (strEdgeIds != null) { - for (String strEdgeId : strEdgeIds) { - EdgeId edgeId = new EdgeId(toUUID(strEdgeId)); - if (!ruleChain.isAssignedToEdge(edgeId)) { - edgeIds.add(edgeId); - } - } - } - - if (edgeIds.isEmpty()) { - return ruleChain; - } else { - RuleChain savedRuleChain = null; - for (EdgeId edgeId : edgeIds) { - savedRuleChain = checkNotNull(ruleChainService.assignRuleChainToEdge(getCurrentUser().getTenantId(), ruleChainId, edgeId)); - ShortEdgeInfo edgeInfo = savedRuleChain.getAssignedEdgeInfo(edgeId); - logEntityAction(ruleChainId, savedRuleChain, - null, - ActionType.ASSIGNED_TO_EDGE, null, strRuleChainId, edgeId.toString(), edgeInfo.getTitle()); - } - return savedRuleChain; - } + RuleChain ruleChain = checkRuleChain(ruleChainId, Operation.WRITE); + ruleChainService.setDefaultRootEdgeRuleChain(getTenantId(), ruleChainId); + return ruleChain; } catch (Exception e) { - - logEntityAction(emptyId(EntityType.RULE_CHAIN), null, + logEntityAction(emptyId(EntityType.RULE_CHAIN), null, - ActionType.ASSIGNED_TO_EDGE, e, strRuleChainId); - + null, + ActionType.UPDATED, e, strRuleChainId); throw handleException(e); } } @PreAuthorize("hasAuthority('TENANT_ADMIN')") - @RequestMapping(value = "/ruleChain/{ruleChainId}/edges/remove", method = RequestMethod.POST) + @RequestMapping(value = "/ruleChain/{ruleChainId}/defaultEdge", method = RequestMethod.POST) @ResponseBody - public RuleChain removeRuleChainEdges(@PathVariable(RULE_CHAIN_ID) String strRuleChainId, - @RequestBody String[] strEdgeIds) throws ThingsboardException { + public RuleChain addDefaultEdgeRuleChain(@PathVariable(RULE_CHAIN_ID) String strRuleChainId) throws ThingsboardException { checkParameter(RULE_CHAIN_ID, strRuleChainId); try { RuleChainId ruleChainId = new RuleChainId(toUUID(strRuleChainId)); - RuleChain ruleChain = checkRuleChain(ruleChainId, Operation.UNASSIGN_FROM_EDGE); - - Set edgeIds = new HashSet<>(); - if (strEdgeIds != null) { - for (String strEdgeId : strEdgeIds) { - EdgeId edgeId = new EdgeId(toUUID(strEdgeId)); - if (ruleChain.isAssignedToEdge(edgeId)) { - edgeIds.add(edgeId); - } - } - } - - if (edgeIds.isEmpty()) { - return ruleChain; - } else { - RuleChain savedRuleChain = null; - for (EdgeId edgeId : edgeIds) { - ShortEdgeInfo edgeInfo = ruleChain.getAssignedEdgeInfo(edgeId); - savedRuleChain = checkNotNull(ruleChainService.unassignRuleChainFromEdge(getCurrentUser().getTenantId(), ruleChainId, edgeId, false)); - logEntityAction(ruleChainId, ruleChain, - null, - ActionType.UNASSIGNED_FROM_EDGE, null, strRuleChainId, edgeId.toString(), edgeInfo.getTitle()); - - } - return savedRuleChain; - } + RuleChain ruleChain = checkRuleChain(ruleChainId, Operation.WRITE); + ruleChainService.addDefaultEdgeRuleChain(getTenantId(), ruleChainId); + return ruleChain; } catch (Exception e) { - - logEntityAction(emptyId(EntityType.RULE_CHAIN), null, + logEntityAction(emptyId(EntityType.RULE_CHAIN), null, - ActionType.UNASSIGNED_FROM_EDGE, e, strRuleChainId); - - throw handleException(e); - } - } - - @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") - @RequestMapping(value = "/edge/{edgeId}/ruleChains", params = { "limit" }, method = RequestMethod.GET) - @ResponseBody - public TimePageData getEdgeRuleChains( - @PathVariable("edgeId") String strEdgeId, - @RequestParam int limit, - @RequestParam(required = false) Long startTime, - @RequestParam(required = false) Long endTime, - @RequestParam(required = false, defaultValue = "false") boolean ascOrder, - @RequestParam(required = false) String offset) throws ThingsboardException { - checkParameter("edgeId", strEdgeId); - try { - TenantId tenantId = getCurrentUser().getTenantId(); - EdgeId edgeId = new EdgeId(toUUID(strEdgeId)); - checkEdgeId(edgeId, Operation.READ); - TimePageLink pageLink = createPageLink(limit, startTime, endTime, ascOrder, offset); - return checkNotNull(ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink).get()); - } catch (Exception e) { + null, + ActionType.UPDATED, e, strRuleChainId); throw handleException(e); } } - @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") - @RequestMapping(value = "/ruleChain/{ruleChainId}/defaultRootEdge", method = RequestMethod.POST) + @PreAuthorize("hasAuthority('TENANT_ADMIN')") + @RequestMapping(value = "/ruleChain/{ruleChainId}/defaultEdge", method = RequestMethod.DELETE) @ResponseBody - public RuleChain setDefaultRootEdgeRuleChain(@PathVariable(RULE_CHAIN_ID) String strRuleChainId) throws ThingsboardException { + public RuleChain removeDefaultEdgeRuleChain(@PathVariable(RULE_CHAIN_ID) String strRuleChainId) throws ThingsboardException { checkParameter(RULE_CHAIN_ID, strRuleChainId); try { RuleChainId ruleChainId = new RuleChainId(toUUID(strRuleChainId)); RuleChain ruleChain = checkRuleChain(ruleChainId, Operation.WRITE); - ruleChainService.setDefaultRootEdgeRuleChain(getTenantId(), ruleChainId); + ruleChainService.removeDefaultEdgeRuleChain(getTenantId(), ruleChainId); return ruleChain; } catch (Exception e) { logEntityAction(emptyId(EntityType.RULE_CHAIN), @@ -642,4 +524,16 @@ public class RuleChainController extends BaseController { throw handleException(e); } } + + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") + @RequestMapping(value = "/ruleChain/defaultEdgeRuleChains", method = RequestMethod.GET) + @ResponseBody + public List getDefaultEdgeRuleChains() throws ThingsboardException { + try { + TenantId tenantId = getCurrentUser().getTenantId(); + return checkNotNull(ruleChainService.findDefaultEdgeRuleChainsByTenantId(tenantId)).get(); + } catch (Exception e) { + throw handleException(e); + } + } } diff --git a/application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java index 56b4e011cb..5227efcfa3 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java @@ -323,14 +323,6 @@ public class CassandraDatabaseUpgradeService extends AbstractCassandraDatabaseUp cluster.getSession().execute("alter table entity_view add edge_id text"); Thread.sleep(2500); } catch (InvalidQueryException e) {} - try { - cluster.getSession().execute("alter table dashboard add assigned_edges text"); - Thread.sleep(2500); - } catch (InvalidQueryException e) {} - try { - cluster.getSession().execute("alter table rule_chain add assigned_edges text"); - Thread.sleep(2500); - } catch (InvalidQueryException e) {} try { cluster.getSession().execute("alter table rule_chain add type text"); Thread.sleep(2500); diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java index bff44636d6..c9966d0cbb 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java @@ -247,12 +247,6 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService try { conn.createStatement().execute("ALTER TABLE entity_view ADD edge_id varchar(31)"); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script } catch (Exception e) {} - try { - conn.createStatement().execute("ALTER TABLE dashboard ADD assigned_edges varchar(10000000)"); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script - } catch (Exception e) {} - try { - conn.createStatement().execute("ALTER TABLE rule_chain ADD assigned_edges varchar(10000000)"); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script - } catch (Exception e) {} try { conn.createStatement().execute("ALTER TABLE rule_chain ADD type varchar(255) DEFAULT 'SYSTEM'"); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script } catch (Exception e) {} diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java index 29e317fe59..8b4baffbae 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java @@ -29,11 +29,13 @@ import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.page.TimePageData; import org.thingsboard.server.common.data.page.TimePageLink; +import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.msg.TbMsg; import java.io.IOException; import java.util.List; import java.util.Optional; +import java.util.UUID; public interface EdgeService { @@ -77,7 +79,9 @@ public interface EdgeService { TimePageData findQueueEvents(TenantId tenantId, EdgeId edgeId, TimePageLink pageLink); - Edge setRootRuleChain(TenantId tenantId, Edge edge, RuleChainId ruleChainId) throws IOException; + Edge setEdgeRootRuleChain(TenantId tenantId, Edge edge, RuleChainId ruleChainId) throws IOException; + + ListenableFuture> findEdgesByTenantIdAndRuleChainId(TenantId tenantId, RuleChainId ruleChainId, TimePageLink pageLink); } 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 8cb77d541b..d0e28a4a57 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 @@ -73,11 +73,16 @@ public interface RuleChainService { void unassignEdgeRuleChains(TenantId tenantId, EdgeId edgeId); - void updateEdgeRuleChains(TenantId tenantId, EdgeId edgeId); - ListenableFuture> findRuleChainsByTenantIdAndEdgeId(TenantId tenantId, EdgeId edgeId, TimePageLink pageLink); RuleChain getDefaultRootEdgeRuleChain(TenantId tenantId); boolean setDefaultRootEdgeRuleChain(TenantId tenantId, RuleChainId ruleChainId); + + boolean addDefaultEdgeRuleChain(TenantId tenantId, RuleChainId ruleChainId); + + boolean removeDefaultEdgeRuleChain(TenantId tenantId, RuleChainId ruleChainId); + + ListenableFuture> findDefaultEdgeRuleChainsByTenantId(TenantId tenantId); + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChain.java b/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChain.java index 49508e2971..cedc458edb 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChain.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChain.java @@ -48,7 +48,6 @@ public class RuleChain extends SearchTextBasedWithAdditionalInfo im private boolean root; private boolean debugMode; private transient JsonNode configuration; - private Set assignedEdges; @JsonIgnore private byte[] configurationBytes; @@ -68,7 +67,6 @@ public class RuleChain extends SearchTextBasedWithAdditionalInfo im this.type = ruleChain.getType(); this.firstRuleNodeId = ruleChain.getFirstRuleNodeId(); this.root = ruleChain.isRoot(); - this.assignedEdges = ruleChain.getAssignedEdges(); this.setConfiguration(ruleChain.getConfiguration()); } @@ -89,29 +87,4 @@ public class RuleChain extends SearchTextBasedWithAdditionalInfo im public void setConfiguration(JsonNode data) { setJson(data, json -> this.configuration = json, bytes -> this.configurationBytes = bytes); } - - - - public boolean isAssignedToEdge(EdgeId edgeId) { - return EdgeUtils.isAssignedToEdge(this.assignedEdges, edgeId); - } - - public ShortEdgeInfo getAssignedEdgeInfo(EdgeId edgeId) { - return EdgeUtils.getAssignedEdgeInfo(this.assignedEdges, edgeId); - } - - public boolean addAssignedEdge(Edge edge) { - if (this.assignedEdges == null) { - this.assignedEdges = new HashSet<>(); - } - return EdgeUtils.addAssignedEdge(this.assignedEdges, edge.toShortEdgeInfo()); - } - - public boolean updateAssignedEdge(Edge edge) { - return EdgeUtils.updateAssignedEdge(this.assignedEdges, edge.toShortEdgeInfo()); - } - - public boolean removeAssignedEdge(Edge edge) { - return EdgeUtils.removeAssignedEdge(this.assignedEdges, edge.toShortEdgeInfo()); - } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/CassandraEdgeDao.java b/dao/src/main/java/org/thingsboard/server/dao/edge/CassandraEdgeDao.java index 80d51303b1..42253db06c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/CassandraEdgeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/CassandraEdgeDao.java @@ -15,16 +15,29 @@ */ package org.thingsboard.server.dao.edge; +import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EntitySubtype; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.RuleChainId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageLink; +import org.thingsboard.server.common.data.page.TimePageLink; +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.dao.model.nosql.EdgeEntity; import org.thingsboard.server.dao.nosql.CassandraAbstractSearchTextDao; +import org.thingsboard.server.dao.relation.RelationDao; import org.thingsboard.server.dao.util.NoSqlDao; +import java.util.ArrayList; import java.util.List; import java.util.Optional; import java.util.UUID; @@ -36,6 +49,9 @@ import static org.thingsboard.server.dao.model.ModelConstants.EDGE_COLUMN_FAMILY @NoSqlDao public class CassandraEdgeDao extends CassandraAbstractSearchTextDao implements EdgeDao { + @Autowired + private RelationDao relationDao; + @Override protected Class getColumnFamilyClass() { return EdgeEntity.class; @@ -91,4 +107,17 @@ public class CassandraEdgeDao extends CassandraAbstractSearchTextDao findByRoutingKey(UUID tenantId, String routingKey) { return Optional.empty(); } + + @Override + public ListenableFuture> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId, TimePageLink pageLink) { + log.debug("Try to find edges by tenantId [{}], ruleChainId [{}] and pageLink [{}]", tenantId, ruleChainId, pageLink); + ListenableFuture> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new RuleChainId(ruleChainId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); + return Futures.transformAsync(relations, input -> { + List> edgeFutures = new ArrayList<>(input.size()); + for (EntityRelation relation : input) { + edgeFutures.add(findByIdAsync(new TenantId(tenantId), relation.getTo().getId())); + } + return Futures.successfulAsList(edgeFutures); + }, MoreExecutors.directExecutor()); + } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java index 34d06b338f..0002b96ad0 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java @@ -20,6 +20,8 @@ import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageLink; +import org.thingsboard.server.common.data.page.TimePageLink; +import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.dao.Dao; import java.util.List; @@ -124,4 +126,13 @@ public interface EdgeDao extends Dao { */ Optional findByRoutingKey(UUID tenantId, String routingKey); + /** + * Find edges by tenantId, ruleChainId and page link. + * + * @param tenantId the tenantId + * @param ruleChainId the ruleChainId + * @param pageLink the page link + * @return the list of rule chain objects + */ + ListenableFuture> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId, TimePageLink pageLink); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java index ef384a9fba..3bfb1d114e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java @@ -62,6 +62,7 @@ import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntitySearchDirection; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; +import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.session.SessionMsgType; import org.thingsboard.server.dao.asset.AssetService; @@ -531,10 +532,20 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic case DataConstants.ENTITY_CREATED: case DataConstants.ENTITY_UPDATED: RuleChain ruleChain = mapper.readValue(tbMsg.getData(), RuleChain.class); - if (ruleChain.getAssignedEdges() != null && !ruleChain.getAssignedEdges().isEmpty()) { - for (ShortEdgeInfo assignedEdge : ruleChain.getAssignedEdges()) { - pushEventToEdge(tenantId, assignedEdge.getEdgeId(), EdgeQueueEntityType.RULE_CHAIN, tbMsg, callback); - } + if (RuleChainType.EDGE.equals(ruleChain.getType())) { + ListenableFuture> future = findEdgesByTenantIdAndRuleChainId(tenantId, ruleChain.getId(), new TimePageLink(Integer.MAX_VALUE)); + Futures.transform(future, edges -> { + if (edges != null && edges.getData() != null && !edges.getData().isEmpty()) { + try { + for (Edge edge : edges.getData()) { + pushEventToEdge(tenantId, edge.getId(), EdgeQueueEntityType.RULE_CHAIN, tbMsg, callback); + } + } catch (IOException e) { + log.error("Can't push event to edge", e); + } + } + return null; + }, MoreExecutors.directExecutor()); } break; default: @@ -628,10 +639,9 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic } @Override - public Edge setRootRuleChain(TenantId tenantId, Edge edge, RuleChainId ruleChainId) throws IOException { + public Edge setEdgeRootRuleChain(TenantId tenantId, Edge edge, RuleChainId ruleChainId) throws IOException { edge.setRootRuleChainId(ruleChainId); Edge savedEdge = saveEdge(edge); - ruleChainService.updateEdgeRuleChains(tenantId, savedEdge.getId()); RuleChain ruleChain = ruleChainService.findRuleChainById(tenantId, ruleChainId); saveEventToEdgeQueue(tenantId, edge.getId(), EdgeQueueEntityType.RULE_CHAIN, DataConstants.ENTITY_UPDATED, mapper.writeValueAsString(ruleChain), new FutureCallback() { @Override @@ -647,6 +657,23 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic return savedEdge; } + @Override + public ListenableFuture> findEdgesByTenantIdAndRuleChainId(TenantId tenantId, RuleChainId ruleChainId, TimePageLink pageLink) { + log.trace("Executing findEdgesByTenantIdAndRuleChainId, tenantId [{}], ruleChainId [{}], pageLink [{}]", tenantId, ruleChainId, pageLink); + Validator.validateId(tenantId, "Incorrect tenantId " + tenantId); + Validator.validateId(ruleChainId, "Incorrect ruleChainId " + ruleChainId); + Validator.validatePageLink(pageLink, "Incorrect page link " + pageLink); + ListenableFuture> edges = edgeDao.findEdgesByTenantIdAndRuleChainId(tenantId.getId(), ruleChainId.getId(), pageLink); + + return Futures.transform(edges, new Function, TimePageData>() { + @Nullable + @Override + public TimePageData apply(@Nullable List edges) { + return new TimePageData<>(edges, pageLink); + } + }, MoreExecutors.directExecutor()); + } + private DataValidator edgeValidator = new DataValidator() { diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index 67ed99cac3..5e5f8398b5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java @@ -347,7 +347,6 @@ public class ModelConstants { public static final String RULE_CHAIN_FIRST_RULE_NODE_ID_PROPERTY = "first_rule_node_id"; public static final String RULE_CHAIN_ROOT_PROPERTY = "root"; public static final String RULE_CHAIN_CONFIGURATION_PROPERTY = "configuration"; - public static final String RULE_CHAIN_ASSIGNED_EDGES_PROPERTY = "assigned_edges"; public static final String RULE_CHAIN_BY_TENANT_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME = "rule_chain_by_tenant_and_search_text"; public static final String RULE_CHAIN_BY_TENANT_BY_TYPE_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME = "rule_chain_by_tenant_by_type_and_search_text"; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/nosql/RuleChainEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/nosql/RuleChainEntity.java index 9a03f5c582..9232a8f9ff 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/nosql/RuleChainEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/nosql/RuleChainEntity.java @@ -20,17 +20,12 @@ import com.datastax.driver.mapping.annotations.ClusteringColumn; import com.datastax.driver.mapping.annotations.Column; import com.datastax.driver.mapping.annotations.PartitionKey; import com.datastax.driver.mapping.annotations.Table; -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.JavaType; import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.ObjectMapper; import lombok.EqualsAndHashCode; import lombok.Getter; import lombok.Setter; import lombok.ToString; import lombok.extern.slf4j.Slf4j; -import org.springframework.util.StringUtils; -import org.thingsboard.server.common.data.ShortEdgeInfo; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.TenantId; @@ -41,14 +36,11 @@ import org.thingsboard.server.dao.model.SearchTextEntity; import org.thingsboard.server.dao.model.type.JsonCodec; import org.thingsboard.server.dao.model.type.RuleChainTypeCodec; -import java.io.IOException; -import java.util.HashSet; import java.util.UUID; import static org.thingsboard.server.dao.model.ModelConstants.ADDITIONAL_INFO_PROPERTY; import static org.thingsboard.server.dao.model.ModelConstants.DEBUG_MODE; import static org.thingsboard.server.dao.model.ModelConstants.ID_PROPERTY; -import static org.thingsboard.server.dao.model.ModelConstants.RULE_CHAIN_ASSIGNED_EDGES_PROPERTY; import static org.thingsboard.server.dao.model.ModelConstants.RULE_CHAIN_COLUMN_FAMILY_NAME; import static org.thingsboard.server.dao.model.ModelConstants.RULE_CHAIN_CONFIGURATION_PROPERTY; import static org.thingsboard.server.dao.model.ModelConstants.RULE_CHAIN_FIRST_RULE_NODE_ID_PROPERTY; @@ -64,10 +56,6 @@ import static org.thingsboard.server.dao.model.ModelConstants.SEARCH_TEXT_PROPER @ToString public class RuleChainEntity implements SearchTextEntity { - private static final ObjectMapper objectMapper = new ObjectMapper(); - private static final JavaType assignedEdgesType = - objectMapper.getTypeFactory().constructCollectionType(HashSet.class, ShortEdgeInfo.class); - @PartitionKey @Column(name = ID_PROPERTY) private UUID id; @@ -93,10 +81,6 @@ public class RuleChainEntity implements SearchTextEntity { @Column(name = ADDITIONAL_INFO_PROPERTY, codec = JsonCodec.class) private JsonNode additionalInfo; - @Getter @Setter - @Column(name = RULE_CHAIN_ASSIGNED_EDGES_PROPERTY) - private String assignedEdges; - public RuleChainEntity() { } @@ -113,13 +97,6 @@ public class RuleChainEntity implements SearchTextEntity { this.debugMode = ruleChain.isDebugMode(); this.configuration = ruleChain.getConfiguration(); this.additionalInfo = ruleChain.getAdditionalInfo(); - if (ruleChain.getAssignedEdges() != null) { - try { - this.assignedEdges = objectMapper.writeValueAsString(ruleChain.getAssignedEdges()); - } catch (JsonProcessingException e) { - log.error("Unable to serialize assigned edges to string!", e); - } - } } @Override @@ -208,13 +185,6 @@ public class RuleChainEntity implements SearchTextEntity { ruleChain.setDebugMode(this.debugMode); ruleChain.setConfiguration(this.configuration); ruleChain.setAdditionalInfo(this.additionalInfo); - if (!StringUtils.isEmpty(assignedEdges)) { - try { - ruleChain.setAssignedEdges(objectMapper.readValue(assignedEdges, assignedEdgesType)); - } catch (IOException e) { - log.warn("Unable to parse assigned edges!", e); - } - } return ruleChain; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/RuleChainEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/RuleChainEntity.java index ce62554fa1..1695bf7370 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/RuleChainEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/RuleChainEntity.java @@ -16,17 +16,12 @@ package org.thingsboard.server.dao.model.sql; import com.datastax.driver.core.utils.UUIDs; -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.JavaType; import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.ObjectMapper; import lombok.Data; import lombok.EqualsAndHashCode; import lombok.extern.slf4j.Slf4j; import org.hibernate.annotations.Type; import org.hibernate.annotations.TypeDef; -import org.springframework.util.StringUtils; -import org.thingsboard.server.common.data.ShortEdgeInfo; import org.thingsboard.server.common.data.UUIDConverter; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; @@ -44,8 +39,6 @@ import javax.persistence.Entity; import javax.persistence.EnumType; import javax.persistence.Enumerated; import javax.persistence.Table; -import java.io.IOException; -import java.util.HashSet; import static org.thingsboard.server.dao.model.ModelConstants.RULE_CHAIN_TYPE_PROPERTY; @@ -57,10 +50,6 @@ import static org.thingsboard.server.dao.model.ModelConstants.RULE_CHAIN_TYPE_PR @Table(name = ModelConstants.RULE_CHAIN_COLUMN_FAMILY_NAME) public class RuleChainEntity extends BaseSqlEntity implements SearchTextEntity { - private static final ObjectMapper objectMapper = new ObjectMapper(); - private static final JavaType assignedEdgesType = - objectMapper.getTypeFactory().constructCollectionType(HashSet.class, ShortEdgeInfo.class); - @Column(name = ModelConstants.RULE_CHAIN_TENANT_ID_PROPERTY) private String tenantId; @@ -91,9 +80,6 @@ public class RuleChainEntity extends BaseSqlEntity implements SearchT @Column(name = ModelConstants.ADDITIONAL_INFO_PROPERTY) private JsonNode additionalInfo; - @Column(name = ModelConstants.RULE_CHAIN_ASSIGNED_EDGES_PROPERTY) - private String assignedEdges; - public RuleChainEntity() { } @@ -112,13 +98,6 @@ public class RuleChainEntity extends BaseSqlEntity implements SearchT this.debugMode = ruleChain.isDebugMode(); this.configuration = ruleChain.getConfiguration(); this.additionalInfo = ruleChain.getAdditionalInfo(); - if (ruleChain.getAssignedEdges() != null) { - try { - this.assignedEdges = objectMapper.writeValueAsString(ruleChain.getAssignedEdges()); - } catch (JsonProcessingException e) { - log.error("Unable to serialize assigned edges to string!", e); - } - } } @Override @@ -145,13 +124,6 @@ public class RuleChainEntity extends BaseSqlEntity implements SearchT ruleChain.setDebugMode(debugMode); ruleChain.setConfiguration(configuration); ruleChain.setAdditionalInfo(additionalInfo); - if (!StringUtils.isEmpty(assignedEdges)) { - try { - ruleChain.setAssignedEdges(objectMapper.readValue(assignedEdges, assignedEdgesType)); - } catch (IOException e) { - log.warn("Unable to parse assigned edges!", e); - } - } return ruleChain; } } 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 ea23781e62..a86d10326e 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 @@ -19,6 +19,7 @@ import com.google.common.base.Function; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; +import jnr.ffi.annotations.In; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; @@ -49,6 +50,7 @@ import org.thingsboard.server.dao.edge.EdgeDao; import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.exception.DataValidationException; +import org.thingsboard.server.dao.relation.RelationDao; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.PaginatedRemover; import org.thingsboard.server.dao.service.TimePaginatedRemover; @@ -62,6 +64,9 @@ import java.util.List; import java.util.Map; import java.util.concurrent.ExecutionException; +import static org.thingsboard.server.dao.service.Validator.validateId; +import static org.thingsboard.server.dao.service.Validator.validateString; + /** * Created by igor on 3/12/18. */ @@ -69,6 +74,8 @@ import java.util.concurrent.ExecutionException; @Slf4j public class BaseRuleChainService extends AbstractEntityService implements RuleChainService { + public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; + @Autowired private RuleChainDao ruleChainDao; @@ -376,11 +383,18 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC if (ruleChain.isRoot()) { throw new DataValidationException("Deletion of Root Tenant Rule Chain is prohibited!"); } - if (ruleChain.getAssignedEdges() != null && !ruleChain.getAssignedEdges().isEmpty()) { - for (ShortEdgeInfo assignedEdge : ruleChain.getAssignedEdges()) { - if (assignedEdge.getRootRuleChainId() != null && assignedEdge.getRootRuleChainId().equals(ruleChainId)) { - throw new DataValidationException("Can't delete rule chain that is root for edge [" + assignedEdge.getTitle() + "]. Please assign another root rule chain first to the edge!"); + if (RuleChainType.EDGE.equals(ruleChain.getType())) { + try { + TimePageData edges = edgeService.findEdgesByTenantIdAndRuleChainId(tenantId, ruleChainId, new TimePageLink(Integer.MAX_VALUE)).get(); + if (edges != null && edges.getData() != null && !edges.getData().isEmpty()) { + for (Edge edge : edges.getData()) { + if (edge.getRootRuleChainId() != null && edge.getRootRuleChainId().equals(ruleChainId)) { + throw new DataValidationException("Can't delete rule chain that is root for edge [" + edge.getName() + "]. Please assign another root rule chain first to the edge!"); + } + } } + } catch (InterruptedException | ExecutionException e) { + log.error("Can't get edges by tenant id [{}] and rule chain id [{}]", tenantId.getId(), ruleChainId.getId(), e); } } } @@ -403,14 +417,11 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC if (!edge.getTenantId().getId().equals(ruleChain.getTenantId().getId())) { throw new DataValidationException("Can't assign ruleChain to edge from different tenant!"); } - if (ruleChain.addAssignedEdge(edge)) { - try { - createRelation(tenantId, new EntityRelation(edgeId, ruleChainId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE)); - } catch (ExecutionException | InterruptedException e) { - log.warn("[{}] Failed to create ruleChain relation. Edge Id: [{}]", ruleChainId, edgeId); - throw new RuntimeException(e); - } - ruleChain = saveRuleChain(ruleChain); + try { + createRelation(tenantId, new EntityRelation(edgeId, ruleChainId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE)); + } catch (ExecutionException | InterruptedException e) { + log.warn("[{}] Failed to create ruleChain relation. Edge Id: [{}]", ruleChainId, edgeId); + throw new RuntimeException(e); } return ruleChain; } @@ -425,17 +436,13 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC if (!remove && edge.getRootRuleChainId() != null && edge.getRootRuleChainId().equals(ruleChainId)) { throw new DataValidationException("Can't unassign root rule chain from edge [" + edge.getName() + "]. Please assign another root rule chain first!"); } - if (ruleChain.removeAssignedEdge(edge)) { - try { - deleteRelation(tenantId, new EntityRelation(edgeId, ruleChainId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE)); - } catch (ExecutionException | InterruptedException e) { - log.warn("[{}] Failed to delete rule chain relation. Edge Id: [{}]", ruleChainId, edgeId); - throw new RuntimeException(e); - } - return saveRuleChain(ruleChain); - } else { - return ruleChain; + try { + deleteRelation(tenantId, new EntityRelation(edgeId, ruleChainId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE)); + } catch (ExecutionException | InterruptedException e) { + log.warn("[{}] Failed to delete rule chain relation. Edge Id: [{}]", ruleChainId, edgeId); + throw new RuntimeException(e); } + return ruleChain; } @Override @@ -449,23 +456,11 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC new EdgeRuleChainsUnassigner(edge).removeEntities(tenantId, edge); } - @Override - public void updateEdgeRuleChains(TenantId tenantId, EdgeId edgeId) { - log.trace("Executing updateEdgeRuleChains, edgeId [{}]", edgeId); - Validator.validateId(edgeId, "Incorrect edgeId " + edgeId); - Edge edge = edgeService.findEdgeById(tenantId, edgeId); - if (edge == null) { - throw new DataValidationException("Can't update ruleChains for non-existent edge!"); - } - new EdgeRuleChainsUpdater(edge).removeEntities(tenantId, edge); - } - - @Override public ListenableFuture> findRuleChainsByTenantIdAndEdgeId(TenantId tenantId, EdgeId edgeId, TimePageLink pageLink) { log.trace("Executing findRuleChainsByTenantIdAndEdgeId, tenantId [{}], edgeId [{}], pageLink [{}]", tenantId, edgeId, pageLink); Validator.validateId(tenantId, "Incorrect tenantId " + tenantId); - Validator.validateId(edgeId, "Incorrect customerId " + edgeId); + Validator.validateId(edgeId, "Incorrect edgeId " + edgeId); Validator.validatePageLink(pageLink, "Incorrect page link " + pageLink); ListenableFuture> ruleChains = ruleChainDao.findRuleChainsByTenantIdAndEdgeId(tenantId.getId(), edgeId.getId(), pageLink); @@ -499,13 +494,45 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC ruleChainDao.save(tenantId, ruleChain); return true; } catch (ExecutionException | InterruptedException e) { - log.warn("[{}] Failed to set default root edge rule chain, ruleChainId: [{}]", ruleChainId); + log.warn("Failed to set default root edge rule chain, ruleChainId: [{}]", ruleChainId, e); throw new RuntimeException(e); } } return false; } + @Override + public boolean addDefaultEdgeRuleChain(TenantId tenantId, RuleChainId ruleChainId) { + try { + createRelation(tenantId, new EntityRelation(tenantId, ruleChainId, + EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE)); + return true; + } catch (ExecutionException | InterruptedException e) { + log.warn("Failed to add default edge rule chain, ruleChainId: [{}]", ruleChainId, e); + throw new RuntimeException(e); + } + } + + @Override + public boolean removeDefaultEdgeRuleChain(TenantId tenantId, RuleChainId ruleChainId) { + try { + deleteRelation(tenantId, new EntityRelation(tenantId, ruleChainId, + EntityRelation.CONTAINS_TYPE, RelationTypeGroup.RULE_CHAIN)); + return true; + } catch (ExecutionException | InterruptedException e) { + log.warn("Failed to remove default edge rule chain, ruleChainId: [{}]", ruleChainId, e); + throw new RuntimeException(e); + } + } + + @Override + public ListenableFuture> findDefaultEdgeRuleChainsByTenantId(TenantId tenantId) { + log.trace("Executing findDefaultEdgeRuleChainsByTenantId, tenantId [{}]", tenantId); + validateId(tenantId, INCORRECT_TENANT_ID + tenantId); + return ruleChainDao.findDefaultEdgeRuleChainsByTenantId(tenantId.getId()); + } + + private void checkRuleNodesAndDelete(TenantId tenantId, RuleChainId ruleChainId) { List nodeRelations = getRuleChainToNodeRelations(tenantId, ruleChainId); for (EntityRelation relation : nodeRelations) { @@ -538,15 +565,6 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC relationService.deleteRelation(tenantId, relation); } - private RuleChain updateAssignedEdge(TenantId tenantId, RuleChainId ruleChainId, Edge edge) { - RuleChain ruleChain = findRuleChainById(tenantId, ruleChainId); - if (ruleChain.updateAssignedEdge(edge)) { - return saveRuleChain(ruleChain); - } else { - return ruleChain; - } - } - private DataValidator ruleChainValidator = new DataValidator() { @Override @@ -616,29 +634,4 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC unassignRuleChainFromEdge(edge.getTenantId(), new RuleChainId(entity.getUuidId()), this.edge.getId(), true); } } - - private class EdgeRuleChainsUpdater extends TimePaginatedRemover { - - private Edge edge; - - EdgeRuleChainsUpdater(Edge edge) { - this.edge = edge; - } - - @Override - protected List findEntities(TenantId tenantId, Edge edge, TimePageLink pageLink) { - try { - return ruleChainDao.findRuleChainsByTenantIdAndEdgeId(edge.getTenantId().getId(), edge.getId().getId(), pageLink).get(); - } catch (InterruptedException | ExecutionException e) { - log.warn("Failed to get ruleChains by tenantId [{}] and edgeId [{}].", edge.getTenantId().getId(), edge.getId().getId()); - throw new RuntimeException(e); - } - } - - @Override - protected void removeEntity(TenantId tenantId, RuleChain entity) { - updateAssignedEdge(edge.getTenantId(), new RuleChainId(entity.getUuidId()), this.edge); - } - - } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/CassandraRuleChainDao.java b/dao/src/main/java/org/thingsboard/server/dao/rule/CassandraRuleChainDao.java index 0a4cec2bf3..8f24b737f2 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/CassandraRuleChainDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/CassandraRuleChainDao.java @@ -22,7 +22,9 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.page.TimePageLink; @@ -104,4 +106,17 @@ public class CassandraRuleChainDao extends CassandraAbstractSearchTextDao> findDefaultEdgeRuleChainsByTenantId(UUID tenantId) { + log.debug("Try to find default edge rule chains by tenantId [{}]", tenantId); + ListenableFuture> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new TenantId(tenantId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); + return Futures.transformAsync(relations, input -> { + List> ruleChainFutures = new ArrayList<>(input.size()); + for (EntityRelation relation : input) { + ruleChainFutures.add(findByIdAsync(new TenantId(tenantId), relation.getTo().getId())); + } + return Futures.successfulAsList(ruleChainFutures); + }, MoreExecutors.directExecutor()); + } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/RuleChainDao.java b/dao/src/main/java/org/thingsboard/server/dao/rule/RuleChainDao.java index b2aa92883c..b883f13fa8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/RuleChainDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/RuleChainDao.java @@ -58,4 +58,12 @@ public interface RuleChainDao extends Dao { * @return the list of rule chain objects */ ListenableFuture> findRuleChainsByTenantIdAndEdgeId(UUID tenantId, UUID edgeId, TimePageLink pageLink); + + /** + * Find default edge rule chains by tenantId. + * + * @param tenantId the tenantId + * @return the list of rule chain objects + */ + ListenableFuture> findDefaultEdgeRuleChainsByTenantId(UUID tenantId); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java index 09d2ba687e..1bdaab82c1 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java @@ -15,7 +15,10 @@ */ package org.thingsboard.server.dao.sql.edge; +import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; +import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.data.domain.PageRequest; import org.springframework.data.repository.CrudRepository; @@ -24,11 +27,18 @@ import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.UUIDConverter; import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageLink; +import org.thingsboard.server.common.data.page.TimePageLink; +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.dao.DaoUtil; import org.thingsboard.server.dao.edge.EdgeDao; import org.thingsboard.server.dao.model.sql.EdgeEntity; +import org.thingsboard.server.dao.relation.RelationDao; import org.thingsboard.server.dao.sql.JpaAbstractSearchTextDao; import org.thingsboard.server.dao.util.SqlDao; @@ -45,11 +55,15 @@ import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID_STR; @Component @SqlDao +@Slf4j public class JpaEdgeDao extends JpaAbstractSearchTextDao implements EdgeDao { @Autowired private EdgeRepository edgeRepository; + @Autowired + private RelationDao relationDao; + @Override protected Class getEntityClass() { return EdgeEntity.class; @@ -132,6 +146,19 @@ public class JpaEdgeDao extends JpaAbstractSearchTextDao imple return Optional.ofNullable(edge); } + @Override + public ListenableFuture> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId, TimePageLink pageLink) { + log.debug("Try to find edges by tenantId [{}], ruleChainId [{}] and pageLink [{}]", tenantId, ruleChainId, pageLink); + ListenableFuture> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new RuleChainId(ruleChainId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); + return Futures.transformAsync(relations, input -> { + List> edgeFutures = new ArrayList<>(input.size()); + for (EntityRelation relation : input) { + edgeFutures.add(findByIdAsync(new TenantId(tenantId), relation.getTo().getId())); + } + return Futures.successfulAsList(edgeFutures); + }, MoreExecutors.directExecutor()); + } + private List convertTenantEdgeTypesToDto(UUID tenantId, List types) { List list = Collections.emptyList(); if (types != null && !types.isEmpty()) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleChainDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleChainDao.java index cac8361df8..c8ca814fe8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleChainDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleChainDao.java @@ -25,7 +25,9 @@ import org.springframework.data.repository.CrudRepository; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.UUIDConverter; +import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.page.TimePageLink; @@ -103,4 +105,17 @@ public class JpaRuleChainDao extends JpaAbstractSearchTextDao> findDefaultEdgeRuleChainsByTenantId(UUID tenantId) { + log.debug("Try to find default edge rule chains by tenantId [{}]", tenantId); + ListenableFuture> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new TenantId(tenantId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); + return Futures.transformAsync(relations, input -> { + List> ruleChainsFutures = new ArrayList<>(input.size()); + for (EntityRelation relation : input) { + ruleChainsFutures.add(findByIdAsync(new TenantId(tenantId), relation.getTo().getId())); + } + return Futures.successfulAsList(ruleChainsFutures); + }, MoreExecutors.directExecutor()); + } } diff --git a/dao/src/main/resources/sql/schema-entities-hsql.sql b/dao/src/main/resources/sql/schema-entities-hsql.sql index 6c40ea439f..3930f428f0 100644 --- a/dao/src/main/resources/sql/schema-entities-hsql.sql +++ b/dao/src/main/resources/sql/schema-entities-hsql.sql @@ -111,7 +111,6 @@ CREATE TABLE IF NOT EXISTS dashboard ( id varchar(31) NOT NULL CONSTRAINT dashboard_pkey PRIMARY KEY, configuration varchar(10000000), assigned_customers varchar(1000000), - assigned_edges varchar(10000000), search_text varchar(255), tenant_id varchar(31), title varchar(255) @@ -228,8 +227,7 @@ CREATE TABLE IF NOT EXISTS rule_chain ( root boolean, debug_mode boolean, search_text varchar(255), - tenant_id varchar(31), - assigned_edges varchar(10000000) + tenant_id varchar(31) ); CREATE TABLE IF NOT EXISTS rule_node ( diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index dfe7079864..10746084ae 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -111,7 +111,6 @@ CREATE TABLE IF NOT EXISTS dashboard ( id varchar(31) NOT NULL CONSTRAINT dashboard_pkey PRIMARY KEY, configuration varchar(10000000), assigned_customers varchar(1000000), - assigned_edges varchar(10000000), search_text varchar(255), tenant_id varchar(31), title varchar(255) @@ -228,8 +227,7 @@ CREATE TABLE IF NOT EXISTS rule_chain ( root boolean, debug_mode boolean, search_text varchar(255), - tenant_id varchar(31), - assigned_edges varchar(10000000) + tenant_id varchar(31) ); CREATE TABLE IF NOT EXISTS rule_node ( diff --git a/ui/src/app/api/rule-chain.service.js b/ui/src/app/api/rule-chain.service.js index 41a526c97f..78c891ba06 100644 --- a/ui/src/app/api/rule-chain.service.js +++ b/ui/src/app/api/rule-chain.service.js @@ -36,14 +36,14 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co resolveTargetRuleChains: resolveTargetRuleChains, testScript: testScript, getLatestRuleNodeDebugInput: getLatestRuleNodeDebugInput, - updateRuleChainEdges: updateRuleChainEdges, - addRuleChainEdges: addRuleChainEdges, - removeRuleChainEdges: removeRuleChainEdges, getEdgeRuleChains: getEdgeRuleChains, getEdgesRuleChains: getEdgesRuleChains, assignRuleChainToEdge: assignRuleChainToEdge, unassignRuleChainFromEdge: unassignRuleChainFromEdge, - setDefaultRootEdgeRuleChain: setDefaultRootEdgeRuleChain + setDefaultRootEdgeRuleChain: setDefaultRootEdgeRuleChain, + addDefaultEdgeRuleChain: addDefaultEdgeRuleChain, + removeDefaultEdgeRuleChain: removeDefaultEdgeRuleChain, + getDefaultEdgeRuleChains: getDefaultEdgeRuleChains }; return service; @@ -64,7 +64,7 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co url += '&type=' + type; } $http.get(url, config).then(function success(response) { - deferred.resolve(prepareRuleChains(response.data)); + deferred.resolve(response.data); }, function fail() { deferred.reject(); }); @@ -75,7 +75,7 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co var deferred = $q.defer(); var url = '/api/ruleChain/' + ruleChainId; $http.get(url, config).then(function success(response) { - deferred.resolve(prepareRuleChain(response.data)); + deferred.resolve(response.data); }, function fail() { deferred.reject(); }); @@ -85,8 +85,8 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co function saveRuleChain(ruleChain) { var deferred = $q.defer(); var url = '/api/ruleChain'; - $http.post(url, cleanRuleChain(ruleChain)).then(function success(response) { - deferred.resolve(prepareRuleChain(response.data)); + $http.post(url, ruleChain).then(function success(response) { + deferred.resolve(response.data); }, function fail() { deferred.reject(); }); @@ -273,7 +273,7 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co var deferred = $q.defer(); getRuleChain(ruleChainId, {ignoreErrors: true}).then( (ruleChain) => { - deferred.resolve(prepareRuleChain(ruleChain)); + deferred.resolve(ruleChain); }, () => { deferred.resolve({ @@ -310,39 +310,6 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co return deferred.promise; } - function updateRuleChainEdges(ruleChainId, edgeIds) { - var deferred = $q.defer(); - var url = '/api/ruleChain/' + ruleChainId + '/edges'; - $http.post(url, edgeIds).then(function success(response) { - deferred.resolve(prepareRuleChain(response.data)); - }, function fail() { - deferred.reject(); - }); - return deferred.promise; - } - - function addRuleChainEdges(ruleChainId, edgeIds) { - var deferred = $q.defer(); - var url = '/api/ruleChain/' + ruleChainId + '/edges/add'; - $http.post(url, edgeIds).then(function success(response) { - deferred.resolve(prepareRuleChain(response.data)); - }, function fail() { - deferred.reject(); - }); - return deferred.promise; - } - - function removeRuleChainEdges(ruleChainId, edgeIds) { - var deferred = $q.defer(); - var url = '/api/ruleChain/' + ruleChainId + '/edges/remove'; - $http.post(url, edgeIds).then(function success(response) { - deferred.resolve(prepareRuleChain(response.data)); - }, function fail() { - deferred.reject(); - }); - return deferred.promise; - } - function getEdgesRuleChains(pageLink, config) { return getRuleChains(pageLink, config, types.edgeRuleChainType); } @@ -354,7 +321,6 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co url += '&offset=' + pageLink.idOffset; } $http.get(url, config).then(function success(response) { - response.data = prepareRuleChains(response.data); if (pageLink.textSearch) { response.data.data = $filter('filter')(response.data.data, {title: pageLink.textSearch}); } @@ -369,7 +335,7 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co var deferred = $q.defer(); var url = '/api/edge/' + edgeId + '/ruleChain/' + ruleChainId; $http.post(url, null).then(function success(response) { - deferred.resolve(prepareRuleChain(response.data)); + deferred.resolve(response.data); }, function fail() { deferred.reject(); }); @@ -380,7 +346,7 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co var deferred = $q.defer(); var url = '/api/edge/' + edgeId + '/ruleChain/' + ruleChainId; $http.delete(url).then(function success(response) { - deferred.resolve(prepareRuleChain(response.data)); + deferred.resolve(response.data); }, function fail() { deferred.reject(); }); @@ -398,35 +364,36 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co return deferred.promise; } - function prepareRuleChains(ruleChainsData) { - if (ruleChainsData.data) { - for (var i = 0; i < ruleChainsData.data.length; i++) { - ruleChainsData.data[i] = prepareRuleChain(ruleChainsData.data[i]); - } - } - return ruleChainsData; + function addDefaultEdgeRuleChain(ruleChainId) { + var deferred = $q.defer(); + var url = '/api/ruleChain/' + ruleChainId + '/defaultEdge'; + $http.post(url, null).then(function success(response) { + deferred.resolve(response.data); + }, function fail() { + deferred.reject(); + }); + return deferred.promise; } - function prepareRuleChain(ruleChain) { - ruleChain.assignedEdgesText = ""; - ruleChain.assignedEdgesIds = []; - - if (ruleChain.assignedEdges && ruleChain.assignedEdges.length) { - var assignedEdgesTitles = []; - for (var j = 0; j < ruleChain.assignedEdges.length; j++) { - var assignedEdge = ruleChain.assignedEdges[j]; - ruleChain.assignedEdgesIds.push(assignedEdge.edgeId.id); - assignedEdgesTitles.push(assignedEdge.title); - } - ruleChain.assignedEdgesText = assignedEdgesTitles.join(', '); - } - - return ruleChain; + function removeDefaultEdgeRuleChain(ruleChainId) { + var deferred = $q.defer(); + var url = '/api/ruleChain/' + ruleChainId + '/defaultEdge'; + $http.delete(url).then(function success(response) { + deferred.resolve(response.data); + }, function fail() { + deferred.reject(); + }); + return deferred.promise; } - function cleanRuleChain(ruleChain) { - delete ruleChain.assignedEdgesText; - delete ruleChain.assignedEdgesIds; - return ruleChain; + function getDefaultEdgeRuleChains(config) { + var deferred = $q.defer(); + var url = '/api/ruleChain/defaultEdgeRuleChains'; + $http.get(url, config).then(function success(response) { + deferred.resolve(response.data); + }, function fail() { + deferred.reject(); + }); + return deferred.promise; } } diff --git a/ui/src/app/rulechain/rulechains.controller.js b/ui/src/app/rulechain/rulechains.controller.js index f3e159566c..ff1d393618 100644 --- a/ui/src/app/rulechain/rulechains.controller.js +++ b/ui/src/app/rulechain/rulechains.controller.js @@ -115,7 +115,7 @@ export default function RuleChainsController(ruleChainService, userService, edge if (vm.ruleChainsScope === 'tenant') { fetchRuleChainsFunction = function (pageLink) { - return fetchRuleChains(pageLink, 'SYSTEM'); + return fetchRuleChains(pageLink, types.systemRuleChainType); }; deleteRuleChainFunction = function (ruleChainId) { return deleteRuleChain(ruleChainId); @@ -176,7 +176,7 @@ export default function RuleChainsController(ruleChainService, userService, edge } else if (vm.ruleChainsScope === 'edges') { fetchRuleChainsFunction = function (pageLink) { - return fetchRuleChains(pageLink, 'EDGE'); + return fetchRuleChains(pageLink, types.edgeRuleChainType); }; deleteRuleChainFunction = function (ruleChainId) { return deleteRuleChain(ruleChainId); @@ -316,7 +316,28 @@ export default function RuleChainsController(ruleChainService, userService, edge } function fetchRuleChains(pageLink, type) { - return ruleChainService.getRuleChains(pageLink, null, type); + if (type === types.systemRuleChainType) { + return ruleChainService.getRuleChains(pageLink, null, type); + } else { + var deferred = $q.defer(); + ruleChainService.getRuleChains(pageLink, null, type).then( + // TODO: deaflynx + // function success(response) { + // ruleChainService.getDefaultEdgeRuleChains().then( + // function success(defaultEdgeRuleChains) { + // // response.data merge with defaultEdgeRuleChains + // deferred.resolve(data); + // }, + // function fail() { + // deferred.reject(); + // } + // ); + // }, + function fail() { + deferred.reject(); + }); + return deferred.promise; + } } function saveRuleChain(ruleChain) {