|
|
|
@ -16,6 +16,7 @@ |
|
|
|
package org.thingsboard.rule.engine.edge; |
|
|
|
|
|
|
|
import com.fasterxml.jackson.core.JsonProcessingException; |
|
|
|
import com.fasterxml.jackson.databind.JsonNode; |
|
|
|
import com.fasterxml.jackson.databind.ObjectMapper; |
|
|
|
import com.google.common.util.concurrent.FutureCallback; |
|
|
|
import com.google.common.util.concurrent.Futures; |
|
|
|
@ -46,6 +47,7 @@ import org.thingsboard.server.common.msg.session.SessionMsgType; |
|
|
|
|
|
|
|
import javax.annotation.Nullable; |
|
|
|
import java.util.List; |
|
|
|
import java.util.UUID; |
|
|
|
|
|
|
|
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; |
|
|
|
|
|
|
|
@ -84,14 +86,10 @@ public class TbMsgPushToEdgeNode implements TbNode { |
|
|
|
Futures.addCallback(getEdgeIdFuture, new FutureCallback<EdgeId>() { |
|
|
|
@Override |
|
|
|
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; |
|
|
|
try { |
|
|
|
edgeEvent = buildEdgeEvent(ctx, msg, edgeId, edgeEventTypeByEntityType); |
|
|
|
edgeEvent = buildEdgeEvent(msg, ctx); |
|
|
|
edgeEvent.setEdgeId(edgeId); |
|
|
|
} catch (JsonProcessingException 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.setTenantId(ctx.getTenantId()); |
|
|
|
edgeEvent.setEdgeId(edgeId); |
|
|
|
edgeEvent.setEdgeEventAction(getActionTypeByMsgType(msg.getType()).name()); |
|
|
|
edgeEvent.setEntityId(msg.getOriginator().getId()); |
|
|
|
edgeEvent.setEdgeEventType(edgeEventTypeByEntityType); |
|
|
|
edgeEvent.setEntityBody(json.readTree(msg.getData())); |
|
|
|
edgeEvent.setTenantId(tenantId); |
|
|
|
edgeEvent.setEdgeEventAction(edgeEventAction.name()); |
|
|
|
edgeEvent.setEntityId(entityId); |
|
|
|
edgeEvent.setEdgeEventType(edgeEventType); |
|
|
|
edgeEvent.setEntityBody(entityBody); |
|
|
|
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) { |
|
|
|
ActionType actionType; |
|
|
|
if (SessionMsgType.POST_TELEMETRY_REQUEST.name().equals(msgType)) { |
|
|
|
@ -161,14 +177,11 @@ public class TbMsgPushToEdgeNode implements TbNode { |
|
|
|
} |
|
|
|
|
|
|
|
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) |
|
|
|
|| DataConstants.ATTRIBUTES_UPDATED.equals(msgType) |
|
|
|
|| DataConstants.ATTRIBUTES_DELETED.equals(msgType)) { |
|
|
|
return true; |
|
|
|
} else { |
|
|
|
return false; |
|
|
|
} |
|
|
|
|| DataConstants.ATTRIBUTES_DELETED.equals(msgType) |
|
|
|
|| DataConstants.ALARM.equals(msgType); |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<EdgeId> getEdgeIdByOriginatorId(TbContext ctx, TenantId tenantId, EntityId originatorId) { |
|
|
|
|