From a7196c4fa7790dd6b00b78037c1e399407284ae9 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 28 Feb 2022 17:15:02 +0200 Subject: [PATCH 1/3] Automatic conversion of rule chain metadata on the back end --- .../thingsboard/server/edge/BaseEdgeTest.java | 2 - ...AbstractRuleEngineFlowIntegrationTest.java | 19 ++++++-- .../common/data/rule/RuleChainMetaData.java | 11 ----- .../server/dao/rule/BaseRuleChainService.java | 43 ++++++++++++++++--- 4 files changed, 53 insertions(+), 22 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java index cd3b14c796..5f7acd0422 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java @@ -612,8 +612,6 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { ruleChainMetaData.addConnectionInfo(0, 2, "fail"); ruleChainMetaData.addConnectionInfo(1, 2, "success"); - ruleChainMetaData.addRuleChainConnectionInfo(2, edge.getRootRuleChainId(), "success", mapper.createObjectNode()); - doPost("/api/ruleChain/metadata", ruleChainMetaData, RuleChainMetaData.class); } diff --git a/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java b/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java index dcc8f35ba3..9d96fe9d9c 100644 --- a/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java @@ -25,6 +25,7 @@ import org.mockito.Mockito; import org.mockito.stubbing.Answer; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.test.util.ReflectionTestUtils; +import org.thingsboard.rule.engine.flow.TbRuleChainInputNodeConfiguration; import org.thingsboard.rule.engine.metadata.TbGetAttributesNodeConfiguration; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.common.data.DataConstants; @@ -35,6 +36,7 @@ import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.rule.NodeConnectionInfo; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.rule.RuleNode; @@ -240,16 +242,27 @@ public abstract class AbstractRuleEngineFlowIntegrationTest extends AbstractRule configuration1.setServerAttributeNames(Collections.singletonList("serverAttributeKey1")); ruleNode1.setConfiguration(mapper.valueToTree(configuration1)); - rootMetaData.setNodes(Collections.singletonList(ruleNode1)); + RuleNode ruleNode12 = new RuleNode(); + ruleNode12.setName("Simple Rule Node 1"); + ruleNode12.setType(org.thingsboard.rule.engine.flow.TbRuleChainInputNode.class.getName()); + ruleNode12.setDebugMode(true); + TbRuleChainInputNodeConfiguration configuration12 = new TbRuleChainInputNodeConfiguration(); + configuration12.setRuleChainId(secondaryRuleChain.getId().getId().toString()); + ruleNode12.setConfiguration(mapper.valueToTree(configuration12)); + + rootMetaData.setNodes(Arrays.asList(ruleNode1, ruleNode12)); rootMetaData.setFirstNodeIndex(0); - rootMetaData.addRuleChainConnectionInfo(0, secondaryRuleChain.getId(), "Success", mapper.createObjectNode()); + NodeConnectionInfo connection = new NodeConnectionInfo(); + connection.setFromIndex(0); + connection.setToIndex(1); + connection.setType("Success"); + rootMetaData.setConnections(Collections.singletonList(connection)); rootMetaData = saveRuleChainMetaData(rootMetaData); Assert.assertNotNull(rootMetaData); rootRuleChain = getRuleChain(rootRuleChain.getId()); Assert.assertNotNull(rootRuleChain.getFirstRuleNodeId()); - RuleChainMetaData secondaryMetaData = new RuleChainMetaData(); secondaryMetaData.setRuleChainId(secondaryRuleChain.getId()); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChainMetaData.java b/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChainMetaData.java index d51725f153..7bd19e2441 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChainMetaData.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleChainMetaData.java @@ -59,15 +59,4 @@ public class RuleChainMetaData { } connections.add(connectionInfo); } - public void addRuleChainConnectionInfo(int fromIndex, RuleChainId targetRuleChainId, String type, JsonNode additionalInfo) { - RuleChainConnectionInfo connectionInfo = new RuleChainConnectionInfo(); - connectionInfo.setFromIndex(fromIndex); - connectionInfo.setTargetRuleChainId(targetRuleChainId); - connectionInfo.setType(type); - connectionInfo.setAdditionalInfo(additionalInfo); - if (ruleChainConnections == null) { - ruleChainConnections = new ArrayList<>(); - } - ruleChainConnections.add(connectionInfo); - } } 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 6916566307..0f62c57741 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 @@ -26,6 +26,7 @@ import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.BaseData; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.edge.Edge; @@ -199,11 +200,42 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC } if (ruleChainMetaData.getRuleChainConnections() != null) { for (RuleChainConnectionInfo nodeToRuleChainConnection : ruleChainMetaData.getRuleChainConnections()) { + RuleChainId targetRuleChainId = nodeToRuleChainConnection.getTargetRuleChainId(); + RuleChain targetRuleChain = findRuleChainById(TenantId.SYS_TENANT_ID, targetRuleChainId); + if (targetRuleChain == null) { + log.info("Skip processing of relation for non existing target rule chain: [{}]", targetRuleChainId); + continue; + } + + RuleNode targetNode = new RuleNode(); + targetNode.setName(targetRuleChain.getName()); + targetNode.setRuleChainId(ruleChain.getId()); + targetNode.setType("org.thingsboard.rule.engine.flow.TbRuleChainInputNode"); + var configuration = JacksonUtil.newObjectNode(); + configuration.put("ruleChainId", targetRuleChain.getId().toString()); + targetNode.setConfiguration(configuration); + ObjectNode layout = (ObjectNode) nodeToRuleChainConnection.getAdditionalInfo(); + layout.remove("description"); + layout.remove("ruleChainNodeId"); + targetNode.setAdditionalInfo(layout); + targetNode.setDebugMode(false); + targetNode = ruleNodeDao.save(tenantId, targetNode); + + EntityRelation sourceRuleChainToRuleNode = new EntityRelation(); + sourceRuleChainToRuleNode.setFrom(ruleChain.getId()); + sourceRuleChainToRuleNode.setTo(targetNode.getId()); + sourceRuleChainToRuleNode.setType(EntityRelation.CONTAINS_TYPE); + sourceRuleChainToRuleNode.setTypeGroup(RelationTypeGroup.RULE_CHAIN); + relationService.saveRelation(tenantId, sourceRuleChainToRuleNode); + + EntityRelation sourceRuleNodeToTargetRuleNode = new EntityRelation(); EntityId from = nodes.get(nodeToRuleChainConnection.getFromIndex()).getId(); - EntityId to = nodeToRuleChainConnection.getTargetRuleChainId(); - String type = nodeToRuleChainConnection.getType(); - createRelation(tenantId, new EntityRelation(from, to, type, RelationTypeGroup.RULE_NODE, nodeToRuleChainConnection.getAdditionalInfo())); - } + sourceRuleNodeToTargetRuleNode.setFrom(from); + sourceRuleNodeToTargetRuleNode.setTo(targetNode.getId()); + sourceRuleNodeToTargetRuleNode.setType(nodeToRuleChainConnection.getType()); + sourceRuleNodeToTargetRuleNode.setTypeGroup(RelationTypeGroup.RULE_NODE); + relationService.saveRelation(tenantId, sourceRuleNodeToTargetRuleNode); + } } } @@ -263,8 +295,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC int toIndex = ruleNodeIndexMap.get(toNodeId); ruleChainMetaData.addConnectionInfo(fromIndex, toIndex, type); } else if (nodeRelation.getTo().getEntityType() == EntityType.RULE_CHAIN) { - RuleChainId targetRuleChainId = new RuleChainId(nodeRelation.getTo().getId()); - ruleChainMetaData.addRuleChainConnectionInfo(fromIndex, targetRuleChainId, type, nodeRelation.getAdditionalInfo()); + log.warn("[{}][{}] Unsupported node relation: {}", tenantId, ruleChainId, nodeRelation.getTo()); } } } From da8bff5209b0396af891a748c5aff56845a9d138 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 28 Feb 2022 17:27:48 +0200 Subject: [PATCH 2/3] Rule Chain Service improvement for invalid rule chain ids --- .../server/dao/rule/BaseRuleChainService.java | 9 ++------- 1 file changed, 2 insertions(+), 7 deletions(-) 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 0f62c57741..097779f33e 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 @@ -202,17 +202,12 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC for (RuleChainConnectionInfo nodeToRuleChainConnection : ruleChainMetaData.getRuleChainConnections()) { RuleChainId targetRuleChainId = nodeToRuleChainConnection.getTargetRuleChainId(); RuleChain targetRuleChain = findRuleChainById(TenantId.SYS_TENANT_ID, targetRuleChainId); - if (targetRuleChain == null) { - log.info("Skip processing of relation for non existing target rule chain: [{}]", targetRuleChainId); - continue; - } - RuleNode targetNode = new RuleNode(); - targetNode.setName(targetRuleChain.getName()); + targetNode.setName(targetRuleChain != null ? targetRuleChain.getName() : "Rule Chain Input"); targetNode.setRuleChainId(ruleChain.getId()); targetNode.setType("org.thingsboard.rule.engine.flow.TbRuleChainInputNode"); var configuration = JacksonUtil.newObjectNode(); - configuration.put("ruleChainId", targetRuleChain.getId().toString()); + configuration.put("ruleChainId", targetRuleChainId.getId().toString()); targetNode.setConfiguration(configuration); ObjectNode layout = (ObjectNode) nodeToRuleChainConnection.getAdditionalInfo(); layout.remove("description"); From 603abff59905150497b6404f31b467b2e1c9fb8b Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Mon, 28 Feb 2022 18:45:33 +0200 Subject: [PATCH 3/3] Handle old rule chain connections while rule chain import --- .../import-export/import-export.service.ts | 54 +++++++++++++++++-- 1 file changed, 51 insertions(+), 3 deletions(-) diff --git a/ui-ngx/src/app/modules/home/components/import-export/import-export.service.ts b/ui-ngx/src/app/modules/home/components/import-export/import-export.service.ts index 2bddd229bc..7f52520bd6 100644 --- a/ui-ngx/src/app/modules/home/components/import-export/import-export.service.ts +++ b/ui-ngx/src/app/modules/home/components/import-export/import-export.service.ts @@ -35,7 +35,7 @@ import { import { MatDialog } from '@angular/material/dialog'; import { ImportDialogComponent, ImportDialogData } from '@home/components/import-export/import-dialog.component'; import { forkJoin, Observable, of } from 'rxjs'; -import { catchError, map, mergeMap, tap } from 'rxjs/operators'; +import {catchError, map, mergeMap, switchMap, tap} from 'rxjs/operators'; import { DashboardUtilsService } from '@core/services/dashboard-utils.service'; import { EntityService } from '@core/http/entity.service'; import { Widget, WidgetSize, WidgetType, WidgetTypeDetails } from '@shared/models/widget.models'; @@ -62,6 +62,7 @@ import { TenantProfileService } from '@core/http/tenant-profile.service'; import { DeviceService } from '@core/http/device.service'; import { AssetService } from '@core/http/asset.service'; import { EdgeService } from '@core/http/edge.service'; +import {RuleNode} from "@shared/models/rule-node.models"; // @dynamic @Injectable() @@ -423,7 +424,7 @@ export class ImportExportService { public importRuleChain(expectedRuleChainType: RuleChainType): Observable { return this.openImportDialog('rulechain.import', 'rulechain.rulechain-file').pipe( - map((ruleChainImport: RuleChainImport) => { + mergeMap((ruleChainImport: RuleChainImport) => { if (!this.validateImportedRuleChain(ruleChainImport)) { this.store.dispatch(new ActionNotificationShow( {message: this.translate.instant('rulechain.invalid-rulechain-file-error'), @@ -435,7 +436,7 @@ export class ImportExportService { type: 'error'})); throw new Error('Invalid rule chain type'); } else { - return ruleChainImport; + return this.processOldRuleChainConnections(ruleChainImport); } }), catchError((err) => { @@ -444,6 +445,53 @@ export class ImportExportService { ); } + private processOldRuleChainConnections(ruleChainImport: RuleChainImport): Observable { + const metadata = ruleChainImport.metadata; + if ((metadata as any).ruleChainConnections) { + const ruleChainNameResolveObservables: Observable[] = []; + for (const ruleChainConnection of (metadata as any).ruleChainConnections) { + if (ruleChainConnection.targetRuleChainId && ruleChainConnection.targetRuleChainId.id) { + const ruleChainNode: RuleNode = { + name: '', + debugMode: false, + type: 'org.thingsboard.rule.engine.flow.TbRuleChainInputNode', + configuration: { + ruleChainId: ruleChainConnection.targetRuleChainId.id + }, + additionalInfo: ruleChainConnection.additionalInfo + }; + ruleChainNameResolveObservables.push(this.ruleChainService.getRuleChain(ruleChainNode.configuration.ruleChainId, + {ignoreErrors: true, ignoreLoading: true}).pipe( + catchError(err => { + return of({name: 'Rule Chain Input'} as RuleChain); + }), + map((ruleChain => { + ruleChainNode.name = ruleChain.name; + return null; + }) + ) + )); + const toIndex = metadata.nodes.length; + metadata.nodes.push(ruleChainNode); + metadata.connections.push({ + toIndex, + fromIndex: ruleChainConnection.fromIndex, + type: ruleChainConnection.type + }); + } + } + if (ruleChainNameResolveObservables.length) { + return forkJoin(ruleChainNameResolveObservables).pipe( + map(() => ruleChainImport) + ); + } else { + return of(ruleChainImport); + } + } else { + return of(ruleChainImport); + } + } + public exportDeviceProfile(deviceProfileId: string) { this.deviceProfileService.getDeviceProfile(deviceProfileId).subscribe( (deviceProfile) => {