Browse Source

Alarm rules CF: reevaluate rule on schedule start

pull/14193/head
VIacheslavKlimov 12 months ago
parent
commit
3ce5215022
  1. 4
      application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
  2. 25
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  3. 2
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  4. 9
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java
  5. 64
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmRuleState.java
  6. 3
      application/src/main/resources/thingsboard.yml
  7. 66
      application/src/test/java/org/thingsboard/server/cf/AlarmRulesTest.java
  8. 4
      common/data/src/main/java/org/thingsboard/server/common/data/alarm/rule/AlarmRule.java
  9. 5
      common/data/src/main/java/org/thingsboard/server/common/data/alarm/rule/condition/AlarmCondition.java
  10. 6
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AlarmCalculatedFieldConfiguration.java
  11. 4
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CalculatedFieldConfiguration.java

4
application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java

@ -664,6 +664,10 @@ public class ActorSystemContext {
@Getter
private long cfCalculationResultTimeout;
@Value("${actors.alarms.reevaluation_interval:120}")
@Getter
private long alarmRulesReevaluationInterval;
@Autowired
@Getter
private MqttClientSettings mqttClientSettings;

25
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java

@ -66,6 +66,8 @@ import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import java.util.function.BiConsumer;
import static org.thingsboard.server.utils.CalculatedFieldUtils.fromProto;
@ -80,6 +82,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
private final Map<EntityId, List<CalculatedFieldCtx>> entityIdCalculatedFields = new HashMap<>();
private final Map<EntityId, List<CalculatedFieldLink>> entityIdCalculatedFieldLinks = new HashMap<>();
private final Map<EntityId, Set<EntityId>> ownerEntities = new HashMap<>();
private ScheduledFuture<?> cfsReevaluationTask;
private final CalculatedFieldProcessingService cfExecService;
private final CalculatedFieldStateService cfStateService;
@ -122,6 +125,10 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
calculatedFields.clear();
entityIdCalculatedFields.clear();
entityIdCalculatedFieldLinks.clear();
if (cfsReevaluationTask != null) {
cfsReevaluationTask.cancel(true);
cfsReevaluationTask = null;
}
ctx.stop(ctx.getSelf());
}
@ -129,6 +136,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
log.debug("[{}] Processing CF actor init message.", msg.getTenantId().getId());
initEntitiesCache();
initCalculatedFields();
scheduleCfsReevaluation();
msg.getCallback().onSuccess();
}
@ -149,6 +157,23 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
ctx.broadcastToChildren(msg, true);
}
private void scheduleCfsReevaluation() {
cfsReevaluationTask = systemContext.getScheduler().scheduleWithFixedDelay(() -> {
try {
calculatedFields.values().forEach(cf -> {
if (cf.isRequiresScheduledReevaluation()) {
applyToTargetCfEntityActors(cf, TbCallback.EMPTY, (entityId, callback) -> {
log.debug("[{}][{}] Pushing scheduled CF reevaluate msg", entityId, cf.getCfId());
getOrCreateActor(entityId).tell(new CalculatedFieldReevaluateMsg(tenantId, cf));
});
}
});
} catch (Exception e) {
log.warn("[{}] Failed to trigger CFs reevaluation", tenantId, e);
}
}, systemContext.getAlarmRulesReevaluationInterval(), systemContext.getAlarmRulesReevaluationInterval(), TimeUnit.SECONDS);
}
public void onEntityLifecycleMsg(CalculatedFieldEntityLifecycleMsg msg) throws CalculatedFieldException {
log.debug("Processing entity lifecycle event: [{}] for entity: [{}]", msg.getData().getEvent(), msg.getData().getEntityId());
var entityType = msg.getData().getEntityId().getEntityType();

2
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java

@ -85,6 +85,7 @@ public class CalculatedFieldCtx {
private Output output;
private String expression;
private boolean useLatestTs;
private boolean requiresScheduledReevaluation;
private ActorSystemContext systemContext;
private TbelInvokeService tbelInvokeService;
@ -163,6 +164,7 @@ public class CalculatedFieldCtx {
if (calculatedField.getConfiguration() instanceof ScheduledUpdateSupportedCalculatedFieldConfiguration scheduledConfig) {
this.scheduledUpdateIntervalMillis = scheduledConfig.isScheduledUpdateEnabled() ? TimeUnit.SECONDS.toMillis(scheduledConfig.getScheduledUpdateInterval()) : -1L;
}
this.requiresScheduledReevaluation = calculatedField.getConfiguration().requiresScheduledReevaluation();
this.systemContext = systemContext;
this.tbelInvokeService = systemContext.getTbelInvokeService();
this.relationService = systemContext.getRelationService();

9
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java

@ -35,6 +35,7 @@ import org.thingsboard.server.common.data.alarm.AlarmCreateOrUpdateActiveRequest
import org.thingsboard.server.common.data.alarm.AlarmSeverity;
import org.thingsboard.server.common.data.alarm.AlarmUpdateRequest;
import org.thingsboard.server.common.data.alarm.rule.AlarmRule;
import org.thingsboard.server.common.data.alarm.rule.condition.AlarmCondition;
import org.thingsboard.server.common.data.alarm.rule.condition.AlarmConditionType;
import org.thingsboard.server.common.data.alarm.rule.condition.AlarmConditionValue;
import org.thingsboard.server.common.data.alarm.rule.condition.expression.AlarmConditionExpression;
@ -142,7 +143,7 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState {
initCurrentAlarm(ctx);
createOrClearAlarms(state -> {
if (state.getCondition().getType() == AlarmConditionType.DURATION) {
AlarmEvalResult evalResult = state.reeval(System.currentTimeMillis());
AlarmEvalResult evalResult = state.reeval(System.currentTimeMillis(), ctx);
if (evalResult.getStatus() == TRUE || evalResult.getStatus() == NOT_YET_TRUE) {
ScheduledFuture<?> future = ctx.scheduleReevaluation(evalResult.getLeftDuration(), actorCtx);
if (future != null) {
@ -161,7 +162,9 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState {
} else {
// when restored
ruleState.setAlarmRule(rule);
if (rule.getCondition().getType() == AlarmConditionType.DURATION && !ruleState.isEmpty()) {
ruleState.setActive(null);
AlarmCondition condition = rule.getCondition();
if (condition.hasSchedule() || (condition.getType() == AlarmConditionType.DURATION && !ruleState.isEmpty())) {
reevalNeeded.set(true);
}
}
@ -199,7 +202,7 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState {
}
return evalResult;
} else {
return state.reeval(System.currentTimeMillis());
return state.reeval(System.currentTimeMillis(), ctx);
}
}, ctx);
return Futures.immediateFuture(AlarmCalculatedFieldResult.builder()

64
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmRuleState.java

@ -58,6 +58,7 @@ public class AlarmRuleState {
private long lastEventTs;
private transient long duration;
private ScheduledFuture<?> durationCheckFuture;
private Boolean active;
public AlarmRuleState(AlarmSeverity severity, AlarmRule alarmRule, AlarmCalculatedFieldState state) {
this.severity = severity;
@ -68,32 +69,42 @@ public class AlarmRuleState {
}
public AlarmEvalResult eval(boolean newEvent, CalculatedFieldCtx ctx) { // on event or config change
boolean active = isActive(state.getLatestTimestamp());
return switch (condition.getType()) {
case SIMPLE -> evalSimple(active, ctx);
case DURATION -> evalDuration(active, ctx);
case REPEATING -> evalRepeating(active, newEvent, ctx);
};
long ts = newEvent ? state.getLatestTimestamp() : System.currentTimeMillis();
active = isActive(ts);
if (!active) {
return AlarmEvalResult.FALSE;
}
return doEval(newEvent, ctx);
}
public AlarmEvalResult reeval(long ts) {
public AlarmEvalResult reeval(long ts, CalculatedFieldCtx ctx) { // on scheduled duration check or periodic re-eval for rules with schedule
boolean active = isActive(ts);
switch (condition.getType()) {
case SIMPLE, REPEATING -> {
return AlarmEvalResult.NOT_YET_TRUE;
if (this.active == null || active != this.active) {
this.active = active;
if (active) {
return doEval(false, ctx);
}
}
if (active) {
return AlarmEvalResult.NOT_YET_TRUE;
} else {
return AlarmEvalResult.FALSE;
}
}
case DURATION -> {
if (!active) {
return AlarmEvalResult.FALSE;
}
long requiredDuration = getRequiredDurationInMs();
if (requiredDuration > 0 && lastEventTs > 0 && ts > lastEventTs) {
duration = ts - firstEventTs;
if (isActive(ts)) {
long leftDuration = requiredDuration - duration;
if (leftDuration <= 0) {
return AlarmEvalResult.TRUE;
} else {
return AlarmEvalResult.notYetTrue(0, leftDuration);
}
long leftDuration = requiredDuration - duration;
if (leftDuration <= 0) {
return AlarmEvalResult.TRUE;
} else {
return AlarmEvalResult.FALSE;
return AlarmEvalResult.notYetTrue(0, leftDuration);
}
}
}
@ -101,13 +112,20 @@ public class AlarmRuleState {
return AlarmEvalResult.FALSE;
}
private AlarmEvalResult evalSimple(boolean active, CalculatedFieldCtx ctx) {
return (active && eval(condition.getExpression(), ctx)) ?
AlarmEvalResult.TRUE : AlarmEvalResult.FALSE;
public AlarmEvalResult doEval(boolean newEvent, CalculatedFieldCtx ctx) {
return switch (condition.getType()) {
case SIMPLE -> evalSimple(ctx);
case DURATION -> evalDuration(ctx);
case REPEATING -> evalRepeating(newEvent, ctx);
};
}
private AlarmEvalResult evalSimple(CalculatedFieldCtx ctx) {
return eval(condition.getExpression(), ctx) ? AlarmEvalResult.TRUE : AlarmEvalResult.FALSE;
}
private AlarmEvalResult evalRepeating(boolean active, boolean newEvent, CalculatedFieldCtx ctx) {
if (active && eval(condition.getExpression(), ctx)) {
private AlarmEvalResult evalRepeating(boolean newEvent, CalculatedFieldCtx ctx) {
if (eval(condition.getExpression(), ctx)) {
if (newEvent) {
eventCount++;
}
@ -123,8 +141,8 @@ public class AlarmRuleState {
}
}
private AlarmEvalResult evalDuration(boolean active, CalculatedFieldCtx ctx) {
if (active && eval(condition.getExpression(), ctx)) {
private AlarmEvalResult evalDuration(CalculatedFieldCtx ctx) {
if (eval(condition.getExpression(), ctx)) {
long eventTs = state.getLatestTimestamp();
if (lastEventTs > 0) {
if (eventTs > lastEventTs) {

3
application/src/main/resources/thingsboard.yml

@ -529,6 +529,9 @@ actors:
configuration: "${ACTORS_CALCULATED_FIELD_DEBUG_MODE_RATE_LIMITS_PER_TENANT_CONFIGURATION:50000:3600}"
# Time in seconds to receive calculation result.
calculation_timeout: "${ACTORS_CALCULATION_TIMEOUT_SEC:5}"
alarms:
# Interval in seconds to re-evaluate Alarm rules that have a time schedule. 2 minutes by default.
reevaluation_interval: "${ACTORS_ALARMS_REEVALUATION_INTERVAL_SEC:120}"
debug:
settings:

66
application/src/test/java/org/thingsboard/server/cf/AlarmRulesTest.java

@ -21,6 +21,7 @@ import lombok.extern.slf4j.Slf4j;
import org.junit.Before;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.test.context.TestPropertySource;
import org.springframework.test.context.bean.override.mockito.MockitoSpyBean;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.action.TbAlarmResult;
@ -45,6 +46,7 @@ import org.thingsboard.server.common.data.alarm.rule.condition.expression.predic
import org.thingsboard.server.common.data.alarm.rule.condition.expression.predicate.StringFilterPredicate;
import org.thingsboard.server.common.data.alarm.rule.condition.expression.predicate.StringFilterPredicate.StringOperation;
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.AlarmSchedule;
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.SpecificTimeSchedule;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.configuration.AlarmCalculatedFieldConfiguration;
@ -63,18 +65,27 @@ import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.dao.event.EventDao;
import org.thingsboard.server.dao.service.DaoSqlTest;
import java.time.Duration;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.temporal.ChronoUnit;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
import java.util.function.Predicate;
import static org.assertj.core.api.Assertions.assertThat;
import static org.testcontainers.shaded.org.awaitility.Awaitility.await;
@Slf4j
@DaoSqlTest
@TestPropertySource(properties = {
"actors.alarms.reevaluation_interval=1"
})
public class AlarmRulesTest extends AbstractControllerTest {
@MockitoSpyBean
@ -209,10 +220,15 @@ public class AlarmRulesTest extends AbstractControllerTest {
assertThat(alarmResult.getConditionRepeats()).isEqualTo(5);
});
for (int i = 0; i < 5; i++) {
for (int i = 0; i < 4; i++) {
postTelemetry(deviceId, "{\"temperature\":50}");
Thread.sleep(10);
}
checkAlarmResult(calculatedField, alarmResult -> alarmResult.getConditionRepeats() == 9, alarmResult -> {
assertThat(alarmResult.isUpdated()).isTrue();
assertThat(alarmResult.getAlarm().getSeverity()).isEqualTo(AlarmSeverity.MAJOR);
});
postTelemetry(deviceId, "{\"temperature\":50}");
checkAlarmResult(calculatedField, alarmResult -> {
assertThat(alarmResult.isSeverityUpdated()).isTrue();
assertThat(alarmResult.getAlarm().getSeverity()).isEqualTo(AlarmSeverity.CRITICAL);
@ -420,7 +436,7 @@ public class AlarmRulesTest extends AbstractControllerTest {
scheduleArgument.setDefaultValue("None");
Map<String, Argument> arguments = Map.of(
"temperature", temperatureArgument,
"schedule", scheduleArgument // fixme:
"schedule", scheduleArgument
);
Map<AlarmSeverity, Condition> createRules = Map.of(
@ -638,11 +654,53 @@ public class AlarmRulesTest extends AbstractControllerTest {
});
}
@Test
public void testCreateAlarm_scheduleStarted() throws Exception {
Argument parkingSpotOccupiedArgument = new Argument();
parkingSpotOccupiedArgument.setRefEntityKey(new ReferencedEntityKey("parkingSpotOccupied", ArgumentType.ATTRIBUTE, AttributeScope.SERVER_SCOPE));
parkingSpotOccupiedArgument.setDefaultValue("false");
Map<String, Argument> arguments = Map.of(
"parkingSpotOccupied", parkingSpotOccupiedArgument
);
SpecificTimeSchedule schedule = new SpecificTimeSchedule();
schedule.setTimezone(ZoneId.systemDefault().getId());
schedule.setDaysOfWeek(Set.of(1, 2, 3, 4, 5, 6, 7));
long startsOn = Duration.between(LocalDate.now().atStartOfDay(), LocalDateTime.now())
.plus(15, ChronoUnit.SECONDS).toMillis();
schedule.setStartsOn(startsOn);
Map<AlarmSeverity, Condition> createRules = Map.of(
AlarmSeverity.CRITICAL, new Condition("return parkingSpotOccupied == true;", null, null, null,
new AlarmConditionValue<>(schedule, null))
);
CalculatedField calculatedField = createAlarmCf(deviceId, "Illegal parking alarm",
arguments, createRules, null);
postAttributes(deviceId, AttributeScope.SERVER_SCOPE, "{\"parkingSpotOccupied\":true}");
Thread.sleep(10000);
assertThat(getLatestAlarmResult(calculatedField.getId())).isNull();
checkAlarmResult(calculatedField, alarmResult -> {
assertThat(alarmResult.isCreated()).isTrue();
assertThat(alarmResult.getAlarm().getSeverity()).isEqualTo(AlarmSeverity.CRITICAL);
assertThat(alarmResult.getAlarm().getStatus()).isEqualTo(AlarmStatus.ACTIVE_UNACK);
});
}
// TODO: MSA tests
private void checkAlarmResult(CalculatedField calculatedField, Consumer<TbAlarmResult> assertion) {
checkAlarmResult(calculatedField, null, assertion);
}
private void checkAlarmResult(CalculatedField calculatedField,
Predicate<TbAlarmResult> waitFor,
Consumer<TbAlarmResult> assertion) {
TbAlarmResult alarmResult = await().atMost(TIMEOUT, TimeUnit.SECONDS)
.until(() -> getLatestAlarmResult(calculatedField.getId()), Objects::nonNull);
.until(() -> getLatestAlarmResult(calculatedField.getId()), result ->
result != null && (waitFor == null || waitFor.test(result)));
assertion.accept(alarmResult);
Alarm alarm = alarmResult.getAlarm();

4
common/data/src/main/java/org/thingsboard/server/common/data/alarm/rule/AlarmRule.java

@ -30,4 +30,8 @@ public class AlarmRule {
private String alarmDetails;
private DashboardId dashboardId;
public boolean requiresScheduledReevaluation() {
return condition.hasSchedule();
}
}

5
common/data/src/main/java/org/thingsboard/server/common/data/alarm/rule/condition/AlarmCondition.java

@ -26,6 +26,7 @@ import lombok.NoArgsConstructor;
import org.jetbrains.annotations.NotNull;
import org.thingsboard.server.common.data.alarm.rule.condition.expression.AlarmConditionExpression;
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.AlarmSchedule;
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.AnyTimeSchedule;
@JsonIgnoreProperties(ignoreUnknown = true)
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = "type")
@ -44,6 +45,10 @@ public abstract class AlarmCondition {
@Valid
private AlarmConditionValue<AlarmSchedule> schedule;
public boolean hasSchedule() {
return schedule != null && !(schedule.getStaticValue() instanceof AnyTimeSchedule);
}
@JsonIgnore
public abstract AlarmConditionType getType();

6
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AlarmCalculatedFieldConfiguration.java

@ -58,4 +58,10 @@ public class AlarmCalculatedFieldConfiguration implements ArgumentsBasedCalculat
}
@Override
public boolean requiresScheduledReevaluation() {
return createRules.values().stream().anyMatch(AlarmRule::requiresScheduledReevaluation) ||
(clearRule != null && clearRule.requiresScheduledReevaluation());
}
}

4
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CalculatedFieldConfiguration.java

@ -72,4 +72,8 @@ public interface CalculatedFieldConfiguration {
.collect(Collectors.toList());
}
default boolean requiresScheduledReevaluation() {
return false;
}
}

Loading…
Cancel
Save