From 2b70cfc95aabb5e354e44d17cc544d3b4995c70c Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 24 Mar 2020 11:44:47 +0200 Subject: [PATCH] Constant usage; ASC order --- .../server/service/edge/rpc/EdgeGrpcSession.java | 2 +- .../rule/engine/edge/TbMsgPushToEdgeNode.java | 13 +++++++++++-- 2 files changed, 12 insertions(+), 3 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 9e56073512..de71e91ace 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -164,7 +164,7 @@ public final class EdgeGrpcSession implements Cloneable { void processHandleMessages() throws ExecutionException, InterruptedException { Long queueStartTs = getQueueStartTs().get(); // TODO: this 100 value must be changed properly - TimePageLink pageLink = new TimePageLink(30, queueStartTs + 1000); + TimePageLink pageLink = new TimePageLink(30, queueStartTs + 1000, null, true); TimePageData pageData; UUID ifOffset = null; do { 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 152c3a386c..92a4103867 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 @@ -25,6 +25,7 @@ import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.session.SessionMsgType; @Slf4j @RuleNode( @@ -39,6 +40,11 @@ import org.thingsboard.server.common.msg.TbMsg; ) public class TbMsgPushToEdgeNode implements TbNode { + private static final String CLOUD_MSG_SOURCE = "cloud"; + private static final String EDGE_MSG_SOURCE = "edge"; + private static final String MSG_SOURCE_KEY = "source"; + private static final String TS_METADATA_KEY = "ts"; + private EmptyNodeConfiguration config; @Override @@ -48,10 +54,13 @@ public class TbMsgPushToEdgeNode implements TbNode { @Override public void onMsg(TbContext ctx, TbMsg msg) { - if ("edge".equalsIgnoreCase(msg.getMetaData().getValue("source"))) { + if (EDGE_MSG_SOURCE.equalsIgnoreCase(msg.getMetaData().getValue(MSG_SOURCE_KEY))) { return; } - msg.getMetaData().putValue("source", "cloud"); + if (msg.getType().equals(SessionMsgType.POST_TELEMETRY_REQUEST.name())) { + msg.getMetaData().putValue(TS_METADATA_KEY, Long.toString(System.currentTimeMillis())); + } + msg.getMetaData().putValue(MSG_SOURCE_KEY, CLOUD_MSG_SOURCE); ctx.getEdgeService().pushEventToEdge(ctx.getTenantId(), msg, new PushToEdgeNodeCallback(ctx, msg)); }