61 changed files with 1811 additions and 473 deletions
@ -0,0 +1,41 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.actors.calculatedField; |
||||
|
|
||||
|
import lombok.Builder; |
||||
|
import lombok.Data; |
||||
|
import org.thingsboard.server.common.data.alarm.Alarm; |
||||
|
import org.thingsboard.server.common.data.audit.ActionType; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.msg.MsgType; |
||||
|
import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; |
||||
|
import org.thingsboard.server.common.msg.queue.TbCallback; |
||||
|
|
||||
|
@Data |
||||
|
@Builder |
||||
|
public class CalculatedFieldAlarmActionMsg implements ToCalculatedFieldSystemMsg { |
||||
|
|
||||
|
private final TenantId tenantId; |
||||
|
private final Alarm alarm; |
||||
|
private final ActionType action; |
||||
|
private final TbCallback callback; |
||||
|
|
||||
|
@Override |
||||
|
public MsgType getMsgType() { |
||||
|
return MsgType.CF_ALARM_ACTION_MSG; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,57 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.actors.calculatedField; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import lombok.Builder; |
||||
|
import lombok.Data; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.audit.ActionType; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.msg.MsgType; |
||||
|
import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; |
||||
|
import org.thingsboard.server.common.msg.queue.TbCallback; |
||||
|
import org.thingsboard.server.common.util.ProtoUtils; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.EntityActionEventProto; |
||||
|
|
||||
|
@Data |
||||
|
@Builder |
||||
|
public class CalculatedFieldEntityActionEventMsg implements ToCalculatedFieldSystemMsg { |
||||
|
|
||||
|
private final TenantId tenantId; |
||||
|
private final EntityId entityId; |
||||
|
private final JsonNode entity; |
||||
|
private final ActionType action; |
||||
|
private final TbCallback callback; |
||||
|
|
||||
|
public static CalculatedFieldEntityActionEventMsg fromProto(EntityActionEventProto proto, |
||||
|
TbCallback callback) { |
||||
|
return CalculatedFieldEntityActionEventMsg.builder() |
||||
|
.tenantId((TenantId) ProtoUtils.fromProto(proto.getTenantId())) |
||||
|
.entityId(ProtoUtils.fromProto(proto.getEntityId())) |
||||
|
.entity(JacksonUtil.toJsonNode(proto.getEntity())) |
||||
|
.action(ActionType.valueOf(proto.getAction())) |
||||
|
.callback(callback) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public MsgType getMsgType() { |
||||
|
return MsgType.CF_ENTITY_ACTION_EVENT_MSG; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,80 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.cf; |
||||
|
|
||||
|
import lombok.Builder; |
||||
|
import lombok.Data; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.rule.engine.action.TbAlarmResult; |
||||
|
import org.thingsboard.server.common.data.DataConstants; |
||||
|
import org.thingsboard.server.common.data.id.CalculatedFieldId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.msg.TbMsgType; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.alarm.AlarmRuleState; |
||||
|
|
||||
|
import java.util.List; |
||||
|
|
||||
|
@Data |
||||
|
@Builder |
||||
|
public class AlarmCalculatedFieldResult implements CalculatedFieldResult { |
||||
|
|
||||
|
private final TbAlarmResult alarmResult; |
||||
|
private final AlarmRuleState alarmRuleState; |
||||
|
|
||||
|
@Override |
||||
|
public TbMsg toTbMsg(EntityId entityId, List<CalculatedFieldId> cfIds) { |
||||
|
TbMsgMetaData metaData = new TbMsgMetaData(); |
||||
|
if (alarmResult.isCreated()) { |
||||
|
metaData.putValue(DataConstants.IS_NEW_ALARM, Boolean.TRUE.toString()); |
||||
|
} else if (alarmResult.isUpdated()) { |
||||
|
metaData.putValue(DataConstants.IS_EXISTING_ALARM, Boolean.TRUE.toString()); |
||||
|
} else if (alarmResult.isSeverityUpdated()) { |
||||
|
metaData.putValue(DataConstants.IS_EXISTING_ALARM, Boolean.TRUE.toString()); |
||||
|
metaData.putValue(DataConstants.IS_SEVERITY_UPDATED_ALARM, Boolean.TRUE.toString()); |
||||
|
} else { |
||||
|
metaData.putValue(DataConstants.IS_CLEARED_ALARM, Boolean.TRUE.toString()); |
||||
|
} |
||||
|
switch (alarmRuleState.getCondition().getType()) { |
||||
|
case REPEATING -> { |
||||
|
metaData.putValue(DataConstants.ALARM_CONDITION_REPEATS, String.valueOf(alarmRuleState.getEventCount())); |
||||
|
} |
||||
|
case DURATION -> { |
||||
|
// TODO: schedule instead of duration
|
||||
|
metaData.putValue(DataConstants.ALARM_CONDITION_DURATION, String.valueOf(alarmRuleState.getDuration())); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
return TbMsg.newMsg() |
||||
|
.type(TbMsgType.ALARM) |
||||
|
.originator(entityId) |
||||
|
.data(JacksonUtil.toString(alarmResult.getAlarm())) |
||||
|
.metaData(metaData) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public String stringValue() { |
||||
|
return alarmResult != null ? JacksonUtil.toString(alarmResult) : null; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public boolean isEmpty() { |
||||
|
return alarmResult == null; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,74 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.cf; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import lombok.Builder; |
||||
|
import lombok.Data; |
||||
|
import org.thingsboard.server.common.data.AttributeScope; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.OutputType; |
||||
|
import org.thingsboard.server.common.data.id.CalculatedFieldId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.msg.TbMsgType; |
||||
|
import org.thingsboard.server.common.msg.TbMsg; |
||||
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
||||
|
|
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
|
||||
|
import static org.thingsboard.server.common.data.DataConstants.SCOPE; |
||||
|
|
||||
|
@Data |
||||
|
@Builder |
||||
|
public final class TelemetryCalculatedFieldResult implements CalculatedFieldResult { |
||||
|
|
||||
|
private final OutputType type; |
||||
|
private final AttributeScope scope; |
||||
|
private final JsonNode result; |
||||
|
|
||||
|
@Override |
||||
|
public TbMsg toTbMsg(EntityId entityId, List<CalculatedFieldId> cfIds) { |
||||
|
TbMsgType msgType = switch (type) { |
||||
|
case ATTRIBUTES -> TbMsgType.POST_ATTRIBUTES_REQUEST; |
||||
|
case TIME_SERIES -> TbMsgType.POST_TELEMETRY_REQUEST; |
||||
|
}; |
||||
|
TbMsgMetaData metaData = switch (type) { |
||||
|
case ATTRIBUTES -> new TbMsgMetaData(Map.of(SCOPE, scope.name())); |
||||
|
case TIME_SERIES -> TbMsgMetaData.EMPTY; |
||||
|
}; |
||||
|
return TbMsg.newMsg() |
||||
|
.type(msgType) |
||||
|
.originator(entityId) |
||||
|
.previousCalculatedFieldIds(cfIds) |
||||
|
.data(stringValue()) |
||||
|
.metaData(metaData) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public String stringValue() { |
||||
|
return result == null ? null : result.toString(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public boolean isEmpty() { |
||||
|
return result == null || result.isMissingNode() || result.isNull() || |
||||
|
(result.isObject() && result.isEmpty()) || |
||||
|
(result.isArray() && result.isEmpty()) || |
||||
|
(result.isTextual() && result.asText().isEmpty()); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,320 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.cf.ctx.state.alarm; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import com.fasterxml.jackson.databind.node.ObjectNode; |
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import lombok.EqualsAndHashCode; |
||||
|
import lombok.Getter; |
||||
|
import lombok.SneakyThrows; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.rule.engine.action.TbAlarmResult; |
||||
|
import org.thingsboard.server.common.data.StringUtils; |
||||
|
import org.thingsboard.server.common.data.alarm.Alarm; |
||||
|
import org.thingsboard.server.common.data.alarm.AlarmApiCallResult; |
||||
|
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.expression.AlarmConditionExpression; |
||||
|
import org.thingsboard.server.common.data.alarm.rule.condition.expression.TbelAlarmConditionExpression; |
||||
|
import org.thingsboard.server.common.data.audit.ActionType; |
||||
|
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.AlarmCalculatedFieldConfiguration; |
||||
|
import org.thingsboard.server.common.data.id.DashboardId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.service.cf.AlarmCalculatedFieldResult; |
||||
|
import org.thingsboard.server.service.cf.CalculatedFieldResult; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; |
||||
|
|
||||
|
import java.util.Comparator; |
||||
|
import java.util.Map; |
||||
|
import java.util.TreeMap; |
||||
|
import java.util.function.Function; |
||||
|
|
||||
|
@EqualsAndHashCode(callSuper = true) |
||||
|
@Slf4j |
||||
|
public class AlarmCalculatedFieldState extends BaseCalculatedFieldState { |
||||
|
|
||||
|
private String alarmType; |
||||
|
private AlarmCalculatedFieldConfiguration configuration; |
||||
|
|
||||
|
@Getter |
||||
|
private final Map<AlarmSeverity, AlarmRuleState> createRuleStates = new TreeMap<>(Comparator.comparing(Enum::ordinal)); |
||||
|
@Getter |
||||
|
private AlarmRuleState clearRuleState; |
||||
|
|
||||
|
@Getter |
||||
|
private Alarm currentAlarm; |
||||
|
private boolean initialFetchDone; |
||||
|
|
||||
|
public AlarmCalculatedFieldState(EntityId entityId) { |
||||
|
super(entityId); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void init(CalculatedFieldCtx ctx) { |
||||
|
super.init(ctx); |
||||
|
|
||||
|
this.alarmType = ctx.getCalculatedField().getName(); |
||||
|
this.configuration = getConfiguration(ctx); |
||||
|
|
||||
|
Map<AlarmSeverity, AlarmRule> createRules = configuration.getCreateRules(); |
||||
|
createRules.forEach((severity, rule) -> { |
||||
|
AlarmRuleState ruleState = createRuleStates.get(severity); |
||||
|
if (ruleState == null) { |
||||
|
ruleState = new AlarmRuleState(severity, rule, this); |
||||
|
createRuleStates.put(severity, ruleState); |
||||
|
} else { // can be null if was restored
|
||||
|
ruleState.setAlarmRule(rule); |
||||
|
// todo: is it enough to just set new alarm rule to alarm rule state? is it ok to leave the state as were??
|
||||
|
} |
||||
|
}); |
||||
|
createRuleStates.keySet().removeIf(severity -> !createRules.containsKey(severity)); |
||||
|
|
||||
|
AlarmRule clearRule = configuration.getClearRule(); |
||||
|
if (clearRule != null) { |
||||
|
if (clearRuleState == null) { |
||||
|
clearRuleState = new AlarmRuleState(null, clearRule, this); |
||||
|
} else { |
||||
|
clearRuleState.setAlarmRule(clearRule); |
||||
|
} |
||||
|
} else { |
||||
|
clearRuleState = null; |
||||
|
} |
||||
|
log.debug("Initialized create rule states {} and clear rule state {} for {}", createRuleStates, clearRuleState, ctx.getCalculatedField()); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void reset(CalculatedFieldCtx ctx) { |
||||
|
super.reset(ctx); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ListenableFuture<CalculatedFieldResult> performCalculation(CalculatedFieldCtx ctx) { |
||||
|
initCurrentAlarm(ctx); |
||||
|
AlarmCalculatedFieldResult result = createOrClearAlarms(state -> state.eval(ctx), ctx); |
||||
|
return Futures.immediateFuture(result); |
||||
|
} |
||||
|
|
||||
|
// TODO: harvesting
|
||||
|
public ListenableFuture<CalculatedFieldResult> performCalculation(long ts, CalculatedFieldCtx ctx) { |
||||
|
initCurrentAlarm(ctx); |
||||
|
AlarmCalculatedFieldResult result = createOrClearAlarms(ruleState -> ruleState.eval(ts), ctx); |
||||
|
return Futures.immediateFuture(result); |
||||
|
} |
||||
|
|
||||
|
@SneakyThrows |
||||
|
public boolean eval(AlarmConditionExpression expression, CalculatedFieldCtx ctx) { |
||||
|
if (expression instanceof TbelAlarmConditionExpression tbelExpression) { |
||||
|
Object result = ctx.evaluateTbelExpression(tbelExpression.getExpression(), this).get(); |
||||
|
if (result instanceof Boolean booleanResult) { |
||||
|
return booleanResult; |
||||
|
} else { |
||||
|
throw new IllegalStateException("Condition expression returned non-boolean value: '" + result + "'"); |
||||
|
} |
||||
|
} else { |
||||
|
throw new UnsupportedOperationException("Simple expressions not supported"); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public void processAlarmAction(Alarm alarm, ActionType action) { |
||||
|
switch (action) { |
||||
|
case ALARM_ACK -> processAlarmAck(alarm); |
||||
|
case ALARM_CLEAR -> processAlarmClear(alarm); |
||||
|
case ALARM_DELETE -> processAlarmDelete(alarm); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private void processAlarmClear(Alarm alarm) { |
||||
|
currentAlarm = null; |
||||
|
createRuleStates.values().forEach(AlarmRuleState::clear); |
||||
|
} |
||||
|
|
||||
|
private void processAlarmAck(Alarm alarm) { |
||||
|
currentAlarm.setAcknowledged(alarm.isAcknowledged()); |
||||
|
currentAlarm.setAckTs(alarm.getAckTs()); |
||||
|
} |
||||
|
|
||||
|
private void processAlarmDelete(Alarm alarm) { |
||||
|
currentAlarm = null; |
||||
|
createRuleStates.values().forEach(AlarmRuleState::clear); |
||||
|
} |
||||
|
|
||||
|
public AlarmCalculatedFieldResult createOrClearAlarms(Function<AlarmRuleState, AlarmEvalResult> evalFunction, CalculatedFieldCtx ctx) { |
||||
|
TbAlarmResult result = null; |
||||
|
AlarmRuleState resultState = null; |
||||
|
for (AlarmRuleState state : createRuleStates.values()) { |
||||
|
AlarmEvalResult evalResult = evalFunction.apply(state); |
||||
|
log.debug("Evaluated create rule {} with args {}. Result: {}", state, arguments, evalResult); |
||||
|
if (AlarmEvalResult.TRUE.equals(evalResult)) { |
||||
|
resultState = state; |
||||
|
break; |
||||
|
} else if (AlarmEvalResult.FALSE.equals(evalResult)) { |
||||
|
clearAlarmState(state); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
if (resultState != null) { |
||||
|
result = calculateAlarmResult(resultState, ctx); |
||||
|
log.debug("Alarm result for state {}: {}", resultState, result); |
||||
|
clearAlarmState(clearRuleState); |
||||
|
} else if (currentAlarm != null && clearRuleState != null) { |
||||
|
AlarmEvalResult evalResult = evalFunction.apply(clearRuleState); |
||||
|
log.debug("Evaluated clear rule {} with args {}. Result: {}", clearRuleState, arguments, evalResult); |
||||
|
if (AlarmEvalResult.TRUE.equals(evalResult)) { |
||||
|
clearAlarmState(clearRuleState); |
||||
|
for (AlarmRuleState state : createRuleStates.values()) { |
||||
|
clearAlarmState(state); |
||||
|
} |
||||
|
AlarmApiCallResult clearResult = ctx.getAlarmService().clearAlarm( |
||||
|
ctx.getTenantId(), currentAlarm.getId(), System.currentTimeMillis(), createDetails(clearRuleState), true |
||||
|
); |
||||
|
if (clearResult.isCleared()) { |
||||
|
result = new TbAlarmResult(false, false, true, clearResult.getAlarm()); |
||||
|
resultState = clearRuleState; |
||||
|
} |
||||
|
currentAlarm = null; |
||||
|
} else if (AlarmEvalResult.FALSE.equals(evalResult)) { |
||||
|
clearAlarmState(clearRuleState); |
||||
|
} |
||||
|
} |
||||
|
return AlarmCalculatedFieldResult.builder() |
||||
|
.alarmResult(result) |
||||
|
.alarmRuleState(resultState) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
private void clearAlarmState(AlarmRuleState state) { |
||||
|
if (state != null) { |
||||
|
state.clear(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private void initCurrentAlarm(CalculatedFieldCtx ctx) { |
||||
|
if (!initialFetchDone) { |
||||
|
Alarm alarm = ctx.getAlarmService().findLatestActiveByOriginatorAndType(ctx.getTenantId(), entityId, alarmType); |
||||
|
if (alarm != null && !alarm.getStatus().isCleared()) { |
||||
|
currentAlarm = alarm; |
||||
|
} |
||||
|
initialFetchDone = true; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private TbAlarmResult calculateAlarmResult(AlarmRuleState ruleState, CalculatedFieldCtx ctx) { |
||||
|
AlarmSeverity severity = ruleState.getSeverity(); |
||||
|
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(); |
||||
|
// Skip update if severity is decreased.
|
||||
|
if (severity.ordinal() <= oldSeverity.ordinal()) { |
||||
|
currentAlarm.setDetails(createDetails(ruleState)); |
||||
|
currentAlarm.setSeverity(severity); |
||||
|
AlarmApiCallResult result = ctx.getAlarmService().updateAlarm(AlarmUpdateRequest.fromAlarm(currentAlarm)); |
||||
|
currentAlarm = result.getAlarm(); |
||||
|
return TbAlarmResult.fromAlarmResult(result); |
||||
|
} else { |
||||
|
return null; |
||||
|
} |
||||
|
} else { |
||||
|
var newAlarm = new Alarm(); |
||||
|
newAlarm.setType(alarmType); |
||||
|
newAlarm.setAcknowledged(false); |
||||
|
newAlarm.setCleared(false); |
||||
|
newAlarm.setSeverity(severity); |
||||
|
long startTs = latestTimestamp; |
||||
|
long currentTime = System.currentTimeMillis(); |
||||
|
if (startTs == 0L || startTs > currentTime) { |
||||
|
startTs = currentTime; |
||||
|
} |
||||
|
newAlarm.setStartTs(startTs); |
||||
|
newAlarm.setEndTs(startTs); |
||||
|
newAlarm.setDetails(createDetails(ruleState)); |
||||
|
newAlarm.setOriginator(entityId); |
||||
|
newAlarm.setTenantId(ctx.getTenantId()); |
||||
|
newAlarm.setPropagate(configuration.isPropagate()); |
||||
|
newAlarm.setPropagateToOwner(configuration.isPropagateToOwner()); |
||||
|
newAlarm.setPropagateToTenant(configuration.isPropagateToTenant()); |
||||
|
if (configuration.getPropagateRelationTypes() != null) { |
||||
|
newAlarm.setPropagateRelationTypes(configuration.getPropagateRelationTypes()); |
||||
|
} |
||||
|
AlarmApiCallResult result = ctx.getAlarmService().createAlarm(AlarmCreateOrUpdateActiveRequest.fromAlarm(newAlarm)); |
||||
|
currentAlarm = result.getAlarm(); |
||||
|
return TbAlarmResult.fromAlarmResult(result); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private JsonNode createDetails(AlarmRuleState ruleState) { |
||||
|
JsonNode alarmDetails; |
||||
|
String alarmDetailsStr = ruleState.getAlarmRule().getAlarmDetails(); |
||||
|
DashboardId dashboardId = ruleState.getAlarmRule().getDashboardId(); |
||||
|
|
||||
|
if (StringUtils.isNotEmpty(alarmDetailsStr) || dashboardId != null) { |
||||
|
ObjectNode newDetails = JacksonUtil.newObjectNode(); |
||||
|
if (StringUtils.isNotEmpty(alarmDetailsStr)) { |
||||
|
for (Map.Entry<String, ArgumentEntry> entry : arguments.entrySet()) { |
||||
|
String key = entry.getKey(); |
||||
|
ArgumentEntry value = entry.getValue(); |
||||
|
alarmDetailsStr = alarmDetailsStr.replaceAll(String.format("\\$\\{%s}", key), String.valueOf(value.getValue())); |
||||
|
} |
||||
|
newDetails.put("data", alarmDetailsStr); |
||||
|
} |
||||
|
if (dashboardId != null) { |
||||
|
newDetails.put("dashboardId", dashboardId.getId().toString()); |
||||
|
} |
||||
|
alarmDetails = newDetails; |
||||
|
} else if (currentAlarm != null) { |
||||
|
alarmDetails = currentAlarm.getDetails(); |
||||
|
} else { |
||||
|
alarmDetails = JacksonUtil.newObjectNode(); |
||||
|
} |
||||
|
|
||||
|
return alarmDetails; |
||||
|
} |
||||
|
|
||||
|
protected SingleValueArgumentEntry getArgument(String key) { |
||||
|
SingleValueArgumentEntry entry = (SingleValueArgumentEntry) arguments.get(key); |
||||
|
if (entry == null) { |
||||
|
throw new IllegalArgumentException("Argument '" + key + "' is missing"); |
||||
|
} |
||||
|
return entry; |
||||
|
} |
||||
|
|
||||
|
private AlarmCalculatedFieldConfiguration getConfiguration(CalculatedFieldCtx ctx) { |
||||
|
return (AlarmCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
protected void validateNewEntry(ArgumentEntry newEntry) { |
||||
|
if (!(newEntry instanceof SingleValueArgumentEntry)) { |
||||
|
throw new IllegalArgumentException("Only single value arguments supported"); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public CalculatedFieldType getType() { |
||||
|
return CalculatedFieldType.ALARM; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,233 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.cf.ctx.state.alarm; |
||||
|
|
||||
|
import lombok.Data; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.thingsboard.common.util.KvUtil; |
||||
|
import org.thingsboard.server.common.adaptor.JsonConverter; |
||||
|
import org.thingsboard.server.common.data.alarm.AlarmSeverity; |
||||
|
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.AlarmConditionValue; |
||||
|
import org.thingsboard.server.common.data.alarm.rule.condition.DurationAlarmCondition; |
||||
|
import org.thingsboard.server.common.data.alarm.rule.condition.RepeatingAlarmCondition; |
||||
|
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.CustomTimeSchedule; |
||||
|
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.CustomTimeScheduleItem; |
||||
|
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.SpecificTimeSchedule; |
||||
|
import org.thingsboard.server.common.data.kv.KvEntry; |
||||
|
import org.thingsboard.server.common.msg.tools.SchedulerUtils; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; |
||||
|
|
||||
|
import java.time.Instant; |
||||
|
import java.time.ZoneId; |
||||
|
import java.time.ZonedDateTime; |
||||
|
import java.util.Optional; |
||||
|
import java.util.function.Function; |
||||
|
|
||||
|
@Data |
||||
|
@Slf4j |
||||
|
public class AlarmRuleState { |
||||
|
|
||||
|
private final AlarmSeverity severity; |
||||
|
private AlarmRule alarmRule; |
||||
|
private AlarmCalculatedFieldState state; |
||||
|
|
||||
|
private AlarmCondition condition; |
||||
|
|
||||
|
private long lastEventTs; |
||||
|
private long duration; |
||||
|
private long eventCount; |
||||
|
|
||||
|
public AlarmRuleState(AlarmSeverity severity, AlarmRule alarmRule, AlarmCalculatedFieldState state) { |
||||
|
this.severity = severity; |
||||
|
if (alarmRule != null) { |
||||
|
setAlarmRule(alarmRule); |
||||
|
} |
||||
|
this.state = state; |
||||
|
} |
||||
|
|
||||
|
public AlarmEvalResult eval(CalculatedFieldCtx ctx) { |
||||
|
boolean active = isActive(state.getLatestTimestamp()); |
||||
|
return switch (condition.getType()) { |
||||
|
case SIMPLE -> (active && eval(condition.getExpression(), ctx)) ? AlarmEvalResult.TRUE : AlarmEvalResult.FALSE; |
||||
|
case DURATION -> evalDuration(active, ctx); |
||||
|
case REPEATING -> evalRepeating(active, ctx); |
||||
|
}; |
||||
|
} |
||||
|
|
||||
|
public AlarmEvalResult eval(long ts) { |
||||
|
switch (condition.getType()) { |
||||
|
case SIMPLE: |
||||
|
case REPEATING: |
||||
|
return AlarmEvalResult.NOT_YET_TRUE; |
||||
|
case DURATION: |
||||
|
long requiredDurationInMs = getRequiredDurationInMs(); |
||||
|
if (requiredDurationInMs > 0 && lastEventTs > 0 && ts > lastEventTs) { |
||||
|
long duration = this.duration + (ts - lastEventTs); |
||||
|
if (isActive(ts)) { |
||||
|
return duration > requiredDurationInMs ? AlarmEvalResult.TRUE : AlarmEvalResult.NOT_YET_TRUE; |
||||
|
} else { |
||||
|
return AlarmEvalResult.FALSE; |
||||
|
} |
||||
|
} |
||||
|
default: |
||||
|
return AlarmEvalResult.FALSE; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private boolean isActive(long eventTs) { |
||||
|
if (condition.getSchedule() == null) { |
||||
|
return true; |
||||
|
} |
||||
|
AlarmSchedule schedule = getValue(condition.getSchedule(), entry -> Optional.ofNullable(KvUtil.getStringValue(entry)) |
||||
|
.map(str -> JsonConverter.parse(str, AlarmSchedule.class)) |
||||
|
.orElse(null)); |
||||
|
return switch (schedule.getType()) { |
||||
|
case ANY_TIME -> true; |
||||
|
case SPECIFIC_TIME -> isActiveSpecific((SpecificTimeSchedule) schedule, eventTs); |
||||
|
case CUSTOM -> isActiveCustom((CustomTimeSchedule) schedule, eventTs); |
||||
|
}; |
||||
|
} |
||||
|
|
||||
|
private boolean isActiveSpecific(SpecificTimeSchedule schedule, long eventTs) { |
||||
|
ZoneId zoneId = SchedulerUtils.getZoneId(schedule.getTimezone()); |
||||
|
ZonedDateTime zdt = ZonedDateTime.ofInstant(Instant.ofEpochMilli(eventTs), zoneId); |
||||
|
if (schedule.getDaysOfWeek().size() != 7) { |
||||
|
int dayOfWeek = zdt.getDayOfWeek().getValue(); |
||||
|
if (!schedule.getDaysOfWeek().contains(dayOfWeek)) { |
||||
|
return false; |
||||
|
} |
||||
|
} |
||||
|
long endsOn = schedule.getEndsOn(); |
||||
|
if (endsOn == 0) { |
||||
|
// 24 hours in milliseconds
|
||||
|
endsOn = 86400000; |
||||
|
} |
||||
|
|
||||
|
return isActive(eventTs, zoneId, zdt, schedule.getStartsOn(), endsOn); |
||||
|
} |
||||
|
|
||||
|
private boolean isActiveCustom(CustomTimeSchedule schedule, long eventTs) { |
||||
|
ZoneId zoneId = SchedulerUtils.getZoneId(schedule.getTimezone()); |
||||
|
ZonedDateTime zdt = ZonedDateTime.ofInstant(Instant.ofEpochMilli(eventTs), zoneId); |
||||
|
int dayOfWeek = zdt.toLocalDate().getDayOfWeek().getValue(); |
||||
|
for (CustomTimeScheduleItem item : schedule.getItems()) { |
||||
|
if (item.getDayOfWeek() == dayOfWeek) { |
||||
|
if (item.isEnabled()) { |
||||
|
long endsOn = item.getEndsOn(); |
||||
|
if (endsOn == 0) { |
||||
|
// 24 hours in milliseconds
|
||||
|
endsOn = 86400000; |
||||
|
} |
||||
|
return isActive(eventTs, zoneId, zdt, item.getStartsOn(), endsOn); |
||||
|
} else { |
||||
|
return false; |
||||
|
} |
||||
|
} |
||||
|
} |
||||
|
return false; |
||||
|
} |
||||
|
|
||||
|
private boolean isActive(long eventTs, ZoneId zoneId, ZonedDateTime zdt, long startsOn, long endsOn) { |
||||
|
long startOfDay = zdt.toLocalDate().atStartOfDay(zoneId).toInstant().toEpochMilli(); |
||||
|
long msFromStartOfDay = eventTs - startOfDay; |
||||
|
if (startsOn <= endsOn) { |
||||
|
return startsOn <= msFromStartOfDay && endsOn > msFromStartOfDay; |
||||
|
} else { |
||||
|
return startsOn < msFromStartOfDay || (0 < msFromStartOfDay && msFromStartOfDay < endsOn); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public void clear() { |
||||
|
eventCount = 0L; |
||||
|
lastEventTs = 0L; |
||||
|
duration = 0L; |
||||
|
} |
||||
|
|
||||
|
private AlarmEvalResult evalRepeating(boolean active, CalculatedFieldCtx ctx) { |
||||
|
if (active && eval(condition.getExpression(), ctx)) { |
||||
|
eventCount++; |
||||
|
long requiredRepeats = getIntValue(((RepeatingAlarmCondition) condition).getCount()); |
||||
|
return eventCount >= requiredRepeats ? AlarmEvalResult.TRUE : AlarmEvalResult.NOT_YET_TRUE; |
||||
|
} else { |
||||
|
return AlarmEvalResult.FALSE; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private AlarmEvalResult evalDuration(boolean active, CalculatedFieldCtx ctx) { |
||||
|
if (active && eval(condition.getExpression(), ctx)) { |
||||
|
if (lastEventTs > 0) { |
||||
|
if (state.getLatestTimestamp() > lastEventTs) { |
||||
|
duration = duration + (state.getLatestTimestamp() - lastEventTs); |
||||
|
lastEventTs = state.getLatestTimestamp(); |
||||
|
} |
||||
|
} else { |
||||
|
lastEventTs = state.getLatestTimestamp(); |
||||
|
duration = 0L; |
||||
|
} |
||||
|
long requiredDurationInMs = getRequiredDurationInMs(); |
||||
|
return duration > requiredDurationInMs ? AlarmEvalResult.TRUE : AlarmEvalResult.NOT_YET_TRUE; |
||||
|
} else { |
||||
|
return AlarmEvalResult.FALSE; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private Integer getIntValue(AlarmConditionValue<Integer> value) { |
||||
|
return getValue(value, entry -> Optional.ofNullable(KvUtil.getLongValue(entry)).map(Long::intValue).orElse(null)); |
||||
|
} |
||||
|
|
||||
|
private long getRequiredDurationInMs() { |
||||
|
return getValue(((DurationAlarmCondition) condition).getValue(), KvUtil::getLongValue); |
||||
|
} |
||||
|
|
||||
|
private boolean eval(AlarmConditionExpression expression, CalculatedFieldCtx ctx) { |
||||
|
return state.eval(expression, ctx); |
||||
|
} |
||||
|
|
||||
|
private <T> T getValue(AlarmConditionValue<T> conditionValue, Function<KvEntry, T> mapper) { |
||||
|
T value = conditionValue.getStaticValue(); |
||||
|
if (value == null) { |
||||
|
String argument = conditionValue.getDynamicValueArgument(); |
||||
|
SingleValueArgumentEntry entry = state.getArgument(argument); |
||||
|
value = mapper.apply(entry.getKvEntryValue()); |
||||
|
if (value == null) { |
||||
|
throw new IllegalArgumentException("No value found for argument " + argument); |
||||
|
} |
||||
|
} |
||||
|
return value; |
||||
|
} |
||||
|
|
||||
|
public void setAlarmRule(AlarmRule alarmRule) { |
||||
|
this.alarmRule = alarmRule; |
||||
|
this.condition = alarmRule.getCondition(); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public String toString() { |
||||
|
return "AlarmRuleState{" + |
||||
|
"severity=" + severity + |
||||
|
", condition=" + condition + |
||||
|
", lastEventTs=" + lastEventTs + |
||||
|
", duration=" + duration + |
||||
|
", eventCount=" + eventCount + |
||||
|
'}'; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,191 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.cf; |
||||
|
|
||||
|
import org.assertj.core.api.Assertions; |
||||
|
import org.junit.Before; |
||||
|
import org.junit.Test; |
||||
|
import org.springframework.beans.factory.annotation.Autowired; |
||||
|
import org.springframework.test.context.bean.override.mockito.MockitoSpyBean; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.rule.engine.action.TbAlarmResult; |
||||
|
import org.thingsboard.server.actors.ActorSystemContext; |
||||
|
import org.thingsboard.server.common.data.Device; |
||||
|
import org.thingsboard.server.common.data.alarm.Alarm; |
||||
|
import org.thingsboard.server.common.data.alarm.AlarmSeverity; |
||||
|
import org.thingsboard.server.common.data.alarm.AlarmStatus; |
||||
|
import org.thingsboard.server.common.data.alarm.rule.AlarmRule; |
||||
|
import org.thingsboard.server.common.data.alarm.rule.condition.SimpleAlarmCondition; |
||||
|
import org.thingsboard.server.common.data.alarm.rule.condition.expression.TbelAlarmConditionExpression; |
||||
|
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; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.Argument; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; |
||||
|
import org.thingsboard.server.common.data.debug.DebugSettings; |
||||
|
import org.thingsboard.server.common.data.event.CalculatedFieldDebugEvent; |
||||
|
import org.thingsboard.server.common.data.event.EventType; |
||||
|
import org.thingsboard.server.common.data.id.CalculatedFieldId; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.EventId; |
||||
|
import org.thingsboard.server.controller.AbstractControllerTest; |
||||
|
import org.thingsboard.server.dao.event.EventDao; |
||||
|
import org.thingsboard.server.dao.service.DaoSqlTest; |
||||
|
|
||||
|
import java.util.HashMap; |
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
import java.util.concurrent.TimeUnit; |
||||
|
import java.util.function.Consumer; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.testcontainers.shaded.org.awaitility.Awaitility.await; |
||||
|
|
||||
|
@DaoSqlTest |
||||
|
public class AlarmRulesTest extends AbstractControllerTest { |
||||
|
|
||||
|
@MockitoSpyBean |
||||
|
private ActorSystemContext actorSystemContext; |
||||
|
|
||||
|
@Autowired |
||||
|
private EventDao eventDao; |
||||
|
|
||||
|
private DeviceId deviceId; |
||||
|
private EventId latestEventId; |
||||
|
|
||||
|
@Before |
||||
|
public void beforeEach() throws Exception { |
||||
|
loginTenantAdmin(); |
||||
|
Device device = createDevice("Device A", "aaa"); |
||||
|
deviceId = device.getId(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void testCreateAndSeverityUpdateAndClear() throws Exception { |
||||
|
Argument temperatureArgument = new Argument(); |
||||
|
temperatureArgument.setRefEntityKey(new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null)); |
||||
|
Map<String, Argument> arguments = Map.of( |
||||
|
"temperature", temperatureArgument |
||||
|
); |
||||
|
|
||||
|
Map<AlarmSeverity, String> createRules = Map.of( |
||||
|
AlarmSeverity.MAJOR, "return temperature >= 50;", |
||||
|
AlarmSeverity.CRITICAL, "return temperature >= 100;" |
||||
|
); |
||||
|
String clearRule = "return temperature <= 25;"; |
||||
|
CalculatedField calculatedField = createAlarmCf(deviceId, "High Temperature Alarm", |
||||
|
arguments, createRules, clearRule); |
||||
|
|
||||
|
postTelemetry(deviceId, "{\"temperature\":50}"); |
||||
|
checkAlarmResult(calculatedField, alarmResult -> { |
||||
|
assertThat(alarmResult.isCreated()).isTrue(); |
||||
|
assertThat(alarmResult.getAlarm().getSeverity()).isEqualTo(AlarmSeverity.MAJOR); |
||||
|
assertThat(alarmResult.getAlarm().getStatus()).isEqualTo(AlarmStatus.ACTIVE_UNACK); |
||||
|
}); |
||||
|
|
||||
|
postTelemetry(deviceId, "{\"temperature\":100}"); |
||||
|
checkAlarmResult(calculatedField, alarmResult -> { |
||||
|
assertThat(alarmResult.isSeverityUpdated()).isTrue(); |
||||
|
assertThat(alarmResult.getAlarm().getSeverity()).isEqualTo(AlarmSeverity.CRITICAL); |
||||
|
assertThat(alarmResult.getAlarm().getStatus()).isEqualTo(AlarmStatus.ACTIVE_UNACK); |
||||
|
}); |
||||
|
|
||||
|
postTelemetry(deviceId, "{\"temperature\":101}"); |
||||
|
checkAlarmResult(calculatedField, alarmResult -> { |
||||
|
assertThat(alarmResult.isUpdated()).isTrue(); |
||||
|
assertThat(alarmResult.getAlarm().getSeverity()).isEqualTo(AlarmSeverity.CRITICAL); |
||||
|
assertThat(alarmResult.getAlarm().getStatus()).isEqualTo(AlarmStatus.ACTIVE_UNACK); |
||||
|
}); |
||||
|
|
||||
|
postTelemetry(deviceId, "{\"temperature\":20}"); |
||||
|
checkAlarmResult(calculatedField, alarmResult -> { |
||||
|
assertThat(alarmResult.isCleared()).isTrue(); |
||||
|
assertThat(alarmResult.getAlarm().getSeverity()).isEqualTo(AlarmSeverity.CRITICAL); |
||||
|
assertThat(alarmResult.getAlarm().getStatus()).isEqualTo(AlarmStatus.CLEARED_UNACK); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
private void checkAlarmResult(CalculatedField calculatedField, Consumer<TbAlarmResult> assertion) { |
||||
|
await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> { |
||||
|
TbAlarmResult alarmResult = getLatestAlarmResult(calculatedField.getId()); |
||||
|
assertThat(alarmResult).isNotNull(); |
||||
|
assertion.accept(alarmResult); |
||||
|
|
||||
|
Alarm alarm = alarmResult.getAlarm(); |
||||
|
assertThat(alarm.getOriginator()).isEqualTo(deviceId); |
||||
|
assertThat(alarm.getType()).isEqualTo(calculatedField.getName()); |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
private TbAlarmResult getLatestAlarmResult(CalculatedFieldId calculatedFieldId) { |
||||
|
List<CalculatedFieldDebugEvent> debugEvents = getDebugEvents(calculatedFieldId, 1); |
||||
|
if (debugEvents.isEmpty()) { |
||||
|
return null; |
||||
|
} |
||||
|
CalculatedFieldDebugEvent debugEvent = debugEvents.get(0); |
||||
|
if (debugEvent.getError() != null) { |
||||
|
System.err.println("CF error: " + debugEvent.getError()); |
||||
|
Assertions.fail(); |
||||
|
} |
||||
|
if (debugEvent.getId().equals(latestEventId)) { |
||||
|
return null; |
||||
|
} |
||||
|
latestEventId = debugEvent.getId(); |
||||
|
return JacksonUtil.fromString(debugEvent.getResult(), TbAlarmResult.class); |
||||
|
} |
||||
|
|
||||
|
private CalculatedField createAlarmCf(EntityId entityId, |
||||
|
String alarmType, |
||||
|
Map<String, Argument> arguments, |
||||
|
Map<AlarmSeverity, String> createConditions, |
||||
|
String clearCondition) { |
||||
|
CalculatedField calculatedField = new CalculatedField(); |
||||
|
calculatedField.setEntityId(entityId); |
||||
|
calculatedField.setName(alarmType); |
||||
|
calculatedField.setType(CalculatedFieldType.ALARM); |
||||
|
AlarmCalculatedFieldConfiguration configuration = new AlarmCalculatedFieldConfiguration(); |
||||
|
configuration.setArguments(arguments); |
||||
|
configuration.setCreateRules(new HashMap<>()); |
||||
|
createConditions.forEach((severity, expression) -> { |
||||
|
configuration.getCreateRules().put(severity, toAlarmRule(expression)); |
||||
|
}); |
||||
|
configuration.setClearRule(toAlarmRule(clearCondition)); |
||||
|
calculatedField.setConfiguration(configuration); |
||||
|
calculatedField.setDebugSettings(DebugSettings.all()); |
||||
|
return saveCalculatedField(calculatedField); |
||||
|
} |
||||
|
|
||||
|
private AlarmRule toAlarmRule(String conditionExpression) { |
||||
|
if (conditionExpression == null) { |
||||
|
return null; |
||||
|
} |
||||
|
AlarmRule rule = new AlarmRule(); |
||||
|
SimpleAlarmCondition condition = new SimpleAlarmCondition(); |
||||
|
TbelAlarmConditionExpression expression = new TbelAlarmConditionExpression(); |
||||
|
expression.setExpression(conditionExpression); |
||||
|
condition.setExpression(expression); |
||||
|
rule.setCondition(condition); |
||||
|
return rule; |
||||
|
} |
||||
|
|
||||
|
private List<CalculatedFieldDebugEvent> getDebugEvents(CalculatedFieldId calculatedFieldId, int limit) { |
||||
|
return eventDao.findLatestEvents(tenantId.getId(), calculatedFieldId.getId(), EventType.DEBUG_CALCULATED_FIELD, limit).stream() |
||||
|
.map(e -> (CalculatedFieldDebugEvent) e).toList(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
Loading…
Reference in new issue