|
|
|
@ -50,9 +50,8 @@ import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
import org.thingsboard.server.common.data.id.UserId; |
|
|
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|
|
|
import org.thingsboard.server.common.data.msg.TbMsgType; |
|
|
|
import org.thingsboard.server.common.data.page.PageData; |
|
|
|
import org.thingsboard.server.common.data.page.PageDataIterable; |
|
|
|
import org.thingsboard.server.common.data.page.PageDataIterableByTenantIdEntityId; |
|
|
|
import org.thingsboard.server.common.data.page.PageLink; |
|
|
|
import org.thingsboard.server.common.data.relation.EntityRelation; |
|
|
|
import org.thingsboard.server.common.data.relation.RelationTypeGroup; |
|
|
|
import org.thingsboard.server.common.data.rule.RuleChain; |
|
|
|
@ -405,15 +404,10 @@ public abstract class BaseEdgeProcessor { |
|
|
|
JsonNode body, EdgeId sourceEdgeId) { |
|
|
|
List<ListenableFuture<Void>> futures = new ArrayList<>(); |
|
|
|
if (TenantId.SYS_TENANT_ID.equals(tenantId)) { |
|
|
|
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); |
|
|
|
PageData<TenantId> tenantsIds; |
|
|
|
do { |
|
|
|
tenantsIds = tenantService.findTenantsIds(pageLink); |
|
|
|
for (TenantId tenantId1 : tenantsIds.getData()) { |
|
|
|
futures.addAll(processActionForAllEdgesByTenantId(tenantId1, type, actionType, entityId, body, sourceEdgeId)); |
|
|
|
} |
|
|
|
pageLink = pageLink.nextPageLink(); |
|
|
|
} while (tenantsIds.hasNext()); |
|
|
|
PageDataIterable<TenantId> tenantIds = new PageDataIterable<>(link -> tenantService.findTenantsIds(link), 1024); |
|
|
|
for (TenantId tenantId1 : tenantIds) { |
|
|
|
futures.addAll(processActionForAllEdgesByTenantId(tenantId1, type, actionType, entityId, body, sourceEdgeId)); |
|
|
|
} |
|
|
|
} else { |
|
|
|
futures = processActionForAllEdgesByTenantId(tenantId, type, actionType, entityId, null, sourceEdgeId); |
|
|
|
} |
|
|
|
@ -426,22 +420,13 @@ public abstract class BaseEdgeProcessor { |
|
|
|
EntityId entityId, |
|
|
|
JsonNode body, |
|
|
|
EdgeId sourceEdgeId) { |
|
|
|
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); |
|
|
|
PageData<Edge> pageData; |
|
|
|
List<ListenableFuture<Void>> futures = new ArrayList<>(); |
|
|
|
do { |
|
|
|
pageData = edgeService.findEdgesByTenantId(tenantId, pageLink); |
|
|
|
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { |
|
|
|
for (Edge edge : pageData.getData()) { |
|
|
|
if (!edge.getId().equals(sourceEdgeId)) { |
|
|
|
futures.add(saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, body)); |
|
|
|
} |
|
|
|
} |
|
|
|
if (pageData.hasNext()) { |
|
|
|
pageLink = pageLink.nextPageLink(); |
|
|
|
} |
|
|
|
PageDataIterable<Edge> edges = new PageDataIterable<>(link -> edgeService.findEdgesByTenantId(tenantId, link), 1024); |
|
|
|
for (Edge edge : edges) { |
|
|
|
if (!edge.getId().equals(sourceEdgeId)) { |
|
|
|
futures.add(saveEdgeEvent(tenantId, edge.getId(), type, actionType, entityId, body)); |
|
|
|
} |
|
|
|
} while (pageData != null && pageData.hasNext()); |
|
|
|
} |
|
|
|
return futures; |
|
|
|
} |
|
|
|
|
|
|
|
@ -548,35 +533,24 @@ public abstract class BaseEdgeProcessor { |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<Void> updateDependentRuleChains(TenantId tenantId, RuleChainId processingRuleChainId, EdgeId edgeId) { |
|
|
|
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); |
|
|
|
PageData<RuleChain> pageData; |
|
|
|
List<ListenableFuture<Void>> futures = new ArrayList<>(); |
|
|
|
do { |
|
|
|
pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink); |
|
|
|
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { |
|
|
|
for (RuleChain ruleChain : pageData.getData()) { |
|
|
|
if (!ruleChain.getId().equals(processingRuleChainId)) { |
|
|
|
List<RuleChainConnectionInfo> connectionInfos = |
|
|
|
ruleChainService.loadRuleChainMetaData(ruleChain.getTenantId(), ruleChain.getId()).getRuleChainConnections(); |
|
|
|
if (connectionInfos != null && !connectionInfos.isEmpty()) { |
|
|
|
for (RuleChainConnectionInfo connectionInfo : connectionInfos) { |
|
|
|
if (connectionInfo.getTargetRuleChainId().equals(processingRuleChainId)) { |
|
|
|
futures.add(saveEdgeEvent(tenantId, |
|
|
|
edgeId, |
|
|
|
EdgeEventType.RULE_CHAIN_METADATA, |
|
|
|
EdgeEventActionType.UPDATED, |
|
|
|
ruleChain.getId(), |
|
|
|
null)); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
PageDataIterable<RuleChain> ruleChains = new PageDataIterable<>(link -> ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, link), 1024); |
|
|
|
for (RuleChain ruleChain : ruleChains) { |
|
|
|
List<RuleChainConnectionInfo> connectionInfos = |
|
|
|
ruleChainService.loadRuleChainMetaData(ruleChain.getTenantId(), ruleChain.getId()).getRuleChainConnections(); |
|
|
|
if (connectionInfos != null && !connectionInfos.isEmpty()) { |
|
|
|
for (RuleChainConnectionInfo connectionInfo : connectionInfos) { |
|
|
|
if (connectionInfo.getTargetRuleChainId().equals(processingRuleChainId)) { |
|
|
|
futures.add(saveEdgeEvent(tenantId, |
|
|
|
edgeId, |
|
|
|
EdgeEventType.RULE_CHAIN_METADATA, |
|
|
|
EdgeEventActionType.UPDATED, |
|
|
|
ruleChain.getId(), |
|
|
|
null)); |
|
|
|
} |
|
|
|
} |
|
|
|
if (pageData.hasNext()) { |
|
|
|
pageLink = pageLink.nextPageLink(); |
|
|
|
} |
|
|
|
} |
|
|
|
} while (pageData != null && pageData.hasNext()); |
|
|
|
} |
|
|
|
return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); |
|
|
|
} |
|
|
|
|
|
|
|
|