Browse Source

Added default rule chain for server side

pull/2436/head
Volodymyr Babak 6 years ago
parent
commit
6ab3714ac5
  1. 4
      application/src/main/java/org/thingsboard/server/controller/EdgeController.java
  2. 202
      application/src/main/java/org/thingsboard/server/controller/RuleChainController.java
  3. 8
      application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java
  4. 6
      application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java
  5. 6
      common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java
  6. 9
      common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java
  7. 27
      common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChain.java
  8. 29
      dao/src/main/java/org/thingsboard/server/dao/edge/CassandraEdgeDao.java
  9. 11
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java
  10. 39
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java
  11. 1
      dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
  12. 30
      dao/src/main/java/org/thingsboard/server/dao/model/nosql/RuleChainEntity.java
  13. 28
      dao/src/main/java/org/thingsboard/server/dao/model/sql/RuleChainEntity.java
  14. 133
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java
  15. 15
      dao/src/main/java/org/thingsboard/server/dao/rule/CassandraRuleChainDao.java
  16. 8
      dao/src/main/java/org/thingsboard/server/dao/rule/RuleChainDao.java
  17. 27
      dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java
  18. 15
      dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleChainDao.java
  19. 4
      dao/src/main/resources/sql/schema-entities-hsql.sql
  20. 4
      dao/src/main/resources/sql/schema-entities.sql
  21. 109
      ui/src/app/api/rule-chain.service.js
  22. 27
      ui/src/app/rulechain/rulechains.controller.js

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

202
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<RuleChain> 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<EdgeId> edgeIds = new HashSet<>();
if (strEdgeIds != null) {
for (String strEdgeId : strEdgeIds) {
edgeIds.add(new EdgeId(toUUID(strEdgeId)));
}
}
Set<EdgeId> addedEdgeIds = new HashSet<>();
Set<EdgeId> removedEdgeIds = new HashSet<>();
for (EdgeId edgeId : edgeIds) {
if (!ruleChain.isAssignedToEdge(edgeId)) {
addedEdgeIds.add(edgeId);
}
}
Set<ShortEdgeInfo> 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<EdgeId> 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<EdgeId> 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<RuleChain> 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<RuleChain> getDefaultEdgeRuleChains() throws ThingsboardException {
try {
TenantId tenantId = getCurrentUser().getTenantId();
return checkNotNull(ruleChainService.findDefaultEdgeRuleChainsByTenantId(tenantId)).get();
} catch (Exception e) {
throw handleException(e);
}
}
}

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

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

6
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<Event> 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<TimePageData<Edge>> findEdgesByTenantIdAndRuleChainId(TenantId tenantId, RuleChainId ruleChainId, TimePageLink pageLink);
}

9
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<TimePageData<RuleChain>> 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<List<RuleChain>> findDefaultEdgeRuleChainsByTenantId(TenantId tenantId);
}

27
common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChain.java

@ -48,7 +48,6 @@ public class RuleChain extends SearchTextBasedWithAdditionalInfo<RuleChainId> im
private boolean root;
private boolean debugMode;
private transient JsonNode configuration;
private Set<ShortEdgeInfo> assignedEdges;
@JsonIgnore
private byte[] configurationBytes;
@ -68,7 +67,6 @@ public class RuleChain extends SearchTextBasedWithAdditionalInfo<RuleChainId> 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<RuleChainId> 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());
}
}

29
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<EdgeEntity, Edge> implements EdgeDao {
@Autowired
private RelationDao relationDao;
@Override
protected Class<EdgeEntity> getColumnFamilyClass() {
return EdgeEntity.class;
@ -91,4 +107,17 @@ public class CassandraEdgeDao extends CassandraAbstractSearchTextDao<EdgeEntity,
public Optional<Edge> findByRoutingKey(UUID tenantId, String routingKey) {
return Optional.empty();
}
@Override
public ListenableFuture<List<Edge>> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId, TimePageLink pageLink) {
log.debug("Try to find edges by tenantId [{}], ruleChainId [{}] and pageLink [{}]", tenantId, ruleChainId, pageLink);
ListenableFuture<List<EntityRelation>> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new RuleChainId(ruleChainId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE);
return Futures.transformAsync(relations, input -> {
List<ListenableFuture<Edge>> edgeFutures = new ArrayList<>(input.size());
for (EntityRelation relation : input) {
edgeFutures.add(findByIdAsync(new TenantId(tenantId), relation.getTo().getId()));
}
return Futures.successfulAsList(edgeFutures);
}, MoreExecutors.directExecutor());
}
}

11
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<Edge> {
*/
Optional<Edge> 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<List<Edge>> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId, TimePageLink pageLink);
}

39
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<TimePageData<Edge>> 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<Void>() {
@Override
@ -647,6 +657,23 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
return savedEdge;
}
@Override
public ListenableFuture<TimePageData<Edge>> 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<List<Edge>> edges = edgeDao.findEdgesByTenantIdAndRuleChainId(tenantId.getId(), ruleChainId.getId(), pageLink);
return Futures.transform(edges, new Function<List<Edge>, TimePageData<Edge>>() {
@Nullable
@Override
public TimePageData<Edge> apply(@Nullable List<Edge> edges) {
return new TimePageData<>(edges, pageLink);
}
}, MoreExecutors.directExecutor());
}
private DataValidator<Edge> edgeValidator =
new DataValidator<Edge>() {

1
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";

30
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<RuleChain> {
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<RuleChain> {
@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<RuleChain> {
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> {
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;
}

28
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<RuleChain> implements SearchTextEntity<RuleChain> {
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<RuleChain> 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<RuleChain> 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<RuleChain> 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;
}
}

133
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<Edge> 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<TimePageData<RuleChain>> 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<List<RuleChain>> 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<List<RuleChain>> 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<EntityRelation> 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<RuleChain> ruleChainValidator =
new DataValidator<RuleChain>() {
@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<Edge, RuleChain> {
private Edge edge;
EdgeRuleChainsUpdater(Edge edge) {
this.edge = edge;
}
@Override
protected List<RuleChain> 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);
}
}
}

15
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<RuleCh
}, MoreExecutors.directExecutor());
}
@Override
public ListenableFuture<List<RuleChain>> findDefaultEdgeRuleChainsByTenantId(UUID tenantId) {
log.debug("Try to find default edge rule chains by tenantId [{}]", tenantId);
ListenableFuture<List<EntityRelation>> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new TenantId(tenantId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE);
return Futures.transformAsync(relations, input -> {
List<ListenableFuture<RuleChain>> ruleChainFutures = new ArrayList<>(input.size());
for (EntityRelation relation : input) {
ruleChainFutures.add(findByIdAsync(new TenantId(tenantId), relation.getTo().getId()));
}
return Futures.successfulAsList(ruleChainFutures);
}, MoreExecutors.directExecutor());
}
}

8
dao/src/main/java/org/thingsboard/server/dao/rule/RuleChainDao.java

@ -58,4 +58,12 @@ public interface RuleChainDao extends Dao<RuleChain> {
* @return the list of rule chain objects
*/
ListenableFuture<List<RuleChain>> 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<List<RuleChain>> findDefaultEdgeRuleChainsByTenantId(UUID tenantId);
}

27
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<EdgeEntity, Edge> implements EdgeDao {
@Autowired
private EdgeRepository edgeRepository;
@Autowired
private RelationDao relationDao;
@Override
protected Class<EdgeEntity> getEntityClass() {
return EdgeEntity.class;
@ -132,6 +146,19 @@ public class JpaEdgeDao extends JpaAbstractSearchTextDao<EdgeEntity, Edge> imple
return Optional.ofNullable(edge);
}
@Override
public ListenableFuture<List<Edge>> findEdgesByTenantIdAndRuleChainId(UUID tenantId, UUID ruleChainId, TimePageLink pageLink) {
log.debug("Try to find edges by tenantId [{}], ruleChainId [{}] and pageLink [{}]", tenantId, ruleChainId, pageLink);
ListenableFuture<List<EntityRelation>> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new RuleChainId(ruleChainId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE);
return Futures.transformAsync(relations, input -> {
List<ListenableFuture<Edge>> 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<EntitySubtype> convertTenantEdgeTypesToDto(UUID tenantId, List<String> types) {
List<EntitySubtype> list = Collections.emptyList();
if (types != null && !types.isEmpty()) {

15
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<RuleChainEntity, R
return Futures.successfulAsList(ruleChainFutures);
}, MoreExecutors.directExecutor());
}
@Override
public ListenableFuture<List<RuleChain>> findDefaultEdgeRuleChainsByTenantId(UUID tenantId) {
log.debug("Try to find default edge rule chains by tenantId [{}]", tenantId);
ListenableFuture<List<EntityRelation>> relations = relationDao.findAllByToAndType(new TenantId(tenantId), new TenantId(tenantId), EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE);
return Futures.transformAsync(relations, input -> {
List<ListenableFuture<RuleChain>> ruleChainsFutures = new ArrayList<>(input.size());
for (EntityRelation relation : input) {
ruleChainsFutures.add(findByIdAsync(new TenantId(tenantId), relation.getTo().getId()));
}
return Futures.successfulAsList(ruleChainsFutures);
}, MoreExecutors.directExecutor());
}
}

4
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 (

4
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 (

109
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;
}
}

27
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) {

Loading…
Cancel
Save