diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java index d9a06655c5..07f7ffb02a 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java @@ -15,14 +15,20 @@ */ package org.thingsboard.server.service.edge.rpc.constructor; +import com.google.gson.Gson; +import com.google.gson.JsonArray; import com.google.gson.JsonElement; +import com.google.gson.JsonObject; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.transport.adaptor.JsonConverter; +import org.thingsboard.server.gen.edge.AttributeDeleteMsg; import org.thingsboard.server.gen.edge.EntityDataProto; +import java.util.List; + @Component @Slf4j public class EntityDataMsgConstructor { @@ -42,13 +48,26 @@ public class EntityDataMsgConstructor { break; case ATTRIBUTES_UPDATED: try { - builder.setPostAttributesMsg(JsonConverter.convertToAttributesProto(entityData)); + JsonObject data = entityData.getAsJsonObject(); + builder.setPostAttributesMsg(JsonConverter.convertToAttributesProto(data.getAsJsonObject("kv"))); + builder.setPostAttributeScope(data.getAsJsonPrimitive("scope").getAsString()); } catch (Exception e) { log.warn("Can't convert to attributes proto, entityData [{}]", entityData, e); } break; - // TODO: voba - add support for attribute delete - // case ATTRIBUTES_DELETED: + case ATTRIBUTES_DELETED: + try { + AttributeDeleteMsg.Builder attributeDeleteMsg = AttributeDeleteMsg.newBuilder(); + attributeDeleteMsg.setScope(entityData.getAsJsonObject().getAsJsonPrimitive("scope").getAsString()); + JsonArray jsonArray = entityData.getAsJsonObject().getAsJsonArray("keys"); + List keys = new Gson().fromJson(jsonArray.toString(), List.class); + attributeDeleteMsg.addAllAttributeNames(keys); + attributeDeleteMsg.build(); + builder.setAttributeDeleteMsg(attributeDeleteMsg); + } catch (Exception e) { + log.warn("Can't convert to AttributeDeleteMsg proto, entityData [{}]", entityData, e); + } + break; } return builder.build(); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java index 8576f96152..1a2a3d3177 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java @@ -71,7 +71,9 @@ import org.thingsboard.server.gen.edge.UserCredentialsRequestMsg; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.UUID; @Service @@ -289,25 +291,29 @@ public class DefaultSyncEdgeService implements SyncEdgeService { public void onSuccess(@Nullable List ssAttributes) { if (ssAttributes != null && !ssAttributes.isEmpty()) { try { - ObjectNode entityNode = mapper.createObjectNode(); + Map entityData = new HashMap<>(); + ObjectNode attributes = mapper.createObjectNode(); for (AttributeKvEntry attr : ssAttributes) { if (attr.getDataType() == DataType.BOOLEAN && attr.getBooleanValue().isPresent()) { - entityNode.put(attr.getKey(), attr.getBooleanValue().get()); + attributes.put(attr.getKey(), attr.getBooleanValue().get()); } else if (attr.getDataType() == DataType.DOUBLE && attr.getDoubleValue().isPresent()) { - entityNode.put(attr.getKey(), attr.getDoubleValue().get()); + attributes.put(attr.getKey(), attr.getDoubleValue().get()); } else if (attr.getDataType() == DataType.LONG && attr.getLongValue().isPresent()) { - entityNode.put(attr.getKey(), attr.getLongValue().get()); + attributes.put(attr.getKey(), attr.getLongValue().get()); } else { - entityNode.put(attr.getKey(), attr.getValueAsString()); + attributes.put(attr.getKey(), attr.getValueAsString()); } } - log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, entityNode); + entityData.put("kv", attributes); + entityData.put("scope", DataConstants.SERVER_SCOPE); + JsonNode entityBody = mapper.valueToTree(entityData); + log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, entityBody); saveEdgeEvent(edge.getTenantId(), edge.getId(), edgeEventType, ActionType.ATTRIBUTES_UPDATED, entityId, - entityNode); + entityBody); } catch (Exception e) { log.error("[{}] Failed to send attribute updates to the edge", edge.getName(), e); } diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index 2cdbd241d1..181ebff531 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -109,9 +109,16 @@ message EntityDataProto { string entityType = 3; transport.PostTelemetryMsg postTelemetryMsg = 4; transport.PostAttributeMsg postAttributesMsg = 5; + string postAttributeScope = 6; + AttributeDeleteMsg attributeDeleteMsg = 7; // transport.ToDeviceRpcRequestMsg ??? } +message AttributeDeleteMsg { + string scope = 1; + repeated string attributeNames = 2; +} + message RuleChainUpdateMsg { UpdateMsgType msgType = 1; int64 idMSB = 2; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java index eef30d5d22..f1134078c1 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java @@ -46,7 +46,9 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.session.SessionMsgType; import javax.annotation.Nullable; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.UUID; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; @@ -137,7 +139,15 @@ public class TbMsgPushToEdgeNode implements TbNode { if (edgeEventTypeByEntityType == null) { return null; } - return buildEdgeEvent(ctx.getTenantId(), getActionTypeByMsgType(msg.getType()), msg.getOriginator().getId(), edgeEventTypeByEntityType, json.readTree(msg.getData())); + ActionType actionType = getActionTypeByMsgType(msg.getType()); + JsonNode entityBody = null; + JsonNode data = json.readTree(msg.getData()); + if (actionType.equals(ActionType.ATTRIBUTES_UPDATED) || actionType.equals(ActionType.ATTRIBUTES_DELETED)) { + entityBody = getAttributeEntityBody(actionType, data, msg.getMetaData().getData()); + } else { + entityBody = data; + } + return buildEdgeEvent(ctx.getTenantId(), actionType, msg.getOriginator().getId(), edgeEventTypeByEntityType, entityBody); } } @@ -151,6 +161,21 @@ public class TbMsgPushToEdgeNode implements TbNode { return edgeEvent; } + private JsonNode getAttributeEntityBody(ActionType actionType, JsonNode data, Map metadata) throws JsonProcessingException { + Map entityData = new HashMap<>(); + switch (actionType) { + case ATTRIBUTES_UPDATED: + entityData.put("kv", data); + break; + case ATTRIBUTES_DELETED: + List keys = json.treeToValue(data.get("attributes"), List.class); + entityData.put("keys", keys); + break; + } + entityData.put("scope", metadata.get("scope")); + return json.valueToTree(entityData); + } + private UUID getUUIDFromMsgData(TbMsg msg) throws JsonProcessingException { JsonNode data = json.readTree(msg.getData()).get("id"); String id = json.treeToValue(data.get("id"), String.class); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java index e776fa4dbd..b6e9e4a292 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java @@ -64,6 +64,7 @@ public class TbMsgAttributesNode implements TbNode { } String src = msg.getData(); Set attributes = JsonConverter.convertToAttributes(new JsonParser().parse(src)); + msg.getMetaData().putValue("scope", config.getScope()); ctx.getTelemetryService().saveAndNotify(ctx.getTenantId(), msg.getOriginator(), config.getScope(), new ArrayList<>(attributes), new TelemetryNodeCallback(ctx, msg)); }