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 3be7524c85..bd57f7da38 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 @@ -155,12 +155,7 @@ public class DefaultTbContext implements TbContext { @Override public void tellNext(TbMsg msg, Set relationTypes) { RuleNode ruleNode = nodeCtx.getSelf(); - DebugStrategy debugStrategy = ruleNode.getDebugStrategy(); - 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); - } + persistDebugOutput(msg, relationTypes); msg.getCallback().onProcessingEnd(ruleNode.getId()); nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(ruleNode.getRuleChainId(), ruleNode.getId(), relationTypes, msg, null)); } @@ -183,13 +178,7 @@ public class DefaultTbContext implements TbContext { if (item == null) { ack(msg); } else { - RuleNode ruleNode = nodeCtx.getSelf(); - 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); - } + persistDebugOutput(msg, relationType); 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(); mainCtx.getClusterService().pushMsgToRuleEngine(tpi, tbMsg.getId(), msg, new SimpleTbQueueCallback( metadata -> { - RuleNode ruleNode = nodeCtx.getSelf(); - DebugStrategy debugStrategy = ruleNode.getDebugStrategy(); - if (debugStrategy.shouldPersistDebugOutputForAllEvents(ruleNode.getLastUpdateTs(), tbMsg.getTs(), getMaxRuleNodeDebugDurationMinutes())) { - mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), tbMsg, TbNodeConnectionType.TO_ROOT_RULE_CHAIN); - } + persistDebugOutput(tbMsg, TbNodeConnectionType.TO_ROOT_RULE_CHAIN); if (onSuccess != null) { onSuccess.run(); } @@ -320,13 +305,7 @@ public class DefaultTbContext implements TbContext { } mainCtx.getClusterService().pushMsgToRuleEngine(tpi, tbMsg.getId(), msg.build(), new SimpleTbQueueCallback( metadata -> { - DebugStrategy debugStrategy = ruleNode.getDebugStrategy(); - 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); - } + persistDebugOutput(tbMsg, relationTypes, null, failureMessage); if (onSuccess != null) { onSuccess.run(); } @@ -343,10 +322,7 @@ public class DefaultTbContext implements TbContext { @Override public void ack(TbMsg tbMsg) { RuleNode ruleNode = nodeCtx.getSelf(); - DebugStrategy debugStrategy = ruleNode.getDebugStrategy(); - if (debugStrategy.shouldPersistDebugOutputForAllEvents(ruleNode.getLastUpdateTs(), tbMsg.getTs(), getMaxRuleNodeDebugDurationMinutes())) { - mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), tbMsg, TbNodeConnectionType.ACK); - } + persistDebugOutput(tbMsg, TbNodeConnectionType.ACK); tbMsg.getCallback().onProcessingEnd(ruleNode.getId()); tbMsg.getCallback().onSuccess(); } @@ -363,9 +339,7 @@ public class DefaultTbContext implements TbContext { @Override public void tellFailure(TbMsg msg, Throwable th) { RuleNode ruleNode = nodeCtx.getSelf(); - if (ruleNode.getDebugStrategy().shouldPersistDebugForFailureEventOnly(ruleNode.getLastUpdateTs(), msg.getTs(), getMaxRuleNodeDebugDurationMinutes())) { - mainCtx.persistDebugOutput(nodeCtx.getTenantId(), ruleNode.getId(), msg, TbNodeConnectionType.FAILURE, th); - } + persistDebugOutput(msg, Set.of(TbNodeConnectionType.FAILURE), th, null); String failureMessage = getFailureMessage(th); nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(ruleNode.getRuleChainId(), ruleNode.getId(), Collections.singleton(TbNodeConnectionType.FAILURE), @@ -1018,6 +992,24 @@ public class DefaultTbContext implements TbContext { return failureMessage; } + private void persistDebugOutput(TbMsg msg, String relationType) { + persistDebugOutput(msg, Set.of(relationType)); + } + + private void persistDebugOutput(TbMsg msg, Set relationTypes) { + persistDebugOutput(msg, relationTypes, null, null); + } + + private void persistDebugOutput(TbMsg msg, Set 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() { if (!DebugStrategy.ALL_EVENTS.equals(nodeCtx.getSelf().getDebugStrategy())) { return 0; diff --git a/application/src/test/java/org/thingsboard/server/actors/rule/DefaultTbContextTest.java b/application/src/test/java/org/thingsboard/server/actors/rule/DefaultTbContextTest.java index e9b344866f..f9ddf7d3fd 100644 --- a/application/src/test/java/org/thingsboard/server/actors/rule/DefaultTbContextTest.java +++ b/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.isNull; import static org.mockito.ArgumentMatchers.notNull; +import static org.mockito.ArgumentMatchers.nullable; import static org.mockito.BDDMockito.given; import static org.mockito.BDDMockito.then; import static org.mockito.Mockito.mock; @@ -161,7 +162,7 @@ class DefaultTbContextTest { // THEN then(nodeCtxMock).should().getChainActor(); 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(); checkTellNextCommonLogic(callbackMock, connections, msg); } @@ -226,7 +227,7 @@ class DefaultTbContextTest { then(nodeCtxMock).shouldHaveNoMoreInteractions(); then(mainCtxMock).should().getTenantProfileCache(); 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(); checkTellNextCommonLogic(callbackMock, connection, msg); } @@ -260,7 +261,7 @@ class DefaultTbContextTest { then(mainCtxMock).should().getMaxRuleNodeDebugModeDurationMinutes(); var nodeConnectionsCaptor = ArgumentCaptor.forClass(String.class); 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(); assertThat(nodeConnectionsCaptor.getAllValues()).hasSize(wantedNumberOfInvocations); assertThat(nodeConnectionsCaptor.getAllValues()).containsExactlyInAnyOrderElementsOf(connections); @@ -288,7 +289,7 @@ class DefaultTbContextTest { // THEN 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(nodeCtxMock).shouldHaveNoMoreInteractions(); } @@ -355,7 +356,7 @@ class DefaultTbContextTest { checkOutputCommonLogic(msgMock, nodeConnection); then(mainCtxMock).should().getTenantProfileCache(); 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(nodeCtxMock).shouldHaveNoMoreInteractions(); } @@ -405,7 +406,7 @@ class DefaultTbContextTest { then(callbackMock).should().onProcessingEnd(RULE_NODE_ID); then(callbackMock).should().onSuccess(); 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 @@ -644,7 +645,7 @@ class DefaultTbContextTest { if (DebugStrategy.ALL_EVENTS.equals(debugStrategy)) { then(mainCtxMock).should().getTenantProfileCache(); 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).shouldHaveNoMoreInteractions(); @@ -676,7 +677,7 @@ class DefaultTbContextTest { then(chainActorMock).should().tell(expectedRuleNodeToRuleChainTellNextMsg); then(chainActorMock).shouldHaveNoMoreInteractions(); 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(nodeCtxMock).shouldHaveNoMoreInteractions(); } @@ -736,7 +737,7 @@ class DefaultTbContextTest { then(chainActorMock).should().tell(expectedRuleNodeToRuleChainTellNextMsg); then(chainActorMock).shouldHaveNoMoreInteractions(); 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().getMaxRuleNodeDebugModeDurationMinutes(); then(mainCtxMock).shouldHaveNoMoreInteractions(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/rule/DebugStrategy.java b/common/data/src/main/java/org/thingsboard/server/common/data/rule/DebugStrategy.java index 333ee9c670..e23b1bea94 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/rule/DebugStrategy.java +++ b/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); } - public boolean shouldPersistDebugForFailureEventOnly(long lastUpdateTs, long msgTs, int debugModeDurationMinutes) { - return shouldPersistDebugOutputForAllEvents(lastUpdateTs, msgTs, debugModeDurationMinutes) || isFailureOnlyStrategy(); - } - public boolean shouldPersistDebugForFailureEventOnly(Set nodeConnections) { return isFailureOnlyStrategy() && nodeConnections.contains(TbNodeConnectionType.FAILURE); }