diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java index 402caed226..b02bc302de 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java @@ -469,6 +469,14 @@ class DefaultTbContext implements TbContext { return mainCtx.getRuleNodeStateService().save(getTenantId(), state); } + @Override + public void clearRuleNodeStates() { + if (log.isDebugEnabled()) { + log.debug("[{}][{}] Going to clear rule node states", getTenantId(), getSelfId()); + } + mainCtx.getRuleNodeStateService().removeByRuleNodeId(getTenantId(), getSelfId()); + } + private TbMsgMetaData getActionMetaData(RuleNodeId ruleNodeId) { TbMsgMetaData metaData = new TbMsgMetaData(); metaData.putValue("ruleNodeId", ruleNodeId.toString()); 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 5c54c0232e..6d80c05c05 100644 --- a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java +++ b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java @@ -164,6 +164,8 @@ public class RuleChainController extends BaseController { RuleChain savedRuleChain = installScripts.createDefaultRuleChain(getCurrentUser().getTenantId(), request.getName()); + tbClusterService.onEntityStateChange(savedRuleChain.getTenantId(), savedRuleChain.getId(), ComponentLifecycleEvent.CREATED); + logEntityAction(savedRuleChain.getId(), savedRuleChain, null, ActionType.ADDED, null); return savedRuleChain; diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleNodeStateService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleNodeStateService.java index 07138a1a11..d5ad9bbbb6 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleNodeStateService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleNodeStateService.java @@ -30,4 +30,5 @@ public interface RuleNodeStateService { RuleNodeState save(TenantId tenantId, RuleNodeState ruleNodeState); + void removeByRuleNodeId(TenantId tenantId, RuleNodeId selfId); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/DynamicValueSourceType.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/DynamicValueSourceType.java index 96734fde92..7da9ab00af 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/DynamicValueSourceType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/DynamicValueSourceType.java @@ -18,5 +18,6 @@ package org.thingsboard.server.common.data.query; public enum DynamicValueSourceType { CURRENT_TENANT, CURRENT_CUSTOMER, - CURRENT_USER + CURRENT_USER, + CURRENT_DEVICE } diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleNodeStateService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleNodeStateService.java index a0b83f0333..1d7e8af03f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleNodeStateService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleNodeStateService.java @@ -68,6 +68,17 @@ public class BaseRuleNodeStateService extends AbstractEntityService implements R return saveOrUpdate(tenantId, ruleNodeState, false); } + @Override + public void removeByRuleNodeId(TenantId tenantId, RuleNodeId ruleNodeId) { + if (tenantId == null) { + throw new DataValidationException("Tenant id should be specified!."); + } + if (ruleNodeId == null) { + throw new DataValidationException("Rule node id should be specified!."); + } + ruleNodeStateDao.removeByRuleNodeId(ruleNodeId.getId()); + } + public RuleNodeState saveOrUpdate(TenantId tenantId, RuleNodeState ruleNodeState, boolean update) { try { if (update) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeStateDao.java b/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeStateDao.java index b12c448aa5..89c2baea7f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeStateDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeStateDao.java @@ -16,6 +16,8 @@ package org.thingsboard.server.dao.rule; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.RuleNodeId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.rule.RuleNodeState; @@ -31,4 +33,6 @@ public interface RuleNodeStateDao extends Dao { PageData findByRuleNodeId(UUID ruleNodeId, PageLink pageLink); RuleNodeState findByRuleNodeIdAndEntityId(UUID ruleNodeId, UUID entityId); + + void removeByRuleNodeId(UUID ruleNodeId); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java index 232a534e88..75f7658a5a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java @@ -472,13 +472,13 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { if (entityFilter.isFetchLastLevelOnly()) { String fromOrTo = (entityFilter.getDirection().equals(EntitySearchDirection.FROM) ? "from" : "to"); StringBuilder notExistsPart = new StringBuilder(); - notExistsPart.append(" NOT EXISTS (SELECT 1 from relation nr where ") + notExistsPart.append(" NOT EXISTS (SELECT 1 from relation nr ") + .append(whereFilter.replaceAll("re\\.", "nr\\.")) + .append(" and ") .append("nr.").append(fromOrTo).append("_id").append(" = re.").append(toOrFrom).append("_id") .append(" and ") .append("nr.").append(fromOrTo).append("_type").append(" = re.").append(toOrFrom).append("_type"); - if (!StringUtils.isEmpty(entityFilter.getRelationType())) { - notExistsPart.append(" and nr.relation_type = :where_relation_type"); - } + notExistsPart.append(")"); whereFilter += " and ( re.lvl = " + entityFilter.getMaxLevel() + " OR " + notExistsPart.toString() + ")"; } @@ -551,12 +551,12 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { StringBuilder notExistsPart = new StringBuilder(); notExistsPart.append(" NOT EXISTS (SELECT 1 from relation nr WHERE "); - notExistsPart.append(whereFilter.toString()); notExistsPart - .append(" and ") .append("nr.").append(fromOrTo).append("_id").append(" = re.").append(toOrFrom).append("_id") .append(" and ") - .append("nr.").append(fromOrTo).append("_type").append(" = re.").append(toOrFrom).append("_type"); + .append("nr.").append(fromOrTo).append("_type").append(" = re.").append(toOrFrom).append("_type") + .append(" and ") + .append(whereFilter.toString().replaceAll("re\\.", "nr\\.")); notExistsPart.append(")"); whereFilter.append(" and ( re.lvl = ").append(entityFilter.getMaxLevel()).append(" OR ").append(notExistsPart.toString()).append(")"); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeStateDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeStateDao.java index 51e487b456..c61e58353c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeStateDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeStateDao.java @@ -19,6 +19,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.data.repository.CrudRepository; import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Transactional; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; @@ -56,4 +57,10 @@ public class JpaRuleNodeStateDao extends JpaAbstractDao entityKeys; private PersistedAlarmRuleState state; private boolean updateFlag; - public AlarmRuleState(AlarmSeverity severity, AlarmRule alarmRule, PersistedAlarmRuleState state) { + AlarmRuleState(AlarmSeverity severity, AlarmRule alarmRule, Set entityKeys, PersistedAlarmRuleState state) { this.severity = severity; this.alarmRule = alarmRule; + this.entityKeys = entityKeys; if (state != null) { this.state = state; } else { @@ -76,6 +82,30 @@ public class AlarmRuleState { this.requiredRepeats = requiredRepeats; } + public boolean validateTsUpdate(Set changedKeys) { + for (EntityKey key : changedKeys) { + if (entityKeys.contains(key)) { + return true; + } + } + return false; + } + + public boolean validateAttrUpdate(Set changedKeys) { + //If the attribute was updated, but no new telemetry arrived - we ignore this until new telemetry is there. + for (EntityKey key : entityKeys) { + if (key.getType().equals(EntityKeyType.TIME_SERIES)) { + return false; + } + } + for (EntityKey key : changedKeys) { + if (entityKeys.contains(key)) { + return true; + } + } + return false; + } + public AlarmConditionSpec getSpec(AlarmRule alarmRule) { AlarmConditionSpec spec = alarmRule.getCondition().getSpec(); if (spec == null) { @@ -93,7 +123,7 @@ public class AlarmRuleState { } } - public boolean eval(DeviceDataSnapshot data) { + public boolean eval(DataSnapshot data) { boolean active = isActive(data.getTs()); switch (spec.getType()) { case SIMPLE: @@ -167,7 +197,7 @@ public class AlarmRuleState { } } - private boolean evalRepeating(DeviceDataSnapshot data, boolean active) { + private boolean evalRepeating(DataSnapshot data, boolean active) { if (active && eval(alarmRule.getCondition(), data)) { state.setEventCount(state.getEventCount() + 1); updateFlag = true; @@ -177,7 +207,7 @@ public class AlarmRuleState { } } - private boolean evalDuration(DeviceDataSnapshot data, boolean active) { + private boolean evalDuration(DataSnapshot data, boolean active) { if (active && eval(alarmRule.getCondition(), data)) { if (state.getLastEventTs() > 0) { if (data.getTs() > state.getLastEventTs()) { @@ -211,45 +241,45 @@ public class AlarmRuleState { } } - private boolean eval(AlarmCondition condition, DeviceDataSnapshot data) { + private boolean eval(AlarmCondition condition, DataSnapshot data) { boolean eval = true; for (KeyFilter keyFilter : condition.getCondition()) { EntityKeyValue value = data.getValue(keyFilter.getKey()); if (value == null) { return false; } - eval = eval && eval(value, keyFilter.getPredicate()); + eval = eval && eval(data, value, keyFilter.getPredicate()); } return eval; } - private boolean eval(EntityKeyValue value, KeyFilterPredicate predicate) { + private boolean eval(DataSnapshot data, EntityKeyValue value, KeyFilterPredicate predicate) { switch (predicate.getType()) { case STRING: - return evalStrPredicate(value, (StringFilterPredicate) predicate); + return evalStrPredicate(data, value, (StringFilterPredicate) predicate); case NUMERIC: - return evalNumPredicate(value, (NumericFilterPredicate) predicate); - case COMPLEX: - return evalComplexPredicate(value, (ComplexFilterPredicate) predicate); + return evalNumPredicate(data, value, (NumericFilterPredicate) predicate); case BOOLEAN: - return evalBoolPredicate(value, (BooleanFilterPredicate) predicate); + return evalBoolPredicate(data, value, (BooleanFilterPredicate) predicate); + case COMPLEX: + return evalComplexPredicate(data, value, (ComplexFilterPredicate) predicate); default: return false; } } - private boolean evalComplexPredicate(EntityKeyValue ekv, ComplexFilterPredicate predicate) { + private boolean evalComplexPredicate(DataSnapshot data, EntityKeyValue ekv, ComplexFilterPredicate predicate) { switch (predicate.getOperation()) { case OR: for (KeyFilterPredicate kfp : predicate.getPredicates()) { - if (eval(ekv, kfp)) { + if (eval(data, ekv, kfp)) { return true; } } return false; case AND: for (KeyFilterPredicate kfp : predicate.getPredicates()) { - if (!eval(ekv, kfp)) { + if (!eval(data, ekv, kfp)) { return false; } } @@ -259,109 +289,55 @@ public class AlarmRuleState { } } - private boolean evalBoolPredicate(EntityKeyValue ekv, BooleanFilterPredicate predicate) { - Boolean value; - switch (ekv.getDataType()) { - case LONG: - value = ekv.getLngValue() > 0; - break; - case DOUBLE: - value = ekv.getDblValue() > 0; - break; - case BOOLEAN: - value = ekv.getBoolValue(); - break; - case STRING: - try { - value = Boolean.parseBoolean(ekv.getStrValue()); - break; - } catch (RuntimeException e) { - return false; - } - case JSON: - try { - value = Boolean.parseBoolean(ekv.getJsonValue()); - break; - } catch (RuntimeException e) { - return false; - } - default: - return false; - } - if (value == null) { + private boolean evalBoolPredicate(DataSnapshot data, EntityKeyValue ekv, BooleanFilterPredicate predicate) { + Boolean val = getBoolValue(ekv); + if (val == null) { return false; } + Boolean predicateValue = getPredicateValue(data, predicate.getValue(), AlarmRuleState::getBoolValue); switch (predicate.getOperation()) { case EQUAL: - return value.equals(predicate.getValue().getDefaultValue()); + return val.equals(predicateValue); case NOT_EQUAL: - return !value.equals(predicate.getValue().getDefaultValue()); + return !val.equals(predicateValue); default: throw new RuntimeException("Operation not supported: " + predicate.getOperation()); } } - private boolean evalNumPredicate(EntityKeyValue ekv, NumericFilterPredicate predicate) { - Double value; - switch (ekv.getDataType()) { - case LONG: - value = ekv.getLngValue().doubleValue(); - break; - case DOUBLE: - value = ekv.getDblValue(); - break; - case BOOLEAN: - value = ekv.getBoolValue() ? 1.0 : 0.0; - break; - case STRING: - try { - value = Double.parseDouble(ekv.getStrValue()); - break; - } catch (RuntimeException e) { - return false; - } - case JSON: - try { - value = Double.parseDouble(ekv.getJsonValue()); - break; - } catch (RuntimeException e) { - return false; - } - default: - return false; - } - if (value == null) { + private boolean evalNumPredicate(DataSnapshot data, EntityKeyValue ekv, NumericFilterPredicate predicate) { + Double val = getDblValue(ekv); + if (val == null) { return false; } - - Double predicateValue = predicate.getValue().getDefaultValue(); + Double predicateValue = getPredicateValue(data, predicate.getValue(), AlarmRuleState::getDblValue); switch (predicate.getOperation()) { case NOT_EQUAL: - return !value.equals(predicateValue); + return !val.equals(predicateValue); case EQUAL: - return value.equals(predicateValue); + return val.equals(predicateValue); case GREATER: - return value > predicateValue; + return val > predicateValue; case GREATER_OR_EQUAL: - return value >= predicateValue; + return val >= predicateValue; case LESS: - return value < predicateValue; + return val < predicateValue; case LESS_OR_EQUAL: - return value <= predicateValue; + return val <= predicateValue; default: throw new RuntimeException("Operation not supported: " + predicate.getOperation()); } } - private boolean evalStrPredicate(EntityKeyValue ekv, StringFilterPredicate predicate) { - String val; - String predicateValue; + private boolean evalStrPredicate(DataSnapshot data, EntityKeyValue ekv, StringFilterPredicate predicate) { + String val = getStrValue(ekv); + if (val == null) { + return false; + } + String predicateValue = getPredicateValue(data, predicate.getValue(), AlarmRuleState::getStrValue); if (predicate.isIgnoreCase()) { - val = ekv.getStrValue().toLowerCase(); - predicateValue = predicate.getValue().getDefaultValue().toLowerCase(); - } else { - val = ekv.getStrValue(); - predicateValue = predicate.getValue().getDefaultValue(); + val = val.toLowerCase(); + predicateValue = predicateValue.toLowerCase(); } switch (predicate.getOperation()) { case CONTAINS: @@ -380,4 +356,100 @@ public class AlarmRuleState { throw new RuntimeException("Operation not supported: " + predicate.getOperation()); } } + + private T getPredicateValue(DataSnapshot data, FilterPredicateValue value, Function transformFunction) { + EntityKeyValue ekv = getDynamicPredicateValue(data, value); + if (ekv != null) { + T result = transformFunction.apply(ekv); + if (result != null) { + return result; + } + } + return value.getDefaultValue(); + } + + private EntityKeyValue getDynamicPredicateValue(DataSnapshot data, FilterPredicateValue value) { + EntityKeyValue ekv = null; + if (value.getDynamicValue() != null) { + ekv = data.getValue(new EntityKey(EntityKeyType.ATTRIBUTE, value.getDynamicValue().getSourceAttribute())); + if (ekv == null) { + ekv = data.getValue(new EntityKey(EntityKeyType.SERVER_ATTRIBUTE, value.getDynamicValue().getSourceAttribute())); + if (ekv == null) { + ekv = data.getValue(new EntityKey(EntityKeyType.SHARED_ATTRIBUTE, value.getDynamicValue().getSourceAttribute())); + if (ekv == null) { + ekv = data.getValue(new EntityKey(EntityKeyType.CLIENT_ATTRIBUTE, value.getDynamicValue().getSourceAttribute())); + } + } + } + } + return ekv; + } + + private static String getStrValue(EntityKeyValue ekv) { + switch (ekv.getDataType()) { + case LONG: + return ekv.getLngValue() != null ? ekv.getLngValue().toString() : null; + case DOUBLE: + return ekv.getDblValue() != null ? ekv.getDblValue().toString() : null; + case BOOLEAN: + return ekv.getBoolValue() != null ? ekv.getBoolValue().toString() : null; + case STRING: + return ekv.getStrValue(); + case JSON: + return ekv.getJsonValue(); + default: + return null; + } + } + + private static Double getDblValue(EntityKeyValue ekv) { + switch (ekv.getDataType()) { + case LONG: + return ekv.getLngValue() != null ? ekv.getLngValue().doubleValue() : null; + case DOUBLE: + return ekv.getDblValue() != null ? ekv.getDblValue() : null; + case BOOLEAN: + return ekv.getBoolValue() != null ? (ekv.getBoolValue() ? 1.0 : 0.0) : null; + case STRING: + try { + return Double.parseDouble(ekv.getStrValue()); + } catch (RuntimeException e) { + return null; + } + case JSON: + try { + return Double.parseDouble(ekv.getJsonValue()); + } catch (RuntimeException e) { + return null; + } + default: + return null; + } + } + + private static Boolean getBoolValue(EntityKeyValue ekv) { + switch (ekv.getDataType()) { + case LONG: + return ekv.getLngValue() != null ? ekv.getLngValue() > 0 : null; + case DOUBLE: + return ekv.getDblValue() != null ? ekv.getDblValue() > 0 : null; + case BOOLEAN: + return ekv.getBoolValue(); + case STRING: + try { + return Boolean.parseBoolean(ekv.getStrValue()); + } catch (RuntimeException e) { + return null; + } + case JSON: + try { + return Boolean.parseBoolean(ekv.getJsonValue()); + } catch (RuntimeException e) { + return null; + } + default: + return null; + } + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceProfileAlarmState.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java similarity index 78% rename from rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceProfileAlarmState.java rename to rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java index f74d88fe62..5fb2c2957c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceProfileAlarmState.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java @@ -17,6 +17,7 @@ package org.thingsboard.rule.engine.profile; import com.fasterxml.jackson.databind.JsonNode; import lombok.Data; +import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.action.TbAlarmResult; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.profile.state.PersistedAlarmRuleState; @@ -27,6 +28,7 @@ import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmStatus; import org.thingsboard.server.common.data.device.profile.DeviceProfileAlarm; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.queue.ServiceQueue; @@ -39,8 +41,10 @@ import java.util.concurrent.ExecutionException; import java.util.function.BiFunction; @Data -class DeviceProfileAlarmState { +@Slf4j +class AlarmState { + private final ProfileState deviceProfile; private final EntityId originator; private DeviceProfileAlarm alarmDefinition; private volatile List createRulesSortedBySeverityDesc; @@ -50,27 +54,33 @@ class DeviceProfileAlarmState { private volatile TbMsgMetaData lastMsgMetaData; private volatile String lastMsgQueueName; - public DeviceProfileAlarmState(EntityId originator, DeviceProfileAlarm alarmDefinition, PersistedAlarmState alarmState) { + AlarmState(ProfileState deviceProfile, EntityId originator, DeviceProfileAlarm alarmDefinition, PersistedAlarmState alarmState) { + this.deviceProfile = deviceProfile; this.originator = originator; this.updateState(alarmDefinition, alarmState); } - public boolean process(TbContext ctx, TbMsg msg, DeviceDataSnapshot data) throws ExecutionException, InterruptedException { + public boolean process(TbContext ctx, TbMsg msg, DataSnapshot data, SnapshotUpdate update) throws ExecutionException, InterruptedException { initCurrentAlarm(ctx); lastMsgMetaData = msg.getMetaData(); lastMsgQueueName = msg.getQueueName(); - return createOrClearAlarms(ctx, data, AlarmRuleState::eval); + return createOrClearAlarms(ctx, data, update, AlarmRuleState::eval); } public boolean process(TbContext ctx, long ts) throws ExecutionException, InterruptedException { initCurrentAlarm(ctx); - return createOrClearAlarms(ctx, ts, AlarmRuleState::eval); + return createOrClearAlarms(ctx, ts, null, AlarmRuleState::eval); } - public boolean createOrClearAlarms(TbContext ctx, T data, BiFunction evalFunction) { + public boolean createOrClearAlarms(TbContext ctx, T data, SnapshotUpdate update, BiFunction evalFunction) { boolean stateUpdate = false; AlarmSeverity resultSeverity = null; + log.debug("[{}] processing update: {}", alarmDefinition.getId(), data); for (AlarmRuleState state : createRulesSortedBySeverityDesc) { + if (!validateUpdate(update, state)) { + log.debug("[{}][{}] Update is not valid for current rule state", alarmDefinition.getId(), state.getSeverity()); + continue; + } boolean evalResult = evalFunction.apply(state, data); stateUpdate |= state.checkUpdate(); if (evalResult) { @@ -81,9 +91,17 @@ class DeviceProfileAlarmState { if (resultSeverity != null) { pushMsg(ctx, calculateAlarmResult(ctx, resultSeverity)); } else if (currentAlarm != null && clearState != null) { + if (!validateUpdate(update, clearState)) { + log.debug("[{}] Update is not valid for current clear state", alarmDefinition.getId()); + return stateUpdate; + } Boolean evalResult = evalFunction.apply(clearState, data); if (evalResult) { stateUpdate |= clearState.checkUpdate(); + for (AlarmRuleState state : createRulesSortedBySeverityDesc) { + state.clear(); + stateUpdate |= state.checkUpdate(); + } ctx.getAlarmService().clearAlarm(ctx.getTenantId(), currentAlarm.getId(), JacksonUtil.OBJECT_MAPPER.createObjectNode(), System.currentTimeMillis()); pushMsg(ctx, new TbAlarmResult(false, false, true, currentAlarm)); currentAlarm = null; @@ -92,6 +110,18 @@ class DeviceProfileAlarmState { return stateUpdate; } + public boolean validateUpdate(SnapshotUpdate update, AlarmRuleState state) { + if (update != null) { + //Check that the update type and that keys match. + if (update.getType().equals(EntityKeyType.TIME_SERIES)) { + return state.validateTsUpdate(update.getKeys()); + } else if (update.getType().equals(EntityKeyType.ATTRIBUTE)) { + return state.validateAttrUpdate(update.getKeys()); + } + } + return true; + } + public void initCurrentAlarm(TbContext ctx) throws InterruptedException, ExecutionException { if (!initialFetchDone) { Alarm alarm = ctx.getAlarmService().findLatestByOriginatorAndType(ctx.getTenantId(), originator, alarmDefinition.getAlarmType()).get(); @@ -137,17 +167,20 @@ class DeviceProfileAlarmState { alarmState.getCreateRuleStates().put(severity, ruleState); } } - createRulesSortedBySeverityDesc.add(new AlarmRuleState(severity, rule, ruleState)); + createRulesSortedBySeverityDesc.add(new AlarmRuleState(severity, rule, + deviceProfile.getCreateAlarmKeys(alarm.getId(), severity), ruleState)); }); createRulesSortedBySeverityDesc.sort(Comparator.comparingInt(state -> state.getSeverity().ordinal())); PersistedAlarmRuleState ruleState = alarmState == null ? null : alarmState.getClearRuleState(); if (alarmDefinition.getClearRule() != null) { - clearState = new AlarmRuleState(null, alarmDefinition.getClearRule(), ruleState); + clearState = new AlarmRuleState(null, alarmDefinition.getClearRule(), deviceProfile.getClearAlarmKeys(alarm.getId()), ruleState); } } private TbAlarmResult calculateAlarmResult(TbContext ctx, AlarmSeverity severity) { if (currentAlarm != null) { + // TODO: In some extremely rare cases, we might miss the event of alarm clear (If one use in-mem queue and restarted the server) or (if one manipulated the rule chain). + // Maybe we should fetch alarm every time? currentAlarm.setEndTs(System.currentTimeMillis()); AlarmSeverity oldSeverity = currentAlarm.getSeverity(); if (!oldSeverity.equals(severity)) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmStateUpdateResult.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmStateUpdateResult.java index de9708dc58..a84cf96710 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmStateUpdateResult.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmStateUpdateResult.java @@ -15,7 +15,7 @@ */ package org.thingsboard.rule.engine.profile; -public enum AlarmStateUpdateResult { +enum AlarmStateUpdateResult { NONE, CREATED, UPDATED, SEVERITY_UPDATED, CLEARED; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceDataSnapshot.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DataSnapshot.java similarity index 70% rename from rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceDataSnapshot.java rename to rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DataSnapshot.java index f1b1067095..0c3abc0ca3 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceDataSnapshot.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DataSnapshot.java @@ -24,7 +24,7 @@ import java.util.Map; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; -public class DeviceDataSnapshot { +class DataSnapshot { private volatile boolean ready; @Getter @@ -33,7 +33,7 @@ public class DeviceDataSnapshot { private final Set keys; private final Map values = new ConcurrentHashMap<>(); - public DeviceDataSnapshot(Set entityKeysToFetch) { + DataSnapshot(Set entityKeysToFetch) { this.keys = entityKeysToFetch; } @@ -56,28 +56,38 @@ public class DeviceDataSnapshot { } } - void putValue(EntityKey key, EntityKeyValue value) { + boolean putValue(EntityKey key, long newTs, EntityKeyValue value) { + boolean updateOfTs = ts != newTs; + boolean result = false; switch (key.getType()) { case ATTRIBUTE: - putIfKeyExists(key, value); - putIfKeyExists(getAttrKey(key, EntityKeyType.CLIENT_ATTRIBUTE), value); - putIfKeyExists(getAttrKey(key, EntityKeyType.SHARED_ATTRIBUTE), value); - putIfKeyExists(getAttrKey(key, EntityKeyType.SERVER_ATTRIBUTE), value); + result |= putIfKeyExists(key, value, updateOfTs); + result |= putIfKeyExists(getAttrKey(key, EntityKeyType.CLIENT_ATTRIBUTE), value, updateOfTs); + result |= putIfKeyExists(getAttrKey(key, EntityKeyType.SHARED_ATTRIBUTE), value, updateOfTs); + result |= putIfKeyExists(getAttrKey(key, EntityKeyType.SERVER_ATTRIBUTE), value, updateOfTs); break; case CLIENT_ATTRIBUTE: case SHARED_ATTRIBUTE: case SERVER_ATTRIBUTE: - putIfKeyExists(key, value); - putIfKeyExists(getAttrKey(key, EntityKeyType.ATTRIBUTE), value); + result |= putIfKeyExists(key, value, updateOfTs); + result |= putIfKeyExists(getAttrKey(key, EntityKeyType.ATTRIBUTE), value, updateOfTs); break; default: - putIfKeyExists(key, value); + result |= putIfKeyExists(key, value, updateOfTs); } + return result; } - private void putIfKeyExists(EntityKey key, EntityKeyValue value) { + private boolean putIfKeyExists(EntityKey key, EntityKeyValue value, boolean updateOfTs) { if (keys.contains(key)) { - values.put(key, value); + EntityKeyValue oldValue = values.put(key, value); + if (updateOfTs) { + return true; + } else { + return oldValue == null || !oldValue.equals(value); + } + } else { + return false; } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceProfileState.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceProfileState.java deleted file mode 100644 index fd9037624e..0000000000 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceProfileState.java +++ /dev/null @@ -1,63 +0,0 @@ -/** - * Copyright © 2016-2020 The Thingsboard Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.thingsboard.rule.engine.profile; - -import lombok.AccessLevel; -import lombok.Getter; -import org.thingsboard.server.common.data.DeviceProfile; -import org.thingsboard.server.common.data.device.profile.AlarmRule; -import org.thingsboard.server.common.data.device.profile.DeviceProfileAlarm; -import org.thingsboard.server.common.data.id.DeviceProfileId; -import org.thingsboard.server.common.data.query.EntityKey; -import org.thingsboard.server.common.data.query.KeyFilter; - -import java.util.List; -import java.util.Set; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.CopyOnWriteArrayList; - - -class DeviceProfileState { - - private DeviceProfile deviceProfile; - @Getter(AccessLevel.PACKAGE) - private final List alarmSettings = new CopyOnWriteArrayList<>(); - @Getter(AccessLevel.PACKAGE) - private final Set entityKeys = ConcurrentHashMap.newKeySet(); - - DeviceProfileState(DeviceProfile deviceProfile) { - updateDeviceProfile(deviceProfile); - } - - void updateDeviceProfile(DeviceProfile deviceProfile) { - this.deviceProfile = deviceProfile; - alarmSettings.clear(); - if (deviceProfile.getProfileData().getAlarms() != null) { - alarmSettings.addAll(deviceProfile.getProfileData().getAlarms()); - for (DeviceProfileAlarm alarm : deviceProfile.getProfileData().getAlarms()) { - for (AlarmRule alarmRule : alarm.getCreateRules().values()) { - for (KeyFilter keyFilter : alarmRule.getCondition().getCondition()) { - entityKeys.add(keyFilter.getKey()); - } - } - } - } - } - - public DeviceProfileId getProfileId() { - return deviceProfile.getId(); - } -} diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java index 078be24f90..6824f56d35 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java @@ -16,6 +16,7 @@ package org.thingsboard.rule.engine.profile; import com.google.gson.JsonParser; +import lombok.extern.slf4j.Slf4j; import org.springframework.util.StringUtils; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.profile.state.PersistedAlarmState; @@ -29,7 +30,6 @@ import org.thingsboard.server.common.data.device.profile.DeviceProfileAlarm; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.RuleNodeStateId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; @@ -53,17 +53,18 @@ import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; +@Slf4j class DeviceState { private final boolean persistState; private final DeviceId deviceId; + private final ProfileState deviceProfile; private RuleNodeState state; - private DeviceProfileState deviceProfile; private PersistedDeviceState pds; - private DeviceDataSnapshot latestValues; - private final ConcurrentMap alarmStates = new ConcurrentHashMap<>(); + private DataSnapshot latestValues; + private final ConcurrentMap alarmStates = new ConcurrentHashMap<>(); - public DeviceState(TbContext ctx, TbDeviceProfileNodeConfiguration config, DeviceId deviceId, DeviceProfileState deviceProfile, RuleNodeState state) { + DeviceState(TbContext ctx, TbDeviceProfileNodeConfiguration config, DeviceId deviceId, ProfileState deviceProfile, RuleNodeState state) { this.persistState = config.isPersistAlarmRulesState(); this.deviceId = deviceId; this.deviceProfile = deviceProfile; @@ -86,7 +87,7 @@ class DeviceState { if (pds != null) { for (DeviceProfileAlarm alarm : deviceProfile.getAlarmSettings()) { alarmStates.computeIfAbsent(alarm.getId(), - a -> new DeviceProfileAlarmState(deviceId, alarm, getOrInitPersistedAlarmState(alarm))); + a -> new AlarmState(deviceProfile, deviceId, alarm, getOrInitPersistedAlarmState(alarm))); } } } @@ -107,14 +108,20 @@ class DeviceState { if (alarmStates.containsKey(alarm.getId())) { alarmStates.get(alarm.getId()).updateState(alarm, getOrInitPersistedAlarmState(alarm)); } else { - alarmStates.putIfAbsent(alarm.getId(), new DeviceProfileAlarmState(deviceId, alarm, getOrInitPersistedAlarmState(alarm))); + alarmStates.putIfAbsent(alarm.getId(), new AlarmState(this.deviceProfile, deviceId, alarm, getOrInitPersistedAlarmState(alarm))); } } } public void harvestAlarms(TbContext ctx, long ts) throws ExecutionException, InterruptedException { - for (DeviceProfileAlarmState state : alarmStates.values()) { - state.process(ctx, ts); + log.debug("[{}] Going to harvest alarms: {}", ctx.getSelfId(), ts); + boolean stateChanged = false; + for (AlarmState state : alarmStates.values()) { + stateChanged |= state.process(ctx, ts); + } + if (persistState && stateChanged) { + state.setStateData(JacksonUtil.toString(pds)); + state = ctx.saveRuleNodeState(state); } } @@ -146,8 +153,8 @@ class DeviceState { boolean stateChanged = false; Alarm alarmNf = JacksonUtil.fromString(msg.getData(), Alarm.class); for (DeviceProfileAlarm alarm : deviceProfile.getAlarmSettings()) { - DeviceProfileAlarmState alarmState = alarmStates.computeIfAbsent(alarm.getId(), - a -> new DeviceProfileAlarmState(deviceId, alarm, getOrInitPersistedAlarmState(alarm))); + AlarmState alarmState = alarmStates.computeIfAbsent(alarm.getId(), + a -> new AlarmState(this.deviceProfile, deviceId, alarm, getOrInitPersistedAlarmState(alarm))); stateChanged |= alarmState.processAlarmClear(ctx, alarmNf); } ctx.tellSuccess(msg); @@ -175,9 +182,9 @@ class DeviceState { EntityKeyType keyType = getKeyTypeFromScope(scope); keys.forEach(key -> latestValues.removeValue(new EntityKey(keyType, key))); for (DeviceProfileAlarm alarm : deviceProfile.getAlarmSettings()) { - DeviceProfileAlarmState alarmState = alarmStates.computeIfAbsent(alarm.getId(), - a -> new DeviceProfileAlarmState(deviceId, alarm, getOrInitPersistedAlarmState(alarm))); - stateChanged |= alarmState.process(ctx, msg, latestValues); + AlarmState alarmState = alarmStates.computeIfAbsent(alarm.getId(), + a -> new AlarmState(this.deviceProfile, deviceId, alarm, getOrInitPersistedAlarmState(alarm))); + stateChanged |= alarmState.process(ctx, msg, latestValues, null); } } ctx.tellSuccess(msg); @@ -192,11 +199,11 @@ class DeviceState { private boolean processAttributesUpdate(TbContext ctx, TbMsg msg, Set attributes, String scope) throws ExecutionException, InterruptedException { boolean stateChanged = false; if (!attributes.isEmpty()) { - merge(latestValues, attributes, scope); + SnapshotUpdate update = merge(latestValues, attributes, scope); for (DeviceProfileAlarm alarm : deviceProfile.getAlarmSettings()) { - DeviceProfileAlarmState alarmState = alarmStates.computeIfAbsent(alarm.getId(), - a -> new DeviceProfileAlarmState(deviceId, alarm, getOrInitPersistedAlarmState(alarm))); - stateChanged |= alarmState.process(ctx, msg, latestValues); + AlarmState alarmState = alarmStates.computeIfAbsent(alarm.getId(), + a -> new AlarmState(this.deviceProfile, deviceId, alarm, getOrInitPersistedAlarmState(alarm))); + stateChanged |= alarmState.process(ctx, msg, latestValues, update); } } ctx.tellSuccess(msg); @@ -206,34 +213,47 @@ class DeviceState { protected boolean processTelemetry(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException { boolean stateChanged = false; Map> tsKvMap = JsonConverter.convertToSortedTelemetry(new JsonParser().parse(msg.getData()), TbMsgTimeseriesNode.getTs(msg)); + // iterate over data by ts (ASC order). for (Map.Entry> entry : tsKvMap.entrySet()) { Long ts = entry.getKey(); List data = entry.getValue(); - merge(latestValues, ts, data); - for (DeviceProfileAlarm alarm : deviceProfile.getAlarmSettings()) { - DeviceProfileAlarmState alarmState = alarmStates.computeIfAbsent(alarm.getId(), - a -> new DeviceProfileAlarmState(deviceId, alarm, getOrInitPersistedAlarmState(alarm))); - stateChanged |= alarmState.process(ctx, msg, latestValues); + SnapshotUpdate update = merge(latestValues, ts, data); + if (update.hasUpdate()) { + for (DeviceProfileAlarm alarm : deviceProfile.getAlarmSettings()) { + AlarmState alarmState = alarmStates.computeIfAbsent(alarm.getId(), + a -> new AlarmState(this.deviceProfile, deviceId, alarm, getOrInitPersistedAlarmState(alarm))); + stateChanged |= alarmState.process(ctx, msg, latestValues, update); + } } } ctx.tellSuccess(msg); return stateChanged; } - private void merge(DeviceDataSnapshot latestValues, Long ts, List data) { - latestValues.setTs(ts); + private SnapshotUpdate merge(DataSnapshot latestValues, Long newTs, List data) { + Set keys = new HashSet<>(); for (KvEntry entry : data) { - latestValues.putValue(new EntityKey(EntityKeyType.TIME_SERIES, entry.getKey()), toEntityValue(entry)); + EntityKey entityKey = new EntityKey(EntityKeyType.TIME_SERIES, entry.getKey()); + if (latestValues.putValue(entityKey, newTs, toEntityValue(entry))) { + keys.add(entityKey); + } } + latestValues.setTs(newTs); + return new SnapshotUpdate(EntityKeyType.TIME_SERIES, keys); } - private void merge(DeviceDataSnapshot latestValues, Set attributes, String scope) { - long ts = latestValues.getTs(); + private SnapshotUpdate merge(DataSnapshot latestValues, Set attributes, String scope) { + long newTs = 0; + Set keys = new HashSet<>(); for (AttributeKvEntry entry : attributes) { - ts = Math.max(ts, entry.getLastUpdateTs()); - latestValues.putValue(new EntityKey(getKeyTypeFromScope(scope), entry.getKey()), toEntityValue(entry)); + newTs = Math.max(newTs, entry.getLastUpdateTs()); + EntityKey entityKey = new EntityKey(getKeyTypeFromScope(scope), entry.getKey()); + if (latestValues.putValue(entityKey, newTs, toEntityValue(entry))) { + keys.add(entityKey); + } } - latestValues.setTs(ts); + latestValues.setTs(newTs); + return new SnapshotUpdate(EntityKeyType.ATTRIBUTE, keys); } private static EntityKeyType getKeyTypeFromScope(String scope) { @@ -248,14 +268,14 @@ class DeviceState { return EntityKeyType.ATTRIBUTE; } - private DeviceDataSnapshot fetchLatestValues(TbContext ctx, EntityId originator) throws ExecutionException, InterruptedException { + private DataSnapshot fetchLatestValues(TbContext ctx, EntityId originator) throws ExecutionException, InterruptedException { Set entityKeysToFetch = deviceProfile.getEntityKeys(); - DeviceDataSnapshot result = new DeviceDataSnapshot(entityKeysToFetch); + DataSnapshot result = new DataSnapshot(entityKeysToFetch); addEntityKeysToSnapshot(ctx, originator, entityKeysToFetch, result); return result; } - private void addEntityKeysToSnapshot(TbContext ctx, EntityId originator, Set entityKeysToFetch, DeviceDataSnapshot result) throws InterruptedException, ExecutionException { + private void addEntityKeysToSnapshot(TbContext ctx, EntityId originator, Set entityKeysToFetch, DataSnapshot result) throws InterruptedException, ExecutionException { Set serverAttributeKeys = new HashSet<>(); Set clientAttributeKeys = new HashSet<>(); Set sharedAttributeKeys = new HashSet<>(); @@ -291,16 +311,16 @@ class DeviceState { if (device != null) { switch (key) { case EntityKeyMapping.NAME: - result.putValue(entityKey, EntityKeyValue.fromString(device.getName())); + result.putValue(entityKey, device.getCreatedTime(), EntityKeyValue.fromString(device.getName())); break; case EntityKeyMapping.TYPE: - result.putValue(entityKey, EntityKeyValue.fromString(device.getType())); + result.putValue(entityKey, device.getCreatedTime(), EntityKeyValue.fromString(device.getType())); break; case EntityKeyMapping.CREATED_TIME: - result.putValue(entityKey, EntityKeyValue.fromLong(device.getCreatedTime())); + result.putValue(entityKey, device.getCreatedTime(), EntityKeyValue.fromLong(device.getCreatedTime())); break; case EntityKeyMapping.LABEL: - result.putValue(entityKey, EntityKeyValue.fromString(device.getLabel())); + result.putValue(entityKey, device.getCreatedTime(), EntityKeyValue.fromString(device.getLabel())); break; } } @@ -312,7 +332,7 @@ class DeviceState { List data = ctx.getTimeseriesService().findLatest(ctx.getTenantId(), originator, latestTsKeys).get(); for (TsKvEntry entry : data) { if (entry.getValue() != null) { - result.putValue(new EntityKey(EntityKeyType.TIME_SERIES, entry.getKey()), toEntityValue(entry)); + result.putValue(new EntityKey(EntityKeyType.TIME_SERIES, entry.getKey()), entry.getTs(), toEntityValue(entry)); } } } @@ -330,13 +350,13 @@ class DeviceState { } } - private void addToSnapshot(DeviceDataSnapshot snapshot, Set commonAttributeKeys, List data) { + private void addToSnapshot(DataSnapshot snapshot, Set commonAttributeKeys, List data) { for (AttributeKvEntry entry : data) { if (entry.getValue() != null) { EntityKeyValue value = toEntityValue(entry); - snapshot.putValue(new EntityKey(EntityKeyType.CLIENT_ATTRIBUTE, entry.getKey()), value); + snapshot.putValue(new EntityKey(EntityKeyType.CLIENT_ATTRIBUTE, entry.getKey()), entry.getLastUpdateTs(), value); if (commonAttributeKeys.contains(entry.getKey())) { - snapshot.putValue(new EntityKey(EntityKeyType.ATTRIBUTE, entry.getKey()), value); + snapshot.putValue(new EntityKey(EntityKeyType.ATTRIBUTE, entry.getKey()), entry.getLastUpdateTs(), value); } } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/EntityKeyValue.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/EntityKeyValue.java index 40ca323307..73a4db63b0 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/EntityKeyValue.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/EntityKeyValue.java @@ -15,9 +15,11 @@ */ package org.thingsboard.rule.engine.profile; +import lombok.EqualsAndHashCode; import lombok.Getter; import org.thingsboard.server.common.data.kv.DataType; +@EqualsAndHashCode class EntityKeyValue { @Getter diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/ProfileState.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/ProfileState.java new file mode 100644 index 0000000000..b758a96c39 --- /dev/null +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/ProfileState.java @@ -0,0 +1,137 @@ +/** + * Copyright © 2016-2020 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.rule.engine.profile; + +import lombok.AccessLevel; +import lombok.Getter; +import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.alarm.AlarmSeverity; +import org.thingsboard.server.common.data.device.profile.AlarmRule; +import org.thingsboard.server.common.data.device.profile.DeviceProfileAlarm; +import org.thingsboard.server.common.data.id.DeviceProfileId; +import org.thingsboard.server.common.data.query.ComplexFilterPredicate; +import org.thingsboard.server.common.data.query.DynamicValue; +import org.thingsboard.server.common.data.query.DynamicValueSourceType; +import org.thingsboard.server.common.data.query.EntityKey; +import org.thingsboard.server.common.data.query.EntityKeyType; +import org.thingsboard.server.common.data.query.FilterPredicateValue; +import org.thingsboard.server.common.data.query.KeyFilter; +import org.thingsboard.server.common.data.query.KeyFilterPredicate; +import org.thingsboard.server.common.data.query.SimpleKeyFilterPredicate; +import org.thingsboard.server.common.data.query.StringFilterPredicate; + +import javax.print.attribute.standard.Severity; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CopyOnWriteArrayList; + + +class ProfileState { + + private DeviceProfile deviceProfile; + @Getter(AccessLevel.PACKAGE) + private final List alarmSettings = new CopyOnWriteArrayList<>(); + @Getter(AccessLevel.PACKAGE) + private final Set entityKeys = ConcurrentHashMap.newKeySet(); + + private final Map>> alarmCreateKeys = new HashMap<>(); + private final Map> alarmClearKeys = new HashMap<>(); + + ProfileState(DeviceProfile deviceProfile) { + updateDeviceProfile(deviceProfile); + } + + void updateDeviceProfile(DeviceProfile deviceProfile) { + this.deviceProfile = deviceProfile; + alarmSettings.clear(); + alarmCreateKeys.clear(); + alarmClearKeys.clear(); + if (deviceProfile.getProfileData().getAlarms() != null) { + alarmSettings.addAll(deviceProfile.getProfileData().getAlarms()); + for (DeviceProfileAlarm alarm : deviceProfile.getProfileData().getAlarms()) { + Map> createAlarmKeys = alarmCreateKeys.computeIfAbsent(alarm.getId(), id -> new HashMap<>()); + alarm.getCreateRules().forEach(((severity, alarmRule) -> { + Set ruleKeys = createAlarmKeys.computeIfAbsent(severity, id -> new HashSet<>()); + for (KeyFilter keyFilter : alarmRule.getCondition().getCondition()) { + entityKeys.add(keyFilter.getKey()); + ruleKeys.add(keyFilter.getKey()); + addDynamicValuesRecursively(keyFilter.getPredicate(), entityKeys, ruleKeys); + } + })); + if (alarm.getClearRule() != null) { + Set clearAlarmKeys = alarmClearKeys.computeIfAbsent(alarm.getId(), id -> new HashSet<>()); + for (KeyFilter keyFilter : alarm.getClearRule().getCondition().getCondition()) { + entityKeys.add(keyFilter.getKey()); + clearAlarmKeys.add(keyFilter.getKey()); + addDynamicValuesRecursively(keyFilter.getPredicate(), entityKeys, clearAlarmKeys); + } + } + } + } + } + + private void addDynamicValuesRecursively(KeyFilterPredicate predicate, Set entityKeys, Set ruleKeys) { + switch (predicate.getType()) { + case STRING: + case NUMERIC: + case BOOLEAN: + DynamicValue value = ((SimpleKeyFilterPredicate) predicate).getValue().getDynamicValue(); + if (value != null && value.getSourceType() == DynamicValueSourceType.CURRENT_DEVICE) { + EntityKey entityKey = new EntityKey(EntityKeyType.ATTRIBUTE, value.getSourceAttribute()); + entityKeys.add(entityKey); + ruleKeys.add(entityKey); + } + break; + case COMPLEX: + for (KeyFilterPredicate child : ((ComplexFilterPredicate) predicate).getPredicates()) { + addDynamicValuesRecursively(child, entityKeys, ruleKeys); + } + break; + } + } + + DeviceProfileId getProfileId() { + return deviceProfile.getId(); + } + + Set getCreateAlarmKeys(String id, AlarmSeverity severity) { + Map> sKeys = alarmCreateKeys.get(id); + if (sKeys == null) { + return Collections.emptySet(); + } else { + Set keys = sKeys.get(severity); + if (keys == null) { + return Collections.emptySet(); + } else { + return keys; + } + } + } + + Set getClearAlarmKeys(String id) { + Set keys = alarmClearKeys.get(id); + if (keys == null) { + return Collections.emptySet(); + } else { + return keys; + } + } +} diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/EntityKeyState.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/SnapshotUpdate.java similarity index 58% rename from rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/EntityKeyState.java rename to rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/SnapshotUpdate.java index 08929bd2a2..e52a784c9c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/EntityKeyState.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/SnapshotUpdate.java @@ -15,8 +15,25 @@ */ package org.thingsboard.rule.engine.profile; -public class EntityKeyState { +import lombok.Getter; +import org.thingsboard.server.common.data.query.EntityKey; +import org.thingsboard.server.common.data.query.EntityKeyType; +import java.util.Set; +class SnapshotUpdate { + @Getter + private final EntityKeyType type; + @Getter + private final Set keys; + + SnapshotUpdate(EntityKeyType type, Set keys) { + this.type = type; + this.keys = keys; + } + + boolean hasUpdate(){ + return !keys.isEmpty(); + } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNode.java index 54a0fa7085..4b0b87043a 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNode.java @@ -23,7 +23,6 @@ import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; -import org.thingsboard.rule.engine.profile.state.PersistedDeviceState; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; @@ -36,11 +35,10 @@ import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.rule.RuleNodeState; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.common.msg.queue.PartitionChangeMsg; import org.thingsboard.server.dao.util.mapping.JacksonUtil; -import java.util.HashMap; import java.util.Map; -import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; @@ -70,11 +68,14 @@ public class TbDeviceProfileNode implements TbNode { this.cache = ctx.getDeviceProfileCache(); scheduleAlarmHarvesting(ctx); if (config.isFetchAlarmRulesStateOnStart()) { + log.info("[{}] Fetching alarm rule state", ctx.getSelfId()); + int fetchCount = 0; PageLink pageLink = new PageLink(1024); while (true) { PageData states = ctx.findRuleNodeStates(pageLink); if (!states.getData().isEmpty()) { for (RuleNodeState rns : states.getData()) { + fetchCount++; if (rns.getEntityId().getEntityType().equals(EntityType.DEVICE) && ctx.isLocalEntity(rns.getEntityId())) { getOrCreateDeviceState(ctx, new DeviceId(rns.getEntityId().getId()), rns); } @@ -86,6 +87,11 @@ public class TbDeviceProfileNode implements TbNode { pageLink = pageLink.nextPageLink(); } } + log.info("[{}] Fetched alarm rule state for {} entities", ctx.getSelfId(), fetchCount); + } + if (!config.isPersistAlarmRulesState() && ctx.isLocalEntity(ctx.getSelfId())) { + log.info("[{}] Going to cleanup rule node states", ctx.getSelfId()); + ctx.clearRuleNodeStates(); } } @@ -114,11 +120,14 @@ public class TbDeviceProfileNode implements TbNode { } } } else if (EntityType.DEVICE_PROFILE.equals(originatorType)) { + log.info("[{}] Received device profile update notification: {}", ctx.getSelfId(), msg.getData()); if (msg.getType().equals("ENTITY_UPDATED")) { DeviceProfile deviceProfile = JacksonUtil.fromString(msg.getData(), DeviceProfile.class); - for (DeviceState state : deviceStates.values()) { - if (deviceProfile.getId().equals(state.getProfileId())) { - state.updateProfile(ctx, deviceProfile); + if (deviceProfile != null) { + for (DeviceState state : deviceStates.values()) { + if (deviceProfile.getId().equals(state.getProfileId())) { + state.updateProfile(ctx, deviceProfile); + } } } } @@ -140,6 +149,12 @@ public class TbDeviceProfileNode implements TbNode { } } + @Override + public void onPartitionChangeMsg(TbContext ctx, PartitionChangeMsg msg) { + // Cleanup the cache for all entities that are no longer assigned to current server partitions + deviceStates.entrySet().removeIf(entry -> !ctx.isLocalEntity(entry.getKey())); + } + @Override public void destroy() { deviceStates.clear(); @@ -150,7 +165,7 @@ public class TbDeviceProfileNode implements TbNode { if (deviceState == null) { DeviceProfile deviceProfile = cache.get(ctx.getTenantId(), deviceId); if (deviceProfile != null) { - deviceState = new DeviceState(ctx, config, deviceId, new DeviceProfileState(deviceProfile), rns); + deviceState = new DeviceState(ctx, config, deviceId, new ProfileState(deviceProfile), rns); deviceStates.put(deviceId, deviceState); } }