diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 54737f84d5..665d2d9489 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -312,23 +312,23 @@ public class ActorSystemContext { return discoveryService.getCurrentServer().getServerAddress().toString(); } - public void persistDebugInput(TenantId tenantId, EntityId entityId, TbMsg tbMsg) { - persistDebug(tenantId, entityId, "IN", tbMsg, null); + public void persistDebugInput(TenantId tenantId, EntityId entityId, TbMsg tbMsg, String relationType) { + persistDebug(tenantId, entityId, "IN", tbMsg, relationType, null); } - public void persistDebugInput(TenantId tenantId, EntityId entityId, TbMsg tbMsg, Throwable error) { - persistDebug(tenantId, entityId, "IN", tbMsg, error); + public void persistDebugInput(TenantId tenantId, EntityId entityId, TbMsg tbMsg, String relationType, Throwable error) { + persistDebug(tenantId, entityId, "IN", tbMsg, relationType, error); } - public void persistDebugOutput(TenantId tenantId, EntityId entityId, TbMsg tbMsg, Throwable error) { - persistDebug(tenantId, entityId, "OUT", tbMsg, error); + public void persistDebugOutput(TenantId tenantId, EntityId entityId, TbMsg tbMsg, String relationType, Throwable error) { + persistDebug(tenantId, entityId, "OUT", tbMsg, relationType, error); } - public void persistDebugOutput(TenantId tenantId, EntityId entityId, TbMsg tbMsg) { - persistDebug(tenantId, entityId, "OUT", tbMsg, null); + public void persistDebugOutput(TenantId tenantId, EntityId entityId, TbMsg tbMsg, String relationType) { + persistDebug(tenantId, entityId, "OUT", tbMsg, relationType, null); } - private void persistDebug(TenantId tenantId, EntityId entityId, String type, TbMsg tbMsg, Throwable error) { + private void persistDebug(TenantId tenantId, EntityId entityId, String type, TbMsg tbMsg, String relationType, Throwable error) { try { Event event = new Event(); event.setTenantId(tenantId); @@ -345,6 +345,7 @@ public class ActorSystemContext { .put("msgId", tbMsg.getId().toString()) .put("msgType", tbMsg.getType()) .put("dataType", tbMsg.getDataType().name()) + .put("relationType", relationType) .put("data", tbMsg.getData()) .put("metadata", metadata); diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index dd270e2757..f5f5848473 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -355,16 +355,15 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso for (Map.Entry> entry : tsData.entrySet()) { JsonObject json = new JsonObject(); - json.addProperty("ts", entry.getKey()); - JsonObject values = new JsonObject(); for (KvEntry kv : entry.getValue()) { - kv.getBooleanValue().ifPresent(v -> values.addProperty(kv.getKey(), v)); - kv.getLongValue().ifPresent(v -> values.addProperty(kv.getKey(), v)); - kv.getDoubleValue().ifPresent(v -> values.addProperty(kv.getKey(), v)); - kv.getStrValue().ifPresent(v -> values.addProperty(kv.getKey(), v)); + kv.getBooleanValue().ifPresent(v -> json.addProperty(kv.getKey(), v)); + kv.getLongValue().ifPresent(v -> json.addProperty(kv.getKey(), v)); + kv.getDoubleValue().ifPresent(v -> json.addProperty(kv.getKey(), v)); + kv.getStrValue().ifPresent(v -> json.addProperty(kv.getKey(), v)); } - json.add("values", values); - TbMsg tbMsg = new TbMsg(UUIDs.timeBased(), SessionMsgType.POST_TELEMETRY_REQUEST.name(), deviceId, defaultMetaData.copy(), TbMsgDataType.JSON, gson.toJson(json), null, null, 0L); + TbMsgMetaData metaData = defaultMetaData.copy(); + metaData.putValue("ts", entry.getKey()+""); + TbMsg tbMsg = new TbMsg(UUIDs.timeBased(), SessionMsgType.POST_TELEMETRY_REQUEST.name(), deviceId, metaData, TbMsgDataType.JSON, gson.toJson(json), null, null, 0L); pushToRuleEngineWithTimeout(context, tbMsg, msgData); } } diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java index 039fdf6127..2e8955bf1f 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java @@ -69,7 +69,7 @@ class DefaultTbContext implements TbContext { @Override public void tellNext(TbMsg msg, String relationType, Throwable th) { if (nodeCtx.getSelf().isDebugMode()) { - mainCtx.persistDebugOutput(nodeCtx.getTenantId(), nodeCtx.getSelf().getId(), msg, th); + mainCtx.persistDebugOutput(nodeCtx.getTenantId(), nodeCtx.getSelf().getId(), msg, relationType, th); } nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(nodeCtx.getSelf().getId(), relationType, msg), nodeCtx.getSelfActor()); } @@ -102,7 +102,7 @@ class DefaultTbContext implements TbContext { @Override public void tellError(TbMsg msg, Throwable th) { if (nodeCtx.getSelf().isDebugMode()) { - mainCtx.persistDebugOutput(nodeCtx.getTenantId(), nodeCtx.getSelf().getId(), msg, th); + mainCtx.persistDebugOutput(nodeCtx.getTenantId(), nodeCtx.getSelf().getId(), msg, "", th); } nodeCtx.getSelfActor().tell(new RuleNodeToSelfErrorMsg(msg, th), nodeCtx.getSelfActor()); } diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java index 9c60327b8b..d069cb0eae 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java @@ -96,12 +96,12 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor ruleNodeList) { for (RuleNode ruleNode : ruleNodeList) { for (TbMsg tbMsg : queue.findUnprocessed(ruleNode.getId().getId(), 0L)) { - pushMsgToNode(nodeActors.get(ruleNode.getId()), tbMsg); + pushMsgToNode(nodeActors.get(ruleNode.getId()), tbMsg, ""); } } if (firstNode != null) { for (TbMsg tbMsg : queue.findUnprocessed(entityId.getId(), 0L)) { - pushMsgToNode(firstNode, tbMsg); + pushMsgToNode(firstNode, tbMsg, ""); } } } @@ -183,13 +183,13 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor pushMsgToNode(firstNode, msg)); + putToQueue(enrichWithRuleChainId(envelope.getTbMsg()), msg -> pushMsgToNode(firstNode, msg, "")); } void onDeviceActorToRuleEngineMsg(DeviceActorToRuleEngineMsg envelope) { checkActive(); putToQueue(enrichWithRuleChainId(envelope.getTbMsg()), msg -> { - pushMsgToNode(firstNode, msg); + pushMsgToNode(firstNode, msg, ""); envelope.getCallbackRef().tell(new RuleEngineQueuePutAckMsg(msg.getId()), self); }); } @@ -197,9 +197,9 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor pushMsgToNode(firstNode, msg)); + putToQueue(enrichWithRuleChainId(envelope.getMsg()), msg -> pushMsgToNode(firstNode, msg, envelope.getFromRelationType())); } else { - pushMsgToNode(firstNode, envelope.getMsg()); + pushMsgToNode(firstNode, envelope.getMsg(), envelope.getFromRelationType()); } } @@ -218,17 +218,17 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor pushMsgToNode(targetNodeCtx, queuedMsg)); + putToQueue(copy, queuedMsg -> pushMsgToNode(targetNodeCtx, queuedMsg, fromRelationType)); } - private void pushToTarget(TbMsg msg, EntityId target) { + private void pushToTarget(TbMsg msg, EntityId target, String fromRelationType) { switch (target.getEntityType()) { case RULE_NODE: - pushMsgToNode(nodeActors.get(new RuleNodeId(target.getId())), msg); + pushMsgToNode(nodeActors.get(new RuleNodeId(target.getId())), msg, fromRelationType); break; case RULE_CHAIN: - parent.tell(new RuleChainToRuleChainMsg(new RuleChainId(target.getId()), entityId, msg, false), self); + parent.tell(new RuleChainToRuleChainMsg(new RuleChainId(target.getId()), entityId, msg, fromRelationType, false), self); break; } } - private void pushMsgToNode(RuleNodeCtx nodeCtx, TbMsg msg) { + private void pushMsgToNode(RuleNodeCtx nodeCtx, TbMsg msg, String fromRelationType) { if (nodeCtx != null) { - nodeCtx.getSelfActor().tell(new RuleChainToRuleNodeMsg(new DefaultTbContext(systemContext, nodeCtx), msg), self); + nodeCtx.getSelfActor().tell(new RuleChainToRuleNodeMsg(new DefaultTbContext(systemContext, nodeCtx), msg, fromRelationType), self); } } diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleChainMsg.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleChainMsg.java index 2b2623bf09..086164680b 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleChainMsg.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleChainMsg.java @@ -31,6 +31,7 @@ public final class RuleChainToRuleChainMsg implements TbActorMsg { private final RuleChainId target; private final RuleChainId source; private final TbMsg msg; + private final String fromRelationType; private final boolean enqueue; @Override diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleNodeMsg.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleNodeMsg.java index e7d866c1eb..abcddc91fc 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleNodeMsg.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleNodeMsg.java @@ -29,6 +29,7 @@ final class RuleChainToRuleNodeMsg implements TbActorMsg { private final TbContext ctx; private final TbMsg msg; + private final String fromRelationType; @Override public MsgType getMsgType() { diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java index ea857dbbb4..f23ba6b9e5 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java @@ -93,7 +93,7 @@ public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor jsEngine.executeFilter(msg)), - filterResult -> ctx.tellNext(msg, Boolean.toString(filterResult)), + filterResult -> ctx.tellNext(msg, filterResult.booleanValue() ? "True" : "False"), t -> ctx.tellError(msg, t)); } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeConfiguration.java index 4c5808b236..671132d8f9 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeConfiguration.java @@ -31,7 +31,7 @@ public class TbJsSwitchNodeConfiguration implements NodeConfiguration { if (future.isSuccess()) { - TbMsg next = ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), msg.getData()); - ctx.tellNext(next, TbRelationTypes.SUCCESS); + ctx.tellNext(msg, TbRelationTypes.SUCCESS); } else { TbMsg next = processException(ctx, msg, future.cause()); ctx.tellNext(next, TbRelationTypes.FAILURE, future.cause()); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java index 69e7c1de3f..993b172ddf 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java @@ -106,7 +106,7 @@ public class TbRabbitMqNode implements TbNode { routingKey, properties, msg.getData().getBytes(UTF8)); - return ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), msg.getData()); + return msg; } private TbMsg processException(TbContext ctx, TbMsg origMsg, Throwable t) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java index 114ca27390..d91d5d41ad 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java @@ -63,12 +63,22 @@ public class TbMsgTimeseriesNode implements TbNode { ctx.tellError(msg, new IllegalArgumentException("Unsupported msg type: " + msg.getType())); return; } - + long ts = -1; + String tsStr = msg.getMetaData().getValue("ts"); + if (!StringUtils.isEmpty(tsStr)) { + try { + ts = Long.parseLong(tsStr); + } catch (NumberFormatException e) {} + } + if (ts == -1) { + ctx.tellError(msg, new IllegalArgumentException("Msg metadata doesn't contain valid ts value: " + msg.getMetaData())); + return; + } String src = msg.getData(); - TelemetryUploadRequest telemetryUploadRequest = JsonConverter.convertToTelemetry(new JsonParser().parse(src)); + TelemetryUploadRequest telemetryUploadRequest = JsonConverter.convertToTelemetry(new JsonParser().parse(src), ts); Map> tsKvMap = telemetryUploadRequest.getData(); if (tsKvMap == null) { - ctx.tellError(msg, new IllegalArgumentException("Msg body us empty: " + src)); + ctx.tellError(msg, new IllegalArgumentException("Msg body is empty: " + src)); return; } List tsKvEntryList = new ArrayList<>(); diff --git a/ui/src/app/event/event-header-debug-rulenode.tpl.html b/ui/src/app/event/event-header-debug-rulenode.tpl.html index 34f4513655..bbf6ae0cef 100644 --- a/ui/src/app/event/event-header-debug-rulenode.tpl.html +++ b/ui/src/app/event/event-header-debug-rulenode.tpl.html @@ -21,7 +21,7 @@
event.entity
event.message-id
event.message-type
-
event.data-type
+
event.relation-type
event.data
event.metadata
event.error
diff --git a/ui/src/app/event/event-row-debug-rulenode.tpl.html b/ui/src/app/event/event-row-debug-rulenode.tpl.html index bb832b15d8..857cd8cbe6 100644 --- a/ui/src/app/event/event-row-debug-rulenode.tpl.html +++ b/ui/src/app/event/event-row-debug-rulenode.tpl.html @@ -21,7 +21,7 @@
{{event.body.entityName}}
{{event.body.msgId}}
{{event.body.msgType}}
-
{{event.body.dataType}}
+
{{event.body.relationType}}