Browse Source

updated rule node upgrade logic

pull/8661/head
ShvaykaD 3 years ago
parent
commit
490ad63ab6
  1. 64
      application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java
  2. 31
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java

64
application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java

@ -27,7 +27,6 @@ import org.springframework.context.annotation.Lazy;
import org.springframework.context.annotation.Profile; import org.springframework.context.annotation.Profile;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.TbVersionedNode; import org.thingsboard.rule.engine.api.TbVersionedNode;
import org.thingsboard.rule.engine.flow.TbRuleChainInputNode; import org.thingsboard.rule.engine.flow.TbRuleChainInputNode;
import org.thingsboard.rule.engine.flow.TbRuleChainInputNodeConfiguration; import org.thingsboard.rule.engine.flow.TbRuleChainInputNodeConfiguration;
@ -86,10 +85,11 @@ import org.thingsboard.server.service.bean.BeanDiscoveryService;
import org.thingsboard.server.service.install.InstallScripts; import org.thingsboard.server.service.install.InstallScripts;
import org.thingsboard.server.service.install.SystemDataLoaderService; import org.thingsboard.server.service.install.SystemDataLoaderService;
import java.lang.reflect.InvocationTargetException;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collections; import java.util.Collections;
import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicLong;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@ -221,35 +221,26 @@ public class DefaultDataUpdateService implements DataUpdateService {
private void upgradeRuleNodes() { private void upgradeRuleNodes() {
try { try {
log.info("Lookup rule nodes to upgrade ..."); log.info("Lookup rule nodes to upgrade ...");
ArrayList<TbVersionedNode> tbVersionedNodes = getTbVersionedNodes(); var nodeClassToVersionMap = getNodeClassToVersionMap();
log.info("Found {} versioned nodes to check for upgrade!", tbVersionedNodes.size()); log.info("Found {} versioned nodes to check for upgrade!", nodeClassToVersionMap.size());
for (TbVersionedNode tbVersionedNode : tbVersionedNodes) { nodeClassToVersionMap.forEach((clazz, toVersion) -> {
String ruleNodeType = tbVersionedNode.getClass().getName(); var ruleNodeType = clazz.getName();
String ruleNodeTypeForLogs = tbVersionedNode.getClass().getSimpleName(); var ruleNodeTypeForLogs = clazz.getSimpleName();
int toVersion = tbVersionedNode.getCurrentVersion();
log.info("Going to check for nodes with type: {} to upgrade to version: {}.", ruleNodeTypeForLogs, toVersion); log.info("Going to check for nodes with type: {} to upgrade to version: {}.", ruleNodeTypeForLogs, toVersion);
var ruleNodesToUpdate = new PageDataIterable<>( var ruleNodesToUpdate = new PageDataIterable<>(
pageLink -> pageLink -> ruleChainService.findAllRuleNodesByTypeAndVersionLessThan(ruleNodeType, toVersion, pageLink), 1024
ruleChainService.findAllRuleNodesByTypeAndVersionLessThan(
ruleNodeType,
toVersion,
pageLink
),
1024
); );
if (Iterables.isEmpty(ruleNodesToUpdate)) { if (Iterables.isEmpty(ruleNodesToUpdate)) {
log.info("There are no active nodes with type: {}, or all nodes with this type already set to latest version!", ruleNodeTypeForLogs); log.info("There are no active nodes with type: {}, or all nodes with this type already set to latest version!", ruleNodeTypeForLogs);
} else { } else {
for (RuleNode ruleNode : ruleNodesToUpdate) { for (var ruleNode : ruleNodesToUpdate) {
RuleNodeId ruleNodeId = ruleNode.getId(); var ruleNodeId = ruleNode.getId();
var oldConfiguration = ruleNode.getConfiguration(); var oldConfiguration = ruleNode.getConfiguration();
int fromVersion = ruleNode.getConfigurationVersion(); int fromVersion = ruleNode.getConfigurationVersion();
log.info("Going to upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {}", log.info("Going to upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {}",
ruleNodeId, ruleNodeId, ruleNodeTypeForLogs, fromVersion, toVersion);
ruleNodeTypeForLogs,
fromVersion,
toVersion);
try { try {
var tbVersionedNode = (TbVersionedNode) clazz.getDeclaredConstructor().newInstance();
TbPair<Boolean, JsonNode> upgradeRuleNodeConfigurationResult = tbVersionedNode.upgrade(fromVersion, oldConfiguration); TbPair<Boolean, JsonNode> upgradeRuleNodeConfigurationResult = tbVersionedNode.upgrade(fromVersion, oldConfiguration);
if (upgradeRuleNodeConfigurationResult.getFirst()) { if (upgradeRuleNodeConfigurationResult.getFirst()) {
ruleNode.setConfiguration(upgradeRuleNodeConfigurationResult.getSecond()); ruleNode.setConfiguration(upgradeRuleNodeConfigurationResult.getSecond());
@ -257,45 +248,34 @@ public class DefaultDataUpdateService implements DataUpdateService {
ruleNode.setConfigurationVersion(toVersion); ruleNode.setConfigurationVersion(toVersion);
ruleChainService.saveRuleNode(TenantId.SYS_TENANT_ID, ruleNode); ruleChainService.saveRuleNode(TenantId.SYS_TENANT_ID, ruleNode);
log.info("Successfully upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {}", log.info("Successfully upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {}",
ruleNodeId, ruleNodeId, ruleNodeTypeForLogs, fromVersion, toVersion);
ruleNodeTypeForLogs, } catch (Exception e) {
fromVersion,
toVersion);
} catch (TbNodeException e) {
log.warn("Failed to upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {} due to: ", log.warn("Failed to upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {} due to: ",
ruleNodeId, ruleNodeId, ruleNodeTypeForLogs, fromVersion, toVersion, e);
ruleNodeTypeForLogs,
fromVersion,
toVersion,
e);
} }
} }
} }
} });
log.info("Finished rule nodes upgrade!"); log.info("Finished rule nodes upgrade!");
} catch (Exception e) { } catch (Exception e) {
log.error("Unexpected error during rule nodes upgrade: ", e); log.error("Unexpected error during rule nodes upgrade: ", e);
} }
} }
private ArrayList<TbVersionedNode> getTbVersionedNodes() { private Map<Class<?>, Integer> getNodeClassToVersionMap() {
var ruleNodeDefinitions = beanDiscoveryService.discoverBeansByAnnotationType( var ruleNodeDefinitions = beanDiscoveryService.discoverBeansByAnnotationType(
org.thingsboard.rule.engine.api.RuleNode.class org.thingsboard.rule.engine.api.RuleNode.class
); );
var tbVersionedNodes = new ArrayList<TbVersionedNode>(); var tbVersionedNodes = new HashMap<Class<?>, Integer>();
for (var def : ruleNodeDefinitions) { for (var def : ruleNodeDefinitions) {
String clazzName = def.getBeanClassName(); String clazzName = def.getBeanClassName();
try { try {
var clazz = Class.forName(clazzName); var clazz = Class.forName(clazzName);
if (TbVersionedNode.class.isAssignableFrom(clazz)) { if (TbVersionedNode.class.isAssignableFrom(clazz)) {
tbVersionedNodes.add((TbVersionedNode) clazz.getDeclaredConstructor().newInstance()); TbVersionedNode tbVersionedNode = (TbVersionedNode) clazz.getDeclaredConstructor().newInstance();
} tbVersionedNodes.put(clazz, tbVersionedNode.getCurrentVersion());
} catch (NoSuchMethodException | }
InstantiationException | } catch (Exception e) {
IllegalAccessException |
InvocationTargetException |
ClassNotFoundException e
) {
log.warn("Failed to create instance of rule node type: {} due to: ", clazzName, e); log.warn("Failed to create instance of rule node type: {} due to: ", clazzName, e);
} }
} }

31
dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java

@ -199,10 +199,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
int toVersion = tbVersionedNode.getCurrentVersion(); int toVersion = tbVersionedNode.getCurrentVersion();
if (fromVersion < toVersion) { if (fromVersion < toVersion) {
log.debug("Going to upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {}", log.debug("Going to upgrade rule node with id: {} type: {} fromVersion: {} toVersion: {}",
ruleNodeId, ruleNodeId, ruleNodeType, fromVersion, toVersion);
ruleNodeType,
fromVersion,
toVersion);
try { try {
TbPair<Boolean, JsonNode> upgradeResult = tbVersionedNode.upgrade(fromVersion, node.getConfiguration()); TbPair<Boolean, JsonNode> upgradeResult = tbVersionedNode.upgrade(fromVersion, node.getConfiguration());
if (upgradeResult.getFirst()) { if (upgradeResult.getFirst()) {
@ -210,33 +207,19 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
} }
node.setConfigurationVersion(toVersion); node.setConfigurationVersion(toVersion);
log.debug("Successfully upgrade rule node with id: {} type: {}, rule chain id: {} fromVersion: {} toVersion: {}", log.debug("Successfully upgrade rule node with id: {} type: {}, rule chain id: {} fromVersion: {} toVersion: {}",
ruleNodeId, ruleNodeId, ruleNodeType, ruleChainId, fromVersion, toVersion);
ruleNodeType,
ruleChainId,
fromVersion,
toVersion);
} catch (TbNodeException e) { } catch (TbNodeException e) {
log.warn("Failed to upgrade rule node with id: {} type: {} rule chain id: {} fromVersion: {} toVersion: {} due to: ", log.warn("Failed to upgrade rule node with id: {} type: {} rule chain id: {} fromVersion: {} toVersion: {} due to: ",
ruleNodeId, ruleNodeId, ruleNodeType, ruleChainId, fromVersion, toVersion, e);
ruleNodeType,
ruleChainId,
fromVersion,
toVersion,
e);
} }
} else { } else {
log.debug("Rule node with id: {} type: {} ruleChainId: {} already set to latest version!", log.debug("Rule node with id: {} type: {} ruleChainId: {} already set to latest version!",
ruleNodeId, ruleNodeId, ruleChainId, ruleNodeType);
ruleChainId,
ruleNodeType);
} }
} }
} catch (ClassNotFoundException | } catch (Exception e) {
InvocationTargetException | log.error("Failed to create instance of rule node with id: {} type: {}, rule chain id: {}",
InstantiationException | ruleNodeId, ruleNodeType, ruleChainId);
IllegalAccessException |
NoSuchMethodException e) {
log.error("Failed to create instance of rule node with id: {} type: {}, rule chain id: {}", ruleNodeId, ruleNodeType, ruleChainId);
} }
RuleNode savedNode = ruleNodeDao.save(tenantId, node); RuleNode savedNode = ruleNodeDao.save(tenantId, node);
relations.add(new EntityRelation(ruleChainMetaData.getRuleChainId(), savedNode.getId(), relations.add(new EntityRelation(ruleChainMetaData.getRuleChainId(), savedNode.getId(),

Loading…
Cancel
Save