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 d85750f7d2..be64f475a2 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 @@ -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); 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..c06093535a 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java @@ -19,12 +19,14 @@ import com.fasterxml.jackson.core.type.TypeReference; import com.google.protobuf.AbstractMessage; import org.junit.Assert; import org.junit.Test; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntityType; 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 +35,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 +45,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 +71,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()); @@ -85,12 +92,28 @@ abstract public class BaseAlarmEdgeTest extends AbstractEdgeTest { Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg); AlarmUpdateMsg alarmUpdateMsg = (AlarmUpdateMsg) latestMessage; Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, alarmUpdateMsg.getMsgType()); + Assert.assertEquals(savedAlarm.getUuidId().getMostSignificantBits(), alarmUpdateMsg.getIdMSB()); + Assert.assertEquals(savedAlarm.getUuidId().getLeastSignificantBits(), alarmUpdateMsg.getIdLSB()); Assert.assertEquals(savedAlarm.getType(), alarmUpdateMsg.getType()); Assert.assertEquals(savedAlarm.getName(), alarmUpdateMsg.getName()); Assert.assertEquals(device.getName(), alarmUpdateMsg.getOriginatorName()); Assert.assertEquals(savedAlarm.getStatus().name(), alarmUpdateMsg.getStatus()); Assert.assertEquals(savedAlarm.getSeverity().name(), alarmUpdateMsg.getSeverity()); + // update alarm + String updatedDetails = "{\"testKey\":\"testValue\"}"; + savedAlarm.setDetails(JacksonUtil.OBJECT_MAPPER.readTree(updatedDetails)); + edgeImitator.expectMessageAmount(1); + savedAlarm = doPost("/api/alarm", savedAlarm, Alarm.class); + Assert.assertTrue(edgeImitator.waitForMessages()); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg); + alarmUpdateMsg = (AlarmUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, alarmUpdateMsg.getMsgType()); + Assert.assertEquals(savedAlarm.getUuidId().getMostSignificantBits(), alarmUpdateMsg.getIdMSB()); + Assert.assertEquals(savedAlarm.getUuidId().getLeastSignificantBits(), alarmUpdateMsg.getIdLSB()); + Assert.assertEquals(updatedDetails, alarmUpdateMsg.getDetails()); + // ack alarm edgeImitator.expectMessageAmount(1); doPost("/api/alarm/" + savedAlarm.getUuidId() + "/ack"); 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(), diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java index 6d8063b8b5..7be7dae2de 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java @@ -74,7 +74,8 @@ public abstract class AbstractTbMsgPushNode entityBody = new HashMap<>(); @@ -107,6 +108,20 @@ public abstract class AbstractTbMsgPushNode('root', 'rulechain.edge-template-root', '60px', + new EntityTableColumn('root', 'rulechain.edge-template-root', '70px', entity => { return checkBoxCell(entity.root); }), - new EntityTableColumn('assignToEdge', 'rulechain.assign-to-edge', '60px', + new EntityTableColumn('assignToEdge', 'rulechain.assign-to-edge', '70px', entity => { return checkBoxCell(this.isAutoAssignToEdgeRuleChain(entity)); })