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