diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java index c695ca84c6..84a768a886 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java @@ -42,51 +42,57 @@ 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 = alarmService.findAlarmById(tenantId, alarmId); + if (alarm == null) { + alarm = new Alarm(); + 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.createAlarm(AlarmCreateOrUpdateActiveRequest.fromAlarm(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())); + alarmService.createAlarm(AlarmCreateOrUpdateActiveRequest.fromAlarm(alarm, null, alarmId)); + 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); diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java index 1cf2cd44bb..4c73128f3b 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java @@ -25,6 +25,7 @@ import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmStatus; +import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; @@ -33,6 +34,7 @@ import org.thingsboard.server.gen.edge.v1.UplinkMsg; import java.util.List; import java.util.Optional; +import java.util.UUID; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @@ -42,8 +44,11 @@ abstract public class BaseAlarmEdgeTest extends AbstractEdgeTest { public void testSendAlarmToCloud() throws Exception { Device device = saveDeviceOnCloudAndVerifyDeliveryToEdge(); + UUID alarmUUID = UUID.randomUUID(); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); AlarmUpdateMsg.Builder alarmUpdateMgBuilder = AlarmUpdateMsg.newBuilder(); + alarmUpdateMgBuilder.setIdMSB(alarmUUID.getMostSignificantBits()); + alarmUpdateMgBuilder.setIdLSB(alarmUUID.getLeastSignificantBits()); alarmUpdateMgBuilder.setName("alarm from edge"); alarmUpdateMgBuilder.setStatus(AlarmStatus.ACTIVE_UNACK.name()); alarmUpdateMgBuilder.setSeverity(AlarmSeverity.CRITICAL.name()); @@ -65,6 +70,7 @@ abstract public class BaseAlarmEdgeTest extends AbstractEdgeTest { Optional foundAlarm = alarms.stream().filter(alarm -> alarm.getType().equals("alarm from edge")).findAny(); Assert.assertTrue(foundAlarm.isPresent()); AlarmInfo alarmInfo = foundAlarm.get(); + Assert.assertEquals(new AlarmId(alarmUUID), alarmInfo.getId()); Assert.assertEquals(device.getId(), alarmInfo.getOriginator()); Assert.assertEquals(AlarmStatus.ACTIVE_UNACK, alarmInfo.getStatus()); Assert.assertEquals(AlarmSeverity.CRITICAL, alarmInfo.getSeverity()); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/alarm/AlarmCreateOrUpdateActiveRequest.java b/common/data/src/main/java/org/thingsboard/server/common/data/alarm/AlarmCreateOrUpdateActiveRequest.java index ecf882e2c9..0f3dd5b636 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/alarm/AlarmCreateOrUpdateActiveRequest.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/alarm/AlarmCreateOrUpdateActiveRequest.java @@ -19,6 +19,7 @@ import com.fasterxml.jackson.databind.JsonNode; import io.swagger.annotations.ApiModelProperty; import lombok.Builder; import lombok.Data; +import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -61,11 +62,17 @@ public class AlarmCreateOrUpdateActiveRequest implements AlarmModificationReques private UserId userId; + private AlarmId edgeAlarmId; + public static AlarmCreateOrUpdateActiveRequest fromAlarm(Alarm a) { return fromAlarm(a, null); } public static AlarmCreateOrUpdateActiveRequest fromAlarm(Alarm a, UserId userId) { + return fromAlarm(a, userId, null); + } + + public static AlarmCreateOrUpdateActiveRequest fromAlarm(Alarm a, UserId userId, AlarmId edgeAlarmId) { return AlarmCreateOrUpdateActiveRequest.builder() .tenantId(a.getTenantId()) .customerId(a.getCustomerId()) @@ -81,6 +88,7 @@ public class AlarmCreateOrUpdateActiveRequest implements AlarmModificationReques .propagateToTenant(a.isPropagateToTenant()) .propagateRelationTypes(a.getPropagateRelationTypes()).build()) .userId(userId) + .edgeAlarmId(edgeAlarmId) .build(); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java index f1abd07d95..321208eb8f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java @@ -245,7 +245,7 @@ public class JpaAlarmDao extends JpaAbstractDao implements A return toAlarmApiResult(alarmRepository.createOrUpdateActiveAlarm( request.getTenantId().getId(), request.getCustomerId() != null ? request.getCustomerId().getId() : CustomerId.NULL_UUID, - UUID.randomUUID(), + request.getEdgeAlarmId() != null ? request.getEdgeAlarmId().getId() : UUID.randomUUID(), System.currentTimeMillis(), request.getOriginator().getId(), request.getOriginator().getEntityType().ordinal(),