Browse Source

Merge branch 'master' of github.com:thingsboard/thingsboard

pull/3576/head
Igor Kulikov 6 years ago
parent
commit
f1cc4fbd04
  1. 8
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  2. 2
      application/src/main/java/org/thingsboard/server/controller/RuleChainController.java
  3. 1
      common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleNodeStateService.java
  4. 3
      common/data/src/main/java/org/thingsboard/server/common/data/query/DynamicValueSourceType.java
  5. 11
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleNodeStateService.java
  6. 4
      dao/src/main/java/org/thingsboard/server/dao/rule/RuleNodeStateDao.java
  7. 14
      dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java
  8. 7
      dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeStateDao.java
  9. 3
      dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeStateRepository.java
  10. 2
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java
  11. 262
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmRuleState.java
  12. 49
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java
  13. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmStateUpdateResult.java
  14. 34
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DataSnapshot.java
  15. 63
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceProfileState.java
  16. 104
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceState.java
  17. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/EntityKeyValue.java
  18. 137
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/ProfileState.java
  19. 19
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/SnapshotUpdate.java
  20. 29
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/TbDeviceProfileNode.java

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

2
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;

1
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);
}

3
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
}

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

4
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<RuleNodeState> {
PageData<RuleNodeState> findByRuleNodeId(UUID ruleNodeId, PageLink pageLink);
RuleNodeState findByRuleNodeIdAndEntityId(UUID ruleNodeId, UUID entityId);
void removeByRuleNodeId(UUID ruleNodeId);
}

14
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(")");

7
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<RuleNodeStateEntity, Rul
public RuleNodeState findByRuleNodeIdAndEntityId(UUID ruleNodeId, UUID entityId) {
return DaoUtil.getData(ruleNodeStateRepository.findByRuleNodeIdAndEntityId(ruleNodeId, entityId));
}
@Transactional
@Override
public void removeByRuleNodeId(UUID ruleNodeId) {
ruleNodeStateRepository.removeByRuleNodeId(ruleNodeId);
}
}

3
dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeStateRepository.java

@ -33,4 +33,7 @@ public interface RuleNodeStateRepository extends PagingAndSortingRepository<Rule
@Query("SELECT e FROM RuleNodeStateEntity e WHERE e.ruleNodeId = :ruleNodeId and e.entityId = :entityId")
RuleNodeStateEntity findByRuleNodeIdAndEntityId(@Param("ruleNodeId") UUID ruleNodeId, @Param("entityId") UUID entityId);
void removeByRuleNodeId(@Param("ruleNodeId") UUID ruleNodeId);
}

2
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java

@ -221,4 +221,6 @@ public interface TbContext {
RuleNodeState findRuleNodeStateForEntity(EntityId entityId);
RuleNodeState saveRuleNodeState(RuleNodeState state);
void clearRuleNodeStates();
}

262
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmRuleState.java

@ -29,6 +29,9 @@ import org.thingsboard.server.common.data.device.profile.SimpleAlarmConditionSpe
import org.thingsboard.server.common.data.device.profile.SpecificTimeSchedule;
import org.thingsboard.server.common.data.query.BooleanFilterPredicate;
import org.thingsboard.server.common.data.query.ComplexFilterPredicate;
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.NumericFilterPredicate;
@ -38,22 +41,25 @@ import org.thingsboard.server.common.msg.tools.SchedulerUtils;
import java.time.Instant;
import java.time.ZoneId;
import java.time.ZonedDateTime;
import java.util.Calendar;
import java.util.Set;
import java.util.function.Function;
@Data
public class AlarmRuleState {
class AlarmRuleState {
private final AlarmSeverity severity;
private final AlarmRule alarmRule;
private final AlarmConditionSpec spec;
private final long requiredDurationInMs;
private final long requiredRepeats;
private final Set<EntityKey> entityKeys;
private PersistedAlarmRuleState state;
private boolean updateFlag;
public AlarmRuleState(AlarmSeverity severity, AlarmRule alarmRule, PersistedAlarmRuleState state) {
AlarmRuleState(AlarmSeverity severity, AlarmRule alarmRule, Set<EntityKey> 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<EntityKey> changedKeys) {
for (EntityKey key : changedKeys) {
if (entityKeys.contains(key)) {
return true;
}
}
return false;
}
public boolean validateAttrUpdate(Set<EntityKey> 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> T getPredicateValue(DataSnapshot data, FilterPredicateValue<T> value, Function<EntityKeyValue, T> transformFunction) {
EntityKeyValue ekv = getDynamicPredicateValue(data, value);
if (ekv != null) {
T result = transformFunction.apply(ekv);
if (result != null) {
return result;
}
}
return value.getDefaultValue();
}
private <T> EntityKeyValue getDynamicPredicateValue(DataSnapshot data, FilterPredicateValue<T> 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;
}
}
}

49
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceProfileAlarmState.java → 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<AlarmRuleState> 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 <T> boolean createOrClearAlarms(TbContext ctx, T data, BiFunction<AlarmRuleState, T, Boolean> evalFunction) {
public <T> boolean createOrClearAlarms(TbContext ctx, T data, SnapshotUpdate update, BiFunction<AlarmRuleState, T, Boolean> 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)) {

2
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;

34
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceDataSnapshot.java → 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<EntityKey> keys;
private final Map<EntityKey, EntityKeyValue> values = new ConcurrentHashMap<>();
public DeviceDataSnapshot(Set<EntityKey> entityKeysToFetch) {
DataSnapshot(Set<EntityKey> 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;
}
}

63
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceProfileState.java

@ -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<DeviceProfileAlarm> alarmSettings = new CopyOnWriteArrayList<>();
@Getter(AccessLevel.PACKAGE)
private final Set<EntityKey> 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();
}
}

104
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<String, DeviceProfileAlarmState> alarmStates = new ConcurrentHashMap<>();
private DataSnapshot latestValues;
private final ConcurrentMap<String, AlarmState> 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<AttributeKvEntry> 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<Long, List<KvEntry>> tsKvMap = JsonConverter.convertToSortedTelemetry(new JsonParser().parse(msg.getData()), TbMsgTimeseriesNode.getTs(msg));
// iterate over data by ts (ASC order).
for (Map.Entry<Long, List<KvEntry>> entry : tsKvMap.entrySet()) {
Long ts = entry.getKey();
List<KvEntry> 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<KvEntry> data) {
latestValues.setTs(ts);
private SnapshotUpdate merge(DataSnapshot latestValues, Long newTs, List<KvEntry> data) {
Set<EntityKey> 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<AttributeKvEntry> attributes, String scope) {
long ts = latestValues.getTs();
private SnapshotUpdate merge(DataSnapshot latestValues, Set<AttributeKvEntry> attributes, String scope) {
long newTs = 0;
Set<EntityKey> 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<EntityKey> 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<EntityKey> entityKeysToFetch, DeviceDataSnapshot result) throws InterruptedException, ExecutionException {
private void addEntityKeysToSnapshot(TbContext ctx, EntityId originator, Set<EntityKey> entityKeysToFetch, DataSnapshot result) throws InterruptedException, ExecutionException {
Set<String> serverAttributeKeys = new HashSet<>();
Set<String> clientAttributeKeys = new HashSet<>();
Set<String> 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<TsKvEntry> 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<String> commonAttributeKeys, List<AttributeKvEntry> data) {
private void addToSnapshot(DataSnapshot snapshot, Set<String> commonAttributeKeys, List<AttributeKvEntry> 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);
}
}
}

2
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

137
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<DeviceProfileAlarm> alarmSettings = new CopyOnWriteArrayList<>();
@Getter(AccessLevel.PACKAGE)
private final Set<EntityKey> entityKeys = ConcurrentHashMap.newKeySet();
private final Map<String, Map<AlarmSeverity, Set<EntityKey>>> alarmCreateKeys = new HashMap<>();
private final Map<String, Set<EntityKey>> 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<AlarmSeverity, Set<EntityKey>> createAlarmKeys = alarmCreateKeys.computeIfAbsent(alarm.getId(), id -> new HashMap<>());
alarm.getCreateRules().forEach(((severity, alarmRule) -> {
Set<EntityKey> 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<EntityKey> 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<EntityKey> entityKeys, Set<EntityKey> 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<EntityKey> getCreateAlarmKeys(String id, AlarmSeverity severity) {
Map<AlarmSeverity, Set<EntityKey>> sKeys = alarmCreateKeys.get(id);
if (sKeys == null) {
return Collections.emptySet();
} else {
Set<EntityKey> keys = sKeys.get(severity);
if (keys == null) {
return Collections.emptySet();
} else {
return keys;
}
}
}
Set<EntityKey> getClearAlarmKeys(String id) {
Set<EntityKey> keys = alarmClearKeys.get(id);
if (keys == null) {
return Collections.emptySet();
} else {
return keys;
}
}
}

19
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/EntityKeyState.java → 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<EntityKey> keys;
SnapshotUpdate(EntityKeyType type, Set<EntityKey> keys) {
this.type = type;
this.keys = keys;
}
boolean hasUpdate(){
return !keys.isEmpty();
}
}

29
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<RuleNodeState> 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);
}
}

Loading…
Cancel
Save