Browse Source

Merge pull request #25 from BohdanSmetanyuk/feature/attribute_removal

attributes removal + fixed bug with incorrect attributes saving
pull/2436/head
VoBa 6 years ago
committed by GitHub
parent
commit
e7e6377407
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 25
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java
  2. 20
      application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java
  3. 7
      common/edge-api/src/main/proto/edge.proto
  4. 27
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java
  5. 1
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java

25
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<String> 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();
}

20
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<AttributeKvEntry> ssAttributes) {
if (ssAttributes != null && !ssAttributes.isEmpty()) {
try {
ObjectNode entityNode = mapper.createObjectNode();
Map<String, Object> 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);
}

7
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;

27
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<String, String> metadata) throws JsonProcessingException {
Map<String, Object> entityData = new HashMap<>();
switch (actionType) {
case ATTRIBUTES_UPDATED:
entityData.put("kv", data);
break;
case ATTRIBUTES_DELETED:
List<String> 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);

1
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<AttributeKvEntry> 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));
}

Loading…
Cancel
Save