diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java index 5883d04fe6..89241f9a21 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java @@ -79,7 +79,12 @@ public class EdgeEventSourcingListener { return; } log.trace("[{}] SaveEntityEvent called: {}", event.getTenantId(), event); - EdgeEventActionType action = Boolean.TRUE.equals(event.getAdded()) ? EdgeEventActionType.ADDED : EdgeEventActionType.UPDATED; + boolean isAdded = Boolean.TRUE.equals(event.getAdded()); + EdgeEventActionType action = isAdded ? EdgeEventActionType.ADDED : EdgeEventActionType.UPDATED; + if (event.getEntity() instanceof AlarmComment) { + processAlarmCommentEvent(event, isAdded); + return; + } tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), null, event.getEntityId(), null, null, action, edgeSynchronizationManager.getEdgeId().get()); } catch (Exception e) { @@ -87,22 +92,33 @@ public class EdgeEventSourcingListener { } } + private void processAlarmCommentEvent(SaveEntityEvent event, boolean added) { + EdgeEventActionType action = added ? EdgeEventActionType.ADDED_COMMENT : EdgeEventActionType.UPDATED_COMMENT; + tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), null, event.getEntityId(), + JacksonUtil.toString(event.getEntity()), EdgeEventType.ALARM_COMMENT, action, edgeSynchronizationManager.getEdgeId().get()); + } + @TransactionalEventListener(fallbackExecution = true) public void handleEvent(DeleteEntityEvent event) { try { log.trace("[{}] DeleteEntityEvent called: {}", event.getTenantId(), event); - EdgeEventType type = null; - if (event.getEntity() instanceof AlarmComment) { - type = EdgeEventType.ALARM_COMMENT; - } + EdgeEventType type = getEdgeEventTypeForEntityEvent(event.getEntity()); + EdgeEventActionType actionType = getEdgeEventActionTypeForEntityEvent(event.getEntity()); tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), null, event.getEntityId(), - JacksonUtil.toString(event.getEntity()), type, EdgeEventActionType.DELETED, + JacksonUtil.toString(event.getEntity()), type, actionType, edgeSynchronizationManager.getEdgeId().get()); } catch (Exception e) { log.error("[{}] failed to process DeleteEntityEvent: {}", event.getTenantId(), event, e); } } + private EdgeEventActionType getEdgeEventActionTypeForEntityEvent(Object entity) { + if (entity instanceof AlarmComment) { + return EdgeEventActionType.DELETED_COMMENT; + } + return EdgeEventActionType.DELETED; + } + @TransactionalEventListener(fallbackExecution = true) public void handleEvent(ActionEntityEvent event) { try { @@ -154,7 +170,7 @@ public class EdgeEventSourcingListener { cleanUpUserAdditionalInfo(user); return !user.equals(oldUser); } - } else if (entity instanceof AlarmApiCallResult || entity instanceof Alarm || entity instanceof AlarmComment) { + } else if (entity instanceof AlarmApiCallResult || entity instanceof Alarm) { return false; } // Default: If the entity doesn't match any of the conditions, consider it as valid. @@ -177,4 +193,18 @@ public class EdgeEventSourcingListener { } } } + + private EdgeEventType getEdgeEventTypeForEntityEvent(Object entity) { + if (entity instanceof AlarmComment) { + return EdgeEventType.ALARM_COMMENT; + } + return null; + } + + private String getBodyMsgForEntityEvent(Object entity) { + if (entity instanceof AlarmComment) { + return JacksonUtil.toString(entity); + } + return null; + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/AlarmEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/AlarmEdgeProcessor.java index e6889a1e6f..f83cbbba9c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/AlarmEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/AlarmEdgeProcessor.java @@ -26,15 +26,11 @@ import org.thingsboard.server.common.data.alarm.AlarmComment; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; -import org.thingsboard.server.common.data.id.AlarmCommentId; import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.page.PageData; -import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.data.page.PageDataIterableByTenantIdEntityId; -import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.gen.edge.v1.AlarmCommentUpdateMsg; import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; @@ -45,7 +41,6 @@ import org.thingsboard.server.service.edge.rpc.constructor.alarm.AlarmMsgConstru import java.util.ArrayList; import java.util.List; -import java.util.Optional; import java.util.UUID; @Slf4j @@ -88,29 +83,22 @@ public abstract class AlarmEdgeProcessor extends BaseAlarmProcessor implements A @Override public DownlinkMsg convertAlarmCommentEventToDownlink(EdgeEvent edgeEvent, EdgeVersion edgeVersion) { - AlarmCommentId alarmCommentId = new AlarmCommentId(edgeEvent.getEntityId()); UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); - AlarmComment alarmComment; switch (edgeEvent.getAction()) { case ADDED_COMMENT: case UPDATED_COMMENT: - alarmComment = alarmCommentService.findAlarmCommentById(edgeEvent.getTenantId(), alarmCommentId); - break; case DELETED_COMMENT: - alarmComment = JacksonUtil.convertValue(edgeEvent.getBody(), AlarmComment.class); - break; + AlarmComment alarmComment = JacksonUtil.convertValue(edgeEvent.getBody(), AlarmComment.class); + if (alarmComment != null) { + return DownlinkMsg.newBuilder() + .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) + .addAlarmCommentUpdateMsg(((AlarmMsgConstructor) alarmMsgConstructorFactory + .getMsgConstructorByEdgeVersion(edgeVersion)).constructAlarmCommentUpdatedMsg(msgType, alarmComment)) + .build(); + } default: return null; } - return Optional.ofNullable(alarmComment).map(comment -> buildAlarmCommentDownlinkMsg(msgType, comment, edgeVersion)).orElse(null); - } - - private DownlinkMsg buildAlarmCommentDownlinkMsg(UpdateMsgType msgType, AlarmComment alarmComment, EdgeVersion edgeVersion) { - return DownlinkMsg.newBuilder() - .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) - .addAlarmCommentUpdateMsg(((AlarmMsgConstructor) alarmMsgConstructorFactory - .getMsgConstructorByEdgeVersion(edgeVersion)).constructAlarmCommentUpdatedMsg(msgType, alarmComment)) - .build(); } public ListenableFuture processAlarmNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { @@ -145,17 +133,14 @@ public abstract class AlarmEdgeProcessor extends BaseAlarmProcessor implements A EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()); AlarmId alarmId = new AlarmId(new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); EdgeId originatorEdgeId = safeGetEdgeId(edgeNotificationMsg.getOriginatorEdgeIdMSB(), edgeNotificationMsg.getOriginatorEdgeIdLSB()); - if (EdgeEventActionType.DELETED.equals(actionType)) { - AlarmComment deletedAlarmComment = JacksonUtil.fromString(edgeNotificationMsg.getBody(), AlarmComment.class); - if (deletedAlarmComment == null) { - return Futures.immediateFuture(null); - } - Alarm alarmById = alarmService.findAlarmById(tenantId, new AlarmId(deletedAlarmComment.getAlarmId().getId())); - List> delFutures = pushEventToAllRelatedEdges(tenantId, alarmById.getOriginator(), - alarmId, actionType, JacksonUtil.valueToTree(deletedAlarmComment), originatorEdgeId, EdgeEventType.ALARM_COMMENT); - return Futures.transform(Futures.allAsList(delFutures), voids -> null, dbCallbackExecutorService); + AlarmComment alarmComment = JacksonUtil.fromString(edgeNotificationMsg.getBody(), AlarmComment.class); + if (alarmComment == null) { + return Futures.immediateFuture(null); } - return Futures.immediateFuture(null); + Alarm alarmById = alarmService.findAlarmById(tenantId, new AlarmId(alarmComment.getAlarmId().getId())); + List> delFutures = pushEventToAllRelatedEdges(tenantId, alarmById.getOriginator(), + alarmId, actionType, JacksonUtil.valueToTree(alarmComment), originatorEdgeId, EdgeEventType.ALARM_COMMENT); + return Futures.transform(Futures.allAsList(delFutures), voids -> null, dbCallbackExecutorService); } private List> pushEventToAllRelatedEdges(TenantId tenantId, EntityId originatorId, AlarmId alarmId, diff --git a/application/src/test/java/org/thingsboard/server/edge/AlarmEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/AlarmEdgeTest.java index 6e10507f1e..acd8cb1444 100644 --- a/application/src/test/java/org/thingsboard/server/edge/AlarmEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/AlarmEdgeTest.java @@ -18,6 +18,7 @@ package org.thingsboard.server.edge; import com.datastax.oss.driver.api.core.uuid.Uuids; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.node.TextNode; +import com.google.protobuf.AbstractMessage; import org.junit.Assert; import org.junit.Test; import org.thingsboard.common.util.JacksonUtil; @@ -44,6 +45,8 @@ import java.util.List; import java.util.Optional; import java.util.UUID; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; + @DaoSqlTest public class AlarmEdgeTest extends AbstractEdgeTest { @@ -80,6 +83,60 @@ public class AlarmEdgeTest extends AbstractEdgeTest { Assert.assertEquals(AlarmSeverity.CRITICAL, alarmInfo.getSeverity()); } + @Test + public void testAlarms() throws Exception { + // create alarm + Device device = findDeviceByName("Edge Device 1"); + Alarm alarm = new Alarm(); + alarm.setOriginator(device.getId()); + alarm.setType("alarm"); + alarm.setSeverity(AlarmSeverity.CRITICAL); + Alarm savedAlarm = doPost("/api/alarm", alarm, Alarm.class); + + // ack alarm + edgeImitator.expectMessageAmount(1); + doPost("/api/alarm/" + savedAlarm.getUuidId() + "/ack"); + Assert.assertTrue(edgeImitator.waitForMessages()); + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg); + AlarmUpdateMsg alarmUpdateMsg = (AlarmUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ALARM_ACK_RPC_MESSAGE, alarmUpdateMsg.getMsgType()); + Alarm alarmMsg = JacksonUtil.fromString(alarmUpdateMsg.getEntity(), Alarm.class, true); + Assert.assertNotNull(alarmMsg); + Assert.assertEquals(savedAlarm.getType(), alarmMsg.getType()); + Assert.assertEquals(savedAlarm.getName(), alarmMsg.getName()); + Assert.assertEquals(AlarmStatus.ACTIVE_ACK, alarmMsg.getStatus()); + + // clear alarm + edgeImitator.expectMessageAmount(1); + doPost("/api/alarm/" + savedAlarm.getUuidId() + "/clear"); + Assert.assertTrue(edgeImitator.waitForMessages()); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg); + alarmUpdateMsg = (AlarmUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ALARM_CLEAR_RPC_MESSAGE, alarmUpdateMsg.getMsgType()); + alarmMsg = JacksonUtil.fromString(alarmUpdateMsg.getEntity(), Alarm.class, true); + Assert.assertNotNull(alarmMsg); + Assert.assertEquals(savedAlarm.getType(), alarmMsg.getType()); + Assert.assertEquals(savedAlarm.getName(), alarmMsg.getName()); + Assert.assertEquals(AlarmStatus.CLEARED_ACK, alarmMsg.getStatus()); + + // delete alarm + edgeImitator.expectMessageAmount(1); + doDelete("/api/alarm/" + savedAlarm.getUuidId()) + .andExpect(status().isOk()); + Assert.assertTrue(edgeImitator.waitForMessages()); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg); + alarmUpdateMsg = (AlarmUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, alarmUpdateMsg.getMsgType()); + alarmMsg = JacksonUtil.fromString(alarmUpdateMsg.getEntity(), Alarm.class, true); + Assert.assertNotNull(alarmMsg); + Assert.assertEquals(savedAlarm.getType(), alarmMsg.getType()); + Assert.assertEquals(savedAlarm.getName(), alarmMsg.getName()); + Assert.assertEquals(AlarmStatus.CLEARED_ACK, alarmMsg.getStatus()); + } + @Test public void testSendAlarmCommentToCloud() throws Exception { Device device = saveDeviceOnCloudAndVerifyDeliveryToEdge(); @@ -125,6 +182,53 @@ public class AlarmEdgeTest extends AbstractEdgeTest { Assert.assertEquals(alarmComment.getAlarmId(), alarmInfo.getAlarmId()); } + @Test + public void testAlarmComments() throws Exception { + Device device = findDeviceByName("Edge Device 1"); + Alarm alarm = new Alarm(); + alarm.setOriginator(device.getId()); + alarm.setType("alarm"); + alarm.setSeverity(AlarmSeverity.MINOR); + Alarm savedAlarm = doPost("/api/alarm", alarm, Alarm.class); + + // create alarm comment + edgeImitator.expectMessageAmount(1); + AlarmComment alarmComment = new AlarmComment(); + alarmComment.setComment(new TextNode("Test")); + alarmComment.setAlarmId(savedAlarm.getId()); + alarmComment = doPost("/api/alarm/" + savedAlarm.getUuidId() + "/comment", alarmComment, AlarmComment.class); + Assert.assertTrue(edgeImitator.waitForMessages()); + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + AlarmCommentUpdateMsg alarmCommentUpdateMsg = (AlarmCommentUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, alarmCommentUpdateMsg.getMsgType()); + AlarmComment alarmCommentMsg = JacksonUtil.fromString(alarmCommentUpdateMsg.getEntity(), AlarmComment.class, true); + Assert.assertNotNull(alarmCommentMsg); + Assert.assertEquals(alarmComment, alarmCommentMsg); + + // update alarm comment + edgeImitator.expectMessageAmount(1); + alarmComment.setComment(JacksonUtil.newObjectNode().set("text", new TextNode("Updated comment"))); + alarmComment = doPost("/api/alarm/" + savedAlarm.getUuidId() + "/comment", alarmComment, AlarmComment.class); + Assert.assertTrue(edgeImitator.waitForMessages()); + latestMessage = edgeImitator.getLatestMessage(); + alarmCommentUpdateMsg = (AlarmCommentUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, alarmCommentUpdateMsg.getMsgType()); + alarmCommentMsg = JacksonUtil.fromString(alarmCommentUpdateMsg.getEntity(), AlarmComment.class, true); + Assert.assertNotNull(alarmCommentMsg); + Assert.assertEquals(alarmComment, alarmCommentMsg); + + // delete alarm + edgeImitator.expectMessageAmount(1); + doDelete("/api/alarm/" + savedAlarm.getUuidId() + "/comment/" + alarmComment.getUuidId()) + .andExpect(status().isOk()); + Assert.assertTrue(edgeImitator.waitForMessages()); + latestMessage = edgeImitator.getLatestMessage(); + alarmCommentUpdateMsg = (AlarmCommentUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, alarmCommentUpdateMsg.getMsgType()); + alarmCommentMsg = JacksonUtil.fromString(alarmCommentUpdateMsg.getEntity(), AlarmComment.class, true); + Assert.assertNotNull(alarmCommentMsg); + } + private Alarm buildAlarmForUplinkMsg(DeviceId deviceId) { Alarm alarm = new Alarm(); alarm.setId(new AlarmId(UUID.randomUUID())); diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmCommentService.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmCommentService.java index b60acc6550..042f0f693b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmCommentService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmCommentService.java @@ -32,6 +32,7 @@ import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; +import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.service.DataValidator; import java.util.UUID; @@ -51,11 +52,18 @@ public class BaseAlarmCommentService extends AbstractEntityService implements Al @Override public AlarmComment createOrUpdateAlarmComment(TenantId tenantId, AlarmComment alarmComment) { alarmCommentDataValidator.validate(alarmComment, c -> tenantId); - if (alarmComment.getId() == null) { - return createAlarmComment(tenantId, alarmComment); + boolean isCreated = alarmComment.getId() == null; + AlarmComment result; + if (isCreated) { + result = createAlarmComment(tenantId, alarmComment); } else { - return updateAlarmComment(tenantId, alarmComment); + result = updateAlarmComment(tenantId, alarmComment); + } + if (result != null) { + eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId).entity(result) + .entityId(result.getAlarmId()).added(isCreated).build()); } + return result; } @Override 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 6ed8a33c2d..31fda250d0 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 @@ -28,7 +28,6 @@ import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.alarm.Alarm; -import org.thingsboard.server.common.data.alarm.AlarmComment; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.TbMsg; @@ -42,8 +41,6 @@ import static org.thingsboard.server.common.data.msg.TbMsgType.ACTIVITY_EVENT; import static org.thingsboard.server.common.data.msg.TbMsgType.ALARM; import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_DELETED; import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_UPDATED; -import static org.thingsboard.server.common.data.msg.TbMsgType.COMMENT_CREATED; -import static org.thingsboard.server.common.data.msg.TbMsgType.COMMENT_UPDATED; import static org.thingsboard.server.common.data.msg.TbMsgType.CONNECT_EVENT; import static org.thingsboard.server.common.data.msg.TbMsgType.DISCONNECT_EVENT; import static org.thingsboard.server.common.data.msg.TbMsgType.INACTIVITY_EVENT; @@ -83,9 +80,6 @@ public abstract class AbstractTbMsgPushNode metadata = msg.getMetaData().getData(); EdgeEventActionType actionType = getEdgeEventActionTypeByMsgType(msg); @@ -139,8 +133,6 @@ public abstract class AbstractTbMsgPushNode getConfigClazz(); @@ -152,11 +144,6 @@ public abstract class AbstractTbMsgPushNode metadata) { String scope = metadata.get(SCOPE); if (StringUtils.isEmpty(scope)) { @@ -179,10 +166,6 @@ public abstract class AbstractTbMsgPushNodeATTRIBUTES_UPDATED" + "
ATTRIBUTES_DELETED" + "
ALARM

" + - "
COMMENT_CREATED" + - "
COMMENT_UPDATED" + "Message will be routed via Failure route if node was not able to save cloud event to database or unsupported originator type/message type arrived. " + "In case successful storage cloud event to database message will be routed via Success route.", uiResources = {"static/rulenode/rulenode-core-config.js"}, @@ -72,11 +70,6 @@ public class TbMsgPushToCloudNode extends AbstractTbMsgPushNodeATTRIBUTES_UPDATED" + "
ATTRIBUTES_DELETED" + "
ALARM

" + - "
COMMENT_CREATED" + - "
COMMENT_UPDATED" + "Message will be routed via Failure route if node was not able to save edge event to database or unsupported message type arrived. " + "In case successful storage edge event to database message will be routed via Success route.", uiResources = {"static/rulenode/rulenode-core-config.js"}, @@ -95,11 +90,6 @@ public class TbMsgPushToEdgeNode extends AbstractTbMsgPushNode> futures = new ArrayList<>(); - EntityId finalOriginatorId = originatorId; PageDataIterableByTenantIdEntityId edgeIds = new PageDataIterableByTenantIdEntityId<>( - ctx.getEdgeService()::findRelatedEdgeIdsByEntityId, ctx.getTenantId(), finalOriginatorId, DEFAULT_PAGE_SIZE); + ctx.getEdgeService()::findRelatedEdgeIdsByEntityId, ctx.getTenantId(), msg.getOriginator(), DEFAULT_PAGE_SIZE); for (EdgeId edgeId : edgeIds) { EdgeEvent edgeEvent = buildEvent(msg, ctx); futures.add(notifyEdge(ctx, edgeEvent, edgeId));