Browse Source

Merge remote-tracking branch 'origin/master' into sparkplug_3_0

pull/12502/head
nick 2 years ago
parent
commit
a787d8b47f
  1. 7
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  2. 134
      application/src/test/java/org/thingsboard/server/actors/rule/DefaultTbContextTest.java
  3. 16
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java

7
application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java

@ -20,6 +20,7 @@ import com.fasterxml.jackson.databind.node.ObjectNode;
import io.netty.channel.EventLoopGroup; import io.netty.channel.EventLoopGroup;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.bouncycastle.util.Arrays; import org.bouncycastle.util.Arrays;
import org.thingsboard.common.util.DebugModeUtil;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ListeningExecutor; import org.thingsboard.common.util.ListeningExecutor;
import org.thingsboard.rule.engine.api.MailService; import org.thingsboard.rule.engine.api.MailService;
@ -64,7 +65,6 @@ import org.thingsboard.server.common.data.msg.TbNodeConnectionType;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.common.data.rule.RuleNode;
import org.thingsboard.common.util.DebugModeUtil;
import org.thingsboard.server.common.data.rule.RuleNodeState; import org.thingsboard.server.common.data.rule.RuleNodeState;
import org.thingsboard.server.common.data.script.ScriptLanguage; import org.thingsboard.server.common.data.script.ScriptLanguage;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorMsg;
@ -178,7 +178,7 @@ public class DefaultTbContext implements TbContext {
.resetRuleNodeId() .resetRuleNodeId()
.build(); .build();
tbMsg.pushToStack(nodeCtx.getSelf().getRuleChainId(), nodeCtx.getSelf().getId()); tbMsg.pushToStack(nodeCtx.getSelf().getRuleChainId(), nodeCtx.getSelf().getId());
TopicPartitionInfo tpi = mainCtx.resolve(ServiceType.TB_RULE_ENGINE, getQueueName(), getTenantId(), tbMsg.getOriginator()); TopicPartitionInfo tpi = resolvePartition(msg);
doEnqueue(tpi, tbMsg, new SimpleTbQueueCallback(md -> ack(msg), t -> tellFailure(msg, t))); doEnqueue(tpi, tbMsg, new SimpleTbQueueCallback(md -> ack(msg), t -> tellFailure(msg, t)));
} }
@ -195,8 +195,7 @@ public class DefaultTbContext implements TbContext {
@Override @Override
public void enqueue(TbMsg tbMsg, Runnable onSuccess, Consumer<Throwable> onFailure) { public void enqueue(TbMsg tbMsg, Runnable onSuccess, Consumer<Throwable> onFailure) {
TopicPartitionInfo tpi = mainCtx.resolve(ServiceType.TB_RULE_ENGINE, getQueueName(), getTenantId(), tbMsg.getOriginator()); enqueue(tbMsg, tbMsg.getQueueName(), onSuccess, onFailure);
enqueue(tpi, tbMsg, onFailure, onSuccess);
} }
@Override @Override

134
application/src/test/java/org/thingsboard/server/actors/rule/DefaultTbContextTest.java

@ -98,6 +98,132 @@ class DefaultTbContextTest {
defaultTbContext = new DefaultTbContext(mainCtxMock, "Test rule chain name", nodeCtxMock); defaultTbContext = new DefaultTbContext(mainCtxMock, "Test rule chain name", nodeCtxMock);
} }
@MethodSource
@ParameterizedTest
public void givenMsgWithQueueName_whenInput_thenVerifyEnqueueWithCorrectTpi(String queueName) {
// GIVEN
var tpi = resolve(queueName);
given(mainCtxMock.resolve(eq(ServiceType.TB_RULE_ENGINE), eq(queueName), eq(TENANT_ID), eq(TENANT_ID))).willReturn(tpi);
var clusterService = mock(TbClusterService.class);
given(mainCtxMock.getClusterService()).willReturn(clusterService);
var callbackMock = mock(TbMsgCallback.class);
given(callbackMock.isMsgValid()).willReturn(true);
var ruleNode = new RuleNode(RULE_NODE_ID);
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(TENANT_ID)
.queueName(queueName)
.copyMetaData(TbMsgMetaData.EMPTY)
.data(TbMsg.EMPTY_STRING)
.callback(callbackMock)
.build();
var ruleChainId = new RuleChainId(UUID.randomUUID());
ruleNode.setRuleChainId(RULE_CHAIN_ID);
ruleNode.setDebugSettings(DebugSettings.failures());
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID);
given(nodeCtxMock.getSelf()).willReturn(ruleNode);
// WHEN
defaultTbContext.input(msg, ruleChainId);
// THEN
then(clusterService).should().pushMsgToRuleEngine(eq(tpi), eq(msg.getId()), any(), any());
}
@MethodSource
@ParameterizedTest
public void givenMsgWithQueueName_whenEnqueue_thenVerifyEnqueueWithCorrectTpi(String queueName) {
// GIVEN
var tpi = resolve(queueName);
given(mainCtxMock.resolve(eq(ServiceType.TB_RULE_ENGINE), eq(queueName), eq(TENANT_ID), eq(TENANT_ID))).willReturn(tpi);
var clusterService = mock(TbClusterService.class);
given(mainCtxMock.getClusterService()).willReturn(clusterService);
var callbackMock = mock(TbMsgCallback.class);
given(callbackMock.isMsgValid()).willReturn(true);
var ruleNode = new RuleNode(RULE_NODE_ID);
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(TENANT_ID)
.queueName(queueName)
.copyMetaData(TbMsgMetaData.EMPTY)
.data(TbMsg.EMPTY_STRING)
.callback(callbackMock)
.build();
ruleNode.setRuleChainId(RULE_CHAIN_ID);
ruleNode.setDebugSettings(DebugSettings.failures());
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID);
// WHEN
defaultTbContext.enqueue(msg, () -> {}, t -> {});
// THEN
then(clusterService).should().pushMsgToRuleEngine(eq(tpi), eq(msg.getId()), any(), any());
}
@MethodSource
@ParameterizedTest
public void givenMsgAndQueueName_whenEnqueue_thenVerifyEnqueueWithCorrectTpi(String queueName) {
// GIVEN
var tpi = resolve(queueName);
given(mainCtxMock.resolve(eq(ServiceType.TB_RULE_ENGINE), eq(queueName), eq(TENANT_ID), eq(TENANT_ID))).willReturn(tpi);
var clusterService = mock(TbClusterService.class);
given(mainCtxMock.getClusterService()).willReturn(clusterService);
var callbackMock = mock(TbMsgCallback.class);
given(callbackMock.isMsgValid()).willReturn(true);
var ruleNode = new RuleNode(RULE_NODE_ID);
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(TENANT_ID)
.copyMetaData(TbMsgMetaData.EMPTY)
.data(TbMsg.EMPTY_STRING)
.callback(callbackMock)
.build();
ruleNode.setRuleChainId(RULE_CHAIN_ID);
ruleNode.setDebugSettings(DebugSettings.failures());
given(nodeCtxMock.getTenantId()).willReturn(TENANT_ID);
// WHEN
defaultTbContext.enqueue(msg, queueName, () -> {}, t -> {});
// THEN
then(clusterService).should().pushMsgToRuleEngine(eq(tpi), eq(msg.getId()), any(), any());
}
private static Stream<String> givenMsgWithQueueName_whenInput_thenVerifyEnqueueWithCorrectTpi() {
return testQueueNames();
}
private static Stream<String> givenMsgWithQueueName_whenEnqueue_thenVerifyEnqueueWithCorrectTpi() {
return testQueueNames();
}
private static Stream<String> givenMsgAndQueueName_whenEnqueue_thenVerifyEnqueueWithCorrectTpi() {
return testQueueNames();
}
private static Stream<String> testQueueNames() {
return Stream.of("Main", "Test", null);
}
private TopicPartitionInfo resolve(String queueName) {
var tpiBuilder = TopicPartitionInfo.builder()
.topic(queueName == null ? "MainQueueTopic" : queueName + "QueueTopic")
.partition(1)
.myPartition(true);
return tpiBuilder.build();
}
@Test @Test
public void givenDebugFailuresEvents_whenTellSuccess_thenVerifyDebugOutputNotPersisted() { public void givenDebugFailuresEvents_whenTellSuccess_thenVerifyDebugOutputNotPersisted() {
// GIVEN // GIVEN
@ -810,10 +936,10 @@ class DefaultTbContextTest {
@MethodSource @MethodSource
@ParameterizedTest @ParameterizedTest
void givenDebugFailuresAndDebugAllAndConnectionAndPersistedResultOptions_whenTellNext_thenVerifyDebugOutputPersistence(boolean debugFailures, void givenDebugFailuresAndDebugAllAndConnectionAndPersistedResultOptions_whenTellNext_thenVerifyDebugOutputPersistence(boolean debugFailures,
long debugAllUntil, long debugAllUntil,
String connection, String connection,
boolean shouldPersist, boolean shouldPersist,
boolean shouldPersistAfterDurationTime) { boolean shouldPersistAfterDurationTime) {
// GIVEN // GIVEN
var callbackMock = mock(TbMsgCallback.class); var callbackMock = mock(TbMsgCallback.class);
var msg = getTbMsgWithCallback(callbackMock); var msg = getTbMsgWithCallback(callbackMock);

16
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java

@ -143,12 +143,19 @@ public interface TbContext {
void tellFailure(TbMsg msg, Throwable th); void tellFailure(TbMsg msg, Throwable th);
/** /**
* Puts new message to queue for processing by the Root Rule Chain * Puts new message to queue from TbMsg for processing by the Root Rule Chain
* *
* @param msg - message * @param msg - message
*/ */
void enqueue(TbMsg msg, Runnable onSuccess, Consumer<Throwable> onFailure); void enqueue(TbMsg msg, Runnable onSuccess, Consumer<Throwable> onFailure);
/**
* Puts new message to custom queue for processing
*
* @param msg - message
*/
void enqueue(TbMsg msg, String queueName, Runnable onSuccess, Consumer<Throwable> onFailure);
/** /**
* Sends message to the nested rule chain. * Sends message to the nested rule chain.
* Fails processing of the message if the nested rule chain is not found. * Fails processing of the message if the nested rule chain is not found.
@ -167,13 +174,6 @@ public interface TbContext {
*/ */
void output(TbMsg msg, String relationType); void output(TbMsg msg, String relationType);
/**
* Puts new message to custom queue for processing
*
* @param msg - message
*/
void enqueue(TbMsg msg, String queueName, Runnable onSuccess, Consumer<Throwable> onFailure);
void enqueueForTellFailure(TbMsg msg, String failureMessage); void enqueueForTellFailure(TbMsg msg, String failureMessage);
void enqueueForTellFailure(TbMsg tbMsg, Throwable t); void enqueueForTellFailure(TbMsg tbMsg, Throwable t);

Loading…
Cancel
Save