From ee6bc5ec4b8c341afff28d2a9104ec3960878321 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Fri, 9 Oct 2020 14:12:52 +0300 Subject: [PATCH 1/4] Fixed issue for last level query in case entities have more relations in the hierarchy below requested level --- .../dao/sql/query/DefaultEntityQueryRepository.java | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) 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..2d5c9cbda9 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,10 +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"); } @@ -556,7 +559,9 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { .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(")"); From ba96409e5d2efd2bfa7b9750a2a24ae46854b113 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Fri, 9 Oct 2020 18:03:09 +0300 Subject: [PATCH 2/4] Code review fixes --- .../server/dao/sql/query/DefaultEntityQueryRepository.java | 5 ----- 1 file changed, 5 deletions(-) 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 2d5c9cbda9..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 @@ -479,9 +479,6 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { .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() + ")"; } @@ -554,9 +551,7 @@ 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") From bebddfd1ae9972eedff35de66d85e3d86a0336a5 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 12 Oct 2020 12:11:53 +0300 Subject: [PATCH 3/4] Improved logic of alarm rules --- .../rule/engine/profile/AlarmRuleState.java | 43 ++++++-- ...ProfileAlarmState.java => AlarmState.java} | 43 ++++++-- .../profile/AlarmStateUpdateResult.java | 2 +- ...iceDataSnapshot.java => DataSnapshot.java} | 34 ++++-- .../rule/engine/profile/DeviceState.java | 104 +++++++++++------- .../rule/engine/profile/EntityKeyValue.java | 2 + ...iceProfileState.java => ProfileState.java} | 52 ++++++++- ...ntityKeyState.java => SnapshotUpdate.java} | 19 +++- .../engine/profile/TbDeviceProfileNode.java | 25 +++-- 9 files changed, 242 insertions(+), 82 deletions(-) rename rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/{DeviceProfileAlarmState.java => AlarmState.java} (81%) rename rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/{DeviceDataSnapshot.java => DataSnapshot.java} (70%) rename rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/{DeviceProfileState.java => ProfileState.java} (51%) rename rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/{EntityKeyState.java => SnapshotUpdate.java} (58%) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmRuleState.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmRuleState.java index e307505fcd..a4618bab63 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmRuleState.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmRuleState.java @@ -29,6 +29,8 @@ 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.KeyFilter; import org.thingsboard.server.common.data.query.KeyFilterPredicate; import org.thingsboard.server.common.data.query.NumericFilterPredicate; @@ -38,22 +40,24 @@ 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; @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 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 +80,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 +121,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 +195,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 +205,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,7 +239,7 @@ 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()); @@ -380,4 +408,5 @@ public class AlarmRuleState { throw new RuntimeException("Operation not supported: " + predicate.getOperation()); } } + } 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 81% 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..a56377d2be 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,6 +91,10 @@ 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(); @@ -92,6 +106,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,12 +163,13 @@ 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); } } 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/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/DeviceProfileState.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/ProfileState.java similarity index 51% rename from rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/DeviceProfileState.java rename to rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/ProfileState.java index fd9037624e..882a7d7c9a 100644 --- 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/ProfileState.java @@ -18,19 +18,25 @@ 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.EntityKey; import org.thingsboard.server.common.data.query.KeyFilter; +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 DeviceProfileState { +class ProfileState { private DeviceProfile deviceProfile; @Getter(AccessLevel.PACKAGE) @@ -38,26 +44,64 @@ class DeviceProfileState { @Getter(AccessLevel.PACKAGE) private final Set entityKeys = ConcurrentHashMap.newKeySet(); - DeviceProfileState(DeviceProfile deviceProfile) { + 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()) { - for (AlarmRule alarmRule : alarm.getCreateRules().values()) { + 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()); + } + })); + 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()); } } } } } - public DeviceProfileId getProfileId() { + 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 97ec088f10..a6fcd8fa64 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; @@ -68,11 +66,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); } @@ -84,6 +85,7 @@ public class TbDeviceProfileNode implements TbNode { pageLink = pageLink.nextPageLink(); } } + log.info("[{}] Fetched alarm rule state for {} entities", ctx.getSelfId(), fetchCount); } } @@ -112,11 +114,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); + } } } } @@ -138,6 +143,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(); @@ -148,7 +159,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); } } From 6d185f1c01090c3180b7312f1fdfbd258f8ba8ba Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 12 Oct 2020 14:29:20 +0300 Subject: [PATCH 4/4] Improvements to device profile rule node --- .../actors/ruleChain/DefaultTbContext.java | 8 + .../controller/RuleChainController.java | 2 + .../server/dao/rule/RuleNodeStateService.java | 1 + .../data/query/DynamicValueSourceType.java | 3 +- .../dao/rule/BaseRuleNodeStateService.java | 11 + .../server/dao/rule/RuleNodeStateDao.java | 4 + .../dao/sql/rule/JpaRuleNodeStateDao.java | 7 + .../dao/sql/rule/RuleNodeStateRepository.java | 3 + .../rule/engine/api/TbContext.java | 2 + .../rule/engine/profile/AlarmRuleState.java | 219 +++++++++++------- .../rule/engine/profile/AlarmState.java | 6 + .../rule/engine/profile/ProfileState.java | 30 +++ .../engine/profile/TbDeviceProfileNode.java | 4 + 13 files changed, 211 insertions(+), 89 deletions(-) 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/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 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: @@ -409,4 +357,99 @@ class AlarmRuleState { } } + 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/AlarmState.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java index a56377d2be..5fb2c2957c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java @@ -98,6 +98,10 @@ class AlarmState { 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; @@ -175,6 +179,8 @@ class AlarmState { 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/ProfileState.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/ProfileState.java index 882a7d7c9a..b758a96c39 100644 --- 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 @@ -22,8 +22,16 @@ 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; @@ -65,6 +73,7 @@ class ProfileState { for (KeyFilter keyFilter : alarmRule.getCondition().getCondition()) { entityKeys.add(keyFilter.getKey()); ruleKeys.add(keyFilter.getKey()); + addDynamicValuesRecursively(keyFilter.getPredicate(), entityKeys, ruleKeys); } })); if (alarm.getClearRule() != null) { @@ -72,12 +81,33 @@ class ProfileState { 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(); } 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 3344af7f21..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 @@ -89,6 +89,10 @@ public class TbDeviceProfileNode implements TbNode { } 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(); + } } /**