Browse Source

bug with alarms fixed

pull/2436/head
Bohdan Smetaniuk 6 years ago
parent
commit
6b153ee2ad
  1. 6
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  2. 51
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java

6
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java

@ -962,18 +962,18 @@ public final class EdgeGrpcSession implements Closeable {
switch (alarmUpdateMsg.getMsgType()) { switch (alarmUpdateMsg.getMsgType()) {
case ENTITY_CREATED_RPC_MESSAGE: case ENTITY_CREATED_RPC_MESSAGE:
case ENTITY_UPDATED_RPC_MESSAGE: case ENTITY_UPDATED_RPC_MESSAGE:
if (existentAlarm == null) { if (existentAlarm == null || existentAlarm.getStatus().isCleared()) {
existentAlarm = new Alarm(); existentAlarm = new Alarm();
existentAlarm.setTenantId(edge.getTenantId()); existentAlarm.setTenantId(edge.getTenantId());
existentAlarm.setType(alarmUpdateMsg.getName()); existentAlarm.setType(alarmUpdateMsg.getName());
existentAlarm.setOriginator(originatorId); existentAlarm.setOriginator(originatorId);
existentAlarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity())); existentAlarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity()));
existentAlarm.setStatus(AlarmStatus.valueOf(alarmUpdateMsg.getStatus()));
existentAlarm.setStartTs(alarmUpdateMsg.getStartTs()); existentAlarm.setStartTs(alarmUpdateMsg.getStartTs());
existentAlarm.setAckTs(alarmUpdateMsg.getAckTs());
existentAlarm.setClearTs(alarmUpdateMsg.getClearTs()); existentAlarm.setClearTs(alarmUpdateMsg.getClearTs());
existentAlarm.setPropagate(alarmUpdateMsg.getPropagate()); existentAlarm.setPropagate(alarmUpdateMsg.getPropagate());
} }
existentAlarm.setStatus(AlarmStatus.valueOf(alarmUpdateMsg.getStatus()));
existentAlarm.setAckTs(alarmUpdateMsg.getAckTs());
existentAlarm.setEndTs(alarmUpdateMsg.getEndTs()); existentAlarm.setEndTs(alarmUpdateMsg.getEndTs());
existentAlarm.setDetails(mapper.readTree(alarmUpdateMsg.getDetails())); existentAlarm.setDetails(mapper.readTree(alarmUpdateMsg.getDetails()));
ctx.getAlarmService().createOrUpdateAlarm(existentAlarm); ctx.getAlarmService().createOrUpdateAlarm(existentAlarm);

51
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java

@ -16,6 +16,7 @@
package org.thingsboard.rule.engine.edge; package org.thingsboard.rule.engine.edge;
import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
@ -46,6 +47,7 @@ import org.thingsboard.server.common.msg.session.SessionMsgType;
import javax.annotation.Nullable; import javax.annotation.Nullable;
import java.util.List; import java.util.List;
import java.util.UUID;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
@ -84,14 +86,10 @@ public class TbMsgPushToEdgeNode implements TbNode {
Futures.addCallback(getEdgeIdFuture, new FutureCallback<EdgeId>() { Futures.addCallback(getEdgeIdFuture, new FutureCallback<EdgeId>() {
@Override @Override
public void onSuccess(@Nullable EdgeId edgeId) { public void onSuccess(@Nullable EdgeId edgeId) {
EdgeEventType edgeEventTypeByEntityType = EdgeUtils.getEdgeEventTypeByEntityType(msg.getOriginator().getEntityType());
if (edgeEventTypeByEntityType == null) {
log.debug("Edge event type is null. Entity Type {}", msg.getOriginator().getEntityType());
ctx.tellFailure(msg, new RuntimeException("Edge event type is null. Entity Type '" + msg.getOriginator().getEntityType() + "'"));
}
EdgeEvent edgeEvent = null; EdgeEvent edgeEvent = null;
try { try {
edgeEvent = buildEdgeEvent(ctx, msg, edgeId, edgeEventTypeByEntityType); edgeEvent = buildEdgeEvent(msg, ctx);
edgeEvent.setEdgeId(edgeId);
} catch (JsonProcessingException e) { } catch (JsonProcessingException e) {
log.error("Failed to build edge event", e); log.error("Failed to build edge event", e);
} }
@ -124,17 +122,35 @@ public class TbMsgPushToEdgeNode implements TbNode {
} }
} }
private EdgeEvent buildEdgeEvent(TbContext ctx, TbMsg msg, EdgeId edgeId, EdgeEventType edgeEventTypeByEntityType) throws JsonProcessingException { private EdgeEvent buildEdgeEvent(TbMsg msg, TbContext ctx) throws JsonProcessingException {
if (DataConstants.ALARM.equals(msg.getType())) {
return buildEdgeEvent(ctx.getTenantId(), ActionType.ADDED, getUUIDFromMsgData(msg), EdgeEventType.ALARM, null);
} else {
EdgeEventType edgeEventTypeByEntityType = EdgeUtils.getEdgeEventTypeByEntityType(msg.getOriginator().getEntityType());
if (edgeEventTypeByEntityType == null) {
log.debug("Edge event type is null. Entity Type {}", msg.getOriginator().getEntityType());
ctx.tellFailure(msg, new RuntimeException("Edge event type is null. Entity Type '" + msg.getOriginator().getEntityType() + "'"));
}
return buildEdgeEvent(ctx.getTenantId(), getActionTypeByMsgType(msg.getType()), msg.getOriginator().getId(), edgeEventTypeByEntityType, json.readTree(msg.getData()));
}
}
private EdgeEvent buildEdgeEvent(TenantId tenantId, ActionType edgeEventAction, UUID entityId, EdgeEventType edgeEventType, JsonNode entityBody) {
EdgeEvent edgeEvent = new EdgeEvent(); EdgeEvent edgeEvent = new EdgeEvent();
edgeEvent.setTenantId(ctx.getTenantId()); edgeEvent.setTenantId(tenantId);
edgeEvent.setEdgeId(edgeId); edgeEvent.setEdgeEventAction(edgeEventAction.name());
edgeEvent.setEdgeEventAction(getActionTypeByMsgType(msg.getType()).name()); edgeEvent.setEntityId(entityId);
edgeEvent.setEntityId(msg.getOriginator().getId()); edgeEvent.setEdgeEventType(edgeEventType);
edgeEvent.setEdgeEventType(edgeEventTypeByEntityType); edgeEvent.setEntityBody(entityBody);
edgeEvent.setEntityBody(json.readTree(msg.getData()));
return edgeEvent; return edgeEvent;
} }
private UUID getUUIDFromMsgData(TbMsg msg) throws JsonProcessingException {
JsonNode data = json.readTree(msg.getData()).get("id");
String id = json.treeToValue(data.get("id"), String.class);
return UUID.fromString(id);
}
private ActionType getActionTypeByMsgType(String msgType) { private ActionType getActionTypeByMsgType(String msgType) {
ActionType actionType; ActionType actionType;
if (SessionMsgType.POST_TELEMETRY_REQUEST.name().equals(msgType)) { if (SessionMsgType.POST_TELEMETRY_REQUEST.name().equals(msgType)) {
@ -161,14 +177,11 @@ public class TbMsgPushToEdgeNode implements TbNode {
} }
private boolean isSupportedMsgType(String msgType) { private boolean isSupportedMsgType(String msgType) {
if (SessionMsgType.POST_TELEMETRY_REQUEST.name().equals(msgType) return SessionMsgType.POST_TELEMETRY_REQUEST.name().equals(msgType)
|| SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType) || SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType)
|| DataConstants.ATTRIBUTES_UPDATED.equals(msgType) || DataConstants.ATTRIBUTES_UPDATED.equals(msgType)
|| DataConstants.ATTRIBUTES_DELETED.equals(msgType)) { || DataConstants.ATTRIBUTES_DELETED.equals(msgType)
return true; || DataConstants.ALARM.equals(msgType);
} else {
return false;
}
} }
private ListenableFuture<EdgeId> getEdgeIdByOriginatorId(TbContext ctx, TenantId tenantId, EntityId originatorId) { private ListenableFuture<EdgeId> getEdgeIdByOriginatorId(TbContext ctx, TenantId tenantId, EntityId originatorId) {

Loading…
Cancel
Save