|
|
@ -110,7 +110,6 @@ import org.thingsboard.server.dao.widget.WidgetTypeService; |
|
|
import org.thingsboard.server.dao.widget.WidgetsBundleService; |
|
|
import org.thingsboard.server.dao.widget.WidgetsBundleService; |
|
|
import org.thingsboard.server.gen.transport.TransportProtos; |
|
|
import org.thingsboard.server.gen.transport.TransportProtos; |
|
|
import org.thingsboard.server.queue.TbQueueCallback; |
|
|
import org.thingsboard.server.queue.TbQueueCallback; |
|
|
import org.thingsboard.server.queue.TbQueueMsgMetadata; |
|
|
|
|
|
import org.thingsboard.server.queue.common.SimpleTbQueueCallback; |
|
|
import org.thingsboard.server.queue.common.SimpleTbQueueCallback; |
|
|
import org.thingsboard.server.service.executors.PubSubRuleNodeExecutorProvider; |
|
|
import org.thingsboard.server.service.executors.PubSubRuleNodeExecutorProvider; |
|
|
import org.thingsboard.server.service.script.RuleNodeJsScriptEngine; |
|
|
import org.thingsboard.server.service.script.RuleNodeJsScriptEngine; |
|
|
@ -177,24 +176,10 @@ class DefaultTbContext implements TbContext { |
|
|
if (!msg.isValid()) { |
|
|
if (!msg.isValid()) { |
|
|
return; |
|
|
return; |
|
|
} |
|
|
} |
|
|
TbMsg inputMsg = msg.copyWithRuleChainId(ruleChainId); |
|
|
TbMsg tbMsg = msg.copyWithRuleChainId(ruleChainId); |
|
|
inputMsg.pushToStack(nodeCtx.getSelf().getRuleChainId(), nodeCtx.getSelf().getId()); |
|
|
tbMsg.pushToStack(nodeCtx.getSelf().getRuleChainId(), nodeCtx.getSelf().getId()); |
|
|
TransportProtos.ToRuleEngineMsg toReMsg = TransportProtos.ToRuleEngineMsg.newBuilder() |
|
|
TopicPartitionInfo tpi = mainCtx.resolve(ServiceType.TB_RULE_ENGINE, getQueueName(), getTenantId(), tbMsg.getOriginator()); |
|
|
.setTenantIdMSB(getTenantId().getId().getMostSignificantBits()) |
|
|
doEnqueue(tpi, tbMsg, new SimpleTbQueueCallback(md -> ack(msg), t -> tellFailure(msg, t))); |
|
|
.setTenantIdLSB(getTenantId().getId().getLeastSignificantBits()) |
|
|
|
|
|
.setTbMsg(TbMsg.toByteString(inputMsg)).build(); |
|
|
|
|
|
TopicPartitionInfo tpi = mainCtx.resolve(ServiceType.TB_RULE_ENGINE, getQueueName(), getTenantId(), inputMsg.getOriginator()); |
|
|
|
|
|
mainCtx.getClusterService().pushMsgToRuleEngine(tpi, inputMsg.getId(), toReMsg, new TbQueueCallback() { |
|
|
|
|
|
@Override |
|
|
|
|
|
public void onSuccess(TbQueueMsgMetadata metadata) { |
|
|
|
|
|
ack(msg); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable error) { |
|
|
|
|
|
tellFailure(msg, error); |
|
|
|
|
|
} |
|
|
|
|
|
}); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
@ -230,14 +215,10 @@ class DefaultTbContext implements TbContext { |
|
|
} |
|
|
} |
|
|
return; |
|
|
return; |
|
|
} |
|
|
} |
|
|
TransportProtos.ToRuleEngineMsg msg = TransportProtos.ToRuleEngineMsg.newBuilder() |
|
|
|
|
|
.setTenantIdMSB(getTenantId().getId().getMostSignificantBits()) |
|
|
|
|
|
.setTenantIdLSB(getTenantId().getId().getLeastSignificantBits()) |
|
|
|
|
|
.setTbMsg(TbMsg.toByteString(tbMsg)).build(); |
|
|
|
|
|
if (nodeCtx.getSelf().isDebugMode()) { |
|
|
if (nodeCtx.getSelf().isDebugMode()) { |
|
|
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), nodeCtx.getSelf().getId(), tbMsg, "To Root Rule Chain"); |
|
|
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), nodeCtx.getSelf().getId(), tbMsg, "To Root Rule Chain"); |
|
|
} |
|
|
} |
|
|
mainCtx.getClusterService().pushMsgToRuleEngine(tpi, tbMsg.getId(), msg, new SimpleTbQueueCallback( |
|
|
doEnqueue(tpi, tbMsg, new SimpleTbQueueCallback( |
|
|
metadata -> { |
|
|
metadata -> { |
|
|
if (onSuccess != null) { |
|
|
if (onSuccess != null) { |
|
|
onSuccess.run(); |
|
|
onSuccess.run(); |
|
|
@ -252,6 +233,14 @@ class DefaultTbContext implements TbContext { |
|
|
})); |
|
|
})); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private void doEnqueue(TopicPartitionInfo tpi, TbMsg tbMsg, TbQueueCallback callback) { |
|
|
|
|
|
TransportProtos.ToRuleEngineMsg msg = TransportProtos.ToRuleEngineMsg.newBuilder() |
|
|
|
|
|
.setTenantIdMSB(getTenantId().getId().getMostSignificantBits()) |
|
|
|
|
|
.setTenantIdLSB(getTenantId().getId().getLeastSignificantBits()) |
|
|
|
|
|
.setTbMsg(TbMsg.toByteString(tbMsg)).build(); |
|
|
|
|
|
mainCtx.getClusterService().pushMsgToRuleEngine(tpi, tbMsg.getId(), msg, callback); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public void enqueueForTellFailure(TbMsg tbMsg, String failureMessage) { |
|
|
public void enqueueForTellFailure(TbMsg tbMsg, String failureMessage) { |
|
|
TopicPartitionInfo tpi = resolvePartition(tbMsg); |
|
|
TopicPartitionInfo tpi = resolvePartition(tbMsg); |
|
|
|