|
|
|
@ -19,19 +19,19 @@ import com.fasterxml.jackson.databind.JsonNode; |
|
|
|
import com.google.common.util.concurrent.Futures; |
|
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.springframework.stereotype.Component; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
import org.thingsboard.server.common.data.EntityType; |
|
|
|
import org.thingsboard.server.common.data.alarm.Alarm; |
|
|
|
import org.thingsboard.server.common.data.alarm.AlarmCreateOrUpdateActiveRequest; |
|
|
|
import org.thingsboard.server.common.data.alarm.AlarmSeverity; |
|
|
|
import org.thingsboard.server.common.data.alarm.AlarmStatus; |
|
|
|
import org.thingsboard.server.common.data.alarm.AlarmUpdateRequest; |
|
|
|
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|
|
|
import org.thingsboard.server.common.data.id.AlarmId; |
|
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; |
|
|
|
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
|
|
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
|
|
|
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
|
|
|
|
|
|
|
import java.util.UUID; |
|
|
|
@ -43,51 +43,59 @@ public abstract class BaseAlarmProcessor extends BaseEdgeProcessor { |
|
|
|
log.trace("[{}] processAlarmMsg [{}]", tenantId, alarmUpdateMsg); |
|
|
|
EntityId originatorId = getAlarmOriginator(tenantId, alarmUpdateMsg.getOriginatorName(), |
|
|
|
EntityType.valueOf(alarmUpdateMsg.getOriginatorType())); |
|
|
|
AlarmId alarmId = new AlarmId(new UUID(alarmUpdateMsg.getIdMSB(), alarmUpdateMsg.getIdLSB())); |
|
|
|
if (originatorId == null) { |
|
|
|
log.warn("Originator not found for the alarm msg {}", alarmUpdateMsg); |
|
|
|
return Futures.immediateFuture(null); |
|
|
|
} |
|
|
|
try { |
|
|
|
Alarm existentAlarm = alarmService.findLatestActiveByOriginatorAndType(tenantId, originatorId, alarmUpdateMsg.getType()); |
|
|
|
switch (alarmUpdateMsg.getMsgType()) { |
|
|
|
case ENTITY_CREATED_RPC_MESSAGE: |
|
|
|
case ENTITY_UPDATED_RPC_MESSAGE: |
|
|
|
if (existentAlarm == null || existentAlarm.getStatus().isCleared()) { |
|
|
|
existentAlarm = new Alarm(); |
|
|
|
existentAlarm.setTenantId(tenantId); |
|
|
|
existentAlarm.setType(alarmUpdateMsg.getName()); |
|
|
|
existentAlarm.setOriginator(originatorId); |
|
|
|
existentAlarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity())); |
|
|
|
existentAlarm.setStartTs(alarmUpdateMsg.getStartTs()); |
|
|
|
existentAlarm.setClearTs(alarmUpdateMsg.getClearTs()); |
|
|
|
existentAlarm.setPropagate(alarmUpdateMsg.getPropagate()); |
|
|
|
} |
|
|
|
Alarm alarm = new Alarm(); |
|
|
|
alarm.setId(alarmId); |
|
|
|
alarm.setTenantId(tenantId); |
|
|
|
alarm.setType(alarmUpdateMsg.getName()); |
|
|
|
alarm.setOriginator(originatorId); |
|
|
|
alarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity())); |
|
|
|
alarm.setStartTs(alarmUpdateMsg.getStartTs()); |
|
|
|
var alarmStatus = AlarmStatus.valueOf(alarmUpdateMsg.getStatus()); |
|
|
|
existentAlarm.setCleared(alarmStatus.isCleared()); |
|
|
|
existentAlarm.setAcknowledged(alarmStatus.isAck()); |
|
|
|
existentAlarm.setAckTs(alarmUpdateMsg.getAckTs()); |
|
|
|
existentAlarm.setEndTs(alarmUpdateMsg.getEndTs()); |
|
|
|
existentAlarm.setDetails(JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); |
|
|
|
alarmService.createOrUpdateAlarm(existentAlarm); |
|
|
|
break; |
|
|
|
alarm.setClearTs(alarmUpdateMsg.getClearTs()); |
|
|
|
alarm.setPropagate(alarmUpdateMsg.getPropagate()); |
|
|
|
alarm.setCleared(alarmStatus.isCleared()); |
|
|
|
alarm.setAcknowledged(alarmStatus.isAck()); |
|
|
|
alarm.setAckTs(alarmUpdateMsg.getAckTs()); |
|
|
|
alarm.setEndTs(alarmUpdateMsg.getEndTs()); |
|
|
|
alarm.setDetails(JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); |
|
|
|
if (UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE.equals(alarmUpdateMsg.getMsgType())) { |
|
|
|
alarmService.createAlarm(AlarmCreateOrUpdateActiveRequest.fromAlarm(alarm, null, alarmId)); |
|
|
|
} else { |
|
|
|
alarmService.updateAlarm(AlarmUpdateRequest.fromAlarm(alarm)); |
|
|
|
} |
|
|
|
return Futures.immediateFuture(null); |
|
|
|
case ALARM_ACK_RPC_MESSAGE: |
|
|
|
if (existentAlarm != null) { |
|
|
|
alarmService.acknowledgeAlarm(tenantId, existentAlarm.getId(), alarmUpdateMsg.getAckTs()); |
|
|
|
Alarm alarmToAck = alarmService.findAlarmById(tenantId, alarmId); |
|
|
|
if (alarmToAck != null) { |
|
|
|
alarmService.acknowledgeAlarm(tenantId, alarmId, alarmUpdateMsg.getAckTs()); |
|
|
|
} |
|
|
|
break; |
|
|
|
return Futures.immediateFuture(null); |
|
|
|
case ALARM_CLEAR_RPC_MESSAGE: |
|
|
|
if (existentAlarm != null) { |
|
|
|
alarmService.clearAlarm(tenantId, existentAlarm.getId(), |
|
|
|
alarmUpdateMsg.getAckTs(), JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); |
|
|
|
Alarm alarmToClear = alarmService.findAlarmById(tenantId, alarmId); |
|
|
|
if (alarmToClear != null) { |
|
|
|
alarmService.clearAlarm(tenantId, alarmId, alarmUpdateMsg.getClearTs(), |
|
|
|
JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); |
|
|
|
} |
|
|
|
break; |
|
|
|
return Futures.immediateFuture(null); |
|
|
|
case ENTITY_DELETED_RPC_MESSAGE: |
|
|
|
if (existentAlarm != null) { |
|
|
|
alarmService.delAlarm(tenantId, existentAlarm.getId()); |
|
|
|
Alarm alarmToDelete = alarmService.findAlarmById(tenantId, alarmId); |
|
|
|
if (alarmToDelete != null) { |
|
|
|
alarmService.delAlarm(tenantId, alarmId); |
|
|
|
} |
|
|
|
break; |
|
|
|
return Futures.immediateFuture(null); |
|
|
|
case UNRECOGNIZED: |
|
|
|
default: |
|
|
|
return handleUnsupportedMsgType(alarmUpdateMsg.getMsgType()); |
|
|
|
} |
|
|
|
return Futures.immediateFuture(null); |
|
|
|
} catch (Exception e) { |
|
|
|
log.error("[{}] Failed to process alarm update msg [{}]", tenantId, alarmUpdateMsg, e); |
|
|
|
return Futures.immediateFailedFuture(e); |
|
|
|
|