Browse Source

Constant usage; ASC order

pull/2818/head
Volodymyr Babak 7 years ago
parent
commit
2b70cfc95a
  1. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  2. 13
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/TbMsgPushToEdgeNode.java

2
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<Event> pageData;
UUID ifOffset = null;
do {

13
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));
}

Loading…
Cancel
Save