Browse Source

minor refactoring

pull/11861/head
YevhenBondarenko 2 years ago
parent
commit
cc977f42b2
  1. 56
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  2. 19
      application/src/test/java/org/thingsboard/server/actors/rule/DefaultTbContextTest.java
  3. 4
      common/data/src/main/java/org/thingsboard/server/common/data/rule/DebugStrategy.java

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

@ -155,12 +155,7 @@ public class DefaultTbContext implements TbContext {
@Override @Override
public void tellNext(TbMsg msg, Set<String> relationTypes) { public void tellNext(TbMsg msg, Set<String> relationTypes) {
RuleNode ruleNode = nodeCtx.getSelf(); RuleNode ruleNode = nodeCtx.getSelf();
DebugStrategy debugStrategy = ruleNode.getDebugStrategy(); persistDebugOutput(msg, relationTypes);
if (debugStrategy.shouldPersistDebugOutputForAllEvents(ruleNode.getLastUpdateTs(), msg.getTs(), getMaxRuleNodeDebugDurationMinutes())) {
relationTypes.forEach(relationType -> mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), msg, relationType));
} else if (debugStrategy.shouldPersistDebugForFailureEventOnly(relationTypes)) {
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), msg, TbNodeConnectionType.FAILURE);
}
msg.getCallback().onProcessingEnd(ruleNode.getId()); msg.getCallback().onProcessingEnd(ruleNode.getId());
nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(ruleNode.getRuleChainId(), ruleNode.getId(), relationTypes, msg, null)); nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(ruleNode.getRuleChainId(), ruleNode.getId(), relationTypes, msg, null));
} }
@ -183,13 +178,7 @@ public class DefaultTbContext implements TbContext {
if (item == null) { if (item == null) {
ack(msg); ack(msg);
} else { } else {
RuleNode ruleNode = nodeCtx.getSelf(); persistDebugOutput(msg, relationType);
DebugStrategy debugStrategy = ruleNode.getDebugStrategy();
if (debugStrategy.shouldPersistDebugOutputForAllEvents(ruleNode.getLastUpdateTs(), msg.getTs(), getMaxRuleNodeDebugDurationMinutes())) {
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), msg, relationType);
} else if (debugStrategy.shouldPersistDebugForFailureEventOnly(relationType)) {
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), msg, relationType);
}
nodeCtx.getChainActor().tell(new RuleChainOutputMsg(item.getRuleChainId(), item.getRuleNodeId(), relationType, msg)); nodeCtx.getChainActor().tell(new RuleChainOutputMsg(item.getRuleChainId(), item.getRuleNodeId(), relationType, msg));
} }
} }
@ -220,11 +209,7 @@ public class DefaultTbContext implements TbContext {
.setTbMsg(TbMsg.toByteString(tbMsg)).build(); .setTbMsg(TbMsg.toByteString(tbMsg)).build();
mainCtx.getClusterService().pushMsgToRuleEngine(tpi, tbMsg.getId(), msg, new SimpleTbQueueCallback( mainCtx.getClusterService().pushMsgToRuleEngine(tpi, tbMsg.getId(), msg, new SimpleTbQueueCallback(
metadata -> { metadata -> {
RuleNode ruleNode = nodeCtx.getSelf(); persistDebugOutput(tbMsg, TbNodeConnectionType.TO_ROOT_RULE_CHAIN);
DebugStrategy debugStrategy = ruleNode.getDebugStrategy();
if (debugStrategy.shouldPersistDebugOutputForAllEvents(ruleNode.getLastUpdateTs(), tbMsg.getTs(), getMaxRuleNodeDebugDurationMinutes())) {
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), tbMsg, TbNodeConnectionType.TO_ROOT_RULE_CHAIN);
}
if (onSuccess != null) { if (onSuccess != null) {
onSuccess.run(); onSuccess.run();
} }
@ -320,13 +305,7 @@ public class DefaultTbContext implements TbContext {
} }
mainCtx.getClusterService().pushMsgToRuleEngine(tpi, tbMsg.getId(), msg.build(), new SimpleTbQueueCallback( mainCtx.getClusterService().pushMsgToRuleEngine(tpi, tbMsg.getId(), msg.build(), new SimpleTbQueueCallback(
metadata -> { metadata -> {
DebugStrategy debugStrategy = ruleNode.getDebugStrategy(); persistDebugOutput(tbMsg, relationTypes, null, failureMessage);
if (debugStrategy.shouldPersistDebugOutputForAllEvents(ruleNode.getLastUpdateTs(), tbMsg.getTs(), getMaxRuleNodeDebugDurationMinutes())) {
relationTypes.forEach(relationType ->
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), tbMsg, relationType, null, failureMessage));
} else if (debugStrategy.shouldPersistDebugForFailureEventOnly(relationTypes)) {
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), tbMsg, TbNodeConnectionType.FAILURE, null, failureMessage);
}
if (onSuccess != null) { if (onSuccess != null) {
onSuccess.run(); onSuccess.run();
} }
@ -343,10 +322,7 @@ public class DefaultTbContext implements TbContext {
@Override @Override
public void ack(TbMsg tbMsg) { public void ack(TbMsg tbMsg) {
RuleNode ruleNode = nodeCtx.getSelf(); RuleNode ruleNode = nodeCtx.getSelf();
DebugStrategy debugStrategy = ruleNode.getDebugStrategy(); persistDebugOutput(tbMsg, TbNodeConnectionType.ACK);
if (debugStrategy.shouldPersistDebugOutputForAllEvents(ruleNode.getLastUpdateTs(), tbMsg.getTs(), getMaxRuleNodeDebugDurationMinutes())) {
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), tbMsg, TbNodeConnectionType.ACK);
}
tbMsg.getCallback().onProcessingEnd(ruleNode.getId()); tbMsg.getCallback().onProcessingEnd(ruleNode.getId());
tbMsg.getCallback().onSuccess(); tbMsg.getCallback().onSuccess();
} }
@ -363,9 +339,7 @@ public class DefaultTbContext implements TbContext {
@Override @Override
public void tellFailure(TbMsg msg, Throwable th) { public void tellFailure(TbMsg msg, Throwable th) {
RuleNode ruleNode = nodeCtx.getSelf(); RuleNode ruleNode = nodeCtx.getSelf();
if (ruleNode.getDebugStrategy().shouldPersistDebugForFailureEventOnly(ruleNode.getLastUpdateTs(), msg.getTs(), getMaxRuleNodeDebugDurationMinutes())) { persistDebugOutput(msg, Set.of(TbNodeConnectionType.FAILURE), th, null);
mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), msg, TbNodeConnectionType.FAILURE, th);
}
String failureMessage = getFailureMessage(th); String failureMessage = getFailureMessage(th);
nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(ruleNode.getRuleChainId(), nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(ruleNode.getRuleChainId(),
ruleNode.getId(), Collections.singleton(TbNodeConnectionType.FAILURE), ruleNode.getId(), Collections.singleton(TbNodeConnectionType.FAILURE),
@ -1018,6 +992,24 @@ public class DefaultTbContext implements TbContext {
return failureMessage; return failureMessage;
} }
private void persistDebugOutput(TbMsg msg, String relationType) {
persistDebugOutput(msg, Set.of(relationType));
}
private void persistDebugOutput(TbMsg msg, Set<String> relationTypes) {
persistDebugOutput(msg, relationTypes, null, null);
}
private void persistDebugOutput(TbMsg msg, Set<String> relationTypes, Throwable error, String failureMessage) {
RuleNode ruleNode = nodeCtx.getSelf();
DebugStrategy debugStrategy = ruleNode.getDebugStrategy();
if (debugStrategy.shouldPersistDebugOutputForAllEvents(ruleNode.getLastUpdateTs(), msg.getTs(), getMaxRuleNodeDebugDurationMinutes())) {
relationTypes.forEach(relationType -> mainCtx.persistDebugOutput(getTenantId(), ruleNode.getId(), msg, relationType, error, failureMessage));
} else if (debugStrategy.shouldPersistDebugForFailureEventOnly(relationTypes)) {
mainCtx.persistDebugOutput(getTenantId(), ruleNode.getId(), msg, TbNodeConnectionType.FAILURE, error, failureMessage);
}
}
private int getMaxRuleNodeDebugDurationMinutes() { private int getMaxRuleNodeDebugDurationMinutes() {
if (!DebugStrategy.ALL_EVENTS.equals(nodeCtx.getSelf().getDebugStrategy())) { if (!DebugStrategy.ALL_EVENTS.equals(nodeCtx.getSelf().getDebugStrategy())) {
return 0; return 0;

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

@ -69,6 +69,7 @@ import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.ArgumentMatchers.isNull; import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.ArgumentMatchers.notNull; import static org.mockito.ArgumentMatchers.notNull;
import static org.mockito.ArgumentMatchers.nullable;
import static org.mockito.BDDMockito.given; import static org.mockito.BDDMockito.given;
import static org.mockito.BDDMockito.then; import static org.mockito.BDDMockito.then;
import static org.mockito.Mockito.mock; import static org.mockito.Mockito.mock;
@ -161,7 +162,7 @@ class DefaultTbContextTest {
// THEN // THEN
then(nodeCtxMock).should().getChainActor(); then(nodeCtxMock).should().getChainActor();
then(nodeCtxMock).shouldHaveNoMoreInteractions(); then(nodeCtxMock).shouldHaveNoMoreInteractions();
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, TbNodeConnectionType.FAILURE); then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, TbNodeConnectionType.FAILURE, null, null);
then(mainCtxMock).shouldHaveNoMoreInteractions(); then(mainCtxMock).shouldHaveNoMoreInteractions();
checkTellNextCommonLogic(callbackMock, connections, msg); checkTellNextCommonLogic(callbackMock, connections, msg);
} }
@ -226,7 +227,7 @@ class DefaultTbContextTest {
then(nodeCtxMock).shouldHaveNoMoreInteractions(); then(nodeCtxMock).shouldHaveNoMoreInteractions();
then(mainCtxMock).should().getTenantProfileCache(); then(mainCtxMock).should().getTenantProfileCache();
then(mainCtxMock).should().getMaxRuleNodeDebugModeDurationMinutes(); then(mainCtxMock).should().getMaxRuleNodeDebugModeDurationMinutes();
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, connection); then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, connection, null, null);
then(mainCtxMock).shouldHaveNoMoreInteractions(); then(mainCtxMock).shouldHaveNoMoreInteractions();
checkTellNextCommonLogic(callbackMock, connection, msg); checkTellNextCommonLogic(callbackMock, connection, msg);
} }
@ -260,7 +261,7 @@ class DefaultTbContextTest {
then(mainCtxMock).should().getMaxRuleNodeDebugModeDurationMinutes(); then(mainCtxMock).should().getMaxRuleNodeDebugModeDurationMinutes();
var nodeConnectionsCaptor = ArgumentCaptor.forClass(String.class); var nodeConnectionsCaptor = ArgumentCaptor.forClass(String.class);
int wantedNumberOfInvocations = connections.size(); int wantedNumberOfInvocations = connections.size();
then(mainCtxMock).should(times(wantedNumberOfInvocations)).persistDebugOutput(eq(TENANT_ID), eq(RULE_NODE_ID), eq(msg), nodeConnectionsCaptor.capture()); then(mainCtxMock).should(times(wantedNumberOfInvocations)).persistDebugOutput(eq(TENANT_ID), eq(RULE_NODE_ID), eq(msg), nodeConnectionsCaptor.capture(), nullable(Throwable.class), nullable(String.class));
then(mainCtxMock).shouldHaveNoMoreInteractions(); then(mainCtxMock).shouldHaveNoMoreInteractions();
assertThat(nodeConnectionsCaptor.getAllValues()).hasSize(wantedNumberOfInvocations); assertThat(nodeConnectionsCaptor.getAllValues()).hasSize(wantedNumberOfInvocations);
assertThat(nodeConnectionsCaptor.getAllValues()).containsExactlyInAnyOrderElementsOf(connections); assertThat(nodeConnectionsCaptor.getAllValues()).containsExactlyInAnyOrderElementsOf(connections);
@ -288,7 +289,7 @@ class DefaultTbContextTest {
// THEN // THEN
checkOutputCommonLogic(msgMock, TbNodeConnectionType.FAILURE); checkOutputCommonLogic(msgMock, TbNodeConnectionType.FAILURE);
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msgMock, TbNodeConnectionType.FAILURE); then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msgMock, TbNodeConnectionType.FAILURE, null, null);
then(mainCtxMock).shouldHaveNoMoreInteractions(); then(mainCtxMock).shouldHaveNoMoreInteractions();
then(nodeCtxMock).shouldHaveNoMoreInteractions(); then(nodeCtxMock).shouldHaveNoMoreInteractions();
} }
@ -355,7 +356,7 @@ class DefaultTbContextTest {
checkOutputCommonLogic(msgMock, nodeConnection); checkOutputCommonLogic(msgMock, nodeConnection);
then(mainCtxMock).should().getTenantProfileCache(); then(mainCtxMock).should().getTenantProfileCache();
then(mainCtxMock).should().getMaxRuleNodeDebugModeDurationMinutes(); then(mainCtxMock).should().getMaxRuleNodeDebugModeDurationMinutes();
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msgMock, nodeConnection); then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msgMock, nodeConnection, null, null);
then(mainCtxMock).shouldHaveNoMoreInteractions(); then(mainCtxMock).shouldHaveNoMoreInteractions();
then(nodeCtxMock).shouldHaveNoMoreInteractions(); then(nodeCtxMock).shouldHaveNoMoreInteractions();
} }
@ -405,7 +406,7 @@ class DefaultTbContextTest {
then(callbackMock).should().onProcessingEnd(RULE_NODE_ID); then(callbackMock).should().onProcessingEnd(RULE_NODE_ID);
then(callbackMock).should().onSuccess(); then(callbackMock).should().onSuccess();
then(nodeCtxMock).should(never()).getChainActor(); then(nodeCtxMock).should(never()).getChainActor();
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msgMock, TbNodeConnectionType.ACK); then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msgMock, TbNodeConnectionType.ACK, null, null);
} }
@Test @Test
@ -644,7 +645,7 @@ class DefaultTbContextTest {
if (DebugStrategy.ALL_EVENTS.equals(debugStrategy)) { if (DebugStrategy.ALL_EVENTS.equals(debugStrategy)) {
then(mainCtxMock).should().getTenantProfileCache(); then(mainCtxMock).should().getTenantProfileCache();
then(mainCtxMock).should().getMaxRuleNodeDebugModeDurationMinutes(); then(mainCtxMock).should().getMaxRuleNodeDebugModeDurationMinutes();
then(mainCtxMock).should().persistDebugOutput(eq(TENANT_ID), eq(RULE_NODE_ID), eq(msg), eq(TbNodeConnectionType.TO_ROOT_RULE_CHAIN)); then(mainCtxMock).should().persistDebugOutput(eq(TENANT_ID), eq(RULE_NODE_ID), eq(msg), eq(TbNodeConnectionType.TO_ROOT_RULE_CHAIN), nullable(Throwable.class), nullable(String.class));
} }
then(mainCtxMock).should().getClusterService(); then(mainCtxMock).should().getClusterService();
then(mainCtxMock).shouldHaveNoMoreInteractions(); then(mainCtxMock).shouldHaveNoMoreInteractions();
@ -676,7 +677,7 @@ class DefaultTbContextTest {
then(chainActorMock).should().tell(expectedRuleNodeToRuleChainTellNextMsg); then(chainActorMock).should().tell(expectedRuleNodeToRuleChainTellNextMsg);
then(chainActorMock).shouldHaveNoMoreInteractions(); then(chainActorMock).shouldHaveNoMoreInteractions();
then(nodeCtxMock).should().getChainActor(); then(nodeCtxMock).should().getChainActor();
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, TbNodeConnectionType.FAILURE, EXCEPTION); then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, TbNodeConnectionType.FAILURE, EXCEPTION, null);
then(mainCtxMock).shouldHaveNoMoreInteractions(); then(mainCtxMock).shouldHaveNoMoreInteractions();
then(nodeCtxMock).shouldHaveNoMoreInteractions(); then(nodeCtxMock).shouldHaveNoMoreInteractions();
} }
@ -736,7 +737,7 @@ class DefaultTbContextTest {
then(chainActorMock).should().tell(expectedRuleNodeToRuleChainTellNextMsg); then(chainActorMock).should().tell(expectedRuleNodeToRuleChainTellNextMsg);
then(chainActorMock).shouldHaveNoMoreInteractions(); then(chainActorMock).shouldHaveNoMoreInteractions();
then(nodeCtxMock).should().getChainActor(); then(nodeCtxMock).should().getChainActor();
then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, TbNodeConnectionType.FAILURE, EXCEPTION); then(mainCtxMock).should().persistDebugOutput(TENANT_ID, RULE_NODE_ID, msg, TbNodeConnectionType.FAILURE, EXCEPTION, null);
then(mainCtxMock).should().getTenantProfileCache(); then(mainCtxMock).should().getTenantProfileCache();
then(mainCtxMock).should().getMaxRuleNodeDebugModeDurationMinutes(); then(mainCtxMock).should().getMaxRuleNodeDebugModeDurationMinutes();
then(mainCtxMock).shouldHaveNoMoreInteractions(); then(mainCtxMock).shouldHaveNoMoreInteractions();

4
common/data/src/main/java/org/thingsboard/server/common/data/rule/DebugStrategy.java

@ -39,10 +39,6 @@ public enum DebugStrategy {
return isAllEventsStrategyAndMsgTsWithinDebugDuration(lastUpdateTs, msgTs, debugModeDurationMinutes); return isAllEventsStrategyAndMsgTsWithinDebugDuration(lastUpdateTs, msgTs, debugModeDurationMinutes);
} }
public boolean shouldPersistDebugForFailureEventOnly(long lastUpdateTs, long msgTs, int debugModeDurationMinutes) {
return shouldPersistDebugOutputForAllEvents(lastUpdateTs, msgTs, debugModeDurationMinutes) || isFailureOnlyStrategy();
}
public boolean shouldPersistDebugForFailureEventOnly(Set<String> nodeConnections) { public boolean shouldPersistDebugForFailureEventOnly(Set<String> nodeConnections) {
return isFailureOnlyStrategy() && nodeConnections.contains(TbNodeConnectionType.FAILURE); return isFailureOnlyStrategy() && nodeConnections.contains(TbNodeConnectionType.FAILURE);
} }

Loading…
Cancel
Save