Browse Source

Merge pull request #34 from BohdanSmetanyuk/fix/incorrect_ts

Fix/incorrect ts in timeseries update msgs from cloud
pull/2436/head
VoBa 6 years ago
committed by GitHub
parent
commit
8f4138632c
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java
  2. 8
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java
  3. 28
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java

4
application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/EntityDataMsgConstructor.java

@ -41,7 +41,9 @@ public class EntityDataMsgConstructor {
switch (actionType) { switch (actionType) {
case TIMESERIES_UPDATED: case TIMESERIES_UPDATED:
try { try {
builder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(entityData)); JsonObject data = entityData.getAsJsonObject();
long ts = data.getAsJsonPrimitive("ts").getAsLong();
builder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(data.getAsJsonObject("data"), ts));
} catch (Exception e) { } catch (Exception e) {
log.warn("Can't convert to telemetry proto, entityData [{}]", entityData, e); log.warn("Can't convert to telemetry proto, entityData [{}]", entityData, e);
} }

8
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java

@ -68,12 +68,16 @@ public class JsonConverter {
private static int maxStringValueLength = 0; private static int maxStringValueLength = 0;
public static PostTelemetryMsg convertToTelemetryProto(JsonElement jsonElement) throws JsonSyntaxException { public static PostTelemetryMsg convertToTelemetryProto(JsonElement jsonElement, long ts) throws JsonSyntaxException {
PostTelemetryMsg.Builder builder = PostTelemetryMsg.newBuilder(); PostTelemetryMsg.Builder builder = PostTelemetryMsg.newBuilder();
convertToTelemetry(jsonElement, System.currentTimeMillis(), null, builder); convertToTelemetry(jsonElement, ts, null, builder);
return builder.build(); return builder.build();
} }
public static PostTelemetryMsg convertToTelemetryProto(JsonElement jsonElement) throws JsonSyntaxException {
return convertToTelemetryProto(jsonElement, System.currentTimeMillis());
}
private static void convertToTelemetry(JsonElement jsonElement, long systemTs, Map<Long, List<KvEntry>> result, PostTelemetryMsg.Builder builder) { private static void convertToTelemetry(JsonElement jsonElement, long systemTs, Map<Long, List<KvEntry>> result, PostTelemetryMsg.Builder builder) {
if (jsonElement.isJsonObject()) { if (jsonElement.isJsonObject()) {
parseObject(systemTs, result, builder, jsonElement.getAsJsonObject()); parseObject(systemTs, result, builder, jsonElement.getAsJsonObject());

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

@ -150,13 +150,7 @@ public class TbMsgPushToEdgeNode implements TbNode {
return null; return null;
} }
ActionType actionType = getActionTypeByMsgType(msg.getType()); ActionType actionType = getActionTypeByMsgType(msg.getType());
JsonNode entityBody = null; JsonNode entityBody = getEntityBody(actionType, msg.getData(), msg.getMetaData().getData());
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); return buildEdgeEvent(ctx.getTenantId(), actionType, msg.getOriginator().getId(), edgeEventTypeByEntityType, entityBody);
} }
} }
@ -171,19 +165,25 @@ public class TbMsgPushToEdgeNode implements TbNode {
return edgeEvent; return edgeEvent;
} }
private JsonNode getAttributeEntityBody(ActionType actionType, JsonNode data, Map<String, String> metadata) throws JsonProcessingException { private JsonNode getEntityBody(ActionType actionType, String data, Map<String, String> metadata) throws JsonProcessingException {
Map<String, Object> entityData = new HashMap<>(); Map<String, Object> entityBody = new HashMap<>();
JsonNode dataJson = json.readTree(data);
switch (actionType) { switch (actionType) {
case ATTRIBUTES_UPDATED: case ATTRIBUTES_UPDATED:
entityData.put("kv", data); entityBody.put("kv", dataJson);
entityBody.put("scope", metadata.get("scope"));
break; break;
case ATTRIBUTES_DELETED: case ATTRIBUTES_DELETED:
List<String> keys = json.treeToValue(data.get("attributes"), List.class); List<String> keys = json.treeToValue(dataJson.get("attributes"), List.class);
entityData.put("keys", keys); entityBody.put("keys", keys);
entityBody.put("scope", metadata.get("scope"));
break;
case TIMESERIES_UPDATED:
entityBody.put("data", dataJson);
entityBody.put("ts", metadata.get("ts"));
break; break;
} }
entityData.put("scope", metadata.get("scope")); return json.valueToTree(entityBody);
return json.valueToTree(entityData);
} }
private UUID getUUIDFromMsgData(TbMsg msg) throws JsonProcessingException { private UUID getUUIDFromMsgData(TbMsg msg) throws JsonProcessingException {

Loading…
Cancel
Save